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

如何让asyncio的tg_fast任务组优先执行,再继续tg_main任务组?

问题描述
  • 需要立即执行任务组tg_fast,之后继续执行任务组tg_main(若无法继续则重新启动),asyncio.gather()效果与TaskGroup类似
  • 当前代码输出顺序为 0 => 2 => 10 => 100,期望输出顺序为 0 => 10 => 100 => ... 或 0 => 100 => 10 => ...,目标是在输出0之后、输出2之前启动10和100对应的任务
  • 要求同时调用another_coro,无需等待它们执行完毕,只需运行到await asyncio.sleep(.1)后即可继续事件循环

当前代码:

import asyncio


async def another_coro(i):
    print(i)
    await asyncio.sleep(.1)


async def coro(i):
    if i == 1:
        async with asyncio.TaskGroup() as tg_fast:
            tg_fast.create_task(another_coro(i * 10))
            tg_fast.create_task(another_coro(i * 100))
        # await asyncio.gather(*[another_coro(i * 10), another_coro(i * 100)])
    else:
        print(i)
        await asyncio.sleep(.1)


async def main():
    async with asyncio.TaskGroup() as tg_main:
        for i in range(0, 3):
            tg_main.create_task(coro(i))


asyncio.run(main(), debug=True)
解决方案

问题根源在于coro(1)中使用async with asyncio.TaskGroup()会等待内部所有任务完成后才会继续执行,导致coro(1)被阻塞,此时事件循环会优先调度coro(2),从而先输出2。

要实现需求,只需让tg_fast的任务后台异步运行,不阻塞coro(1)的执行,这样事件循环会优先处理another_coro的任务,再调度coro(2)。可以通过直接创建独立任务(或使用asyncio.gather批量创建任务)来实现:

修改后的基础代码

import asyncio


async def another_coro(i):
    print(i)
    await asyncio.sleep(.1)


async def coro(i):
    if i == 1:
        # 直接创建异步任务,无需等待完成
        asyncio.create_task(another_coro(i * 10))
        asyncio.create_task(another_coro(i * 100))
        # 也可以用asyncio.gather批量创建任务(需包装成独立任务)
        # asyncio.create_task(asyncio.gather(another_coro(i*10), another_coro(i*100)))
    else:
        print(i)
        await asyncio.sleep(.1)


async def main():
    async with asyncio.TaskGroup() as tg_main:
        for i in range(0, 3):
            tg_main.create_task(coro(i))


asyncio.run(main(), debug=True)

效果说明

  • 去掉async with TaskGroup()的等待逻辑后,coro(1)会立即创建两个another_coro任务,然后自身执行完成
  • 事件循环在coro(0)打印0并进入sleep后,会优先调度刚创建的another_coro任务,输出10和100(顺序由事件循环调度决定)
  • 最后才会调度coro(2),输出2,完全符合需求

增加tg_main异常重启逻辑

如果需要保证tg_fast的任务出现异常时自动重启tg_main,可以在执行逻辑中添加异常捕获与递归重启:

import asyncio


async def another_coro(i):
    print(i)
    await asyncio.sleep(.1)
    # 模拟异常场景
    if i == 100:
        raise ValueError("Test error")


async def coro(i):
    if i == 1:
        asyncio.create_task(another_coro(i * 10))
        asyncio.create_task(another_coro(i * 100))
    else:
        print(i)
        await asyncio.sleep(.1)


async def run_main():
    try:
        async with asyncio.TaskGroup() as tg_main:
            for i in range(0, 3):
                tg_main.create_task(coro(i))
    except Exception as e:
        print(f"tg_main 异常: {e},正在重启...")
        await run_main()


asyncio.run(run_main(), debug=True)

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.04 03:25:35