ARTICLE DETAIL

资讯详情

深耕编程入门与网站建设的一线实战洞察。

LangGraph 的 Checkpointer 机制:给 Agent 加上断点续跑能力

LangGraph 的 Checkpointer 机制:给 Agent 加上断点续跑能力 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、线程 IDthread_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单调递增的时间戳或版本 UUIDchannel_values当前状态字典中所有字段的真实序列化数据next_nodes从当前快照出发下一步应当被执行的目标节点集合。生产级 PostgresSaver 持久化实战在本地开发时可以使用内存MemorySaver或轻量SqliteSaver但在生产分布式多副本容器环境下必须使用支持高可用连接池的AsyncPostgresSaverimport 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:passwordpg-master.local:5432/agent_db async with AsyncConnectionPool(conninfodb_uri, max_size20) 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( checkpointercheckpointer, 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_nodeapproval) # 传入 None 表示从上次中断的断点直接向下续跑 async for event in app.astream(None, config): print(event) # 执行流程 # asyncio.run(run_resumable_agent())Checkpointer 带来的架构跃迁引入 Checkpointer 后多 Agent 系统获得了三个质的飞跃零丢单的高可用韧性服务随时被重启只要重新拉起 Worker 传入相同的thread_idAgent 能分毫不差地从上一个成功节点的快照恢复执行原生的人机协同Human-in-the-loop通过interrupt_before与interrupt_after可以轻松在任意业务节点插入人工审核流管理员修改状态后即可一键恢复流转可追溯的时间旅行Time Travel与回滚通过查看 Postgres 中的历史快照运维人员可以随意回放 Agent 在任意历史时刻的完整思维链甚至可以修改历史节点的数据后开辟一条全新分支重新跑分支测试。掌握了 Checkpointer你的多 Agent 应用才算真正走出了玩具 Demo 阶段具备了在企业严苛生产环境中长效、稳定运转的工业级硬实力。
返回列表