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

使用asyncio.gather处理千级任务时的并发控制问题咨询

解决asyncio.gather启动大量任务时的卡顿问题,实现固定并发数执行

你的问题核心在于一次性创建数千个asyncio任务,哪怕这些任务被Semaphore阻塞,它们也会进入事件循环的等待队列,占用大量内存和调度资源,导致启动阶段长时间挂起。要实现固定数量的任务同时执行,关键是控制任务的创建速率,而不是一次性初始化所有任务。

方案1:控制任务创建速率的并发调度

通过维护一个活跃任务列表,保持同时运行的任务数不超过设定的限制,避免一次性创建所有任务。示例代码如下:

import asyncio

async def your_task_logic(task_id):
    # 替换为你的实际任务逻辑(兼具IO和CPU操作)
    await asyncio.sleep(0.5)  # 模拟IO等待
    # CPU密集操作建议放到线程/进程池执行(见方案2)
    result = task_id * 2
    return result

async def run_with_fixed_concurrency(task_items, max_concurrent):
    semaphore = asyncio.Semaphore(max_concurrent)
    
    async def bounded_task(item):
        async with semaphore:
            return await your_task_logic(item)
    
    active_tasks = []
    for item in task_items:
        # 创建单个任务并加入活跃列表
        task = asyncio.create_task(bounded_task(item))
        active_tasks.append(task)
        
        # 当活跃任务数达到上限时,等待至少一个任务完成再继续
        if len(active_tasks) >= max_concurrent:
            done, _ = await asyncio.wait(active_tasks, return_when=asyncio.FIRST_COMPLETED)
            # 移除已完成的任务,保持活跃列表大小不超限
            active_tasks = [t for t in active_tasks if not t.done()]
    
    # 等待剩余所有任务完成
    await asyncio.gather(*active_tasks)

async def main():
    # 模拟2000+个任务
    all_task_items = list(range(2500))
    # 固定10个任务同时执行
    await run_with_fixed_concurrency(all_task_items, max_concurrent=10)

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

这个方案的优势是:

  • 不会一次性创建数千个任务,启动阶段资源占用极低
  • 通过Semaphore严格控制并发执行的任务数量
  • 任务完成一个就补充一个,始终保持设定的并发量

方案2:处理CPU密集型任务的优化

由于你的任务兼具高CPU和IO消耗,而asyncio是单线程事件循环,CPU密集型代码会阻塞整个事件循环,导致所有任务变慢。建议把CPU密集的部分放到线程池或进程池执行:

import asyncio
import concurrent.futures

async def your_task_logic(task_id):
    await asyncio.sleep(0.5)  # IO操作部分
    
    # 将CPU密集操作委托给线程池(进程池用ProcessPoolExecutor)
    loop = asyncio.get_running_loop()
    with concurrent.futures.ThreadPoolExecutor(max_workers=5) as pool:
        cpu_result = await loop.run_in_executor(pool, lambda: sum(range(1000000)) + task_id)
    
    return cpu_result

如果是多核CPU的场景,用ProcessPoolExecutor能更好地利用多核资源,但要注意任务数据必须是可序列化的(比如不能传递复杂的非picklable对象)。

为什么之前的方法会卡顿?

当你用asyncio.gather(*list_of_tasks)一次性传入2000个任务时,asyncio会立即初始化所有任务并加入事件循环的等待队列。哪怕这些任务被Semaphore阻塞,它们依然占用内存,并且事件循环需要处理大量的任务调度逻辑,导致启动阶段出现长时间的卡顿。而上面的方案通过延迟创建任务,只维护少量活跃任务,从根源上避免了这个问题。

内容的提问来源于stack exchange,提问作者M. Liver

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.07 05:24:59