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

如何用Python asyncio实现类似ThreadPoolExecutor的并发任务(无需gather)

用Python asyncio实现固定并发数的任务调度(类似ThreadPoolExecutor)

如果你需要用asyncio实现类似ThreadPoolExecutor的固定并发控制——限制同时运行的任务数量,且某个任务完成后立即补充新任务,而不是等整批任务全部结束再启动下一批,可以通过以下方式实现:

核心思路是维护一个正在运行的任务集合,结合asyncio.wait的FIRST_COMPLETED模式,循环等待任务完成并动态补充新任务,直到所有任务都执行完毕。

基础实现(不带结果收集)

import asyncio

async def bounded_gather(tasks, concurrency_limit):
    task_iter = iter(tasks)
    running_tasks = set()

    # 启动第一批并发任务
    for _ in range(min(concurrency_limit, len(tasks))):
        try:
            task = next(task_iter)
            running_tasks.add(task)
        except StopIteration:
            break

    while running_tasks:
        # 等待任意一个任务完成
        done, pending = await asyncio.wait(running_tasks, return_when=asyncio.FIRST_COMPLETED)
        # 移除已完成的任务
        running_tasks.difference_update(done)

        # 补充新任务,维持并发数上限
        while len(running_tasks) < concurrency_limit:
            try:
                task = next(task_iter)
                running_tasks.add(task)
            except StopIteration:
                break

带结果收集的版本(含异常处理)

如果需要收集所有任务的执行结果(包括异常),可以用这个版本:

import asyncio

async def bounded_gather_with_results(tasks, concurrency_limit):
    task_iter = iter(tasks)
    running_tasks = set()
    results = []

    # 初始化第一批任务
    for _ in range(min(concurrency_limit, len(tasks))):
        task = next(task_iter)
        running_tasks.add(task)

    while running_tasks:
        # 等待首个完成的任务
        done, pending = await asyncio.wait(running_tasks, return_when=asyncio.FIRST_COMPLETED)
        
        # 处理完成任务的结果或异常
        for task in done:
            try:
                results.append(task.result())
            except Exception as e:
                # 可根据需求修改异常处理逻辑,比如记录日志
                results.append(e)
        
        running_tasks = pending

        # 补充新任务,保持并发数
        while len(running_tasks) < concurrency_limit:
            try:
                task = next(task_iter)
                running_tasks.add(task)
            except StopIteration:
                break

    return results

使用示例

async def main():
    # 生成模拟任务:每个任务耗时i秒
    tasks = [asyncio.create_task(asyncio.sleep(i)) for i in range(10)]
    # 指定并发数为3,即同时最多运行3个任务
    results = await bounded_gather_with_results(tasks, 3)
    print(results)  # 输出:[0,1,2,3,4,5,6,7,8,9]

if __name__ == "__main__":
    asyncio.run(main())

关键说明

  • asyncio.wait(..., return_when=asyncio.FIRST_COMPLETED):这是实现动态补充任务的核心,它会在任意一个任务完成时立即返回,而不是等待所有任务结束。
  • 用迭代器遍历任务列表:避免一次性加载所有任务到内存,适合任务数量极大的场景。
  • 并发数控制:通过running_tasks集合的长度维持指定的并发上限,确保不会同时运行过多任务。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.11 23:15:36