如何基于asyncio实现可控最大并发数的动态任务补位执行器
asyncio 自定义限流动态补位执行器实现方案
方案一:基于asyncio.wait实现(最简适配)
该方案直接解决原有as_completed实现的缺陷,逻辑直观无额外依赖,完全符合运行规则:
import asyncio # 模拟业务协程,可替换为你的实际逻辑 async def job(param): await asyncio.sleep(1) print(f"执行完成,参数:{param}") return param async def custom_executor(async_level: int, params: list): # 转迭代器,大参数列表也不会额外占用内存,取完自动抛出StopIteration param_iter = iter(params) running_tasks = set() # 初始填充并发池到最大限制 for _ in range(async_level): try: param = next(param_iter) running_tasks.add(asyncio.create_task(job(param))) except StopIteration: # 参数数量小于并发数的场景直接跳出 break # 循环处理直到所有任务跑完 while running_tasks: # 等待任意一个任务完成,无需等待整批结束 done, running_tasks = await asyncio.wait( running_tasks, return_when=asyncio.FIRST_COMPLETED ) # 处理已完成的任务,可在此处捕获异常、收集返回结果 for task in done: # 示例:打印返回结果,异常可通过try-except task.result()捕获 print(f"任务返回结果:{task.result()}") # 补入新任务 try: new_param = next(param_iter) running_tasks.add(asyncio.create_task(job(new_param))) except StopIteration: # 无剩余参数,无需补位 pass # 测试运行 if __name__ == "__main__": asyncio.run(custom_executor(async_level=4, params=[i for i in range(1, 11)]))
逻辑说明
- 完全匹配三条运行规则:参数为空直接终止,每次取新参数创建协程,同一时间最多
async_level个协程运行,任务完成立即补位 - 不会预先创建所有协程对象,同时存在的协程数最高等于
async_level,无参数量大时的资源过载问题 - 可灵活扩展异常处理、结果收集等自定义逻辑
方案二:基于asyncio.Queue实现(消费者模式)
适合需要动态追加参数、多生产者场景的更灵活方案:
import asyncio async def job(param): await asyncio.sleep(1) print(f"执行完成,参数:{param}") return param # 消费者协程,数量和最大并发数一一对应 async def consumer(queue: asyncio.Queue): while True: param = await queue.get() # 收到结束信号就退出 if param is None: queue.task_done() break try: await job(param) finally: queue.task_done() async def custom_executor(async_level: int, params: list): queue = asyncio.Queue(maxsize=async_level) # 启动对应数量的消费者 consumers = [asyncio.create_task(consumer(queue)) for _ in range(async_level)] # 往队列塞参数,队列满会自动阻塞,不会一次性加载所有参数 for param in params: await queue.put(param) # 塞结束信号,通知消费者退出 for _ in range(async_level): await queue.put(None) # 等待所有任务执行完成 await queue.join() # 等待所有消费者退出 await asyncio.gather(*consumers) # 测试运行 if __name__ == "__main__": asyncio.run(custom_executor(async_level=4, params=[i for i in range(1, 11)]))
内容的提问来源于stack exchange,提问作者Evgen
相关产品推荐
相关产品推荐

