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

asyncio中run_until_complete早于run_in_executor任务完成结束问题

问题原因

你修改后的run_query函数没有等待asyncio.gather(*writers)执行完成,直接将gather返回的Future对象作为返回值抛出,导致上层任务认为run_query已经执行完毕,整个主流程提前结束。

核心原理说明
  • loop.run_in_executor本身会返回一个Future对象,对应线程池里的阻塞写入任务的执行状态
  • asyncio.gather(*writers)会将多个Future打包成一个新的Future,只有当所有传入的Future都执行完成时,这个打包的Future才会进入完成状态
  • 你没有对asyncio.gather(*writers)加await,run_query协程执行到return语句就直接结束,上层的asyncio.gather(*tasks)只会等待你传入的run_query协程完成,不会自动等待run_query返回的新Future执行,因此会直接返回结果,主流程走到log.info('Done.')逻辑。此时线程池里的写入任务还在后台执行,所以你会先看到Done日志,再看到写入日志。
修复方案

只需要修改run_query的返回逻辑,添加await等待打包的写入任务全部完成即可,同时这种写法相比原来串行await每个写入的效率更高,所有写入任务会并行在线程池中执行:

async def run_query(query: List[Dict[str, Any]]) -> Union[asyncio.Future, List[asyncio.Future], Exception]:
    _loop = asyncio.get_event_loop()
    try:
        result = await asyncio.create_task(query['func']())
        writers = []
        for fname in result.keys():
            writers.append(_loop.run_in_executor(None, write_query_result, query['folder'], fname+query['file_extension'], result))
        # 新增await,等待所有写入任务完成再返回
        return await asyncio.gather(*writers)
    except Exception as e:
        log.error('Error getting results. Query [%s]. Error [%s]' % (
            query['name'], e))
        raise e

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.03 02:24:03