Day 37 / 共 60 天 · 阶段 6 持久化与记忆

PostgresSaver:同样套路的工业强度实现

Sqlite 是单机文件柜、串行化。今天看生产环境真正在用PostgresSaver。它做了 Sqlite 偷懒没做的事:把通道值拆进独立的 checkpoint_blobs 表按版本去重(回到 InMemory 的思路)、用 MIGRATIONS 列表管理表结构演进、用一条带 JSONB + JOIN 的 SELECT 在数据库里就把通道值拼好、用 Pipeline 批量提交 + 连接池扛并发。看懂它,你就理解了"一个持久化后端要上生产,除了存取还要操心哪些工程问题"。

📍 阶段 6 · 持久化与记忆(8 天)你在这里
D33 概念 D34 接口 D35 InMemory D36 Sqlite D37 Postgres D38 id体系 D39 serde D40 Store
💡 用一个类比先兜住今天 如果 Sqlite 是"办公室角落的单人铁皮柜",Postgres 就是"专业档案馆":
· 档案本体、通道值、半成品、迁移记录分四个库房(四张表)分门别类;
· 通道值单独存一个库房,并按版本去重——没变的值多份档共享一份,省地方;
· 有一本"装修施工日志"(MIGRATIONS),记录档案馆每次改造到第几版,换了新版软件能从上次停的地方接着施工
· 取档时档案馆后台自动把散落的通道值 JOIN 拼齐再一次性交给你,不用你自己跑腿捞。
L01

四张表:档 / 通道值 / 半成品 / 迁移

🤔 痛点:Sqlite 两张表就够了,为什么 Postgres 要四张? 因为它要做 Sqlite 省掉的两件事:通道值去重(多一张 blobs 表)和表结构版本管理(多一张 migrations 表)。看建表定义(checkpoint/postgres/base.py:43-91):
-- 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 硬扛,改结构就麻烦)。
数据结构:Postgres 的四张表分工 checkpoints档本体 JSONB(不含大通道值) checkpoint_blobs通道值按版本去重PK 含 version checkpoint_writes半成品写入五元组主键 ..migrations迁移版本记录单列 v 取档时 SELECT 用 JSONB 里的 channel_versions JOIN blobs 表把通道值拼回(L04) 对比 Sqlite:只有前两张的合并版(通道值内联进 checkpoint 列,无 blobs/migrations)
四张表各司其职;blobs 表按 (channel, version) 去重是比 Sqlite 多做的关键一步
🅰 设计取舍①:Postgres 为何愿意付出"四张表 + JOIN"的复杂度换 blob 去重,而 Sqlite 不愿? 因为定位不同。Sqlite 面向 demo/小项目,数据量小、图省事,整包存的冗余无所谓。Postgres 面向生产:可能有成千上万个长对话线程,messages 这类大通道若每份 checkpoint 都整包重存,磁盘和 IO 成本会爆炸。此时"多一张 blobs 表 + 取档时 JOIN 拼"的复杂度,换来的是没变的通道跨 checkpoint 只存一份——在长对话场景能省下数量级的存储。这是"实现复杂度 vs 规模化存储成本"的取舍,随规模上升,天平必然倒向去重。
L02

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 值兜底。
💡 为什么"版本化迁移"是生产持久化的刚需?软件会升级,表结构会变。Sqlite 靠 IF NOT EXISTS 能建新表,但改已有表(加列、加索引)就没法幂等地"只对没改过的库执行"。MIGRATIONS + migrations 表记录版本号,让每个数据库都能"从我当前的版本,把之后的迁移依次补跑一遍"——无论它停在哪一版。这是所有严肃数据库应用(Django/Rails/Flyway…)的标准范式,LangGraph 用最朴素的"列表下标当版本"实现了它。
L03

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 从当前版本断点续跑迁移 查 migrations 表当前版本 = 6 range(7, len(MIGRATIONS))只跑 7,8,9 三条 跑一条记一版INSERT v=7,8,9 全新库:version=-1 → 从 0 全跑;老库:只补跑没跑过的 → 幂等、可中途崩溃续跑 MIGRATIONS 列表只增不改,下标即版本号(v=5 是历史 no-op 占位)
"查当前版本 → 跑之后的迁移 → 记账"三步,让任何版本的库都能平滑升级
⚠ 边界:setup 必须由用户显式调用一次,且迁移顺序不可逆 docstring 明确(postgres/__init__.py:89-90):setup() MUST be called directly by the user the first time——不像 Sqlite 懒建表,Postgres 要你主动调一次。更重要的坑:MIGRATIONS 列表是"只增不改"的历史。你绝不能删除或修改中间某条已发布的迁移——因为线上数据库可能已经跑过它、记了版本号。改了会导致"新库和老库的表结构分叉"。想改结构?只能在列表末尾追加新迁移。这也是 v=5 那条 no-op 存在的原因:历史包袱一旦背上,就只能小心翼翼往后走。
L04

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 顺序稳定可复现。
🅰 设计取舍②:为什么把"拼通道值"下推到数据库 SQL 里做,而不是像 Sqlite/InMemory 在 Python 里拼? 因为减少网络往返和数据传输。Sqlite 在同进程、Python 里拼没成本;但 Postgres 是远程数据库,每次查询都要走网络。如果先查档、再回 Python 逐个通道发查询捞 blob,就是"N+1 次网络往返"——延迟灾难。把 JOIN 下推到一条 SQL,一次往返拿回全部所需数据,让数据库(离数据最近、最擅长 JOIN)干这活。这是分布式系统"能在数据端算的就别拉回来算"的通用智慧。代价是这条 SQL 相当复杂、可读性差——但性能收益在生产规模下完全值得。
L05

put:小值内联、大值拆进 blobs

存档时它比 Sqlite 多一个聪明判断——不是所有通道值都拆 blob,小的基本类型直接内联进 JSONBcheckpoint/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 才能查它内部字段。
💡 "小值内联、大值外拆"是分层存储的精髓这是 Postgres 版比 InMemory/Sqlite 更进一步的优化:不搞一刀切。小标量内联省去 JOIN、大对象外拆享受去重,两头的好处都占。判断标准就是简单的 isinstance 基本类型检查。这体现了成熟系统的特征——针对数据的实际形态做差异化处理,而不是用一个统一但次优的策略糊弄所有情况
L06

Pipeline + 连接池:扛并发的两件武器

Sqlite 用"一把锁串行化",Postgres 用真正的并发机制。核心在 _cursorcheckpoint/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(好取字段)。都是面向性能和易用的生产细节。
对比记忆:Sqlite = 一把锁串行 + 不支持异步;Postgres = 连接池并发 + Pipeline 批量 + 真异步(aio.py 里有 AsyncPostgresSaver)。这正是 Day 36 结尾说的"有真 IO 时,异步要认真实现"——Postgres 就认真实现了。
L07

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
💡 一条主线:幂等永远交给"主键约束 + 冲突策略"从 InMemory 的"字典键判存在"、Sqlite 的"INSERT OR IGNORE"、到 Postgres 的"ON CONFLICT DO NOTHING"——三个后端实现不同,但幂等的思想完全一致:用唯一键定义"什么算重复",用冲突策略定义"重复了怎么办",绝不用应用层"查了再改"。这就是为什么 Day 35 说 put_writes 幂等是"持久执行的地基"——它在每个后端都被一丝不苟地实现。

🧠 今日小结自测

  • 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
🔮 明日预告 · Day 38 thread / checkpoint_ns / 超步三个 saver 都围着 (thread_id, checkpoint_ns, checkpoint_id) 三元组转。明天专门拆这套坐标体系thread_id 怎么隔离不同对话、checkpoint_ns 如何用 | 分隔符表达主图/子图的层级、checkpoint_id 为什么用 uuid6(时间有序 UUID)保证单调递增,以及"超步(super-step)"和 checkpoint 是怎么对应的。你会看到 pregel/_checkpoint.py 里 checkpoint 是怎么被"造"出来的。
← Day 36 SqliteSaver Day 38 · id 体系与超步 →