You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何让异步生成器/迭代器在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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.14 17:45:38