Class BufferedReader
A proxy that organizes incoming data by collecting records of the same type (using values in a subset of fields) to release them downstream together.
-
Nested Class Summary
Nested classes/interfaces inherited from class com.northconcepts.datapipeline.core.DataEndpoint
DataEndpoint.State -
Field Summary
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
ConstructorsConstructorDescriptionBufferedReader(DataReader reader, int queueSize, FieldList bufferByFields) Creates a reader that groups records by the values of the given fields; a background thread reads up toqueueSizerecords ahead from the nested reader.BufferedReader(DataReader reader, int queueSize, String... bufferByFields) Creates a reader that groups records by the values of the given fields; a background thread reads up toqueueSizerecords ahead from the nested reader. -
Method Summary
Modifier and TypeMethodDescriptionaddExceptionProperties(DataException exception) Adds this endpoint's current state to aDataException.voidclose()Indicates that this endpoint has finished reading or writing.longlongReturns the strategy that decides when an open buffer is closed and its records released (defaults to closing 200 milliseconds after the last record was added or once it holds 100 records).booleanisDebug()Indicates if buffer activity (buffers opened and closed, records added) is logged at debug level (default is false).voidopen()Makes this endpoint ready for reading or writing.protected RecordreadImpl()Overridden by subclasses to read the next record from thisDataReader.setBufferStrategy(BufferStrategy bufferStrategy) Sets the strategy that decides when an open buffer is closed and its records released; null restores the default (200 milliseconds after the last record was added or 100 records).setDebug(boolean debug) Indicates if buffer activity (buffers opened and closed, records added) is logged at debug level (default is false).toString()Methods inherited from class com.northconcepts.datapipeline.core.ProxyReader
available, getNestedReader, interceptRecord, map, map, setNestedDataReader, setNestedDataReaderMethods inherited from class com.northconcepts.datapipeline.core.DataReader
addLineage, getBufferSize, getNestedEndpoint, 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, resetRecordCountMethods 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
-
Constructor Details
-
BufferedReader
Creates a reader that groups records by the values of the given fields; a background thread reads up toqueueSizerecords ahead from the nested reader. Records whose grouping values are all null are dropped. -
BufferedReader
Creates a reader that groups records by the values of the given fields; a background thread reads up toqueueSizerecords ahead from the nested reader. Records whose grouping values are all null are dropped.
-
-
Method Details
-
getBufferByFields
-
isDebug
public boolean isDebug()Indicates if buffer activity (buffers opened and closed, records added) is logged at debug level (default is false). -
setDebug
Indicates if buffer activity (buffers opened and closed, records added) is logged at debug level (default is false). -
getBufferStrategy
Returns the strategy that decides when an open buffer is closed and its records released (defaults to closing 200 milliseconds after the last record was added or once it holds 100 records). -
setBufferStrategy
Sets the strategy that decides when an open buffer is closed and its records released; null restores the default (200 milliseconds after the last record was added or 100 records). Must be called beforeopen(). -
getBuffersCreated
public long getBuffersCreated() -
getBuffersClosed
public long getBuffersClosed() -
open
public void open()Description copied from class:DataEndpointMakes this endpoint ready for reading or writing.- Overrides:
openin classProxyReader
-
close
Description copied from class:DataEndpointIndicates that this endpoint has finished reading or writing.- Overrides:
closein classProxyReader- 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).
- Overrides:
readImplin classProxyReader- Throws:
Throwable
- Return
-
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
-
toString
- Overrides:
toStringin classDataEndpoint
-