Class AsyncMultiReader
java.lang.Object
com.northconcepts.datapipeline.core.DataObject
com.northconcepts.datapipeline.core.Endpoint
com.northconcepts.datapipeline.core.DataEndpoint
com.northconcepts.datapipeline.core.DataReader
com.northconcepts.datapipeline.core.AsyncMultiReader
AsyncMultiReader reads from one or more source DataReaders asynchronously using a separate thread for each one. When this reader's
DataReader.read() method is called, it will either return the next record that was put into the buffer asynchronously from a source reader
or block until data is available in the buffer. This reader makes no guarantees as to the order of records returned.-
Nested Class Summary
Nested ClassesModifier and TypeClassDescriptionclassThread that opens one source reader, copies its records into theAsyncMultiReader's buffer and closes it.Nested classes/interfaces inherited from class com.northconcepts.datapipeline.core.DataEndpoint
DataEndpoint.State -
Field Summary
FieldsModifier and TypeFieldDescriptionprotected final LinkedBlockingQueue<Record> protected final List<AsyncMultiReader.ReaderThread> Fields inherited from class com.northconcepts.datapipeline.core.DataReader
fieldLineage, recordLineageFields 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 TypeMethodDescriptionadd(DataReader... readers) Adds source readers, each read on its own thread that also opens and closes it; must be called beforeopen().addExceptionProperties(DataException exception) Adds this endpoint's current state to aDataException.voidclose()Indicates that this endpoint has finished reading or writing.Returns the first exception thrown by a source reader or null if none failed.booleanIndicates if a source reader failure stops all reading and is rethrown byDataReader.read()(default is true) instead of only being kept forgetException().protected voidonTheadFinished(AsyncMultiReader.ReaderThread finishedThread) Called by eachAsyncMultiReader.ReaderThreadas it ends; signals the end of input once all have finished, or at the first failure whenisFailOnException()is true.voidopen()Makes this endpoint ready for reading or writing.protected RecordreadImpl()Overridden by subclasses to read the next record from thisDataReader.setFailOnException(boolean failOnException) Indicates if a source reader failure stops all reading and is rethrown byDataReader.read()(default is true) instead of only being kept forgetException().Methods inherited from class com.northconcepts.datapipeline.core.DataReader
addLineage, available, getBufferSize, getNestedEndpoint, getNestedReader, getReader, getRootEndpoint, getRootReader, isExhausted, isLineageSupported, isSaveLineage, peek, pop, push, read, setSaveLineage, skipMethods inherited from class com.northconcepts.datapipeline.core.DataEndpoint
decrementRecordCount, enableJmx, getLastRecord, getRecordCount, getRecordCountAsBigInteger, getRecordCountAsString, incrementRecordCount, isRecordCountBigInteger, 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
-
buffer
-
threads
-
-
Constructor Details
-
AsyncMultiReader
-
-
Method Details
-
add
Adds source readers, each read on its own thread that also opens and closes it; must be called beforeopen(). -
isFailOnException
public boolean isFailOnException()Indicates if a source reader failure stops all reading and is rethrown byDataReader.read()(default is true) instead of only being kept forgetException(). -
setFailOnException
Indicates if a source reader failure stops all reading and is rethrown byDataReader.read()(default is true) instead of only being kept forgetException(). -
getException
Returns the first exception thrown by a source reader or null if none failed. -
open
Description copied from class:DataEndpointMakes this endpoint ready for reading or writing.- Overrides:
openin classDataEndpoint- Throws:
DataException
-
close
Description copied from class:DataEndpointIndicates that this endpoint has finished reading or writing.- Overrides:
closein classDataEndpoint- Throws:
DataException
-
readImpl
Description copied from class:DataReaderOverridden by subclasses to read the next record from thisDataReader. The default implementation ofDataReader.read()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 —
DataReader.read()is the template method that wraps exceptions, increments the record count, and tracks lineage. - Do not call
DataEndpoint.incrementRecordCount()here;DataReader.read()already does. - Throw raw exceptions;
DataReader.read()wraps them viaexception(throwable).
- Specified by:
readImplin classDataReader- Throws:
Throwable
- Return
-
onTheadFinished
Called by eachAsyncMultiReader.ReaderThreadas it ends; signals the end of input once all have finished, or at the first failure whenisFailOnException()is true. -
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 classDataReader
-