【问题标题】:Spring Kafka consumer polling to specific offsets at runtimeSpring Kafka 消费者在运行时轮询特定偏移量
【发布时间】:2020-04-27 21:33:54
【问题描述】:

在我们的 kafka 消费者中使用 spring kafka 时。根据我的业务需求,我需要再次轮询同一批次的记录,以防处理该批次失败。 根据https://kafka.apache.org/22/javadoc/org/apache/kafka/clients/consumer/KafkaConsumer.html,部分:-“偏移量和消费者位置”说 Kafka 为分区中的每条记录维护一个数字偏移量。此偏移量充当该分区内记录的唯一标识符,并且还表示消费者在分区中的位置。例如,位于位置 5 的消费者已经消费了偏移量为 0 到 4 的记录,接下来将接收偏移量为 5 的记录。实际上有两个与消费者的用户相关的位置概念:-

  • 消费者的位置给出了将要发出的下一条记录的偏移量。它将比消费者在该分区中看到的最高偏移量大一。每次消费者在调用 poll(Duration) 中收到消息时,它都会自动前进。

  • 已提交位置是已安全存储的最后一个偏移量。如果进程失败并重新启动,这是消费者将恢复到的偏移量。消费者可以定期自动提交偏移量;或者它可以选择通过调用其中一个提交 API(例如 commitSync 和 commitAsync)手动控制这个提交位置。

对于我的用例,我想控制第一个。有什么办法吗?

【问题讨论】:

    标签: kafka-consumer-api spring-kafka


    【解决方案1】:

    SeekToCurrentErrorHandler 将重新定位消费者,以便在侦听器抛出异常时重新传递失败的记录。

    实现ConsumerSeekAware 以在启动期间寻找开始。

    【讨论】:

    • 我的要求是在应用程序启动时寻求最早。在运行时,如果前一个批次在应用程序级别出现故障,则后续轮询应提供前一个批次,如果成功则应提供新批次。我检查了“docs.spring.io/spring-kafka/docs/2.4.5.RELEASE/reference/html/…”,我想知道我必须在哪个选项中为我的用例注册回调。
    • 两者都需要 - 实现 ConsumerSeekAware 以在启动期间寻求开始,并添加 SeekToCurrentErrorHander 以在失败后重播。
    • 谢谢 Gary,我可以在失败的情况下使用 acknowledgment.nack 和 acknowledgment.acknowledge 吗?
    猜你喜欢
    • 2017-08-17
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2015-03-25
    • 2017-07-22
    • 1970-01-01
    • 2015-11-30
    相关资源
    最近更新 更多