Day 46 / 共 60 天 · 阶段 7 中断与人在环

replay 与幂等:重跑一遍,为什么结果不会乱

阶段 7 收官。前面反复出现一个词——重跑:interrupt 恢复要重跑节点、时间旅行要从旧档重放、崩溃恢复要重放已完成的超步。重跑的最大风险是做重复的事:同一个写被应用两遍、reducer 累加两次、状态翻倍。今天看 LangGraph 如何用一套"重放 + 幂等"机制,做到"至少执行一次,但效果等价于恰好一次"

📍 阶段 7 · 中断与人在环(6 天)—— 收官日
D41 interrupt原理 D42 Command恢复 D43 静态断点 D44 时间旅行 D45 durability D46 replay幂等
L01

痛点:恢复时前面的节点会不会又跑一遍?

🤔 痛点一个超步里有节点 A、B、C 并行。A、B 都跑完了、C 还在跑时进程崩了。恢复时,框架读回这个超步的存档——问题来了:A、B 明明跑完了,恢复会不会把 A、B 再跑一遍?如果 A 里写了 count += 1(reducer 累加),重跑就变成加两次,状态就错了!
💡 本质:区分"结果已存的" vs "确实要重跑的"LangGraph 的策略很聪明:崩溃时,已完成任务的写(A、B 的产出)早已通过 put_writes 落进了 pending_writes(D45 讲的每步存)。恢复时,框架不重新执行 A、B,而是把它们存好的写直接复原到内存里——等于"认账"它们已经干完了。只有真正没完成的 C(写为空)才重新执行。这样:完成的不重做,没完成的补上,合起来效果 = 每个任务恰好执行一次。
类比:搬家搬到一半停电。来电后你不会把已经搬上车的箱子再搬一遍——你看一眼清单(pending_writes),已装车的打勾跳过,只搬还在地上的。清单就是幂等的钥匙。
L02

is_replaying:这一 tick 是"重放"还是"新算"

核心开关是 is_replaying,在 loop 初始化时按"有没有指定 checkpoint_id"决定(_loop.py:315):

# _loop.py:173
    is_replaying: bool
# _loop.py:315
        self.is_replaying = CONFIG_KEY_CHECKPOINT_ID in config[CONF]

它只在第一个 tick 为真,之后立刻关掉(_loop.py:715-716,在 after_tick 里):

# _loop.py:715
        # only replay (re-execute) done tasks on the first tick
        self.is_replaying = False
CONFIG_KEY_CHECKPOINT_ID in config恢复/时间旅行时,config 里会带上"要从哪份 checkpoint 起"的 id → is_replaying=True。全新运行没这个 id → False。
只第一个 tick 为真关键:重放只发生在恢复后的第一步——那一步要"消化"存档里已完成的写。从第二步起就是全新计算了。
after_tick 里置 False第一个超步处理完,立刻 is_replaying = False。注释写得明白:"only replay done tasks on the first tick"。
💡 本质:重放是"一次性的开机自检"恢复不是"把整段历史重跑一遍",而是只在接续点那一步做一次特殊处理:把已完成的写复原、把没完成的补跑。过了这一步,图就回到正常前进模式。所以"重放"开销极小——不是重演历史,是接上断点。
💡 设计取舍②:为什么重放只在第一个 tick、而不是整个恢复过程持续?因为"需要甄别已完成 vs 未完成"的只有接续点那一步。恢复时框架从某个 checkpoint 起步,那份 checkpoint 里可能夹着"半完成的超步"(部分任务写了、部分没写,如 L01 的 A/B 完成 C 未完成)——只有这一步要做"复原完成的、重跑未完成的"精细活。一旦这步理清、状态归位,后面全是从干净状态出发的全新计算,和首次运行没区别,再带着 replay 标志只会徒增判断开销、甚至误把新产出的写当成"旧的已完成写"跳过。所以 after_tick 里第一步一结束就 is_replaying=False,把特殊模式的作用域精确限制在唯一需要它的那一步。这是"特殊逻辑作用域最小化"的典范——能用普通路径的地方绝不让特殊标志多停留一个 tick。
L03

复用已成功的写:不重跑,直接认账

tick 开始时,如果不是重放且有遗留写,就把它们复原到任务上(_loop.py:661-664):

# _loop.py:661
        # 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()

复原逻辑 _loop.py:736-749

# _loop.py:736
    def _reapply_writes_to_succeeded_nodes(self, tasks) -> None:
        """Restore successful channel writes from checkpoint to in-memory tasks.
        Skips control signals (ERROR, ERROR_SOURCE_NODE, INTERRUPT, RESUME)
        so that failed/interrupted tasks remain with empty writes and will be
        re-executed (or routed to error handlers) by the runner."""
        for tid, k, v in self.checkpoint_pending_writes:
            if k in (ERROR, ERROR_SOURCE_NODE, INTERRUPT, RESUME):
                continue                                    # ① 控制信号跳过
            if task := tasks.get(tid):
                task.writes.append((k, v))                  # ② 普通写复原到任务
遍历 pending_writes存档里记着"上次各任务写了什么"。逐条看。
if k in (ERROR,..,INTERRUPT,RESUME): continue只复原"普通通道写",跳过控制信号。为什么?下一讲专门说——这是"完成的复用、没完成的重跑"的分界线。
task.writes.append((k,v))把普通写塞回任务的 writes。有 writes 的任务,runner 视为"已完成",不会再执行它
runner 只跑空 writes 的任务关键机制:runner 只执行 not t.writes 的任务。复原了写 = 非空 = 跳过;A、B 因此不重跑。
💡 设计取舍①:为什么用"writes 非空就跳过"来表达"已完成",而不是单独存一个 completed 标记?朴素做法:给每个任务存个 done: bool问题:done 和 writes 是两份信息,得同步维护,还可能不一致(标了 done 但 writes 没存)。源码做法"有没有写"本身就是"完成没完成"的唯一真相——任务完成必然产出写(哪怕是 NO_WRITES 哨兵,见 D41 runner),写落了盘就是"干完了"的铁证。恢复时把写复原 → 任务自动"变成已完成态"。用同一份数据(writes)同时表达"产出了什么"和"完成没有",消除了双写不一致的可能。runner 的判据 not t.writes 简单到无懈可击。
L04

为什么跳过 ERROR/INTERRUPT/RESUME?

上一讲那句 if k in (ERROR, ERROR_SOURCE_NODE, INTERRUPT, RESUME): continue 是幂等的灵魂。函数注释(_loop.py:739-743)解释了动机:

🍼 翻译那段注释"跳过控制信号(ERROR/INTERRUPT/RESUME),好让失败或被中断的任务保持空 writes,这样它们会被 runner 重新执行(或路由到错误处理器)。"
普通写 → 复原 → 跳过重跑成功完成的任务,产出的是普通通道写。复原它 = 认账已完成 = 不重跑。(L03)
INTERRUPT 写 → 不复原被中断的任务,存档里是 (INTERRUPT, ...)(D41)。故意不复原 → 该任务 writes 仍为空 → runner 重新执行它。这正是 D42 说的"恢复时节点从头重跑"!中断任务必须重跑,才能让 interrupt() 这次拿到恢复值、走完后续逻辑。
ERROR 写 → 不复原失败的任务,存的是 (ERROR, 异常)。不复原 → 空 writes → 重跑(给它再来一次的机会,或交给错误处理器)。
RESUME 写 → 不复原恢复值是要被 interrupt() 重新消费的(D42 的 get_null_resume),不能当"已完成的普通产出"复原掉。
💡 本质:一条 continue 划出了"重放"与"重跑"的分水岭同样是 pending_writes 里的记录,普通写代表"已成功的成果"→ 复原、不重跑;控制信号(ERROR/INTERRUPT/RESUME)代表"未了结的状态"→ 不复原、逼其重跑。幂等的精髓不是"什么都不重跑",而是精确区分哪些该重跑、哪些不该。成功的绝不重做(防翻倍),失败/中断的必须重做(才能推进)。一行 if...continue 把这两类干净利落地分开。
L05

put_writes:存的时候就去重

幂等的另一半在写入端。put_writes_loop.py:415-437)存写时会清掉同一任务的旧写再存新的:

# _loop.py:415
    def put_writes(self, task_id: str, writes: WritesT) -> None:
        if not writes:
            return
        # deduplicate writes to special channels, last write wins
        if all(w[0] in WRITES_IDX_MAP for w in writes):
            writes = list({w[0]: w for w in writes}.values())   # ① 特殊通道去重
        if task_id == NULL_TASK_ID:
            ...  # null 任务的写是累积的(如多个 RESUME)
        else:
            # remove existing writes for this task
            self.checkpoint_pending_writes = [
                w for w in self.checkpoint_pending_writes if w[0] != task_id  # ② 先删旧
            ]
            writes_to_save = writes                                          # ③ 再存新
WRITES_IDX_MAP 去重特殊通道(ERROR/RESUME 等)的写"最后一次为准"——用字典按通道名去重,同通道只留最后一条。
remove existing for this task关键:普通任务存写前,先把 pending_writes 里该 task_id 的旧写全清掉。所以一个任务重跑并重新 put_writes 时,不会新旧并存、不会累积两份。
NULL_TASK_ID 累积例外:null 任务(如恢复值)的写是累积的——因为可能要攒多个 RESUME 值。用不同策略。
回顾 D35 L05:InMemorySaver 的 put_writes 还有一层"普通写 (task_id, idx) 已存在则跳过、负数特殊写允许覆盖"的幂等(checkpoint/memory/__init__.py:494)。两层幂等:loop 层 put_writes 管"内存 pending 的新旧替换",checkpointer 层 put_writes 管"落盘时按 (task,idx) 去重"。双保险。
控制流:崩溃恢复的第一个 tick(幂等重放) 读回 pending_writes A:普通写 B:普通写 C: (无,未完成) A/B 普通写 → 复原 writes 非空 → runner 跳过 C 无写 → 不复原 writes 空 → runner 重跑 INTERRUPT/ERROR 跳过 → 逼其重跑 结果 每任务恰好一次 效果等价
图注:A/B 已完成(有写)→ 复原跳过;C 未完成 + 中断/错误 → 空写重跑。合起来"恰好一次"。
L06

子图的重放:ReplayState 再看一眼

D44 已见过 ReplayState_replay.py:14),这里从"幂等"角度再理解它。子图重放的难点:同一个逻辑子图在循环里会跑多次,重放时该给它哪一份存档?

# _replay.py:34
    def _is_first_visit(self, checkpoint_ns: str) -> bool:
        # "sub_node:task_id" -> "sub_node"  剥掉 task_id 后缀
        stable_ns = (checkpoint_ns.rsplit(NS_END, 1)[0]
                     if NS_END in checkpoint_ns else checkpoint_ns)
        if stable_ns in self._visited_ns:
            return False                          # 已访问过 → 不是第一次
        self._visited_ns.add(stable_ns)
        return True                               # 第一次 → 记下
rsplit(NS_END,1)[0]"sub:task123" 里的 task_id 后缀剥掉,还原成稳定的逻辑名 "sub"——这样循环里每轮虽 task_id 不同,仍认得出"是同一个子图"。
_visited_ns 集合记"哪些子图已经做过一次性回溯了"。第一次访问 → 加进集合、返回 True;再访问 → 已在集合、返回 False。
第一次取旧档、之后取新档配合 get_checkpoint(D44 L07):第一次回到"重放点之前"的档;之后用正常最新档。保证循环里子图只回溯一次,不会每轮都拽回远古。
⚠️ 边界:幂等只保护"框架管理的通道写",管不了你的外部副作用贯穿 D42、D44、D46 的铁律,这里正式点破:LangGraph 的幂等作用于通道写(状态更新)——它能保证 A 的 {"count": +1} 不被应用两遍。但你在节点里手写的外部副作用——发邮件、调支付、写外部数据库、发 MQ 消息——框架看不见也管不着。恢复重跑一个"interrupt 之前发过邮件"的节点,邮件会再发一遍。正确姿势:① 把副作用放在 interrupt 之后(interrupt 前的代码才会重放);② 或副作用自带幂等键(如支付用唯一 idempotency-key);③ 或用 @task(D47-48,其结果会被 checkpoint 缓存,重放时不重执行)把副作用包起来。别指望框架替你的副作用兜底。
L07

幂等心法 + 阶段 7 收官

💡 一句话总结幂等机制崩溃/中断/时间旅行都要"重放接续点那一步"。框架靠三件事保证"至少一次执行、效果等价恰好一次":① 已完成任务的普通写被复原(writes 非空 → 不重跑,L03);② 中断/失败任务的控制信号不复原(writes 空 → 重跑,L04);③ put_writes 存时清旧存新去重(不累积,L05)。子图另有 ReplayState 管"回溯一次"(L06)。
数据结构:pending_writes 里两类记录的命运 普通通道写 (task_id, "count", 5) = 已成功的成果 → 复原到任务 writes 非空 → 不重跑 控制信号 (task_id, INTERRUPT, ..) (task_id, ERROR, ..) (NULL, RESUME, ..) → 跳过不复原 writes 空 → 重跑
图注:同在 pending_writes,普通写=复原跳过,控制信号=跳过重跑。这就是 _reapply 那句 continue 的全部意义。

👶 小白:整个阶段 7 讲的中断、恢复、时间旅行、durability、幂等,能不能串成一句话?

👨‍🏫 能:因为每一步都可以落盘成 checkpoint(D45 决定多勤),图就能在任意点暂停(D41 interrupt / D43 静态断点)、把答案喂回去继续(D42 Command)、翻看或修改任意历史点(D44 时间旅行),而这一切"停了再续、回到过去再跑"之所以不会把状态搞乱,是因为重放时严格区分"已完成的复用、未完成的重跑"(D46 幂等)。一句话:持久化是地基,中断/时间旅行是能力,幂等是让这些能力安全的保险丝。

🧠 今天你应该能回答

  • 恢复时为什么已完成的节点不会重跑?(其普通写被复原、writes 非空,runner 只跑空 writes 的任务)
  • is_replaying 什么时候为真、持续多久?(config 带 checkpoint_id 时为真,只第一个 tick)
  • _reapply 为什么跳过 ERROR/INTERRUPT/RESUME?(让失败/中断任务保持空 writes 从而重跑)
  • 为什么"writes 非空"就能表示"已完成"?(写是完成的唯一真相,避免双标记不一致)
  • put_writes 如何防止写累积?(普通任务存新写前先删该 task_id 的旧写)
  • 幂等保护不了什么?(节点里手写的外部副作用——发邮件、扣款,需自行幂等或放 interrupt 后)

✋ 10 分钟动手

# 1. 读幂等三处核心
sed -n '661,664p' libs/langgraph/langgraph/pregel/_loop.py     # 复原触发
sed -n '736,749p' libs/langgraph/langgraph/pregel/_loop.py     # 复原+跳过控制信号
sed -n '415,437p' libs/langgraph/langgraph/pregel/_loop.py     # put_writes 去重
# 2. 验证"中断恢复不重复副作用"——反例(它会重复!)
python - <<'PY'
from langgraph.graph import StateGraph, START
from langgraph.checkpoint.memory import InMemorySaver
from langgraph.types import interrupt, Command
from typing_extensions import TypedDict
hits=[]
class S(TypedDict):
    x: int
def node(s):
    hits.append("side-effect")   # 假装这是"发邮件",放在 interrupt 之前
    ans = interrupt("确认?")
    return {"x": 1}
g=StateGraph(S); g.add_node("n",node); g.add_edge(START,"n")
app=g.compile(checkpointer=InMemorySaver())
cfg={"configurable":{"thread_id":"t"}}
list(app.stream({"x":0}, cfg))            # 停,副作用已发生 1 次
list(app.stream(Command(resume="ok"), cfg))  # 恢复→节点重跑,副作用又发生
print("副作用次数:", len(hits))            # 2 !! 证明 interrupt 前的副作用会重放
PY
阶段 7 完结 · 明天进入阶段 8:Day 47 讲函数式 API——用 @entrypoint / @task 装饰器写图,不用显式建 StateGraph。而上面动手题里"副作用重放"的痛点,正是 @task 要解决的(task 结果会被缓存,重放时不重执行)。承上启下,敬请期待。
← Day 45 durability Day 47 · @entrypoint / @task →