You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何在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文档

有两种常用的序列化方式:

  1. 使用Python标准库的json.dumps,通过default=str参数处理特殊类型;
  2. 使用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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.05.29 07:36:50