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

使用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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.12 09:24:54