基于confluent_kafka的生产者消息无法投递,但Kafka可正常访问
解决Confluent Kafka生产者回调不触发&flush挂起的问题
嘿,我来帮你排查下这个问题,你的代码里有几个容易踩的小坑,咱们一步步来解决:
1. 移除生产者配置里的多余参数
group.id是消费者专属的配置项,生产者完全不需要这个参数——你把它加到生产者配置里,会导致客户端逻辑混乱,这大概率是消息发送异常、回调不触发的核心原因之一。
修改你的配置部分:
config = { 'metadata.broker.list': ','.join(kafka_hosts), # 删掉这一行无效配置:'group.id': 'mygroup', }
2. 给flush设置超时,避免无限挂起
默认情况下flush()会一直等待所有消息发送确认,如果消息因为各种原因发送失败,程序就会一直卡在这里。给它加个超时时间,既能终止挂起,还能帮你判断有没有消息发送失败:
# 替换原来的flush调用 flushed_count = mys.producer.flush(timeout=5) if flushed_count > 0: print(f"⚠️ 有{flushed_count}条消息未成功发送")
3. 优化消息发送逻辑,先触发元数据刷新
Confluent Kafka的生产者是异步的,produce()只是把消息加到本地队列,真正的发送和回调触发需要依赖poll()或flush()驱动。对于新创建的主题,生产者可能还没同步到最新的元数据,先调用一次poll(0)能帮它快速获取主题信息:
@staticmethod def main(): topic = 'my-topic' message = json.dumps({'measurement': [1, 2, 3]}) mys = MySource(['kafka'], topic) # 先触发一次元数据刷新,避免新主题同步延迟 mys.producer.poll(0) mys.producer.produce(topic, message, on_delivery=delivery_report) flushed_count = mys.producer.flush(timeout=5) if flushed_count > 0: print(f"⚠️ 有{flushed_count}条消息未成功发送")
4. 验证旧主题的消息是否真的发送失败
当你使用已存在的主题时,Kafka默认不会打印主题创建日志,但不代表消息没发送。你可以用Kafka的命令行消费者工具验证:
kafka-console-consumer.sh --bootstrap-server kafka:9092 --topic my-topic --from-beginning
如果修改代码后还是看不到消息,检查下你的broker地址是否正确——['kafka']这个主机名需要能被你的Python程序解析,不确定的话换成['localhost:9092']或者具体的IP+端口试试。
按照这些步骤调整后,你的回调函数应该能正常触发,flush也不会无限挂起了。
内容的提问来源于stack exchange,提问作者Yushin Washio
相关产品推荐
相关产品推荐

