BaseCheckpointSaver:持久化接口的形状
职责
在 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_values、channel_versions、versions_seen、id / ts / v 这几样。put 收到的不只是一个 Checkpoint,还有 metadata 和 new_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 路径走原生驱动。 - 元数据带来源标签:
CheckpointMetadata里source字段取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 重建会静默返回空。
关键文件
BaseCheckpointSaver 类:176— 抽象基类本体,规定 put/get/list/put_writes 等接口。Checkpoint TypedDict:92-123— 一个 checkpoint 的字段形状:channel_values/channel_versions/versions_seen/id/ts。CheckpointTuple:139-146—get_tuple返回值,把 checkpoint、metadata、parent_config、pending_writes 打包成一个 NamedTuple。CheckpointMetadata:38-86—source/step/parents/run_id/ DeltaChannel 用的counters_since_delta_snapshot。get 方法:227-237—get默认委托get_tuple,只返回checkpoint字段。put 方法:277-298— 存整份 checkpoint + metadata,返回更新后的 config。put_writes 方法:300-318— 存 task 在某 checkpoint 下的增量写列表。delete_thread / delete_for_runs:320-348— 线程级清理,带 DeltaChannel 警告。WRITES_IDX_MAP:795—ERROR/SCHEDULED/INTERRUPT/RESUME这四个控制 channel 用负 idx 占坑,避免和真任务写撞 primary key。
数据流
PregelLoop 每跑完一个超步,会调 saver 把当前所有通道值和这步新版本号写下去。这里看 put 的默认签名和 InMemorySaver 的实现就够清楚——引擎调 put,saver 把 channel_values 按 (thread_id, ns, channel, version) 拆出去存 blob,主体进主表,返回新 config 指向刚落的 checkpoint_id。
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_MAP把ERROR/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不能丢:list的filter参数就是按 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 + SqliteSaver 和 PostgresSaver;运行时怎么决定落盘节奏见 Pregel 引擎 和 PregelLoop。
对照官方资料:LangGraph 文档 · README。