如何在不同时调度协程的情况下等待指定子协程列表完成?
解决方法:等待队列中指定协程执行完成
要实现仅等待指定子列表的协程完成,核心是跟踪目标协程的执行状态,以下是两种可行方案:
方案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
相关产品推荐
相关产品推荐

