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
相关产品推荐
相关产品推荐

