发布日期

第 08 讲 · 编译机制:StateGraph → Pregel

StateGraph 编译生成 PregelNode / ChannelWrite / Branch 的完整流程

学习目标

  • 看懂 compile() 把蓝图变成 CompiledStateGraph(即 Pregel)的全过程。
  • 理解节点如何变成 PregelNode,边/分支如何变成 ChannelWrite + 触发通道。
  • 掌握"分支触发通道(branch:to:X)"这一关键内部机制。

compile 的产物:CompiledStateGraph 就是 Pregel

graph/state.pycompile()(1164 行)。它返回 CompiledStateGraph,而后者直接继承 Pregel

1394:libs/langgraph/langgraph/graph/state.py
class CompiledStateGraph(
    Pregel[StateT, ContextT, InputT, OutputT],
    Generic[StateT, ContextT, InputT, OutputT],
):

也就是说:StateGraph 是面向用户的"友好建图 API",Pregel 是底层执行引擎;compile 就是 把前者翻译成后者。 模块四会专讲 Pregel 怎么跑,本讲只看"翻译"。

compile 干了哪几件事

compile() 主体(1218–1388 行)按顺序:

  1. 校验 checkpointer、(可选)serde 白名单、图结构 validate()
  2. 确定输出/流式通道output_channels / stream_channels(1256–1273 行)。
  3. 应用 per-node 默认值:retry / cache / timeout / 错误处理节点(1277–1331 行)。
  4. 构造 CompiledStateGraph,把 channels 组装好:
1346:libs/langgraph/langgraph/graph/state.py
            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 讲:只活一步),用户输入就写进它。

  1. attach 节点 / 边 / 分支(1360–1386 行)——这是翻译的核心:
1386:libs/langgraph/langgraph/graph/state.py
        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)。关键看普通节点分支:

1533:libs/langgraph/langgraph/graph/state.py
            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:节点输出怎么落通道

1492:libs/langgraph/langgraph/graph/state.py
        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 的触发通道写:

1545:libs/langgraph/langgraph/graph/state.py
    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 往它写自己的名字:

1561:libs/langgraph/langgraph/graph/state.py
        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 讲)。

1582:libs/langgraph/langgraph/graph/state.py
        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。

动手实验

  1. 编译第 7 讲的图,app.get_graph().draw_mermaid() 打印,对照源码理解节点与边。
  2. attach_node 打断点,观察每个节点的 triggerswriters 列表。
  3. 给一个节点设 defer=True,对比它的触发通道从 EphemeralValue 变成 LastValueAfterFinish, 理解"延迟执行"的实现。

阅读作业

  • 精读 attach_node / attach_edge / attach_branch(1431–1640 行)。
  • 浏览 pregel/_read.pyPregelNode,为模块四预热。

小结

  • compile() 把 StateGraph 翻译成 Pregel:节点→PregelNode、边/分支→ChannelWrite + 触发通道。
  • 每个节点有专属触发通道 branch:to:X;"边"= 往目标触发通道写信号。
  • 节点输出经 ChannelWrite 落状态通道(reducer 合并)或落控制流(Command)。

模块三完结。下一模块是全课最硬核的部分——Pregel 执行内核。