Class DataReader

Direct Known Subclasses:
AbstractReader, AsyncMultiReader, CollectionReader, DataReaderIterator, Dataset.ColumnsDataReader, DatasetReader, EventBusReader, FileReader, FileSystemInfoReader, IntegrationReader, JdbcReader, JsonRecordReader, MemoryReader, NullReader, PipedReader, ProxyReader, SequenceReader, SimpleJsonReader, SimpleXmlReader, SplitReader, XmlReader, XmlRecordReader

public abstract class DataReader extends DataEndpoint
Abstract super-class for reading records. The only method that a subclass must implement is readImpl(), however, most subclasses will also override DataEndpoint.open() and DataEndpoint.close().
  • Field Details

  • Constructor Details

    • DataReader

      public DataReader()
  • Method Details

    • getNestedEndpoint

      public DataReader getNestedEndpoint()
      Description copied from class: Endpoint
      Returns the Endpoint held inside this one or null if there isn't one.
      Overrides:
      getNestedEndpoint in class DataEndpoint
    • getRootEndpoint

      public DataReader getRootEndpoint()
      Description copied from class: Endpoint
      Returns the deepest, nested Endpoint held inside this one, otherwise this instance is returned if there aren't any nested Endpoints.
      Overrides:
      getRootEndpoint in class DataEndpoint
    • getNestedReader

      public DataReader getNestedReader()
      Returns the DataReader held inside this one or null if there isn't one.
    • getReader

      public <T extends DataReader> T getReader(Class<T> type)
      Returns the first DataReader of the specified type held inside this one or null if there isn't one.
    • getRootReader

      public DataReader getRootReader()
      Returns the deepest, nested DataReader held inside this one, otherwise this instance is returned if there aren't any nested DataReaders.
    • available

      public int available() throws DataException
      Returns the number of records that can probably be read without blocking.
      Throws:
      DataException
    • getBufferSize

      public final int getBufferSize()
      Returns the number of records in this DataReaders buffer.
    • isExhausted

      public boolean isExhausted()
      Returns true if this stream has already returned a null indicating that no further reads are possible.
    • push

      public void push(Record record)
      Adds a record to this DataReaders buffer. Records in the buffer will be returned by read() before attempting to read from the underlying implementation.
      See Also:
    • pop

      protected Record pop()
      Removes and returns the next record in this DataReaders buffer or null if it is empty.
      See Also:
    • peek

      public final Record peek(int index)
      Looks ahead to return the next record at the specified index without reading it or null if one does not exist. The buffer of look ahead records are added to this DataReader using push(Record).
      Parameters:
      index - specifies how far to look ahead (0 refers to the next record, 1--to the one after that, etc.)
      Returns:
      the record at the look ahead position or null if one does not exist
    • readImpl

      protected abstract Record readImpl() throws Throwable
      Overridden by subclasses to read the next record from this DataReader. The default implementation of read() now insures that this method will not be called again after it returns a null.

      If no record is available, null will be returned.

      Contract for subclasses (see also docs/authoring/DataReader.md):

      • Return null exactly once to signal end-of-stream.
      • Do not call this method directly — read() is the template method that wraps exceptions, increments the record count, and tracks lineage.
      • Do not call DataEndpoint.incrementRecordCount() here; read() already does.
      • Throw raw exceptions; read() wraps them via exception(throwable).
      Throws:
      Throwable
    • read

      public Record read() throws DataException
      Reads the next record from this DataReader and increases the record-count by 1. This method will first read any pushed (push(Record)) records before reading from the underlying source.

      If no record is available, null will be returned. This method blocks until a record is available, the end of the stream is reached, or an exception is thrown.

      Any exception raised while reading will be converted to a DataException using DataObject.exception(Throwable).

      Subclasses generally do not need to override this method, instead they should implement readImpl().

      Throws:
      DataException
      See Also:
    • addLineage

      protected Record addLineage(Record record)
      Called by read() for each record from readImpl() while lineage is saved; the default copies recordLineage into every field along with its original index and name. Overrides set their source details on recordLineage first and end with super.addLineage(record).
    • skip

      public int skip(int count) throws DataException
      Skips over the specified number of records. Since there may be fewer records than the specified number, this method returns the actual number skipped.
      Parameters:
      count - the number of records to skip
      Returns:
      The actual number of records skipped
      Throws:
      DataException
    • isLineageSupported

      public boolean isLineageSupported()
      Indicates if this reader can capture record and field lineage (false unless overridden by a reader that supports it).
    • isSaveLineage

      public boolean isSaveLineage()
      Indicates if record and field lineage is captured for each record read (default is false).
    • setSaveLineage

      public DataReader setSaveLineage(boolean saveLineage)
      Indicates if record and field lineage is captured for each record read (default is false); enabling it throws if isLineageSupported() is false or the product edition does not include lineage.
    • addExceptionProperties

      public DataException addExceptionProperties(DataException exception)
      Description copied from class: Endpoint
      Adds this endpoint's current state to a DataException. Since this method is called whenever an exception is thrown, subclasses should override it to add their specific information.
      Overrides:
      addExceptionProperties in class DataEndpoint