Skip to content

流式输出:astream 与 stream_mode

源码版本1.2.9

职责

Pregel.astream / Pregel.stream 是编译图对外暴露的执行入口。你把输入丢进来,它启动 PregelLoop 一圈圈跑,过程中把每一超步 (superstep) 产生的数据按调用方指定的 stream_mode 投影成流事件,塞进一个队列,再由调用方迭代消费。ainvoke / invoke 本质上就是 astream / stream 收敛成单值——Pregel.invoke 的实现体就是 for chunk in self.stream(...) 取最后一片(Pregel.invoke body:3891)。

这一层的责任是把「执行」和「输出形状」解耦:执行还是 Pregel 那套 BSP 循环,输出形状由 stream_mode 决定——同一个图跑一次,你可以要 values(每步完整状态)、updates(每步增量)、messages(LLM token 流)、custom(节点内自定义输出)、checkpoints(检查点事件)、tasks(任务起止)、debug(全部)等不同投影,甚至传一个 list 同时拿多种模式。

设计动机

  • 不阻塞执行:输出走队列,runner.tick / atick 每完成一个任务就往里塞事件,外层 for 循环边迭代边拿,而不是跑完整张图再一次性返回。这对 LLM token 流(messages)是刚需——用户得边生成边看。
  • 多投影复用同一份执行:同一个 runner.atick(...) 内部产生多种事件,_output 函数(_output:4184)只做过滤/分发的动作,不重新跑图。传 stream_mode=["values","updates"] 不会跑两次。
  • 子图命名空间统一:开 subgraphs=True 后,事件带 namespace 前缀((ns, mode, payload)),父图能看见子图内部事件,靠 _output 在出队时打标签,而不是子图自己额外开一条流。
  • v2 类型化:新版本通过 version="v2" 把事件变成 {"type": mode, "ns": ..., "data": ..., "interrupts": ...} 字典,便于下游程序化处理;values 模式还会把 __interrupt__ 单独抽出成 interrupts 字段。

关键文件

  • Pregel class:450Pregel 类定义,所有流式入口都在这里。
  • Pregel.stream:2655 — 同步流式入口,签名里列出全部 stream_mode 选项。
  • Pregel.astream:3063 — 异步流式入口,真正的异步主循环写在它的实现体里。
  • Pregel.invoke:3836 — 同步收敛入口,内部 for chunk in self.stream(...) 取最后一个 values chunk。
  • sync main loop:2964 — 同步主循环:while loop.tick()runner.tickloop.after_tick()
  • async main loop:3437 — 异步主循环:三段式同上,runner.atick 替换同步版。
  • _output:4184 — 出队过滤函数,按 stream_mode / print_mode / subgraphs 决定每个事件怎么 yield。
  • GraphRunStream:31 — 同步调用方驱动的流包装,for 循环就是 pump,没有后台线程。
  • AsyncGraphRunStream:304 — 异步对应物,多个投影(run.values / run.messages)共用一个 pump。
  • stream_mode attr:709Pregel 默认 stream_mode="values",可被 stream(stream_mode=...) 覆盖。

数据流

异步主循环骨架在 astream 实现体里(async main loop:3437),三段式和 Pregel 引擎页讲的一样,差别在每段中间夹了 _output 把队列里的东西吐给调用方:

python
while loop.tick():
    for task in await loop.amatch_cached_writes():
        loop.output_writes(task.id, task.writes, cached=True)
    async for _ in runner.atick(
        [t for t in loop.tasks.values() if not t.writes],
        timeout=self.step_timeout,
        get_waiter=get_waiter,
        schedule_task=loop.aaccept_push,
    ):
        # emit output
        for o in _output(
            stream_mode,
            print_mode,
            subgraphs,
            stream.get_nowait,
            asyncio.QueueEmpty,
            version,
            _output_mapper,
            _state_mapper,
        ):
            yield o
    loop.after_tick()
    await aemit_graph_lifecycle_events(loop)
    # wait for checkpoint
    if durability_ == "sync":
        await cast(asyncio.Future, loop._put_checkpoint_fut)

runner.atick 内部把每个任务产生的 writes 写回 loop,同时把 (ns, mode, payload) 三元组塞进 stream(asyncio.Queue);外层 async for 拿到执行权时,_outputstream.get_nowait() 把队列里所有就绪事件抽干,按 stream_mode 过滤后再 yield 给用户。_output 的核心逻辑(_output:4184):拿到 (ns, mode, payload),如果 mode in print_modeprint 一份,如果 mode in stream_mode 才向外 yield——所以 print_mode 是「只看不发」的调试开关,不影响真正 yield 出去的内容。

下面这张图把 stream 调用从入口到事件出队走一遍:

边界与失败

  • 默认 stream_mode:stream_mode=None 时,如果是被作为子图调用(CONFIG_KEY_TASK_ID 在 config 里),默认 values,否则用 self.stream_mode(stream_mode default:2740),避免被子图默认 updates 模式污染父图输出。
  • 步数超限:主循环退出后看 loop.status,如果是 out_of_stepsGraphRecursionError,提示调高 recursion_limit(out_of_steps:3002);draining 状态抛 GraphDrained,把控制权交还给外部 RunControl
  • 多模式同时输出:stream_mode 传 list 时,每个事件变成 (mode, payload) 元组;再叠 subgraphs=True 则是 (ns, mode, payload)(tuple mode:4240),调用方要按元组解构。
  • messages 模式的继承处理:stream_mode="messages"version="v1" 时,会剥掉继承自上层的 v2 messages handler,避免 v1 流被路由到 content-block 事件协议(strip v2 handler:2773),但保留 v1 handler 以支持 subgraphs=True 时内层 messages 事件被外层观察。
  • durability 与 checkpoint 同步:durability="sync" 时主循环每步结束都 await loop._put_checkpoint_fut,保证 checkpoint 落盘后才进入下一步;"async" 是默认值,持久化和下一步并行;"exit" 只在退出时落盘。
  • v3 实验性 stream_events:stream_events(version="v3") 不接受 stream_modesubgraphs 参数,这两个被 mux 内部接管,显式传会被 _reject_v3_invariant_kwargs 拒掉(_reject_v3_invariant_kwargs:387)。

小结

astream / stream 把 Pregel 的 BSP 循环包成可迭代接口,stream_mode 选择决定输出投影,_output 是出队过滤的单一咽喉。要追执行细节看 /pregel/pregel/pregel/loop,要追多模式 mux 看 libs/langgraph/langgraph/stream/_mux.py。对照官方资料:LangGraph 文档 · README