使用Python向Kafka发送JSON消息时值为空的问题求助
解决Kafka发送JSON消息后值为空的问题
你的问题出在produce方法的参数顺序错误,把消息key和value的位置搞反了。
多数Kafka Python客户端(比如confluent-kafka)的produce方法签名是 produce(topic, value=None, key=None, ...),第二个参数是消息值,第三个才是消息key。你之前的代码把record_key放到了value的位置,真正的JSON值被当成了key,自然消费时看不到期望的value内容。
修正后的代码:
# 原错误代码 # self._producer.produce(self._output_topic, record_key, json.dumps(json).encode('utf-8')) # 调整参数顺序,将JSON值作为第二个参数传入 self._producer.produce(self._output_topic, json.dumps(json).encode('utf-8'), record_key)
额外注意事项:
- 发送消息后记得调用
flush(),确保消息从生产者缓冲区提交到Kafka集群:
self._producer.flush()
- 确认你的消费者代码是从
value字段读取内容,而不是错误地读取key字段。
内容的提问来源于stack exchange,提问作者NikNik
相关产品推荐
相关产品推荐

