流式输出:astream 与 stream_mode
职责
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:450—Pregel类定义,所有流式入口都在这里。Pregel.stream:2655— 同步流式入口,签名里列出全部stream_mode选项。Pregel.astream:3063— 异步流式入口,真正的异步主循环写在它的实现体里。Pregel.invoke:3836— 同步收敛入口,内部for chunk in self.stream(...)取最后一个valueschunk。sync main loop:2964— 同步主循环:while loop.tick()→runner.tick→loop.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:709—Pregel默认stream_mode="values",可被stream(stream_mode=...)覆盖。
数据流
异步主循环骨架在 astream 实现体里(async main loop:3437),三段式和 Pregel 引擎页讲的一样,差别在每段中间夹了 _output 把队列里的东西吐给调用方:
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 拿到执行权时,_output 调 stream.get_nowait() 把队列里所有就绪事件抽干,按 stream_mode 过滤后再 yield 给用户。_output 的核心逻辑(_output:4184):拿到 (ns, mode, payload),如果 mode in print_mode 先 print 一份,如果 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_steps抛GraphRecursionError,提示调高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_mode和subgraphs参数,这两个被 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。