Class KafkaMessageProcessor
java.lang.Object
io.eventuate.messaging.kafka.basic.consumer.KafkaMessageProcessor
Processes a Kafka message and tracks the message offsets that have been successfully processed and can be committed
-
Constructor Summary
ConstructorsConstructorDescriptionKafkaMessageProcessor(String subscriberId, EventuateKafkaConsumerMessageHandler handler) -
Method Summary
Modifier and TypeMethodDescriptionintbacklog()voidnoteOffsetsCommitted(Map<org.apache.kafka.common.TopicPartition, org.apache.kafka.clients.consumer.OffsetAndMetadata> offsetsToCommit) Map<org.apache.kafka.common.TopicPartition,org.apache.kafka.clients.consumer.OffsetAndMetadata> void
-
Constructor Details
-
KafkaMessageProcessor
-
-
Method Details
-
process
-
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
-
backlog
public int backlog()
-