PostgresSaver:生产级关系型检查点
职责
PostgresSaver 是 BaseCheckpointSaver 在 Postgres 上的落地实现,也是 LangGraph 官方推荐的生产 saver。代码在 libs/checkpoint-postgres/langgraph/checkpoint/postgres/,同步版入口是 libs/checkpoint-postgres/langgraph/checkpoint/postgres/__init__.py 里的 PostgresSaver (PostgresSaver:40),异步版在 aio.py 的 AsyncPostgresSaver (AsyncPostgresSaver:40),共用基类 BasePostgresSaver (BasePostgresSaver:297),后者负责 SQL 模板和迁移脚本。
跟 SQLite 单表 + 主键去重的玩法不一样,PostgresSaver 把数据拆成三张表:checkpoints(主表,JSONB)、checkpoint_blobs(channel value 二进制)、checkpoint_writes(task 增量写)(MIGRATIONS:43-91)。这种拆法让读取能用单条 SELECT + lateral 子查询一次性把主表 + blobs + writes 全捞回来 (SELECT_SQL:93-118),避免 SQLite 那种 N+1 查询;同时 blob 列独立用 BYTEA,不参与 JSONB 反序列化,大对象不会拖累主表读取。
PostgresSaver 额外支持 psycopg3 的 Pipeline 模式(supports_pipeline:60)——多语句打包下发,服务端不立即返回结果,适合一次 put 内部的 executemany(blob) + execute(checkpoint) 这种组合。
设计动机
- 三表分离,JSONB 存主体:
checkpoints主表用 JSONB 列装CheckpointTypedDict,checkpoint_blobs独立存大对象 (checkpoint_blobs:57-65)。这避免了把整个 channel_values 序列化塞进 JSONB——那种做法在通道多、value 大时性能崩盘。 - 单 SELECT 把 tuple 拼全:
SELECT_SQL用两个 lateral 子查询 (array_aggovercheckpoint_blobs+array_aggovercheckpoint_writes) 一次拿回所有数据 (lateral subqueries:101-117),_load_checkpoint_tuple在 Python 端解 array 即可,不需要二次查询。 - UPSERT 用
ON CONFLICT DO UPDATE/NOTHING区分语义:UPSERT_CHECKPOINT_BLOBS_SQL用DO NOTHING(UPSERT_BLOBS:131-135),UPSERT_CHECKPOINTS_SQL用DO UPDATE(UPSERT_CHECKPOINTS:137-144),UPSERT_CHECKPOINT_WRITES_SQL用DO UPDATE,而INSERT_CHECKPOINT_WRITES_SQL用DO NOTHING(INSERT_WRITES:155-159)。控制 channel 写要覆盖最新状态(用 UPDATE),普通业务写撞了忽略(用 NOTHING),语义跟 InMemory 的if inner_key[1] >= 0 ... continue对齐。 - 迁移列表追加版本号:
MIGRATIONS是 list (MIGRATIONS:43-91),setup()读checkpoint_migrations表的最大 v,只跑后面的迁移 (setup:85-110)。加列、加索引、改约束都通过追加 migration 字符串完成,不改老迁移。 - Pipeline 模式让单连接真并发:同步
PostgresSaver._cursor在pipeline=True时拿conn.pipeline()上下文 (pipeline branch:424-432),一条连接上的多个 cursor 不阻塞互相互;而pipe模式 (从from_conn_string(pipeline=True)) 让一个连接在多线程下复用 (pipe branch:413-423),配合self.lock串行 cursor 但不串行语句。
关键文件
MIGRATIONS:43-91— 三张表的 CREATE 语句 + WAL/索引/task_path 迁移。SELECT_SQL:93-118— 主查询,lateral 子查询把 blobs 和 writes 一起 array_agg 出来。SELECT_PENDING_SENDS_SQL:120-129— 取待发送的 PUSH 任务 writes,按 task_path / task_id / idx 排序。UPSERT / INSERT SQL:131-159— 四条写入语句,区分DO UPDATE/DO NOTHING。BasePostgresSaver:297— 共用基类,持有 SQL 模板常量。PostgresSaver 类:40-60— 同步实现,接Connection或ConnectionPool,带threading.Lock。from_conn_string:62-83— 上下文管理器工厂,pipeline=True走conn.pipeline()。setup:85-110— 读 migration 版本,顺序跑未应用的迁移并记 v。put:263-345— 把_DeltaSnapshot/ 非原始值挪到 blob_values,主体进 JSONB。_cursor:405-442— 按是否 pipeline / 是否 pipe 三种模式挑上下文。AsyncPostgresSaver:40— 异步版,aput/aget_tuple/alist原生 async。
数据流
PostgresSaver.put 是看 saver 怎么拆 Checkpoint TypedDict 的好例子——它先决定哪些 channel value 内联进 JSONB 主体、哪些挪到 checkpoint_blobs,然后一次 cursor 内跑 executemany(blob) + execute(checkpoint):
# inline primitive values in checkpoint table
# others are stored in blobs table
blob_values = {}
for k, v in checkpoint["channel_values"].items():
if isinstance(v, _DeltaSnapshot):
blob_values[k] = copy["channel_values"].pop(k)
copy["channel_values"][k] = True
elif v is None or isinstance(v, (str, int, float, bool)):
pass
else:
blob_values[k] = copy["channel_values"].pop(k)
with self._cursor(pipeline=True) as cur:
if blob_versions := {
k: v for k, v in new_versions.items() if k in blob_values
}:
cur.executemany(
self.UPSERT_CHECKPOINT_BLOBS_SQL,
self._dump_blobs(
thread_id,
checkpoint_ns,
blob_values,
blob_versions,
),
)
cur.execute(
self.UPSERT_CHECKPOINTS_SQL,
(
thread_id,
checkpoint_ns,
checkpoint["id"],
checkpoint_id,
Jsonb(copy),
Jsonb(get_serializable_checkpoint_metadata(config, metadata)),
),
)(put body:309-344) 注意 _DeltaSnapshot 走特殊路径——通道值在主表里只放一个 True 占位,真实数据进 blob,这是 DeltaChannel 重建时反向走祖先链的依据。读回来时 get_tuple 跑 SELECT_SQL,psycopg 把 JSONB 自动转 dict、BYTEA 转 bytes,Python 端再调 _load_checkpoint_tuple 组装 CheckpointTuple (_load_checkpoint_tuple:552)。
边界与失败
from_conn_string(pipeline=True)不能跟 ConnectionPool 混用:__init__明确 raise (pool vs pipe check:52-55),pipeline 是单连接模式,pool 已经是多连接。setup()必须显式调一次:docstring 写明「MUST be called directly by the user the first time checkpointer is used」(setup docstring:85-91)。它建表 + 跑迁移,跑前要给数据库建好权限。- MIGRATIONS 改顺序会炸:迁移按 list 位置当版本号 (
MIGRATIONS:43-91),已经发布过的迁移字符串不能改;改了checkpoint_migrations表里 v 跟代码里位置就错位。 CREATE INDEX CONCURRENTLY在事务里会失败:MIGRATIONS里两条并发建索引语句 (CONCURRENTLY indexes:82-89) 跟其它迁移放同一个setup()上下文里跑,如果连接不在 autocommit 模式会报错;from_conn_string显式autocommit=True是为这个 (autocommit:76-78)。supports_pipeline运行时检测:Capabilities().has_pipeline()(has_pipeline:60) 在__init__里查 psycopg 版本和服务端能力,老版 psycopg 或某些托管 Postgres 退化为普通事务模式 (transaction fallback:434-439)。- DeltaChannel 有专门的快速路径:
get_delta_channel_history在PostgresSaver里被重写 (get_delta_channel_history:444-463) 走两阶段查询(Stage 1 metadata 分页 + Stage 2 UNION ALL 拉 writes),不走基类默认实现,自定义 saver 必须同步实现这一接口才能跑 DeltaChannel。
小结
PostgresSaver 是 LangGraph 推荐生产 saver,三表拆分让「主表 JSONB + 大对象 BYTEA + 增量 writes」各自走最合适的存储;Pipeline 模式和 async 原生实现给到生产规模的吞吐。接口形状见 BaseCheckpointSaver,轻量场景见 InMemorySaver + SqliteSaver。
对照官方资料:LangGraph 文档 · README。