14406 字
72 分钟

LangGraph 源码深潜:StateGraph 与 Workflow Agent 运行机制解剖

核心问题: LangGraph 如何把用户定义的 State、节点函数和边关系,编译成一个可执行、可流式、可批量、可持久化的 CompiledStateGraph,并在运行时通过 partial state update 推进 Workflow Agent?

源码主线: StateGraph(...) → add_node() / add_edge() → compile() → CompiledStateGraph(Pregel) → invoke() / stream() → node returns partial update → channel reducer 合并 state

前置文章: 第 1 篇 Runnable 源码解剖;第 2 篇 Prompt / Message / ChatModel 源码解剖;第 3 篇 OutputParser 与结构化输出;第 4 篇 Tool Calling 源码解剖;第 5 篇 create_agent 与 ReAct 运行机制源码解剖

依赖基线: langgraph==1.2.7langchain-core 使用该版本依赖解析结果

源码基线: https://github.com/langchain-ai/langgraph/tree/1.2.7;PyPI sdist langgraph-1.2.7.tar.gz SHA256:dcdf5b441bf8c7c7c154e603b302c9dbfbfd2d11e1b7ae7d93a5aba979dc87bd

阅读边界: 本文覆盖 StateGraph 的构建期、编译期、运行时主链、partial update 合并、边与终止条件、扩展机制;不深入 add_conditional_edges 的复杂路由、不展开 checkpoint 持久化、不展开 interrupt / human-in-the-loop,这些放到后续专题。


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

前几篇已经完成 LangChain / LangGraph Agent 基础组件的源码认知:

Runnable 统一执行协议
Prompt / Message / ChatModel 消息调用链
OutputParser / structured output 结构化结果
Tool Calling 工具调用协议
create_agent 预构建 ReAct Agent 图

本篇开始进入 LangGraph 的底层编排核心:

Workflow Agent 范式
StateGraph 构建状态图
compile 生成 CompiledStateGraph
Pregel runtime 执行节点、边和状态更新

本篇只解决:

  • StateGraph 为什么是 builder,而不是执行器。
  • State schema 如何被解析成 channels 与 reducer。
  • add_node() 如何保存节点函数、输入 schema 和节点策略。
  • add_edge() 如何保存固定图结构。
  • compile() 如何把 builder 转换为可执行 CompiledStateGraph
  • CompiledStateGraph 为什么能 invoke / stream / batch / ainvoke
  • 节点返回 partial update 后,LangGraph 如何合并到全局 state。

本篇不展开:

  • add_conditional_edges() 的动态路由与 BranchSpec 细节。
  • Command(goto=...)Send、子图和多分支并行高级语义。
  • checkpoint、durable execution、interrupt、time travel。
  • LangGraph Platform 部署和 LangSmith 可观测性。

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

1.1 核心问题#

LangGraph 如何把“用户写的一组普通 Python 节点函数 + 固定边关系 + 共享状态 schema”,编译成一个可以按 Workflow 顺序执行、自动合并节点返回值、并在到达 END 后产出最终 state 的运行时对象?

这个问题对应你提出的第 6 篇主线:

StateGraph 源码解剖
Workflow Agent
StateGraph + Node + Edge + State

1.2 学习目标#

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

  1. 解释 StateGraph 的构建期对象模型,包括 nodesedgeschannelsschemaswaiting_edgesbranches
  2. 解释 State schema 如何映射成 state key、channel、默认覆盖 reducer 或自定义 reducer。
  3. 解释 add_node() 如何把普通函数、Runnable 或节点名归一化成 StateNodeSpec
  4. 解释 add_edge() 如何注册 START → nodenode → nodenode → END 的固定转移。
  5. 解释 compile() 为什么必须存在,以及它如何校验图、构造 CompiledStateGraph、attach 节点与边。
  6. 解释 CompiledStateGraph 为什么继承 / 基于 Pregel runtime,并因此实现 Runnable 风格调用。
  7. 解释一次 graph.invoke(initial_state) 的运行时主链:输入 state 进入图、节点运行、partial update 写入 channel、state 合并、边触发下一节点、最终到达 END
  8. 区分公共 API、扩展接口和内部实现,知道哪些对象可以依赖,哪些只能作为源码理解。

1.3 能力边界#

能力本篇是否覆盖说明
StateGraph 构建期重点讲 __init___add_schemaadd_nodeadd_edgecompile
CompiledStateGraph 运行时讲它作为 Pregel 应用如何执行图,但不逐行复现 Pregel 全部调度细节。
partial state update重点讲节点返回 dict 后如何通过 channel / reducer 合并。
固定边覆盖普通 add_edge,含 STARTEND
条件边只在边界说明中提及,后续单独写 Conditional Edge 与 Router
checkpoint / persistence只解释 compile 可接收 checkpointer,不展开持久化。
interrupt / HITL后续单独写。
Pregel 完整算法部分只解释与 StateGraph 运行相关的 super-step、channels、active node、halt。

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

2.1 一句话定义#

StateGraph 是 LangGraph 的图构建器,负责把 state schema、节点函数和边关系登记成可编译的图定义;它不负责直接执行图,真正执行发生在 compile() 返回的 CompiledStateGraph / Pregel runtime 中。

官方 API Reference 明确说明:StateGraph 是 builder class,不能直接执行,必须先调用 .compile() 生成支持 invoke()stream()astream()ainvoke() 等方法的可执行图。

2.2 最小心智模型#

TypedDict / Pydantic State schema
StateGraph builder
add_node 注册节点函数
add_edge 注册控制流
compile 生成 CompiledStateGraph
invoke 初始化 state
节点返回 partial update
channel / reducer 合并 state
到达 END 输出最终 state

2.3 核心术语#

术语源码对象语义不要误解为
StateTypedDict / BaseModel / dataclass schema图运行过程中的共享状态结构普通全局变量
StateGraphStateGraph构建器,保存节点、边、schema、channel 定义执行器
NodeStateNode / StateNodeSpec接收 state,返回 partial update 的计算单元必须是 LLM 调用
Edgeedges / waiting_edges节点之间的固定转移关系业务数据流本身
STARTSTART / __start__图入口虚拟节点用户自定义节点
ENDEND / __end__图终止虚拟节点一个会执行的 Python 函数
ChannelBaseChannel / LastValue / BinaryOperatorAggregatestate key 的读写与合并通道单纯 dict 字段
partial updatedict[str, Any]节点返回的局部状态更新完整 state 替换
CompiledStateGraphCompiledStateGraph编译后的可执行图仍可随意修改的 builder
PregelPregel底层 message-passing graph runtime只针对 LLM 的循环器

2.4 与相邻抽象的边界#

对象负责什么不负责什么与本篇对象的关系
StateGraph定义 state、节点、边、schema直接运行、持久化、调度每个 super-step本篇构建期主角
CompiledStateGraph暴露 invoke / stream / batch / ainvoke让用户继续随意添加节点本篇运行时入口
Pregel底层图运行、channels、super-step 调度业务语义解释CompiledStateGraph 基于它执行
Runnable统一调用协议图结构定义CompiledStateGraph 通过 Pregel 具备 Runnable 风格接口
create_agent预构建 model/tools Agent 图定义任意 Workflow 细节它内部也依赖 LangGraph;本篇是更底层的图 API

3. 完整执行链路#

3.1 高层链路#

观察目标:先不看源码细节,只看用户 API 到最终输出之间的对象流转。

from typing_extensions import TypedDict
from langgraph.graph import StateGraph, START, END
class TravelState(TypedDict):
user_request: str
destination: str | None
days: int | None
plan: str | None
def parse_request(state: TravelState):
return {
"destination": "东京",
"days": 5,
}
def generate_plan(state: TravelState):
return {
"plan": f"为 {state['destination']} 生成 {state['days']} 天旅行计划"
}
builder = StateGraph(TravelState)
builder.add_node("parse_request", parse_request)
builder.add_node("generate_plan", generate_plan)
builder.add_edge(START, "parse_request")
builder.add_edge("parse_request", "generate_plan")
builder.add_edge("generate_plan", END)
graph = builder.compile()
result = graph.invoke({
"user_request": "我想去东京玩 5 天",
"destination": None,
"days": None,
"plan": None,
})
print(result)

高层链路:

用户初始 state
CompiledStateGraph.invoke()
START 激活 parse_request
parse_request 读取 state,返回 partial update
合并 destination / days
固定边触发 generate_plan
generate_plan 读取合并后的 state,返回 plan
合并 plan
固定边到 END
输出最终 state

3.2 对象流转#

阶段输入类型核心函数输出类型状态变化
构造 buildertype[TravelState]StateGraph.__init__()StateGraph解析 schema,初始化 channels / nodes / edges。
注册节点str + CallableStateGraph.add_node()Self保存 StateNodeSpecself.nodes
注册边START/node + node/ENDStateGraph.add_edge()Self保存固定边到 self.edges 或等待边结构。
编译StateGraphStateGraph.compile()CompiledStateGraph校验结构,生成 Pregel channels / nodes / triggers。
调用dictCompiledStateGraph.invoke()dict初始化输入 channel,执行图,返回输出 state。
节点执行TravelStateparse_request()dict返回 partial update。
状态合并partial updatechannel update / reducerchannel value更新 destinationdaysplan

3.3 时序链路#

Caller
│ builder = StateGraph(TravelState)
StateGraph builder
│ add_node / add_edge
Graph definition
│ compile()
CompiledStateGraph / Pregel
│ invoke(initial_state)
Pregel runtime
│ execute nodes and merge partial updates
Final state dict
│ return
Caller

3.4 正常结束条件#

一次 StateGraph 执行的正常结束不是“某个函数 return 了最终答案”,而是图运行时满足终止语义:

所有需要执行的节点都已经完成
没有新的边消息需要传递
执行到 END 或所有节点 inactive 且无消息在传递
Pregel runtime 产出 output channel 对应的最终 state

在简单固定边 Workflow 中,可以把它理解为:

START → parse_request → generate_plan → END

到达 END 后,最终输出 schema 对应的 state 被返回。


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

4.1 核心目录#

langgraph/
├── graph/
│ ├── state.py # StateGraph / CompiledStateGraph 主入口
│ ├── graph.py # 更基础的 Graph 抽象与边语义
│ ├── _node.py # StateNode / StateNodeSpec
│ └── _branch.py # BranchSpec,条件边相关
├── channels/
│ ├── base.py # BaseChannel
│ ├── last_value.py # LastValue 默认覆盖语义
│ ├── binop.py # BinaryOperatorAggregate reducer 语义
│ └── ephemeral_value.py
├── pregel/
│ ├── main.py # Pregel runtime 主类
│ ├── _read.py # PregelNode / ChannelRead
│ └── _write.py # ChannelWrite / ChannelWriteEntry
└── constants.py # START / END 等常量

4.2 关键文件#

优先级文件核心对象阅读目的
1langgraph/graph/state.pyStateGraphCompiledStateGraph理解 builder、schema、compile、attach node/edge。
2langgraph/channels/last_value.pyLastValue理解默认 state key 覆盖语义。
3langgraph/channels/binop.pyBinaryOperatorAggregate理解 Annotated[..., reducer] 的聚合语义。
4langgraph/pregel/main.pyPregel理解编译图运行时为什么能 invoke / stream / batch。
5langgraph/pregel/_read.pyPregelNode理解节点如何读取 channel。
6langgraph/pregel/_write.pyChannelWrite理解节点输出如何写入 channel。
7langgraph/constants.pySTARTEND理解虚拟入口和终止符。

4.3 推荐阅读顺序#

1. StateGraph.__init__
2. StateGraph._add_schema
3. _get_channels / _is_field_binop / _is_field_channel
4. StateGraph.add_node
5. StateGraph.add_edge
6. StateGraph.compile
7. CompiledStateGraph.attach_node / attach_edge
8. Pregel.invoke / Pregel.stream
9. Channel update / reducer 实现

4.4 不建议的阅读顺序#

不建议直接从 Pregel.stream() 或 Pregel loop 细节开始读。

原因是 Pregel runtime 是通用 message-passing 执行器,里面有大量并发、checkpoint、interrupt、stream、retry、debug、callback 逻辑。若没有先理解 StateGraph 如何把 state key 转成 channels、如何把 node 转成 PregelNode、如何把 edge 转成 channel trigger,就会误以为 LangGraph 是一个普通 while 循环。


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

5.1 核心对象关系#

用户 State schema
StateGraph builder
CompiledStateGraph
Pregel runtime
Runnable-style invoke / stream / batch

更贴近源码的关系:

StateGraph
保存:nodes / edges / branches / schemas / channels / managed
↓ compile()
CompiledStateGraph
保存:Pregel nodes / channels / input_channels / output_channels / stream_channels
↓ 运行时
Pregel
调度:super-step / channel reads / channel writes / active nodes / termination

5.2 对象职责#

对象生命周期输入输出核心职责
StateGraph构建期state schema、节点、边builder 自身收集图定义,不执行。
StateNodeSpec构建期到编译期节点函数 / Runnable节点规格保存 node runnable、metadata、input_schema、retry/cache/timeout 等策略。
BaseChannel编译期到运行时state key updatechannel value管理状态字段如何更新。
LastValue运行时单个新值最新值默认覆盖式 state 更新。
BinaryOperatorAggregate运行时多个 update聚合值reducer 聚合,例如 list 追加。
CompiledStateGraph运行时初始 input state输出 state可执行图对象。
Pregel运行时channels / nodes / triggersstream / final output底层图调度引擎。

5.3 协议边界#

StateGraph builder 协议
负责:注册节点、注册边、解析 schema、编译图。
不负责:每次运行中的节点调度、checkpoint 写入、stream 输出。
Node 协议
负责:接收当前 state,返回 partial state update。
不负责:自己决定全局 state 如何合并,除非返回 Command 等高级对象。
Channel 协议
负责:决定某个 state key 如何应用更新。
不负责:理解业务含义。
Pregel runtime 协议
负责:按 channel 消息驱动节点执行,进行 super-step 调度。
不负责:定义业务节点和状态字段。

5.4 稳定接口与内部实现#

类型对象文章中的使用原则
公共 APIStateGraphadd_node()add_edge()compile()invoke()可以用于工程示例和项目代码。
扩展接口reducer annotation、context_schemainput_schemaoutput_schema、node retry/cache/timeout可用于工程扩展,但要遵守官方契约。
内部实现StateNodeSpec_get_channels()_get_updates()attach_node()、Pregel 内部 loop用于理解源码,不建议业务代码依赖。

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

6.1 本篇阅读策略#

先读公开入口 StateGraph
确认 builder 保存哪些字段
读 schema 如何变成 channels
读 add_node / add_edge 如何收集图结构
读 compile 如何构造 CompiledStateGraph
读运行时 invoke 如何进入 Pregel
读节点返回 partial update 如何写入 channel
回到 Workflow Agent 设计目的

6.2 证据等级#

标记含义写作要求
源码事实可由 langgraph==1.2.7 源码证明给出 tag / 文件 / 符号。
官方契约官方文档或 API Reference 明确承诺给出官方文档链接。
简化伪代码对真实控制流的压缩表达明确标注“伪代码”。
作者推断根据调用链得出的设计理解使用“从调用关系可以推断”。
工程建议面向项目实践的建议说明适用条件。

6.3 本篇证据清单#

结论证据类型文件或文档定位
StateGraph 是 builder,需 compile 后执行官方契约 / 源码文档API Reference;state.py docstringStateGraph 类说明。
图由 State、Nodes、Edges 三部分组成官方契约Graph API overviewGraphs / StateGraph 章节。
编译会进行结构校验并可配置运行时参数官方契约Graph API overviewCompiling your graph。
StateGraph 编译后自动创建 Pregel application官方契约LangGraph runtime docsHigh-level API / StateGraph。
langgraph==1.2.7 是当前 PyPI 最新稳定版官方发布信息PyPIRelease / file details。
state key 默认按 channel 更新,reducer 可由 schema 注解指定官方契约 / 源码事实Graph API overview;state.pyState / Reducers;_add_schema

7. 构建期源码解剖#

本章回答:

用户传入的 State schema、节点函数和边关系,如何被转换、校验、归一化并保存为可编译的图定义?

7.1 构建期职责#

输入归一化动作构建结果
TravelState解析 state keys、channels、managed valuesself.schemasself.channelsself.managed
parse_request推断节点名、coerce 为 runnable、保存输入 schemaself.nodes["parse_request"]
START → parse_request校验起点和终点,保存固定边self.edges
generate_plan → END识别终止边self.edges 中的终止转移

7.2 构建期总链路#

StateGraph(TravelState)
初始化 builder 字段
_add_schema(state_schema)
_get_channels(schema)
保存 channels / managed
add_node(name, action)
coerce_to_runnable(action)
保存 StateNodeSpec
add_edge(start, end)
保存 fixed edges

7.3 StateGraph.__init__() 源码解剖#

职责与所处阶段#

StateGraph.__init__() 是构建期入口,负责把用户传入的 state schema、context schema、input schema、output schema 转换为 builder 内部结构。

真实源码签名#

以下签名来自 langgraph==1.2.7langgraph/graph/state.py

def __init__(
self,
state_schema: type[StateT],
context_schema: type[ContextT] | None = None,
*,
input_schema: type[InputT] | None = None,
output_schema: type[OutputT] | None = None,
**kwargs: Unpack[DeprecatedKwargs],
) -> None:
...

调用方与被调用方#

用户代码 StateGraph(TravelState)
StateGraph.__init__
StateGraph._add_schema(state_schema)
_get_channels(schema)

输入、输出与状态变化#

项目类型说明
输入type[StateT]用户定义的 TypedDict / dataclass / Pydantic state。
输出None构造对象本身,不直接返回图。
状态变化self.nodesself.edgesself.schemasself.channels初始化 builder 内部容器。
副作用deprecated 参数 warning老参数如 config_schema 会触发警告。

细粒度伪代码#

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

def __init__(state_schema, context_schema=None,
input_schema=None, output_schema=None, **kwargs):
# 1. 处理 deprecated 参数
if "config_schema" in kwargs and context_schema is None:
warn_deprecated("config_schema")
context_schema = kwargs["config_schema"]
if "input" in kwargs and input_schema is None:
warn_deprecated("input")
input_schema = kwargs["input"]
if "output" in kwargs and output_schema is None:
warn_deprecated("output")
output_schema = kwargs["output"]
# 2. 初始化 builder 容器
self.nodes = {}
self.edges = set()
self.branches = defaultdict(dict)
self.schemas = {}
self.channels = {}
self.managed = {}
self.waiting_edges = set()
self.compiled = False
# 3. 保存 schema 边界
self.state_schema = state_schema
self.input_schema = input_schema or state_schema
self.output_schema = output_schema or state_schema
self.context_schema = context_schema
# 4. 初始化节点默认策略
self._node_defaults = _NodeDefaults()
# 5. 解析 state / input / output schema
self._add_schema(self.state_schema)
self._add_schema(self.input_schema, allow_managed=False)
self._add_schema(self.output_schema, allow_managed=False)

逐段解释#

第 1 段处理历史参数兼容。它不影响本篇主链,但说明 LangGraph 在正式版中推荐使用 context_schemainput_schemaoutput_schema,不再推荐旧的 config_schemainputoutput

第 2 段初始化 builder 容器。这里是判断 StateGraph 不是执行器的关键证据:它只是保存节点、边、分支、schema、channel 等定义,没有启动任何运行时循环。

第 3 段确定 schema 边界。如果用户没有显式提供 input_schemaoutput_schema,默认输入输出 schema 都等于 state schema。这解释了最小示例里 graph.invoke(initial_state) 和最终 result 都使用 TravelState 结构。

第 5 段调用 _add_schema(),这是 state schema 被转成 channels 的入口。

关键分支#

条件分支动作构建结果
未提供 input_schema使用 state_schema输入必须符合完整 state schema。
未提供 output_schema使用 state_schema输出返回完整 state。
传入 config_schema警告并映射到 context_schema兼容旧代码。
input/output schema 含 managed values抛错Input/Output schema 不允许 managed channels。

设计原因与工程影响#

构建期必须先解析 schema,因为后续 add_node()compile() 和运行时 state 合并都依赖 state key 对应的 channel。若等到运行时才解析 schema,每次调用都会重复处理类型信息,并且无法在 compile 阶段做结构校验。

源码证据#

  • langgraph/graph/state.py::StateGraph.__init__
  • https://github.com/langchain-ai/langgraph/blob/1.2.7/libs/langgraph/langgraph/graph/state.py

7.4 _add_schema() 与 channel 生成源码解剖#

职责与所处阶段#

_add_schema() 负责把 TypedDict / Pydantic / dataclass 等 schema 解析成 LangGraph 内部的 channels 和 managed values。

真实源码签名#

def _add_schema(
self,
schema: type[Any],
/,
allow_managed: bool = True,
) -> None:
...

调用方与被调用方#

StateGraph.__init__
StateGraph._add_schema
_get_channels(schema)
_get_channel / _is_field_binop / _is_field_channel

输入、输出与状态变化#

项目类型说明
输入schema: type[Any]state / input / output schema。
输出None不返回值。
状态变化self.schemasself.channelsself.managed保存 state key 到 channel 的映射。
副作用schema 冲突时抛错同名 channel 类型不兼容会失败。

细粒度伪代码#

def _add_schema(schema, allow_managed=True):
# 1. 避免重复解析同一个 schema
if schema in self.schemas:
return
# 2. 对非法 schema 给出警告
warn_if_schema_not_type_or_annotated(schema)
# 3. 从 schema 提取 channels 与 managed values
channels, managed, type_hints = _get_channels(schema)
# 4. input/output schema 不允许 managed values
if managed and not allow_managed:
raise ValueError("Managed channels are not permitted")
# 5. 记录 schema 对应的所有字段定义
self.schemas[schema] = {**channels, **managed}
# 6. 合并普通 channels
for key, channel in channels.items():
if key in self.channels:
if self.channels[key] != channel:
if isinstance(channel, LastValue):
# 默认 LastValue 可以兼容已有定义
pass
else:
raise ValueError("Channel already exists with different type")
else:
self.channels[key] = channel
# 7. 合并 managed values
for key, managed_value in managed.items():
if key in self.managed and self.managed[key] != managed_value:
raise ValueError("Managed value already exists with different type")
self.managed[key] = managed_value

逐段解释#

第 1 段说明 schema 是构建期静态信息,重复 schema 不需要重复解析。

第 3 段是核心:_get_channels(schema) 会读取类型注解,并决定每个 state key 使用哪种 channel。例如普通字段通常对应 LastValue,带 reducer 的 Annotated 字段可能对应 BinaryOperatorAggregate

第 4 段体现 schema 边界:managed values 是运行时管理值,不应作为外部输入或最终输出 schema 的普通字段。

第 6 段解决多 schema 场景下同名字段的 channel 合并问题。如果不同 schema 对同一字段定义了不兼容 channel,构建期会失败,避免运行时出现隐式状态冲突。

正常路径#

TravelState TypedDict
_get_channels
user_request → LastValue(str)
destination → LastValue(str | None)
days → LastValue(int | None)
plan → LastValue(str | None)
self.channels 保存字段通道

关键分支与异常路径#

条件行为结果
schema 首次出现解析并保存channels 可用于 compile。
schema 已解析直接返回避免重复工作。
managed 出现在 input/output schemaValueError保持外部 IO schema 简洁。
同名 channel 类型冲突ValueError防止 state 合并语义不一致。

设计原因与工程影响#

LangGraph 的 state 不是普通 dict。每个字段背后都有 channel。channel 决定了节点返回 partial update 时,是覆盖旧值、追加列表、聚合多个并行结果,还是使用特殊 managed value。因此 _add_schema() 是理解 partial update 的第一入口。

源码证据#

  • langgraph/graph/state.py::StateGraph._add_schema
  • langgraph/graph/state.py::_get_channels
  • https://github.com/langchain-ai/langgraph/blob/1.2.7/libs/langgraph/langgraph/graph/state.py

7.5 add_node() 源码解剖#

职责与所处阶段#

add_node() 负责把用户传入的普通函数、Runnable 或带名称的 action 归一化成 StateNodeSpec,并保存到 builder 的 self.nodes 中。

真实源码签名#

add_node() 在源码中有多个 overload,实际实现可以概括为:

def add_node(
self,
node: str | StateNode[NodeInputT, ContextT],
action: StateNode[NodeInputT, ContextT] | None = None,
*,
defer: bool = False,
metadata: dict[str, Any] | None = None,
input_schema: type[NodeInputT] | None = None,
retry_policy: RetryPolicy | Sequence[RetryPolicy] | None = None,
cache_policy: CachePolicy | None = None,
error_handler: StateNode[Any, ContextT] | None = None,
destinations: dict[str, str] | tuple[str, ...] | None = None,
timeout: float | timedelta | TimeoutPolicy | None = None,
**kwargs: Unpack[DeprecatedKwargs],
) -> Self:
...

调用方与被调用方#

用户代码 builder.add_node("parse_request", parse_request)
StateGraph.add_node
coerce_to_runnable(action)
StateNodeSpec(...)
self.nodes[name] = spec

输入、输出与状态变化#

项目类型说明
输入str + CallableCallable节点名和节点逻辑。
输出Self支持链式调用。
状态变化self.nodes[node_name]保存节点规格。
副作用compiled 后再添加会 warning已编译图不会自动同步更新。

细粒度伪代码#

def add_node(node, action=None, *, input_schema=None,
retry_policy=None, cache_policy=None,
error_handler=None, defer=False,
metadata=None, destinations=None, timeout=None, **kwargs):
# 1. 处理 deprecated 参数
if "retry" in kwargs and retry_policy is None:
warn_deprecated("retry")
retry_policy = kwargs["retry"]
if "input" in kwargs and input_schema is None:
warn_deprecated("input")
input_schema = kwargs["input"]
# 2. 规范化 timeout policy
timeout = coerce_timeout_policy(timeout)
# 3. 如果第一个参数不是字符串,说明用户写的是 add_node(func)
if not isinstance(node, str):
action = node
if isinstance(action, Runnable):
node_name = action.get_name()
else:
node_name = getattr(action, "__name__", action.__class__.__name__)
else:
node_name = node
# 4. 基本校验
if node_name is None:
raise ValueError("Node name must be provided")
if self.compiled:
logger.warning("Adding node after compile will not affect compiled graph")
if node_name in {START, END}:
raise ValueError("START/END cannot be used as normal node name")
if node_name in self.nodes:
raise ValueError("Node already exists")
if node_name in self.channels:
raise ValueError("Node name conflicts with state channel")
# 5. 推断或注册节点输入 schema
resolved_input_schema = input_schema or self.state_schema
self._add_schema(resolved_input_schema)
# 6. 将函数 / Runnable 归一化为 Runnable
runnable = coerce_to_runnable(action, name=node_name)
# 7. 保存节点规格
self.nodes[node_name] = StateNodeSpec(
runnable=runnable,
metadata=metadata,
input_schema=resolved_input_schema,
retry_policy=retry_policy,
cache_policy=cache_policy,
error_handler=error_handler,
defer=defer,
destinations=destinations,
timeout=timeout,
)
return self

逐段解释#

第 3 段解释了两种常见写法:

builder.add_node(parse_request)

会从函数名推断节点名:

parse_request

而:

builder.add_node("parse_request", parse_request)

显式指定节点名。

第 4 段是构建期校验。STARTEND 是虚拟节点,不允许作为普通业务节点名。节点名也不能和 state channel 同名,否则运行时 channel 和 node 触发器会混淆。

第 6 段说明节点函数会被转成 Runnable 风格对象。这承接第 1 篇 Runnable 源码认知:LangGraph 节点不一定是裸函数,也可以是 Runnable、chain、subgraph。

第 7 段将所有节点信息保存为节点规格,而不是立即执行函数。

正常路径#

add_node("parse_request", parse_request)
node_name = "parse_request"
input_schema = TravelState
coerce_to_runnable(parse_request)
StateNodeSpec(...)
self.nodes["parse_request"] = spec

关键分支与异常路径#

条件行为结果
add_node(func)func.__name__ 推断节点名简化 API。
add_node("name", func)使用显式节点名更适合稳定图结构。
节点名为 START / END抛错保护虚拟节点语义。
节点名重复抛错防止覆盖。
已 compile 后添加节点warning旧 compiled graph 不自动改变。

设计原因与工程影响#

add_node() 把节点函数封装成规格对象,是为了把“用户业务逻辑”与“图运行时调度”分离。节点函数只需要遵守 State -> Partial[State] 协议,运行时如何读取 state、如何写 channel、如何触发下一节点,全部交给编译后的图处理。

源码证据#

  • langgraph/graph/state.py::StateGraph.add_node
  • langgraph/graph/_node.py::StateNodeSpec
  • https://github.com/langchain-ai/langgraph/blob/1.2.7/libs/langgraph/langgraph/graph/state.py

7.6 add_edge() 源码解剖#

职责与所处阶段#

add_edge() 负责保存固定控制流边,表示某个节点完成后,下一步应该激活哪个节点或到达 END

真实源码签名#

def add_edge(
self,
start_key: str | list[str],
end_key: str,
) -> Self:
...

调用方与被调用方#

用户代码 builder.add_edge(START, "parse_request")
StateGraph.add_edge
self.edges.add((START, "parse_request"))

输入、输出与状态变化#

项目类型说明
输入start_keyend_key起点和终点。
输出Self支持链式调用。
状态变化self.edgesself.waiting_edges保存固定边。
副作用非法边抛错如从 END 出发、指向 START

细粒度伪代码#

def add_edge(start_key, end_key):
# 1. 已编译后添加边只影响 builder,不影响旧 compiled graph
if self.compiled:
logger.warning("Adding edge after compile will not affect compiled graph")
# 2. 校验特殊节点方向
if start_key == END:
raise ValueError("END cannot be a start node")
if end_key == START:
raise ValueError("START cannot be an end node")
# 3. 多起点边:等待多个节点都完成后再触发 end
if isinstance(start_key, list):
for key in start_key:
if key == END:
raise ValueError("END cannot be a start node")
self.waiting_edges.add((tuple(start_key), end_key))
return self
# 4. 普通固定边
self.edges.add((start_key, end_key))
return self

逐段解释#

第 2 段保护 START / END 的方向语义:

START 只能作为入口起点
END 只能作为终点

所以:

builder.add_edge(END, "x")

是非法的。

第 3 段是高级语义:多个起点组成等待边,表示多个上游节点都到达后再触发下游节点。这通常用于 join / barrier 场景,本篇只说明边界,不深入展开。

第 4 段是最小代码使用的普通固定边。

正常路径#

add_edge(START, "parse_request")
add_edge("parse_request", "generate_plan")
add_edge("generate_plan", END)
self.edges = {
(START, "parse_request"),
("parse_request", "generate_plan"),
("generate_plan", END)
}

关键分支与异常路径#

条件行为结果
start_key == END抛错不能从终点继续执行。
end_key == START抛错不能回到入口虚拟节点。
start_key 是 list写入 waiting_edges等待多个上游。
普通字符串边写入 edges固定转移。

设计原因与工程影响#

固定边让 Workflow Agent 具备可控执行顺序。与 ReAct 的模型自由循环不同,Workflow Agent 更适合把确定业务流程写成显式边,降低循环不稳定性。

源码证据#

  • langgraph/graph/state.py::StateGraph.add_edge
  • langgraph/constants.py::START / END
  • https://github.com/langchain-ai/langgraph/blob/1.2.7/libs/langgraph/langgraph/graph/state.py

7.7 compile() 构建期产物源码解剖#

职责与所处阶段#

compile() 是构建期和运行时之间的分界线。它把 StateGraph builder 中保存的 schema、nodes、edges、branches 编译成可执行 CompiledStateGraph

真实源码签名#

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[StateT, ContextT, InputT, OutputT]:
...

调用方与被调用方#

用户代码 graph = builder.compile()
StateGraph.compile
validate graph
CompiledStateGraph(...)
compiled.attach_node / attach_edge / attach_branch
compiled.validate()

输入、输出与状态变化#

项目类型说明
输入builder 内部状态 + runtime options已注册的 nodes / edges / schema。
输出CompiledStateGraph可执行图。
状态变化self.compiled = True标记 builder 已被编译。
副作用结构校验失败时抛错缺少入口、孤立节点、非法边等。

细粒度伪代码#

def compile(checkpointer=None, *, cache=None, store=None,
interrupt_before=None, interrupt_after=None,
debug=False, name=None):
# 1. 规范化 interrupt 配置
interrupt_before = interrupt_before or []
interrupt_after = interrupt_after or []
# 2. 校验图结构
self.validate(
interrupt_before=interrupt_before,
interrupt_after=interrupt_after,
)
# 3. 标记 builder 已编译
self.compiled = True
# 4. 创建 CompiledStateGraph / Pregel 应用
compiled = CompiledStateGraph(
builder=self,
channels={
**self.channels,
START: EphemeralValue(self.input_schema),
},
input_channels=START,
output_channels=output_keys_from_schema(self.output_schema),
stream_channels=stream_keys_from_state_channels(self.channels),
checkpointer=checkpointer,
cache=cache,
store=store,
interrupt_before=interrupt_before,
interrupt_after=interrupt_after,
debug=debug,
name=name,
)
# 5. 附加 START 虚拟节点
compiled.attach_node(START, None)
# 6. 附加用户节点
for node_name, node_spec in self.nodes.items():
compiled.attach_node(node_name, node_spec)
# 7. 附加固定边
for start, end in self.edges:
compiled.attach_edge(start, end)
# 8. 附加等待边
for starts, end in self.waiting_edges:
compiled.attach_edge(starts, end)
# 9. 附加条件分支,本篇不展开
for start, branches in self.branches.items():
for name, branch in branches.items():
compiled.attach_branch(start, name, branch)
# 10. 校验 compiled graph 并返回
return compiled.validate()

逐段解释#

第 2 段是 compile() 必须存在的核心原因之一:builder 阶段允许用户逐步添加节点和边,而编译阶段必须集中检查图结构是否可运行。

第 4 段创建 CompiledStateGraph。这里会把 state schema 解析出的 self.channels 带入运行时,同时额外加入 START 入口 channel。

第 6 段将构建期的 StateNodeSpec 附加到 compiled graph。此时节点不再只是 builder 字典中的函数,而会被包装成 Pregel 可调度节点。

第 7 段把固定边转成 Pregel channel trigger / write 关系。

第 10 段返回的是 compiled graph,不是 builder 本身。这解释了为什么 StateGraph 不能直接 invoke()

正常路径#

builder.compile()
validate builder
CompiledStateGraph(...)
attach START
attach parse_request
attach generate_plan
attach START → parse_request
attach parse_request → generate_plan
attach generate_plan → END
return executable graph

关键分支与异常路径#

条件行为结果
缺少入口边validate 失败图不可执行。
边指向不存在节点validate 失败防止运行时找不到节点。
设置 checkpointer保存到 compiled graph后续支持持久化。
设置 interrupt保存到 compiled graph后续支持中断。
有 branchesattach_branch条件边后续处理。

设计原因与工程影响#

compile() 的设计使 LangGraph 能把用户友好的构建 API 与底层 Pregel runtime 解耦。构建期可以使用 StateGraph 这种接近业务语言的抽象;运行期则使用 channels、triggers、PregelNode 等更底层的执行结构。

源码证据#

  • langgraph/graph/state.py::StateGraph.compile
  • langgraph/graph/state.py::CompiledStateGraph
  • https://github.com/langchain-ai/langgraph/blob/1.2.7/libs/langgraph/langgraph/graph/state.py

7.8 构建期产物#

产物保存的信息运行时用途
self.channelsstate key 与 channel 类型决定 partial update 如何合并。
self.nodes节点名到 StateNodeSpeccompile 时生成 PregelNode。
self.edges固定边集合compile 时生成触发关系。
self.waiting_edges多起点等待边compile 时生成 barrier 触发。
CompiledStateGraph.nodesPregel 可调度节点运行时实际执行。
CompiledStateGraph.channelsPregel channel保存和传递 state / control messages。
input_channelsSTART接收图输入。
output_channelsoutput schema keys决定最终返回哪些 state 字段。

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

本章回答:

构建完成后,一次 graph.invoke() 如何进入底层执行逻辑,如何调用节点,如何合并 partial update,并最终输出 state?

8.1 运行时入口#

调用方式公开入口核心内部入口返回类型
同步调用graph.invoke(input)Pregel.invoke()Pregel.stream()dict / output schema
异步调用graph.ainvoke(input)Pregel.ainvoke()Pregel.astream()dict / output schema
流式调用graph.stream(input)Pregel.stream()iterator of chunks
批量调用graph.batch(inputs)Runnable/Pregel batch 逻辑list[Output]

8.2 运行时总链路#

CompiledStateGraph.invoke(initial_state)
Pregel.invoke
Pregel.stream 收集最终 chunk
初始化 input channel: START
START 边激活第一个业务节点
节点读取 state channels
节点函数返回 partial update
ChannelWrite 将 update 写入对应 state channels
reducer / LastValue 合并 state
边 channel 激活下一个节点
重复直到 END / halt
返回 output_channels 对应 state

8.3 CompiledStateGraph 作为运行时对象源码解剖#

职责与所处阶段#

CompiledStateGraphStateGraph.compile() 的输出对象。它负责承接 builder 编译结果,并通过 Pregel runtime 提供 invoke / stream / batch / ainvoke 能力。

真实源码签名#

CompiledStateGraph 类定义位于 state.py,可以概括为:

class CompiledStateGraph(Pregel[...]):
builder: StateGraph
...

调用方与被调用方#

StateGraph.compile
CompiledStateGraph(...)
CompiledStateGraph.attach_node / attach_edge
Pregel.invoke / stream

输入、输出与状态变化#

项目类型说明
输入compiled graph configchannels、nodes、input/output channels。
输出executable graph可调用对象。
状态变化self.nodesself.channels保存 Pregel 执行结构。
副作用无业务副作用只是构造可执行对象。

细粒度伪代码#

class CompiledStateGraph(Pregel):
def attach_node(self, key, node_spec):
# 1. START 是特殊入口节点
if key == START:
self.nodes[START] = PregelNode(
triggers=[START],
channels=[START],
writers=[...],
)
return
# 2. 普通节点读取 state channels
input_schema = node_spec.input_schema
read_channels = keys_from_schema(input_schema)
# 3. 包装用户 runnable
runnable = node_spec.runnable
# 4. 生成写入器:只允许写 state 中合法字段
writers = [
ChannelWrite(
entries=[ChannelWriteTupleEntry(...)]
)
]
# 5. 保存 PregelNode
self.nodes[key] = PregelNode(
triggers=[key],
channels=read_channels,
mapper=state_mapper,
writers=writers,
) | runnable | writers

逐段解释#

第 1 段说明 START 在 compiled graph 中也会变成一个特殊 Pregel 节点 / channel。用户输入并不是直接交给第一个业务节点,而是先写入入口 channel,再通过边触发业务节点。

第 2 段说明节点读取的是 schema 对应的 channels,而不是复制一份全局 dict。运行时会根据 channels 构造当前节点可见 state。

第 4 段是 partial update 合并的关键前置:节点输出会通过 writer 写入 state channels,而不是直接修改 Python dict。

正常路径#

CompiledStateGraph.attach_node("parse_request", node_spec)
生成 PregelNode
PregelNode 订阅 parse_request trigger
读取 TravelState channels
执行 parse_request runnable
将返回 dict 写入 destination / days channels

关键分支与异常路径#

条件行为结果
key == START创建入口节点处理输入 state。
普通节点包装 runnable + channel read/write可被 Pregel 调度。
节点输出非法字段写入阶段抛错防止污染 state。
节点没有返回 dict / Command 等合法对象InvalidUpdateError保持 state update 协议。

设计原因与工程影响#

Compiled graph 把节点执行拆成 read、run、write 三段,可以让 LangGraph 在同一套 runtime 中支持普通 Workflow、并行分支、条件路由、checkpoint、stream、retry、interrupt 等能力。

源码证据#

  • langgraph/graph/state.py::CompiledStateGraph
  • langgraph/pregel/_read.py::PregelNode
  • langgraph/pregel/_write.py::ChannelWrite
  • https://github.com/langchain-ai/langgraph/blob/1.2.7/libs/langgraph/langgraph/graph/state.py

8.4 invoke()stream() 的运行时入口源码解剖#

职责与所处阶段#

invoke() 是同步调用入口。对于 Pregel 风格 runtime,它通常会复用 stream() 主链,并收集最终输出。

真实源码签名#

实际签名在 Pregel 中较长,可概括为:

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(initial_state)
CompiledStateGraph.invoke
Pregel.invoke
Pregel.stream
运行时循环

输入、输出与状态变化#

项目类型说明
输入dict初始 graph input。
输出dictoutput channels 对应的最终 state。
状态变化channels 更新每个节点返回的 update 写入 channels。
副作用可选 checkpoint / stream events本篇不展开。

细粒度伪代码#

def invoke(input, config=None, *, context=None,
stream_mode="values", output_keys=None,
interrupt_before=None, interrupt_after=None, **kwargs):
# 1. 同步 invoke 复用 stream 主链
chunks = []
for chunk in self.stream(
input,
config=config,
context=context,
stream_mode=stream_mode,
output_keys=output_keys,
interrupt_before=interrupt_before,
interrupt_after=interrupt_after,
**kwargs,
):
# 2. 收集流式输出
chunks.append(chunk)
# 3. 根据 stream_mode 整理最终输出
if stream_mode == "values":
return last_value_chunk(chunks)
if stream_mode == "updates":
return merge_or_return_updates(chunks)
return normalize_stream_result(chunks)

逐段解释#

第 1 段说明 invoke() 不是一套完全独立的执行器,它通常复用 stream()。这也是为什么 LangGraph 能在同一个 runtime 中同时支持最终结果和流式观察。

第 2 段收集 stream chunks。对于 stream_mode="values",chunk 往往是当前完整 state 值;对于 updates,chunk 更偏向节点增量更新。

第 3 段将 stream 输出规整为 invoke() 的返回值。

正常路径#

graph.invoke(initial_state)
for chunk in graph.stream(...)
记录最后一个 values chunk
返回最终 state

关键分支与异常路径#

条件行为结果
stream_mode="values"返回最终 state 值最常见。
stream_mode="updates"返回节点更新流适合调试。
图执行中异常stream 抛出invoke 同步失败。
recursion limit 触发抛图递归错误防止无限循环。

设计原因与工程影响#

复用 stream() 可以让同步调用和流式调用共享同一套 runtime 语义,避免出现“invoke 正常但 stream 行为不同”的两套实现。

源码证据#

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

8.5 Pregel super-step 主链源码解剖#

职责与所处阶段#

Pregel runtime 负责按 message-passing 语义执行图:节点收到消息后变为 active,执行后写出更新,更新通过 channel 触发下一批节点。官方 Graph API 文档说明,LangGraph 底层受 Pregel 启发,按 super-step 执行;同一 super-step 中可并行运行节点,图在所有节点 inactive 且没有消息传递时终止。

真实源码签名#

Pregel 内部调度函数较多,本节使用 Pregel.stream() 作为运行主链入口。

调用方与被调用方#

Pregel.invoke
Pregel.stream
create loop / runtime context
while loop.tick(): execute active tasks
apply writes to channels
yield stream chunks

输入、输出与状态变化#

项目类型说明
输入initial input + config初始 state 和运行配置。
输出stream chunks / final outputvalues 或 updates。
状态变化channel values节点写入会更新 channel。
副作用callbacks、checkpoint、debug events运行时可观测与恢复能力。

细粒度伪代码#

def stream(input, config=None, *, stream_mode="values", output_keys=None,
interrupt_before=None, interrupt_after=None, context=None, **kwargs):
# 1. 归一化 config / context / stream 参数
config = ensure_config(config)
output_keys = resolve_output_keys(output_keys)
stream_modes = normalize_stream_mode(stream_mode)
# 2. 创建运行时 loop,包含 channels、checkpoint、tasks、step 计数
with create_pregel_loop(
graph=self,
input=input,
config=config,
context=context,
output_keys=output_keys,
interrupt_before=interrupt_before,
interrupt_after=interrupt_after,
) as loop:
# 3. 将初始 input 写入 input channel
loop.initialize_input(input)
# 4. Pregel 主循环:每次 tick 对应一个 super-step
while loop.tick():
# 4.1 找到当前 super-step 中 active tasks
tasks = loop.prepare_next_tasks()
# 4.2 执行 active nodes,可并行
results = runner.tick(tasks)
# 4.3 收集每个节点的 writes
writes = collect_writes(results)
# 4.4 将 writes 应用到 channels
loop.apply_writes(writes)
# 4.5 根据 stream_mode 产出 updates / values
for chunk in loop.emit_stream_chunks(stream_modes):
yield chunk
# 5. loop 结束后产出最终状态或清理
yield from loop.emit_final_if_needed()

逐段解释#

第 2 段创建 runtime loop。此处保存了执行状态:当前 step、channels、tasks、checkpoint、stream writer、interrupt 设置等。

第 3 段把用户输入写入 input channel。对 StateGraph 来说,input channel 通常是 START,然后 START 的边会激活入口业务节点。

第 4 段是主链。每个 super-step 中,runtime 找到当前 active 节点,执行节点函数,收集写入,统一应用到 channels,再决定下一步哪些节点 active。

第 4.4 段是 partial update 合并的关键。节点不是直接修改全局 dict,而是把输出转成 channel writes,由 channel 决定如何合并。

正常路径#

step 0:
START channel 收到 input
parse_request 被激活
step 1:
parse_request 执行
写入 destination / days
generate_plan 被激活
step 2:
generate_plan 执行
写入 plan
END 被触发
step 3:
无 active tasks
输出最终 state

关键分支与异常路径#

条件行为结果
同一 super-step 多个节点 active并行执行多个 writes 同时进入 channel。
多个节点写同一 LastValue channel可能冲突需要 reducer 或避免并行写同键。
节点抛异常进入 retry / error 机制或抛出本篇不展开 retry 细节。
达到 recursion limit抛错误防止循环不终止。
interrupt before/after 命中暂停图后续 HITL 篇展开。

设计原因与工程影响#

Pregel super-step 机制使 LangGraph 不只是线性 DAG 执行器,而能支持循环、并行、动态路由、持久化和中断。对于 Workflow Agent,普通固定边只是一种简化用法;底层 runtime 可以表达更复杂的 agent 控制流。

源码证据#

  • langgraph/pregel/main.py::Pregel.stream
  • LangGraph Graph API overview 关于 super-step、active/inactive、halt 的说明
  • https://github.com/langchain-ai/langgraph/blob/1.2.7/libs/langgraph/langgraph/pregel/main.py

8.6 节点执行与 partial update 源码解剖#

职责与所处阶段#

节点执行阶段负责把当前 state 交给用户函数,拿到返回值,并把返回值解释为 state partial update。

真实源码签名#

用户节点签名由用户定义,本篇最小示例为:

def parse_request(state: TravelState) -> dict[str, Any]:
...

Compiled graph 内部会围绕该函数构造 PregelNode / Runnable。

调用方与被调用方#

Pregel runner
PregelNode
ChannelRead 读取当前 state
用户节点 runnable.invoke(state)
ChannelWrite 写入 partial update

输入、输出与状态变化#

项目类型说明
输入TravelState当前 state 快照。
输出dictpartial update。
状态变化state channels只更新返回 dict 中包含的 key。
副作用取决于用户节点节点内部可能调用模型、工具、数据库。

细粒度伪代码#

def execute_node(node_spec, channel_values, config, runtime):
# 1. 根据节点 input_schema 从 channels 构造 state view
state = read_state_from_channels(
schema=node_spec.input_schema,
channels=channel_values,
)
# 2. 调用用户节点逻辑
raw_output = node_spec.runnable.invoke(
state,
config=config,
)
# 3. 解析节点返回值
if isinstance(raw_output, Command):
updates, goto = parse_command(raw_output)
elif isinstance(raw_output, dict):
updates = raw_output
goto = None
else:
raise InvalidUpdateError("Expected dict or Command")
# 4. 只允许更新 state schema 中存在的字段
valid_updates = {}
for key, value in updates.items():
if key not in known_state_channels:
raise InvalidUpdateError(f"Unknown state key: {key}")
valid_updates[key] = value
# 5. 将更新转成 channel writes
writes = []
for key, value in valid_updates.items():
writes.append(ChannelWriteEntry(channel=key, value=value))
# 6. 如果有 goto / edge trigger,也写入控制 channel
if goto is not None:
writes.append(make_branch_or_goto_write(goto))
return writes

逐段解释#

第 1 段说明节点看到的是 state view,而不是直接访问内部 channel 对象。对业务代码来说,它就是 TravelState dict。

第 2 段执行用户逻辑。节点可以是普通函数,也可以是 Runnable 链,这取决于 add_node() 时的归一化结果。

第 3 段解释了返回值协议。最简单路径是返回 dict;高级路径可以返回 Command 等对象,本篇只聚焦 dict partial update。

第 4 段保护 state schema。节点不能随意写入 schema 中不存在的字段,否则会破坏图状态契约。

第 5 段将 dict update 转成 channel writes。真正合并发生在 channel 层。

正常路径#

parse_request 读取:
{
"user_request": "我想去东京玩 5 天",
"destination": None,
"days": None,
"plan": None
}
parse_request 返回:
{
"destination": "东京",
"days": 5
}
运行时写入:
destination channel ← "东京"
days channel ← 5

关键分支与异常路径#

条件行为结果
返回 dict解释为 partial update正常更新 state。
返回空 dict不更新 state只推进控制流。
返回未知字段InvalidUpdateError防止 state 污染。
返回非 dict / 非 CommandInvalidUpdateError保持节点协议。
节点内部异常向 runtime 抛出可能进入 retry/error handler。

设计原因与工程影响#

partial update 让节点只关心自己负责的字段,降低节点之间的耦合。例如 parse_request 不需要返回完整 TravelState,只返回 destinationdays。这也是 Workflow Agent 中节点可以小而清晰的原因。

源码证据#

  • langgraph/graph/state.py::CompiledStateGraph.attach_node
  • langgraph/pregel/_read.py::PregelNode
  • langgraph/pregel/_write.py::ChannelWrite
  • langgraph/errors.py::InvalidUpdateError

8.7 state 合并与 channel reducer 源码解剖#

职责与所处阶段#

state 合并阶段负责把节点返回的 partial update 应用到 state channels。不同 channel 有不同合并语义。

真实源码签名#

核心 channel 类分布在:

langgraph/channels/last_value.py
langgraph/channels/binop.py
langgraph/channels/base.py

调用方与被调用方#

ChannelWriteEntry(key, value)
Pregel loop apply_writes
channel.update(values)
channel.value changed

输入、输出与状态变化#

项目类型说明
输入list[update_value]同一 super-step 中写入某 state key 的值。
输出channel 内部 value更新后的字段值。
状态变化channel value覆盖或聚合。
副作用update 冲突可能抛错取决于 channel 类型。

细粒度伪代码:默认 LastValue#

class LastValue(BaseChannel):
def update(self, values):
# 1. 没有更新,不改变当前值
if len(values) == 0:
return False
# 2. 默认 channel 每个 step 只接受一个值
if len(values) != 1:
raise InvalidUpdateError(
"Can receive only one value per step"
)
# 3. 覆盖为最新值
self.value = values[-1]
return True

逐段解释#

默认 state 字段通常使用 LastValue。这意味着同一字段在正常顺序执行中会被最新 partial update 覆盖:

destination: None
↓ parse_request returns {"destination": "东京"}
destination: "东京"

但如果同一 super-step 中多个并行节点同时写入同一个 LastValue 字段,可能触发冲突。此时应显式定义 reducer。

细粒度伪代码:BinaryOperatorAggregate#

class BinaryOperatorAggregate(BaseChannel):
def __init__(self, typ, operator):
self.typ = typ
self.operator = operator
self.value = empty_initial_value()
def update(self, values):
if len(values) == 0:
return False
for value in values:
if self.value is empty:
self.value = value
else:
self.value = self.operator(self.value, value)
return True

逐段解释#

如果 state 字段写成:

from typing import Annotated
from operator import add
class State(TypedDict):
notes: Annotated[list[str], add]

那么多个节点写入:

{"notes": ["美食"]}
{"notes": ["景点"]}

会合并为:

{"notes": ["美食", "景点"]}

而不是互相覆盖或冲突。

正常路径#

节点返回 partial update
按 key 分组 writes
找到 key 对应 channel
channel.update(values)
更新 channel 内部值
后续节点读取新 state

关键分支与异常路径#

条件行为结果
单节点写 LastValue覆盖旧值默认最常见。
多节点同 step 写 LastValue抛错或冲突需要 reducer。
reducer 字段多值写入聚合支持并行分支汇聚。
update 字段不存在InvalidUpdateError保持 state schema 稳定。

设计原因与工程影响#

LangGraph 用 channel 抽象取代“直接修改 dict”,是为了让状态更新在并行、循环和恢复场景下有明确语义。你写的节点只是返回 partial update;真正的合并规则由 schema 和 channel 决定。

源码证据#

  • langgraph/channels/last_value.py::LastValue
  • langgraph/channels/binop.py::BinaryOperatorAggregate
  • langgraph/graph/state.py::_get_channels

8.8 Callback、事件与可观测性#

时机事件携带数据失败行为
图开始graph run startinput、config、run id初始化失败则抛出。
节点开始task/node startnode name、state slice、config节点失败进入异常路径。
节点结束task/node endpartial update、耗时写入 channel。
写入完成channel updatekey、value、step冲突可能抛错。
图结束graph run endoutput state返回给 caller。
图失败graph run errorexception、step、node抛给 caller 或 checkpoint。

8.9 运行时主链总结#

initial state dict
START input channel
entry node activated
node reads state channels
node returns partial update
ChannelWrite writes updates
channels merge by reducer
edge activates next node
END / no active tasks
output state

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

本章回答:

当输入、图结构、节点返回值、状态更新或运行循环偏离正常主链时,LangGraph 如何分流、恢复、终止或失败?

9.1 分支矩阵#

分支类型触发条件核心函数结果
构建后继续修改self.compiled == True 后 add_node/add_edgeadd_node() / add_edge()warning,旧 compiled graph 不变。
schema 冲突同名 state key channel 类型不兼容_add_schema()抛错。
非法节点名节点名是 START / END 或重复add_node()抛错。
非法边方向END → nodenode → STARTadd_edge()抛错。
节点返回非法对象非 dict / Command 等runtime write pathInvalidUpdateError
partial update 字段未知节点写 schema 外字段write pathInvalidUpdateError
LastValue 并发冲突同 step 多节点写同 keychannel update抛错,建议 reducer。
循环不终止图持续产生新任务Pregel runtimerecursion limit 触发。

9.2 同步与异步分支#

维度同步路径异步路径
入口invoke() / stream()ainvoke() / astream()
调度方式同步 runner / executorasync runner / awaitable tasks
资源语义同步节点可能阻塞线程async 节点可协作式调度
超时语义同步节点不易安全取消async timeout 更可控
工程建议适合简单 CPU / 快速函数适合 LLM、HTTP、DB 等 I/O

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

StateGraph builder 本身不直接执行 batch、stream、parallel。

这些能力由编译后的 CompiledStateGraph / Pregel runtime 提供:

StateGraph
负责定义图结构
CompiledStateGraph / Pregel
负责 invoke / stream / batch / async / super-step 调度

并行不是由 add_edge() 自动表达的所有场景都支持。通常需要多个节点在同一 super-step 被激活,或者使用条件边 / Send / waiting edges 等机制。对于普通线性 Workflow:

START → A → B → END

它就是顺序执行。

9.4 异常分类#

异常类别抛出位置是否可恢复处理策略是否反馈上层
构建期结构错误add_node / add_edge / compile修改图定义
节点返回格式错误runtime write path修正节点返回协议
state channel 冲突channel update视情况增加 reducer 或避免并行写同 key
节点业务异常用户节点函数视 retry/error handlerretry、fallback、抛出
递归上限Pregel runtime视业务增加停止条件或调高 limit

9.5 异常路径源码解剖#

细粒度伪代码#

def apply_node_output(raw_output, known_channels):
try:
# 1. 解析节点返回值
if isinstance(raw_output, dict):
updates = raw_output
elif isinstance(raw_output, Command):
updates = raw_output.update
else:
raise InvalidUpdateError("Expected dict or Command")
# 2. 校验字段是否合法
for key in updates:
if key not in known_channels:
raise InvalidUpdateError(f"Unknown state key: {key}")
# 3. 写入 channels
grouped = group_updates_by_channel(updates)
for channel_name, values in grouped.items():
channel = known_channels[channel_name]
channel.update(values)
except InvalidUpdateError:
# 4. 协议错误,继续抛出
raise
except Exception:
# 5. 用户节点或 channel 内部错误,交给 runtime 处理
raise

逐段解释#

第 1 段区分合法节点返回值。最小 Workflow 中,节点应该返回 dict。如果返回字符串、list 或完整 state 之外的对象,运行时无法知道如何合并。

第 2 段保护 state schema。所有 update key 都必须来自 schema,否则 state 会变成不可预测的动态字典。

第 3 段应用 channel update。默认 LastValue 只接受单个值,自定义 reducer 可以接受多个值。

第 4 段不吞掉协议错误。节点返回协议错误通常是代码 bug,不应该让 Agent 静默继续。

9.6 Retry、Fallback 与恢复边界#

机制适用条件不适用条件幂等要求
RetryLLM 调用超时、临时网络错误、可重试工具节点非幂等写操作、参数协议错误节点必须可重复执行。
Fallback主模型或服务失败业务规则不确定fallback 输出需符合同一 state schema。
Repair结构化输出解析失败节点函数代码错误repair 后仍必须返回合法 partial update。
Error handler节点业务异常可转成状态安全风险或数据损坏handler 不能捕获自身异常形成循环。

9.7 停止条件与保护上限#

正常结束:执行到 END,或所有节点 inactive 且无消息在传递。
提前结束:节点返回 Command / 分支路由直接到 END。
人工中断:interrupt_before / interrupt_after 命中。
框架保护:recursion_limit 防止无限循环。
异常失败:构建错误、节点错误、channel 冲突、非法 update。

9.8 能力边界#

容易误判的能力实际提供者本篇对象的真实职责
durable executioncheckpointer / Pregel runtimeStateGraph.compile() 只接收 checkpointer 并交给 compiled graph。
human-in-the-loopinterrupt / checkpointStateGraph 只定义图结构。
条件路由BranchSpec / add_conditional_edges本篇只讲固定边。
工具调用ToolNode / BaseToolStateGraph 可承载工具节点,但不定义工具协议。
LLM prompt用户节点内部StateGraph 不抽象 prompt。

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

本章回答:

StateGraph 允许在哪里插入自定义行为,它如何与 Pregel、Runnable、Channel、LangChain 组件协作?

10.1 扩展点总览#

扩展点扩展方式执行时机可修改内容约束
State schemaTypedDict / dataclass / Pydantic构建期state 字段、输入输出 schema字段必须可解析为 channels。
ReducerAnnotated[T, reducer]state 合并期同 key 多更新合并方式reducer 需确定、可重复。
Node函数 / Runnable / subgraph运行时业务逻辑、LLM、工具调用返回合法 partial update。
Edgeadd_edge构建期固定控制流节点必须存在,方向合法。
Node policiesretry/cache/timeout/error_handler节点执行期失败恢复、缓存、超时注意幂等性。
Runtime contextcontext_schema每次运行只读上下文,如 user_id / db不应当用作可变 state。
Checkpointercompile(checkpointer=...)super-step 边界持久化状态后续专题展开。

10.2 Reducer 扩展源码解剖#

职责与所处阶段#

Reducer 让同一个 state key 能够接收多个节点的更新,并按用户定义的函数聚合。

真实源码签名#

用户侧 reducer 函数形态:

def reducer(old_value: T, new_value: T) -> T:
...

schema 使用:

class State(TypedDict):
notes: Annotated[list[str], operator.add]

调用方与被调用方#

_get_channels(schema)
识别 Annotated[field_type, reducer]
BinaryOperatorAggregate(field_type, reducer)
运行时 channel.update(values)

输入、输出与状态变化#

项目类型说明
输入old value + new valuechannel 当前值和节点更新值。
输出merged value合并后的 state 字段。
状态变化channel value被 reducer 更新。
副作用不应有reducer 应保持纯函数。

细粒度伪代码#

def parse_field_annotation(field_type):
# 1. 普通字段:默认 LastValue
if not is_annotated(field_type):
return LastValue(field_type)
# 2. Annotated 字段:提取 metadata
base_type, *metadata = get_annotated_args(field_type)
# 3. 如果 metadata 中有 reducer 函数
reducer = find_binary_operator(metadata)
if reducer is not None:
return BinaryOperatorAggregate(base_type, reducer)
# 4. 如果 metadata 中显式给了 channel
channel = find_channel(metadata)
if channel is not None:
return channel
# 5. fallback
return LastValue(base_type)

逐段解释#

第 1 段解释默认行为:普通字段使用覆盖语义。

第 3 段解释 reducer 行为:Annotated[list[str], add] 会把 state key 变成可聚合 channel,而不是默认 LastValue。

第 4 段说明 LangGraph 也允许更底层 channel 自定义,但业务代码通常只需要 reducer annotation。

正常路径#

class State(TypedDict):
research_notes: Annotated[list[str], add]
_get_channels
research_notes → BinaryOperatorAggregate(list[str], add)

关键分支与异常路径#

条件行为结果
无 AnnotatedLastValue覆盖更新。
Annotated + reducerBinaryOperatorAggregate聚合更新。
reducer 输入类型不匹配运行时异常修正 schema 或节点输出。
reducer 有副作用行为不可预测不建议。

设计原因与工程影响#

Reducer 是 Workflow Agent 支持并行信息收集的关键。例如旅行助手中可以并行搜索美食和景点,两个节点都写 research_notes,再用 reducer 聚合结果。

源码证据#

  • langgraph/graph/state.py::_get_channels
  • langgraph/channels/binop.py::BinaryOperatorAggregate

10.3 扩展调用链#

StateGraph.add_node(Runnable chain)
coerce_to_runnable
CompiledStateGraph.attach_node
PregelNode executes Runnable
Runnable internally calls Prompt / Model / Parser / Tool
returns partial state update
ChannelWrite updates state

这说明 LangGraph 不替代 LangChain Core,而是作为编排层承载任意 Runnable 或普通函数。

10.4 与相邻框架组件的协作#

相邻组件输入协议输出协议协作边界
LangChain Runnableinvoke(input, config)任意对象节点可以是 Runnable,但必须最终返回合法 partial update。
ChatModelmessages / PromptValueAIMessage通常封装在节点内部。
ToolNodemessages with tool_callsToolMessage update后续工具节点专题展开。
Checkpointergraph state snapshotscheckpointcompile 时注入,runtime 使用。
LangSmith tracingrun eventstracesruntime 产生可观测事件。

10.5 公共扩展接口与内部实现#

业务代码可以依赖:
StateGraph
add_node
add_edge
compile
Annotated reducer
context_schema
input_schema / output_schema
retry_policy / cache_policy / timeout
业务代码避免依赖:
_get_channels
StateNodeSpec 内部字段
CompiledStateGraph.attach_node
Pregel internal loop
ChannelWriteTupleEntry

10.6 自定义扩展示例#

示例只展示扩展契约,不重新实现框架:

from typing import Annotated
from operator import add
from typing_extensions import TypedDict
class TravelState(TypedDict):
user_request: str
research_notes: Annotated[list[str], add]
plan: str | None
def search_food(state: TravelState):
return {"research_notes": ["大阪美食:章鱼烧、串炸"]}
def search_attractions(state: TravelState):
return {"research_notes": ["大阪景点:大阪城、道顿堀"]}

说明:

  1. 扩展点是 state schema 的 reducer annotation。
  2. 节点仍然返回 partial update。
  3. reducer 决定并行写入如何合并。
  4. 节点不需要知道其他节点是否也写同一字段。
  5. 如果 reducer 不满足幂等/确定性,checkpoint replay 时可能难以推理。

10.7 选择扩展还是重写流程#

条件选择扩展点选择更底层框架
只需要改变状态合并方式
需要固定业务流程
需要复杂动态路由是,使用条件边
需要自定义 Pregel 调度算法
需要跨节点人工审批视情况通常继续用 LangGraph + interrupt
需要完全不同状态存储模型可能需要更底层实现

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

11.1 适用场景#

场景是否推荐原因
旅行规划 Workflow需求解析、规划、查询、生成可拆成显式节点。
订单履约流程状态机、规则判断、工具执行、校验都适合显式图。
多步骤审批可以把人工确认、风险检查、执行动作拆成节点。
简单单轮问答普通 chain 更轻量。
完全开放式探索 Agent视情况create_agent 或 ReAct loop 可能更快。
高并发低延迟纯 API 转发图 runtime 可能增加不必要开销。

11.2 工程决策表#

决策点推荐选择前提风险
State schemaTypedDict 起步追求性能和简洁缺少运行时深度校验。
复杂校验Pydantic state 或节点内部 Pydantic 输出需要字段校验性能略低。
节点返回partial update节点只负责局部字段返回未知字段会失败。
并行写同字段显式 reducer有多个节点写同 keyreducer 设计不当会污染状态。
固定流程add_edge顺序明确灵活性低。
动态流程add_conditional_edges需要 Router后续专题学习。
长任务compile 时加 checkpointer需要恢复需要 thread/config 管理。

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

性能:StateGraph 本身很轻,主要开销来自节点内部 LLM / 工具 / IO;但 graph runtime 会增加调度、状态合并、tracing 成本。
可靠性:显式 State + Node + Edge 比开放式 ReAct 更可控;但节点必须遵守 partial update 协议。
安全:StateGraph 不自动做权限校验;高风险节点仍需要 guardrail、policy、interrupt 或人工确认。
可观测性:必须记录每个节点输入 state、partial update、异常、耗时、下一节点。

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

12.1 误区:StateGraph 本身可以执行#

错误原因:名字里有 Graph,容易以为 builder 就是运行图。

源码事实:

StateGraph 是 builder class。
必须 compile 生成 CompiledStateGraph 后,才支持 invoke / stream / batch / ainvoke。

工程影响:

如果在 builder 上寻找 invoke(),会误解 LangGraph 的构建期和运行期边界。

12.2 误区:节点必须返回完整 State#

错误原因:很多状态机框架要求每一步返回完整状态。

源码事实:

LangGraph 节点签名是 State -> Partial[State]。
节点返回 dict partial update,由 channel 合并到 state。

工程影响:

要求每个节点返回完整 state 会增加耦合,容易覆盖其他节点更新。

12.3 误区:state 就是普通 dict#

错误原因:业务代码里看到的是 dict-like state。

源码事实:

schema 会被解析成 channels。
每个 state key 背后有 channel 合并语义。

工程影响:

并行写同一字段时,如果没有 reducer,可能发生冲突。

12.4 误区:add_edge() 传递的是业务数据#

错误原因:边看起来像数据流。

源码事实:

edge 表示控制流触发关系。
业务数据通过 state channels 共享。

工程影响:

不要试图在 edge 上携带业务 payload;应写入 state。

12.5 误区:compile() 只是形式化步骤#

错误原因:简单例子里 compile 后立即 invoke,看起来只是固定写法。

源码事实:

compile 会校验图结构,构造 CompiledStateGraph,attach nodes / edges / branches,并创建 Pregel 应用。

工程影响:

compile 后再修改 builder,不会影响已经拿到的 compiled graph。

12.6 误区:STARTEND 是普通节点#

错误原因:它们像节点名一样传给 add_edge()

源码事实:

START / END 是特殊常量,用于入口和终止语义。
不能作为普通业务节点名。

工程影响:

不要定义名为 __start____end__ 的业务节点。

12.7 误区:Workflow Agent 不需要 recursion limit#

错误原因:线性 workflow 通常不会循环。

源码事实:

LangGraph 支持循环;一旦图有回边或动态路由,就可能不终止。
Pregel runtime 使用 recursion limit 作为保护。

工程影响:

复杂 Agent 图应显式设计停止条件,并设置合理 recursion limit。


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

13.1 构建期心智模型#

State schema
_get_channels 解析 state key 与 reducer
add_node 保存 StateNodeSpec
add_edge 保存固定控制流
compile 校验并构造 CompiledStateGraph

13.2 运行时心智模型#

initial state
START channel
entry node reads state
node returns partial update
ChannelWrite writes updates
channels merge by LastValue / reducer
edge activates next node
END outputs final state

13.3 分支与异常心智模型#

正常路径 → 节点返回 dict,channel 合并,边触发下一节点。
分支路径 → waiting_edges / conditional_edges / Command 等高级机制改变控制流。
可恢复异常 → retry_policy / error_handler / fallback 处理。
不可恢复异常 → 构建错误、非法 update、未知字段继续抛出。
保护上限 → recursion_limit 防止循环不终止。

13.4 一句话总结#

StateGraph 通过 state schema 生成 channels,通过 add_node() 收集节点规格,通过 add_edge() 收集控制流,通过 compile() 生成基于 Pregel 的 CompiledStateGraph;运行时节点只返回 partial update,真正的状态合并由 channel / reducer 完成,最终由图执行到 END 或 halt 后输出 state。

13.5 掌握检查#

  • 能说清 StateGraphCompiledStateGraph 的区别。
  • 能解释为什么 StateGraph 是 builder,而不是执行器。
  • 能画出 StateGraph.__init__ → _add_schema → add_node → add_edge → compile 构建链。
  • 能解释 State schema 如何映射为 channels。
  • 能说明默认 LastValue 合并语义。
  • 能说明什么时候需要 Annotated[..., reducer]
  • 能解释 add_node() 如何保存节点函数。
  • 能解释 add_edge() 如何保存固定控制流。
  • 能解释 compile() 做了哪些事情。
  • 能画出一次 invoke() 的节点执行和 partial update 合并过程。
  • 能说明节点返回未知字段为什么会失败。
  • 能说明 START / END 为什么不是普通节点。
  • 能区分固定边、等待边、条件边的学习边界。
  • 能根据业务场景判断是否适合用 StateGraph

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

14.1 官方概念文档#

  1. LangGraph Overview
    https://docs.langchain.com/oss/python/langgraph/overview

  2. LangGraph Graph API Overview
    https://docs.langchain.com/oss/python/langgraph/graph-api

  3. LangGraph Runtime / Pregel
    https://docs.langchain.com/oss/python/langgraph/pregel

  4. LangGraph PyPI Project
    https://pypi.org/project/langgraph/

14.2 官方 API Reference#

  1. StateGraph
    https://reference.langchain.com/python/langgraph/graph/state/StateGraph

  2. StateGraph.compile()
    https://reference.langchain.com/python/langgraph/graph/state/StateGraph/compile

  3. Graphs Reference
    https://reference.langchain.com/python/langgraph/graphs

14.3 官方源码#

  1. langgraph/graph/state.py::StateGraph
    https://github.com/langchain-ai/langgraph/blob/1.2.7/libs/langgraph/langgraph/graph/state.py

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

  3. langgraph/channels/last_value.py::LastValue
    https://github.com/langchain-ai/langgraph/blob/1.2.7/libs/langgraph/langgraph/channels/last_value.py

  4. langgraph/channels/binop.py::BinaryOperatorAggregate
    https://github.com/langchain-ai/langgraph/blob/1.2.7/libs/langgraph/langgraph/channels/binop.py

14.4 下一篇衔接#

下一篇进入:

第 7 篇:Conditional Edge 与 Router 实现机制源码解剖

需要继续回答:

add_conditional_edges 如何保存路由函数?
BranchSpec 如何表示条件边?
route function 返回 node name、END、list 或 Send 时如何处理?
条件边如何在 Pregel runtime 中转成 branch channel?
Router 范式如何落到 LangGraph 图执行机制?
LangGraph 源码深潜:StateGraph 与 Workflow Agent 运行机制解剖
https://jupiter-ws.cn/posts/agent-frameworks/langgraph-stategraph-source-deep-dive/
作者
Jupiter
发布于
2026-03-10
许可协议
CC BY-NC-SA 4.0