asyncio+pymongo双游标并行迭代无效率提升问题排查
并行查询MongoDB耗时异常问题分析
我需要从src和dst两个MongoDB数据库查询数据,尝试了两种写法,总耗时却完全相同,无法理解原因,希望实现真正的并行查询来减少耗时。
方法1:使用asyncio.gather并行查询
此方法中单个src与dst查询耗时相近,但总耗时接近两者耗时之和,未体现并行优势。
代码实现
async def get_src_cursor_data(self, src_coll, values, start_idx): start4 = time.time() src_cursor = src_coll.find({"_id": {"$in": values}}) if not isinstance(src_cursor, list): while src_cursor.alive: doc = await src_cursor.next() end4 = time.time() self.log_info('[{}] get_src_cursor_data: {}'.format(start_idx, end4 - start4)) async def get_dst_cursor_data(self, dst_coll, values, start_idx): start5 = time.time() dst_cursor = dst_coll.find({"_id": {"$in": values}}) if not isinstance(dst_cursor, list): while dst_cursor.alive: dst = await dst_cursor.next() end5 = time.time() self.log_info('[{}] get_dst_cursor_data: {}'.format(start_idx, end5 - start5)) async def process_document(self, src_coll, dst_coll, values, start_idx): start = time.time() res = await asyncio.gather(self.get_src_cursor_data(src_coll, values, start_idx), self.get_dst_cursor_data(dst_coll, values, start_idx)) end0 = time.time() self.log_info( "[{}] process {} docs total time :{}".format( start_idx, len(values), end0 - start0))
运行日志
[20] get_src_cursor_data: 5.120812177658081 [60] get_src_cursor_data: 5.62122654914856 [0] get_src_cursor_data: 6.334164142608643 [80] get_src_cursor_data: 6.916379451751709 [40] get_src_cursor_data: 7.283424139022827 [20] get_dst_cursor_data: 5.805237531661987 [20] process 20 docs total time :10.929065227508545, get src data time :5.120812177658081, get dst data time: 5.805237531661987
方法2:串行查询
此方法中单个src与dst查询耗时为方法1的两倍,总耗时与方法1基本一致。
代码实现
async def process_sample_document(self, src_coll, dst_coll, values, start_idx): start4 = time.time() src_cursor = src_coll.find({"_id": {"$in": values}}) if not isinstance(src_cursor, list): while src_cursor.alive: doc = await src_cursor.next() end4 = time.time() self.log_info('[{}] get_src_cursor_data: {}'.format(start_idx, end4 - start4)) start5 = time.time() dst_cursor = dst_coll.find({"_id": {"$in": values}}) if not isinstance(dst_cursor, list): while dst_cursor.alive: dst = await dst_cursor.next() end5 = time.time() self.log_info('[{}] get_dst_cursor_data: {}'.format(start_idx, end5 - start5)) self.log_info( "[{}] process {} docs total time :{}, get src data time :{}, get dst data time: {}".format(start_idx, len(values), end5 - start4, end4 - start4, end5 - start5))
运行日志
[0] get_dst_cursor_data: 9.8009352684021 [40] get_dst_cursor_data: 10.02169156074524 [20] get_dst_cursor_data: 10.243099212646484 [60] get_dst_cursor_data: 10.4615159034729 [80] get_dst_cursor_data: 10.679930925369263 [0] get_src_cursor_data: 10.909341096878052 [40] get_src_cursor_data: 11.130756378173828 [80] get_src_cursor_data: 11.351166486740112 [20] get_src_cursor_data: 11.579554557800293 [60] get_src_cursor_data: 11.807490825653076 [0] process 20 docs total time :11.81546807289123
问题分析与解决方案
核心原因
- MongoDB异步驱动连接池限制:默认情况下,MongoDB的async客户端可能只维护少量连接,并行查询时会争抢连接资源,导致每个查询的等待时间增加,最终总耗时和串行接近。
- 逐行获取文档的低效性:使用
while cursor.alive + await cursor.next()逐个获取文档,会产生大量IO交互,放大了并行时的资源竞争问题。 - 数据库端瓶颈:如果src和dst是同一MongoDB实例或集群,数据库本身可能会串行处理来自同一客户端的请求,无法真正并行。
优化方案
- 调整连接池配置:创建MongoDB客户端时,增大
maxPoolSize参数,确保有足够的连接支持并行查询:from motor.motor_asyncio import AsyncIOMotorClient src_client = AsyncIOMotorClient(src_uri, maxPoolSize=10) dst_client = AsyncIOMotorClient(dst_uri, maxPoolSize=10) - 批量获取文档:改用
await cursor.to_list(length=None)一次性获取所有文档,减少IO交互次数:async def get_src_cursor_data(self, src_coll, values, start_idx): start4 = time.time() src_docs = await src_coll.find({"_id": {"$in": values}}).to_list(length=None) end4 = time.time() self.log_info('[{}] get_src_cursor_data: {}'.format(start_idx, end4 - start4)) - 确认数据库部署:如果src和dst是独立的MongoDB实例,确保它们的网络带宽充足,且数据库服务器有足够的CPU/内存资源处理并行请求。
内容的提问来源于stack exchange,提问作者yuancheng hu
相关产品推荐
相关产品推荐

