- 发布日期
第 10 讲 · 超步循环(上):Plan 阶段
学习目标
- 看懂
PregelLoop.tick()如何规划本超步的任务。 - 理解 版本触发机制:
channel_versionsvsversions_seen。 - 理解
prepare_next_tasks如何枚举 PULL(边触发)与 PUSH(Send)任务。
PregelLoop:超步状态机
libs/langgraph/langgraph/pregel/_loop.py,类 PregelLoop(156 行起), 分同步 SyncPregelLoop 与异步 AsyncPregelLoop。它是超步循环的状态机,持有:
channels:当前所有通道(活的内存对象)。checkpoint:含channel_versions(每个通道的版本)与versions_seen(每个节点见过的版本)。tasks:本超步要执行的任务。step/stop:当前步号与上限。
tick():Plan 阶段主函数
tick()(592–674 行)。它的核心就一句:调 prepare_next_tasks 算出本步任务,没任务就结束。
精简流程:
- 若
step > stop→out_of_steps,返回False(防止无限循环,对应recursion_limit)。 prepare_next_tasks(...)→ 填充self.tasks。- 若没有任务 →
status="done",返回False(图跑完了)。 - 处理 pending writes(恢复/replay 场景,第 18 讲)。
- 检查
interrupt_before→ 可能抛GraphInterrupt(第 18 讲)。 - 输出 debug/updates 流。
- 返回
True,进入 Execute。
返回值的含义:True = 还有任务要跑,继续循环;False = 结束。 这就是第 9 讲 while loop.tick() 的判定来源。
核心:版本触发机制
这是整个 LangGraph 调度的灵魂。回顾:
- 每个通道有一个单调递增的版本号,存在
checkpoint["channel_versions"]。 - 每个节点记录它上次见过的各通道版本,存在
checkpoint["versions_seen"][节点名]。 - 触发判定:若某节点订阅的某个通道,其当前版本 > 该节点见过的版本 → 该节点被触发。
flowchart LR
subgraph 通道
CV["channel_versions<br/>messages: 5"]
end
subgraph 节点agent
VS["versions_seen[agent]<br/>messages: 3"]
end
CV -->|5 > 3| TRIG[触发 agent 执行]
TRIG -.执行后.-> UPD["versions_seen[agent]<br/>messages: 5"]
判定逻辑在 _algo.py 的 _triggers(约 1260–1277 行):比较节点订阅通道的版本与 versions_seen。 执行后,apply_writes(第 11 讲)会更新 versions_seen,避免同一份更新反复触发。
这解释了两类常见现象:
- "节点没执行":它订阅的触发通道版本没涨(上游没往它的
branch:to:X写)。- "节点重复执行":每步都有新写让版本涨(典型是循环图,符合预期)。
prepare_next_tasks:枚举任务
_algo.py,prepare_next_tasks(392–513 行)。它返回两类任务的并集:
PUSH 任务(动态,来自 Send)
先消费内部 TASKS(Topic)通道里的 Send 包,每个 Send 变成一个 PUSH 任务 (fan-out,第 19 讲)。Send 携带自己的输入 arg,作为目标节点的输入。
PULL 任务(静态,来自边触发)
枚举候选节点,对每个调 prepare_single_task 判断是否被触发:
flowchart TD
A[prepare_next_tasks] --> B{有 updated_channels<br/>且有反向索引?}
B -->|是| C[只检查相关节点<br/>trigger_to_nodes 优化]
B -->|否| D[首步/兜底: 扫全部节点]
C --> E[对每个候选 prepare_single_task]
D --> E
E --> F{_triggers 版本比较<br/>有更新?}
F -->|是| G[生成 PULL 任务]
F -->|否| H[跳过]
优化点:trigger_to_nodes 反向索引(第 9 讲)让 Plan 不必每步扫全图—— 只检查"上一步真正变化的通道"所订阅的节点。这对大图性能很关键。
prepare_single_task:把节点实例化成可执行任务
prepare_single_task(524–761 行)。被触发的节点会被包装成 PregelExecutableTask, 其中关键是给任务的 config 注入几个回调(第 13 讲细讲):
CONFIG_KEY_SEND:节点写操作的落点(写进task.writes缓冲)。CONFIG_KEY_READ:节点读状态的入口(local_read,支持读到本任务刚写的值)。CONFIG_KEY_SCRATCHPAD:interrupt/resume 计数器(第 18 讲)。
使用场景
- 排查调度问题:在
prepare_next_tasks/_triggers打断点,打印channel_versions与versions_seen,立刻看清"为什么这个节点这步跑了/没跑"。 - 理解 recursion_limit:
step > stop的判定就是它,循环图跑太多步会触发GraphRecursionError——这时该检查路由条件是否能正常收敛到 END。
动手实验
- 用第 1 讲的循环图,在
prepare_next_tasks打断点,每步打印self.tasks的节点名和checkpoint["channel_versions"],观察版本如何驱动触发。 - 故意写一个永不返回 END 的路由,观察
step > stop触发的GraphRecursionError, 并尝试用config={"recursion_limit": N}调整上限。
阅读作业
- 精读
_loop.py的tick()(592–674 行)。 - 精读
_algo.py的prepare_next_tasks(392–513 行)与_triggers(1260–1277 行)。
小结
tick()= Plan:调prepare_next_tasks算任务,无任务则结束循环。- 版本触发:
channel_versions > versions_seen⇒ 节点被触发,这是调度的灵魂。 - 任务 = PUSH(Send 动态)+ PULL(边触发);
trigger_to_nodes反向索引做性能优化。
下一讲:after_tick() 的 Update 阶段——写如何合并、版本如何 bump。