LangGraph Checkpoint Postgres:基于 PostgreSQL 的 Agent 状态持久化实现与实战
本篇技术文章围绕 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.0与psycopg-pool>=3.2.0:Psycopg 3 数据库驱动及其连接池;orjson>=3.11.5:高性能 JSON 编解码。
包内除了 saver 还附带了一个跨线程长期记忆库的 Postgres 实现(PostgresStore 与 AsyncPostgresStore),本文聚焦 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,且每个测试前会清空 checkpoints、checkpoint_blobs、checkpoint_writes、checkpoint_migrations 等表——这也侧面印证了后文会讲到的完整表结构。
首次使用必须调用 .setup()
README 中用醒目的 IMPORTANT 标注:第一次使用 Postgres checkpointer 时,必须调用其 .setup() 方法创建必需的表。从 base.py 的 MIGRATIONS 列表可以看到 .setup() 具体做了哪些事情:
checkpoint_migrations表:记录迁移版本号v(整数主键);checkpoints表:主键(thread_id, checkpoint_ns, checkpoint_id),checkpoint与metadata两列都是 JSONB,另有parent_checkpoint_id维护父子链;checkpoint_blobs表:主键(thread_id, checkpoint_ns, channel, version),存BYTEA类型的序列化通道值;checkpoint_writes表:主键(thread_id, checkpoint_ns, checkpoint_id, task_id, idx),存任务级中间写入;- 三张表各自补建
thread_id上的并发索引(CREATE INDEX CONCURRENTLY); - 给
checkpoint_writes追加task_path列。
setup() 的执行逻辑(见 init.py 与 aio.py)是幂等的:先执行 MIGRATIONS[0] 建迁移表,然后读取 SELECT v FROM checkpoint_migrations ORDER BY v DESC LIMIT 1 得到当前版本(不存在则为 -1),再把列表中高版本号的迁移逐条执行并记录。这意味着升级包之后重新调用 setup() 即可安全地补齐新迁移,不会重复建表。
连接参数要求:autocommit 与 dict_row
README 中第二个 IMPORTANT 是本包最常见的踩坑点:当你手动创建 Postgres 连接并传给 PostgresSaver 或 AsyncPostgresSaver 时,必须带上 autocommit=True 和 row_factory=dict_row(from psycopg.rows import dict_row)。原因在 README 中已给出解释,源码可以进一步印证:
autocommit=True:.setup()需要把 checkpoint 表直接提交到数据库。若连接默认处于隐式事务中且你不显式提交,建表可能不会被持久化;row_factory=dict_row:PostgresSaver的实现全程用字典语法访问行(如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中为None或str/int/float/bool的原始值继续内联在checkpoints.checkpointJSONB 里,其余复杂值(包括_DeltaSnapshot标记类型)被弹出并写入checkpoint_blobs表;最后用UPSERT_CHECKPOINTS_SQL(ON 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):一次性删除checkpoints、checkpoint_blobs、checkpoint_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_ENABLED,libs/checkpoint/langgraph/checkpoint/serde/jsonplus.py 的 JsonPlusSerializer 构造函数接受 allowed_msgpack_modules 参数来决定允许反序列化的模块白名单。
实践建议是二选一:
- 在部署环境统一设置
LANGGRAPH_STRICT_MSGPACK=true(全局生效); - 或仅在特定 checkpointer 上构造
JsonPlusSerializer(allowed_msgpack_modules=[...])并通过serde=参数传入PostgresSaver/AsyncPostgresSaver构造器(两者都接受serde: SerializerProtocol | None)。
延伸阅读:ShallowPostgresSaver 与版本兼容
- shallow.py 提供
ShallowPostgresSaver/AsyncShallowPostgresSaver:checkpoints表主键退化为(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.py、test_async.py、test_store.py、test_conformance_delta.py 等用例,tests/conftest.py 的 clear_test_db 夹具会在每个测试前清空全部 checkpoint 表,保证用例互不干扰。
小结
langgraph-checkpoint-postgres 是 LangGraph 生态中将 Agent 状态落到 Postgres 的官方实现:setup() 负责幂等的建表与迁移,put/get/list/put_writes/delete_thread 构成完整的状态生命周期,同步与异步两套接口共享同一套表结构与 SQL。落地时抓住三个要点即可——手动建连必须带 autocommit=True 与 row_factory=dict_row、首次使用先 .setup()、生产环境开启 LANGGRAPH_STRICT_MSGPACK 或显式 msgpack 白名单——就可以为长时运行的工作流提供可靠、可回溯的持久化底座。
atomcodeClaude Code 的开源替代方案。连接任意大模型,编辑代码,运行命令,自动验证 — 全自动执行。用 Rust 构建,极致性能。 | An open-source alternative to Claude Code. Connect any LLM, edit code, run commands, and verify changes — autonomously. Built in Rust for speed. Get StartedRust0623
Hy4-previewHy4 preview 是由腾讯混元团队研发的新一代混合专家(MoE)旗舰模型。模型总参数量 770B,每个 token 激活 49B,主干共包含78层,第一层采用标准 FFN,其余 77 层均为 MoE 结构,每层包含 256 个路由专家与 1 个共享专家,每个 token 激活 top-8 路由专家及共享专家。主干之外原生内置 1 层 MTP(总参数量 10B,激活 0.7B)以支持投机解码。Python00
GLM-5.3GLM-5.3 与 GLM-5.2 使用相同的基座模型——所有提升均来自后训练。与 GLM-5.2 相比,它在复杂编程和长程任务上的表现显著提升。Jinja00
GLM-5.3-FlashGLM-5.3-Flash (320B-A18B),是GLM-5系列的首个原生多模态模型。320B总参数,能力超过GLM-5.2Jinja00
Spark-X2.5-4BSpark-X2.5-4B 旨在让强大的 AI 更实用、更高效、更易获得。在广泛日常任务中表现强劲,涵盖对话、写作、翻译、推理、编码、工具调用以及智能体工作流,并在同等规模的开源模型中取得领先成绩。Spark-X2.5 将面向效率的架构与最高 1M tokens 的原生上下文窗口相结合,并支持 200 多种语言。Python00
Spark-X2.5-1.7BSpark-X2.5-1.7B 旨在让强大的 AI 更加实用、高效且易于获取。这些模型在广泛的日常任务中表现出色,涵盖对话、写作、翻译、推理、编程、工具调用和智能体工作流,并在同等规模的开源模型中取得领先结果。Spark-X2.5 将面向效率的架构与最高 1M tokens 的原生上下文窗口相结合,并支持 200 多种语言。Python00