【问题标题】:Restarting a Kafka python MultiProcessConsumer consumes all the messages in the queue again重新启动 Kafka python MultiProcessConsumer 再次消耗队列中的所有消息
【发布时间】: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


    【解决方案1】:

    这取决于消费者请求 kafka 代理进行偏移的方式。很可能您正在 Java 中做与此等价的事情

    readOffset = getLastOffset(consumer,topic, partition, kafka.api.OffsetRequest.EarliestTime(), clientName);
    

    试试这样的

    readOffset = getLastOffset(consumer,topic, partition, kafka.api.OffsetRequest.LatestTime(), clientName);
    

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2014-04-03
      • 2016-04-13
      • 1970-01-01
      • 1970-01-01
      • 2023-01-07
      • 1970-01-01
      相关资源
      最近更新 更多