【问题标题】:Reduce kafka offset by one message size将 kafka 偏移量减少一个消息大小
【发布时间】:2016-02-13 09:35:17
【问题描述】:

我有一个 kafka 集群设置。大多数时候,消费者在无人看管的情况下运行。它读取消息并调用外部 API。但是如果外部 API 出现故障(这种情况很少发生),我需要在固定的时间内重试消息。如果重试失败,我需要停止消费者并在固定时间后重新开始。

问题是我需要处理来自主题的最后一条读取消息。有没有一种方法可以在阅读消息之前将 zookeeper 偏移量重置为一个?(即减少一条消息的偏移量)。所以下次我启动消费者时,我可以再次阅读消息。

这可以通过使用低级消费者来完成。但是有没有办法对高级消费者做到这一点?

我正在使用基于 Java 的消费者客户端。

【问题讨论】:

  • 如果适用,一个非常简单的解决方案是在停止时将主题中的消息推回。你甚至可以有一个专门的主题,并从两者中消费。

标签: java apache-kafka apache-zookeeper kafka-consumer-api


【解决方案1】:

您需要向 zookeeper 提交偏移量。 查看下一个消费者配置: http://kafka.apache.org/documentation.html#consumerconfigs

auto.commit.enable:如果为 true,则定期向 ZooKeeper 提交 消费者已经获取的消息的偏移量。这承诺 当进程失败时,将使用偏移量作为起始位置 新的消费者将开始。

auto.commit.interval.ms:频率 ms 消费者偏移量已提交给 zookeeper。

如果您想在需要的每条消息之后提交偏移量:

auto.commit.enable=false

并在操作成功后提交偏移量:

consumer.commitOffsets(true)

不建议在每条消息之后执行此操作或设置较小的提交间隔,因为它会增加 zookeeper 的读/写负载。

您还可以查看新的偏移管理(偏移存储在经纪人而不是 zookeper 中):

http://kafka.apache.org/documentation.html

Kafka 提供了存储给定所有偏移量的选项 指定代理(针对该组)中的消费者组称为 偏移管理器。即,该消费者组中的任何消费者实例 应该将其偏移提交和获取发送到该偏移管理器 (经纪人)

【讨论】:

  • 谢谢。但是新的消费者发布了吗? Kafka 站点中的示例仍在使用 0.8.0 。我还读到 0.9 将是下一个主要版本。虽然文档提到了 0.8.2..
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2021-02-14
  • 2013-04-15
  • 2014-09-22
  • 1970-01-01
  • 1970-01-01
  • 2019-06-08
相关资源
最近更新 更多