在异步生成器的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
这个版本解决了所有问题:
pull_task和receive_channel的处理并行执行,不会出现无限等待的情况pull_task结束后会关闭发送端,确保接收端循环能正常终止- 如果迭代提前终止(比如调用
aclose()),nursery会自动取消pull_task,避免资源泄漏
内容的提问来源于stack exchange,提问作者Anders E. Andersen
相关产品推荐
相关产品推荐

