LangGraph 的 Checkpointer 机制:给 Agent 加上断点续跑能力
当多 Agent 协作系统从简单的“一问一答”走向长链路的复杂任务(例如:全自动生成数十页行业研报、跨多系统的多步骤代码重构与自动化运维发布)时,执行时间往往长达数分钟甚至数小时。
在这漫长的流转过程中,任何现实世界的物理意外都可能发生:
- 服务器 Kubernetes Pod 被驱逐或突发 OOM 重启;
- 某个外部第三方 API 瞬时抖动返回 503 错误;
- 执行到关键敏感节点(例如向生产数据库执行 SQL UPDATE 或触发扣款接口),必须停下来等待管理员在前端点击“确认授权”。
如果 Agent 的状态全存放在内存变量里,一旦进程中断,前面跑了 10 分钟的所有中间思考、检索证据和推理成果将瞬间化为乌有,只能从头再来。
LangGraph 的Checkpointer(状态快照持久化)机制,正是为解决长链路 Agent 的**断点续跑(Fault-Tolerant Resumption)与人机协同中断(Human-in-the-loop)**而生的工业级架构基石。
Checkpointer 的底层状态快照原理
Checkpointer 的核心思想非常清晰:在状态图(StateGraph)每走完一个节点、发生一次状态转移时,自动将当前的完整 State、线程 ID(thread_id)以及当前节点的版本指纹序列化落盘到持久化存储中。
[Node A: 检索] ---> 写入快照 Checkpoint 1 (thread_id: 101, checkpoint_id: v1) | v [Node B: 质量评估] -> 写入快照 Checkpoint 2 (thread_id: 101, checkpoint_id: v2) | v [系统崩溃重启 / 人工中断] x [系统恢复] --------> 读取 Checkpoint 2 快照,直接从 Node C 恢复执行! | v [Node C: 生成报告] -> 写入快照 Checkpoint 3 (thread_id: 101, checkpoint_id: v3)每个 Checkpoint 包含以下核心元数据:
thread_id:标识某一个独立的用户会话或任务流水线实例;checkpoint_id:单调递增的时间戳或版本 UUID;channel_values:当前状态字典中所有字段的真实序列化数据;next_nodes:从当前快照出发,下一步应当被执行的目标节点集合。
生产级 PostgresSaver 持久化实战
在本地开发时可以使用内存MemorySaver或轻量SqliteSaver,但在生产分布式多副本容器环境下,必须使用支持高可用连接池的AsyncPostgresSaver:
import asyncio from typing import TypedDict, List from langgraph.graph import StateGraph, END from langgraph.checkpoint.postgres.aio import AsyncPostgresSaver from psycopg_pool import AsyncConnectionPool # 1. 定义业务状态 class ReportAgentState(TypedDict): topic: str outline: List[str] draft_sections: List[str] review_approved: bool final_report: str # 2. 节点逻辑定义 async def outline_node(state: ReportAgentState): print(">>> 正在生成大纲...") await asyncio.sleep(1) return {"outline": ["1. 架构总览", "2. 存储选型", "3. 压测数据"]} async def drafting_node(state: ReportAgentState): print(">>> 正在起草详细章节...") await asyncio.sleep(1) return {"draft_sections": ["详细章节正文内容..."]} async def human_approval_node(state: ReportAgentState): # 模拟人工介入审核节点 print(">>> 等待人工审核...") return state async def publish_node(state: ReportAgentState): print(">>> 审核通过,正式发布报告!") return {"final_report": "完整已发布研报"}组装带持久化与断点续跑的状态图
async def run_resumable_agent(): # 建立 PostgreSQL 连接池 db_uri = "postgresql://agent_user:password@pg-master.local:5432/agent_db" async with AsyncConnectionPool(conninfo=db_uri, max_size=20) as pool: # 初始化异步 Checkpointer checkpointer = AsyncPostgresSaver(pool) # 第一次启动需初始化数据库表结构(自动建表) await checkpointer.setup() # 构建图 workflow = StateGraph(ReportAgentState) workflow.add_node("outline", outline_node) workflow.add_node("drafting", drafting_node) workflow.add_node("approval", human_approval_node) workflow.add_node("publish", publish_node) workflow.set_entry_point("outline") workflow.add_edge("outline", "drafting") workflow.add_edge("drafting", "approval") workflow.add_edge("approval", "publish") workflow.add_edge("publish", END) # 关键配置:指定在 approval 节点前自动挂起等待人工介入 app = workflow.compile( checkpointer=checkpointer, interrupt_before=["approval"] ) # 唯一任务标识 config = {"configurable": {"thread_id": "report_task_20260901_001"}} # 第一阶段执行:生成大纲与草稿,随后在 approval 节点前自动挂起 print("=== 启动第一阶段任务 ===") async for event in app.astream({"topic": "向量数据库运维实践"}, config): print(event) # 此时任务安全停在 approval 节点前,哪怕重启服务,状态也完好保存在 Postgres 中 print("\n--- 任务已在 approval 节点前安全挂起 ---") # 模拟人工在管理后台审核完成,注入审核状态并唤醒继续执行 print("\n=== 管理员审批通过,唤醒继续执行 ===") # 更新状态字段 await app.aupdate_state(config, {"review_approved": True}, as_node="approval") # 传入 None 表示从上次中断的断点直接向下续跑 async for event in app.astream(None, config): print(event) # 执行流程 # asyncio.run(run_resumable_agent())Checkpointer 带来的架构跃迁
引入 Checkpointer 后,多 Agent 系统获得了三个质的飞跃:
- 零丢单的高可用韧性:服务随时被重启,只要重新拉起 Worker 传入相同的
thread_id,Agent 能分毫不差地从上一个成功节点的快照恢复执行; - 原生的人机协同(Human-in-the-loop):通过
interrupt_before与interrupt_after,可以轻松在任意业务节点插入人工审核流,管理员修改状态后即可一键恢复流转; - 可追溯的时间旅行(Time Travel)与回滚:通过查看 Postgres 中的历史快照,运维人员可以随意回放 Agent 在任意历史时刻的完整思维链,甚至可以修改历史节点的数据后开辟一条全新分支重新跑分支测试。
掌握了 Checkpointer,你的多 Agent 应用才算真正走出了玩具 Demo 阶段,具备了在企业严苛生产环境中长效、稳定运转的工业级硬实力。