← blog 技术原理与实践 · 2026-08-13

LangGraph Checkpoint 深度解析:存储原理、常见坑点与生产实践

基于 LangChain 官方文档与源码,详解 LangGraph checkpoint 的持久化原理(版本号驱动的状态机)、四大核心能力、内置后端选型,以及 thread_id、并发安全、存储膨胀等 21 个生产常见坑点。

28 min read

一、Checkpoint 是什么

Checkpointer(检查点器) 是 LangGraph 的持久化组件:它在图的每个 super-step(超级步)边界保存一次图状态的快照(snapshot),按 thread(线程) 组织。编译图时传入 checkpointer 即可启用,调用时必须指定 thread_id,否则无法保存状态或恢复中断。

from langgraph.checkpoint.memory import InMemorySaver

checkpointer = InMemorySaver()
graph = builder.compile(checkpointer=checkpointer)

# 调用时必须指定 thread_id
result = graph.invoke(
    {"messages": [{"role": "user", "content": "Hi, my name is Bob."}]},
    {"configurable": {"thread_id": "thread-1"}},
)

LangGraph 有两套互补的持久化体系:

  • Checkpointer:保存线程内的图状态快照 → 短期的、线程级记忆(对话连续性、human-in-the-loop、时间旅行、容错)
  • Store:保存应用自定义的跨线程数据 → 长期的、跨线程记忆(用户偏好、事实、共享知识)

四大核心能力

能力说明
Human-in-the-loop人类可以在任意时刻查看/编辑图状态,批准或中断节点执行,之后从该状态恢复
记忆(Memory)多轮对话:同一 thread_id 的后续消息可以访问之前的状态
时间旅行(Time travel)回放到任意历史 checkpoint 调试,或在任意 checkpoint 分叉(fork)出替代路径
容错(Fault tolerance)某一步失败可以从最后成功的 super-step 重启,不重跑已成功的节点(pending writes)

二、Checkpoint 的原理

一句话总结:Checkpoint = 把“图状态”序列化后按 super-step 追加写进一张以 (thread_id, checkpoint_id) 为键的表里;而“图状态”之所以能恢复,靠的是三样东西——channel_values(值)、channel_versions(版本号)、versions_seen(谁看到了哪个版本)。

版本号系统是整个机制的灵魂。它把“谁该在下一步执行”变成纯数据,因此可以持久化、重放、跳过。

1. 数据模型:一个 checkpoint 里到底存了什么

源码里 Checkpoint 是一个 TypedDict,核心字段:

字段含义
v格式版本(当前 1)
idcheckpoint 唯一 ID,用 uuid6(clock_seq=step) 生成,时间有序、可排序、可复现
tsISO 时间戳
channel_values每个 channel(如 messages)此刻的值
channel_versions每个 channel 的单调递增版本号
versions_seen节点名 → {channel → 版本号},记录每个节点已经消费过哪些版本
updated_channels本步更新了哪些 channel

存储层(不管是内存、SQLite 还是 Postgres)都是两张表的抽象:

  • checkpoints:每 super-step 一行,行 = (checkpoint_id, checkpoint blob, metadata, parent_id)
  • checkpoint_writes(源码里叫 writes):每节点输出一行,键是 (thread_id, ns, checkpoint_id, task_id, channel_idx)

2. 核心机制:版本号如何驱动执行

Checkpoint 本身是“死的快照”,让它活起来的是版本号。原理:

  1. 每个 channel 有一个整数版本号,被写入时 get_next_version 就 +1。

  2. 节点运行时,它看到的 channel 版本被记进 versions_seen[node]。

  3. 每一步开始时,pregel 循环对每个节点问一个问题:

    “你订阅的 channel 里,有没有版本号 > 你上次看到(versions_seen)的?”

    有 → 节点被调度执行;没有 → 跳过。这就是 _algo.py 里 whats_new 的判断逻辑(version > seen.get(chan, null_version))。

  4. 节点执行完,产出写入 channel,版本号递增,versions_seen 更新。

所以图的执行进度完全编码在 checkpoint 里——恢复执行 = 把 checkpoint 读出来,比较一遍版本,自然就知道该从哪个节点继续。

3. 写入时序:先写日志、再提交快照(WAL 式)

LangGraph 的持久化是先写“节点输出日志”,后提交“步快照”(一个 super-step 内):

┌─ 第 N 步 ─────────────────────────────────────────────┐
│ 1. 读取最新 checkpoint,算出哪些节点要跑(版本比较)      │
│ 2. 并行执行节点                                        │
│ 3. 每个节点一完成 → put_writes() 立即把它的输出          │
│    写入 writes 表(键含 task_id)   ← 写前日志          │
│ 4. 该步所有节点完成 → apply_writes() 合并到内存状态,     │
│    版本号递增,versions_seen 更新                       │
│ 5. put() 把新 checkpoint 追加进 checkpoints 表         │
│    (新行,不动旧行;id 用 uuid6(clock_seq=step))      │
└───────────────────────────────────────────────────────┘

先 put_writes 再 put 是容错的根基:第 3 步里别的节点崩了,已成功节点的输出已经在 writes 表里,不用重算。恢复时把最后成功的 checkpoint + 它挂着的 pending writes 一起读出来,重放即可。重放是幂等的(put_writes 里有 WRITES_IDX_MAP 索引,同 (task_id, channel) 已存在就跳过)。没输出的节点会写 null write 标记“我跑过了”,否则恢复时会被误判为没执行。

4. 恢复、时间旅行、分叉的底层原理

三者本质是同一个操作:按 checkpoint_id 定位到历史某行。

  • 恢复/续跑:读最新 checkpoint → 反序列化 → 应用 pending writes → 版本比较决定下一步节点 → 继续循环。step 编号从 -1(输入)开始。
  • 时间旅行:get_state_history 沿着 parent_config 链表把整条历史列出来(list() 按 checkpoint_id 排序,最新在前);回放 = 指定老的 checkpoint_id 重新 invoke。
  • 分叉(fork):不修改原 checkpoint,而是新建一个 source="fork" 的 checkpoint 作为分支起点——append-only,永不原地修改,天然支持多分支。

5. ID 为什么用 uuid6(clock_seq=step)

源码:id=str(uuid6(clock_seq=step))。uuid6 是时间有序 UUID(前 48 位是时间戳),保证 id 本身就能排序,get_state_history 按 id 排序即得时间序,无需额外时间列索引。clock_seq=step 让同一步的所有变体可复现、且同 step 内唯一。

6. 存储层抽象

所有后端都实现 BaseCheckpointSaver 接口,需实现以下五个方法(异步版对应 aput / aput_writes / aget_tuple / alist / aget_next_version):

  • put(config, checkpoint, metadata, new_versions) → 提交快照,返回带新 checkpoint_id 的 config
  • put_writes(config, writes, task_id) → 写节点输出日志
  • get_tuple(config) → 按 (thread_id, checkpoint_id) 取 CheckpointTuple(checkpoint + metadata + parent + pending_writes)
  • list(...) → 遍历历史
  • get_next_version(current, channel) → 生成下一个版本号(同步/异步版共用)

thread_id 是分区键,checkpoint_ns 区分父子图(子图在 "node:uuid" 命名空间下独立存)。序列化走可插拔的 serde(默认 JsonPlusSerializer)。

新版本的一个细节:InMemorySaver.put 里 channel_values 被从 checkpoint 行里剥出来,单独存进 blobs 表,键是 (thread_id, ns, channel, version)——checkpoint 行只留 channel_versions,取值时按版本查 blob。这就是 DeltaChannel 的存储级基础:值按版本去重,重复的增量只存一份。

7. 一个例子串起全流程

图 START → A → B → END,线程 t1:

checkpoint_id  step  source   channel_versions        versions_seen
--------------------------------------------------------------------
c0             -1    input    {}                      {}
c1             0     loop     {x:1}     (A 写入)       A:{x:1}
c2             1     loop     {x:2}     (B 写入)       A:{x:1}, B:{x:2}
c3             1     exit     {x:2}     终态
  • 二次 invoke 时读到 c3:A 的 versions_seen 是 {x:1},当前 x=2 → 不变,跳过;没有节点需要跑 → 直接结束。这就是“多轮对话不吃掉历史状态”的原理。
  • 若 c2 存盘前崩溃:读 c1 + pending writes(B 的输出已在 writes 表)→ 重放 B 的写入即可,A 白干不了。

三、核心 API 与持久化级别

# 获取最新状态快照(StateSnapshot)
snap = graph.get_state(config)

# 获取指定 checkpoint_id 的历史状态
snap = graph.get_state({"configurable": {"thread_id": "1", "checkpoint_id": "..."}})

# 获取该线程全部历史(最新在前)
history = list(graph.get_state_history(config))

# 在任意 checkpoint 上分叉/修改状态(生成新 checkpoint,不修改原 checkpoint;带 reducer 的通道会累加而非覆盖)
graph.update_state(config, {"messages": [...]}, as_node="node_a")

# 回放:给旧 checkpoint_id 重新 invoke,之前的节点跳过,之后的节点重跑
result = graph.invoke(None, {"configurable": {"thread_id": "1", "checkpoint_id": "..."}})

调用 graph.stream(..., durability=...) 可选三档持久化,选型取决于你对性能和崩溃安全性的权衡:

模式行为典型适用场景
"exit"只在图退出时(成功/报错/interrupt)落盘高吞吐、能接受进程崩溃丢中间状态的应用;最适合频繁调用、量大的生产路径
"async"异步落盘,与下一步执行并行性能与持久性折中的默认推荐;崩溃时可能丢少量最近 checkpoint,但绝大多数历史仍在
"sync"每步开始前同步落盘崩溃安全要求极高、单步不可重算或代价昂贵的场景(如已产生外部副作用);有性能开销

四、内置 Checkpointer 后端

后端包用途
InMemorySaverlanggraph-checkpoint(内置)实验/单进程,重启即丢
SqliteSaver / AsyncSqliteSaverlanggraph-checkpoint-sqlite本地开发、单机生产
PostgresSaver / AsyncPostgresSaverlanggraph-checkpoint-postgres生产首选(LangSmith 内部即用它)
CosmosDB、Redis、MongoDB、DynamoDB、ScyllaDB、CockroachDB、Aerospike、Tigris、TypeDB、Supabase 等各自独立包按云环境选型

生产警示:PostgresSaver 的 thread_id 列有长度限制(<255 字符);MemorySaver 重启即丢,生产必须换持久化后端;长对话 checkpoint 会无限增长,需要定期 prune。

序列化与加密

  • 默认 JsonPlusSerializer(基于 msgpack/JSON),支持 LangChain 原生类型、pydantic v2、datetime、enum、numpy 等;不支持的(如 Pandas DataFrame)可开 pickle_fallback=True。
  • 可用 EncryptedSerializer(AES)加密落盘,密钥从 LANGGRAPH_AES_KEY 环境变量读取。

存储优化:DeltaChannel(beta,需 langgraph>=1.2)

默认每个 checkpoint 都写全量通道值,多轮长对话会疯狂膨胀。DeltaChannel 只存增量(checkpoint blob 里只存一个 MISSING 哨兵),读取时沿祖先链回放 writes 通过 reducer 重建状态,使 checkpoint 大小从 O(N) 降到 O(1)。

五、常见坑点清单

thread_id 相关(最高频)

1. 忘了传 thread_id → 直接报错 编译图时挂了 checkpointer,invoke 的 config 里就必须带 {"configurable": {"thread_id": ...}}。不传就无法保存状态、无法在 interrupt 后恢复。

2. thread_id 过长 → Postgres 索引爆炸,且殃及所有请求(langgraph#6239)

there is no unique or exclusion constraint matching the ON CONFLICT specification

官方解决方案:约束 thread_id 的尺寸(postgres 列有长度限制,<255 字符)。坑的是:一个超长 thread_id 不仅自己失败,连短 thread_id 的调用也会失败。

3. thread_id 类型混用(int vs str)→ “幽灵”无历史 thread_id=1 和 thread_id="1" 在部分后端是不同的键。存的时候用 int,取的时候用 str,就会得到空状态、以为“记忆丢了”。

深度变体(langgraph#6623):Postgres 下 interrupt 后 checkpoint_writes 表里出现两种形式的 thread_id(普通字符串 vs 十六进制编码),用字符串形式查询时部分状态缺失。排查时直接 SELECT DISTINCT thread_id FROM checkpoint_writes 看是否有两种形态。

4. 自定义 context schema 里给 thread_id 塞默认值

thread_id: str = Field(default_factory=lambda: str(uuid.uuid4()))  # ❌ 会出错

正确做法:让平台自动注入,或 thread_id: str | None = Field(default=None)。

并发问题(生产事故高发区)

5. msgpack 非线程安全 + 浅拷贝 → checkpoint 损坏(langgraph#4168)

AsyncPostgresSaver 在 asyncio.to_thread 里打包 checkpoint 时,checkpoint.copy() 只是浅拷贝——如果另一个线程正在往 messages 里 append,msgpack 算出的数组长度头就和实际内容对不上,报 unpack(b) received extra data,写进去的就是损坏的 checkpoint,之后所有 invoke 全部失败。消息量大时必现。

教训:不要在节点外部并发地边写状态边跑图;生产用 Postgres 时给线程加锁或避免同一线程并发写。

6. 共享单例 agent + 并发 invoke → 跨线程数据污染(langgraphjs#2040)

生产多租户场景:单例 agent + worker concurrency: 2,线程 A 的对话数据出现在线程 B 的 checkpoint 和 LLM 响应里。根因指向 AsyncLocalStorageProviderSingleton(静态初始化、进程内全局传播 async context)在并发调用间泄漏上下文。冒烟证据:某 LLM 调用 cache_read_input_tokens: 18,878 > prompt_tokens: 9,719,缓存内容来自别的线程。

修复状态:该问题源于 LangGraphJS 的 AsyncLocalStorage 上下文泄漏,已在 langgraphjs#2552 修复,属特定版本缺陷而非当前仍存在的设计问题。引用的 issue(#2040)于 2024 年报告,修复已随后续版本合入。建议现在做两件事:确认你使用的 langgraphjs 版本已包含 #2552 的修复(升级到涵盖该 PR 的 release);在升级前的高危场景仍建议每 invocation 新建 agent 或 concurrency: 1 作为临时规避。

7. 同一 thread_id 并发 invoke → 状态竞争/丢数据 两个请求同时往一个线程写,checkpoint 互相覆盖。需要业务层加锁(如 Redis NX lock)。

8. 后端本身的线程安全限制

  • InMemorySaver 不是线程安全的,多线程并发读写会出问题。
  • SqliteSaver 官方明说 “for lightweight, synchronous use cases… does not scale to multiple threads”——多进程/多线程下会撞 database is locked。并发场景必须 Postgres/Redis 等。

存储膨胀与 TTL(运维头号坑)

9. checkpoint 无限增长,没有内置 TTL

每次 super-step 都存全量通道值,多轮长对话指数级膨胀。实测数据(Scaling LangGraph Postgres Checkpointer):平均一次对话在 4 张表里产生 ~93 条记录,staging 一周 checkpoint_blobs 就涨到 56MB / 1.8 万条。OSS PostgresSaver 没有开箱即用的 TTL,必须自己写 cron DELETE(建议分区删除避免锁表)。

10. TTL 只对启用后新建的 checkpoint 生效 启用 TTL 之前的历史数据不会自动清理,要单独跑一次清理。

11. 大对象(图片/PDF/base64)塞进 state → checkpoint 爆炸

MongoDB 有 16MB document 上限(DocumentTooLarge),Postgres 单字段上限 1GB。正确姿势:大对象存外部存储(S3),state 里只存引用;或放 LangGraph Store。

12. 缓解手段排序

  • DeltaChannel(langgraph>=1.2,beta):只存增量,checkpoint 从 O(N) 降到 O(1),但读时沿祖先链重放,有延迟代价。
  • durability="exit":只在退出时落盘,省存储但进程崩溃中间状态全丢。
  • 官方 prune 扩展能力(aprune)做历史修剪。

连接与资源

13. PostgresSaver 直接持有连接 → 长任务超时

用 from_conn_string 的方式会在整个 run 期间占用一个数据库连接,长 workflow 直接连接超时。正确做法:用 ConnectionPool。

14. 自建表与迁移冲突

checkpointer.setup() 建的表和团队标准 Alembic 迁移冲突。建议把官方建表 SQL 抽成自己的迁移,跳过 setup(),并持续关注上游变更。

15. 给子图也挂 checkpointer → 存储翻倍 + 语义混乱

每个子图有独立 checkpoint_ns,状态各自存一份。建议:只在 supervisor 图挂 saver,评估子图前问三个问题(是否需要独立暂停/恢复?是否需要该粒度的时间旅行?子图无状态是否可接受?)。

语义误解(代码没错但行为不对)

16. 父图 get_state 看不到子图状态 子图状态在 "node:uuid" 命名空间下,父图快照里没有。想拿到子图实时状态要用 stream(..., subgraphs=True);跨图共享数据用 Store。

17. update_state 是 fork 不是修改 它会生成一个 source="update" 的新 checkpoint,原 checkpoint 不动;带 reducer 的 channel 会累加而非覆盖。

18. replay 时 interrupt 会重新触发 指定旧 checkpoint_id 回放,checkpoint 之前的节点跳过、之后的节点全部重跑(包括 LLM 调用、interrupt)。恢复中断必须用 Command(resume=...) 且保持同一 thread_id,否则读到的是另一条链。

19. 自定义 checkpointer 的隐藏要求(官方文档 “Build a custom checkpointer”)

BaseCheckpointSaver 的以下 5 个方法全都要实现,缺一个运行时抛 NotImplementedError:

# 同步接口
def put(self, config, checkpoint, metadata, new_versions) -> config
def put_writes(self, config, writes, task_id) -> None
def get_tuple(self, config) -> Optional[CheckpointTuple]
def list(self, config, *, filter=None, before=None, limit=None, cursor=None) -> Iterator[CheckpointTuple]
def get_next_version(self, current, channel) -> Any

# 异步接口(异步后端需实现)
async def aput(self, config, checkpoint, metadata, new_versions) -> config
async def aput_writes(self, config, writes, task_id) -> None
async def aget_tuple(self, config) -> Optional[CheckpointTuple]
async def alist(self, config, *, filter=None, before=None, limit=None, cursor=None) -> AsyncIterator[CheckpointTuple]
async def aget_next_version(self, current, channel) -> Any

其他隐藏要求:

  • metadata 不能裁剪未知 key(如 DeltaChannel 的 counters_since_delta_snapshot,静默丢弃会破坏特性)。
  • get_tuple 的指定 checkpoint_id 路径必须正确——它坏了会静默损坏 DeltaChannel 状态(每次 invoke 都要用它重建)。
  • 行键设计:checkpoint_id 必须是可排序主键(ULID),“按 id 精确查”和“取最新”都要 O(1),靠扫描不扩展。
  • 官方提供 langgraph-checkpoint-conformance 一致性测试套件——这是一套由官方发布的、专用于校验自定义 checkpointer 是否符合 LangGraph 存储契约的一致性测试,内置覆盖 put/get/list/版本号/父链/DeltaChannel 等行为的断言集。写自己的 checkpointer 时必须跑且建议纳入 CI,否则存储契约一旦偏差,线上会在状态恢复、时间旅行等路径上静默出错,裸测很难发现。

序列化

20. JsonPlusSerializer 不认识的对象 → 直接炸 Pandas DataFrame 等默认不支持,需要 JsonPlusSerializer(pickle_fallback=True)。自定义 metadata 也别直接 pickle,用 self.serde。

21. 加密序列化器忘配密钥 EncryptedSerializer.from_pycryptodome_aes() 从 LANGGRAPH_AES_KEY 读密钥,没配环境变量直接失败。

六、一句话总结

checkpoint 的坑集中在四类——thread_id 的规范(长度/类型/必传)、并发安全(msgpack 损坏、上下文泄漏、后端线程限制)、存储膨胀(无 TTL、全量快照、大对象入 state),以及语义误解(update_state 是 fork、子图隔离、replay 重跑)。

checkpointer 本身 = 基于版本号的持久化状态机:channel_versions 提供单调递增的“世界时钟”,versions_seen 记录每个节点的消费进度,writes 表做写前日志保证容错,checkpoint_id(uuid6) + parent_config 链表支撑时间旅行与分叉。执行进度、恢复点、历史分支全部编码成可持久化的数据,这是它区别于普通 KV 缓存的核心。

参考来源

Sources

No external sources for this entry.

Related