12811 字
64 分钟

LangGraph 源码深潜:Streaming、Events 与 Observability 机制解剖

核心问题: LangChain / LangGraph 如何把一次 Agent / Workflow 执行过程拆成可观察的状态更新、模型 token、节点事件、自定义事件和外部 trace?

源码主线: graph.stream()Pregel.stream()StreamProtocolstream_mode 分发 → values / updates / messages / custom / eventscallbacks / tracersLangSmith

前置文章: 第 1 篇 Runnable、第 2 篇 Prompt / Message / ChatModel、第 6 篇 StateGraph、第 8 篇 Reducer、第 11 篇 Checkpoint、第 12 篇 Interrupt / HITL

依赖基线: langgraph==1.2.7langchain-core==1.4.8langchain==1.3.11

源码基线: https://github.com/langchain-ai/langgraph/tree/1.2.7https://github.com/langchain-ai/langchain,以当前正式发布版本对应源码为准

阅读边界: 本文分析 LangGraph 本地进程内 stream / astream / stream_events / astream_events、LangChain Runnable callbacks、metadata、LangSmith trace 接入;不展开 LangGraph Platform 远程 Streaming API、前端 SSE/WebSocket 实现、LangSmith 后端存储和 UI 内部实现。


0. 本篇在源码学习主线中的位置#

前面的文章已经解释了:

Runnable
解释了为什么组件可以 invoke / stream / batch。
StateGraph
解释了 Workflow 如何被编译成可执行图。
Conditional Edge
解释了路由和分支如何在运行时发生。
Reducer
解释了节点 update 如何合并进 state。
Checkpoint
解释了状态如何持久化。
Interrupt
解释了图如何暂停并恢复。

本篇进入生产调试与可观测性:

Graph 执行
Streaming / Events / Metadata / Callbacks
Agent Trace / Debugging / Observability

本篇只解决:

  • stream()invoke() 的执行路径有什么关系。
  • stream_mode="values"stream_mode="updates" 如何观察 state。
  • stream_mode="messages" 如何观察 LLM token 与 message chunks。
  • stream_mode="custom" 如何从节点或工具中发业务进度事件。
  • event streaming 与低层 stream_mode 的区别。
  • metadata / tags / callbacks 如何随 RunnableConfig 传播。
  • LangSmith trace 如何接入 LangChain / LangGraph 执行链。
  • 旅行规划助手如何设计可观测 trace。

本篇不展开:

  • LangGraph Platform 的远程 agent server streaming protocol。
  • 前端 SSE、WebSocket、反向代理缓冲、浏览器渲染细节。
  • LangSmith 服务端 trace 存储模型。
  • OpenTelemetry GenAI 语义约定的完整映射。
  • 第三方监控平台接入实现。

1. 本篇问题、学习目标与能力边界#

1.1 核心问题#

如何从 LangGraph 的一次黑盒执行中观察每个节点、每次模型调用、每个 token、每次工具进度和完整 trace?

1.2 学习目标#

完成本篇后,读者必须能够:

  1. 解释 invoke()stream() 在结果形态上的区别,以及为什么 invoke() 可以理解为对 streaming 执行结果的聚合。
  2. 区分 stream_mode="values""updates""messages""custom""checkpoints""tasks""debug" 的观察对象。
  3. 解释 stream_mode 如何影响 Pregel runtime 向外发出的 chunk。
  4. 解释 token streaming 与 node update streaming 的本质区别。
  5. 解释 get_stream_writer() 如何把节点内部业务进度写入 custom stream。
  6. 解释 stream_events() / astream_events() 为什么比低层 stream_mode 更适合应用代码。
  7. 解释 metadatatagscallbacks 如何通过 RunnableConfig 传递到子调用。
  8. 解释 LangSmith trace 如何接入,以及它与本地 stream 的边界。
  9. 能为旅行规划助手设计一套调试与线上观测方案。

1.3 能力边界#

能力本篇是否覆盖说明
graph.stream()覆盖 sync streaming 与 stream modes
graph.astream()覆盖 async streaming 边界
stream_mode="values"观察每步后的完整 state
stream_mode="updates"观察每个节点返回的 partial update
stream_mode="messages"观察模型 token / message chunk
stream_mode="custom"观察节点或工具自定义进度
event streaming覆盖 typed projections 和 stream_events()
callbacks覆盖 RunnableConfig 中的 callbacks 传播
metadata / tags覆盖过滤、trace、debug 用法
LangSmith trace覆盖启用方式与边界
OpenTelemetry只作为后续方向,不展开
前端流式传输不展开 SSE / WebSocket
LangGraph Platform server streaming不展开远程部署 API

2. 核心概念与最小心智模型#

2.1 一句话定义#

Streaming 是 LangGraph 把图运行时内部事件投影为外部可消费输出的机制,负责观察执行过程,不负责改变图的业务逻辑。

2.2 最小心智模型#

graph.invoke(input)
执行完整图
返回最终 state
graph.stream(input, stream_mode=...)
执行同一张图
边执行边产出观察事件
最后也能得到最终结果或最后一次 state
graph.stream_events(input)
基于底层 stream 事件
转换成 typed projections
messages / values / interrupts / output / subgraphs

2.3 核心术语#

术语源码对象语义不要误解为
流式执行Pregel.stream()图执行过程中逐步 yield chunk另一个完全不同的执行器
异步流式执行Pregel.astream()async iterator 版本自动并发所有节点
状态快照流stream_mode="values"每步后的完整 state节点返回值
状态增量流stream_mode="updates"每个节点返回的 partial update完整 state
消息流stream_mode="messages"LLM token / message chunk + metadata节点 update
自定义流stream_mode="custom"节点或工具主动写出的业务事件callback 事件
事件流stream_events()typed projections API简单 chunk wrapper
回调callbacksLangChain Runnable 生命周期 hooksLangGraph branch
元数据metadatatrace 和过滤用辅助信息业务 state
标签tags调用链标记和过滤节点名
TraceLangSmith run tree执行步骤记录streaming output

2.4 与相邻抽象的边界#

对象负责什么不负责什么与本篇对象的关系
invoke()返回最终结果不逐步暴露过程stream() 使用同一图运行语义
stream()低层 chunk 输出不提供 typed projection适合调试 runtime 事件
stream_events()typed projections不替代 trace 存储更适合应用代码
callbacks记录 Runnable 生命周期不直接改变 state可与 LangSmith trace 协作
LangSmithtrace、debug、monitor、eval不负责业务状态合并线上观测平台
Checkpointer保存 thread state不输出 token与 checkpoints stream / get_state 协作
get_stream_writer()发 custom 数据不修改 state适合工具进度和业务日志

3. 完整执行链路#

3.1 高层链路#

用户输入
graph.stream(input, stream_mode="updates")
Pregel runtime 执行 super-step
节点返回 partial update
runtime 应用 update 到 state
stream 输出该节点 update
继续下一 super-step
到达 END

3.2 最小代码骨架#

观察目标:最小代码只展示 updates 模式如何观察每个节点返回的增量更新。

for chunk in graph.stream(
{"user_request": "东京 5 天旅行"},
stream_mode="updates",
):
print(chunk)

version="v2" 输出格式中,chunk 通常包含:

{
"type": "updates",
"data": {
"node_name": {
"state_key": "updated_value"
}
},
"ns": (),
}

如果不使用 version="v2",不同版本和模式下返回形态可能不同。生产建议明确传入:

graph.stream(
input,
stream_mode="updates",
version="v2",
)

3.3 常见 stream 模式理解#

模式观察内容适合场景
values每步后的完整 state状态调试、回放对比
updates每个节点的增量更新节点输出调试、定位哪个节点写错
messages模型 token / message chunks + metadataChat UI、模型输出实时展示
custom节点或工具主动发出的自定义数据工具进度、业务日志
checkpointscheckpoint eventsDurable execution 调试
taskstask start / finish / error节点执行耗时、失败定位
debug多种底层调试信息深度排障,不适合直接给用户
events / stream_eventstyped projections应用层多消费者观测

3.4 对象流转#

阶段输入类型核心函数输出类型状态变化
构建图StateGraphcompile()CompiledStateGraph图获得 stream 能力
同步流式调用dictgraph.stream()Iterator[StreamPart]state 随节点更新
异步流式调用dictgraph.astream()AsyncIterator[StreamPart]state 随节点更新
模型调用messagesChatModel.stream/astreamAIMessageChunk不一定写 state
节点返回partial updateapply_writes新 statereducer 合并
自定义事件writer(data)get_stream_writer()custom chunk不修改 state
事件投影raw stream eventsstream_events()typed projections不修改 state
Trace 记录RunnableConfigcallbacks / tracersLangSmith run不修改业务 state

3.5 时序链路#

Caller
│ graph.stream(input, stream_mode=["updates", "messages", "custom"])
CompiledStateGraph / Pregel
│ 初始化 config、stream protocol、callback manager
Pregel Loop
│ 执行 node
Node Function
├─ return {"plan": "..."} ───────────────► updates stream
├─ model.stream(...) ───────────────────► messages stream
└─ writer({"progress": 50}) ────────────► custom stream
Pregel Loop
│ 合并 state,调度下一节点
Caller receives chunks

3.6 正常结束条件#

一次流式执行正常结束的条件是:

1. 图到达 END。
2. 所有已调度 task 完成。
3. stream iterator 被耗尽。
4. 最终 state 可以从最后一次 values、stream.output 或 invoke 聚合结果中获得。

如果图触发 interrupt,流式输出会包含中断信息,执行不会继续到 END,直到外部使用 Command(resume=...) 恢复。


4. 源码地图、关键文件与阅读顺序#

4.1 核心目录#

langgraph/
├── pregel/
│ ├── main.py
│ ├── loop.py
│ ├── runner.py
│ ├── io.py
│ └── messages.py
├── graph/
│ └── state.py
├── config.py
└── types.py
langchain_core/
├── runnables/
│ ├── base.py
│ └── config.py
├── callbacks/
│ └── manager.py
└── tracers/

4.2 关键文件#

优先级文件核心对象阅读目的
1langgraph/pregel/main.pyPregel.stream() / Pregel.astream()理解 stream 主入口
2langgraph/pregel/loop.pyPregel loop理解 super-step 与事件发出
3langgraph/pregel/runner.pytask runner理解节点执行和错误事件
4langgraph/config.pyget_stream_writer()理解 custom stream writer
5langgraph/types.pyStreamPart / Command / Interrupt理解流式输出对象
6langchain_core/runnables/base.pyRunnable.stream() / astream_events()理解 Runnable 级 streaming
7langchain_core/runnables/config.pyRunnableConfig理解 tags、metadata、callbacks
8langchain_core/callbacks/manager.pycallback manager理解 callback 生命周期
9langchain_core/tracers/tracer理解 LangSmith trace 基础

4.3 推荐阅读顺序#

1. 官方 Streaming 文档:先理解 stream modes。
2. `Pregel.stream()`:确认公开参数和输出。
3. `Pregel.astream()`:对比异步路径。
4. `get_stream_writer()`:理解 custom 事件。
5. `stream_events()` 文档:理解 typed projections。
6. `RunnableConfig`:理解 callbacks / metadata / tags 传播。
7. `callbacks/manager.py`:理解 lifecycle event。
8. LangSmith docs:理解 trace 如何启用。

4.4 不建议的阅读顺序#

不建议直接从 langchain_core/tracers 开始。原因是 tracer 只解决“记录到哪里”,不解释“图运行时何时产生什么事件”。

不建议直接从前端 SSE 代码开始。前端只消费 stream,不能解释 stream_mode 的数据来源。

不建议把所有 stream modes 混在一个示例里开始学习。应该先分别理解:

updates:节点输出
values:完整 state
messages:模型 token
custom:业务进度
events:typed projections

5. 对象模型、继承关系与协议边界#

5.1 核心对象关系#

Runnable protocol
CompiledStateGraph / Pregel
invoke / stream / astream
Pregel loop
StreamProtocol / callback manager
StreamPart / typed event projections

5.2 对象职责#

对象生命周期输入输出核心职责
CompiledStateGraph编译后graph inputfinal state / stream chunks暴露 Runnable 接口
Pregel.stream()运行时input + config + stream_modeiterator执行图并 yield chunks
Pregel.astream()运行时input + config + stream_modeasync iterator异步执行并 yield chunks
stream_mode调用时配置string / list[str]影响输出 projection决定观察内容
get_stream_writer()节点运行时custom datacustom stream chunk发送业务进度
RunnableConfig调用时配置callbacks/tags/metadata子调用继承观测上下文传播
CallbackManager运行时lifecycle eventcallback calls触发回调和 tracer
LangSmith tracer运行时callback eventstrace run tree上传可视化 trace

5.3 协议边界#

stream 协议
负责:把执行过程逐步 yield 给调用方。
不负责:存储长期 trace。
callback 协议
负责:在 Runnable 生命周期节点触发 hooks。
不负责:决定 LangGraph 路由。
metadata / tags 协议
负责:携带过滤和 trace 辅助信息。
不负责:作为业务 state 使用。
LangSmith 协议
负责:接收和展示 trace。
不负责:让本地 stream 更快或改变业务结果。

5.4 稳定接口与内部实现#

类型对象文章中的使用原则
公共 APIgraph.stream()可用于工程代码
公共 APIgraph.astream()可用于异步服务
公共 APIgraph.stream_events()推荐应用层 typed streaming
公共 APIget_stream_writer()可用于 custom stream
公共 APIconfig={"tags": ..., "metadata": ...}可用于 trace 标记
扩展接口callbacks可接自定义 callback
扩展接口LangSmith env vars / tracing_context可用于 trace 控制
内部实现Pregel loop event plumbing只解释,不建议业务依赖
内部实现chunk 内部私有字段不建议长期依赖

6. 源码阅读策略与证据标准#

6.1 本篇阅读策略#

先读文档确认 stream modes
读 Pregel.stream 参数和输出
观察 stream_mode 分发
读 state updates 如何变成 values / updates
读 LLM token 如何变成 messages stream
读 get_stream_writer 如何变成 custom stream
读 stream_events typed projection
读 callbacks / metadata / tags 传播
读 LangSmith trace 启用方式

6.2 证据等级#

标记含义写作要求
源码事实当前正式版源码可以证明附源码链接
官方契约官方文档或 API Reference 明确说明附官方链接
简化伪代码压缩真实控制流标注不是源码逐字复制
作者推断根据调用链得出的理解标注“从调用关系可以推断”
工程建议实践选型建议说明适用条件

6.3 本篇证据清单#

结论证据类型文件或文档定位
stream() / astream() 会 yield streamed outputs官方契约LangGraph Streaming docsBasic usage
updates 是每步 state update官方契约LangGraph Streaming docsGraph state
values 是每步完整 state官方契约LangGraph Streaming docsGraph state
messages 是 LLM token + metadata官方契约LangGraph Streaming docsLLM tokens
custom 来自 get_stream_writer()官方契约LangGraph Streaming docsCustom data
event streaming 提供 typed projections官方契约Event streaming docsWhat event streaming provides
metadata/tags 可加入 trace官方契约LangGraph Observability docsAdd metadata to traces
LangSmith trace 通过环境变量启用官方契约LangSmith docsEnable tracing

7. 构建期源码解剖#

本章回答:

为了让图运行过程可观察,构建期需要准备哪些对象和协议?

7.1 构建期职责#

输入归一化动作构建结果
StateGraph编译为 Pregel获得 stream / astream 能力
Node 函数包装为 Pregel node可在运行时产生 update
ChatModel作为 Runnable 子调用可产生 message chunks
get_stream_writer()构建期不执行运行时从 context 读取 writer
callbacks构建期不固定运行时从 config 注入
metadata/tags构建期不固定运行时从 config 或 model config 注入

7.2 构建期总链路#

StateGraph(State)
add_node(...)
add_edge(...)
compile()
CompiledStateGraph / Pregel
继承 Runnable stream / invoke 协议

7.3 StateGraph.compile() 与 streaming 能力来源源码解剖#

职责与所处阶段#

compile() 在构建期把 builder 转成 CompiledStateGraph。Streaming 能力不是单独开启的功能,而是 compiled graph 作为 Runnable / Pregel runtime 的基本能力。

真实源码签名#

以下签名来自当前正式版公开 API,具体泛型以源码为准:

def compile(
self,
checkpointer: Checkpointer = None,
*,
cache: BaseCache | None = None,
store: BaseStore | None = None,
interrupt_before: All | list[str] | None = None,
interrupt_after: All | list[str] | None = None,
debug: bool = False,
name: str | None = None,
) -> CompiledStateGraph:
...

调用方与被调用方#

业务代码
builder.compile()
StateGraph.compile
CompiledStateGraph / Pregel
graph.stream / graph.astream

输入、输出与状态变化#

项目类型说明
输入builder nodes / edges / state schema构建期注册内容
输出CompiledStateGraph运行时可执行对象
状态变化compiled graph 持有节点、边、channel后续 stream 使用
副作用无业务副作用只是构建对象

细粒度伪代码#

以下为保留关键控制流的简化伪代码,不是源码逐字复制:

def compile(self, **options):
# 1. 校验图结构
self.validate()
# 2. 构造 CompiledStateGraph
graph = CompiledStateGraph(
nodes={},
channels=self.channels,
input_channels=self.input_channels,
output_channels=self.output_channels,
stream_channels=self.stream_channels,
checkpointer=options.get("checkpointer"),
store=options.get("store"),
debug=options.get("debug"),
name=options.get("name"),
)
# 3. 附加节点
for node_name, node_spec in self.nodes.items():
graph.attach_node(node_name, node_spec)
# 4. 附加边和分支
for edge in self.edges:
graph.attach_edge(edge)
for branch in self.branches:
graph.attach_branch(branch)
# 5. 返回可执行图
return graph.validate()

逐段解释#

第一段校验 graph,保证 stream 运行时不会遇到不存在的节点或无效边。

第二段构造 compiled graph。这里的 stream_channelsoutput_channelschannels 决定了 valuesupdates 能看到哪些 state key。

第三段和第四段把节点、边、分支转成 Pregel 可调度结构。Streaming 不是节点自己主动把返回值 print 出去,而是 Pregel runtime 在应用 node update 时生成 stream chunk。

第五段返回 compiled graph。这个对象既支持 invoke(),也支持 stream()astream()

正常路径#

StateGraph builder
compile
CompiledStateGraph
stream capable runtime

关键分支与异常路径#

条件行为结果
图结构合法compile 成功可 stream
节点名缺失校验失败构建期报错
checkpointer 缺失不影响 values/updates/messages/custom影响 checkpoints/tasks
debug=True增强调试输出可能增加输出噪声

设计原因与工程影响#

Streaming 能力被放在 compiled runtime,而不是每个 node 自己实现,意味着所有节点天然可以被统一观测,不需要为每个节点写一套日志协议。

源码证据#

7.4 RunnableConfig 中 metadata / tags / callbacks 构建边界源码解剖#

职责与所处阶段#

RunnableConfig 是调用时配置,但它的结构属于 LangChain Core 的稳定协议。它把 callbacks、tags、metadata、run_name 等观测信息传给当前调用和子调用。

真实源码签名#

典型结构:

class RunnableConfig(TypedDict, total=False):
tags: list[str]
metadata: dict[str, Any]
callbacks: Callbacks
run_name: str
configurable: dict[str, Any]
run_id: UUID
max_concurrency: int
recursion_limit: int

调用方与被调用方#

业务代码
graph.stream(input, config={"tags": ..., "metadata": ..., "callbacks": ...})
Pregel runtime
节点内 Runnable 子调用
ChatModel / Tool / Chain
callbacks / tracers

输入、输出与状态变化#

项目类型说明
输入RunnableConfig调用级观测配置
输出merged config传给子 Runnable
状态变化无业务 state 修改只影响观测上下文
副作用callbacks 可能记录日志 / trace不应改变业务结果

细粒度伪代码#

以下为保留关键控制流的简化伪代码,不是源码逐字复制:

def ensure_config(config):
# 1. 从父调用 ContextVar 中读取 inherited config
parent_config = get_current_runnable_config()
# 2. 创建默认 config
merged = {
"tags": [],
"metadata": {},
"callbacks": None,
"configurable": {},
"recursion_limit": DEFAULT_RECURSION_LIMIT,
}
# 3. 合并父 config
if parent_config:
merged["tags"].extend(parent_config.get("tags", []))
merged["metadata"].update(parent_config.get("metadata", {}))
merged["callbacks"] = merge_callbacks(
parent_config.get("callbacks"),
merged.get("callbacks"),
)
# 4. 合并当前 config
if config:
merged["tags"].extend(config.get("tags", []))
merged["metadata"].update(config.get("metadata", {}))
merged["callbacks"] = merge_callbacks(
merged.get("callbacks"),
config.get("callbacks"),
)
merged["configurable"].update(config.get("configurable", {}))
return merged

逐段解释#

第一段读取父级 Runnable config,这解释了为什么模型子调用可以继承 graph 调用传入的 tags 和 metadata。

第二段创建默认值,让没有配置 callbacks 的调用也能正常运行。

第三段合并父级配置。父级 tags 通常用于标识会话、环境、agent 版本。

第四段合并当前配置。当前配置可以增加更细粒度的 node、model、user_id 等 metadata。

正常路径#

config
ensure_config
callbacks / tags / metadata
传递给子调用
trace 可过滤、可分组

关键分支与异常路径#

条件行为结果
没有 callbacks正常运行只是无外部 callback
有 LangSmith tracer记录 trace可视化运行过程
metadata 不可序列化trace 可能失败生产中需 JSON-friendly
tags 过多trace 噪声增加应规范命名
子调用未传 configPython 低版本 async 可能丢 streaming context显式传 config

设计原因与工程影响#

观测信息不应污染业务 state,因此通过 config 传播,而不是写进 state。这也是区分“业务数据”和“观测数据”的关键。

源码证据#

7.5 构建期产物#

产物保存的信息运行时用途
CompiledStateGraph节点、边、channel、stream channel提供 stream / astream
RunnableConfigtags、metadata、callbacks、configurable传递观测上下文
CallbackManagercallback handlers生命周期事件记录
StreamProtocolstream mode、writer、event sink输出 chunks
LangSmith tracertrace sink上传 run tree
get_stream_writer() context当前 stream writer节点内发 custom 数据

8. 运行时主链源码解剖#

本章回答:

构建完成后,一次 stream / astream / stream_events 如何进入核心执行逻辑并产生可观察输出?

8.1 运行时入口#

调用方式公开入口核心内部入口返回类型
同步最终结果invoke()stream() 聚合 / Pregel loopfinal state
异步最终结果ainvoke()astream() 聚合 / async loopfinal state
同步流式stream()Pregel.stream()Iterator[StreamPart]
异步流式astream()Pregel.astream()AsyncIterator[StreamPart]
事件流stream_events()event router + transformersrun stream object
异步事件流astream_events()async event routerasync run stream object

8.2 运行时总链路#

公开调用入口
输入归一化
Config / callbacks / stream mode 准备
Pregel loop 启动
节点执行
模型 token / 节点 update / custom data / checkpoint / task event 产生
根据 stream_mode 投影为 chunk
yield 给调用方
图到 END 后 stream 结束

8.3 stream() 主入口源码解剖#

职责与所处阶段#

stream() 是同步流式入口,负责执行图并根据 stream_mode 持续产出 chunk。

真实源码签名#

公开 API 形态如下,具体参数以正式版源码为准:

def stream(
self,
input: InputT | Command | None,
config: RunnableConfig | None = None,
*,
context: ContextT | None = None,
stream_mode: StreamMode | Sequence[StreamMode] | None = None,
print_mode: StreamMode | Sequence[StreamMode] = (),
output_keys: str | Sequence[str] | None = None,
interrupt_before: All | Sequence[str] | None = None,
interrupt_after: All | Sequence[str] | None = None,
durability: Durability | None = None,
subgraphs: bool = False,
debug: bool | None = None,
version: Literal["v1", "v2"] | None = None,
**kwargs: Any,
) -> Iterator[dict[str, Any] | Any]:
...

调用方与被调用方#

业务代码
graph.stream(input, stream_mode="updates")
Pregel.stream
PregelLoop
PregelRunner
StreamProtocol
yield chunk

输入、输出与状态变化#

项目类型说明
输入graph input / Command初始 state 或恢复命令
ConfigRunnableConfigtags、metadata、callbacks、thread_id
输出iterator按 stream_mode yield chunks
状态变化graph state节点执行后更新
副作用可能触发 callbacks / tracer / checkpoint观测和持久化

细粒度伪代码#

以下为保留关键控制流的简化伪代码,不是源码逐字复制:

def stream(input_value, config=None, *, stream_mode=None, version=None, **kwargs):
# 1. 规范化 config
config = ensure_config(config)
# 2. 规范化 stream_mode
modes = normalize_stream_modes(
stream_mode or self.default_stream_mode
)
# 3. 初始化 callback manager / run manager
callback_manager = get_callback_manager_for_config(config)
run_manager = callback_manager.on_chain_start(
serialized=self,
inputs=input_value,
name=config.get("run_name") or self.get_name(),
)
# 4. 初始化 stream protocol
stream_protocol = StreamProtocol(
modes=modes,
version=version or "v2",
subgraphs=kwargs.get("subgraphs", False),
print_modes=normalize_print_modes(kwargs.get("print_mode", ())),
)
# 5. 初始化 Pregel loop
with PregelLoop(
graph=self,
input=input_value,
config=config,
stream=stream_protocol,
run_manager=run_manager,
interrupt_before=kwargs.get("interrupt_before"),
interrupt_after=kwargs.get("interrupt_after"),
durability=kwargs.get("durability"),
) as loop:
# 6. 在 loop 中逐步执行 task
while loop.tick():
for task in loop.ready_tasks:
runner.submit(task)
# 7. 读取 runtime 产生的 stream chunks
while stream_protocol.has_chunks():
chunk = stream_protocol.pop_chunk()
yield format_chunk(
chunk,
version=stream_protocol.version,
)
# 8. 成功结束
run_manager.on_chain_end(loop.output)

逐段解释#

第一段规范化 config。callbacks、tags、metadata、thread_id 都在这里进入调用链。

第二段规范化 stream mode。单个字符串和多个 mode 列表都会转成统一结构。

第三段启动 callback 生命周期。LangSmith trace 也是通过 callback/tracer 机制记录每次 run。

第四段初始化 stream protocol。它决定运行时哪些事件要被收集、如何变成输出 chunk。

第五段创建 Pregel loop。图执行、checkpoint、interrupt、state update、task scheduling 都在 loop 中发生。

第六段逐步执行 task。每个 super-step 中可能有多个节点并行执行。

第七段不断从 stream protocol 取出 chunk 并 yield 给调用方。调用方每迭代一次,就看到一部分执行过程。

第八段 run 成功结束,触发 callback end。

正常路径#

stream()
init config / callbacks / stream protocol
Pregel loop
node execution
stream_protocol receives events
yield chunks
END

关键分支与异常路径#

条件行为结果
stream_mode="updates"只 yield node update适合调试节点输出
stream_mode="values"yield 完整 state适合状态快照
stream_mode=["updates","messages"]多模式输出chunk 需按 type 分发
subgraphs=True包含子图事件chunk 带 namespace
节点异常callback on_errorstream 抛异常
interruptyield interrupt 相关数据等待 resume
callbacks 后台提交失败可能影响 trace不应影响业务结果

设计原因与工程影响#

stream() 没有改变图执行语义,只改变结果的暴露方式。这意味着同一张图可以既用于 API 最终响应,也用于调试、前端流式体验和线上观测。

源码证据#

8.4 invoke()stream() 关系源码解剖#

职责与所处阶段#

invoke() 是最终结果入口;从执行语义上,它可以理解为运行图并收集最终输出,而 stream() 是边运行边暴露中间结果。

真实源码签名#

def invoke(
self,
input: InputT | Command | None,
config: RunnableConfig | None = None,
*,
context: ContextT | None = None,
stream_mode: StreamMode = "values",
print_mode: StreamMode | Sequence[StreamMode] = (),
output_keys: str | Sequence[str] | None = None,
interrupt_before: All | Sequence[str] | None = None,
interrupt_after: All | Sequence[str] | None = None,
durability: Durability | None = None,
**kwargs: Any,
) -> dict[str, Any] | Any:
...

调用方与被调用方#

graph.invoke(input)
graph.stream(input, stream_mode="values")
遍历 chunks
保存最后一个 output
return latest

输入、输出与状态变化#

项目类型说明
输入graph input普通 state 或 Command
输出final output最后结果
状态变化与 stream 相同同一运行语义
副作用callbacks / checkpoint与 stream 类似

细粒度伪代码#

以下为保留关键控制流的简化伪代码,不是源码逐字复制:

def invoke(input_value, config=None, **kwargs):
latest_output = None
for chunk in self.stream(
input_value,
config=config,
stream_mode=kwargs.pop("stream_mode", "values"),
**kwargs,
):
# 1. 根据 output mode 识别真正的 state/output
if is_output_chunk(chunk):
latest_output = extract_output(chunk)
# 2. 忽略纯 debug / custom / messages chunks
else:
continue
return latest_output

逐段解释#

invoke() 不需要把每个 chunk 返回给调用方,它只关心最终输出。

如果底层 stream mode 包含多个类型,invoke() 需要选择最终 output。

这解释了为什么调试时 stream() 更透明,而线上 API 如果只要最终结果,可以用 invoke()

正常路径#

invoke
stream execution
consume all chunks
return latest output

关键分支与异常路径#

条件行为结果
图到 END返回最终 state正常
图 interrupt返回 interrupt 或暂停信息需要 resume
节点异常抛出异常调用失败
stream 被调用方提前停止图执行可能被取消取决于 runtime

设计原因与工程影响#

理解 invoke()stream() 的关系后,就不会把二者看成两套逻辑。真正的差异是“是否暴露中间事件”。

源码证据#

  • langgraph/pregel/main.py::Pregel.invoke
  • langgraph/pregel/main.py::Pregel.stream

8.5 stream_mode="updates" 源码解剖#

职责与所处阶段#

updates 用于观察每个节点返回的 partial state update。

真实源码签名#

graph.stream(input, stream_mode="updates", version="v2")

调用方与被调用方#

node returns {"plan": "..."}
Pregel apply writes
updates stream formatter
yield {"type": "updates", "data": {"node": {"plan": "..."}}}

输入、输出与状态变化#

项目类型说明
输入node partial update节点返回值
输出updates chunk节点名到 update 的映射
状态变化update 被 reducer 合并同时也写入 state
副作用无业务副作用只是观察

细粒度伪代码#

以下为保留关键控制流的简化伪代码,不是源码逐字复制:

def emit_updates(task_result, stream_protocol):
# 1. 取得节点名和节点返回值
node_name = task_result.name
update = task_result.writes
# 2. 过滤不可作为 state update 的内部写入
public_update = filter_state_update_keys(update)
if not public_update:
return
# 3. 生成 updates stream part
stream_protocol.write(
mode="updates",
data={
node_name: public_update
},
namespace=task_result.namespace,
metadata=task_result.metadata,
)

逐段解释#

第一段从 task result 中取出节点名和写入内容。

第二段过滤内部 channel 写入,因为并非所有 Pregel write 都是用户 state update。

第三段发出 updates chunk。它的核心价值是定位“哪个节点写了什么”。

正常路径#

node update
filter public keys
yield updates chunk

关键分支与异常路径#

条件行为结果
节点返回空 dict不输出或输出空 update无业务变化
并行多个节点同 step 分别输出多个 updates可定位每个节点
节点异常无正常 update进入 error/task/debug
reducer 合并后值变化updates 仍显示原始节点写入不等于完整 state

设计原因与工程影响#

updates 适合定位节点输出问题,例如旅行规划助手中 parse_requestdays 写错,updates 能直接指出哪个节点输出了错误字段。

源码证据#

  • langgraph/pregel/main.py::Pregel.stream
  • LangGraph Streaming docs: updates streams state updates after each graph step

8.6 stream_mode="values" 源码解剖#

职责与所处阶段#

values 用于观察每个 super-step 后完整 state 快照。

真实源码签名#

graph.stream(input, stream_mode="values", version="v2")

调用方与被调用方#

node updates
reducers apply writes
state channels updated
read output channels
yield full state snapshot

输入、输出与状态变化#

项目类型说明
输入当前 state channels合并后的 state
输出values chunk完整可见 state
状态变化已经完成values 是更新后的结果
副作用只读快照

细粒度伪代码#

以下为保留关键控制流的简化伪代码,不是源码逐字复制:

def emit_values(loop, stream_protocol):
# 1. 当前 super-step 的 writes 已应用到 channels
channels = loop.channels
# 2. 根据 output_keys / stream_channels 读取可见 state
values = read_channels(
channels=channels,
keys=loop.output_keys,
)
# 3. 生成 values stream part
stream_protocol.write(
mode="values",
data=values,
namespace=loop.namespace,
metadata=loop.metadata,
)

逐段解释#

第一段强调 values 发生在更新后。它不是节点原始返回值,而是 reducer 合并后的完整 state。

第二段根据 output keys 读取 state。不是所有内部 channel 都对用户可见。

第三段输出 values chunk。

正常路径#

apply writes
read state
yield values

关键分支与异常路径#

条件行为结果
初始 state可能先输出初始 values便于观察起点
每步后 state输出完整 state适合回放
state 很大输出成本高生产慎用
包含敏感字段可能泄露需要过滤 output_keys

设计原因与工程影响#

values 适合调试 reducer、branch、checkpoint 和最终 state,但线上高频流式 UI 不应盲目暴露完整 state。

源码证据#

  • langgraph/pregel/main.py::Pregel.stream
  • LangGraph Streaming docs: values streams full state after each step

8.7 stream_mode="messages" 源码解剖#

职责与所处阶段#

messages 用于观察图中任何 LangChain ChatModel 调用产生的 token / message chunk,并携带 metadata。

真实源码签名#

graph.stream(input, stream_mode="messages", version="v2")

输出形态:

(message_chunk, metadata)

调用方与被调用方#

node calls ChatModel.invoke / stream
ChatModel emits AIMessageChunk
LangChain callbacks / stream hooks
LangGraph messages stream
yield (message_chunk, metadata)

输入、输出与状态变化#

项目类型说明
输入ChatModel token chunks模型输出流
输出(message_chunk, metadata)token + 调用元数据
状态变化不一定token 流不等于 state update
副作用callbacks / trace可记录模型输出

细粒度伪代码#

以下为保留关键控制流的简化伪代码,不是源码逐字复制:

def node(state, config):
# 1. 节点内部调用模型
response = model.invoke(
[{"role": "user", "content": state["user_request"]}],
config=config,
)
# 2. 节点最终仍然返回 state update
return {"answer": response.content}
def on_chat_model_stream(message_chunk, run_metadata):
# 3. ChatModel 每产生一个 chunk,callback / stream hook 捕获
metadata = {
"langgraph_node": run_metadata.node_name,
"tags": run_metadata.tags,
"run_id": run_metadata.run_id,
"model_name": run_metadata.model_name,
}
# 4. 写入 messages stream
stream_protocol.write(
mode="messages",
data=(message_chunk, metadata),
)

逐段解释#

节点本身最终仍然返回 update,例如 {"answer": response.content}

messages 观察的是模型生成过程中的 token,不是节点完成后的 update。

metadata 让调用方知道这个 token 来自哪个节点、哪个模型调用、哪些 tags。

所以在一个图里有多个模型节点时,前端可以通过:

metadata["langgraph_node"]
metadata["tags"]

过滤要展示的 token。

正常路径#

ChatModel streaming
AIMessageChunk
metadata attached
messages chunk

关键分支与异常路径#

条件行为结果
模型支持 token streaming输出 chunks正常流式
模型不支持 streaming可能只输出最终 message体验下降
使用 nostream tag不输出 token内部模型调用不暴露
多模型并行tokens 交错需用 metadata 过滤
Python < 3.11 async 未传 configcontext 可能丢失messages 流异常

设计原因与工程影响#

token streaming 是用户体验层能力;node update streaming 是流程调试能力。二者不要混为一谈。

源码证据#

  • langgraph/pregel/messages.py
  • langchain_core/callbacks/manager.py
  • LangGraph Streaming docs: messages streams (message_chunk, metadata)

8.8 stream_mode="custom"get_stream_writer() 源码解剖#

职责与所处阶段#

custom 用于从节点或工具内部主动发出业务进度、日志或阶段事件。

真实源码签名#

def get_stream_writer() -> StreamWriter:
...

使用方式:

writer = get_stream_writer()
writer({"status": "searching", "progress": 30})

调用方与被调用方#

node / tool
get_stream_writer()
current runtime context
writer(data)
custom stream chunk

输入、输出与状态变化#

项目类型说明
输入JSON-friendly data自定义业务事件
输出custom chunk调用方可消费
状态变化不写入 state
副作用stream 输出可用于 UI 进度

细粒度伪代码#

以下为保留关键控制流的简化伪代码,不是源码逐字复制:

def get_stream_writer():
# 1. 从当前运行上下文读取 writer
writer = current_runtime_context.get("stream_writer")
# 2. 如果当前不是 stream 执行,可能返回 no-op writer 或抛错
if writer is None:
return noop_writer_or_raise()
return writer
def node(state):
writer = get_stream_writer()
# 3. 发出自定义进度
writer({
"type": "progress",
"stage": "search_food",
"message": "正在检索美食推荐",
"progress": 30,
})
# 4. 执行业务逻辑
notes = search_food(state["user_request"])
writer({
"type": "progress",
"stage": "search_food",
"message": "美食检索完成",
"progress": 100,
})
# 5. 返回真正的 state update
return {"research_notes": notes}

逐段解释#

custom 数据不会进入 state,除非节点同时返回 state update。

custom 适合 UI 进度和业务日志,例如“正在查询天气”“已召回 20 条景点”。

不要把 custom 当成状态持久化机制。如果后续节点需要读取数据,必须返回到 state。

正常路径#

get_stream_writer()
writer(data)
custom chunk
node return update

关键分支与异常路径#

条件行为结果
stream_mode 包含 custom输出 custom chunk正常
stream_mode 不含 customcustom 数据可能不被消费不影响业务结果
async Python < 3.11contextvar 可能不可用需手动传 writer
custom 数据不可序列化UI / trace 失败应使用 JSON-friendly
把业务数据只写 custom下游节点读不到应写 state

设计原因与工程影响#

custom 把“过程提示”和“业务状态”分开。前者用于用户体验和调试,后者用于图执行逻辑。

源码证据#

  • langgraph/config.py::get_stream_writer
  • LangGraph Streaming docs: Custom data

8.9 stream_events() typed projections 源码解剖#

职责与所处阶段#

stream_events() 是事件流入口,把底层 raw stream events 通过 transformer 转成 typed projections,如 messages、values、subgraphs、interrupts、output。

真实源码签名#

公开用法:

stream = graph.stream_events(input, version="v3")

调用方与被调用方#

graph.stream_events(input)
底层 Pregel raw events
event router
stream transformers
stream.messages / stream.values / stream.output / stream.interrupts

输入、输出与状态变化#

项目类型说明
输入graph input与 stream 类似
输出run stream objecttyped projections
状态变化与 graph 执行一致projections 不改变 state
副作用callback / tracer可记录 trace

细粒度伪代码#

以下为保留关键控制流的简化伪代码,不是源码逐字复制:

def stream_events(input_value, config=None, version="v3"):
# 1. 创建底层 raw event stream
raw_stream = self.stream(
input_value,
config=config,
stream_mode=[
"values",
"updates",
"messages",
"custom",
"tasks",
"checkpoints",
],
version="v2",
)
# 2. 创建 event router
router = EventRouter()
# 3. 注册内置 transformers
router.add_transformer("messages", MessagesTransformer())
router.add_transformer("values", ValuesTransformer())
router.add_transformer("subgraphs", SubgraphsTransformer())
router.add_transformer("interrupts", InterruptsTransformer())
router.add_transformer("output", OutputTransformer())
# 4. 将 raw events 路由到 projections
for event in raw_stream:
router.route(event)
# 5. 返回可并发消费的 run stream object
return router.run_stream

逐段解释#

第一段从底层 stream modes 获取 raw graph execution events。

第二段和第三段创建 transformer pipeline。不同 projection 关心不同事件。

第四段把 raw event 送入 router,每个 transformer 选择是否消费该事件。

第五段暴露 typed projections。应用代码可以同时读取 stream.messagesstream.values,互不消耗彼此数据。

正常路径#

raw Pregel events
event router
transformers
typed projections

关键分支与异常路径#

条件行为结果
应用只关心 token读取 stream.messages简洁
应用关心状态读取 stream.values简洁
应用关心中断读取 stream.interruptsHITL 友好
多消费者同时消费 projections互不影响
需要底层细节使用 stream() modes更底层

设计原因与工程影响#

低层 stream_mode 更像“原始事件总线”,stream_events() 更像“应用层投影 API”。新应用优先使用 typed projections,可以减少分支判断和 chunk 形状处理。

源码证据#

  • LangGraph Event Streaming docs
  • langgraph/pregel/main.py
  • event streaming transformer 相关源码

8.10 Callback、事件与可观测性源码解剖#

职责与所处阶段#

Callbacks 负责在 Runnable 生命周期关键点触发事件,tracer 可以把这些事件记录为 trace。

真实源码签名#

典型 callback 方法包括:

on_chain_start(...)
on_chain_end(...)
on_chain_error(...)
on_llm_start(...)
on_llm_new_token(...)
on_llm_end(...)
on_tool_start(...)
on_tool_end(...)

调用方与被调用方#

Runnable / Graph / Model / Tool
CallbackManager
CallbackHandler / Tracer
Console / LangSmith / Custom sink

输入、输出与状态变化#

项目类型说明
输入lifecycle eventstart/end/error/token
输出handler side effectlog / trace / metrics
状态变化无业务 state 修改只记录观测
副作用上传 trace / 打日志可能异步后台执行

细粒度伪代码#

以下为保留关键控制流的简化伪代码,不是源码逐字复制:

def runnable_invoke(input_value, config):
config = ensure_config(config)
callback_manager = get_callback_manager_for_config(config)
run_manager = callback_manager.on_chain_start(
serialized=self,
inputs=input_value,
tags=config.get("tags"),
metadata=config.get("metadata"),
)
try:
output = self._invoke(input_value, config)
except Exception as exc:
run_manager.on_chain_error(exc)
raise
else:
run_manager.on_chain_end(output)
return output
def chat_model_stream(messages, config):
run_manager = callback_manager.on_llm_start(
serialized=model,
prompts_or_messages=messages,
)
for token in provider_stream(messages):
chunk = AIMessageChunk(content=token)
run_manager.on_llm_new_token(token, chunk=chunk)
yield chunk
run_manager.on_llm_end(final_generation)

逐段解释#

第一段展示 Runnable 调用生命周期:start、error、end。

第二段展示模型 token 生命周期:模型开始、每个 token、模型结束。

LangGraph messages stream 和 LangSmith trace 都可以利用这些事件,但二者目标不同:

messages stream
面向实时 UI / token 展示。
LangSmith trace
面向调试、评估、监控和历史记录。

正常路径#

Runnable start
sub-runs start/end
LLM token events
tool events
Runnable end
trace tree

关键分支与异常路径#

条件行为结果
handler 抛异常通常不应破坏业务需要安全处理
callbacks=None不记录外部 trace业务仍执行
LangSmith enabled自动上传 trace可视化
metadata 含敏感信息可能泄露需要 anonymizer / filtering
token streaming disabled无 token event仍可能有 final trace

设计原因与工程影响#

Callbacks 是跨 LangChain / LangGraph / Model / Tool 的统一观测协议。Streaming 更关注“现在给调用方看什么”,callbacks 更关注“完整执行链记录什么”。

源码证据#

  • langchain_core/callbacks/manager.py
  • langchain_core/runnables/config.py
  • langchain_core/tracers/
  • LangSmith Observability docs

8.11 运行时主链总结#

stream(input)
ensure_config(callbacks/tags/metadata)
init stream protocol
Pregel loop executes tasks
state writes → updates / values
LLM chunks → messages
writer(data) → custom
checkpoint/task events → checkpoints/tasks/debug
callbacks → traces
yield chunks / event projections

9. 关键分支、异常与边界#

本章回答:

当输出模式、执行方式或观测目标不同,框架如何分流、恢复、终止或失败?

9.1 分支矩阵#

分支类型触发条件核心函数结果
完整状态流stream_mode="values"Pregel.stream()每步完整 state
节点增量流stream_mode="updates"Pregel.stream()每节点 update
token 流stream_mode="messages"ChatModel callbacks / stream hookstoken + metadata
自定义流stream_mode="custom"get_stream_writer()业务事件
多模式流stream_mode=[...]stream protocol多类 chunk
typed event flowstream_events()event routerprojections
子图流subgraphs=TruePregel namespaceroot + subgraph chunks
checkpoint 流stream_mode="checkpoints"checkpointerStateSnapshot-like events
task 流stream_mode="tasks"task runnerstart/finish/error
debug 流stream_mode="debug"combined debug高噪声细节

9.2 同步与异步分支#

维度同步路径异步路径
入口stream()astream()
返回IteratorAsyncIterator
节点函数sync callableasync callable / mixed
模型调用invoke() / stream()ainvoke() / astream()
配置传播ContextVarPython < 3.11 需注意显式传 config
适用场景CLI、脚本、简单服务FastAPI async、WebSocket、并发服务

9.3 Batch、Stream、Parallel 或路由分支#

batch()stream() 是不同观察维度:

batch
多个输入并发执行,返回多个最终结果。
stream
一个输入执行过程中逐步输出中间事件。
parallel graph step
一个 super-step 内多个节点并行执行,updates 可能分多条输出。
messages stream
多个模型节点并行时 token 可能交错,需要 metadata 过滤。

9.4 异常分类#

异常类别抛出位置是否可恢复处理策略是否反馈上层
节点业务异常node function视情况retry/fallback/fail
模型 streaming 异常ChatModel stream视情况fallback model
callback 异常callback handler通常不应影响主流程handler 内部兜底视实现
custom 数据不可序列化writer / transport改成 JSON-friendly
消费方过慢外部 iterator consumer视部署backpressure / buffer
stream 提前中断调用方停止消费视 runtime取消或清理
LangSmith 上传失败tracer通常可忽略后台重试 / flush不应破坏业务

9.5 异常路径源码解剖#

细粒度伪代码#

以下为保留关键控制流的简化伪代码,不是源码逐字复制:

def stream_with_errors(input_value, config):
run_manager = callback_manager.on_chain_start(...)
try:
for chunk in run_pregel_loop(input_value, config):
yield chunk
except GraphInterrupt as interrupt:
# 中断不是普通失败,它是可恢复暂停
run_manager.on_chain_end({"__interrupt__": interrupt})
yield build_interrupt_chunk(interrupt)
except Exception as exc:
# 普通异常进入 error callback
run_manager.on_chain_error(exc)
raise
else:
run_manager.on_chain_end(final_output)

逐段解释#

GraphInterrupt 属于可恢复暂停,不应和普通异常混为一谈。

普通异常需要向上抛出,否则调用方以为图成功完成。

callback error 和业务 error 的处理边界不同。生产中 callback 失败不应该导致退款、下单、审批等业务失败。

9.6 Retry、Fallback 与恢复边界#

机制适用条件不适用条件幂等要求
Retry模型或工具临时失败已产生不可重复副作用节点幂等
Fallback模型流式失败安全边界降低fallback 输出需校验
Repair输出结构错误业务动作已执行不应重复副作用
Resumeinterrupt 暂停普通异常同 thread_id
Trace replay复盘失败改写线上状态只读或隔离环境
Stream reconnect前端断线图已结束且无 checkpoint需要 thread/checkpoint

9.7 停止条件与保护上限#

正常结束:
图到达 END,stream iterator 耗尽。
提前结束:
调用方停止消费 stream,运行时可能取消任务或释放资源。
人工中断:
interrupt 触发,stream 输出中断信息,等待 Command(resume=...)。
框架保护:
recursion_limit 防止图循环无限执行。
异常失败:
节点异常、模型异常、序列化异常或 callback 严重异常。

9.8 能力边界#

容易误判的能力实际提供者本篇对象的真实职责
前端打字机效果messages stream + UILangGraph 只输出 token chunks
节点调试updates / tasks / tracestream 暴露数据,不分析原因
线上监控LangSmith / metrics systemstream 不是长期监控数据库
持久恢复checkpointerstream 不保存状态
日志审计logger / trace / DBcallbacks 只是事件入口
成本统计model usage metadata / LangSmithstream 不自动计费
用户可见输出应用层选择不要直接展示所有 chunks

10. 扩展机制与框架协作#

本章回答:

框架允许在哪里插入自定义观测行为,它如何与相邻模块协作,哪些内部实现不应该被业务代码依赖?

10.1 扩展点总览#

扩展点扩展方式执行时机可修改内容约束
stream_mode调用参数graph 运行时输出类型不改变业务 state
get_stream_writer()节点 / 工具调用节点内部custom stream data数据应可序列化
callbacksRunnableConfigstart/end/error/token记录日志/trace不应改变业务结果
metadataRunnableConfig调用开始trace 属性避免敏感信息
tagsRunnableConfig / model config子调用传播过滤 token/trace命名规范
LangSmith env vars环境变量全局运行时trace 上传注意成本和隐私
tracing_contextcontext manager局部代码块启停 trace / project只影响观测
custom stream transformersevent streaming 扩展事件路由自定义 projection需维护协议稳定

10.2 get_stream_writer() 自定义进度源码解剖#

职责与所处阶段#

在节点或工具内部输出不进入 state 的过程信息。

真实源码签名#

def get_stream_writer() -> Callable[[Any], None]:
...

调用方与被调用方#

node/tool
get_stream_writer()
writer(payload)
custom stream
UI / log consumer

输入、输出与状态变化#

项目类型说明
输入custom payload业务进度
输出custom chunk调用方消费
状态变化不写 state
副作用输出事件影响 UI / debug

细粒度伪代码#

def search_node(state):
writer = get_stream_writer()
writer({
"stage": "rewrite_query",
"status": "start",
"message": "正在改写检索 query",
})
query = rewrite_query(state["user_request"])
writer({
"stage": "retrieval",
"status": "start",
"query": query,
})
docs = retriever.invoke(query)
writer({
"stage": "retrieval",
"status": "end",
"doc_count": len(docs),
})
return {"search_results": docs}

逐段解释#

每次 writer(payload) 都变成一个 custom event。

这些 payload 不参与 reducer,不影响路由,不写入 checkpoint。

如果业务下游需要 docs,就必须通过 return {"search_results": docs} 写入 state。

正常路径#

writer(progress)
custom chunk
return state update

关键分支与异常路径#

条件行为结果
stream_mode 包含 custom输出进度UI 可展示
不包含 custom进度无人消费业务不受影响
payload 大增加传输压力应精简
payload 含敏感信息可能泄露必须脱敏

设计原因与工程影响#

custom stream 是 Agent 产品体验的重要基础,可以让用户看到“正在查天气 / 正在检索景点 / 正在合并计划”,但它不是 trace 或 state。

源码证据#

  • langgraph/config.py::get_stream_writer
  • LangGraph Streaming docs: Custom data

10.3 callbacks 与 LangSmith trace 协作源码解剖#

职责与所处阶段#

callbacks 负责把执行生命周期事件交给 handler;LangSmith tracer 是其中一种 handler / tracing sink。

真实源码签名#

典型用法:

graph.invoke(
input,
config={
"tags": ["travel-agent", "prod"],
"metadata": {
"user_id": "u_123",
"session_id": "s_456",
},
},
)

启用环境变量:

Terminal window
export LANGSMITH_TRACING=true
export LANGSMITH_API_KEY=<your-api-key>
export LANGSMITH_PROJECT=travel-agent

调用方与被调用方#

graph.invoke / stream
RunnableConfig
CallbackManager
LangSmith tracer
LangSmith project trace

输入、输出与状态变化#

项目类型说明
输入tags / metadata / callbackstrace 属性
输出trace run treeLangSmith 可视化
状态变化无业务 state 修改只观测
副作用网络上传 trace需要 API key

细粒度伪代码#

def traceable_graph_call(input_value, config):
config = ensure_config(config)
if langsmith_tracing_enabled():
config = attach_langsmith_tracer(config)
return graph.invoke(
input_value,
config=config,
)
def attach_langsmith_tracer(config):
tracer = LangSmithTracer(
project_name=os.getenv("LANGSMITH_PROJECT", "default"),
api_key=os.getenv("LANGSMITH_API_KEY"),
)
config["callbacks"] = merge_callbacks(
config.get("callbacks"),
[tracer],
)
return config

逐段解释#

LangSmith 可以通过环境变量自动接入,不需要重写业务图。

tagsmetadata 会进入 trace,便于按照环境、版本、用户、场景过滤。

生产中 metadata 必须避免敏感数据,或者使用 anonymizer。

正常路径#

set env vars
run graph
callbacks receive events
LangSmith displays trace tree

关键分支与异常路径#

条件行为结果
未启用 tracing不上传 trace本地业务正常
API key 错误trace 上传失败业务不应失败
metadata 含敏感数据trace 泄露风险必须脱敏
project 未设置进入 default project难以管理
selective tracing只记录指定调用降低成本与隐私风险

设计原因与工程影响#

stream 适合实时观察,LangSmith trace 适合事后分析。线上问题排查通常需要二者结合:用户当场看到 progress,工程师事后看 trace。

源码证据#

  • LangSmith Observability docs
  • LangSmith trace-with-LangChain docs
  • langchain_core/callbacks/manager.py
  • langchain_core/tracers/

10.4 与 checkpoint 的协作#

checkpointer
保存状态快照和恢复点
stream_mode="checkpoints"
观察 checkpoint events
get_state()
读取最新状态
LangSmith trace
记录执行步骤和调用链

边界:

checkpoint 是恢复机制;
trace 是观测机制;
stream 是实时输出机制。

不要用 trace 代替 checkpoint,也不要用 checkpoint 当调试 UI。

10.5 与 subgraph 的协作#

parent graph stream(subgraphs=True)
root chunks ns=()
subgraph chunks ns=("parent_node:<task_id>", ...)

使用场景:

Plan-and-Execute 中 executor subgraph 调试
Multi-Agent 中子 Agent 执行过程观察
复杂旅行规划中搜索子图和评估子图的输出拆分

注意:

如果不打开 subgraphs=True,只看父图输出,可能误以为子图没有执行细节。

10.6 与 Agent Trace 的协作#

Agent Trace 通常需要同时记录:

用户输入
state updates
LLM messages/tokens
tool calls
tool results
custom progress
interrupts
checkpoints
errors
metadata
usage
latency
cost

LangGraph stream 提供实时过程,LangSmith 提供运行历史和可视化 trace,自定义业务日志提供审计闭环。

10.7 选择 stream 还是 events 还是 callbacks#

条件选择 stream选择 stream_events选择 callbacks / LangSmith
调试节点 update
前端展示 token
多消费者同时读取 messages / values
线上问题事后排查部分
自定义业务进度
长期监控与评估
需要最底层 runtime 事件部分

11. 工程决策与适用场景#

11.1 适用场景#

场景是否推荐原因
本地调试节点输出updates能看到每个节点写了什么
调试状态合并values能看到 reducer 后完整 state
Chat UI 打字机效果messagesstream_events.messages可流式展示 token
工具进度显示custom节点/工具主动发进度
复杂子图排障subgraphs=True可观察子图内部
HITL 前端stream_events.interrupts直接拿 pending interrupt
线上 traceLangSmith持久化、可搜索、可评估
成本/延迟分析LangSmith + usage metadatastream 只实时输出,不做聚合
普通后端接口最终结果invoke()不需要中间过程
高敏数据生产环境谨慎 stream values/debug可能泄露完整 state

11.2 工程决策表#

决策点推荐选择前提风险
调试 stateupdates 优先想看哪个节点写错不是完整 state
调试 reducervaluesstate 不太大可能泄露敏感字段
前端 tokenmessages / stream_events.messages模型支持 streaming多模型 token 交错
业务进度custom节点有明确阶段不要写过大数据
应用层事件stream_events()需要多个 typed projections学习成本略高
线上 traceLangSmith可接受数据上传隐私与成本
metadata放 user/session/env/versionJSON-friendly不要放原文敏感数据
callbacks自定义 handler需要内部日志/指标handler 不能破坏主流程

11.3 性能、可靠性与安全边界#

性能:
values/debug 输出大 state 成本高;messages token 高频输出会增加传输压力;LangSmith trace 上传也有额外成本。
可靠性:
streaming consumer 过慢可能造成背压;callback handler 不应影响主流程;前端断线需要配合 checkpoint 或任务状态恢复。
安全:
不要把完整 state、工具结果、用户隐私直接通过 values/debug/metadata 发给前端或第三方 trace。
可观测性:
至少记录 trace_id、thread_id、node_name、tool_name、latency、model、token usage、error、route decision。

11.4 旅行规划助手观测设计#

观测目标推荐机制记录内容
用户看到进度custom正在解析需求、正在检索景点、正在生成计划
用户看到回答生成messagesfinal_response 节点 token
开发调试节点updatesintent、search_results、draft_plan、final_plan
调试 reducervaluesresearch_notes、messages、step_results
调试子图subgraphs=Trueplanner / executor subgraph
审批 / interruptstream_events.interrupts待确认计划或工具动作
线上排错LangSmith完整 run tree、metadata、错误
成本分析LangSmith + usage_metadatamodel、tokens、latency

12. 常见误区与源码纠正#

12.1 误区:stream 是另一套执行逻辑#

错误原因:

stream() 的返回形态和 invoke() 完全不同。

源码事实:

stream 和 invoke 使用同一张 compiled graph 和 Pregel runtime;
差异主要是是否边执行边 yield 中间事件。

工程影响:

不要为 stream 和 invoke 写两套业务图。应复用同一 graph。

12.2 误区:updates 就是完整 state#

错误原因:

updates 看起来也是 dict。

源码事实:

updates 是节点返回的 partial update;
values 才是每步后的完整 state。

工程影响:

如果要调试 reducer 合并后结果,用 values;如果要定位哪个节点写错,用 updates。

12.3 误区:messages 是节点输出#

错误原因:

messages 模式名字容易和 state[“messages”] 混淆。

源码事实:

stream_mode="messages" 输出 LLM token / message chunks + metadata;
它不等于 LangGraph state 中的 messages 字段。

工程影响:

不要用 messages stream 当作节点最终 state。节点仍需 return state update。

12.4 误区:custom stream 会写入 state#

错误原因:

writer(data) 很像 return update。

源码事实:

get_stream_writer().write(data) 只发 custom event;
不会进入 reducer,也不会被下游节点读取。

工程影响:

下游需要的数据必须返回 state update。

12.5 误区:metadata 可以随便放业务数据#

错误原因:

metadata 传递方便,LangSmith 也能展示。

源码事实:

metadata 会进入 callbacks 和 trace;
可能被外部观测系统记录。

工程影响:

不要放身份证、手机号、订单敏感字段、API key、完整用户隐私。

12.6 误区:LangSmith trace 等于 stream#

错误原因:

二者都能观察运行过程。

源码事实:

stream 是实时输出机制;
LangSmith trace 是持久化运行记录和调试监控平台。

工程影响:

实时 UI 用 stream;线上排查和评估用 trace。

12.7 误区:debug stream 可以直接给前端#

错误原因:

debug 信息很完整。

源码事实:

debug 可能包含 checkpoints、tasks、metadata、内部状态。

工程影响:

debug 只适合开发和受控运维,不适合用户前端。


13. 最终心智模型与掌握检查#

13.1 构建期心智模型#

StateGraph builder
compile
CompiledStateGraph / Pregel
具备 invoke / stream / astream / stream_events 能力

13.2 运行时心智模型#

graph.stream(input, stream_mode)
ensure_config(callbacks/tags/metadata)
Pregel loop
node execution
state updates / LLM tokens / custom data / task events
mode-specific chunks
caller consumes chunks

13.3 分支与异常心智模型#

values → 完整 state 快照
updates → 节点 partial update
messages → LLM token + metadata
custom → 节点/工具业务事件
checkpoints → checkpoint snapshot
tasks → task lifecycle
debug → 高噪声调试
events → typed projections
普通异常 → on_error + stream 抛错
interrupt → pending interrupt + resume
callback 失败 → 应隔离,不破坏业务

13.4 一句话总结#

Streaming 把 LangGraph 的 Pregel 执行过程投影为状态、节点、模型 token、自定义进度和调试事件;callbacks 与 metadata 把同一执行链交给 tracer,LangSmith 则把这些运行步骤持久化为可搜索、可复盘、可评估的 Agent trace。

13.5 掌握检查#

  • 能说清 invoke()stream() 的关系。
  • 能解释 updatesvalues 的区别。
  • 能解释 messages stream 输出的不是 state[“messages”]。
  • 能写出 custom progress event。
  • 能用 metadata 过滤某个节点或模型调用的 token。
  • 能解释 stream_events() 相比 stream() 的优势。
  • 能解释 callbacks 与 LangSmith trace 的关系。
  • 能说明 metadata / tags 不应放敏感数据。
  • 能为旅行规划助手设计实时 UI stream 和后台 trace。
  • 能根据调试目标选择正确 stream mode。

14. 参考资料与下一篇衔接#

14.1 官方概念文档#

  1. LangGraph Streaming
    https://docs.langchain.com/oss/python/langgraph/streaming

  2. LangGraph Event Streaming
    https://docs.langchain.com/oss/python/langgraph/event-streaming

  3. LangGraph Observability
    https://docs.langchain.com/oss/python/langgraph/observability

  4. LangSmith Trace with LangChain
    https://docs.langchain.com/langsmith/trace-with-langchain

14.2 官方 API Reference#

  1. Pregel.stream
    https://reference.langchain.com/python/langgraph/pregel/main/Pregel/stream

  2. Pregel.astream
    https://reference.langchain.com/python/langgraph/pregel/main/Pregel/astream

  3. Runnable.astream_events
    https://reference.langchain.com/python/langchain-core/runnables/base/Runnable/astream_events

  4. RunnableConfig
    https://reference.langchain.com/python/langchain-core/runnables/config/RunnableConfig

  5. get_stream_writer
    https://reference.langchain.com/python/langgraph/config/get_stream_writer

14.3 官方源码#

  1. langgraph/pregel/main.py::Pregel.stream
    https://github.com/langchain-ai/langgraph/blob/1.2.7/libs/langgraph/langgraph/pregel/main.py

  2. langgraph/pregel/main.py::Pregel.astream
    https://github.com/langchain-ai/langgraph/blob/1.2.7/libs/langgraph/langgraph/pregel/main.py

  3. langgraph/config.py::get_stream_writer
    https://github.com/langchain-ai/langgraph/blob/1.2.7/libs/langgraph/langgraph/config.py

  4. langchain_core/runnables/base.py::Runnable.astream_events
    https://github.com/langchain-ai/langchain/blob/master/libs/core/langchain_core/runnables/base.py

  5. langchain_core/runnables/config.py::RunnableConfig
    https://github.com/langchain-ai/langchain/blob/master/libs/core/langchain_core/runnables/config.py

  6. langchain_core/callbacks/manager.py
    https://github.com/langchain-ai/langchain/blob/master/libs/core/langchain_core/callbacks/manager.py

  7. langchain_core/tracers/
    https://github.com/langchain-ai/langchain/tree/master/libs/core/langchain_core/tracers

14.4 下一篇衔接#

下一篇进入:

第 14 篇:Memory / Store 与长期记忆源码解剖

需要继续回答:

checkpoint 与 store 的边界是什么?
thread-scoped state 与 cross-thread memory 有什么区别?
Store 如何被 node / tool 读取?
长期记忆写入策略如何工程化?
记忆如何避免污染上下文?
LangGraph 源码深潜:Streaming、Events 与 Observability 机制解剖
https://jupiter-ws.cn/posts/agent-frameworks/langgraph-streaming-events-observability-deep-dive/
作者
Jupiter
发布于
2026-03-17
许可协议
CC BY-NC-SA 4.0