【问题标题】:seekToEnd of all partitions and survive automatic rebalancing of Kafka consumers所有分区的 seekToEnd 并在 Kafka 消费者的自动重新平衡中幸存下来
【发布时间】:2017-12-17 21:54:18
【问题描述】:

当消费者组 A 的 Kafka 消费者连接到 Kafka 代理时,我想寻找所有分区的末尾,即使在代理端存储了偏移量。如果更多额外的消费者正在为同一个消费者组连接,他们应该获取最新存储的偏移量。 我正在执行以下操作:

consumer.poll(timeout) 
consumer.seekToEnd(emptyList())

while(true) {
  val records = consumer.poll(timeout)
  if(records.isNotEmpty()) {
    //print records
    consumer.commitSync()
  }
}

问题是当我连接消费者组 A 的第一个消费者 c1 时,一切都按预期工作,如果我连接消费者组 A 的另一个消费者 c2,则该组正在重新平衡,c1 将消耗跳过的偏移量。

有什么想法吗?

【问题讨论】:

    标签: java apache-kafka kotlin


    【解决方案1】:

    你可以创建一个实现ConsumerRebalanceListener的类,如下所示:

    public class AlwaysSeekToEndListener<K, V> implements ConsumerRebalanceListener {
    
        private Consumer<K, V> consumer;
    
        public AlwaysSeekToEndListener(Consumer consumer) {
            this.consumer = consumer;
        }
    
        @Override
        public void onPartitionsRevoked(Collection<TopicPartition> partitions) {
    
        }
    
        @Override
        public void onPartitionsAssigned(Collection<TopicPartition> partitions) {
            consumer.seekToEnd(partitions);
        }
    }
    

    然后在订阅主题时使用这个监听器:

    consumer.subscribe(Collections.singletonList("test"), new AlwaysSeekToEndListener<String, String>(consumer));
    

    【讨论】:

    • 谢谢。那应该这样做。但是,如果分区被撤销,是否还有提交当前偏移量的选项,所以我并不总是需要从最后开始寻找?还是不行?
    • onPartitionsRevoked 就是这样做的。确保在重新平衡操作开始之前调用此方法,您可以处理其中的偏移提交。
    猜你喜欢
    • 1970-01-01
    • 2021-01-26
    • 1970-01-01
    • 2020-08-04
    • 1970-01-01
    • 2018-05-23
    • 2020-01-10
    • 2017-06-19
    • 2017-04-20
    相关资源
    最近更新 更多