【问题标题】:Kafka get the partition id for a corresponding message IDKafka 获取对应消息 ID 的分区 ID
【发布时间】:2021-10-21 00:08:05
【问题描述】:

我正在编写一个 Kafka 消费者(在 python 中),它需要从特定时间戳开始流式传输给定密钥的所有消息。

集群使用默认的分区策略,每条消息都被赋予了一个键,所以我知道我想要的所有消息都在一个分区上。

问题是,我怎样才能找到我的消息写入了哪个分区?

即我需要这样的东西……(伪代码)

import kafka
from kafka.structs import TopicPartition

target_msg_key = "3a11d08b-d635-490a-aa4e-16b282a599e6"

consumer = kafka.KafkaConsumer([...])

# This doesn't exist?
partition_id = consumer.getPartitionIDforKey(target_msg_key)

tp = TopicPartition("my-topic",partition_id)

consumer.assign([tp])

for msg in consumer:
    ....

基本上,我需要一种将消息键转换为分区 id 的方法(在 python 中)。

【问题讨论】:

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


    【解决方案1】:

    您可以在DefaultPartitioner的源代码中找到该逻辑

    # given: all_partitions as a sorted list of numbers for the topic's partitions
    idx = murmur2(key)
    idx &= 0x7fffffff
    idx %= len(all_partitions)
    return all_partitions[idx]
    

    【讨论】:

    • 谢谢,这真的很有用,但它需要我知道在写入第一条消息(用于密钥)时存在多少个分区。由于分区的数量不固定,所以当我来阅读消息时,我不能使用这种方法。
    • 好吧,我已经意识到我完全错了,这是正确的答案......简而言之:不,你不能(不应该)改变分区的数量......
    • 嗯,您可以更改数字(但只能增加)。从消费者那里,有一个方法partitions_for_topic,或者类似的东西
    • 是的,一切都很好,只是您将获得写入新消息的分区 ID,而不是原始消息所在的分区。但它仍然适合我需要的东西。谢谢!
    猜你喜欢
    • 2019-11-28
    • 2016-11-16
    • 2017-08-17
    • 1970-01-01
    • 1970-01-01
    • 2021-06-21
    • 2020-10-20
    • 2014-08-24
    • 2022-07-23
    相关资源
    最近更新 更多