重新启动Kafka(python)使用者会再次消耗队列中的所有消息

use*_*335 7 python apache-kafka

我正在使用Kafka 0.8.1和Kafka python-0.9.0.在我的设置中,我有2个kafka经纪人设置.当我运行我的kafka消费者时,我可以看到它从队列中检索消息并跟踪两个经纪人的偏移量.一切都很棒!

我的问题是,当我重新启动消费者时,它从一开始就开始消费消息.我所期待的是,在重新启动时,消费者会开始在消息丢失前从消失的地方消费消息.

我确实尝试跟踪Redis中的消息偏移,然后在从队列中读取消息之前调用consumer.seek,以确保我只获取了之前从未见过的消息.虽然这有用,但在部署这个解决方案之前,我想跟你们一起检查......也许我对Kafka或python-Kafka客户端有些误解.似乎消费者能够从它停止的地方重新开始阅读是非常基本的功能.

谢谢!

小智 5

小心kafka-python库.它有一些小问题.

如果速度对您的消费者来说不是真正的问题,您可以在每条消息中设置自动提交.它应该有效.

SimpleConsumer提供了一种seek方法(https://github.com/mumrah/kafka-python/blob/master/kafka/consumer/simple.py#L174-L185),允许您在任何需要的位置开始使用消息.

最常见的电话是:

  • consumer.seek(0, 0) 从队列的开头开始阅读.
  • consumer.seek(0, 1) 开始从当前偏移量读取.
  • consumer.seek(0, 2) 跳过所有待处理的消息并开始只读取新消息.

第一个参数是这些位置的偏移量.这样,如果你打电话,consumer.seek(5, 0)你将跳过队列中的前5条消息.

另外,不要忘记,为消费者组存储偏移量.确保你一直使用同一个.


Dmi*_*try 1

Kafka 消费者能够在 Zookeeper 中存储偏移量。在 Java API 中,我们有两个选项:高级消费者,它为我们管理状态并在重新启动后从其离开的位置开始消费;以及无状态低级消费者,没有这种超能力。

根据我对Python消费者代码(https://github.com/mumrah/kafka-python/blob/master/kafka/consumer.py)的理解,SimpleConsumerMultiProcessConsumer都是有状态的,并跟踪Zookeeper中的当前偏移量,所以很奇怪你有这个重新消耗的问题。

确保在重新启动时具有相同的消费者组 ID(可能是随机设置的?)并检查以下选项:

auto_commit: default True. Whether or not to auto commit the offsets
auto_commit_every_n: default 100. How many messages to consume
                     before a commit
auto_commit_every_t: default 5000. How much time (in milliseconds) to
                     wait before commit
Run Code Online (Sandbox Code Playgroud)

您可能消耗 < 100 条消息或 < 5000 毫秒?