Class KafkaMessageProcessor

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

public class KafkaMessageProcessor extends Object
Processes a Kafka message and tracks the message offsets that have been successfully processed and can be committed
  • Constructor Details

  • Method Details

    • process

      public void process(org.apache.kafka.clients.consumer.ConsumerRecord<String,byte[]> record)
    • 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 OffsetTracker getPending()
    • backlog

      public int backlog()