首次使用kafka-python实现从JSON文件逐条读取数据的Kafka Producer报错求助
解决Kafka生产者发送JSON数据时的TypeError问题
咱们先看你遇到的错误:
Producer created..............
Traceback (most recent call last):
File "prodECJson.py", line 10, in
producer.send('ECJson', json.dump(state))
TypeError: dump() missing 1 required positional argument: 'fp'
问题出在json.dump()的用法上——这个函数是用来把Python对象写入文件的,必须传入第二个参数(文件指针fp)指定要写入的文件。但你这里需要的是把字典转成JSON字符串发给Kafka,应该用json.dumps()(注意末尾的s,代表string),它会直接返回序列化后的JSON字符串,而不是写入文件。
另外,Kafka生产者的send()方法默认只接受字节类型的数据,所以还需要把JSON字符串编码成字节(比如用encode('utf-8'))。
这里给你两种修正方案:
方案1:使用生产者的value_serializer(推荐)
可以在初始化KafkaProducer时指定value_serializer,这样每次发送数据时会自动帮你序列化并编码,代码更简洁:
import json from kafka import KafkaProducer # 配置value_serializer,自动把Python对象转成JSON字节 producer = KafkaProducer( bootstrap_servers='localhost:9092', value_serializer=lambda m: json.dumps(m).encode('utf-8') ) print('Producer created..............') with open('/home/ravi/test.json') as f: data = json.load(f) for state in data['states']: producer.send('ECJson', state) producer.flush() # 确保消息立即发送,避免批量延迟
方案2:手动处理序列化和编码
如果你不想用序列化器,也可以在发送时手动处理:
import json from kafka import KafkaProducer producer = KafkaProducer(bootstrap_servers='localhost:9092') print('Producer created..............') with open('/home/ravi/test.json') as f: data = json.load(f) for state in data['states']: # 先转成JSON字符串,再编码成字节 json_str = json.dumps(state) producer.send('ECJson', json_str.encode('utf-8')) producer.flush()
简单总结下:
json.dump()→ 写数据到文件,需要文件指针json.dumps()→ 把对象转成JSON字符串,适合网络传输(比如发给Kafka)- Kafka需要字节类型的消息,所以必须把字符串编码成字节
内容的提问来源于stack exchange,提问作者Sun
相关产品推荐
相关产品推荐

