如何异步扁平化多个异步生成器,实现结果即时获取?
我来帮你理清楚问题所在,然后给出完全符合需求的解决方案~
先分析你遇到的问题根源
第一个
run_test的同步陷阱
你最开始的代码是逐个迭代每个test(i):for i in range(10): async for j in test(i): ls.append(j)这相当于先把
test(0)的所有元素处理完(等0.5秒+迭代10次),再轮到test(1),以此类推,10个test加起来就是5秒,完全没用到异步并行的特性,结果自然是固定的同步顺序。asyncio.as_completed/gather报错的原因
这两个API要求传入可等待对象(asyncio.Future、协程或其他awaitable类型),但你传入的test(i)是异步生成器——它不能被await,只能用async for迭代,所以会抛出TypeError: An asyncio.Future, a coroutine or an awaitable is required。
符合需求的解决方案
你的核心诉求是:让所有异步生成器并行运行,一旦某个生成器产出元素就立刻收集,不需要等整个生成器完成,结果顺序可以随机。
这里提供两种实用的实现方式:
方法1:任务包装迭代(简单直接)
把每个异步生成器的迭代逻辑包装成协程,用asyncio.create_task启动所有任务并行运行,共享一个结果列表(asyncio是单线程环境,列表append操作是原子的,不会有竞争问题):
import asyncio async def test(n): await asyncio.sleep(0.5) for i in range(1, 11): yield n * i # 包装异步生成器的迭代逻辑,将元素添加到结果列表 async def process_generator(gen, result_list): async for item in gen: result_list.append(item) async def run_test(): result = [] tasks = [] # 为每个异步生成器创建并行任务 for i in range(10): gen = test(i) task = asyncio.create_task(process_generator(gen, result)) tasks.append(task) # 等待所有任务完成 await asyncio.gather(*tasks) return result # 运行测试 output = asyncio.run(run_test()) print(output)
效果:总耗时约0.5秒(所有test并行sleep 0.5秒),结果是各个生成器产出元素的交错排列,每次运行顺序可能不同,完全满足“即时收集”的需求。
方法2:队列实现实时处理(灵活扩展)
如果需要对产出的元素做实时处理(比如一拿到就打印),可以用asyncio.Queue传递元素,主协程实时消费队列中的内容:
import asyncio async def test(n): await asyncio.sleep(0.5) for i in range(1, 11): yield n * i async def run_test(): result = [] queue = asyncio.Queue() task_count = 10 # 每个任务负责迭代生成器,将元素放入队列,完成后发送结束标记 async def process_gen(gen): async for item in gen: await queue.put(item) await queue.put(None) # 用None标记该生成器处理完成 # 启动所有并行任务 tasks = [asyncio.create_task(process_gen(test(i))) for i in range(task_count)] # 从队列取元素,直到所有生成器都完成 completed = 0 while completed < task_count: item = await queue.get() if item is None: completed += 1 else: result.append(item) # 这里可以添加实时处理逻辑,比如: # print(f"实时收到元素: {item}") # 等待所有任务彻底结束 await asyncio.gather(*tasks) return result output = asyncio.run(run_test()) print(output)
效果:同样是0.5秒左右完成,而且能在元素产出的第一时间进行处理,不需要等所有任务结束。
结果验证
运行上述代码你会发现:
- 总耗时只有0.5秒,证明所有
test是并行运行的; - 每次运行的结果顺序可能不同,因为元素是一产出就被收集的,完全符合你想要的“即时获取”效果。
备注:内容来源于stack exchange,提问作者Ξένη Γήινος

