如何并行链式多个异步生成器/迭代器?是否存在异步版itertools.chain?
并行处理异步迭代器并实时获取数据
你问到的这个需求很常见,但确实没有直接对应itertools.chain的异步工具——因为itertools.chain是顺序遍历迭代器,而你需要的是并行监听多个异步迭代器,一旦有数据就立即返回,还要能追踪数据来源。
下面我给你一个实用的实现方案,完全满足你的核心需求:
核心实现代码
import asyncio from typing import AsyncIterable, Tuple, Any async def chain_parallel(*async_iterables: AsyncIterable[Any]) -> AsyncIterable[Tuple[AsyncIterable[Any], Any]]: # 初始化任务集合:为每个异步迭代器创建第一个迭代任务 tasks = { asyncio.create_task(ait.__anext__(), name=repr(ait)): ait for ait in async_iterables } while tasks: # 等待任意一个迭代任务完成(谁先有数据就先处理谁) done, pending = await asyncio.wait(tasks, return_when=asyncio.FIRST_COMPLETED) for completed_task in done: source_iter = tasks.pop(completed_task) try: # 获取迭代结果 value = completed_task.result() # 返回(来源迭代器,数据值),方便你识别哪个生成器产出的数据 yield source_iter, value # 给这个迭代器创建下一个迭代任务,继续监听 next_task = asyncio.create_task(source_iter.__anext__(), name=repr(source_iter)) tasks[next_task] = source_iter except StopAsyncIteration: # 该迭代器已经耗尽,不再添加新任务 continue # 测试示例 async def async_gen1(): for i in range(3): await asyncio.sleep(0.5) # 模拟较慢的异步生成 yield f"gen1: {i}" async def async_gen2(): for i in range(3): await asyncio.sleep(0.3) # 模拟较快的异步生成 yield f"gen2: {i}" async def main(): async for source, value in chain_parallel(async_gen1(), async_gen2()): print(f"Received from {repr(source)}: {value}") if __name__ == "__main__": asyncio.run(main())
代码说明
- 并行监听逻辑:我们给每个异步迭代器的下一次迭代创建一个任务,然后用
asyncio.wait的FIRST_COMPLETED模式等待——只要有一个迭代器产出数据,就立即处理,不用等其他迭代器。 - 来源识别:返回的是
(来源迭代器, 数据值)的元组,你可以通过repr(source)或者自定义标识来区分不同的生成器。 - 迭代器耗尽处理:当某个迭代器抛出
StopAsyncIteration时,我们就不再给它创建新任务,直到所有任务都处理完毕,函数退出。 - 兼容性:如果你的Python版本低于3.10,不能用
anext()函数,代码里用的__anext__()是通用写法,所有支持异步迭代的版本都能运行。
扩展调整
- 如果不需要识别来源,只需要数据本身,把
yield source_iter, value改成yield value即可。 - 如果需要限制同时处理的任务数量(比如防止任务过多),可以在添加新任务时做数量判断,不过大部分场景下这个基础实现已经足够。
内容的提问来源于stack exchange,提问作者Chris
相关产品推荐
相关产品推荐

