学习笔记 · Obsidian
持久化、记忆、时间旅行与容错
/durable-execution当前规范跳转到/persistence;耐久执行依赖的 checkpoint、thread、replay 与副作用边界均以本页持久化语义为准。
核心结论
LangGraph 的持久化分为两个互补系统:
| 系统 | 保存内容 | 范围 | 典型用途 |
|---|---|---|---|
| Checkpointer | 完整 graph state 快照与节点 writes | 单个 thread | 多轮对话、HITL、失败恢复、time travel |
| Store | 应用自定义 key-value 数据 | 跨 thread | 用户偏好、事实、共享知识、长期记忆 |
大多数生产应用两者都需要:checkpointer 记住“这条会话执行到哪里”,Store 记住“这个用户或业务长期知道什么”。
Agent Server 会管理持久化基础设施;直接运行 OSS 图时,需要自行选择、配置和迁移 checkpointer 与 Store。
Thread 与 Checkpoint
thread_id 是持久化主键
图配置 checkpointer 后,每次 invoke、stream 或恢复都必须在 config.configurable 中传入 thread_id。它决定:
- 本次输入写到哪个会话;
- 从哪个最新 checkpoint 加载;
- interrupt 恢复哪个暂停点;
- 后续多轮调用共享哪组短期记忆。
新 thread_id 表示新会话。重复使用旧 thread_id 表示继续同一持久状态,不等于“重新跑一个全新请求”。
生产 thread_id 应:
- 长度小于持久化后端列限制,PostgresSaver 建议不超过 255 字符;
- 使用 UUID 或稳定 hash,而不是拼接超长业务文本;
- 做租户隔离,不允许用户任意猜测其他 thread;
- 不包含敏感信息。
checkpoint 是 super-step 快照
LangGraph 在每个 super-step 边界保存 StateSnapshot。对 START → A → B → END 的顺序图,通常会得到:
- 空状态,下一步 START;
- 应用输入,下一步 A;
- A 的输出,下一步 B;
- B 的输出,图完成。
并行节点属于同一个 super-step。time travel 只能从完整 checkpoint 边界开始,不能从普通节点函数内部任意代码行恢复。
StateSnapshot 字段
| 字段 | 含义 |
|---|---|
| values | 当前所有 state channel 值 |
| next | 下一批将执行的节点;空元组表示完成 |
| config | thread_id、checkpoint_ns、checkpoint_id |
| metadata | source、writes、step 等执行信息 |
| created_at | checkpoint 时间 |
| parent_config | 前一 checkpoint 配置 |
| tasks | 待执行 task、错误、interrupt 与可选 subgraph state |
checkpoint_ns 标识父图或子图。根图通常为空字符串;子图 namespace 包含节点名与运行 ID,嵌套层级用分隔符连接。
pending writes
并行 super-step 中,如果一个节点失败而其他节点已成功:
- 本轮完整 StateSnapshot 不提交;
- 成功节点的 task writes 已写入 checkpoint_writes;
- 恢复时成功节点不重复执行,只重试失败节点。
这既减少成本,也意味着自定义 checkpointer 不能只存完整 snapshot,必须正确保存节点级 writes。
状态读取、历史与更新
get_state
读取指定 thread 的最新 StateSnapshot;提供 checkpoint_id 时读取历史指定点。
get_state_history
返回该 thread 的 checkpoint 历史,默认最新在前。可以按 step、source、writes、interrupt 等条件定位目标 checkpoint。
update_state
update_state 不会修改原 checkpoint,而是创建一个新 checkpoint 分支。更新仍通过 reducer:
- 有 reducer 的字段会累积;
- 没有 reducer 的字段会覆盖;
- as_node 决定这次更新视为哪个节点产生,并影响下一步去向。
并行分支无法自动判断最后节点、全新 thread 没有历史,或测试中要跳过某节点时,应显式传 as_node。
Durability mode
| 模式 | 持久化时机 | 优点 | 风险 |
|---|---|---|---|
| exit | 图结束、报错或 interrupt 时 | 最高吞吐、最少写入 | 进程中途崩溃会丢失中间进度 |
| async | 下一步执行时异步写上一步 | 性能与持久性平衡,默认思路 | 极端崩溃窗口可能少一个 checkpoint |
| sync | 写成功后才进入下一步 | 每步最强持久性 | 增加延迟与存储压力 |
支付、外部不可重放动作或严格审计链路需要评估 sync;可安全重算的推理工作可以使用 async 或 exit。不能把持久化模式当作副作用幂等的替代品。
Checkpointer 选型
| 实现 | 场景 |
|---|---|
| InMemorySaver / MemorySaver | 单元测试、临时实验;进程重启即丢失 |
| SqliteSaver / AsyncSqliteSaver | 单机本地开发 |
| PostgresSaver / AsyncPostgresSaver | 常规生产环境 |
| CosmosDB saver | Azure 生产环境 |
| 其他集成 | 依据官方 checkpointer integrations |
数据库实现首次使用通常要运行 setup 或等价 migration。应把 migration 作为部署步骤,不要每个请求动态建表。
Checkpointer 接口
BaseCheckpointSaver 的核心方法:
- put / aput:保存完整 checkpoint;
- put_writes / aput_writes:保存节点级 pending writes;
- get_tuple / aget_tuple:按 thread、namespace 与可选 checkpoint_id 读取;
- list / alist:按条件列出历史;
- delete_thread / adelete_thread:删除 thread 全部 checkpoint 与 writes。
异步图会调用 async 版本。自定义后端若 async 方法只是阻塞同步 I/O 的薄包装,会拖住事件循环。
自定义后端的数据模型
至少需要:
- checkpoints 表:每个 super-step 一行;
- writes 表:每个节点输出一行或一组行;
- 直接按 thread_id + checkpoint_ns + checkpoint_id 查找;
- 按 checkpoint_id 倒序取最新;
- checkpoint 到 parent checkpoint 的链;
- 同时删除 checkpoint 与 writes。
checkpoint_id 是可排序 ULID。按 ID 直接查找必须接近 O(1),不能每次扫描整个 thread 历史。time travel 和 DeltaChannel 重建都依赖这一点。
Serializer 与加密
默认 JsonPlusSerializer 基于 ormsgpack 与 JSON,支持 LangGraph、LangChain、datetime、enum、dataclass、Pydantic 等常见类型。无法编码的对象可启用 pickle fallback,但 pickle 增加安全与跨版本风险,不应对不可信数据反序列化。
EncryptedSerializer 可加密持久 State。加密密钥应由部署环境注入,不能写入仓库或笔记。加密只保护静态存储,仍需控制日志、trace、stream 与 API 返回。
DeltaChannel 与 checkpoint 体积
默认每个 checkpoint 会序列化所有 channel 的完整值。长对话 messages 持续增长时,存储量可能随 thread 长度快速膨胀。
DeltaChannel 只保存每步增量,读取时沿 parent checkpoint 回放 writes:
- checkpoint blob 可从 O(N) 降到每步近似 O(1);
- snapshot_frequency 每 K 步写完整快照,限制重建深度;
- reducer 必须可批量、结合且纯函数;
- prune 时不能删除存活 checkpoint 仍依赖的祖先 writes;
- copy_thread 必须复制完整依赖链或先生成快照;
- 老版本无法读取新 delta 格式,回滚必须先迁移。
自定义 checkpointer 应运行 langgraph-checkpoint-conformance,并覆盖 delta history、prune、copy 与指定 checkpoint_id 查找。
Store 与长期记忆
namespace 设计
Store 以 tuple namespace + key 组织数据。常见形式:
- tenant_id / user_id / memories;
- tenant_id / project_id / facts;
- tenant_id / shared / policies。
namespace 是权限边界的一部分。只用 user_id 而不带 tenant_id,容易在多租户环境产生跨租户读取。
基础操作
- put:保存或覆盖 item;
- get:按精确 namespace 与 key 读取;
- delete:删除;
- search:按 namespace prefix 查询;
- list_namespaces:发现 namespace。
Item 包含 value、key、namespace、created_at、updated_at。
search 的三个易错点
- namespace_prefix 是前缀匹配,不是精确匹配;
- 超过 limit 的数据会静默截断,应分页;
- 不同后端默认排序不同,顺序重要时按 updated_at 在应用层排序。
语义搜索
Store 可配置 embedding,把指定字段写入索引,并用自然语言 query 按语义相似度检索。不是每个值都必须 embed;敏感或纯元数据字段可关闭 index。
语义搜索仍需:
- namespace 与租户过滤先于向量相似度;
- 写入和删除时同步维护向量索引;
- 记录 embedding 模型版本;
- 处理模型升级后的重建;
- 对无向量能力的自定义 Store 明确抛出 NotImplementedError。
生产 Store
InMemoryStore 只适合开发。生产可用 Postgres、MongoDB、Redis、Upstash、Oracle 等集成,具体以 Store integrations 为准。
BaseStore 自定义实现至少需要五个 async 方法:aput、aget、adelete、asearch、alist_namespaces。值应是 JSON 可序列化 dict,不存任意 Python 对象。
Short-term 与 Long-term memory
Short-term memory
短期记忆就是 thread State + checkpointer。多轮调用使用相同 thread_id,图自动加载之前消息与状态。
长对话会超过模型窗口,常见治理方式:
- trim:调用模型前按 token 数保留最近消息;
- delete:用 RemoveMessage 从 State 永久删除;
- summarize:把旧消息压缩到 summary,再删除原消息;
- checkpoint retention:定期清理旧历史;
- 业务过滤:只保留任务相关消息。
删除消息时要保持合法对话结构:
- 某些模型要求以 user 消息开始;
- AI tool call 后必须保留对应 ToolMessage;
- 使用 RemoveMessage 要求 messages 字段使用 add_messages reducer。
总结策略应保存原始 summary 数据并按需生成 prompt,避免把格式化提示词写回 State。
Long-term memory
长期记忆通过 Store 跨 thread 保存。图 compile 时传 store,节点通过 Runtime.store 读写;user_id 等 namespace 信息通过 runtime context 传入。
同一个 user_id 即使换 thread 仍能访问长期记忆。这里必须由认证身份映射 user_id,不能信任客户端任意传入。
Subgraph 中的持久化
父图配置 checkpointer 后,默认会向子图传播。子图也可以明确选择:
- 每次调用隔离但继承父 checkpointer;
- checkpointer=True,跨同一 thread 多次调用持续保存;
- checkpointer=False,完全无持久化。
具体并发与 namespace 风险见 04-流式HITL与子图。
Time travel
Replay
从历史 checkpoint 的 config 再次 invoke:
- checkpoint 之前的节点不执行;
- checkpoint 之后的节点重新执行;
- LLM、API、工具与 interrupt 都可能再次触发;
- 从最终、next 为空的 checkpoint replay 是 no-op。
replay 不是结果缓存回放 后续节点是真实重执行。任何外部写操作都需要幂等,否则调试一次可能重复发邮件、重复扣费或重复建单。
Fork
在历史 checkpoint 上调用 update_state 会创建新分支:
- 选择历史 checkpoint;
- 更新一个或多个 State 字段;
- 必要时指定 as_node;
- 用新返回的 config 继续 invoke。
原历史保留,不发生“回滚覆盖”。Fork 适合:
- 比较不同人工答案;
- 修正模型输出后继续;
- 调试另一条路由;
- 从中间状态构造测试。
Interrupt 与 time travel
time travel 经过 interrupt 时,interrupt 会重新触发并等待新的 Command.resume。多个 interrupt 可以从两次暂停之间 fork,从而保留前一个答案并重答后一个。
Subgraph 粒度
- 默认继承父 checkpointer 的子图,在父图看来是一个 super-step,只能从子图前后 travel;
- checkpointer=True 的子图有内部 checkpoint 历史,可以从子图内部节点 fork;
- 使用 get_state(..., subgraphs=True) 获取内部 config。
容错生命周期
节点失败后的固定顺序是:
- 当前 attempt 运行;
- timeout 或其他异常产生;
- retry_policy 判断是否重试;
- 重试耗尽后才进入 error_handler;
- handler 可更新 State,并用 Command 路由补偿分支。
interrupt 不属于普通错误,使用 GraphBubbleUp 暂停,不进入 retry 或 error_handler。
RetryPolicy
默认:
- max_attempts 为 3,包含首次执行;
- 初始间隔 0.5 秒;
- backoff_factor 为 2;
- max_interval 为 128 秒;
- 默认启用 jitter。
默认不重试常见程序错误,例如 ValueError、TypeError、ArithmeticError、ImportError、LookupError、NameError、SyntaxError、RuntimeError、ReferenceError、StopIteration、OSError 等;requests 和 httpx 通常只重试 5xx。NodeTimeoutError 默认可重试。
生产上应按异常类型与动作语义定制,不能对所有异常无限重试。写操作需要幂等键,重试预算要计入下游限流。
TimeoutPolicy
timeout 仅支持 async node、task 与 entrypoint:
- run_timeout:单 attempt 绝对墙钟上限,不会刷新;
- idle_timeout:没有可观察进展时才触发;
- 两者同时设置时先到者生效。
idle timeout 的自动进展信号包括 State write、stream chunk、child task、stream writer 与 LangChain callback。需要严格定义空闲时,可用 heartbeat-only 模式,由节点显式调用 runtime.heartbeat。
超时后:
- 抛出 NodeTimeoutError;
- 丢弃该 attempt 的 buffered writes;
- retry_policy 决定是否重试;
- 每次 retry 重新计时。
同步阻塞代码要放入 asyncio.to_thread 或改用真正异步驱动;框架无法安全取消进程内同步函数。
Error handler
error_handler 在重试耗尽后接收:
- 当前 State;
- NodeError,包含失败节点名与原异常;
- 可选 Runtime 或 RunnableConfig。
它可以返回普通 State 更新,也可以用 Command 进入 Saga 补偿路径。失败来源会 checkpoint;进程在节点失败后、handler 完成前崩溃,恢复后仍会看到相同 NodeError。
handler 自己失败会向外抛出。每个节点最多一个 handler。
Graph defaults
set_node_defaults 可统一设置 retry、timeout、cache 和 handler:
- 节点显式值优先;
- compile 时解析;
- retry 与 timeout 也适用于 handler 节点;
- handler 节点不会继承 error_handler,也不会缓存;
- 父图默认值不传给子图。
Graceful shutdown
RunControl.request_drain 用于收到 SIGTERM 等信号后,在当前 super-step 完成处停止:
- 不抢占已运行节点;
- 当前 retry 会先成功或耗尽;
- 若仍有后续步骤,抛出 GraphDrained;
- checkpoint 已保存,可用同一 thread_id 与 invoke(None, config) 恢复;
- 子图 drain 会向父图冒泡;
- 它不会强行取消 asyncio task 或线程。
生产 supervisor 仍需设置总的优雅退出期限,超时后再执行平台级取消。
数据保留与运维
- InMemorySaver 重启丢数据,不可作为生产持久化;
- 为 checkpoint 和 Store 配置 TTL、归档或 prune;
- prune DeltaChannel 历史前确认保留 checkpoint 的祖先依赖;
- 删除 thread 时同时清除 checkpoint 与 writes;
- 数据库 migration 独立执行并可重复;
- 监控 checkpoint 大小、写延迟、历史数量、恢复失败和 Store 查询延迟;
- trace、stream、checkpoint、Store 分别做敏感字段保护;
- 定期演练进程崩溃、节点超时、重试耗尽、drain 与恢复。
生产审查清单
- thread_id 是否有租户边界和长度限制?
- 业务需要 exit、async 还是 sync durability?
- 外部副作用是否对 retry、replay、interrupt 恢复都幂等?
- InMemory 实现是否只存在于测试环境?
- setup migration 是否纳入部署步骤?
- checkpoint retention 是否有上限?
- 长消息字段是否评估 DeltaChannel 与回滚策略?
- Store namespace 是否包含 tenant 与 user/project?
- search 是否分页且不依赖后端默认排序?
- 敏感 State 是否加密,并从日志与 stream 中脱敏?
- SIGTERM 是否先 drain,再由平台执行最终超时取消?