Asyncio技术问题:如何在主Future中并行运行任务状态检查Future
异步并行处理任务提交与状态检查的解决方案
我看了你这段异步代码的问题啦——你已经创建了检查任务状态的future,但不知道怎么让它和主流程并行跑起来对吧?其实核心是要管理这些异步任务的生命周期,既让它们在后台自主执行,又要确保主流程结束时不会丢失结果或者留下未完成的任务。
核心改动点
- 用一个列表保存所有正在运行的检查任务
future,避免它们被GC回收或者中途终止 - 在主循环结束后,等待所有检查任务完成,确保结果都能被收集
- 优化
ClientSession的创建逻辑,减少重复创建连接的开销
修改后的完整代码
import asyncio from time import time from aiohttp import ClientSession async def post_async_recognize_timeout(in_url, filepath, outer_ids, timeout, exp_time, token, check_timeout): start_time = time() # 用来保存所有正在运行的检查任务future pending_futures = [] # 把ClientSession移到循环外,避免每次创建新连接 async with ClientSession() as session: while time() - start_time <= exp_time: with filepath.open('rb') as file: in_files = {'image': file} async with session.post(in_url, data=in_files) as response: body = await response.json() task_id = body['task_id'] # 创建检查任务的future并加入待处理列表 future = asyncio.create_task( check_task_status(in_url, token, outer_ids, check_timeout, {'id': task_id, 'created_at': time()}) ) pending_futures.append(future) await asyncio.sleep(timeout) # 主循环结束后,等待所有检查任务完成,return_exceptions=True避免单个任务报错导致整体失败 await asyncio.gather(*pending_futures, return_exceptions=True) # 此时outer_ids里已经收集了所有完成的任务结果 async def check_task_status(in_url, token, result, timeout, task): while True: async with ClientSession() as session: async with session.get(f"{in_url}{task['id']}?token={token}") as response: body = await response.json() if body['message'] != 'Async task not done': # 处理任务完成结果 task['status'] = 'Done' if body['message'] is None else f"ERROR: {body['message']}" task['done_at'] = time() result.append(task) break await asyncio.sleep(timeout)
关键细节说明
- 并行任务管理:创建的
future(这里用更推荐的asyncio.create_task替代ensure_future)会被自动注册到事件循环后台执行,加入pending_futures列表是为了追踪所有任务,避免它们被意外终止。 - 清理与等待:主循环结束后用
asyncio.gather等待所有任务完成,确保所有任务的结果都能被写入outer_ids。return_exceptions=True参数能保证即使某个检查任务报错,其他任务也能正常完成。 - 连接优化:把
ClientSession移到主循环外面,避免每次提交任务都创建新的HTTP连接池,有效提升性能。
内容的提问来源于stack exchange,提问作者MrLalatg
相关产品推荐
相关产品推荐

