Class ProxyReader
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.ProxyReader
- Direct Known Subclasses:
AggregateReader,AsyncReader,AsyncTaskReader,BufferedReader,DebugReader,FilteringReader,GroupByReader,IntegrationProxyReader,LimitReader,MeteredReader,RemoveDuplicatesReader,RetryingReader,SortingReader,TailProxyReader,TeeReader,TimedReader,TransformingReader
Abstract super-class for obtaining records from another
DataReader,
possibly transforming them along the way. The only method that a subclass
should implement is interceptRecord(Record).-
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
ConstructorsConstructorDescriptionProxyReader(DataReader nestedDataReader) Creates a proxy for the given (non-null) reader, which is opened and closed along with this proxy. -
Method Summary
Modifier and TypeMethodDescriptionaddExceptionProperties(DataException exception) Adds this endpoint's current state to aDataException.intReturns the number of records that can probably be read without blocking.voidclose()Indicates that this endpoint has finished reading or writing.Returns theDataReaderheld inside this one ornullif there isn't one.protected RecordinterceptRecord(Record record) Overridden by subclasses to transform, drop, or observe a record on its way through this proxy.static ProxyReadermap(DataReader reader, BiFunction<ProxyReader, Record, Record> mapper) Wraps the reader in a proxy that passes itself and each record to the mapper; a null result drops the record.static ProxyReadermap(DataReader reader, Function<Record, Record> mapper) Wraps the reader in a proxy that passes each record through the mapper; a null result drops the record.voidopen()Makes this endpoint ready for reading or writing.protected RecordreadImpl()Overridden by subclasses to read the next record from thisDataReader.protected voidsetNestedDataReader(DataReader nestedDataReader) Make sure to close old target and open the new one or call withtrue.protected voidsetNestedDataReader(DataReader nestedDataReader, boolean manageLifecycle) Replaces the nested reader; withmanageLifecyclethe old reader is closed and, if this proxy is open, the new one is opened.Methods 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, 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
-
Constructor Details
-
ProxyReader
Creates a proxy for the given (non-null) reader, which is opened and closed along with this proxy.
-
-
Method Details
-
map
Wraps the reader in a proxy that passes each record through the mapper; a null result drops the record.- Throws:
Throwable
-
map
public static ProxyReader map(DataReader reader, BiFunction<ProxyReader, Record, throws ThrowableRecord> mapper) Wraps the reader in a proxy that passes itself and each record to the mapper; a null result drops the record.- Throws:
Throwable
-
getNestedReader
Description copied from class:DataReaderReturns theDataReaderheld inside this one ornullif there isn't one.- Overrides:
getNestedReaderin classDataReader
-
setNestedDataReader
Make sure to close old target and open the new one or call withtrue. -
setNestedDataReader
protected void setNestedDataReader(DataReader nestedDataReader, boolean manageLifecycle) throws DataException Replaces the nested reader; withmanageLifecyclethe old reader is closed and, if this proxy is open, the new one is opened.- Throws:
DataException
-
available
Description copied from class:DataReaderReturns the number of records that can probably be read without blocking.- Overrides:
availablein classDataReader- Throws:
DataException
-
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
-
interceptRecord
Overridden by subclasses to transform, drop, or observe a record on its way through this proxy.Contract for subclasses (see also
docs/authoring/ProxyReader.md):- Return the record (possibly modified) to pass it downstream.
- Return
nullto drop the record; the framework reads the next one. - Throw to fail the pipeline; the framework wraps the throwable in a
DataExceptionwith.setRecord(record). - The nested
DataReaderis opened and closed for you; do not read from it here. - Do not call
DataEndpoint.incrementRecordCount()here.
- Throws:
Throwable
-
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
-
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
-