【发布时间】:2016-06-23 02:03:09
【问题描述】:
典型的 kafka 消费者如下所示:
kafka-broker ---> kafka-consumer ----> 像 Elastic-Search 这样的下游消费者
根据Kafka High Level Consumer的文档:
“auto.commit.interval.ms”设置是多久更新一次 消耗的偏移量被写入 ZooKeeper
如果发生以下两种情况,似乎可能会丢失消息:
- 在从 kafka 代理检索到一些消息后立即提交偏移量。
- 下游消费者(例如 Elastic-Search)无法处理最近一批消息,或者消费者进程本身被终止。
如果偏移量不根据时间间隔自动提交,但它们由 API 提交,这可能是最理想的。这将确保 kafka-consumer 只有在收到来自下游-consumer 已成功使用消息的确认后才能发出提交偏移的信号。可能会有一些消息重播(如果 kafka-consumer 在提交偏移量之前死亡),但至少不会丢失消息。
请让我知道高级消费者中是否存在这样的 API。
注意:我知道 Kafka 的 0.8.x 版本中有低级消费者 API,但我不想自己管理所有东西,因为我只需要高级消费者中的一个简单 API。
参考:
- AutoCommitTask.run(),寻找 commitOffsetsAsync
- SubscriptionState.allConsumed()
【问题讨论】:
标签: apache-kafka kafka-consumer-api