LangGraph 源码深潜:Streaming、Events 与 Observability 机制解剖
核心问题: LangChain / LangGraph 如何把一次 Agent / Workflow 执行过程拆成可观察的状态更新、模型 token、节点事件、自定义事件和外部 trace?
源码主线:
graph.stream()→Pregel.stream()→StreamProtocol→stream_mode分发 →values / updates / messages / custom / events→callbacks / tracers→LangSmith前置文章: 第 1 篇
Runnable、第 2 篇Prompt / Message / ChatModel、第 6 篇StateGraph、第 8 篇Reducer、第 11 篇Checkpoint、第 12 篇Interrupt / HITL依赖基线:
langgraph==1.2.7、langchain-core==1.4.8、langchain==1.3.11源码基线:
https://github.com/langchain-ai/langgraph/tree/1.2.7、https://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 学习目标
完成本篇后,读者必须能够:
- 解释
invoke()与stream()在结果形态上的区别,以及为什么invoke()可以理解为对 streaming 执行结果的聚合。 - 区分
stream_mode="values"、"updates"、"messages"、"custom"、"checkpoints"、"tasks"、"debug"的观察对象。 - 解释
stream_mode如何影响 Pregel runtime 向外发出的 chunk。 - 解释 token streaming 与 node update streaming 的本质区别。
- 解释
get_stream_writer()如何把节点内部业务进度写入customstream。 - 解释
stream_events()/astream_events()为什么比低层stream_mode更适合应用代码。 - 解释
metadata、tags、callbacks如何通过RunnableConfig传递到子调用。 - 解释 LangSmith trace 如何接入,以及它与本地 stream 的边界。
- 能为旅行规划助手设计一套调试与线上观测方案。
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 / subgraphs2.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 |
| 回调 | callbacks | LangChain Runnable 生命周期 hooks | LangGraph branch |
| 元数据 | metadata | trace 和过滤用辅助信息 | 业务 state |
| 标签 | tags | 调用链标记和过滤 | 节点名 |
| Trace | LangSmith run tree | 执行步骤记录 | streaming output |
2.4 与相邻抽象的边界
| 对象 | 负责什么 | 不负责什么 | 与本篇对象的关系 |
|---|---|---|---|
invoke() | 返回最终结果 | 不逐步暴露过程 | 与 stream() 使用同一图运行语义 |
stream() | 低层 chunk 输出 | 不提供 typed projection | 适合调试 runtime 事件 |
stream_events() | typed projections | 不替代 trace 存储 | 更适合应用代码 |
callbacks | 记录 Runnable 生命周期 | 不直接改变 state | 可与 LangSmith trace 协作 |
LangSmith | trace、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 ↓到达 END3.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 + metadata | Chat UI、模型输出实时展示 |
custom | 节点或工具主动发出的自定义数据 | 工具进度、业务日志 |
checkpoints | checkpoint events | Durable execution 调试 |
tasks | task start / finish / error | 节点执行耗时、失败定位 |
debug | 多种底层调试信息 | 深度排障,不适合直接给用户 |
events / stream_events | typed projections | 应用层多消费者观测 |
3.4 对象流转
| 阶段 | 输入类型 | 核心函数 | 输出类型 | 状态变化 |
|---|---|---|---|---|
| 构建图 | StateGraph | compile() | CompiledStateGraph | 图获得 stream 能力 |
| 同步流式调用 | dict | graph.stream() | Iterator[StreamPart] | state 随节点更新 |
| 异步流式调用 | dict | graph.astream() | AsyncIterator[StreamPart] | state 随节点更新 |
| 模型调用 | messages | ChatModel.stream/astream | AIMessageChunk | 不一定写 state |
| 节点返回 | partial update | apply_writes | 新 state | reducer 合并 |
| 自定义事件 | writer(data) | get_stream_writer() | custom chunk | 不修改 state |
| 事件投影 | raw stream events | stream_events() | typed projections | 不修改 state |
| Trace 记录 | RunnableConfig | callbacks / tracers | LangSmith 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 chunks3.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 关键文件
| 优先级 | 文件 | 核心对象 | 阅读目的 |
|---|---|---|---|
| 1 | langgraph/pregel/main.py | Pregel.stream() / Pregel.astream() | 理解 stream 主入口 |
| 2 | langgraph/pregel/loop.py | Pregel loop | 理解 super-step 与事件发出 |
| 3 | langgraph/pregel/runner.py | task runner | 理解节点执行和错误事件 |
| 4 | langgraph/config.py | get_stream_writer() | 理解 custom stream writer |
| 5 | langgraph/types.py | StreamPart / Command / Interrupt | 理解流式输出对象 |
| 6 | langchain_core/runnables/base.py | Runnable.stream() / astream_events() | 理解 Runnable 级 streaming |
| 7 | langchain_core/runnables/config.py | RunnableConfig | 理解 tags、metadata、callbacks |
| 8 | langchain_core/callbacks/manager.py | callback manager | 理解 callback 生命周期 |
| 9 | langchain_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:完整 statemessages:模型 tokencustom:业务进度events:typed projections5. 对象模型、继承关系与协议边界
5.1 核心对象关系
Runnable protocol ↓CompiledStateGraph / Pregel ↓invoke / stream / astream ↓Pregel loop ↓StreamProtocol / callback manager ↓StreamPart / typed event projections5.2 对象职责
| 对象 | 生命周期 | 输入 | 输出 | 核心职责 |
|---|---|---|---|---|
CompiledStateGraph | 编译后 | graph input | final state / stream chunks | 暴露 Runnable 接口 |
Pregel.stream() | 运行时 | input + config + stream_mode | iterator | 执行图并 yield chunks |
Pregel.astream() | 运行时 | input + config + stream_mode | async iterator | 异步执行并 yield chunks |
stream_mode | 调用时配置 | string / list[str] | 影响输出 projection | 决定观察内容 |
get_stream_writer() | 节点运行时 | custom data | custom stream chunk | 发送业务进度 |
RunnableConfig | 调用时配置 | callbacks/tags/metadata | 子调用继承 | 观测上下文传播 |
CallbackManager | 运行时 | lifecycle event | callback calls | 触发回调和 tracer |
LangSmith tracer | 运行时 | callback events | trace run tree | 上传可视化 trace |
5.3 协议边界
stream 协议 负责:把执行过程逐步 yield 给调用方。 不负责:存储长期 trace。
callback 协议 负责:在 Runnable 生命周期节点触发 hooks。 不负责:决定 LangGraph 路由。
metadata / tags 协议 负责:携带过滤和 trace 辅助信息。 不负责:作为业务 state 使用。
LangSmith 协议 负责:接收和展示 trace。 不负责:让本地 stream 更快或改变业务结果。5.4 稳定接口与内部实现
| 类型 | 对象 | 文章中的使用原则 |
|---|---|---|
| 公共 API | graph.stream() | 可用于工程代码 |
| 公共 API | graph.astream() | 可用于异步服务 |
| 公共 API | graph.stream_events() | 推荐应用层 typed streaming |
| 公共 API | get_stream_writer() | 可用于 custom stream |
| 公共 API | config={"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 docs | Basic usage |
updates 是每步 state update | 官方契约 | LangGraph Streaming docs | Graph state |
values 是每步完整 state | 官方契约 | LangGraph Streaming docs | Graph state |
messages 是 LLM token + metadata | 官方契约 | LangGraph Streaming docs | LLM tokens |
custom 来自 get_stream_writer() | 官方契约 | LangGraph Streaming docs | Custom data |
| event streaming 提供 typed projections | 官方契约 | Event streaming docs | What event streaming provides |
| metadata/tags 可加入 trace | 官方契约 | LangGraph Observability docs | Add metadata to traces |
| LangSmith trace 通过环境变量启用 | 官方契约 | LangSmith docs | Enable 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_channels、output_channels、channels 决定了 values 和 updates 能看到哪些 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 自己实现,意味着所有节点天然可以被统一观测,不需要为每个节点写一套日志协议。
源码证据
langgraph/graph/state.py::StateGraph.compilelanggraph/pregel/main.py::Pregel.stream- https://github.com/langchain-ai/langgraph/blob/1.2.7/libs/langgraph/langgraph/pregel/main.py
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 噪声增加 | 应规范命名 |
| 子调用未传 config | Python 低版本 async 可能丢 streaming context | 显式传 config |
设计原因与工程影响
观测信息不应污染业务 state,因此通过 config 传播,而不是写进 state。这也是区分“业务数据”和“观测数据”的关键。
源码证据
langchain_core/runnables/config.py::RunnableConfiglangchain_core/runnables/config.py::ensure_config- https://github.com/langchain-ai/langchain/blob/master/libs/core/langchain_core/runnables/config.py
7.5 构建期产物
| 产物 | 保存的信息 | 运行时用途 |
|---|---|---|
CompiledStateGraph | 节点、边、channel、stream channel | 提供 stream / astream |
RunnableConfig | tags、metadata、callbacks、configurable | 传递观测上下文 |
CallbackManager | callback handlers | 生命周期事件记录 |
StreamProtocol | stream mode、writer、event sink | 输出 chunks |
LangSmith tracer | trace sink | 上传 run tree |
get_stream_writer() context | 当前 stream writer | 节点内发 custom 数据 |
8. 运行时主链源码解剖
本章回答:
构建完成后,一次
stream / astream / stream_events如何进入核心执行逻辑并产生可观察输出?
8.1 运行时入口
| 调用方式 | 公开入口 | 核心内部入口 | 返回类型 |
|---|---|---|---|
| 同步最终结果 | invoke() | stream() 聚合 / Pregel loop | final state |
| 异步最终结果 | ainvoke() | astream() 聚合 / async loop | final state |
| 同步流式 | stream() | Pregel.stream() | Iterator[StreamPart] |
| 异步流式 | astream() | Pregel.astream() | AsyncIterator[StreamPart] |
| 事件流 | stream_events() | event router + transformers | run stream object |
| 异步事件流 | astream_events() | async event router | async 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 或恢复命令 |
| Config | RunnableConfig | tags、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_error | stream 抛异常 |
| interrupt | yield interrupt 相关数据 | 等待 resume |
| callbacks 后台提交失败 | 可能影响 trace | 不应影响业务结果 |
设计原因与工程影响
stream() 没有改变图执行语义,只改变结果的暴露方式。这意味着同一张图可以既用于 API 最终响应,也用于调试、前端流式体验和线上观测。
源码证据
langgraph/pregel/main.py::Pregel.streamlanggraph/pregel/loop.pylanggraph/pregel/runner.py- https://github.com/langchain-ai/langgraph/blob/1.2.7/libs/langgraph/langgraph/pregel/main.py
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.invokelanggraph/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_request 把 days 写错,updates 能直接指出哪个节点输出了错误字段。
源码证据
langgraph/pregel/main.py::Pregel.stream- LangGraph Streaming docs:
updatesstreams 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:
valuesstreams 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 未传 config | context 可能丢失 | messages 流异常 |
设计原因与工程影响
token streaming 是用户体验层能力;node update streaming 是流程调试能力。二者不要混为一谈。
源码证据
langgraph/pregel/messages.pylangchain_core/callbacks/manager.py- LangGraph Streaming docs:
messagesstreams(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 不含 custom | custom 数据可能不被消费 | 不影响业务结果 |
| async Python < 3.11 | contextvar 可能不可用 | 需手动传 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 object | typed 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.messages 和 stream.values,互不消耗彼此数据。
正常路径
raw Pregel events ↓event router ↓transformers ↓typed projections关键分支与异常路径
| 条件 | 行为 | 结果 |
|---|---|---|
| 应用只关心 token | 读取 stream.messages | 简洁 |
| 应用关心状态 | 读取 stream.values | 简洁 |
| 应用关心中断 | 读取 stream.interrupts | HITL 友好 |
| 多消费者 | 同时消费 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 event | start/end/error/token |
| 输出 | handler side effect | log / 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.pylangchain_core/runnables/config.pylangchain_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 projections9. 关键分支、异常与边界
本章回答:
当输出模式、执行方式或观测目标不同,框架如何分流、恢复、终止或失败?
9.1 分支矩阵
| 分支类型 | 触发条件 | 核心函数 | 结果 |
|---|---|---|---|
| 完整状态流 | stream_mode="values" | Pregel.stream() | 每步完整 state |
| 节点增量流 | stream_mode="updates" | Pregel.stream() | 每节点 update |
| token 流 | stream_mode="messages" | ChatModel callbacks / stream hooks | token + metadata |
| 自定义流 | stream_mode="custom" | get_stream_writer() | 业务事件 |
| 多模式流 | stream_mode=[...] | stream protocol | 多类 chunk |
| typed event flow | stream_events() | event router | projections |
| 子图流 | subgraphs=True | Pregel namespace | root + subgraph chunks |
| checkpoint 流 | stream_mode="checkpoints" | checkpointer | StateSnapshot-like events |
| task 流 | stream_mode="tasks" | task runner | start/finish/error |
| debug 流 | stream_mode="debug" | combined debug | 高噪声细节 |
9.2 同步与异步分支
| 维度 | 同步路径 | 异步路径 |
|---|---|---|
| 入口 | stream() | astream() |
| 返回 | Iterator | AsyncIterator |
| 节点函数 | sync callable | async callable / mixed |
| 模型调用 | invoke() / stream() | ainvoke() / astream() |
| 配置传播 | ContextVar | Python < 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 | 输出结构错误 | 业务动作已执行 | 不应重复副作用 |
| Resume | interrupt 暂停 | 普通异常 | 同 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 + UI | LangGraph 只输出 token chunks |
| 节点调试 | updates / tasks / trace | stream 暴露数据,不分析原因 |
| 线上监控 | LangSmith / metrics system | stream 不是长期监控数据库 |
| 持久恢复 | checkpointer | stream 不保存状态 |
| 日志审计 | logger / trace / DB | callbacks 只是事件入口 |
| 成本统计 | model usage metadata / LangSmith | stream 不自动计费 |
| 用户可见输出 | 应用层选择 | 不要直接展示所有 chunks |
10. 扩展机制与框架协作
本章回答:
框架允许在哪里插入自定义观测行为,它如何与相邻模块协作,哪些内部实现不应该被业务代码依赖?
10.1 扩展点总览
| 扩展点 | 扩展方式 | 执行时机 | 可修改内容 | 约束 |
|---|---|---|---|---|
stream_mode | 调用参数 | graph 运行时 | 输出类型 | 不改变业务 state |
get_stream_writer() | 节点 / 工具调用 | 节点内部 | custom stream data | 数据应可序列化 |
callbacks | RunnableConfig | start/end/error/token | 记录日志/trace | 不应改变业务结果 |
metadata | RunnableConfig | 调用开始 | trace 属性 | 避免敏感信息 |
tags | RunnableConfig / model config | 子调用传播 | 过滤 token/trace | 命名规范 |
| LangSmith env vars | 环境变量 | 全局运行时 | trace 上传 | 注意成本和隐私 |
tracing_context | context manager | 局部代码块 | 启停 trace / project | 只影响观测 |
| custom stream transformers | event 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", }, },)启用环境变量:
export LANGSMITH_TRACING=trueexport LANGSMITH_API_KEY=<your-api-key>export LANGSMITH_PROJECT=travel-agent调用方与被调用方
graph.invoke / stream ↓RunnableConfig ↓CallbackManager ↓LangSmith tracer ↓LangSmith project trace输入、输出与状态变化
| 项目 | 类型 | 说明 |
|---|---|---|
| 输入 | tags / metadata / callbacks | trace 属性 |
| 输出 | trace run tree | LangSmith 可视化 |
| 状态变化 | 无业务 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 可以通过环境变量自动接入,不需要重写业务图。
tags 和 metadata 会进入 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.pylangchain_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 updatesLLM messages/tokenstool callstool resultscustom progressinterruptscheckpointserrorsmetadatausagelatencycostLangGraph 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 打字机效果 | messages 或 stream_events.messages | 可流式展示 token |
| 工具进度显示 | custom | 节点/工具主动发进度 |
| 复杂子图排障 | subgraphs=True | 可观察子图内部 |
| HITL 前端 | stream_events.interrupts | 直接拿 pending interrupt |
| 线上 trace | LangSmith | 持久化、可搜索、可评估 |
| 成本/延迟分析 | LangSmith + usage metadata | stream 只实时输出,不做聚合 |
| 普通后端接口最终结果 | invoke() | 不需要中间过程 |
| 高敏数据生产环境 | 谨慎 stream values/debug | 可能泄露完整 state |
11.2 工程决策表
| 决策点 | 推荐选择 | 前提 | 风险 |
|---|---|---|---|
| 调试 state | updates 优先 | 想看哪个节点写错 | 不是完整 state |
| 调试 reducer | values | state 不太大 | 可能泄露敏感字段 |
| 前端 token | messages / stream_events.messages | 模型支持 streaming | 多模型 token 交错 |
| 业务进度 | custom | 节点有明确阶段 | 不要写过大数据 |
| 应用层事件 | stream_events() | 需要多个 typed projections | 学习成本略高 |
| 线上 trace | LangSmith | 可接受数据上传 | 隐私与成本 |
| metadata | 放 user/session/env/version | JSON-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 | 正在解析需求、正在检索景点、正在生成计划 |
| 用户看到回答生成 | messages | final_response 节点 token |
| 开发调试节点 | updates | intent、search_results、draft_plan、final_plan |
| 调试 reducer | values | research_notes、messages、step_results |
| 调试子图 | subgraphs=True | planner / executor subgraph |
| 审批 / interrupt | stream_events.interrupts | 待确认计划或工具动作 |
| 线上排错 | LangSmith | 完整 run tree、metadata、错误 |
| 成本分析 | LangSmith + usage_metadata | model、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 chunks13.3 分支与异常心智模型
values → 完整 state 快照updates → 节点 partial updatemessages → LLM token + metadatacustom → 节点/工具业务事件checkpoints → checkpoint snapshottasks → task lifecycledebug → 高噪声调试events → typed projections
普通异常 → on_error + stream 抛错interrupt → pending interrupt + resumecallback 失败 → 应隔离,不破坏业务13.4 一句话总结
Streaming 把 LangGraph 的 Pregel 执行过程投影为状态、节点、模型 token、自定义进度和调试事件;callbacks 与 metadata 把同一执行链交给 tracer,LangSmith 则把这些运行步骤持久化为可搜索、可复盘、可评估的 Agent trace。
13.5 掌握检查
- 能说清
invoke()和stream()的关系。 - 能解释
updates与values的区别。 - 能解释
messagesstream 输出的不是 state[“messages”]。 - 能写出
customprogress event。 - 能用 metadata 过滤某个节点或模型调用的 token。
- 能解释
stream_events()相比stream()的优势。 - 能解释 callbacks 与 LangSmith trace 的关系。
- 能说明 metadata / tags 不应放敏感数据。
- 能为旅行规划助手设计实时 UI stream 和后台 trace。
- 能根据调试目标选择正确 stream mode。
14. 参考资料与下一篇衔接
14.1 官方概念文档
-
LangGraph Streaming
https://docs.langchain.com/oss/python/langgraph/streaming -
LangGraph Event Streaming
https://docs.langchain.com/oss/python/langgraph/event-streaming -
LangGraph Observability
https://docs.langchain.com/oss/python/langgraph/observability -
LangSmith Trace with LangChain
https://docs.langchain.com/langsmith/trace-with-langchain
14.2 官方 API Reference
-
Pregel.stream
https://reference.langchain.com/python/langgraph/pregel/main/Pregel/stream -
Pregel.astream
https://reference.langchain.com/python/langgraph/pregel/main/Pregel/astream -
Runnable.astream_events
https://reference.langchain.com/python/langchain-core/runnables/base/Runnable/astream_events -
RunnableConfig
https://reference.langchain.com/python/langchain-core/runnables/config/RunnableConfig -
get_stream_writer
https://reference.langchain.com/python/langgraph/config/get_stream_writer
14.3 官方源码
-
langgraph/pregel/main.py::Pregel.stream
https://github.com/langchain-ai/langgraph/blob/1.2.7/libs/langgraph/langgraph/pregel/main.py -
langgraph/pregel/main.py::Pregel.astream
https://github.com/langchain-ai/langgraph/blob/1.2.7/libs/langgraph/langgraph/pregel/main.py -
langgraph/config.py::get_stream_writer
https://github.com/langchain-ai/langgraph/blob/1.2.7/libs/langgraph/langgraph/config.py -
langchain_core/runnables/base.py::Runnable.astream_events
https://github.com/langchain-ai/langchain/blob/master/libs/core/langchain_core/runnables/base.py -
langchain_core/runnables/config.py::RunnableConfig
https://github.com/langchain-ai/langchain/blob/master/libs/core/langchain_core/runnables/config.py -
langchain_core/callbacks/manager.py
https://github.com/langchain-ai/langchain/blob/master/libs/core/langchain_core/callbacks/manager.py -
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 读取?长期记忆写入策略如何工程化?记忆如何避免污染上下文?