【发布时间】:2015-01-12 17:23:22
【问题描述】:
参考:Restarting a Kafka (python) consumer consumes all the messages in the queue again
我是 kafka 的新手,我也在尝试处理偏移管理。
使用 Apache-Kafka (0.8.1.1.) 的最新版本和从 pypi 安装的 kafka-python 0.9.2(最后一次上传于 2014 年 8 月 27 日),这与 github 上的当前 master 分支不同。
当使用“SimpleConsumer”进行测试时 => 崩溃和重新启动脚本会消耗来自最后一个已知偏移量的消息。
当使用“MultiProcessConsumer”进行测试时 => 崩溃并重新启动脚本会从偏移量“0”重新开始消耗
我的小脚本(MultiProcessConsumer):
from kafka import KafkaClient, MultiProcessConsumer
KFK = KafkaClient("localhost:9092")
consumer = MultiProcessConsumer(KFK, "my-group1", "my-topic", num_procs=2)
我可以通过以下方式检查偏移量:
consumer.offsets
{0: 0, 1: 0}
然后,我运行:
A = consumer.get_messages(count=1235)
consumer.offsets
{0: 1235, 1: 0}
再次崩溃并重新启动脚本后,第一次调用“consumer.offsets”返回“{0: 1235, 1: 0}”,这很好。但是运行:
A.consumer.get_messages(count=388)
consumer.offsets
{0: 388, 1: 0}
关于如何处理这个问题的任何想法?此外,是否有正确更改 MultiProcessConsumer 偏移量以从定义的位置开始?
感谢您的帮助。
编辑: 在深入了解 kafka-python 库源并在 GitHub 上检查问题后, 见:https://github.com/mumrah/kafka-python/issues/173
所以问题在于,当 master multiprocessconsumer 启动子进程时,它在主题的每个分区上将它们的偏移量初始化为“0”(因为子进程的 autocommit 设置为 false),而不是给它们正确的值。
请参阅 GitHub 上的“mahall”评论。
【问题讨论】:
标签: offset apache-kafka