【发布时间】:2020-11-07 00:32:15
【问题描述】:
我有以下 Kafka 消费者,如果将 group_id 分配给 None 效果很好 - 它收到了所有历史消息和我新测试的消息。
consumer = KafkaConsumer(
topic,
bootstrap_servers=bootstrap_servers,
auto_offset_reset=auto_offset_reset,
enable_auto_commit=enable_auto_commit,
group_id=group_id,
value_deserializer=lambda x: json.loads(x.decode('utf-8'))
)
for m in consumer:
但是,如果我将group_id 设置为某个值,它不会收到任何信息。我尝试运行测试生产者发送新消息,但没有收到任何消息。
消费者控制台确实显示以下消息:
2020-11-07 00:56:01 INFO ThreadPoolExecutor-0_0 base.py(重新)加入组 my_group 2020-11-07 00:56:07 INFO ThreadPoolExecutor-0_0 base.py 成功加入组 my_group 与第 497 代 2020-11-07 00:56:07 INFO ThreadPoolExecutor-0_0 subscription_state.py 更新了分区分配:[] 2020-11-07 00:56:07 INFO ThreadPoolExecutor-0_0 consumer.py 为组 my_group 设置新分配的分区 set()【问题讨论】:
标签: python apache-kafka kafka-consumer-api kafka-python