使用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
相关产品推荐
相关产品推荐

