Class EventBusReader


public class EventBusReader extends DataReader
Reads the records published to an EventBus on the given topics (all topics if none are given), for example by an EventBusWriter.
  • Constructor Details

    • EventBusReader

      public EventBusReader(EventBus eventBus, Object... topics)
      Creates a reader with an unbounded queue for the records published on any of the given topics.
    • EventBusReader

      public EventBusReader(EventBus eventBus, int queueCapacity, Object... topics)
      Creates a reader that queues at most queueCapacity records; a full queue blocks the bus thread delivering to it.
  • Method Details

    • getEventBus

      public EventBus getEventBus()
    • getTopics

      public Object[] getTopics()
      Returns the topics to read records for; empty or null means all topics.
    • setTopics

      public EventBusReader setTopics(Object... topics)
      Sets the topics to read records for (empty or null for all topics); must be called before opening.
    • getTopicFieldName

      public String getTopicFieldName()
      Returns the field name used to add the delivery topic to each record or null if the topic should not be added to records (defaults to null).
    • setTopicFieldName

      public EventBusReader setTopicFieldName(String topicFieldName)
      Assigns the field name to use when adding the delivery topic to each record or null if the topic should not be added to records (defaults to null).
    • getEventFilter

      public EventFilter getEventFilter()
      Returns the extra filter events must pass besides the topic check or null if there is none.
    • setEventFilter

      public EventBusReader setEventFilter(EventFilter eventFilter)
      Sets an extra filter events must pass besides the topic check (null for 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

      public EventBusReader 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).
    • open

      public void open() throws DataException
      Description copied from class: DataEndpoint
      Makes this endpoint ready for reading or writing.
      Overrides:
      open in class DataEndpoint
      Throws:
      DataException
    • close

      public void close() throws DataException
      Description copied from class: DataEndpoint
      Indicates that this endpoint has finished reading or writing.
      Overrides:
      close in class DataEndpoint
      Throws:
      DataException
    • available

      public int available() throws DataException
      Description copied from class: DataReader
      Returns the number of records that can probably be read without blocking.
      Overrides:
      available in class DataReader
      Throws:
      DataException
    • readImpl

      protected Record readImpl() throws Throwable
      Description copied from class: DataReader
      Overridden by subclasses to read the next record from this DataReader. The default implementation of DataReader.read() now insures that this method will not be called again after it returns a null.

      If no record is available, null will be returned.

      Contract for subclasses (see also docs/authoring/DataReader.md):

      Specified by:
      readImpl in class DataReader
      Throws:
      Throwable
    • addExceptionProperties

      public DataException addExceptionProperties(DataException exception)
      Description copied from class: Endpoint
      Adds this endpoint's current state to a DataException. Since this method is called whenever an exception is thrown, subclasses should override it to add their specific information.
      Overrides:
      addExceptionProperties in class DataReader
    • toString

      public String toString()
      Overrides:
      toString in class DataEndpoint