Python异步编程:如何处理多任务提交并限制最大2个并行任务?
异步任务队列实现(限制最大并发数为2)
问题根源
你的核心问题有两个:一是CPU密集型任务会阻塞asyncio事件循环,导致后续API请求被卡住;二是任务提交逻辑没有做到真正的异步非阻塞,导致提交第二个任务时必须等待第一个完成。
解决方案思路
- 用
asyncio.Queue存储待执行任务,通过固定数量的工作协程控制并发数(设为2) - 将CPU密集计算任务放到线程池/进程池执行,避免阻塞事件循环
- 提交任务的API接口异步处理,加入队列后立即返回响应
代码实现示例(基于FastAPI)
import asyncio from concurrent.futures import ThreadPoolExecutor from fastapi import FastAPI app = FastAPI() # 初始化异步任务队列 task_queue = asyncio.Queue() MAX_CONCURRENT_TASKS = 2 # 线程池大小根据CPU核心数调整,CPU密集型任务建议设为核心数*2 executor = ThreadPoolExecutor(max_workers=4) # 工作协程:负责从队列取任务并执行 async def task_worker(): while True: task_id, compute_func, args = await task_queue.get() try: # 把CPU密集任务委托给线程池执行,不阻塞事件循环 await asyncio.get_event_loop().run_in_executor(executor, compute_func, *args) print(f"任务 {task_id} 执行完成") except Exception as e: print(f"任务 {task_id} 执行失败: {str(e)}") finally: task_queue.task_done() # 启动指定数量的工作协程 @app.on_event("startup") async def start_workers(): for _ in range(MAX_CONCURRENT_TASKS): asyncio.create_task(task_worker()) # 任务提交API @app.post("/submit/{task_id}") async def submit_task(task_id: int): # 模拟CPU密集型计算任务 def heavy_computation(task_id): import time # 替换为实际的计算逻辑 for _ in range(10**7): pass print(f"任务 {task_id} 正在计算") # 将任务加入队列,立即返回响应 await task_queue.put((task_id, heavy_computation, (task_id,))) return {"status": "queued", "task_id": task_id, "message": "任务已加入等待队列"}
关键细节说明
- 线程池/进程池选择:CPU密集型任务优先用
ProcessPoolExecutor(规避GIL限制),IO密集型用ThreadPoolExecutor即可 - 并发数控制:启动的工作协程数量等于
MAX_CONCURRENT_TASKS,保证同时最多处理2个任务 - 非阻塞提交:API接口在把任务放入队列后立刻返回,不会等待任务执行,解决你提交第二个任务时无法立即得到响应的问题
- 异常处理:工作协程里捕获任务执行异常,避免单个任务失败导致整个worker崩溃
额外建议
- 可以给队列设置
maxsize参数,防止任务堆积过多耗尽内存 - 增加任务状态查询接口,方便追踪任务执行情况
- 若需要持久化任务队列(比如服务重启后不丢失任务),可以用Redis等外部队列替代
asyncio.Queue
内容的提问来源于stack exchange,提问作者Lisa
相关产品推荐
相关产品推荐

