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

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)

关键细节说明

  1. 并行任务管理:创建的future(这里用更推荐的asyncio.create_task替代ensure_future)会被自动注册到事件循环后台执行,加入pending_futures列表是为了追踪所有任务,避免它们被意外终止。
  2. 清理与等待:主循环结束后用asyncio.gather等待所有任务完成,确保所有任务的结果都能被写入outer_ids。return_exceptions=True参数能保证即使某个检查任务报错,其他任务也能正常完成。
  3. 连接优化:把ClientSession移到主循环外面,避免每次提交任务都创建新的HTTP连接池,有效提升性能。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.09 19:57:34