PregelLoop:超步循环驱动器
职责
PregelLoop 是 LangGraph 运行时的「节拍器」。Pregel 本身只负责组装资源、暴露 invoke / astream 入口,真正一圈圈跑图的是一个 PregelLoop 实例:它持有当前检查点 (checkpoint)、通道 (channel) 集合、pending_writes、step 计数器、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 的活;它只决定「这一步该看谁、跑完怎么收尾、下一步还要不要继续」。
设计动机
为什么把循环逻辑拆成 tick 和 after_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_writesfuture,_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_nodes、should_interrupt。out_of_steps:607-609—self.step > self.stop时进入out_of_steps状态并返回False,外层while退出。done branch:653-655—self.tasks为空时进done状态,这是图的「自然死亡」。draining branch:657-659—control.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列表里,命中就 raiseGraphInterrupt。after_tick method:683-726— 收 writes →apply_writes→ emit values → 清 pending →_put_checkpoint→interrupt_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_writesfuture。SyncPregelLoop:1469/AsyncPregelLoop:1722— 两个子类,覆盖__enter__/__exit__、accept_push、schedule_error_handler等同步/异步差异部分。
数据流
tick 的入口先做步数上限检查,接着调 prepare_next_tasks 把「下一步要跑哪些 task」算出来。这步是整个超步的种子——后面所有动作都是对这批 task 的加工。
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-629。self.step > self.stop 是「步数上限」退出条件——stop 在 PregelLoop.__init__ 时被设成 recursion_limit,递归深度超了就停在 out_of_steps。prepare_next_tasks 拿当前 checkpoint、pending_writes、节点表、通道表算出本步要执行哪些 task,塞进 self.tasks。for_execution=True 表示这些 task 是要真跑的(而非 dry-run 用于流式预览)。
紧接着是三个并列的退出条件——空 task / drain / 续跑贴回 writes——再走 should_interrupt:
# 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-671。interrupt_before 的检查只在有 pending_writes 续跑之后才做——因为续跑时通道版本号已经反映了上次的写入,should_interrupt 才能正确判断「自上次 interrupt 以来有没有新更新」。
after_tick 在 runner 把 task 全部跑完后调用,负责收尾:
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_replaying 置 False、_put_checkpoint({"source": "loop"}) 落检查点,然后跑 interrupt_after 判断(after_tick 尾部:714-724)。
整个超步的生命周期可以画成这样:
边界与失败
out_of_steps不是错误:step > stop直接返回False,status 被设成out_of_steps,外层while自然退出。Pregel上层会根据 status 判断是 throwRecursionError还是静默收尾——这是「软上限」语义,不是异常路径(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_writesfuture,下个 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。