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

如何基于asyncio实现可控最大并发数的动态任务补位执行器

asyncio 自定义限流动态补位执行器实现方案

方案一:基于asyncio.wait实现(最简适配)

该方案直接解决原有as_completed实现的缺陷,逻辑直观无额外依赖,完全符合运行规则:

import asyncio

# 模拟业务协程,可替换为你的实际逻辑
async def job(param):
    await asyncio.sleep(1)
    print(f"执行完成,参数:{param}")
    return param

async def custom_executor(async_level: int, params: list):
    # 转迭代器,大参数列表也不会额外占用内存,取完自动抛出StopIteration
    param_iter = iter(params)
    running_tasks = set()

    # 初始填充并发池到最大限制
    for _ in range(async_level):
        try:
            param = next(param_iter)
            running_tasks.add(asyncio.create_task(job(param)))
        except StopIteration:
            # 参数数量小于并发数的场景直接跳出
            break
    
    # 循环处理直到所有任务跑完
    while running_tasks:
        # 等待任意一个任务完成,无需等待整批结束
        done, running_tasks = await asyncio.wait(
            running_tasks,
            return_when=asyncio.FIRST_COMPLETED
        )

        # 处理已完成的任务,可在此处捕获异常、收集返回结果
        for task in done:
            # 示例:打印返回结果,异常可通过try-except task.result()捕获
            print(f"任务返回结果:{task.result()}")

            # 补入新任务
            try:
                new_param = next(param_iter)
                running_tasks.add(asyncio.create_task(job(new_param)))
            except StopIteration:
                # 无剩余参数,无需补位
                pass

# 测试运行
if __name__ == "__main__":
    asyncio.run(custom_executor(async_level=4, params=[i for i in range(1, 11)]))

逻辑说明

  1. 完全匹配三条运行规则:参数为空直接终止,每次取新参数创建协程,同一时间最多async_level个协程运行,任务完成立即补位
  2. 不会预先创建所有协程对象,同时存在的协程数最高等于async_level,无参数量大时的资源过载问题
  3. 可灵活扩展异常处理、结果收集等自定义逻辑

方案二:基于asyncio.Queue实现(消费者模式)

适合需要动态追加参数、多生产者场景的更灵活方案:

import asyncio

async def job(param):
    await asyncio.sleep(1)
    print(f"执行完成,参数:{param}")
    return param

# 消费者协程,数量和最大并发数一一对应
async def consumer(queue: asyncio.Queue):
    while True:
        param = await queue.get()
        # 收到结束信号就退出
        if param is None:
            queue.task_done()
            break
        try:
            await job(param)
        finally:
            queue.task_done()

async def custom_executor(async_level: int, params: list):
    queue = asyncio.Queue(maxsize=async_level)
    # 启动对应数量的消费者
    consumers = [asyncio.create_task(consumer(queue)) for _ in range(async_level)]

    # 往队列塞参数,队列满会自动阻塞,不会一次性加载所有参数
    for param in params:
        await queue.put(param)
    
    # 塞结束信号,通知消费者退出
    for _ in range(async_level):
        await queue.put(None)
    
    # 等待所有任务执行完成
    await queue.join()
    # 等待所有消费者退出
    await asyncio.gather(*consumers)

# 测试运行
if __name__ == "__main__":
    asyncio.run(custom_executor(async_level=4, params=[i for i in range(1, 11)]))

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.24 14:36:08