超步主循环:PregelLoop.tick() 的心跳
昨天知道了"图一波一波跑"。今天找到那个摇动发条的手——_loop.py 的 PregelLoop。它的 tick() 负责每个超步的"要不要再来一步 + 计划本步任务 + 该不该中断",after_tick() 负责"合并写入 + 存档"。而 invoke/stream 里那句 while loop.tick(): ... loop.after_tick() 就是整台引擎的心跳。读 _loop.py:599 和 main.py:2964。
谁在摇发条:PregelLoop 是什么
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 逻辑。
驱动循环: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。while 串起来反复跑。整个复杂的 Pregel 引擎,主干骨架就这十几行。剩下的复杂度都藏在 tick/runner/after_tick 各自内部。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 这段优化)。用一点记账换来"候选集从全图缩到极小",这是引擎能扛大图的关键。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 交给执行器。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 对称,命中就在"节点跑完、写入已合并"处停下。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: 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 个超步"。_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。_first 里被抹平——之后每次 tick 都一视同仁。👶 小白:为什么不加个参数 resume=True 让用户明说,非要"猜"?
👨🏫 老师:因为"传 None / 传 Command / 传新输入"这三种写法本身就已经把意图表达清楚了,再加一个布尔参数反而容易和输入自相矛盾(传了新输入又说 resume=True 该听谁的?)。让"输入的形态"自己表达意图,比多一个易错的开关更干净。当然它也留了 CONFIG_KEY_RESUMING 这个内部旗标给子图显式传递——猜不准的场景仍有兜底。
今日小结 + 动手 + 明日预告
🧠 今天你应该能回答
- 驱动引擎的主循环长什么样?(
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)
"
prepare_next_tasks(...) 和 after_tick 里那句 apply_writes(...) 都被我们"跳过内部"了。明天进 _algo.py 把这两个 Pregel 引擎最核心的函数逐行拆开——看它怎么用"版本号 > 已见版本号"决定谁被触发(PULL),怎么处理 Send 动态任务(PUSH),以及 apply_writes 怎么按 reducer 合并、升版本、算下一步触发。