【问题标题】:Preventing message loss with Kafka High Level Consumer 0.8.x使用 Kafka High Level Consumer 0.8.x 防止消息丢失
【发布时间】:2016-06-23 02:03:09
【问题描述】:

典型的 kafka 消费者如下所示:

kafka-broker ---> kafka-consumer ----> 像 Elastic-Search 这样的下游消费者

根据Kafka High Level Consumer的文档:

“auto.commit.interval.ms”设置是多久更新一次 消耗的偏移量被写入 ZooKeeper

如果发生以下两种情况,似乎可能会丢失消息:

  1. 在从 kafka 代理检索到一些消息后立即提交偏移量。
  2. 下游消费者(例如 Elastic-Search)无法处理最近一批消息,或者消费者进程本身被终止。

如果偏移量不根据时间间隔自动提交,但它们由 API 提交,这可能是最理想的。这将确保 kafka-consumer 只有在收到来自下游-consumer 已成功使用消息的确认后才能发出提交偏移的信号。可能会有一些消息重播(如果 kafka-consumer 在提交偏移量之前死亡),但至少不会丢失消息。

请让我知道高级消费者中是否存在这样的 API。

注意:我知道 Kafka 的 0.8.x 版本中有低级消费者 API,但我不想自己管理所有东西,因为我只需要高级消费者中的一个简单 API。

参考:

  1. AutoCommitTask.run(),寻找 commitOffsetsAsync
  2. SubscriptionState.allConsumed()

【问题讨论】:

    标签: apache-kafka kafka-consumer-api


    【解决方案1】:

    High Level Consumer API 中有一个 commitOffsets() API 可以用来解决这个问题。

    还将选项“auto.commit.enable”设置为“false”,以便 kafka 消费者在任何时候都自动提交偏移量。

    【讨论】:

      猜你喜欢
      • 2014-04-12
      • 1970-01-01
      • 2019-05-14
      • 2019-12-06
      • 1970-01-01
      • 1970-01-01
      • 2012-01-23
      • 1970-01-01
      • 2018-11-28
      相关资源
      最近更新 更多