Class EventuateKafkaConsumer
java.lang.Object
io.eventuate.messaging.kafka.basic.consumer.EventuateKafkaConsumer
A Kafka consumer that manually commits offsets and supports asynchronous message processing
-
Constructor Summary
ConstructorsConstructorDescriptionEventuateKafkaConsumer(String subscriberId, EventuateKafkaConsumerMessageHandler handler, List<String> topics, String bootstrapServers, EventuateKafkaConsumerConfigurationProperties eventuateKafkaConsumerConfigurationProperties, KafkaConsumerFactory kafkaConsumerFactory) -
Method Summary
Modifier and TypeMethodDescriptiongetState()booleanvoidsetCloseConsumerOnStop(boolean closeConsumerOnStop) voidsetConsumerCallbacks(Optional<ConsumerCallbacks> consumerCallbacks) voidstart()voidstop()static List<org.apache.kafka.common.PartitionInfo>verifyTopicExistsBeforeSubscribing(KafkaMessageConsumer consumer, String topic)
-
Constructor Details
-
EventuateKafkaConsumer
public EventuateKafkaConsumer(String subscriberId, EventuateKafkaConsumerMessageHandler handler, List<String> topics, String bootstrapServers, EventuateKafkaConsumerConfigurationProperties eventuateKafkaConsumerConfigurationProperties, KafkaConsumerFactory kafkaConsumerFactory)
-
-
Method Details
-
setConsumerCallbacks
-
isCloseConsumerOnStop
public boolean isCloseConsumerOnStop() -
setCloseConsumerOnStop
public void setCloseConsumerOnStop(boolean closeConsumerOnStop) -
verifyTopicExistsBeforeSubscribing
public static List<org.apache.kafka.common.PartitionInfo> verifyTopicExistsBeforeSubscribing(KafkaMessageConsumer consumer, String topic) -
start
public void start() -
stop
public void stop() -
getState
-