LangGraph 节点触发机制通俗解读
基于 LangGraph 1.2.9 源码,用白话解释“节点为什么会被触发”。
节点触发机制
LangGraph 的调度模型并非“通道数据变化时即时触发回调/事件通知节点(Push-based)”,而是基于 BSP(块同步并行)的拉取式调度(Pull-based):
- 超步同步与版本递增:每个超步(Superstep)结束时,调度引擎统一应用本轮活跃节点的写入操作,并为所有被更新的通道分配全局单调递增的新版本号。
- 索引定位与触发检测:下一个超步开始时,Pregel 调度引擎根据更新通道的反向索引(
trigger_to_nodes)定位到相关的候选节点,并代替节点检查其订阅的通道——只要订阅的通道数据可用且当前版本号大于该节点上次已见的版本水位,该节点即被激活。 - 任务打包派发:所有被激活的节点会被统一封装为
PULL Task派发执行。
本质上,这是一种基于“通道版本号比对与水位线控制”的拉取式批同步状态机,而非传统的事件监听与回调驱动。
【编译期 compile()】
• 创建控制通道: “branch:to:B” (EphemeralValue 类型)
• 节点 A.writers 登记: 写入 “branch:to:B”
• 节点 B.triggers 登记: 监听 “branch:to:B”【超步 N:节点 A 运行完毕】
• 执行节点 A 的 writers,将信号投递给 “branch:to:B”
• 超步同步栅栏处:更新通道,“branch:to:B” 版本号自增 (version++)【超步 N+1:任务分发】
• 调度引擎发现 “branch:to:B” 版本更新,通过 trigger_to_nodes 索引命中节点 B
• 检查版本号: version > seen_version,判定节点 B 激活
• 封装为 Task,调用节点 B 对应的函数
0. 前置:Pregel / BSP 模型
LangGraph 借鉴了 Google Pregel 的**超步(superstep)**思想:
- 每一轮(超步),所有被触发的节点并行执行;
- 节点执行中产生的写入(包括 state 更新和边信号)不立即生效,而是收集起来;
- 等到本轮所有节点执行完毕,统一应用所有写入,然后进入下一轮;
- 下一轮开始前,根据写入情况决定哪些节点在下一轮被触发。
这种模型天然适合图计算、分布式和 checkpoint 恢复。
1. 编译期:边 → 订阅(triggers)
在StateGraph.compile()时,你写的add_edge("A", "B")并不会生成A → B的直接函数调用。它被翻译成:
- 为每个节点B创建一个专属的虚拟 channel,名字叫
branch:to:B; - 节点B的
triggers列表里放入这个 channel —— 表示“我订阅了branch:to:B,只要它被写入,我就可能被触发”; - 节点A的
writers列表里被追加一条写入指令:当 A 执行完毕后,向branch:to:B写入一个空信号(None)。
扇入(多个前驱 → 一个后继)稍有不同:add_edge([A, B], C)会创建一个NamedBarrierValue类型的 join channel,A 和 B 各自往里面写入自己的名字,只有当所有前驱都写过了,该 channel 才变为“可用”(is_available() == True)。这样 C 只有在 A 和 B都执行完毕后才会被触发。
条件边类似,只是路由函数决定往哪个branch:to:目标写信号。
编译期最后会建立一张倒排索引trigger_to_nodes:每个 channel → 订阅了它的节点列表。
例如branch:to:B→[B],join_channel→[C]。这张表在运行期快速查找“哪些节点可能被触发”。
2. 节点执行时:写入只是“暂存”
当一个节点(比如 A)执行完毕,它的返回值(新的 state)以及编译期挂载的ChannelWrite指令会生成一批写入条目。
但这些写入不会立即修改 channel 内容,而是被追加到当前任务的task.writes队列里(一个deque)。
为什么这样?因为同一超步内可能有多个节点并行执行,如果直接修改共享 channel,会导致竞争和不确定性。所以所有写入都暂存,等本轮所有任务结束后统一合并。
3. 超步结束:apply_writes—— 统一应用,版本号升级
每轮超步结束后(所有节点执行完毕),主循环调用apply_writes,它做四件关键事:
记账:对于本轮已执行的每个节点,把它的
triggers中每个 channel 的当前版本号记录到versions_seen[node]里。这相当于“此节点已经见过这些 channel 的当前版本,下一次别因为同一个版本再触发它”。消费一次性 channel:对于某些 channel(如
EphemeralValue),如果被读取过,就清空其内容(consume())。真正写入并升版本:
遍历所有暂存的写入,按 channel 分组,调用各 channel 的update()方法。
如果update()返回True(表示 channel 确实发生了变化),则给该 channel 分配一个新的全局版本号(next_version,全局递增整数),并把该 channel 加入updated_channels集合。
注意:版本号是每一步全局递增,不是每个 channel 独立计数。所以即使写入的值和上次一样,只要update()返回True,版本号也会增加,下游就会认为“发生了变化”。处理“空写入”(保证临时 channel 的过期):
对每一个 channel,如果它没有被本轮显式写入,但它的类型要求“每步清空”(如EphemeralValue),则调用update(EMPTY_SEQ)将其清空,同样也会触发版本号递增。这就是边信号branch:to:X只存活一个超步的原因——下一步它就会被清空,从而不会反复触发下游。
apply_writes最终返回updated_channels—— 本轮哪些 channel 被更新了(包括被清空的)。
4. 下一超步开始:prepare_next_tasks—— 谁被触发?
主循环的tick()调用prepare_next_tasks,它做两件事:
先消费 PUSH 任务:如果之前有代码通过
Send显式推送任务到TASKS队列,则直接从队列里取出,生成任务,不走版本号判定(优先级高)。再筛选 PULL 候选节点(即普通边触发的节点):
拿着updated_channels,通过倒排索引trigger_to_nodes找出所有可能被触发的节点(候选集合)。
然后对每个候选节点调用_triggers()函数做最终裁定。
_triggers()函数的逻辑是:遍历节点订阅的所有 trigger channel,如果该 channel当前可用(is_available()为 True),且它的当前版本号大于该节点上次见过的版本号(记录在versions_seen中),则触发该节点;否则不触发。
这里的is_available()对于普通边信号(EphemeralValue)只要存在值就为 True;对于 join channel 则是所有前驱都到齐才为 True;对于LastValue(持久 state)只要有值就为 True。
如果某个 channel 被写入了但节点还没执行过(seen为 None),那么只要 channel 可用就触发(因为版本比较时 null 版本 < 当前版本)。
触发后,会从节点订阅的channels(注意和triggers区分,channels是输入来源,即 state keys)读取当前 state,组装成任务的输入,然后调度执行。
5. 完整示例:A → B 的时序
假设图中只有一条边A → B,初始状态 A 未执行过。
| 阶段 | 发生了什么 |
|---|---|
| 编译后 | A 的 writers 里有向branch:to:B写入的指令;B 的 triggers 包含branch:to:B。倒排索引:branch:to:B → [B]。 |
| 超步 0 开始 | prepare_next_tasks发现 A 是起始节点(或通过add_edge(START, A)设定),触发 A。A 开始执行。 |
| 超步 0 执行中 | A 返回新 state,并产生写入:向branch:to:B写None,以及可能向 state keys 写新值。这些写入暂存在task.writes。 |
| 超步 0 结束 | apply_writes:① 记录versions_seen[A](把branch:to:A当前版本记下);② 应用写入:branch:to:B的update(None)返回 True,版本号从 0 升至 1,updated_channels= {branch:to:B};同时 state key 的LastValue也更新,版本号同样升至 1(但 state keys 通常不在 triggers 中,所以不会触发任何节点);③ 清空临时 channel(本步无)。 |
| 超步 1 开始 | prepare_next_tasks:通过updated_channels查倒排,候选节点 = {B};_triggers(B)发现branch:to:B版本 1 >versions_seen[B](空,视为 -∞),成立 → 触发 B。B 从 state keys 读取输入并执行。 |
| 超步 1 执行中 | B 产生写入(可能没有新边)。 |
| 超步 1 结束 | apply_writes:记录versions_seen[B](记下branch:to:B版本 1);然后对branch:to:B调用update(EMPTY_SEQ)(因为它是EphemeralValue,每步末尾会被清空),清空后is_available()变为 False,版本升至 2,但_triggers要求is_available()为 True,所以后续不会再触发 B。 |
| 后续 | 如果 B 没有产生新的边信号,则所有 channel 在后续超步中会被finish()处理,图结束。 |
6. 不同 Channel 类型及其触发语义
LangGraph 通过多态 Channel 类实现不同触发行为,核心方法有:
update(values):应用一批写入,返回bool表示是否有变化(用于决定是否升版本)。is_available():当前是否有可读数据(用于_triggers判定)。consume():读取后清空(一次性 channel)。finish():图结束时调用,用于延迟释放。
常见类型:
| Channel 类 | 用途 | 触发特点 |
|---|---|---|
LastValue | 存储 state 的 key(默认) | 持久保存;只要有值就is_available();update任何非空值都会返回 True,所以即使值相同也会升版本,因此会触发订阅它的节点(但通常 state keys 不在 triggers 中,所以不触发节点)。 |
EphemeralValue | 边信号branch:to:X | 写入后is_available()为 True;每步结束后自动update(EMPTY_SEQ)清空,清空后is_available()为 False,实现一次性触发。 |
NamedBarrierValue | 扇入(join) | 需要所有名前驱都写入后才is_available()为 True,实现同步屏障。 |
LastValueAfterFinish | 用于defer=True的节点 | 只有在图即将结束(finish()被调用)时才变为可用,实现“最后执行”的效果。 |
Topic | 用于Send队列 | 队列语义,prepare_next_tasks直接逐条取出作为 PUSH 任务,不依赖版本号。 |
7. 为什么用版本号而不是值比较?
- 确定性:版本号随每步递增,与具体值无关,保证同样的执行历史一定产生相同的版本序列,便于 checkpoint 恢复、时间旅行(replay)。
- 简化并发:多节点并行写入同一个 channel 时,只需比较版本号即可决定是否触发,无需深度比较值是否变化。
- 兼容中断(interrupt):
should_interrupt也基于版本号比较,统一了调度和中断逻辑。
总结
- LangGraph 的节点触发是编译期订阅 + 运行期版本号拉取。
- 边不是函数调用,而是发布信号到 channel。
- 写入是批量异步的,每步结束统一应用。
- 版本号是全局递增的,
is_available()和版本号共同决定触发。 - Channel 多态实现了边信号、扇入、延迟等丰富语义。
理解了这个模型,你就掌握了 LangGraph 调度的核心。调试时如果发现节点意外触发或未触发,可以检查:
- 它的 trigger channel 是否在
updated_channels中? - 当前版本号是否大于
versions_seen中的记录? - 该 channel 的
is_available()是否为 True?