Class DataReader
java.lang.Object
com.northconcepts.datapipeline.core.DataObject
com.northconcepts.datapipeline.core.Endpoint
com.northconcepts.datapipeline.core.DataEndpoint
com.northconcepts.datapipeline.core.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
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().-
Nested Class Summary
Nested classes/interfaces inherited from class com.northconcepts.datapipeline.core.DataEndpoint
DataEndpoint.State -
Field Summary
FieldsFields inherited from class com.northconcepts.datapipeline.core.DataEndpoint
lastRecord, PRODUCT, PRODUCT_VERSION, VENDOR, XML_INPUT_FACTORY_KEYFields inherited from class com.northconcepts.datapipeline.core.Endpoint
BUFFER_SIZE, captureElapsedTime, DEFAULT_READ_BUFFER_SIZEFields inherited from class com.northconcepts.datapipeline.core.DataObject
id, log, name, TIMESTAMP_FORMAT -
Constructor Summary
Constructors -
Method Summary
Modifier and TypeMethodDescriptionaddExceptionProperties(DataException exception) Adds this endpoint's current state to aDataException.protected RecordaddLineage(Record record) Called byread()for each record fromreadImpl()while lineage is saved; the default copiesrecordLineageinto every field along with its original index and name.intReturns the number of records that can probably be read without blocking.final intReturns the number of records in thisDataReaders buffer.Returns theEndpointheld inside this one ornullif there isn't one.Returns theDataReaderheld inside this one ornullif there isn't one.<T extends DataReader>
TReturns the firstDataReaderof the specified type held inside this one ornullif there isn't one.Returns the deepest, nestedEndpointheld inside this one, otherwise this instance is returned if there aren't any nestedEndpoints.Returns the deepest, nestedDataReaderheld inside this one, otherwise this instance is returned if there aren't any nestedDataReaders.booleanReturnstrueif this stream has already returned anullindicating that no further reads are possible.booleanIndicates if this reader can capture record and field lineage (false unless overridden by a reader that supports it).booleanIndicates if record and field lineage is captured for each record read (default is false).final Recordpeek(int index) Looks ahead to return the next record at the specified index without reading it ornullif one does not exist.protected Recordpop()Removes and returns the next record in thisDataReaders buffer ornullif it is empty.voidAdds a record to thisDataReaders buffer.read()Reads the next record from thisDataReaderand increases the record-count by 1.protected abstract RecordreadImpl()Overridden by subclasses to read the next record from thisDataReader.setSaveLineage(boolean saveLineage) Indicates if record and field lineage is captured for each record read (default is false); enabling it throws ifisLineageSupported()is false or the product edition does not include lineage.intskip(int count) Skips over the specified number of records.Methods inherited from class com.northconcepts.datapipeline.core.DataEndpoint
close, decrementRecordCount, enableJmx, getLastRecord, getRecordCount, getRecordCountAsBigInteger, getRecordCountAsString, incrementRecordCount, isRecordCountBigInteger, open, resetRecordCount, toStringMethods inherited from class com.northconcepts.datapipeline.core.Endpoint
addElapsedtime, assertClosed, assertNotOpened, assertOpened, finalize, getClosedOn, getDescription, getElapsedTime, getElapsedTimeAsString, getOpenedOn, getOpenElapsedTime, getOpenElapsedTimeAsString, getSelfTime, getSelfTimeAsString, getState, isCaptureElapsedTime, isClosed, isOpen, setCaptureElapsedTime, setDescription
-
Field Details
-
recordLineage
-
fieldLineage
-
-
Constructor Details
-
DataReader
public DataReader()
-
-
Method Details
-
getNestedEndpoint
Description copied from class:EndpointReturns theEndpointheld inside this one ornullif there isn't one.- Overrides:
getNestedEndpointin classDataEndpoint
-
getRootEndpoint
Description copied from class:EndpointReturns the deepest, nestedEndpointheld inside this one, otherwise this instance is returned if there aren't any nestedEndpoints.- Overrides:
getRootEndpointin classDataEndpoint
-
getNestedReader
Returns theDataReaderheld inside this one ornullif there isn't one. -
getReader
Returns the firstDataReaderof the specified type held inside this one ornullif there isn't one. -
getRootReader
Returns the deepest, nestedDataReaderheld inside this one, otherwise this instance is returned if there aren't any nestedDataReaders. -
available
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 thisDataReaders buffer. -
isExhausted
public boolean isExhausted()Returnstrueif this stream has already returned anullindicating that no further reads are possible. -
push
Adds a record to thisDataReaders buffer. Records in the buffer will be returned byread()before attempting to read from the underlying implementation.- See Also:
-
pop
Removes and returns the next record in thisDataReaders buffer ornullif it is empty.- See Also:
-
peek
Looks ahead to return the next record at the specified index without reading it ornullif one does not exist. The buffer of look ahead records are added to thisDataReaderusingpush(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
nullif one does not exist
-
readImpl
Overridden by subclasses to read the next record from thisDataReader. The default implementation ofread()now insures that this method will not be called again after it returns anull.If no record is available,
nullwill be returned.Contract for subclasses (see also
docs/authoring/DataReader.md):- Return
nullexactly 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 viaexception(throwable).
- Throws:
Throwable
- Return
-
read
Reads the next record from thisDataReaderand 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,
nullwill 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
DataExceptionusingDataObject.exception(Throwable).Subclasses generally do not need to override this method, instead they should implement
readImpl().- Throws:
DataException- See Also:
-
addLineage
Called byread()for each record fromreadImpl()while lineage is saved; the default copiesrecordLineageinto every field along with its original index and name. Overrides set their source details onrecordLineagefirst and end withsuper.addLineage(record). -
skip
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
Indicates if record and field lineage is captured for each record read (default is false); enabling it throws ifisLineageSupported()is false or the product edition does not include lineage. -
addExceptionProperties
Description copied from class:EndpointAdds this endpoint's current state to aDataException. Since this method is called whenever an exception is thrown, subclasses should override it to add their specific information.- Overrides:
addExceptionPropertiesin classDataEndpoint
-