Class OffsetTracker

java.lang.Object
io.eventuate.messaging.kafka.basic.consumer.OffsetTracker

public class OffsetTracker extends Object
Keeps track of message offsets that are (a) being processed and (b) have been processed and can be committed
  • Constructor Details

    • OffsetTracker

      public OffsetTracker()
  • Method Details

    • toString

      public String toString()
      Overrides:
      toString in class Object
    • offsetsToCommit

      public Map<org.apache.kafka.common.TopicPartition,org.apache.kafka.clients.consumer.OffsetAndMetadata> offsetsToCommit()
    • noteOffsetsCommitted

      public void noteOffsetsCommitted(Map<org.apache.kafka.common.TopicPartition,org.apache.kafka.clients.consumer.OffsetAndMetadata> offsetsToCommit)
    • getPending

      public Map<org.apache.kafka.common.TopicPartition,Set<Long>> getPending()