Interface KafkaMessageConsumer

All Known Implementing Classes:
DefaultKafkaMessageConsumer

public interface KafkaMessageConsumer
  • Method Summary

    Modifier and Type
    Method
    Description
    void
    assign(Collection<org.apache.kafka.common.TopicPartition> topicPartitions)
     
    void
     
    void
    close(Duration duration)
     
    void
    commitOffsets(Map<org.apache.kafka.common.TopicPartition,org.apache.kafka.clients.consumer.OffsetAndMetadata> offsets)
     
    List<org.apache.kafka.common.PartitionInfo>
     
    void
    pause(Set<org.apache.kafka.common.TopicPartition> partitions)
     
    org.apache.kafka.clients.consumer.ConsumerRecords<String,byte[]>
    poll(Duration duration)
     
    long
    position(org.apache.kafka.common.TopicPartition topicPartition)
     
    void
    resume(Set<org.apache.kafka.common.TopicPartition> partitions)
     
    void
    seek(org.apache.kafka.common.TopicPartition topicPartition, long position)
     
    void
    seekToEnd(Collection<org.apache.kafka.common.TopicPartition> topicPartitions)
     
    void
    subscribe(Collection<String> topics, org.apache.kafka.clients.consumer.ConsumerRebalanceListener callback)
     
    void
     
  • Method Details

    • assign

      void assign(Collection<org.apache.kafka.common.TopicPartition> topicPartitions)
    • seekToEnd

      void seekToEnd(Collection<org.apache.kafka.common.TopicPartition> topicPartitions)
    • position

      long position(org.apache.kafka.common.TopicPartition topicPartition)
    • seek

      void seek(org.apache.kafka.common.TopicPartition topicPartition, long position)
    • subscribe

      void subscribe(List<String> topics)
    • commitOffsets

      void commitOffsets(Map<org.apache.kafka.common.TopicPartition,org.apache.kafka.clients.consumer.OffsetAndMetadata> offsets)
    • partitionsFor

      List<org.apache.kafka.common.PartitionInfo> partitionsFor(String topic)
    • poll

      org.apache.kafka.clients.consumer.ConsumerRecords<String,byte[]> poll(Duration duration)
    • pause

      void pause(Set<org.apache.kafka.common.TopicPartition> partitions)
    • resume

      void resume(Set<org.apache.kafka.common.TopicPartition> partitions)
    • close

      void close()
    • close

      void close(Duration duration)
    • subscribe

      void subscribe(Collection<String> topics, org.apache.kafka.clients.consumer.ConsumerRebalanceListener callback)