Skip to content

PregelLoop:超步循环驱动器

源码版本1.2.9

职责

PregelLoop 是 LangGraph 运行时的「节拍器」。Pregel 本身只负责组装资源、暴露 invoke / astream 入口,真正一圈圈跑图的是一个 PregelLoop 实例:它持有当前检查点 (checkpoint)、通道 (channel) 集合、pending_writesstep 计数器、status 状态机,以及一份 tasks 字典——这些都是单个超步 (superstep) 内的全部上下文。外层 while loop.tick(): ... loop.after_tick() 的形态就建立在它身上。

它的两个核心方法是 tick(tick:599)和 after_tick(after_tick:683)。前者负责「跑之前」:看是否超出步数上限、调 prepare_next_tasks 把这一步要执行的节点装进 self.tasks、检查中断、处理 drain_requested、把上次断点续跑留下的 pending_writes 重新贴回成功节点上。后者负责「跑之后」:收齐所有 task.writes,调 apply_writes 把它们合进通道,清掉 pending,落检查点,再判断要不要 interrupt_after。两者一前一后,把「读取输入 → 调度 → 执行 → 写回 → 持久化」这个 BSP 循环夹在中间。

换句话说,PregelLoop 不亲自跑节点——那是 PregelRunner 的活;它只决定「这一步该看谁、跑完怎么收尾、下一步还要不要继续」。

设计动机

为什么把循环逻辑拆成 tickafter_tick 两段,而不是一个 step() 方法搞定?

  • 执行权交还给调用者:tick 返回 True 后,调用方(Pregel.astream 的实现体,见 sync main loop:2964-2984)自己负责驱动 runner.tick 把这一批 task 跑完,再回头调 loop.after_tick。这样 PregelLoop 就不需要持有 PregelRunner 的引用,两个对象之间是松耦合的协作关系,而不是嵌套。
  • 状态机显式化:status 字段有 "input" / "pending" / "done" / "draining" / "interrupt_before" / "interrupt_after" / "out_of_steps" 七种取值(status Literal:256-264)。每个状态都对应一种退出原因,外部代码靠它判断该 resume、该 raise、还是该 break,不用猜返回值语义。
  • 断点续跑语义清晰:is_replaying 只在第一步 tick 为真——表示当前这一步是「重放」上次中断前已经跑完的 task,等 after_tick 跑完就关掉(is_replaying = False:716)。配合 _reapply_writes_to_succeeded_nodes 把检查点里残留的成功 writes 贴回内存 task,就能让成功节点不重跑、失败节点进 error handler。
  • Drain 是外部干预入口:RunControl.drain_requested 由 stream 消费方设置(比如客户端要中途停),tick 一看到就进 draining 状态并返回 False,让外层 while 自然退出,而不是在节点中途硬杀。
  • Exit-mode delta 写入与检查点时序:_delta_write_futs(_delta_write_futs:207)收集所有 delta 通道的 put_writes future,_checkpointer_put_after_previous 必须先 drain 它们再写下个 checkpoint——保证「写产生 checkpoint,而不是 checkpoint 盖过写」的因果顺序。

关键文件

  • class PregelLoop:158 — 类定义,字段几乎全是裸 attribute,没有 property 包装,便于子类 SyncPregelLoop / AsyncPregelLoop 直接覆盖。
  • status Literal:256-264 — 七种状态枚举,"out_of_steps" 是步数上限触发,"draining" 是外部请求停止,"done" 是正常结束。
  • tick method:599-681 — 一个完整 tick 的全部逻辑:步数检查、prepare_next_tasks、空任务判 done、drain、_reapply_writes_to_succeeded_nodesshould_interrupt
  • out_of_steps:607-609self.step > self.stop 时进入 out_of_steps 状态并返回 False,外层 while 退出。
  • done branch:653-655self.tasks 为空时进 done 状态,这是图的「自然死亡」。
  • draining branch:657-659control.drain_requested 为真时进 draining 状态。
  • reapply + resume handlers:662-664 — 续跑时先把 pending writes 贴回成功节点,再调 _resume_error_handlers_if_applicable 给失败节点准备 handler task。
  • interrupt_before:667-671 — 调 should_interrupt 检查节点是否落在 interrupt_before 列表里,命中就 raise GraphInterrupt
  • after_tick method:683-726 — 收 writes → apply_writes → emit values → 清 pending → _put_checkpointinterrupt_after 检查。
  • _reapply_writes_to_succeeded_nodes:736-749 — 把 pending writes 贴回内存 task,但跳过 ERROR / ERROR_SOURCE_NODE / INTERRUPT / RESUME 这四种控制信号。
  • _resume_error_handlers_if_applicable:751-816 — 扫 ERROR_SOURCE_NODE 标记,给失败节点构造 error-handler task,让 runner 跳过原任务直接跑 handler。
  • put_writes:415-508 — task 产出 writes 的入口,去重控制通道、累积 null task、起 checkpointer 的 put_writes future。
  • SyncPregelLoop:1469 / AsyncPregelLoop:1722 — 两个子类,覆盖 __enter__/__exit__accept_pushschedule_error_handler 等同步/异步差异部分。

数据流

tick 的入口先做步数上限检查,接着调 prepare_next_tasks 把「下一步要跑哪些 task」算出来。这步是整个超步的种子——后面所有动作都是对这批 task 的加工。

python
def tick(self) -> bool:
    """Execute a single iteration of the Pregel loop.

    Returns:
        True if more iterations are needed.
    """

    # check if iteration limit is reached
    if self.step > self.stop:
        self.status = "out_of_steps"
        return False

    # prepare next tasks
    self.tasks = prepare_next_tasks(
        self.checkpoint,
        self.checkpoint_pending_writes,
        self.nodes,
        self.channels,
        self.managed,
        self.config,
        self.step,
        self.stop,
        for_execution=True,
        manager=self.manager,
        store=self.store,
        checkpointer=self.checkpointer,
        trigger_to_nodes=self.trigger_to_nodes,
        updated_channels=self.updated_channels,
        retry_policy=self.retry_policy,
        cache_policy=self.cache_policy,
    )

这段来自 tick 头部:599-629self.step > self.stop 是「步数上限」退出条件——stopPregelLoop.__init__ 时被设成 recursion_limit,递归深度超了就停在 out_of_stepsprepare_next_tasks 拿当前 checkpoint、pending_writes、节点表、通道表算出本步要执行哪些 task,塞进 self.tasksfor_execution=True 表示这些 task 是要真跑的(而非 dry-run 用于流式预览)。

紧接着是三个并列的退出条件——空 task / drain / 续跑贴回 writes——再走 should_interrupt:

python
# if no more tasks, we're done
if not self.tasks:
    self.status = "done"
    return False

if self.control is not None and self.control.drain_requested:
    self.status = "draining"
    return False

# if there are pending writes from a previous loop, apply them
if not self.is_replaying and self.checkpoint_pending_writes:
    self._reapply_writes_to_succeeded_nodes(self.tasks)
    self._resume_error_handlers_if_applicable()

# before execution, check if we should interrupt
if self.interrupt_before and should_interrupt(
    self.checkpoint, self.interrupt_before, self.tasks.values()
):
    self.status = "interrupt_before"
    raise GraphInterrupt()

这段来自 tick 中段:652-671interrupt_before 的检查只在有 pending_writes 续跑之后才做——因为续跑时通道版本号已经反映了上次的写入,should_interrupt 才能正确判断「自上次 interrupt 以来有没有新更新」。

after_tick 在 runner 把 task 全部跑完后调用,负责收尾:

python
def after_tick(self) -> None:
    # finish superstep
    writes = [w for t in self.tasks.values() for w in t.writes]
    self._delta_channels_with_overwrite.update(
        ch
        for ch, v in writes
        if isinstance(self.specs.get(ch), DeltaChannel) and _get_overwrite(v)[0]
    )
    # all tasks have finished
    self.updated_channels = apply_writes(
        self.checkpoint,
        self.channels,
        self.tasks.values(),
        self.checkpointer_get_next_version,
        self.trigger_to_nodes,
    )

这段来自 after_tick 头部:683-698。先把所有 task 的 writes 摊平,顺手统计哪些 delta 通道发生了 overwrite(影响 sparse replay 的起始值);再调 apply_writes 把这些 writes 真正合进通道,返回值 updated_channels 会被下一步的 prepare_next_tasks 用来加速「找下一个该触发的节点」。after_tick 的尾部还会清空 checkpoint_pending_writes、把 is_replayingFalse_put_checkpoint({"source": "loop"}) 落检查点,然后跑 interrupt_after 判断(after_tick 尾部:714-724)。

整个超步的生命周期可以画成这样:

边界与失败

  • out_of_steps 不是错误:step > stop 直接返回 False,status 被设成 out_of_steps,外层 while 自然退出。Pregel 上层会根据 status 判断是 throw RecursionError 还是静默收尾——这是「软上限」语义,不是异常路径(out_of_steps:607-609)。
  • draining 不能 emit lifecycle events:_emit_graph_lifecycle_event 显式拒绝在 draining 状态下被调用(draining guard:384-385),因为 drain 是用户主动停,不算 interrupt / resume 语义。
  • 续跑时 pending writes 含控制信号要跳过:_reapply_writes_to_succeeded_nodes 必须跳过 ERROR / ERROR_SOURCE_NODE / INTERRUPT / RESUME(skip control signals:746-747),否则失败过的 task 会被当成成功的,跑过的 handler 又被重新触发。
  • Resume 时的 error handler 重新调度:_resume_error_handlers_if_applicable 只对带 error_handler_node 的节点生效(handler_node check:791-793),没有 handler 的失败节点保持原样,由 runner 的 _should_stop_others 走 panic 路径。
  • put_writes 的 null task 累积语义:NULL_TASK_ID 的 writes 不覆盖而是累积(null task accumulate:422-431),用于 input writes 这种「不属于任何节点但要进 checkpoint」的写入。
  • Delta channel 的 future 必须先于下个 checkpoint 落盘:_delta_write_futs 收集了所有 delta 通道的 put_writes future,下个 checkpoint 落盘前会先 drain 它(_delta_write_futs:201-207)——颠倒顺序会让 sparse replay 看不到产生 checkpoint 的那些写入。

小结

PregelLoop 是一个状态机驱动的超步节拍器,把「准备 task → 跑 → 收尾」三段夹住 PregelRunner,自己不碰执行。理解了 tick / after_tick 两段、status 七种状态、以及 is_replaying + _reapply_writes_to_succeeded_nodes 的续跑机制,就抓住了 LangGraph 运行模型的主轴。Runner 怎么把 task 真正跑起来见 /pregel/runner;算法层 prepare_next_tasks / apply_writes / should_interrupt 的内部细节见 /pregel/algo;整体怎么组装见 /pregel/pregel。通道的写回语义见 /channel/base-channel

对照官方资料:LangGraph 文档 · README