Asyncio动态添加任务至运行中事件循环及shutdown_asyncgens解析
Asyncio动态任务数量处理及代码行解释
原场景与实现代码
我们有12个任务,要求同一时间仅能运行3个任务:初始启动3个任务,待其中一个完成后再启动下一个。以下是用Asyncio结合Semaphore实现的代码:
import asyncio import random max_tasks = 12 sem = asyncio.Semaphore(3) async def counter(n): print(f'counter with argument {n} has been launched') for i in range(n): for j in range(n): for k in range(n): pass await asyncio.sleep(1) print(f'counter with argument {n} has FINISHED') async def safe_calc(n): async with sem: await counter(n) async def main(): tasks = [asyncio.ensure_future(safe_calc(random.randint(100, 600))) for _ in range(max_tasks)] await asyncio.gather(*tasks) loop = asyncio.get_event_loop() loop.run_until_complete(main()) loop.run_until_complete(loop.shutdown_asyncgens()) loop.close()
一、动态处理任务数量(运行中新增任务)
要在事件循环运行过程中动态新增任务,核心是放弃一次性创建所有任务的方式,转而维护动态的任务集合,结合信号监听逻辑实现任务的按需添加,同时保留Semaphore的并发控制。具体实现思路如下:
- 用任务集合替代固定列表:不再提前生成所有任务,而是用集合追踪当前运行/待运行的任务,避免内存泄漏。
- 新增任务触发逻辑:通过协程监听外部信号(比如队列、状态变量),一旦需要新增任务,就创建对应的任务并加入集合。
- 自动清理完成任务:给每个任务添加回调,任务完成后自动从集合中移除。
示例代码:
import asyncio import random sem = asyncio.Semaphore(3) task_set = set() # 用队列模拟动态任务参数的来源 task_queue = asyncio.Queue() async def counter(n): print(f'counter with argument {n} has been launched') for i in range(n): for j in range(n): for k in range(n): pass await asyncio.sleep(1) print(f'counter with argument {n} has FINISHED') async def safe_calc(n): async with sem: await counter(n) async def task_generator(): """模拟动态生成任务的逻辑:初始发3个,之后每2秒发1个,共发12个""" for _ in range(3): await task_queue.put(random.randint(100, 600)) for _ in range(9): await asyncio.sleep(2) await task_queue.put(random.randint(100, 600)) # 发送结束信号 await task_queue.put(None) async def task_manager(): while True: n = await task_queue.get() if n is None: break # 创建任务并加入集合 task = asyncio.create_task(safe_calc(n)) task_set.add(task) # 任务完成后自动从集合移除 task.add_done_callback(task_set.discard) async def main(): # 启动任务生成器和管理器 await asyncio.gather(task_generator(), task_manager()) # 等待剩余任务全部完成 await asyncio.gather(*task_set) loop = asyncio.get_event_loop() loop.run_until_complete(main()) loop.run_until_complete(loop.shutdown_asyncgens()) loop.close()
这个方案中,只要向task_queue中传入任务参数,就能动态触发新任务的创建,Semaphore依然保证最多3个任务同时运行。
二、loop.run_until_complete(loop.shutdown_asyncgens())的作用
这行代码的核心作用是安全关闭所有异步生成器,确保它们能执行完自身的清理逻辑(比如finally块、async with的退出操作)。
异步生成器是指用async def定义、包含yield语句的函数,这类函数在运行过程中可能持有文件句柄、网络连接等资源。如果不调用该方法,事件循环直接关闭时会强制终止这些生成器,导致资源泄漏或清理代码未执行。
你的代码里虽然没有直接使用异步生成器,但加上这行是通用的良好实践——后续代码扩展引入异步生成器时,能避免潜在的资源问题。另外,Python 3.7+的asyncio.run()会自动处理异步生成器的关闭,但手动管理事件循环时,就需要手动调用这行完成清理。
内容的提问来源于stack exchange,提问作者user20426821
相关产品推荐
相关产品推荐

