【问题标题】:Read only new messages in a kafka topic只读 kafka 主题中的新消息
【发布时间】:2021-06-17 19:52:14
【问题描述】:

我正在 python 中使用 confluent-kafka 创建一个消费者,我想以这样一种方式创建它,如果消费者重新启动,它会从主题中的最后一条可用消息(每个分区)开始,它不会不管它是否在没有提交的情况下留下消息。

这是为了避免处理数以百万计的消息,这些消息在消费者关闭时生成并且不再需要处理。

我尝试设置参数 auto.offset.reset 的不同选项,但最多从上次提交的偏移量开始。这是我的配置:

consumer = Consumer({"bootstrap.servers": "localhost:9092",
                     "group.id": group_id,
                     "auto.offset.reset": "latest",
                     "isolation.level": "read_committed",
                     "default.topic.config": {"enable.auto.commit": False}})

有没有办法实现这种行为?

注意:我可能有多个消费者,但没有一个手动分配给特定分区

【问题讨论】:

    标签: python-3.x apache-kafka confluent-kafka-python


    【解决方案1】:

    auto.offset.reset 配置仅在没有提交的偏移量时应用。

    如果您想始终从头开始,您可以使用enable.auto.commit=false 禁用自动提交(并确保也不要显式提交),并将auto.offset.reset 设置为latest

    另一种选择是使用get_watermark_offsets()seek() 的组合将分区分配给消费者(使用on_assigned())时显式搜索到末尾

    【讨论】:

    • 嗨,谢谢你的回答,这给我带来了几个问题:如果自动提交设置为 false 并且我没有明确提交,这是否意味着我不能再实现一个行为?另外,如果我使用 seek 和 watermak 选项,我必须明确告诉消费者要查看哪个分区,并且由于我有多个消费者,我可能会动态创建新的消费者,所以我不知道这是否是个问题
    猜你喜欢
    • 2015-12-07
    • 1970-01-01
    • 1970-01-01
    • 2021-03-01
    • 2018-01-12
    • 2015-07-17
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多