Class EventBusReader
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.eventbus.EventBusReader
Reads the records published to an
EventBus on the given topics (all topics if none are given), for example
by an EventBusWriter.-
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
ConstructorsConstructorDescriptionEventBusReader(EventBus eventBus, int queueCapacity, Object... topics) Creates a reader that queues at mostqueueCapacityrecords; a full queue blocks the bus thread delivering to it.EventBusReader(EventBus eventBus, Object... topics) Creates a reader with an unbounded queue for the records published on any of the given topics. -
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 the extra filter events must pass besides the topic check ornullif there is none.Returns the field name used to add the delivery topic to each record ornullif the topic should not be added to records (defaults tonull).Object[]Returns the topics to read records for; empty ornullmeans all topics.booleanIndicates if this reader should stop reading from the event bus when any event bus writer on the same topic finishes (default to true).voidopen()Makes this endpoint ready for reading or writing.protected RecordreadImpl()Overridden by subclasses to read the next record from thisDataReader.setEventFilter(EventFilter eventFilter) Sets an extra filter events must pass besides the topic check (nullfor none); must be called before opening.setStopOnWriterEOF(boolean stopOnReaderEOF) Indicates if this reader should stop reading from the event bus when any event bus writer on the same topic finishes (default to true).setTopicFieldName(String topicFieldName) Assigns the field name to use when adding the delivery topic to each record ornullif the topic should not be added to records (defaults tonull).Sets the topics to read records for (empty ornullfor all topics); must be called before opening.toString()Methods inherited from class com.northconcepts.datapipeline.core.DataReader
addLineage, 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, 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
-
EventBusReader
Creates a reader with an unbounded queue for the records published on any of the given topics. -
EventBusReader
Creates a reader that queues at mostqueueCapacityrecords; a full queue blocks the bus thread delivering to it.
-
-
Method Details
-
getEventBus
-
getTopics
Returns the topics to read records for; empty ornullmeans all topics. -
setTopics
Sets the topics to read records for (empty ornullfor all topics); must be called before opening. -
getTopicFieldName
Returns the field name used to add the delivery topic to each record ornullif the topic should not be added to records (defaults tonull). -
setTopicFieldName
Assigns the field name to use when adding the delivery topic to each record ornullif the topic should not be added to records (defaults tonull). -
getEventFilter
Returns the extra filter events must pass besides the topic check ornullif there is none. -
setEventFilter
Sets an extra filter events must pass besides the topic check (nullfor none); must be called before opening. -
isStopOnWriterEOF
public boolean isStopOnWriterEOF()Indicates if this reader should stop reading from the event bus when any event bus writer on the same topic finishes (default to true). -
setStopOnWriterEOF
Indicates if this reader should stop reading from the event bus when any event bus writer on the same topic finishes (default to true). -
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
-
available
Description copied from class:DataReaderReturns the number of records that can probably be read without blocking.- Overrides:
availablein classDataReader- 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
-
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
-
toString
- Overrides:
toStringin classDataEndpoint
-