如何提升从MongoDB取数并写入Kafka的速率?当前速率357条/秒
嘿,针对你当前每秒357条的传输速度,咱们从MongoDB读取的核心环节入手,给你几个立竿见影的优化点,结合你的代码来调整:
1. 开启批量读取,减少网络往返开销
你现在是逐条遍历db.Response.find()的结果,MongoDB默认的批量拉取量很小,会导致频繁和数据库建立网络请求。通过设置batch_size一次性拉取更多文档,能大幅减少网络IO次数:
# 每次从MongoDB拉取1000条数据(可根据服务器内存情况调整,比如2000) cursor = db.Response.find().batch_size(1000) for response in cursor: # 后续业务逻辑不变
2. 只读取需要的字段(投影优化)
如果Response集合里有很多你不需要推送到Kafka的字段,一定要用**投影(projection)**过滤掉,这样返回的文档体积更小,传输和处理效率都会提升:
# 只保留业务需要的字段,比如_id、content、create_time,排除其他字段 cursor = db.Response.find({}, {"_id": 1, "content": 1, "create_time": 1}).batch_size(1000)
3. 优化数据序列化方式
你当前用bytes(response)转换MongoDB文档的效率很低。MongoDB原生用BSON格式存储数据,直接用BSON序列化比转JSON快得多:
import bson # 用BSON原生序列化,比标准json.dumps快数倍 future = producer.send('collect-production-response', bson.encode(response))
如果必须用JSON格式,建议用ujson库代替Python标准库的json,序列化速度会有明显提升。
4. 批量推送Kafka,避免频繁IO
你现在逐条发送数据到Kafka,即使有flush操作,也会频繁触发Kafka的网络请求。可以攒一批数据再批量发送,既能提升Kafka的发送效率,也能让MongoDB的读取过程更连续:
batch_size = 500 # 每攒500条推送一次 batch_data = [] for idx, response in enumerate(cursor, 1): try: batch_data.append(bson.encode(response)) # 攒够批量就发送 if idx % batch_size == 0: for data in batch_data: producer.send('collect-production-response', data) producer.flush() batch_data = [] except Exception as e: print(e) if idx >= 100000: print("Done 100k") # 发送剩余的未批量数据 if batch_data: for data in batch_data: producer.send('collect-production-response', data) producer.flush() sys.exit()
5. 调整MongoDB连接池配置
默认的pymongo连接池大小可能不足以支撑高并发读取,你可以在创建MongoClient时增大连接池,为后续的多线程/多进程优化预留空间:
client = MongoClient('my_uri', maxPoolSize=20, minPoolSize=5)
6. 索引与遍历优化(全量读取场景)
如果你是全量读取整个集合,MongoDB默认会用_id索引来遍历(默认自带)。如果你的读取需要按特定字段排序(比如create_time),一定要给该字段建立索引,避免全表扫描:
# 先给排序字段建索引(只需执行一次) # db.Response.create_index("create_time") # 按时间顺序批量读取 cursor = db.Response.find().sort("create_time", pymongo.ASCENDING).batch_size(1000)
把这些优化点结合起来,你的MongoDB读取速度应该能提升数倍,轻松突破每秒数千条的传输量。
内容的提问来源于stack exchange,提问作者hasherBaba

