【问题标题】:Properly Seek and Consume Kafka Messages on Multipartition Topic在多分区主题上正确查找和使用 Kafka 消息
【发布时间】:2021-03-23 07:53:47
【问题描述】:

我最近发现我一直在使用的一个主题是多分区而不是单分区。我需要重新配置我的消费者类来处理多个分区,但我有点困惑。我目前正在使用一个偏移组,我们称它为test_offset_group 以用于下面的示例。正常情况下,我会一直线性解析,及时继续前行;随着消息被添加到主题中,我将解析它们并继续前进,但如果发生崩溃或需要返回并重新运行前一天的提要,我需要能够按时间戳进行搜索。 Kafka 在这个项目中是强制性的,所以我无法更改我正在使用的流数据服务的类型。

我这样配置我的消费者:

test_consumer = KafkaConsumer("test_topic", bootstrap_servers="bootstrap_string", enable_auto_commit=False, group_id="test_offset_group"

如果我需要寻找时间戳,我会提供一个时间戳,然后使用以下方法寻找:

test_consumer.poll()

tp = TopicPartition("test_topic", 0)

needed_date = datetime.timestamp(timestamp)

rec_in = test_consumer.offsets_for_times({tp: needed_date * 1000})

test_consumer.seek(tp, rec_in[tp].offset)

上述功能非常适合单个分区使用者,但是当您考虑多个分区时,这感觉非常笨重和困难。我想我可以使用获取分区数 test_consumer.partitions_for_topic('test_topic") 然后遍历它们中的每一个,但是再一次,这似乎违背了 Kafka 的原则,我觉得应该有一种更简单的方法来做到这一点。

总而言之:我想了解如何利用 offset_group 功能寻找具有多个分区的多个偏移量,并且我想确认,通过执行上述操作,我实际上忽略了所有0以外的分区?

【问题讨论】:

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


    【解决方案1】:

    你的逻辑是正确的,你只需要在分配给这个消费者实例的所有分区上执行它。

    您可以使用assignment() 检索当前分配。

    【讨论】:

      猜你喜欢
      • 2017-04-29
      • 2016-04-18
      • 1970-01-01
      • 2020-01-07
      • 1970-01-01
      • 1970-01-01
      • 2020-06-24
      • 2017-10-17
      • 1970-01-01
      相关资源
      最近更新 更多