【问题标题】:How do I get the most recent message from a topic using Pykafka?如何使用 Pykafka 从主题中获取最新消息?
【发布时间】:2017-11-29 15:35:18
【问题描述】:

我不断地使用 pykafka 为某个主题生成消息

producer.produce('test')

我想收到最新消息。我在 pykafka Github 页面上找到了一个解决方案,它建议:

client = KafkaClient(hosts="xxxxxxx")
topic = client.topics['mytopic']
consumer = topic.get_simple_consumer(
    auto_offset_reset=OffsetType.LATEST,
    reset_offset_on_start=True)
LAST_N_MESSAGES = 2
offsets = [(p, op.next_offset - LAST_N_MESSAGES) for p, op in consumer._partitions.iteritems()]
consumer.reset_offsets(offsets)
consumer.consume()

但是,我真的不明白这里发生了什么,如果那里至少有两条消息,它只会获取最新消息。

有更强大的解决方案吗?

【问题讨论】:

    标签: python apache-kafka kafka-consumer-api pykafka


    【解决方案1】:

    准确定义“最新消息”的含义很重要。在具有多个分区的 Kafka 主题中,如果不检查消息内容,实际上不可能知道每个分区上的哪些最新消息是全局最新消息。定义何时获取最新消息也很重要——您现在想要它们一次吗?您想从最近的消息开始消费,然后在消息添加到主题时继续消费吗?您想定期只获取最新的 N 条消息吗?

    您在上面包含的配方(我为 PyKafka 文档编写的基础)为您提供了每个分区的最后 N 条消息,供您选择 N。如果您只想获取最后一条消息,您可以简单地设置 @987654322 @ 到 1。本质上,配方检查每个分区消耗的最新偏移量,然后在此之前将消费者的偏移量重置为精确的 LAST_N_MESSAGES。当您从这一点开始消费时,您只会获得分区的最后 N 条消息。

    综上所述,如果你只是对从话题结尾开始消费感兴趣,你可以使用这个:

    consumer = topic.get_simple_consumer(
        auto_offset_reset=OffsetType.LATEST,
        reset_offset_on_start=True)
    

    然后开始正常消费。

    【讨论】:

    猜你喜欢
    • 2019-06-08
    • 2021-05-12
    • 2015-08-01
    • 1970-01-01
    • 1970-01-01
    • 2019-02-10
    • 1970-01-01
    • 2022-07-14
    • 1970-01-01
    相关资源
    最近更新 更多