PostgresSaver:同样套路的工业强度实现
Sqlite 是单机文件柜、串行化。今天看生产环境真正在用的 PostgresSaver。它做了 Sqlite 偷懒没做的事:把通道值拆进独立的 checkpoint_blobs 表按版本去重(回到 InMemory 的思路)、用 MIGRATIONS 列表管理表结构演进、用一条带 JSONB + JOIN 的 SELECT 在数据库里就把通道值拼好、用 Pipeline 批量提交 + 连接池扛并发。看懂它,你就理解了"一个持久化后端要上生产,除了存取还要操心哪些工程问题"。
· 档案本体、通道值、半成品、迁移记录分四个库房(四张表)分门别类;
· 通道值单独存一个库房,并按版本去重——没变的值多份档共享一份,省地方;
· 有一本"装修施工日志"(MIGRATIONS),记录档案馆每次改造到第几版,换了新版软件能从上次停的地方接着施工;
· 取档时档案馆后台自动把散落的通道值 JOIN 拼齐再一次性交给你,不用你自己跑腿捞。
四张表:档 / 通道值 / 半成品 / 迁移
-- postgres/base.py:47 checkpoints:档本体(注意 checkpoint 列是 JSONB!)
CREATE TABLE IF NOT EXISTS checkpoints (
thread_id TEXT NOT NULL,
checkpoint_ns TEXT NOT NULL DEFAULT '',
checkpoint_id TEXT NOT NULL,
parent_checkpoint_id TEXT,
checkpoint JSONB NOT NULL, -- 档(不含大通道值),JSONB 可被 SQL 内部查询
metadata JSONB NOT NULL DEFAULT '{}',
PRIMARY KEY (thread_id, checkpoint_ns, checkpoint_id)
);
-- postgres/base.py:57 checkpoint_blobs:通道值单独存,按 (channel, version) 去重
CREATE TABLE IF NOT EXISTS checkpoint_blobs (
thread_id TEXT NOT NULL, checkpoint_ns TEXT NOT NULL DEFAULT '',
channel TEXT NOT NULL, version TEXT NOT NULL,
type TEXT NOT NULL, blob BYTEA,
PRIMARY KEY (thread_id, checkpoint_ns, channel, version) -- ← 版本去重的关键
);
-- postgres/base.py:66 checkpoint_writes:半成品写入(同 Sqlite 的 writes)
-- postgres/base.py:44 checkpoint_migrations:只有一列 v,记录已迁移到第几版
checkpoint 列是 JSONB不是 Sqlite 的 BLOB(不透明字节),而是 JSONB——Postgres 能直接在 SQL 里查询/索引 JSON 内部字段。这让"取档时用 SQL JOIN 拼通道值"成为可能(L04)。checkpoint_blobs 独立表回到 InMemory 的做法:通道值从档里抠出来,按 (thread, ns, channel, version) 存。主键含 version → 同一通道同一版本只存一份,多份档共享,跨档去重。checkpoint_migrations 表只有一列 v INTEGER PRIMARY KEY——记录"这个数据库已经迁移到第几版表结构"。这是 Sqlite 完全没有的(Sqlite 靠 IF NOT EXISTS 硬扛,改结构就麻烦)。MIGRATIONS:用列表下标当版本号
Postgres 用一个字符串列表管理所有历史表结构变更(checkpoint/postgres/base.py:43-91):
# postgres/base.py:40 注释:
# "To add a new migration, add a new string to the MIGRATIONS list.
# The position of the migration in the list is the version number."
MIGRATIONS = [
"CREATE TABLE IF NOT EXISTS checkpoint_migrations (v INTEGER PRIMARY KEY);", # v=0
"CREATE TABLE IF NOT EXISTS checkpoints (...);", # v=1
"CREATE TABLE IF NOT EXISTS checkpoint_blobs (...);", # v=2
"CREATE TABLE IF NOT EXISTS checkpoint_writes (...);", # v=3
"ALTER TABLE checkpoint_blobs ALTER COLUMN blob DROP not null;", # v=4
"SELECT 1;", # v=5 故意的 no-op(历史遗留,占位保证后续版本号对齐)
"CREATE INDEX CONCURRENTLY IF NOT EXISTS checkpoints_thread_id_idx ...;", # v=6
"... checkpoint_blobs_thread_id_idx ...", # v=7
"... checkpoint_writes_thread_id_idx ...", # v=8
"ALTER TABLE checkpoint_writes ADD COLUMN IF NOT EXISTS task_path TEXT ...;", # v=9
]
列表下标 = 版本号精妙:迁移在列表里的位置就是它的版本号。第 0 个建 migrations 表本身,第 1 个建 checkpoints……新增变更只需往列表末尾追加一条 SQL,版本号自动是新下标。v=5 的 "SELECT 1;"一个故意的空操作——注释说因为历史上曾加过一条空迁移,为保证"老数据库里记的版本号"和"列表下标"仍对齐,留一个 no-op 占位。这是"迁移一旦发布就不能删改、只能往后加"的活教材。CREATE INDEX CONCURRENTLY后几条加索引用 CONCURRENTLY——建索引时不锁表,生产环境在线加索引不阻塞读写。这是 Sqlite 完全不需要操心的生产级细节。ADD COLUMN IF NOT EXISTSv=9 给 writes 表加 task_path 列——表结构演进的典型:新功能需要新列,老数据库跑迁移补上,老数据用 DEFAULT 值兜底。IF NOT EXISTS 能建新表,但改已有表(加列、加索引)就没法幂等地"只对没改过的库执行"。MIGRATIONS + migrations 表记录版本号,让每个数据库都能"从我当前的版本,把之后的迁移依次补跑一遍"——无论它停在哪一版。这是所有严肃数据库应用(Django/Rails/Flyway…)的标准范式,LangGraph 用最朴素的"列表下标当版本"实现了它。setup:从当前版本断点续跑迁移
setup 的核心就是"查当前版本 → 把之后的迁移补跑"(checkpoint/postgres/__init__.py:85-110):
# postgres/__init__.py:85
def setup(self) -> None:
with self._cursor() as cur:
cur.execute(self.MIGRATIONS[0]) # ① 先确保 migrations 表存在(第 0 条)
results = cur.execute(
"SELECT v FROM checkpoint_migrations ORDER BY v DESC LIMIT 1")
row = results.fetchone()
version = -1 if row is None else row["v"] # ② 查"已迁到第几版",没有则 -1
for v, migration in zip(
range(version + 1, len(self.MIGRATIONS)), # ③ 从"下一版"到"最新版"
self.MIGRATIONS[version + 1 :],
strict=False,
):
cur.execute(migration) # 跑这条迁移
cur.execute("INSERT INTO checkpoint_migrations (v) VALUES (%s)", (v,)) # 记账
if self.pipe:
self.pipe.sync()
先跑 MIGRATIONS[0]第 0 条建 migrations 表本身——鸡生蛋问题:得先有记账的表,才能查"迁到哪了"。它自带 IF NOT EXISTS,重复跑无害。SELECT v ... ORDER BY v DESC LIMIT 1查已记录的最大版本号 = "这个库当前在第几版"。全新库查不到 → version=-1(表示啥都没迁)。range(version+1, len(MIGRATIONS))核心:只跑"当前版本之后"的迁移。新库从 0 全跑;已在 v=6 的老库,只补跑 7/8/9。这就是"断点续跑"。跑一条就 INSERT 记一版每成功执行一条迁移,就往 migrations 表插一行版本号。这样中途崩了,下次 setup 会从崩溃点继续,不会重跑已完成的。setup() MUST be called directly by the user the first time——不像 Sqlite 懒建表,Postgres 要你主动调一次。更重要的坑:MIGRATIONS 列表是"只增不改"的历史。你绝不能删除或修改中间某条已发布的迁移——因为线上数据库可能已经跑过它、记了版本号。改了会导致"新库和老库的表结构分叉"。想改结构?只能在列表末尾追加新迁移。这也是 v=5 那条 no-op 存在的原因:历史包袱一旦背上,就只能小心翼翼往后走。SELECT_SQL:一条查询里 JOIN 拼好通道值
Sqlite 取档要查两次(档 + 半成品)并在 Python 里组装。Postgres 更狠——一条 SQL 用子查询把通道值和半成品全 JOIN 拼好(checkpoint/postgres/base.py:93-118):
-- postgres/base.py:93
select
thread_id, checkpoint, checkpoint_ns, checkpoint_id, parent_checkpoint_id, metadata,
( -- ① 子查询:按 channel_versions 把 blobs 表里的通道值 JOIN 出来
select array_agg(array[bl.channel::bytea, bl.type::bytea, bl.blob])
from jsonb_each_text(checkpoint -> 'channel_versions') -- 展开档里的"通道→版本"
inner join checkpoint_blobs bl
on bl.thread_id = checkpoints.thread_id
and bl.checkpoint_ns = checkpoints.checkpoint_ns
and bl.channel = jsonb_each_text.key -- 通道名对上
and bl.version = jsonb_each_text.value -- 版本号也对上
) as channel_values,
( -- ② 子查询:把这份档的 pending_writes 也一起捞出来
select array_agg(array[cw.task_id::text::bytea, cw.channel::bytea, cw.type::bytea, cw.blob]
order by cw.task_id, cw.idx)
from checkpoint_writes cw
where cw.thread_id = checkpoints.thread_id and ... = checkpoints.checkpoint_id
) as pending_writes
from checkpoints
jsonb_each_text(checkpoint -> 'channel_versions')灵魂操作:直接在 SQL 里展开档 JSONB 里的 channel_versions 字段(通道名→版本号)。这就是为什么 checkpoint 列要用 JSONB 而非不透明 BLOB——SQL 能读进去。inner join checkpoint_blobs on channel + version用展开出的每个 (通道, 版本) 去 blobs 表精确 JOIN 出对应那份通道值。等价于 Day 35 InMemory 的 _load_blobs,但整个过程在数据库内完成、一次返回。array_agg(...)把 JOIN 出的多个通道值聚合成一个数组,作为 channel_values 列随档一起返回。Python 侧拿到就是"档 + 已拼好的通道值 + 半成品"一整包。order by cw.task_id, cw.idx半成品子查询里排序——和 Sqlite 一样保证 pending_writes 顺序稳定可复现。put:小值内联、大值拆进 blobs
存档时它比 Sqlite 多一个聪明判断——不是所有通道值都拆 blob,小的基本类型直接内联进 JSONB(checkpoint/postgres/__init__.py:299-345):
# postgres/__init__.py:299
copy = checkpoint.copy()
copy["channel_values"] = copy["channel_values"].copy()
blob_values = {}
for k, v in checkpoint["channel_values"].items():
if isinstance(v, _DeltaSnapshot): # DeltaChannel 快照 → 拆 blob
blob_values[k] = copy["channel_values"].pop(k)
copy["channel_values"][k] = True # 档里留个占位 True
elif v is None or isinstance(v, (str, int, float, bool)):
pass # ← 小基本类型:留在档 JSONB 里,不拆
else:
blob_values[k] = copy["channel_values"].pop(k) # 其他(大对象)→ 拆进 blobs
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))))
小基本类型 pass(不拆)关键优化:None/str/int/float/bool 这些小值直接留在 checkpoint JSONB 里,不进 blobs 表。因为拆一个小整数进独立表反而更费——多一行、多一次 JOIN。只有大对象(list/dict 等)才值得拆。_DeltaSnapshot → 拆 blob + 占位 TrueDay 29 的 DeltaChannel 快照特殊处理:拆进 blobs,档里留个 True 当"这里有个 delta 快照"的标记。只对 blob_versions 写 blobs和 InMemory 一样:只把"本次变了版本 & 需要拆 blob"的通道写进 blobs 表——没变的复用旧 blob,去重就发生在这。Jsonb(copy)档以 JSONB 形式存(psycopg 的 Jsonb 包装),这样 L04 的 SQL 才能查它内部字段。isinstance 基本类型检查。这体现了成熟系统的特征——针对数据的实际形态做差异化处理,而不是用一个统一但次优的策略糊弄所有情况。Pipeline + 连接池:扛并发的两件武器
Sqlite 用"一把锁串行化",Postgres 用真正的并发机制。核心在 _cursor(checkpoint/postgres/__init__.py:404-443):
# postgres/__init__.py:404
@contextmanager
def _cursor(self, *, pipeline: bool = False):
with self.lock, _internal.get_connection(self.conn) as conn: # 从连接池取连接
if self.pipe:
# 连接已处于 pipeline 模式:可被多协程并发使用,但一次一个游标
try:
with conn.cursor(binary=True, row_factory=dict_row) as cur:
yield cur
finally:
if pipeline:
self.pipe.sync() # 一次性把批量攒的命令发出去
elif pipeline:
if self.supports_pipeline:
with conn.pipeline(), conn.cursor(binary=True, row_factory=dict_row) as cur:
yield cur # 临时开 pipeline,批量发送减少往返
else:
with conn.transaction(), conn.cursor(...) as cur: # 不支持则退回普通事务
yield cur
else:
with conn.cursor(binary=True, row_factory=dict_row) as cur:
yield cur
get_connection(self.conn)它可以接一个 ConnectionPool 连接池——多个请求各拿一条连接真并发操作数据库,而非 Sqlite 那样全挤一把锁。这是生产吞吐的关键。Pipeline 模式psycopg 的 pipeline:把多条 SQL 攒成一批一次性发给数据库,最后 pipe.sync() 统一同步。省掉每条 SQL 一次网络往返——put 里同时写 blobs 和 checkpoints 两条,pipeline 下一次往返搞定。supports_pipeline 降级健壮性:检测数据库/驱动是否支持 pipeline,不支持就退回普通 conn.transaction()。保证在各种 Postgres 版本上都能跑,只是慢一点。binary=True, row_factory=dict_row用二进制协议(更快)+ 行返回成 dict(好取字段)。都是面向性能和易用的生产细节。ON CONFLICT:SQL 级的幂等 + 小结
最后看幂等——Postgres 用标准 ON CONFLICT(Sqlite 是 INSERT OR REPLACE/IGNORE)(checkpoint/postgres/base.py:131-159):
-- postgres/base.py:131 通道值:主键冲突 → 什么都不做(版本相同即同值,无需覆盖)
INSERT INTO checkpoint_blobs (...) VALUES (...)
ON CONFLICT (thread_id, checkpoint_ns, channel, version) DO NOTHING;
-- postgres/base.py:137 档:主键冲突 → 更新(同一 checkpoint_id 重存,覆盖)
INSERT INTO checkpoints (...) VALUES (...)
ON CONFLICT (thread_id, checkpoint_ns, checkpoint_id)
DO UPDATE SET checkpoint = EXCLUDED.checkpoint, metadata = EXCLUDED.metadata;
-- put_writes 两条(:146 UPSERT / :155 INSERT ... DO NOTHING)
-- 全特殊写 → DO UPDATE(允许覆盖);含普通写 → DO NOTHING(幂等去重)
blobs: DO NOTHING通道值主键含 version,版本相同 = 值相同,冲突了不用管——直接 DO NOTHING。这既是幂等、也是去重。checkpoints: DO UPDATE档冲突则用新值覆盖(EXCLUDED 指本次要插的值)。同一 checkpoint_id 重存就更新,对应 Sqlite 的 INSERT OR REPLACE。put_writes 两种和 Sqlite 完全同构:全特殊写用 DO UPDATE 覆盖(如 RESUME)、含普通写用 DO NOTHING 幂等去重。思想跨数据库一致,只是语法从 OR IGNORE 换成 ON CONFLICT DO NOTHING。🧠 今日小结自测
- Postgres 四张表分别是什么?(checkpoints 档 / checkpoint_blobs 通道值去重 / checkpoint_writes 半成品 / checkpoint_migrations 迁移记录)
- MIGRATIONS 列表的版本号怎么来?(列表下标就是版本号,只增不改,v=5 是历史遗留 no-op 占位)
- setup 如何做到"断点续跑"?(查 migrations 表当前版本,只跑之后的迁移,每跑一条记一版)
- 为什么 checkpoint 列用 JSONB?(SQL 能查其内部 channel_versions,从而 JOIN blobs 拼通道值)
- put 为什么不拆所有通道值?(小基本类型内联进 JSONB 更省,只拆大对象和 Delta 快照)
- Postgres 靠什么扛并发?(连接池真并发 + Pipeline 批量减少往返 + 真异步实现,而非 Sqlite 的一把锁)
✋ 10 分钟动手
# 1. 四张表定义 + MIGRATIONS 列表
sed -n '43,91p' libs/checkpoint-postgres/langgraph/checkpoint/postgres/base.py
# 2. 那条"JOIN 拼通道值"的 SELECT_SQL
sed -n '93,118p' libs/checkpoint-postgres/langgraph/checkpoint/postgres/base.py
# 3. setup 断点续跑 + put 小值内联
sed -n '85,110p' libs/checkpoint-postgres/langgraph/checkpoint/postgres/__init__.py
sed -n '308,320p' libs/checkpoint-postgres/langgraph/checkpoint/postgres/__init__.py
# 4. 对照 Sqlite/Postgres 的幂等语法差异
grep -n "OR IGNORE\|OR REPLACE" libs/checkpoint-sqlite/langgraph/checkpoint/sqlite/__init__.py
grep -n "ON CONFLICT" libs/checkpoint-postgres/langgraph/checkpoint/postgres/base.py
(thread_id, checkpoint_ns, checkpoint_id) 三元组转。明天专门拆这套坐标体系:thread_id 怎么隔离不同对话、checkpoint_ns 如何用 | 分隔符表达主图/子图的层级、checkpoint_id 为什么用 uuid6(时间有序 UUID)保证单调递增,以及"超步(super-step)"和 checkpoint 是怎么对应的。你会看到 pregel/_checkpoint.py 里 checkpoint 是怎么被"造"出来的。