首页
/ LangGraph Checkpoint Postgres:基于 PostgreSQL 的 Agent 状态持久化实现与实战

LangGraph Checkpoint Postgres:基于 PostgreSQL 的 Agent 状态持久化实现与实战

2026-09-05 14:31:35作者:昌雅子Ethen

本篇技术文章围绕 libs/checkpoint-postgres 包展开,完整讲解 langgraph-checkpoint-postgres 的安装配置、PostgresSaver / AsyncPostgresSaver 的同步与异步用法、首次使用时必须调用的 .setup() 迁移机制,以及 README 中强调的连接参数要求与安全加固要点。读完本文,你将能够独立把 LangGraph 的长期工作流状态持久化到 PostgreSQL,理解 checkpoint 数据在数据库中的真实存储结构,并知道如何运行该包自带的测试来验证你的本地环境。

这个包是什么

langgraph-checkpoint-postgres 提供 LangGraph checkpoint saver 的 Postgres 实现。当你需要为 LangGraph 的状态持久化提供 Postgres 后端,以支撑可恢复、长时运行的工作流和 Agent 时,就使用这个库。它是 langgraph-checkpoint 抽象接口的一个具体落地:图每次 super-step 执行后,checkpoint、通道值与中间写入(pending writes)都会被序列化并落到 Postgres 的几张表中,后续可以按 thread_id 精确恢复、列出或回退到任意历史节点。

pyproject.toml 可以看到当前包版本为 3.1.2,Python 要求 >=3.10,核心依赖为:

  • langgraph-checkpoint>=4.1.0,<5.0.0:checkpoint 的基类与 serde 序列化层;
  • psycopg>=3.2.0psycopg-pool>=3.2.0:Psycopg 3 数据库驱动及其连接池;
  • orjson>=3.11.5:高性能 JSON 编解码。

包内除了 saver 还附带了一个跨线程长期记忆库的 Postgres 实现(PostgresStoreAsyncPostgresStore),本文聚焦 checkpoint 部分。

安装与环境准备

按 README 的 Quick Install,使用 uv 安装即可:

uv add langgraph-checkpoint-postgres

默认情况下 langgraph-checkpoint-postgres 安装的是不带任何 extras 的 psycopg(Psycopg 3)。你可以按 Psycopg 官方安装文档自行选择最合适的安装方式,例如加上预编译二进制轮子的 psycopg[binary](该包自己的测试依赖组在 pyproject.toml 中正是使用 psycopg[binary])。

本地开发时可以用仓库自带的 docker compose 拉起测试数据库。Makefile 提供了 start-postgres / stop-postgres 目标,其背后是 tests/compose-postgres.yml

  • 使用 pgvector/pgvector:pg16 镜像(POSTGRES_VERSION 默认 16,也可取 15);
  • 端口映射为 5441:5432,账号/密码均为 postgres,并加载 vector 扩展;
  • make test 会依次在 PostgreSQL 15 和 16 两个版本上跑完整测试。

测试夹具 tests/conftest.py 中使用的默认连接串为 postgres://postgres:postgres@localhost:5441/postgres?sslmode=disable,且每个测试前会清空 checkpointscheckpoint_blobscheckpoint_writescheckpoint_migrations 等表——这也侧面印证了后文会讲到的完整表结构。

首次使用必须调用 .setup()

README 中用醒目的 IMPORTANT 标注:第一次使用 Postgres checkpointer 时,必须调用其 .setup() 方法创建必需的表。从 base.pyMIGRATIONS 列表可以看到 .setup() 具体做了哪些事情:

  1. checkpoint_migrations 表:记录迁移版本号 v(整数主键);
  2. checkpoints 表:主键 (thread_id, checkpoint_ns, checkpoint_id)checkpointmetadata 两列都是 JSONB,另有 parent_checkpoint_id 维护父子链;
  3. checkpoint_blobs 表:主键 (thread_id, checkpoint_ns, channel, version),存 BYTEA 类型的序列化通道值;
  4. checkpoint_writes 表:主键 (thread_id, checkpoint_ns, checkpoint_id, task_id, idx),存任务级中间写入;
  5. 三张表各自补建 thread_id 上的并发索引(CREATE INDEX CONCURRENTLY);
  6. checkpoint_writes 追加 task_path 列。

setup() 的执行逻辑(见 init.pyaio.py)是幂等的:先执行 MIGRATIONS[0] 建迁移表,然后读取 SELECT v FROM checkpoint_migrations ORDER BY v DESC LIMIT 1 得到当前版本(不存在则为 -1),再把列表中高版本号的迁移逐条执行并记录。这意味着升级包之后重新调用 setup() 即可安全地补齐新迁移,不会重复建表。

连接参数要求:autocommit 与 dict_row

README 中第二个 IMPORTANT 是本包最常见的踩坑点:当你手动创建 Postgres 连接并传给 PostgresSaverAsyncPostgresSaver 时,必须带上 autocommit=Truerow_factory=dict_rowfrom psycopg.rows import dict_row。原因在 README 中已给出解释,源码可以进一步印证:

  • autocommit=True.setup() 需要把 checkpoint 表直接提交到数据库。若连接默认处于隐式事务中且你不显式提交,建表可能不会被持久化;
  • row_factory=dict_rowPostgresSaver 的实现全程用字典语法访问行(如 row["v"]value["thread_id"])。默认的 tuple_row 只支持下标访问,一旦 checkpointer 按列名取值就会抛出 TypeError: tuple indices must be integers or slices, not str

README 给出的错误示例:

# ❌ This will fail with TypeError during checkpointer operations
with psycopg.connect(DB_URI) as conn:  # Missing autocommit=True and row_factory=dict_row
    checkpointer = PostgresSaver(conn)
    checkpointer.setup()  # May not persist tables properly
    # Any operation that reads from database will fail with:
    # TypeError: tuple indices must be integers or slices, not str

如果你不想手写这些参数,直接使用类方法 from_conn_string 即可,它在 init.py 内部已经替你配好了 autocommit=True, prepare_threshold=0, row_factory=dict_row

@classmethod
@contextmanager
def from_conn_string(cls, conn_string: str, *, pipeline: bool = False):
    ...
    with Connection.connect(
        conn_string, autocommit=True, prepare_threshold=0, row_factory=dict_row
    ) as conn:
        if pipeline:
            with conn.pipeline() as pipe:
                yield cls(conn, pipe)
        else:
            yield cls(conn)

注意构造器还有一个约束:若传入的是 ConnectionPool,则不能再同时使用 Pipeline__init__ 中会直接抛出 ValueError),pipeline 只能用于单条连接。

同步用法:PostgresSaver

下面是 README 中给出的完整同步示例,演示了存、取、列三个核心操作:

from langgraph.checkpoint.postgres import PostgresSaver

write_config = {"configurable": {"thread_id": "1", "checkpoint_ns": ""}}
read_config = {"configurable": {"thread_id": "1"}}

DB_URI = "postgres://postgres:postgres@localhost:5432/postgres?sslmode=disable"
with PostgresSaver.from_conn_string(DB_URI) as checkpointer:
    # call .setup() the first time you're using the checkpointer
    checkpointer.setup()
    checkpoint = {
        "v": 4,
        "ts": "2024-07-31T20:14:19.804150+00:00",
        "id": "1ef4f797-8335-6428-8001-8a1503f9b875",
        "channel_values": {"my_key": "meow", "node": "node"},
        "channel_versions": {"__start__": 2, "my_key": 3, "start:node": 3, "node": 3},
        "versions_seen": {
            "__input__": {},
            "__start__": {"__start__": 1},
            "node": {"start:node": 2},
        },
    }

    # store checkpoint
    checkpointer.put(write_config, checkpoint, {}, {})

    # load checkpoint
    checkpointer.get(read_config)

    # list checkpoints
    list(checkpointer.list(read_config))

结合 init.py 的源码,各方法的实际行为如下:

  • put(config, checkpoint, metadata, new_versions):先对 config["configurable"] 中的 thread_id / checkpoint_ns / checkpoint_id 解包,然后执行"值分拆"——channel_values 中为 Nonestr/int/float/bool 的原始值继续内联在 checkpoints.checkpoint JSONB 里,其余复杂值(包括 _DeltaSnapshot 标记类型)被弹出并写入 checkpoint_blobs 表;最后用 UPSERT_CHECKPOINTS_SQLON CONFLICT ... DO UPDATE)写主表。返回值是补全了 checkpoint_id 的新 config;
  • get_tuple(config):config 中带 checkpoint_id 时按主键精确取;不带时 ORDER BY checkpoint_id DESC LIMIT 1 取该 thread_id + checkpoint_ns 的最新 checkpoint。对老版本(v < 4)数据还会执行 pending sends 的在线迁移;
  • list(config, *, filter, before, limit):基于 _search_where 构造 WHERE 条件(支持按 metadata 过滤、按 before 指定时间点截断),按 checkpoint_id 降序返回 CheckpointTuple 迭代器;
  • put_writes(config, writes, task_id, task_path):把任务中间写入落到 checkpoint_writes;全部写入通道都在 WRITES_IDX_MAP 内时用 upsert 语义(可覆盖重试),否则用 INSERT ... ON CONFLICT DO NOTHING
  • delete_thread(thread_id):一次性删除 checkpointscheckpoint_blobscheckpoint_writes 三张表中该线程的所有数据。

内部游标管理 _cursor(pipeline=True)init.py)值得注意:写操作(put / put_writes / delete_thread)都会走 pipeline 模式批量提交;若当前 Psycopg 版本不支持 pipeline(Capabilities().has_pipeline() 为假),会自动退回到 conn.transaction() 事务上下文,保证正确性。

异步用法:AsyncPostgresSaver

异步入口在 aio.py,README 给出的示例:

from langgraph.checkpoint.postgres.aio import AsyncPostgresSaver

async with AsyncPostgresSaver.from_conn_string(DB_URI) as checkpointer:
    # 首次使用同样需要先 await checkpointer.setup()
    checkpoint = {
        "v": 4,
        "ts": "2024-07-31T20:14:19.804150+00:00",
        "id": "1ef4f797-8335-6428-8001-8a1503f9b875",
        "channel_values": {"my_key": "meow", "node": "node"},
        "channel_versions": {"__start__": 2, "my_key": 3, "start:node": 3, "node": 3},
        "versions_seen": {
            "__input__": {},
            "__start__": {"__start__": 1},
            "node": {"start:node": 2},
        },
    }

    # store checkpoint
    await checkpointer.aput(write_config, checkpoint, {}, {})

    # load checkpoint
    await checkpointer.aget(read_config)

    # list checkpoints
    [c async for c in checkpointer.alist(read_config)]

AsyncPostgresSaver.from_conn_string 内部使用 AsyncConnection.connect(conn_string, autocommit=True, prepare_threshold=0, row_factory=dict_row),参数要求与同步版一致;序列化这类 CPU 密集步骤(如 _dump_blobs_dump_writes_load_writes)被放到 asyncio.to_thread 中执行,避免阻塞事件循环。

一个容易忽略的细节:AsyncPostgresSaver 同时提供了同步外观(put / get_tuple / list / delete_thread)。从源码(aio.py)看,这些同步方法通过 asyncio.run_coroutine_threadsafe 把协程调度回创建实例时的事件循环执行;如果你在主线程(即运行着事件循环的那个线程)直接调用,会抛出 asyncio.InvalidStateError,提示改用 checkpointer.alist(...)await graph.ainvoke(...) 等异步接口。也就是说:异步实例的同步包装方法只能在后台线程中使用,主线程请一律走 a* 前缀接口。

序列化模型与查询路径

base.py 的 SQL 常量可以看到这套实现的查询策略:

  • SELECT_SQL:从 checkpoints 出发,用两个标量子查询分别按 channel_versions 关联 checkpoint_blobs(重建 channel_values)、按主键关联 checkpoint_writes(重建 pending_writes),一次查询拼出完整 CheckpointTuple 所需的行数据;
  • 三个 UPSERT_*_SQL 分别对应 blob、主表、writes 表的幂等写入,全部以 ON CONFLICT 处理并发重复写。

另外,get_delta_channel_history 为 delta 通道(增量存储的消息类通道)实现了两阶段快速查询:Stage 1 以 1024 行(_DELTA_PAGE_SIZE)为页大小沿父链分页扫描 checkpoints 元数据,Stage 2 用按通道拆分的 UNION ALL 一次性取回 writes 与种子 blob。base.py 中的注释还保留了作者的基准记录:该动态列设计相比"整包回传 JSONB 由 Python 挑选"的替代方案,端到端延迟约快 3 倍、网络载荷减小 61%。这是理解"为什么 checkpoint 表用 JSONB 内联、大值拆 blob"这一设计的直接依据。

安全加固:限制 checkpoint 反序列化

README 的 Security 一节要求:创建 checkpointer 时设置 LANGGRAPH_STRICT_MSGPACK=true 或显式传入 allowed_msgpack_modules 列表,把 checkpoint 反序列化限制到已知安全的类型,防止数据库被攻破时发生代码执行。这一机制实现在 langgraph-checkpoint 的 serde 层:libs/checkpoint/langgraph/checkpoint/serde/_msgpack.py 中通过环境变量开关 STRICT_MSGPACK_ENABLEDlibs/checkpoint/langgraph/checkpoint/serde/jsonplus.pyJsonPlusSerializer 构造函数接受 allowed_msgpack_modules 参数来决定允许反序列化的模块白名单。

实践建议是二选一:

  • 在部署环境统一设置 LANGGRAPH_STRICT_MSGPACK=true(全局生效);
  • 或仅在特定 checkpointer 上构造 JsonPlusSerializer(allowed_msgpack_modules=[...]) 并通过 serde= 参数传入 PostgresSaver / AsyncPostgresSaver 构造器(两者都接受 serde: SerializerProtocol | None)。

延伸阅读:ShallowPostgresSaver 与版本兼容

  • shallow.py 提供 ShallowPostgresSaver / AsyncShallowPostgresSavercheckpoints 表主键退化为 (thread_id, checkpoint_ns),即只保留每个命名空间的最新 checkpoint,适合只需要断点恢复、不需要完整历史回溯的场景;
  • base.py 顶部有一段版本兼容检查:检测到 langgraph 小于 0.5 时会发出 DeprecationWarning,提醒升级主库以避免意外行为;
  • 该包同时导出 Conn 别名(指向 _internal.Conn / _ainternal.Conn)用于向后兼容旧版类型注解。

验证你的环境

libs/checkpoint-postgres 目录下:

make start-postgres   # docker compose 拉起 pgvector/pgvector:pg16(端口 5441)
uv run pytest          # 运行 tests/ 下全部用例(含同步/异步、store、delta 迁移等)
make stop-postgres    # 用完后释放容器

tests/ 目录包含 test_sync.pytest_async.pytest_store.pytest_conformance_delta.py 等用例,tests/conftest.pyclear_test_db 夹具会在每个测试前清空全部 checkpoint 表,保证用例互不干扰。

小结

langgraph-checkpoint-postgres 是 LangGraph 生态中将 Agent 状态落到 Postgres 的官方实现:setup() 负责幂等的建表与迁移,put/get/list/put_writes/delete_thread 构成完整的状态生命周期,同步与异步两套接口共享同一套表结构与 SQL。落地时抓住三个要点即可——手动建连必须带 autocommit=Truerow_factory=dict_row、首次使用先 .setup()、生产环境开启 LANGGRAPH_STRICT_MSGPACK 或显式 msgpack 白名单——就可以为长时运行的工作流提供可靠、可回溯的持久化底座。

登录后查看全文
热门项目推荐
相关项目推荐

项目优选

收起
kernelkernel
deepin linux kernel
C
33
18
ops-transformerops-transformer
本项目是CANN提供的transformer类大模型算子库,实现网络在NPU上加速计算。
C++
1.12 K
2.72 K
kernelkernel
openEuler内核是openEuler操作系统的核心,既是系统性能与稳定性的基石,也是连接处理器、设备与服务的桥梁。
C
528
588
ops-nnops-nn
本项目是CANN提供的神经网络类计算算子库,实现网络在NPU上加速计算。
C++
906
1.82 K
pytorchpytorch
作为 Ascend for PyTorch 社区的核心组件,TorchNPU 是昇腾专为 PyTorch 打造的深度学习适配插件,使 PyTorch 框架能够直接调用昇腾 NPU,为开发者提供昇腾 AI 处理器的超强算力。
Python
854
1.34 K
docsdocs
暂无描述
Markdown
891
5.78 K
jiuwenswarmjiuwenswarm
JiuwenSwarm 是一款基于openJiuwen开发的智能AI Agent,它能够将大语言模型的强大能力,通过你日常使用的各类通讯应用,直接延伸至你的指尖。
Python
3.53 K
1.01 K
ops-mathops-math
本项目是CANN提供的数学类基础计算算子库,实现网络在NPU上加速计算。
C++
1.34 K
1.45 K
cann-learning-hubcann-learning-hub
CANN 学习中心仓,支持在线互动运行、边学边练,提供教程、示例与优化方案,一站式助力昇腾开发者快速上手。
Jupyter Notebook
987
504
AscendNPU-IRAscendNPU-IR
AscendNPU-IR是基于MLIR(Multi-Level Intermediate Representation)构建的,面向昇腾亲和算子编译时使用的中间表示,提供昇腾完备表达能力,通过编译优化提升昇腾AI处理器计算效率,支持通过生态框架使能昇腾AI处理器与深度调优
C++
540
384