【问题标题】:Consumer.poll() returns new records even without committing offsets?Consumer.poll() 即使没有提交偏移量也会返回新记录?
【发布时间】:2017-09-16 01:32:47
【问题描述】:

如果我有一个enable.auto.commit=false 并且我调用consumer.poll() 之后没有调用consumer.commitAsync(),为什么consumer.poll() 返回 下次调用时有新记录吗?

由于我没有提交我的偏移量,我希望poll() 会返回最新的偏移量,这应该是相同的记录。

我之所以这么问,是因为我试图在处理过程中处理失败场景。我希望在不提交偏移量的情况下,poll() 会再次返回相同的记录,以便我可以再次重新处理那些失败的记录。

public class MyConsumer implements Runnable {
    @Override
    public void run() {
        while (true) {
            ConsumerRecords<String, LogLine> records = consumer.poll(Long.MAX_VALUE);
            for (ConsumerRecord record : records) {
                try {
                   //process record
                   consumer.commitAsync();
                } catch (Exception e) {
                }
                /**
                If exception happens above, I was expecting poll to return new records so I can re-process the record that caused the exception. 
                **/
            }

        }
    }
}

【问题讨论】:

    标签: apache-kafka kafka-consumer-api


    【解决方案1】:

    轮询的起始偏移量不是由代理决定的,而是由消费者决定的。消费者跟踪最后收到的偏移量,并在下一次轮询期间请求以下一堆消息。

    当消费者停止或失败并且不知道最后消费的偏移量的另一个实例开始使用分区时,偏移量提交开始发挥作用。

    KafkaConsumer 有相当丰富的 Javadoc,值得一读。

    【讨论】:

    • 有道理。但是poll() doc 说“最后消耗的偏移量可以通过 seek(TopicPartition, long) 手动设置或自动设置为订阅的分区列表的最后提交的偏移量” 后者不符合我的问题 - 哪个将消耗的偏移量设置为最后提交的偏移量 - 如果我从不提交新的偏移量,这将导致 poll() 不返回新记录。我的理解正确吗?
    • 我认为文档的部分仅指消耗的起点/偏移量。所以你要么从任何地方开始使用 seek 要么使用提交的偏移量。
    • 你的意思是它只会在第一次调用poll() 时使用提交的偏移量?
    • @Glide 是的,它仅在第一次轮询时使用存储的偏移量。你也可以在调用poll之前seek到所需的偏移量。
    【解决方案2】:

    如果重新平衡,消费者将从上次提交偏移量中读取(意味着如果任何消费者离开组或添加了新消费者),因此在 kafka 中处理重复数据删除不会直接进行,因此您必须将最后一个进程偏移量存储在外部存储,当发生重新平衡或应用重启时,您应该寻找该偏移量并开始处理,或者您应该检查消息中的某些唯一键以查找数据库是否重复

    【讨论】:

      【解决方案3】:

      我想分享一些代码,您可以如何在 Java 代码中解决这个问题。

      这种方法是轮询记录,尝试处理它们,如果发生异常,则寻求主题分区的最小值。之后,您执行commitAsync()

      public class MyConsumer implements Runnable {
          @Override
          public void run() {
              while (true) {
                  List<ConsumerRecord<String, LogLine>> records = StreamSupport
                      .stream( consumer.poll(Long.MAX_VALUE).spliterator(), true )
                      .collect( Collectors.toList() );
      
                  boolean exceptionRaised = false;
                  for (ConsumerRecord<String, LogLine> record : records) {
                      try {
                          // process record
                      } catch (Exception e) {
                          exceptionRaised = true;
                          break;
                      }
                  }
      
                  if( exceptionRaised ) {
                      Map<TopicPartition, Long> offsetMinimumForTopicAndPartition = records
                          .stream()
                          .collect( Collectors.toMap( r -> new TopicPartition( r.topic(), r.partition() ),
                              ConsumerRecord::offset,
                              Math::min
                          ) );
      
                      for( Map.Entry<TopicPartition, Long> entry : offsetMinimumForTopicAndPartition.entrySet() ) {
                          consumer.seek( entry.getKey(), entry.getValue() );
                      }
                  }
      
                  consumer.commitAsync();
              }
          }
      }
      

      使用此设置,您可以一次又一次地轮询消息,直到成功处理一次轮询的所有消息。

      请注意,您的代码应该能够处理a poison pill。否则,您的代码将陷入死循环。

      【讨论】:

        猜你喜欢
        • 1970-01-01
        • 2022-10-17
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 2019-09-20
        • 2017-06-24
        • 1970-01-01
        相关资源
        最近更新 更多