如何让异步生成器/迭代器在trio中实现真正异步并行运行?
异步生成器并行执行的实现方案
你的测试代码里,async for会逐个迭代异步生成器的产出值——每次只有当前一个await trio.sleep(...)完成、生成器yield出值后,才会进入下一次迭代。所以本质上还是串行执行,所有任务不会同时触发。
方案一:用Trio Nursery直接并行执行任务
这是实现“所有任务立即触发、并行执行”最直接的方式,不需要依赖异步生成器:
import trio import random async def do_task(i): await trio.sleep(random.random()) print(i) async def main(): async with trio.Nursery() as nursery: for i in range(10): nursery.start_soon(do_task, i) trio.run(main)
这段代码会同时启动10个任务,每个任务独立sleep随机时间,完成后立即打印结果,输出顺序是随机的,总耗时等于最长的那个sleep时间(必然小于1秒)。
方案二:结合异步生成器与队列收集并行结果
如果需要用异步生成器来迭代并行任务的结果,可以借助Trio的MemorySendChannel实现——任务完成后把结果发送到通道,异步生成器从通道接收并产出值:
import trio import random async def task_producer(sender, num_tasks): async with trio.Nursery() as nursery: for i in range(num_tasks): async def task(i=i): await trio.sleep(random.random()) await sender.send(i) nursery.start_soon(task) async def result_generator(num_tasks): async with trio.open_memory_channel(0) as (sender, receiver): # 启动生产者任务,并行执行所有任务 await trio.lowlevel.spawn_system_task(task_producer, sender, num_tasks) # 从通道接收结果,作为生成器产出值 async for result in receiver: yield result async def main(): async for idx in result_generator(10): print(idx) trio.run(main)
这里result_generator是异步生成器,内部启动nursery并行执行所有任务,任务结果通过通道发送,生成器再逐个产出。虽然async for还是逐个处理结果,但任务本身是并行执行的,总耗时同样小于1秒,输出顺序随机。
关于BigQuery Async Storage Write API的补充
该API中的异步生成器通常用于流式处理待写入数据(比如逐个读取数据源),而非并行执行写入操作。要实现并行写入,你需要用Trio Nursery启动多个写入任务,每个任务处理一部分数据,避免依赖异步生成器的串行迭代逻辑。
内容的提问来源于stack exchange,提问作者Mike Williamson
相关产品推荐
相关产品推荐

