- 发布日期
第 08 讲 · 编译机制:StateGraph → Pregel
StateGraph 编译生成 PregelNode / ChannelWrite / Branch 的完整流程
学习目标
- 看懂
compile()把蓝图变成CompiledStateGraph(即Pregel)的全过程。 - 理解节点如何变成
PregelNode,边/分支如何变成ChannelWrite+ 触发通道。 - 掌握"分支触发通道(branch:to:X)"这一关键内部机制。
compile 的产物:CompiledStateGraph 就是 Pregel
graph/state.py,compile()(1164 行)。它返回 CompiledStateGraph,而后者直接继承 Pregel:
class CompiledStateGraph(
Pregel[StateT, ContextT, InputT, OutputT],
Generic[StateT, ContextT, InputT, OutputT],
):
也就是说:StateGraph 是面向用户的"友好建图 API",Pregel 是底层执行引擎;compile 就是 把前者翻译成后者。 模块四会专讲 Pregel 怎么跑,本讲只看"翻译"。
compile 干了哪几件事
compile() 主体(1218–1388 行)按顺序:
- 校验 checkpointer、(可选)serde 白名单、图结构
validate()。 - 确定输出/流式通道:
output_channels/stream_channels(1256–1273 行)。 - 应用 per-node 默认值:retry / cache / timeout / 错误处理节点(1277–1331 行)。
- 构造
CompiledStateGraph,把 channels 组装好:
channels={
**self.channels,
**self.managed,
START: EphemeralValue(self.input_schema),
},
input_channels=START,
stream_mode="updates",
output_channels=output_channels,
stream_channels=stream_channels,
注意:START 是一个 EphemeralValue 通道(回顾第 4 讲:只活一步),用户输入就写进它。
- attach 节点 / 边 / 分支(1360–1386 行)——这是翻译的核心:
compiled.attach_node(START, None)
for key, node in self.nodes.items():
compiled.attach_node(key, node)
...
for start, end in self.edges:
compiled.attach_edge(start, end)
for starts, end in self.waiting_edges:
compiled.attach_edge(starts, end)
for start, branches in self.branches.items():
for name, branch in branches.items():
compiled.attach_branch(start, name, branch)
attach_node:节点 → PregelNode
attach_node(1431 行)把每个用户节点翻译成一个 PregelNode(Actor)。关键看普通节点分支:
branch_channel = _CHANNEL_BRANCH_TO.format(key)
self.channels[branch_channel] = (
LastValueAfterFinish(Any)
if node.defer
else EphemeralValue(Any, guard=False)
)
self.nodes[key] = PregelNode(
triggers=[branch_channel],
# read state keys and managed values
channels=("__root__" if is_single_input else input_channels),
# coerce state dict to schema class (eg. pydantic model)
mapper=mapper,
# publish to state keys
writers=[ChannelWrite(write_entries)],
...
bound=node.runnable,
)
划重点:
- 每个节点 X 都有一个专属触发通道
branch:to:X(_CHANNEL_BRANCH_TO)。 默认是EphemeralValue(只活一步:被触发后就消失,避免重复触发); 若节点defer=True则用LastValueAfterFinish(延迟到收尾)。 triggers=[branch_channel]:节点订阅自己的触发通道——谁往branch:to:X写东西,X 下一步就跑。channels=input_channels:节点从这些状态通道读输入。writers=[ChannelWrite(write_entries)]:节点的输出经ChannelWrite写回状态通道。bound=node.runnable:你写的那个函数本体。
write_entries:节点输出怎么落通道
write_entries: tuple[ChannelWriteEntry | ChannelWriteTupleEntry, ...] = (
ChannelWriteTupleEntry(
mapper=_get_root if output_keys == ["__root__"] else _get_updates
),
ChannelWriteTupleEntry(
mapper=_control_branch,
static=_control_static(node.ends) ...
),
)
两条写规则:
_get_updates:把节点返回的 dict 拆成(channel, value)写入对应状态通道(走 reducer)。_control_branch:处理节点返回的Command(goto=...)——把控制流跳转翻译成对目标节点 触发通道的写(第 18 讲 Command 细讲)。
attach_edge:边 → 往触发通道写
attach_edge(1537 行)。普通单边:让 start 节点多挂一个 writer,往 end 的触发通道写:
def attach_edge(self, starts: str | Sequence[str], end: str) -> None:
if isinstance(starts, str):
# subscribe to start channel
if end != END:
self.nodes[starts].writers.append(
ChannelWrite(
(ChannelWriteEntry(_CHANNEL_BRANCH_TO.format(end), None),)
)
)
多起点边(join):生成 NamedBarrierValue 屏障通道,end 订阅它,每个 start 往它写自己的名字:
elif end != END:
channel_name = f"join:{'+'.join(starts)}:{end}"
if self.builder.nodes[end].defer:
self.channels[channel_name] = NamedBarrierValueAfterFinish(str, set(starts))
else:
self.channels[channel_name] = NamedBarrierValue(str, set(starts))
self.nodes[end].triggers.append(channel_name)
for start in starts:
self.nodes[start].writers.append(
ChannelWrite((ChannelWriteEntry(channel_name, start),))
)
这下第 5 讲的屏障与第 7 讲的"多对一边"完全闭环了。
attach_branch:条件边 → 动态写触发通道
attach_branch(1563 行)把 add_conditional_edges 的路由函数包装成一个 writer: 运行时执行路由函数 → 得到目标名字/Send 列表 → get_writes 把它们翻译成对目标触发通道的写 (或写入 TASKS 通道以触发 Send,第 19 讲)。
def get_writes(
packets: Sequence[str | Send], static: bool = False
) -> Sequence[ChannelWriteEntry | Send]:
writes = [
(
ChannelWriteEntry(
p if p == END else _CHANNEL_BRANCH_TO.format(p), None
)
if not isinstance(p, Send)
else p
)
for p in packets
if (True if static else p != END)
]
整体翻译对照表
| 用户 API | 编译产物 | 机制 |
|---|---|---|
add_node("X", fn) | PregelNode(triggers=[branch:to:X], bound=fn, writers=...) | Actor,订阅自己的触发通道 |
| 状态字段 | channel(按 reducer 选型) | 第 7 讲映射 |
add_edge(A, B) | A 多一个 writer → 写 branch:to:B | 单边触发 |
add_edge([A,B], C) | NamedBarrierValue join 通道 | 屏障同步 |
add_conditional_edges(A, route) | A 多一个 writer,运行时执行 route → 写目标触发通道 / Send | 动态路由 |
| 节点返回 dict | _get_updates → 写状态通道 | reducer 合并 |
节点返回 Command(goto=) | _control_branch → 写目标触发通道 | 控制流 |
全景图:一切都是"往触发通道写"
flowchart LR
START[(branch:to:agent)] --> AG[PregelNode agent]
AG -->|写状态| MSG[(messages channel)]
AG -->|条件边 route| W{route 结果}
W -->|tools| T1[(branch:to:tools)]
W -->|END| OUT[(输出)]
T1 --> TN[PregelNode tools]
TN -->|写状态| MSG
TN -->|静态边| BA[(branch:to:agent)]
BA --> AG
核心顿悟:LangGraph 里"边"不是数据通道,而是"谁该被触发"的信号—— 全部通过往目标节点的 branch:to:X 触发通道写来实现。节点是否执行,取决于其触发通道版本是否变化(第 10 讲)。
使用场景
- 排查"节点没执行":用
get_graph().draw_mermaid()看编译后的边;或在attach_edge/attach_branch打断点,确认 writer 是否挂上了目标触发通道。 - 理解循环为什么成立:tools→agent 的回边,本质是 tools 节点 writer 写
branch:to:agent,DAG 框架做不到。 - 自定义高级路由:返回
Send列表实现运行时 fan-out。
动手实验
- 编译第 7 讲的图,
app.get_graph().draw_mermaid()打印,对照源码理解节点与边。 - 在
attach_node打断点,观察每个节点的triggers与writers列表。 - 给一个节点设
defer=True,对比它的触发通道从EphemeralValue变成LastValueAfterFinish, 理解"延迟执行"的实现。
阅读作业
- 精读
attach_node/attach_edge/attach_branch(1431–1640 行)。 - 浏览
pregel/_read.py的PregelNode,为模块四预热。
小结
compile()把 StateGraph 翻译成 Pregel:节点→PregelNode、边/分支→ChannelWrite + 触发通道。- 每个节点有专属触发通道
branch:to:X;"边"= 往目标触发通道写信号。 - 节点输出经
ChannelWrite落状态通道(reducer 合并)或落控制流(Command)。
模块三完结。下一模块是全课最硬核的部分——Pregel 执行内核。