如何让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
相关产品推荐
相关产品推荐

