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

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

问题分析与解决方案

核心原因

  1. MongoDB异步驱动连接池限制:默认情况下,MongoDB的async客户端可能只维护少量连接,并行查询时会争抢连接资源,导致每个查询的等待时间增加,最终总耗时和串行接近。
  2. 逐行获取文档的低效性:使用while cursor.alive + await cursor.next()逐个获取文档,会产生大量IO交互,放大了并行时的资源竞争问题。
  3. 数据库端瓶颈:如果src和dst是同一MongoDB实例或集群,数据库本身可能会串行处理来自同一客户端的请求,无法真正并行。

优化方案

  1. 调整连接池配置:创建MongoDB客户端时,增大maxPoolSize参数,确保有足够的连接支持并行查询:
    from motor.motor_asyncio import AsyncIOMotorClient
    src_client = AsyncIOMotorClient(src_uri, maxPoolSize=10)
    dst_client = AsyncIOMotorClient(dst_uri, maxPoolSize=10)
    
  2. 批量获取文档:改用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))
    
  3. 确认数据库部署:如果src和dst是独立的MongoDB实例,确保它们的网络带宽充足,且数据库服务器有足够的CPU/内存资源处理并行请求。

内容的提问来源于stack exchange,提问作者yuancheng hu

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.24 13:27:02