如何在Kafka中发布字典?从MongoDB取数至Kafka遇类型错误求助
解决Python3中MongoDB文档转Kafka消息的TypeError问题
这个错误的核心原因是Python3和Python2中bytes()函数的行为完全不同:
- 在Python2里,
bytes(obj)等价于str(obj),会把MongoDB返回的文档对象直接转成字符串形式的字节流; - 但在Python3中,
bytes()的参数要么是整数(指定字节数),要么是可迭代的整数序列,要么是字符串+编码参数。你直接传入MongoDB的dict类型文档时,它会先被转成字符串(比如"{'_id': ObjectId('...'), ...}"),然后bytes()会把这个字符串当成要生成的字节数,自然就会抛出"str不能解释为整数"的错误。
正确的解决方案
我们需要把MongoDB的文档先序列化为JSON字符串,再编码成UTF-8字节流发送给Kafka,同时还要处理MongoDB特有的数据类型(比如ObjectId、datetime),这些类型默认的json.dumps无法序列化。
步骤1:修复KafkaProducer初始化
你代码里初始化了两次KafkaProducer,第二次会覆盖第一次的bootstrap_servers配置,这是冗余且错误的,应该合并成一个配置。
步骤2:序列化MongoDB文档
有两种常用的序列化方式:
- 使用Python标准库的
json.dumps,通过default=str参数处理特殊类型; - 使用pymongo自带的
json_util,可以更精准地序列化MongoDB的BSON类型。
修改后的完整代码
方式1:使用标准json库处理
from kafka import KafkaProducer from kafka.errors import KafkaError import json import pymongo from pymongo import MongoClient import sys import datetime try: client = MongoClient('mongodb://A.B.C.D:27017/prod-production') db = client["prod-production"] except Exception as e: print("Error occurred while connecting to DB") print(e) sys.exit() # 连接失败后直接退出,避免后续报错 # 合并KafkaProducer配置,添加序列化器 producer = KafkaProducer( bootstrap_servers=['localhost:9092'], retries=5, value_serializer=lambda v: json.dumps(v, default=str).encode('utf-8') ) print("Initial time:") print(datetime.datetime.now()) count = 1 for response in db.Response.find(): if count >= 20: producer.flush() sys.exit() count += 1 print(count) # 直接传入文档对象,序列化器会自动处理 producer.send('example-topic', response) print("Final time") print(datetime.datetime.now())
方式2:使用pymongo的json_util(更推荐)
from kafka import KafkaProducer from kafka.errors import KafkaError import json from bson import json_util # 导入json_util import pymongo from pymongo import MongoClient import sys import datetime try: client = MongoClient('mongodb://A.B.C.D:27017/prod-production') db = client["prod-production"] except Exception as e: print("Error occurred while connecting to DB") print(e) sys.exit() producer = KafkaProducer( bootstrap_servers=['localhost:9092'], retries=5, value_serializer=lambda v: json_util.dumps(v).encode('utf-8') ) print("Initial time:") print(datetime.datetime.now()) count = 1 for response in db.Response.find(): if count >= 20: producer.flush() sys.exit() count += 1 print(count) producer.send('example-topic', response) print("Final time") print(datetime.datetime.now())
额外说明
value_serializer是KafkaProducer的内置参数,会自动帮你把消息序列化为字节流,不需要手动调用bytes();json_util.dumps可以正确处理MongoDB的ObjectId、Decimal128、datetime等特殊类型,序列化后的JSON会保留这些类型的原始信息,方便消费者解析;- 代码中添加了
sys.exit()在数据库连接失败时,避免后续代码因未连接数据库而抛出更多错误。
内容的提问来源于stack exchange,提问作者hasherBaba
相关产品推荐
相关产品推荐

