Class KafkaReader


public class KafkaReader extends IntegrationReader
Read Records from an Apache Kafka distributed messaging system.
  • Constructor Details

    • KafkaReader

      public KafkaReader(Properties properties, String topic, long pollTimeout)
      Creates a reader for the topic that waits up to pollTimeout milliseconds for each poll. Opening it puts auto.offset.reset=earliest and the key/value deserializers into the given properties.
  • Method Details

    • isKeepPolling

      public boolean isKeepPolling()
      Indicates if the reader keeps polling while no messages arrive instead of ending after an empty poll (default is true).
    • setKeepPolling

      public KafkaReader setKeepPolling(boolean keepPolling)
      Indicates if the reader keeps polling while no messages arrive instead of ending after an empty poll (default is true).
    • open

      public void open() throws DataException
      Description copied from class: DataEndpoint
      Makes this endpoint ready for reading or writing.
      Overrides:
      open in class IntegrationReader
      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
    • 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