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

durability:写得多勤、写得多稳,你自己拍板

中断、恢复、时间旅行全靠"每步落盘"。但落盘不是免费的——每个超步都往数据库写一次,慢且贵。durability 给你三档旋钮:"sync"(每步同步写完再走,最稳最慢)、"async"(边写边跑,默认,平衡)、"exit"(只在图结束时写一次,最快最省但中途崩了没档)。今天看这三档在主循环里到底改变了什么。

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

痛点:每步都写数据库,慢在哪、值不值

🤔 痛点我的图有 20 个超步,用 PostgresSaver。开着 checkpoint 后明显变慢——每个超步结束都要往 Postgres 写一次(序列化、网络往返、事务提交)。可我这个图不需要中断、也不做时间旅行,只是想要"万一崩了能恢复"。为了这点保险,20 次数据库写值得吗?能不能少写点?
💡 本质:持久化频率是一个可调的"安全 vs 性能"旋钮落盘越勤 → 崩溃时丢的进度越少(安全),但写得越多(慢)。落盘越懒 → 越快,但崩了丢得越多。LangGraph 把这个权衡交给你,三档:
"sync":每超步同步写——写完落盘确认了才开始下一步。最稳,崩溃只会丢"正在跑的那一步"。
"async"(默认):每超步异步写——一边把上一步往盘上写、一边已经跑下一步。近乎每步都存,但把写的耗时和计算重叠了。
"exit"只在图退出(结束/中断/报错)时写一次。跑得最快,但中途崩溃=从头再来(没中间档)。
类比:写文档的自动保存。sync=每敲一个字存一次(绝不丢,但卡);async=后台每几秒悄悄存(几乎不丢、不卡);exit=只在关闭时存(最快,但没保存就崩了全没)。
L02

三档的类型定义与官方说明

类型定义在 types.py:87-88,就一行字面量:

# types.py:87
Durability = Literal["sync", "async", "exit"]
"""Durability mode for the graph execution."""

官方对三档的一句话说明(stream 文档 main.py:2705-2711):

# main.py:2705
    durability: The durability mode for the graph execution, defaults to "async".
        Options are:
        - "sync":  Changes are persisted synchronously before the next step starts.
        - "async": Changes are persisted asynchronously while the next step executes.
        - "exit":  Changes are persisted only when the graph exits.
档位什么时候写崩溃能恢复到速度
sync每超步,同步(写完才继续)上一个完成的超步最慢
async(默认)每超步,异步(与下一步并行)上一个完成的超步(几乎)
exit只在图退出时一次什么都没有→从头最快
还有个历史遗留参数 checkpoint_during(布尔)已废弃(main.py:2727-2735):True 映射到 "async"False 映射到 "exit"。新代码统一用 durability
L03

默认 async,且可逐次运行覆盖

默认值和读取逻辑在 main.py:2602-2603

# main.py:2602
        if durability is None:
            durability = config.get(CONF, {}).get(CONFIG_KEY_DURABILITY, "async")

没有 checkpointer 时给 durability 会警告(main.py:2802-2804):

# main.py:2802
            if checkpointer is None and durability is not None:
                warn("`durability` has no effect when no checkpointer is present.", ...)
durability is None → "async"你不传就是 async。这是"既要几乎每步可恢复、又不想太慢"的合理默认。
从 config 读也能通过 config 的 CONFIG_KEY_DURABILITY 设——方便统一配置或子图继承。
无 checkpointer 就警告durability 只管"往 checkpointer 写多勤"。没 checkpointer 就无处可写,此时设 durability 毫无意义,框架好心提醒你。
逐次可传同一个编译好的 app,这次 stream(..., durability="sync")、下次 "exit" 都行——按业务重要性临场调。
L04

do_checkpoint:这一档到底存不存

每次要落盘时,主循环先算一个布尔 do_checkpoint_loop.py:1132-1135),它是三档差异的总闸门

# _loop.py:1132
        # do checkpoint?
        do_checkpoint = self._checkpointer_put_after_previous is not None and (
            exiting or self.durability != "exit"
        )
_checkpointer_put_after_previous is not None前提:得有 checkpointer(这个函数是它的写入句柄)。没有就永远不存。
exiting如果正在退出(图跑完/中断/报错)——不管哪一档都要存。这保证 exit 模式"退出时那一次"能落盘。
or self.durability != "exit"核心:只要不是 exit 档(即 sync 或 async),每个超步都 do_checkpoint=True,正常落盘。
合起来翻译:exit 档 → 只有退出时才存;sync/async 档 → 每步都存。一个布尔表达式区分了"存不存"这一层。
💡 本质:exit 与 sync/async 的分界是"存不存",sync 与 async 的分界是"等不等"durability 其实是两个正交问题合成的:
要不要每步存?——exit 说"不,只退出存";sync/async 说"要,每步存"。这一层由 do_checkpoint 决定(本讲)。
存的时候要不要等它写完?——sync 说"等",async 说"不等,边写边跑"。这一层由 L06 的 future 等待决定。
看懂这两层,三档就不再是三个孤立的魔法词,而是 2×2 里有意义的三个组合。
数据结构:三档 = 两个正交问题的组合 纵轴:每步都存吗? 横轴:存时要等写完吗? 每步存 只退出存 不等(异步) 等(同步) async(默认) 每步存 + 不等 sync 每步存 + 等 exit 不每步存("等不等"无意义 → 占整行)
图注:exit 在"不每步存"这一行(等不等不适用);async/sync 都每步存,只差"等不等落盘"。
L05

exit 模式:把落盘攒到退出那一刻

exit 档下,中间超步都不存,直到图退出时在"消音器"里统一落盘(_loop.py:1324-1334,还记得吗,这就是 D41 见过的 _suppress_interrupt):

# _loop.py:1324
        if self.durability == "exit" and (
            not self.is_nested                               # ① 顶层图
            or exc_value is not None                          # ② 或子图出错/中断
            or all(NS_END not in part for part in self.checkpoint_ns)  # ③ 或独立 checkpointer 子图
        ):
            self._put_exit_delta_writes()                     # 补落"增量通道"的写
            self._put_checkpoint(self.checkpoint_metadata)    # 落最终 checkpoint
            self._put_pending_writes()                        # 落待写
durability == "exit"只有 exit 档走这里——因为 sync/async 中间步已经存过了,退出时不必特意补。
__exit__ 时机这段在主循环上下文管理器的 __exit__ 里。图无论正常结束、还是中断/报错退出,都会经过这,保证 exit 档"退出时那一次"必然发生。
三种触发条件顶层图退出必存;子图带错误/中断退出也存(否则错误现场丢了);带独立 checkpointer 的子图也存。覆盖了各种退出路径。
三个 _put_*一次性把最终 checkpoint + delta 写 + pending 写全落盘。退出前的最后一份完整现场。
⚠️ 边界:exit 档下中途崩溃 = 白跑,且中断/时间旅行基本失效exit 只在退出时存。如果图跑到第 10 步进程被 kill(OOM、断电、部署重启),前 9 步一份档都没有——恢复只能从头再来。更关键:既然中间不落盘,静态断点、时间旅行的"逐步历史"在 exit 档下就没有了(get_state_history 看不到中间步)。所以 exit 适合:短、快、不需要中途恢复/观察的一次性任务(如一个无中断的批处理)。只要你的图用到 interrupt/断点/时间旅行,就别用 exit。注意 interrupt 触发退出时那次会存(因为 exiting 为真),所以"用 exit + interrupt"能停一次,但停之前的中间步历史是缺的。
L06

sync vs async:差别在"等不等上一次写完"

落盘不是直接调 put,而是提交到线程池,串成一条链(_loop.py:1200-1209):

# _loop.py:1200
            # if there's a previous checkpoint save in progress, wait for it
            # ensuring checkpointers receive checkpoints in order
            self._put_checkpoint_fut = self.submit(
                self._checkpointer_put_after_previous,
                getattr(self, "_put_checkpoint_fut", None),   # ← 上一次写的 future
                self.checkpoint_config,
                copy_checkpoint(self.checkpoint),
                self.checkpoint_metadata,
                new_versions,
            )

_checkpointer_put_after_previous_loop.py:1530-1547)保证"按顺序写":

# _loop.py:1530
    def _checkpointer_put_after_previous(self, prev, config, checkpoint, metadata, new_versions):
        if self._delta_write_futs:
            futs, self._delta_write_futs = self._delta_write_futs, []
            concurrent.futures.wait(futs)
        try:
            if prev is not None:
                prev.result()                                  # ① 等上一次写完
        finally:
            cast(BaseCheckpointSaver, self.checkpointer).put(  # ② 再写这一次
                config, checkpoint, metadata, new_versions)
submit(...)把"写这份 checkpoint"提交给执行器,拿到一个 future(写操作的凭据),存进 _put_checkpoint_fut
把上次的 fut 传进去每次写都携带上一次写的 future。函数内部 prev.result() 会阻塞到上一次写完——保证 checkpointer 收到的顺序严格是 step1, step2, step3……不乱序。
prev.result()等上一次真正落盘完成再执行本次 put。这是"有序写"的关键,无论 sync/async 都要保序。
sync 与 async 的区别主循环等不等这个 futasync——提交完就往下跑下一超步,写在后台并行进行(快);sync——提交后立即等 fut 完成、确认落盘了才进下一步(稳)。
💡 设计取舍①:为什么落盘要串成 future 链保序,而不是各写各的?看似"每步各自 submit 一个写、谁先写完算谁"更简单、更并行。但那会毁掉恢复的正确性:如果 step3 的写比 step2 先落盘(线程调度谁快说不准),此刻崩溃,checkpointer 里就有 step3 却缺 step2——恢复读到一个"跳步"的残缺历史,父子链断裂。所以源码让每次写携带上一次写的 future 并 prev.result() 等它先完成,强制 checkpointer 严格按 step1→step2→step3 收档。代价是写与写之间被串行化(不能任意乱序并行),但换来"存档链永远连续、恢复永远读到完整前缀"。注意这不影响 async 的性能优势——串行的是"写与写",而"写"仍与"下一步计算"并行,计算才是耗时大头。
💡 设计取舍②:为什么默认是 async 而不是最稳的 sync?因为 async 在绝大多数场景下"安全性几乎等于 sync,但快得多"。关键洞察:async 也是每步都写的,只是"写"和"下一步计算"并行——由于有 prev.result() 保序,checkpoint 落盘顺序和 sync 完全一致。唯一差别是:async 下,"第 N 步写盘"可能和"第 N+1 步计算"同时进行,若正好在这个重叠窗口崩溃,可能丢的是第 N 步的档(而 sync 会确保第 N 步落定才算第 N+1 步)。但这个窗口极短,且节点计算(调模型)通常远比写盘慢,重叠几乎总能让写先完成。于是 async 用"一个极窄的丢档窗口"换来"写盘耗时被计算完全掩盖"的吞吐。只有当你要求"绝对不丢任何一步、宁可慢"(如金融强一致场景)才值得上 sync。这是"默认给 95% 场景的最优解、极端场景留后门"的务实取舍。
L07

选型指南 + 今日小结

控制流:三档在超步时间线上的落盘节奏 sync step1 step2 计算与存交替,存完才下一步 async step1 step2 step3 存1 存2 存在后台并行(与下一步重叠) exit step1 step2 step3 全程不存,退出才存一次 时间 →
图注:sync 存挡在计算之间;async 存与下一步并行;exit 全程不存、末尾一次。
🍼 怎么选(大白话) 用 async(默认):绝大多数带 checkpointer 的应用——要中断、要恢复、要时间旅行,又不想太慢。
用 sync:状态极其宝贵、绝不允许丢任何一步,宁慢求稳(如涉及资金、法务的强一致流程)。
用 exit:短平快的一次性任务,不需要中途恢复/中断/历史,只想跑得最快(如离线批量处理)。

👶 小白:async 是异步写,那我用同步的 invoke() 也能用 async 档吗?会不会要 asyncio?

👨‍🏫 能用,不需要你写 asyncio。这里的"async"指的是写盘操作被提交到线程池后台执行submit 到 executor),主循环不阻塞地继续算下一步——靠的是线程而非 asyncio 协程。所以同步 invoke() 也享受 async 档的并行写盘。真正的异步图(ainvoke)则用异步版 _checkpointer_put_after_previous_loop.py:1783),道理一样,只是等待用 await

🧠 今天你应该能回答

  • durability 三档分别什么时候写?(sync 每步同步 / async 每步异步 / exit 只退出时)
  • 默认是哪档、为什么?(async;每步都存但写盘与计算重叠,安全近似 sync 却快)
  • do_checkpoint 这个布尔区分了什么?(exit vs 非exit:中间步存不存)
  • sync 和 async 的真正区别在哪?(主循环等不等落盘 future:sync 等、async 不等)
  • 为什么多次落盘要传上一次的 future?(prev.result() 保证 checkpointer 按超步顺序收到)
  • exit 档为什么不能配中断/时间旅行?(中间步没落盘,历史缺失、崩溃从头)

✋ 10 分钟动手

# 1. 读三档的分岔点
sed -n '2602,2603p' libs/langgraph/langgraph/pregel/main.py    # 默认 async
sed -n '1132,1135p' libs/langgraph/langgraph/pregel/_loop.py   # do_checkpoint 闸门
sed -n '1530,1547p' libs/langgraph/langgraph/pregel/_loop.py   # 有序写 / prev.result
# 2. 对比三档产生的历史档数
python - <<'PY'
from langgraph.graph import StateGraph, START, END
from langgraph.checkpoint.memory import InMemorySaver
from typing_extensions import TypedDict
class S(TypedDict):
    x: int
def build():
    g=StateGraph(S)
    g.add_node("a",lambda s:{"x":s["x"]+1}); g.add_node("b",lambda s:{"x":s["x"]+1})
    g.add_edge(START,"a"); g.add_edge("a","b"); g.add_edge("b",END)
    return g.compile(checkpointer=InMemorySaver())
for mode in ["sync","async","exit"]:
    app=build(); cfg={"configurable":{"thread_id":"t"}}
    app.invoke({"x":0}, cfg, durability=mode)
    print(mode, "历史档数:", len(list(app.get_state_history(cfg))))
    # 你会看到 exit 档的历史档数明显更少
PY
明天预告 · Day 46:本阶段收官。中断/时间旅行都涉及"重跑"——从存档恢复会重新执行已跑过的部分。那如何保证"重跑一遍结果不出错、副作用不翻倍"?看 _replay.py 的重放控制、以及 put_writes 的幂等去重如何做到"至少执行一次,效果等价恰好一次"。
← Day 44 时间旅行 Day 46 · replay 与幂等 →