Class DefaultKafkaMessageConsumer
java.lang.Object
io.eventuate.messaging.kafka.basic.consumer.DefaultKafkaMessageConsumer
- All Implemented Interfaces:
KafkaMessageConsumer
-
Method Summary
Modifier and TypeMethodDescriptionvoidassign(Collection<org.apache.kafka.common.TopicPartition> topicPartitions) voidclose()voidvoidcommitOffsets(Map<org.apache.kafka.common.TopicPartition, org.apache.kafka.clients.consumer.OffsetAndMetadata> offsets) static KafkaMessageConsumercreate(Properties properties) List<org.apache.kafka.common.PartitionInfo>partitionsFor(String topic) voidorg.apache.kafka.clients.consumer.ConsumerRecords<String,byte[]> longposition(org.apache.kafka.common.TopicPartition topicPartition) voidvoidseek(org.apache.kafka.common.TopicPartition topicPartition, long position) voidseekToEnd(Collection<org.apache.kafka.common.TopicPartition> topicPartitions) voidsubscribe(Collection<String> topics, org.apache.kafka.clients.consumer.ConsumerRebalanceListener callback) void
-
Method Details
-
create
-
assign
- Specified by:
assignin interfaceKafkaMessageConsumer
-
seekToEnd
- Specified by:
seekToEndin interfaceKafkaMessageConsumer
-
position
public long position(org.apache.kafka.common.TopicPartition topicPartition) - Specified by:
positionin interfaceKafkaMessageConsumer
-
seek
public void seek(org.apache.kafka.common.TopicPartition topicPartition, long position) - Specified by:
seekin interfaceKafkaMessageConsumer
-
subscribe
- Specified by:
subscribein interfaceKafkaMessageConsumer
-
subscribe
public void subscribe(Collection<String> topics, org.apache.kafka.clients.consumer.ConsumerRebalanceListener callback) - Specified by:
subscribein interfaceKafkaMessageConsumer
-
commitOffsets
public void commitOffsets(Map<org.apache.kafka.common.TopicPartition, org.apache.kafka.clients.consumer.OffsetAndMetadata> offsets) - Specified by:
commitOffsetsin interfaceKafkaMessageConsumer
-
partitionsFor
- Specified by:
partitionsForin interfaceKafkaMessageConsumer
-
poll
- Specified by:
pollin interfaceKafkaMessageConsumer
-
pause
- Specified by:
pausein interfaceKafkaMessageConsumer
-
resume
- Specified by:
resumein interfaceKafkaMessageConsumer
-
close
public void close()- Specified by:
closein interfaceKafkaMessageConsumer
-
close
- Specified by:
closein interfaceKafkaMessageConsumer
-