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

如何在不同时调度协程的情况下等待指定子协程列表完成?

解决方法:等待队列中指定协程执行完成

要实现仅等待指定子列表的协程完成,核心是跟踪目标协程的执行状态,以下是两种可行方案:

方案1:将协程转为Task后放入队列(最优解)

直接把协程包装成asyncio.Task,保留Task的引用,后续通过asyncio.gather等待指定子列表的Task完成。Task会在被工作者await时自动执行,无需额外处理。

import asyncio

# 示例协程
async def coro1():
    await asyncio.sleep(1)
    print("coro1 done")

async def coro2():
    await asyncio.sleep(2)
    print("coro2 done")

async def coro3():
    await asyncio.sleep(0.5)
    print("coro3 done")

async def coro4():
    await asyncio.sleep(1.5)
    print("coro4 done")

async def coro5():
    await asyncio.sleep(0.8)
    print("coro5 done")

async def coro6():
    await asyncio.sleep(1.2)
    print("coro6 done")

# 队列工作者:负责读取并执行Task
async def queue_worker(queue):
    while True:
        task = await queue.get()
        try:
            await task
        finally:
            queue.task_done()

async def main():
    coro_list = [coro1(), coro2(), coro3(), coro4(), coro5(), coro6()]
    # 将所有协程转为Task,保留引用以便后续等待
    task_list = [asyncio.create_task(coro) for coro in coro_list]
    
    queue = asyncio.Queue()
    # 将Task放入队列
    for task in task_list:
        queue.put_nowait(task)
    
    # 启动工作者任务
    worker = asyncio.create_task(queue_worker(queue))
    
    # 等待指定子列表的Task完成(示例:索引2到6,对应coro3到coro6)
    await asyncio.gather(*task_list[2:6])
    print("指定的协程已全部执行完成")
    
    # 等待队列所有任务处理完毕,关闭工作者
    await queue.join()
    worker.cancel()
    await worker

asyncio.run(main())

方案2:包装协程并通过Future跟踪状态(必须放原协程到队列时用)

如果必须将原始协程放入队列,可以为目标协程创建asyncio.Future,用包装器协程在执行完成后更新Future状态,最后等待这些Future完成。

import asyncio

# 示例协程同上,省略重复定义
async def coro1():
    await asyncio.sleep(1)
    print("coro1 done")

async def coro2():
    await asyncio.sleep(2)
    print("coro2 done")

async def coro3():
    await asyncio.sleep(0.5)
    print("coro3 done")

async def coro4():
    await asyncio.sleep(1.5)
    print("coro4 done")

async def coro5():
    await asyncio.sleep(0.8)
    print("coro5 done")

async def coro6():
    await asyncio.sleep(1.2)
    print("coro6 done")

async def queue_worker(queue):
    while True:
        coro = await queue.get()
        try:
            await coro
        finally:
            queue.task_done()

async def main():
    coro_list = [coro1(), coro2(), coro3(), coro4(), coro5(), coro6()]
    target_indices = range(2, 6)  # 指定要等待的协程索引范围
    tracker_futures = []
    
    # 包装目标协程,绑定Future跟踪完成状态
    for idx in target_indices:
        original_coro = coro_list[idx]
        fut = asyncio.Future()
        
        # 包装器:执行原协程后更新Future
        async def wrapped_coro(coro, future):
            try:
                result = await coro
                future.set_result(result)
            except Exception as e:
                future.set_exception(e)
        
        coro_list[idx] = wrapped_coro(original_coro, fut)
        tracker_futures.append(fut)
    
    queue = asyncio.Queue()
    for coro in coro_list:
        queue.put_nowait(coro)
    
    worker = asyncio.create_task(queue_worker(queue))
    
    # 等待所有跟踪的Future完成
    await asyncio.gather(*tracker_futures)
    print("指定的协程已全部执行完成")
    
    await queue.join()
    worker.cancel()
    await worker

asyncio.run(main())

关键说明

  • 方案1更简洁高效,因为Task本身就是可等待的对象,无需额外的包装逻辑。
  • 两种方案都不影响队列工作者的正常执行,仅在指定协程完成后触发后续逻辑。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.21 18:33:15