寻找AsyncGenerator的asyncio.create_task等效方案:后台并行执行
让AsyncGenerator并行执行的实现方案
要让异步生成器像asyncio.create_task那样在后台并行执行,核心思路是用队列作为中间缓存,让生成器在后台任务中持续产出元素,同时主逻辑从队列中实时获取结果。这样既实现了并行,又无需提前等待所有结果生成。
实现代码
import asyncio import time async def g(): for i in range(3): await asyncio.sleep(1.0) yield i # 包装生成器,将产出的元素推入队列 async def _gen_to_queue(gen, queue): try: async for item in gen: await queue.put((gen, item)) finally: # 用None标记当前生成器执行结束 await queue.put((gen, None)) async def main(): g1 = g() g2 = g() result_queue = asyncio.Queue() t = time.time() # 启动后台任务,让两个生成器并行运行 tasks = [ asyncio.create_task(_gen_to_queue(g1, result_queue)), asyncio.create_task(_gen_to_queue(g2, result_queue)) ] active_generators = {g1, g2} while active_generators: gen, item = await result_queue.get() if item is None: active_generators.remove(gen) continue # 实时处理生成器产出的元素 print(item, time.time() - t) # 等待所有后台任务收尾 await asyncio.gather(*tasks) asyncio.run(main())
执行效果
这段代码的总耗时约3秒,输出顺序会因并行调度略有差异,示例输出如下:
0 1.0011980533600 0 1.0012540817261 1 2.0023450851440 1 2.0024039745331 2 3.0034260749817 2 3.0034830570221
方案说明
- 队列缓存:通过
asyncio.Queue实现生成器与主逻辑的解耦,生成器在后台产出元素后直接推入队列,主逻辑无需等待生成器执行完毕。 - 后台任务:用
asyncio.create_task启动包装生成器的协程,让多个生成器同时运行。 - 结束标记:在生成器执行完毕后推入
None,主逻辑通过这个标记跟踪所有生成器的执行状态,确保不会遗漏结果。
如果不需要实时处理元素,也可以先收集每个生成器的所有结果再统一处理,这种方式代码更简洁,但无法实时获取产出:
async def collect_gen(gen): return [item async for item in gen] async def main(): g1 = g() g2 = g() t = time.time() # 并行收集两个生成器的所有结果 results1, results2 = await asyncio.gather(collect_gen(g1), collect_gen(g2)) for i in results1: print(i, time.time() - t) for i in results2: print(i, time.time() - t) asyncio.run(main())
这种方式总耗时同样约3秒,但所有结果会在生成器全部执行完毕后才开始打印。
内容的提问来源于stack exchange,提问作者Piotr Dabkowski
相关产品推荐
相关产品推荐

