Python asyncio:如何惰性链式处理as_completed()返回的结果?
嘿,这个需求用异步生成器就能完美解决!我给你写个完整的实现,再拆解下思路:
解决方案
import asyncio as aio import random async def produce_numbers(x): await aio.sleep(random.uniform(0, 3.0)) return [x * x, x * x + 1, x * x + 2] async def flatten_completed(future_generator): # 遍历as_completed返回的已完成Future对象 for fut in future_generator: # 等待拿到当前协程的结果列表 result_list = await fut # 逐个yield列表中的元素,实现惰性返回 for item in result_list: yield item async def main(): # Step 1) 创建未等待的协程列表。 coros = [produce_numbers(i) for i in range(0, 10)] # Step 2) 创建按完成顺序返回结果的生成器。 futs = aio.as_completed(coros) # type: generator # Step 3) 创建惰性生成器,逐个返回单个元素 single_item_gen = flatten_completed(futs) # 测试迭代这个生成器 async for num in single_item_gen: print(f"拿到单个元素: {num}") if __name__ == "__main__": aio.run(main())
思路拆解
- 为什么用异步生成器?因为
as_completed返回的是包含Future对象的生成器,我们需要用await来获取每个Future的结果,而普通生成器没法处理异步操作,所以必须用async def定义的异步生成器,搭配async for来迭代。 flatten_completed函数的核心逻辑:- 遍历
as_completed给出的每个已完成Future; - await拿到该协程返回的列表;
- 逐个yield列表中的元素——这就实现了惰性:只有当你迭代生成器时,才会处理下一个完成的协程结果,并且不会一次性把所有结果都加载到内存里。
- 遍历
- 在
main函数里,用async for遍历这个生成器,就能按协程完成的顺序,逐个拿到每个协程返回列表里的单个元素了。
内容的提问来源于stack exchange,提问作者Rotareti
相关产品推荐
相关产品推荐

