Day 20 / 共 60 天 · 阶段 4 Pregel 执行引擎(核心深水区)

超步主循环:PregelLoop.tick() 的心跳

昨天知道了"图一波一波跑"。今天找到那个摇动发条的手——_loop.pyPregelLoop。它的 tick() 负责每个超步的"要不要再来一步 + 计划本步任务 + 该不该中断",after_tick() 负责"合并写入 + 存档"。而 invoke/stream 里那句 while loop.tick(): ... loop.after_tick() 就是整台引擎的心跳。读 _loop.py:599main.py:2964

📍 你在 60 天里的位置
①入门 D1-6· ②状态 D7-12· ③控制流 D13-18· ④Pregel D19-26· ⑤通道 D27-32· ⑥持久化 D33-40· ⑦中断 D41-46· ⑧函数式 D47-52· ⑨预制件 D53-58· ⑩收官 D59-60
D19 BSP模型 D20 主循环 D21 任务准备 D22 执行器 D23 读写 D24 IO映射 D25 重试超时 D26 调试画图
L01

谁在摇发条:PregelLoop 是什么

🤔 痛点:昨天说"重复超步直到没节点可跑",可这个"重复"到底是谁写的 for/while? 昨天讲的是模型(应该怎么跑),今天要看的是实现(代码里谁真的在循环)。这个循环体必须做很多事:判断该不该停、给这一步挑任务、在合适时机中断、把写入合并、存档……如果全塞进一个函数会乱成一团。LangGraph 的做法是拆成两半——tick() 管"执行之前",after_tick() 管"执行之后",中间夹着真正跑节点的执行器(Day 22)。

PregelLoop_loop.py:158)是一次 invoke/stream 的"运行时上下文":它持有当前 checkpoint、channels、待执行的 tasks、当前步号 step、上限 stop、状态 status 等。它有同步、异步两个子类(_loop.py:1469 SyncPregelLoop / _loop.py:1722 AsyncPregelLoop),共享同一份 tick/after_tick 逻辑。

💡 本质:tick = "还要不要再来一步 + 这一步跑什么",after_tick = "这一步的账怎么结" 记住这个分工,后面所有代码都好懂:tick() 返回 True 表示"还有活干、这些任务已备好",返回 False 表示"到头了"。执行器把 tick 备好的任务跑完后,after_tick() 把大家的写入合并进 channel、存一个 checkpoint,然后回到 while 再 tick。就是"备菜 → 炒菜 → 装盘上桌"三步循环。
类比:tick 是发牌员,after_tick 是记分员 打一局牌:tick = 发牌员看牌局判断"这局还能不能打,能打就发好这轮每个人的手牌";玩家各自出牌(执行器并行跑节点);after_tick = 记分员把这轮所有人出的牌统一算分、记进账本(合并写入 + 存档)。发牌员说"没牌可发了"(tick 返回 False),一局结束。
L02

驱动循环:invoke 里那句 while loop.tick()

先看"心跳"本身。在 Pregel.stream() 里(invoke 内部就是收集 stream 的最后一份),主循环长这样(main.py:2964-2988):

# main.py:2960 —— 注释直接说明了 BSP:N 步的写入只在 N+1 步可见
# computation proceeds in steps, while there are channel updates.
# Channel updates from step N are only visible in step N+1
while loop.tick():                                  # 2964 备好本步任务;无任务→False→退出
    for task in loop.match_cached_writes():         # 2965 命中缓存的任务直接产出
        loop.output_writes(task.id, task.writes, cached=True)
    for _ in runner.tick(                           # 2967 执行器:并行跑本步任务(Day 22)
        [t for t in loop.tasks.values() if not t.writes],
        timeout=self.step_timeout,                  # 2969 单步墙钟超时
        get_waiter=get_waiter,
        schedule_task=loop.accept_push,             # 2971 支持运行中动态追加 Send 任务
    ):
        yield from _output(stream_mode, ...)        # 2974 每有任务完成就流式产出
    loop.after_tick()                               # 2984 合并写入 + 存档
    if durability_ == "sync":                       # 2987 sync 模式:等 checkpoint 落盘
        loop._put_checkpoint_fut.result()
while loop.tick()整台引擎的主循环。tick() 内部完成 Plan(计划任务),返回 True 就继续、False 就结束。
runner.tick(...)Execution 阶段。把 tick 备好的、还没写入的任务丢进执行器并行跑(Day 22)。它是个生成器——每完成一个任务就 yield 一次,好让外层流式吐出。
timeout=self.step_timeoutDay 19 讲的"木桶效应"的解药:给整个超步设墙钟上限。到点还没跑完,执行器就停等,避免一个卡死节点拖垮整个流程。
loop.after_tick()Update 阶段。这一步所有任务跑完后调用一次,合并写入、算触发、存 checkpoint。
💡 本质:Day 19 三阶段在这里"落地成代码" 看这段就明白了:tick() = Plan、runner.tick() = Execution、after_tick() = Update,三者被一个朴素的 while 串起来反复跑。整个复杂的 Pregel 引擎,主干骨架就这十几行。剩下的复杂度都藏在 tick/runner/after_tick 各自内部。
L03

tick() 逐行①:先判停,再计划任务

进入 tick() 本体(_loop.py:599),前半段:

def tick(self) -> bool:                     # _loop.py:599
    """Execute a single iteration of the Pregel loop.
    Returns: True if more iterations are needed."""
    if self.step > self.stop:               # 607 超过步数上限
        self.status = "out_of_steps"
        return False

    self.tasks = prepare_next_tasks(        # 612 Plan:算出本步要跑哪些任务(Day 21)
        self.checkpoint, self.checkpoint_pending_writes,
        self.nodes, self.channels, self.managed, self.config,
        self.step, self.stop, for_execution=True,
        trigger_to_nodes=self.trigger_to_nodes,
        updated_channels=self.updated_channels,   # 626 上一步更新了哪些 channel
        retry_policy=self.retry_policy, cache_policy=self.cache_policy,
    )
    ...
    if not self.tasks:                      # 653 没有任务了 = 图跑完了
        self.status = "done"
        return False
if self.step > self.stop第一道停止闸:步号超过上限,立刻停并标记 out_of_steps(对应 GraphRecursionError,L06 细讲)。防死循环。
prepare_next_tasks(...)Plan 的核心。传入"上一步更新了哪些 channel"(updated_channels)和"channel→节点"反查表,算出这一步该激活哪些节点、把它们的输入读好,打包成 tasks。Day 21 逐行拆。
if not self.tasks: return False第二道、也是正常的停止条件:Plan 完发现没有任何任务要跑——说明没有 channel 触发任何节点,图自然结束(status="done")。这就是 Day 19 说的"没有 actor 被选中就停"。
🎯 设计取舍①:为什么把 updated_channels 一路传进 Plan,而不是每步重新扫全图? 朴素实现:每个超步遍历所有节点,逐个问"你订阅的 channel 有没有新版本?"——图大了就慢。源码的做法是增量触发:只记住"上一步更新了哪几个 channel",配合 trigger_to_nodes 反查表,直接得到"这几个 channel 会触发哪些节点"这一小撮候选,只检查它们(Day 21 会看到 _algo.py:475 这段优化)。用一点记账换来"候选集从全图缩到极小",这是引擎能扛大图的关键。
L04

tick() 逐行②:执行前中断 + 回放

tick() 后半段——在真正放行任务之前,还要处理"回放已有写入"和"静态断点"(_loop.py:661-681):

    # 有上次遗留的 pending writes(比如从 checkpoint 恢复),先补上
    if not self.is_replaying and self.checkpoint_pending_writes:   # 662
        self._reapply_writes_to_succeeded_nodes(self.tasks)
        self._resume_error_handlers_if_applicable()

    # 执行前静态断点:interrupt_before 命中就抛 GraphInterrupt 暂停
    if self.interrupt_before and should_interrupt(                 # 667
        self.checkpoint, self.interrupt_before, self.tasks.values()
    ):
        self.status = "interrupt_before"
        raise GraphInterrupt()

    self._emit("tasks", map_debug_tasks, self.tasks.values())      # 674 debug 流产出(Day 26)

    for task in self.tasks.values():        # 677 已经有缓存写入的任务,直接产出
        if task.writes:
            self.output_writes(task.id, task.writes, cached=True)
    return True                             # 681 一切就绪,放行给执行器
_reapply_writes_to_succeeded_nodes从 checkpoint 恢复时,某些节点上次已经成功、写入还在 pending_writes 里。这里把它们补回去,避免重复执行已完成的节点(幂等恢复,Day 46 会深入)。
interrupt_before + should_interrupt执行前断点:如果配置了 interrupt_before=["tool"],且本步的任务里有 tool,就抛 GraphInterrupt 把图停在"即将执行 tool"处,交回给你审查(人在环,Day 43)。注意:断点是在 tick 里、执行器之前触发的。
return True走到这,说明本步有任务、没被断点拦下——返回 True,主循环随即把 loop.tasks 交给执行器。
💡 本质:tick 不只是"计划",还是所有"执行前拦截"的统一关卡 停止判断、恢复补写、静态断点、debug 事件——凡是"节点真正跑起来之前"要做的事,全收拢在 tick 里。这样执行器(Day 22)就能专心干一件事:把交给它的任务跑完。关注点分离:tick 决策,runner 执行,after_tick 结算。
L05

after_tick():把这一步的账结清

执行器把本步任务都跑完后,主循环调 after_tick()_loop.py:683)——这就是 Day 19 的 Update 阶段:

def after_tick(self) -> None:               # _loop.py:683
    writes = [w for t in self.tasks.values() for w in t.writes]   # 685 收齐本步所有写入
    ...
    self.updated_channels = apply_writes(   # 692 统一合并进 channel(Day 21)
        self.checkpoint, self.channels, self.tasks.values(),
        self.checkpointer_get_next_version, self.trigger_to_nodes,
    )
    # 如果 output 通道被更新了,产出一份 values 流
    if not self.updated_channels.isdisjoint(...):                 # 700
        self._emit("values", map_output_values, self.output_keys, writes, self.channels)
    self.checkpoint_pending_writes.clear()  # 714 清空本步 pending
    self.is_replaying = False               # 716 只有第一步会 replay
    self._put_checkpoint({"source": "loop"})# 718 存档:这一步的 checkpoint
    # 执行后静态断点
    if self.interrupt_after and should_interrupt(                 # 720
        self.checkpoint, self.interrupt_after, self.tasks.values()
    ):
        self.status = "interrupt_after"
        raise GraphInterrupt()
    self.config[CONF].pop(CONFIG_KEY_RESUMING, None)             # 726
writes = [... t.writes ...]把本步所有任务各自攒的写入汇总成一个大列表。注意:执行期间它们只是攒在每个 task 身上(Day 22),到这里才拿出来。这正是 BSP "写入延迟到步末统一生效"的落点。
apply_writes(...)Update 的核心(Day 21 逐行):按每个 channel 的 reducer 合并写入、给 channel 升版本号,返回"这一步到底更新了哪些 channel"——这个 updated_channels 又喂给下一次 tick 做增量触发。闭环了!
_put_checkpoint每个超步末尾存一个 checkpoint(阶段 6)。这就是为什么 LangGraph 能"时间旅行"回到任意一步——每一步的状态都存了档。
interrupt_after执行后断点:和 tick 里的 interrupt_before 对称,命中就在"节点跑完、写入已合并"处停下。
控制流:一次超步内 tick → 执行器 → after_tick 的闭环 tick() 判停·计划·断点 runner.tick() 并行跑任务·攒写入 after_tick() 合并写入·存档 updated_channels 回喂下一次 tick → 增量触发下一步的节点 tick 返回 False(无任务 / 超上限)→ 跳出 while,图结束
图注:tick→runner→after_tick 三段循环,after_tick 产出的 updated_channels 决定下一步谁被触发。
L06

step / stop 与递归上限从哪来

tick 第一行 if self.step > self.stop 里的两个数是怎么来的?初始化时都是 0(_loop.py:301-302):

# _loop.py:301
self.step = 0
self.stop = 0

而真正决定上限,是在从 checkpoint 建立起点时(SyncPregelLoop_loop.py:1700-1701):

# _loop.py:1700
self.step = self.checkpoint_metadata["step"] + 1          # 从上次存档的下一步开始
self.stop = self.step + self.config["recursion_limit"] + 1 # 上限 = 起点 + 递归上限 + 1
step = metadata["step"] + 1step 不总是从 0 开始——如果是从 checkpoint 恢复,就从"上次存到第几步"的下一步接着跑。这就是断点续跑的数字基础。
stop = step + recursion_limit + 1recursion_limit(默认 25)就是你在 config 里能调的那个"最多跑多少步"。它被换算成绝对步号上限 stop。tick 里 step > stop 就触发 out_of_steps
⚠️ 边界/坑:GraphRecursionError 不是"你的代码错了",而是"步数用光了" 很多人第一次撞见 GraphRecursionError: Recursion limit of 25 reached 会以为图写错了。其实它就是这里 step > stop 的结果——通常是你的图有环、但退出条件没写对,一直在超步里打转(比如 agent↔tools 反复横跳停不下来)。解法两个:① 修条件边让它能到 END;② 确实需要更多步就 config={"recursion_limit": 100} 抬高上限。它是防死循环的保险丝,不是 bug。递归上限单位是"超步数",不是"节点调用次数"——一个超步并行跑 5 个节点也只算 1 步。
📝 真实值走一遍 默认 recursion_limit=25,全新运行:step=1, stop=1+25+1=27。第 1、2、…、27 步都能跑;跑到第 28 步开头 28 > 27 为真 → out_of_steps → 抛 GraphRecursionError。所以"25"大致等于"最多 25~27 个超步"。
L07

_first:第一个超步的特殊处理

第一次 tick 之前,引擎要先把你的输入"喂进图",这活由 _first()_loop.py:848)干。它最重要的判断是"这是全新运行,还是从中断恢复?"_loop.py:861-872):

# _loop.py:861
input_is_command = isinstance(self.input, Command)
is_resuming = bool(self.checkpoint["channel_versions"]) and bool(
    configurable.get(
        CONFIG_KEY_RESUMING,
        self.input is None                  # invoke(None, config) = 恢复
        or input_is_command                 # 传 Command 也是对已有状态操作
        or (not self.is_nested and 同一个 run_id)  # stream 重连
    )
)
self.input is None经典恢复姿势:graph.invoke(None, config)。输入是 None + 已有 checkpoint = "别喂新输入,接着上次跑"。这是人在环恢复的常见写法(Day 42)。
input_is_command传了 Command(resume=...)Command(goto=...),也表示对已有状态操作而非全新开始。
bool(channel_versions)前置门槛:channel 有版本记录说明确实存过档。全新图这里是空的,怎么都不算 resume。
💡 本质:一个 invoke 入口,区分"新跑 vs 续跑"就靠这几条线索 用户不需要显式说"我要恢复",引擎从"输入是不是 None/Command + 有没有旧 checkpoint + 是不是同一次 run"这几条线索自动推断。推断对了,就能无缝续跑;这套判断是阶段 7(中断与人在环)的地基。今天只需知道:第一个超步和后续超步的区别,全在 _first 里被抹平——之后每次 tick 都一视同仁。

👶 小白:为什么不加个参数 resume=True 让用户明说,非要"猜"?

👨‍🏫 老师:因为"传 None / 传 Command / 传新输入"这三种写法本身就已经把意图表达清楚了,再加一个布尔参数反而容易和输入自相矛盾(传了新输入又说 resume=True 该听谁的?)。让"输入的形态"自己表达意图,比多一个易错的开关更干净。当然它也留了 CONFIG_KEY_RESUMING 这个内部旗标给子图显式传递——猜不准的场景仍有兜底。

L08

今日小结 + 动手 + 明日预告

🧠 今天你应该能回答

  • 驱动引擎的主循环长什么样?(while loop.tick(): runner.tick(...); loop.after_tick()
  • tick() 和 after_tick() 各管什么?(tick=判停/计划/执行前断点;after_tick=合并写入/存档/执行后断点)
  • tick 返回 False 的两种情况?(step>stop 超步数上限;prepare 完没有任务=图跑完)
  • step/stop 怎么算?(step 从上次存档下一步起;stop=step+recursion_limit+1)
  • GraphRecursionError 本质是什么?(step>stop,步数用光,多半是环没写对退出条件)
  • _first 怎么区分新跑和续跑?(输入是 None/Command + 有旧 checkpoint + 同一 run → resume)
  • updated_channels 的闭环作用?(after_tick 算出它 → 下次 tick 用它增量触发节点)

✋ 10 分钟动手

# 1. 读 tick 主体(判停 + 计划 + 断点)
sed -n '599,681p' libs/langgraph/langgraph/pregel/_loop.py

# 2. 读 after_tick(合并写入 + 存档)
sed -n '683,726p' libs/langgraph/langgraph/pregel/_loop.py

# 3. 看 step/stop 上限怎么定
sed -n '1700,1701p' libs/langgraph/langgraph/pregel/_loop.py

# 4. 亲手触发一次 GraphRecursionError(体会保险丝)
python -c "
from langgraph.graph import StateGraph, START
g = StateGraph(dict); g.add_node('loop', lambda s: {'n': s.get('n',0)+1})
g.add_edge(START, 'loop'); g.add_edge('loop', 'loop')   # 故意死循环
try: g.compile().invoke({'n':0}, {'recursion_limit':5})
except Exception as e: print(type(e).__name__, e)
"
明天预告 · Day 21:今天 tick 里那句 prepare_next_tasks(...) 和 after_tick 里那句 apply_writes(...) 都被我们"跳过内部"了。明天进 _algo.py 把这两个 Pregel 引擎最核心的函数逐行拆开——看它怎么用"版本号 > 已见版本号"决定谁被触发(PULL),怎么处理 Send 动态任务(PUSH),以及 apply_writes 怎么按 reducer 合并、升版本、算下一步触发。
← Day 19 BSP 模型 Day 21 · 任务准备 _algo.py →