durability:写得多勤、写得多稳,你自己拍板
中断、恢复、时间旅行全靠"每步落盘"。但落盘不是免费的——每个超步都往数据库写一次,慢且贵。durability 给你三档旋钮:"sync"(每步同步写完再走,最稳最慢)、"async"(边写边跑,默认,平衡)、"exit"(只在图结束时写一次,最快最省但中途崩了没档)。今天看这三档在主循环里到底改变了什么。
痛点:每步都写数据库,慢在哪、值不值
•
"sync":每超步同步写——写完落盘确认了才开始下一步。最稳,崩溃只会丢"正在跑的那一步"。•
"async"(默认):每超步异步写——一边把上一步往盘上写、一边已经跑下一步。近乎每步都存,但把写的耗时和计算重叠了。•
"exit":只在图退出(结束/中断/报错)时写一次。跑得最快,但中途崩溃=从头再来(没中间档)。类比:写文档的自动保存。sync=每敲一个字存一次(绝不丢,但卡);async=后台每几秒悄悄存(几乎不丢、不卡);exit=只在关闭时存(最快,但没保存就崩了全没)。
三档的类型定义与官方说明
类型定义在 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。默认 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" 都行——按业务重要性临场调。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 说"要,每步存"。这一层由
do_checkpoint 决定(本讲)。② 存的时候要不要等它写完?——sync 说"等",async 说"不等,边写边跑"。这一层由 L06 的 future 等待决定。
看懂这两层,三档就不再是三个孤立的魔法词,而是 2×2 里有意义的三个组合。
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 写全落盘。退出前的最后一份完整现场。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 的区别在主循环等不等这个 fut:async——提交完就往下跑下一超步,写在后台并行进行(快);sync——提交后立即等 fut 完成、确认落盘了才进下一步(稳)。prev.result() 等它先完成,强制 checkpointer 严格按 step1→step2→step3 收档。代价是写与写之间被串行化(不能任意乱序并行),但换来"存档链永远连续、恢复永远读到完整前缀"。注意这不影响 async 的性能优势——串行的是"写与写",而"写"仍与"下一步计算"并行,计算才是耗时大头。prev.result() 保序,checkpoint 落盘顺序和 sync 完全一致。唯一差别是:async 下,"第 N 步写盘"可能和"第 N+1 步计算"同时进行,若正好在这个重叠窗口崩溃,可能丢的是第 N 步的档(而 sync 会确保第 N 步落定才算第 N+1 步)。但这个窗口极短,且节点计算(调模型)通常远比写盘慢,重叠几乎总能让写先完成。于是 async 用"一个极窄的丢档窗口"换来"写盘耗时被计算完全掩盖"的吞吐。只有当你要求"绝对不丢任何一步、宁可慢"(如金融强一致场景)才值得上 sync。这是"默认给 95% 场景的最优解、极端场景留后门"的务实取舍。选型指南 + 今日小结
用 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
_replay.py 的重放控制、以及 put_writes 的幂等去重如何做到"至少执行一次,效果等价恰好一次"。