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

在异步生成器的nursery内部yield是否不安全?求问题复现代码

让我来帮你拆解这个问题,并且用实际例子展示这段代码的问题所在。

核心问题:你的代码逻辑存在阶段割裂

这段代码的关键问题在于:async with trio.open_nursery()块会等待pull_task完全结束之后,才会执行后面的receive_channel处理逻辑。这直接导致了两种典型的异常场景:


场景1:处理无限流时完全无法产出任何元素

假设我们有一个持续生成数据的无限异步数据源:

import trio
import math
from typing import AsyncIterable, TypeVar, aclosing

T = TypeVar('T')

async def infinite_source() -> AsyncIterable[int]:
    i = 0
    while True:
        yield i
        i += 1
        await trio.sleep(1)

async def delay(interval: float, source: AsyncIterable[T]) -> AsyncIterable[T]:
    send_channel, receive_channel = trio.open_memory_channel(math.inf)
    async def pull_task():
        async with aclosing(source) as agen:
            async for item in agen:
                send_channel.send_nowait((item, trio.current_time() + interval))
    async with trio.open_nursery() as nursery:
        nursery.start_soon(pull_task)
    async with receive_channel:
        async for item, timestamp in receive_channel:
            now = trio.current_time()
            if timestamp > now:
                await trio.sleep(timestamp - now)
            yield item

async def main():
    async for num in delay(2, infinite_source()):
        print(f"Received delayed num: {num}")

if __name__ == "__main__":
    trio.run(main)

运行这段代码你会发现:控制台完全没有输出。
原因很直白:pull_task在处理infinite_source()时会无限循环,async with trio.open_nursery()会一直等待pull_task完成,永远不会进入后面的receive_channel处理逻辑,自然也就不会yield任何元素。


场景2:处理有限流时会在最后卡住

再看一个有限数据源的例子:

async def finite_source() -> AsyncIterable[int]:
    for i in range(3):
        yield i
        await trio.sleep(0.5)

async def main():
    print("Starting to iterate delayed source...")
    async for num in delay(1, finite_source()):
        print(f"Received delayed num: {num}")
    print("Iteration finished!")

if __name__ == "__main__":
    trio.run(main)

运行后你会看到前3个元素被正常输出,但程序会卡住,永远不会打印Iteration finished!。
这是因为:

  • pull_task完成后,我们没有关闭send_channel的发送端
  • receive_channel的async for会一直等待新的数据,直到所有发送端都被关闭
  • 而send_channel还在delay函数的作用域内存在,没有被关闭,所以循环永远不会终止

为什么这违反了Trio的异步生成器规范?

Trio的设计规范明确提到:在异步生成器中,不能在nursery关闭后再执行yield操作。因为这会把生成器的执行拆成两个完全独立的阶段(nursery运行阶段和yield产出阶段),两者之间没有同步机制,很容易引发资源泄漏、任务无法取消、数据流中断等问题。

修复后的正确版本

要实现“无背压延迟异步序列”的目标,应该把pull_task的运行和receive_channel的处理放在同一个nursery上下文里,同时确保任务结束或生成器被取消时能正确清理资源:

async def delay(interval: float, source: AsyncIterable[T]) -> AsyncIterable[T]:
    send_channel, receive_channel = trio.open_memory_channel(math.inf)
    
    async def pull_task():
        try:
            async with aclosing(source) as agen:
                async for item in agen:
                    send_channel.send_nowait((item, trio.current_time() + interval))
        finally:
            await send_channel.aclose()  # 任务结束后主动关闭发送端
    
    async with trio.open_nursery() as nursery:
        nursery.start_soon(pull_task)
        async with receive_channel:
            async for item, timestamp in receive_channel:
                now = trio.current_time()
                if timestamp > now:
                    await trio.sleep(timestamp - now)
                yield item
        # 当receive_channel循环结束/被取消,nursery会自动取消pull_task

这个版本解决了所有问题:

  1. pull_task和receive_channel的处理并行执行,不会出现无限等待的情况
  2. pull_task结束后会关闭发送端,确保接收端循环能正常终止
  3. 如果迭代提前终止(比如调用aclose()),nursery会自动取消pull_task,避免资源泄漏

内容的提问来源于stack exchange,提问作者Anders E. Andersen

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.08 19:33:09