【发布时间】: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