replay 与幂等:重跑一遍,为什么结果不会乱
阶段 7 收官。前面反复出现一个词——重跑:interrupt 恢复要重跑节点、时间旅行要从旧档重放、崩溃恢复要重放已完成的超步。重跑的最大风险是做重复的事:同一个写被应用两遍、reducer 累加两次、状态翻倍。今天看 LangGraph 如何用一套"重放 + 幂等"机制,做到"至少执行一次,但效果等价于恰好一次"。
痛点:恢复时前面的节点会不会又跑一遍?
count += 1(reducer 累加),重跑就变成加两次,状态就错了!类比:搬家搬到一半停电。来电后你不会把已经搬上车的箱子再搬一遍——你看一眼清单(pending_writes),已装车的打勾跳过,只搬还在地上的。清单就是幂等的钥匙。
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"。after_tick 里第一步一结束就 is_replaying=False,把特殊模式的作用域精确限制在唯一需要它的那一步。这是"特殊逻辑作用域最小化"的典范——能用普通路径的地方绝不让特殊标志多停留一个 tick。复用已成功的写:不重跑,直接认账
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 因此不重跑。done: bool。问题:done 和 writes 是两份信息,得同步维护,还可能不一致(标了 done 但 writes 没存)。源码做法:"有没有写"本身就是"完成没完成"的唯一真相——任务完成必然产出写(哪怕是 NO_WRITES 哨兵,见 D41 runner),写落了盘就是"干完了"的铁证。恢复时把写复原 → 任务自动"变成已完成态"。用同一份数据(writes)同时表达"产出了什么"和"完成没有",消除了双写不一致的可能。runner 的判据 not t.writes 简单到无懈可击。为什么跳过 ERROR/INTERRUPT/RESUME?
上一讲那句 if k in (ERROR, ERROR_SOURCE_NODE, INTERRUPT, RESUME): continue 是幂等的灵魂。函数注释(_loop.py:739-743)解释了动机:
普通写 → 复原 → 跳过重跑成功完成的任务,产出的是普通通道写。复原它 = 认账已完成 = 不重跑。(L03)INTERRUPT 写 → 不复原被中断的任务,存档里是 (INTERRUPT, ...)(D41)。故意不复原 → 该任务 writes 仍为空 → runner 重新执行它。这正是 D42 说的"恢复时节点从头重跑"!中断任务必须重跑,才能让 interrupt() 这次拿到恢复值、走完后续逻辑。ERROR 写 → 不复原失败的任务,存的是 (ERROR, 异常)。不复原 → 空 writes → 重跑(给它再来一次的机会,或交给错误处理器)。RESUME 写 → 不复原恢复值是要被 interrupt() 重新消费的(D42 的 get_null_resume),不能当"已完成的普通产出"复原掉。if...continue 把这两类干净利落地分开。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 值。用不同策略。put_writes 还有一层"普通写 (task_id, idx) 已存在则跳过、负数特殊写允许覆盖"的幂等(checkpoint/memory/__init__.py:494)。两层幂等:loop 层 put_writes 管"内存 pending 的新旧替换",checkpointer 层 put_writes 管"落盘时按 (task,idx) 去重"。双保险。子图的重放: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):第一次回到"重放点之前"的档;之后用正常最新档。保证循环里子图只回溯一次,不会每轮都拽回远古。{"count": +1} 不被应用两遍。但你在节点里手写的外部副作用——发邮件、调支付、写外部数据库、发 MQ 消息——框架看不见也管不着。恢复重跑一个"interrupt 之前发过邮件"的节点,邮件会再发一遍。正确姿势:① 把副作用放在 interrupt 之后(interrupt 前的代码才会重放);② 或副作用自带幂等键(如支付用唯一 idempotency-key);③ 或用 @task(D47-48,其结果会被 checkpoint 缓存,重放时不重执行)把副作用包起来。别指望框架替你的副作用兜底。幂等心法 + 阶段 7 收官
👶 小白:整个阶段 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
@entrypoint / @task 装饰器写图,不用显式建 StateGraph。而上面动手题里"副作用重放"的痛点,正是 @task 要解决的(task 结果会被缓存,重放时不重执行)。承上启下,敬请期待。