发布日期

第 10 讲 · 超步循环(上):Plan 阶段

学习目标

  • 看懂 PregelLoop.tick() 如何规划本超步的任务。
  • 理解 版本触发机制channel_versions vs versions_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 算出本步任务,没任务就结束。

精简流程:

  1. step > stopout_of_steps,返回 False(防止无限循环,对应 recursion_limit)。
  2. prepare_next_tasks(...) → 填充 self.tasks
  3. 若没有任务 → status="done",返回 False(图跑完了)。
  4. 处理 pending writes(恢复/replay 场景,第 18 讲)。
  5. 检查 interrupt_before → 可能抛 GraphInterrupt(第 18 讲)。
  6. 输出 debug/updates 流。
  7. 返回 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.pyprepare_next_tasks(392–513 行)。它返回两类任务的并集

PUSH 任务(动态,来自 Send)

先消费内部 TASKSTopic)通道里的 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_versionsversions_seen,立刻看清"为什么这个节点这步跑了/没跑"。
  • 理解 recursion_limitstep > stop 的判定就是它,循环图跑太多步会触发 GraphRecursionError——这时该检查路由条件是否能正常收敛到 END。

动手实验

  1. 用第 1 讲的循环图,在 prepare_next_tasks 打断点,每步打印 self.tasks 的节点名和 checkpoint["channel_versions"],观察版本如何驱动触发。
  2. 故意写一个永不返回 END 的路由,观察 step > stop 触发的 GraphRecursionError, 并尝试用 config={"recursion_limit": N} 调整上限。

阅读作业

  • 精读 _loop.pytick()(592–674 行)。
  • 精读 _algo.pyprepare_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。