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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.13 01:16:11