使用asyncio.TaskGroup时协程未被等待问题及信号量封装优化咨询
问题分析与解决方案
问题描述
当前runner函数支持传入可选的信号量参数,传入时会把所有协程包装进信号量的上下文管理器。但当其中一个协程抛出异常时,会触发coroutine 'my_task' was never awaited的运行时警告。我猜测原因是:当某个任务抛出异常时,TaskGroup会取消所有剩余任务,如果任务还没进入信号量上下文作用域就被取消,原协程就从未被await过。请问这个理解对吗?如果正确,有没有更优的信号量封装方式?
复现代码
import asyncio from typing import Any, Coroutine, TypeVar T = TypeVar("T") def _with_sem(a: Coroutine[Any, Any, T], sem: asyncio.Semaphore): async def wrapper(): async with sem: return await a return wrapper() async def runner(coros: list[Coroutine[Any, Any, T]], sem: asyncio.Semaphore | None = None): if sem: coros = [_with_sem(c, sem) for c in coros] async with asyncio.TaskGroup() as tgp: for c in coros: tgp.create_task(c) async def my_task(i: int): print("in task", i) if i == 27: raise Exception("test") await asyncio.sleep(0) async def main(): await runner([my_task(i) for i in range(100)], sem=asyncio.BoundedSemaphore(10)) if __name__ == "__main__": asyncio.run(main())
运行输出
in task 0 in task 1 in task 2 in task 3 in task 4 in task 5 in task 6 in task 7 in task 8 in task 9 in task 10 in task 11 in task 12 in task 13 in task 14 in task 15 in task 16 in task 17 in task 18 in task 19 in task 20 in task 21 in task 22 in task 23 in task 24 in task 25 in task 26 in task 27 in task 28 in task 29 in task 30 /usr/local/lib/python3.11/asyncio/base_events.py:1921: RuntimeWarning: coroutine 'my_task' was never awaited handle = self._ready.popleft() RuntimeWarning: Enable tracemalloc to get the object allocation traceback /usr/local/lib/python3.11/asyncio/base_events.py:1937: RuntimeWarning: coroutine 'my_task' was never awaited handle = None # Needed to break cycles when an exception occurs. RuntimeWarning: Enable tracemalloc to get the object allocation traceback + Exception Group Traceback (most recent call last): | File "/code/src/test2.py", line 38, in <module> | asyncio.run(main()) | File "/usr/local/lib/python3.11/asyncio/runners.py", line 190, in run | return runner.run(main) | ^^^^^^^^^^^^^^^^ | File "/usr/local/lib/python3.11/asyncio/runners.py", line 118, in run | return self._loop.run_until_complete(task) | ^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^ | File "/usr/local/lib/python3.11/asyncio/base_events.py", line 654, in run_until_complete | return future.result() | ^^^^^^^^^^^^^^^ | File "/code/src/test2.py", line 34, in main | await runner([my_task(i) for i in range(100)], sem=asyncio.BoundedSemaphore(10)) | File "/code/src/test2.py", line 21, in runner | async with asyncio.TaskGroup() as tgp: | File "/usr/local/lib/python3.11/asyncio/taskgroups.py", line 145, in __aexit__ | raise me from None | ExceptionGroup: unhandled errors in a TaskGroup (1 sub-exception) +-+---------------- 1 ---------------- | Traceback (most recent call last): | File "/code/src/test2.py", line 13, in wrapper | return await a | ^^^^^^^ | File "/code/src/test2.py", line 29, in my_task | raise Exception("test") | Exception: test +------------------------------------
解答
你的理解完全正确
当TaskGroup中的某个任务抛出未处理异常时,它会立即触发所有未完成任务的取消操作。对于那些还在等待信号量的包装任务来说,取消会发生在async with sem:的等待阶段——此时原协程my_task还没被await,而协程对象一旦创建就必须被await,否则就会触发never awaited警告。
更优的信号量封装方式
核心思路是:确保原协程无论是否被取消,都会被正确await。推荐以下两种简洁可靠的实现方式:
方式一:动态在任务创建环节包装信号量逻辑
不在外部提前批量包装协程,而是在runner创建任务时,为每个协程动态添加信号量上下文,避免提前创建大量未执行的协程对象:
async def runner(coros: list[Coroutine[Any, Any, T]], sem: asyncio.Semaphore | None = None): async with asyncio.TaskGroup() as tgp: for c in coros: # 定义包装函数,捕获当前循环的协程对象 async def wrapped(coro): if sem: async with sem: return await coro return await coro tgp.create_task(wrapped(c))
这种方式从根源上避免了未await的协程问题:每个原协程只会在任务内部被包装,且一定会进入await流程,即使任务被取消,也不会留下孤立的协程对象。
方式二:修改包装器确保原协程始终被处理
如果需要保留提前包装的逻辑,可以修改_with_sem,在finally块中确保原协程被await:
def _with_sem(a: Coroutine[Any, Any, T], sem: asyncio.Semaphore): async def wrapper(): try: async with sem: return await a finally: # 确保原协程被处理,避免未await警告 try: await a except (asyncio.CancelledError, Exception): # 吞取消或其他异常,避免二次抛出干扰TaskGroup的异常处理 pass return wrapper()
最优推荐
方式一更符合TaskGroup的设计逻辑,代码更简洁,也避免了提前创建大量协程对象带来的潜在问题,是优先选择的实现方案。
内容的提问来源于stack exchange,提问作者cjp94
相关产品推荐
相关产品推荐

