Class AsyncReader
A proxy that reads data asynchronously using a separate thread.
-
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.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 TypeMethodDescriptionfinal DataExceptionaddExceptionProperties(DataException exception) Adds this endpoint's current state to aDataException.protected DataExceptionaddExceptionPropertiesImpl(DataException exception) Adds this reader's own state, but not the nested reader's (it runs on another thread), to the exception; subclasses override this instead of the finaladdExceptionProperties(DataException).intReturns the number of records that can probably be read without blocking.voidclose()Indicates that this endpoint has finished reading or writing.final DataExceptionCreates an exception with the specified message and containing this endpoint's properties (name-value pairs).final DataExceptionConverts an exception to aDataException, prefixing the original message with the specified text and adding this endpoint's properties (name-value pairs).final DataExceptionConverts an exception to aDataExceptionand adds this endpoint's properties (name-value pairs).protected voidRuns on the background thread, reading the nested reader into the buffer (pausing while it is full) until the input ends, this reader is closed or a read fails.longReturns the estimated size in bytes of the records buffered but not yet read.Returns the first exception thrown by the background thread or null if there was none.longReturns the estimated size in bytes of buffered records at which the background thread pauses reading (defaults to 512 KB).longReturns the largest valuegetBufferSizeInBytes()has reached.intReturns the priority given to the background reading thread when this reader is opened (defaults toThread.NORM_PRIORITY).voidopen()Makes this endpoint ready for reading or writing.protected Recordpop()Removes and returns the next record in thisDataReaders buffer ornullif it is empty.voidAdds a record to thisDataReaders buffer.protected RecordreadImpl()Overridden by subclasses to read the next record from thisDataReader.voidRethrows the exception thrown by the internal thread or returns silently if no exception was thrown.setMaxBufferSizeInBytes(long maxBufferSize) Sets the estimated size in bytes of buffered records at which the background thread pauses reading (defaults to 512 KB, minimumMINIMUM_BUFFER_SIZE).setPriority(int priority) Sets the priority given to the background reading thread when this reader is opened, fromThread.MIN_PRIORITYtoThread.MAX_PRIORITY(defaults toThread.NORM_PRIORITY).Methods inherited from class com.northconcepts.datapipeline.core.ProxyReader
getNestedReader, interceptRecord, map, map, setNestedDataReader, setNestedDataReaderMethods inherited from class com.northconcepts.datapipeline.core.DataReader
addLineage, getBufferSize, getNestedEndpoint, getReader, getRootEndpoint, getRootReader, isExhausted, isLineageSupported, isSaveLineage, peek, 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, setDescriptionMethods inherited from class com.northconcepts.datapipeline.core.DataObject
getId, getName, resetID
-
Field Details
-
MINIMUM_BUFFER_SIZE
public static final int MINIMUM_BUFFER_SIZE- See Also:
-
lock
-
-
Constructor Details
-
AsyncReader
-
-
Method Details
-
getMaxBufferSizeInBytes
public long getMaxBufferSizeInBytes()Returns the estimated size in bytes of buffered records at which the background thread pauses reading (defaults to 512 KB). -
setMaxBufferSizeInBytes
Sets the estimated size in bytes of buffered records at which the background thread pauses reading (defaults to 512 KB, minimumMINIMUM_BUFFER_SIZE). -
getBufferSizeInBytes
public long getBufferSizeInBytes()Returns the estimated size in bytes of the records buffered but not yet read. -
getPeakBufferSizeInBytes
public long getPeakBufferSizeInBytes()Returns the largest valuegetBufferSizeInBytes()has reached. -
getPriority
public int getPriority()Returns the priority given to the background reading thread when this reader is opened (defaults toThread.NORM_PRIORITY). -
setPriority
Sets the priority given to the background reading thread when this reader is opened, fromThread.MIN_PRIORITYtoThread.MAX_PRIORITY(defaults toThread.NORM_PRIORITY). -
getException
Returns the first exception thrown by the background thread or null if there was none. -
open
Description copied from class:DataEndpointMakes this endpoint ready for reading or writing.- Overrides:
openin classProxyReader- Throws:
DataException
-
close
Description copied from class:DataEndpointIndicates that this endpoint has finished reading or writing.- Overrides:
closein classProxyReader- Throws:
DataException
-
push
Description copied from class:DataReaderAdds a record to thisDataReaders buffer. Records in the buffer will be returned byDataReader.read()before attempting to read from the underlying implementation.- Overrides:
pushin classDataReader- See Also:
-
pop
Description copied from class:DataReaderRemoves and returns the next record in thisDataReaders buffer ornullif it is empty.- Overrides:
popin classDataReader- See Also:
-
available
Description copied from class:DataReaderReturns the number of records that can probably be read without blocking.- Overrides:
availablein classProxyReader- Throws:
DataException
-
fillCache
protected void fillCache()Runs on the background thread, reading the nested reader into the buffer (pausing while it is full) until the input ends, this reader is closed or a read fails. -
rethrowAsyncException
public void rethrowAsyncException()Rethrows the exception thrown by the internal thread or returns silently if no exception was thrown. -
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).
- Overrides:
readImplin classProxyReader- Throws:
Throwable
- Return
-
exception
Description copied from class:DataObjectConverts an exception to aDataException, prefixing the original message with the specified text and adding this endpoint's properties (name-value pairs).The returned exception can then be thrown in the normal way.
If the supplied exception is an instanceof
DataException, it will be returned, otherwise it will be nested inside aDataException. In either case the result message will contain the new prefix along with this endpoint properties.- Overrides:
exceptionin classDataObject
-
exception
Description copied from class:DataObjectConverts an exception to aDataExceptionand adds this endpoint's properties (name-value pairs).The returned exception can then be thrown in the normal way.
If the supplied exception is an instanceof
DataException, it will be returned, otherwise it will be nested inside aDataException. In either case the result message will contain this endpoint properties.- Overrides:
exceptionin classDataObject
-
exception
Description copied from class:DataObjectCreates an exception with the specified message and containing this endpoint's properties (name-value pairs).The returned exception can then be thrown in the normal way.
- Overrides:
exceptionin classDataObject
-
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 classProxyReader
-
addExceptionPropertiesImpl
Adds this reader's own state, but not the nested reader's (it runs on another thread), to the exception; subclasses override this instead of the finaladdExceptionProperties(DataException).
-