Skip to content

BaseCheckpointSaver:持久化接口的形状

源码版本1.2.9

职责

在 LangGraph 的运行时里,Pregel 引擎每跑完一个超步 (superstep) 都要把自己当前所有通道 (channel) 的值连同这一步产生的写 (writes) 一起存下来,这样下次进来同一个 thread_id 才能接着上一次的状态往下跑,也能在被打断 (interrupt) 之后恢复、回放历史、做时间旅行调试。BaseCheckpointSaver 就是这层「持久化」的抽象接口 (BaseCheckpointSaver:176)。它本身不关心你把数据存到哪——内存 dict、SQLite、Postgres、Redis 都行——只规定 saver 必须实现的几个方法:put / get / get_tuple / list / put_writes (put / get / list / put_writes:227-318)。

它同时把检查点 (checkpoint) 这个概念的形状钉死成 Checkpoint TypedDict:Checkpoint TypedDict:92-123,里面装着 channel_valueschannel_versionsversions_seenid / ts / v 这几样。put 收到的不只是一个 Checkpoint,还有 metadatanew_versions,saver 决定哪些字段进主表、哪些进 blob 表。put_writes 则单独存「这条任务在这步产出的写」,跟 checkpoint 解耦,这样写可以增量落盘,不用每写一次就重写整份快照。

换句话说,BaseCheckpointSaver 是 PregelLoop 和具体存储后端之间的合同:引擎只认这套接口,后端换实现不用改引擎;后端写自己的 saver 只要照着五个方法填,引擎就自动获得该后端的并发和持久化特性。

设计动机

  • 状态和写分开存:put 存整份快照,put_writes 存单个 task 在某步内的增量写 (put_writes:300-318)。拆开是为了让中断恢复时能精确还原「这一步哪个 task 已经写过、哪个还没」,而不是只能整体回滚到上一个 checkpoint。
  • thread_id 当主键:put 返回的 config 里只装 thread_id / checkpoint_ns / checkpoint_id (put:277-298),跨调用、跨进程都以这三个字段定位一个 checkpoint。checkpoint_ns 是给子图用的命名空间,空串就是顶层图。
  • 同步 + 异步双套接口:每个写方法都有 aput / aget / alist / aput_writes 的 async 版本 (async methods:417-509),默认实现抛 NotImplementedError。saver 子类可以只实现同步版,异步版默认走 asyncio.to_thread 转发,也可以自己写真正的 async 路径走原生驱动。
  • 元数据带来源标签:CheckpointMetadatasource 字段取 input / loop / update / fork (source:41-48)。list 和过滤都靠 metadata,所以 saver 必须把 metadata 跟 checkpoint 一起存,不能丢。
  • DeltaChannel 旁路:prune / delete_for_runs / copy_thread 的 docstring 反复警告 (DeltaChannel-aware ops:320-415),如果图用了 DeltaChannel,自定义 saver 不能只删最新 checkpoint,得把祖先链上的 checkpoint_writes 一并保留,否则 delta 重建会静默返回空。

关键文件

数据流

PregelLoop 每跑完一个超步,会调 saver 把当前所有通道值和这步新版本号写下去。这里看 put 的默认签名和 InMemorySaver 的实现就够清楚——引擎调 put,saver 把 channel_values 按 (thread_id, ns, channel, version) 拆出去存 blob,主体进主表,返回新 config 指向刚落的 checkpoint_id。

python
def put(
    self,
    config: RunnableConfig,
    checkpoint: Checkpoint,
    metadata: CheckpointMetadata,
    new_versions: ChannelVersions,
) -> RunnableConfig:
    """Store a checkpoint with its configuration and metadata."""
    raise NotImplementedError

这是基类里的纯接口 (put method:277-298),不带默认实现。每个子类得自己决定哪些字段进主表、哪些进 blob、用不用 ON CONFLICT DO UPDATE 之类。读回来时 get_tuple 要负责把 blob 反序列化拼回 channel_values,并把同一 checkpoint 下的 pending_writes 一起装进 CheckpointTuple (get_tuple:239-251)。

边界与失败

  • put_writes 用负 idx 占坑控制 channel:WRITES_IDX_MAPERROR / SCHEDULED / INTERRUPT / RESUME 映射到 -1 / -2 / -3 / -4 (WRITES_IDX_MAP:795),自定义 saver 的主键 (task_id, idx) 必须能接受负值,否则控制写会撞主键冲突。
  • get 默认调 get_tuple:子类只实现 get_tuple 就够了,get 默认实现帮你取 .checkpoint (get:227-237)。但 get_tuple 不实现就是 raise NotImplementedError,不能 silently 返回 None
  • async 默认是 NotImplementedError:同步 saver 不实现 aput 时,默认 async 路径不会自动转线程——要么子类自己写 asyncio.to_thread 转发,要么调用方调同步路径 (async defaults:468-509)。
  • metadata 不能丢:listfilter 参数就是按 metadata 字段过滤的 (list method:253-275),saver 实现如果不存或存错 metadata,过滤和分页都会哑掉。
  • DeltaChannel 下不能 naive 删:基类 prune 的 docstring 明确说 (prune:374-415)——只留最新 checkpoint 会让 delta channel 重建静默返回空,不报错。自定义 saver 实现删除前要走祖先链。
  • checkpoint_ns 默认空串:子图调用会带非空 ns,顶层调用是 ""。saver 的主键必须把 ns 算进去,不然不同子图的同名 checkpoint 会互相覆盖。

小结

BaseCheckpointSaver 把「检查点怎么存、怎么读、怎么列、怎么写增量」这几件事一次性钉成五方法 + 一份 TypedDict 形状,Pregel 引擎只跟接口打交道,后端实现任意换。具体两种参考实现见兄弟页面 InMemorySaver + SqliteSaverPostgresSaver;运行时怎么决定落盘节奏见 Pregel 引擎PregelLoop

对照官方资料:LangGraph 文档 · README