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

