拓十年匠心定制 · 商业建站与技术教学双线并行 咨询热线:400-886-1026 service@lmnt.cn
ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

LangGraph + PostgreSQL Checkpoint:构建可中断恢复的 Agent Runtime

LangGraph + PostgreSQL Checkpoint:构建可中断恢复的 Agent Runtime 1. 为什么手写 Loop 撑不过第三个需求我最早做对话式 Agent 的时候和大多数人一样写了一个while True循环调模型、解析工具调用、执行工具、把结果塞回消息列表、再调模型直到模型不再要求调用工具为止。这套东西跑 Demo 特别爽二十行代码就能让模型查天气、算数学、读文件。但只要产品经理提三个需求它立刻崩盘。第一个需求是用户点停止按钮要能中断。手写 Loop 里模型调用是阻塞的工具执行也是阻塞的你想中断只能靠抛异常或者杀线程状态全丢。第二个需求是中断之后要能接着聊。用户说刚才那个查询继续你的 Loop 早就退出了上下文栈、已执行到哪一步、中间结果是什么全在内存里进程一重启就没了。第三个需求是人工审批高风险操作。比如 Agent 要删数据库记录你得先暂停、等人点确认、再继续。手写 Loop 里这个暂停根本没有落脚点你只能把整个函数拆成状态机拆到最后代码变成一坨意大利面。这三个需求的本质是同一个东西执行过程需要可持久化、可中断、可恢复。手写 Loop 把执行状态和函数调用栈绑死了而函数调用栈是没法序列化到数据库的。所以正确的做法不是把 Loop 写得更复杂而是把 Loop 换成一个显式的状态图每一步的状态都能存下来每一步的跳转都由数据决定而不是由代码位置决定。这就是 LangGraph 要解决的问题。它把 Agent 的执行建模成一张有向图节点是计算步骤边是跳转条件整张图的当前状态是一个可序列化的字典。状态一旦可序列化就能存进 PostgreSQL一旦能存就能中断一旦能存能读就能恢复。再配上 AG-UI 这类前端协议把图执行到哪了、要不要等人实时推给界面一个真正可用的 Runtime 就成型了。这篇东西我打算按我自己踩坑的顺序讲先讲清楚 LangGraph 的状态模型到底和手写 Loop 差在哪再讲 PostgreSQL Checkpoint 是怎么把状态落地的、有哪些坑然后讲中断恢复的完整链路怎么跑通最后讲 AG-UI 怎么把这一切暴露给前端。全程给可复现的代码和参数不讲空话。提示本文假设你已经能跑通一个最基础的 LangGraph 图。如果你还没接触过 LangGraph建议先看官方文档里 StateGraph 和 Checkpointer 两节把StateGraph、compile、invoke这三个概念过一遍再回来不然下面的内容会有点跳。2. LangGraph 的状态模型把执行到哪变成数据2.1 手写 Loop 的状态藏在调用栈里LangGraph 把它摊平成一个字典先看手写 Loop 的典型结构。你有一个messages列表循环里判断最后一条消息有没有tool_calls有就执行工具、追加ToolMessage没有就退出。这个messages列表就是全部状态。听起来挺干净问题在于循环的进度不在这个列表里而在代码执行到第几行。你没法从messages反推出现在是在等工具返回还是在等模型返回还是已经结束了。LangGraph 的做法是引入一个显式的 State。用TypedDict定义比如from typing import Annotated, TypedDict from langgraph.graph.message import add_messages class AgentState(TypedDict): messages: Annotated[list, add_messages] pending_action: str | None retry_count: int关键在Annotated[list, add_messages]这个写法。add_messages是一个 reducer它规定了当多个节点都想更新messages时怎么合并。默认行为是覆盖而add_messages是追加并按 ID 去重。这个设计非常重要因为 LangGraph 的节点是返回增量而不是返回全量的节点只返回自己新增的那几条消息框架负责用 reducer 合并进全局状态。我一开始不理解为什么要这么设计觉得直接返回全量列表不就行了。后来做并行分支才明白如果两个节点并行执行、各自返回全量messages合并时就会互相覆盖。有了 reducer两个分支各自追加自己的消息合并逻辑由 reducer 统一处理不会丢数据。这是手写 Loop 根本做不到的——你没法让两个while循环并行跑还共享一个列表。2.2 节点返回增量、边决定跳转执行进度第一次变成可观测的LangGraph 里一个节点就是一个普通函数签名是state - dict返回的是要合并进状态的增量。边分两种普通边add_edge(a, b)表示 a 执行完无条件去 b条件边add_conditional_edges(a, router, {...})表示 a 执行完后调用router(state)根据返回值决定去哪个节点。这套东西的价值在于执行进度变成了图上的一个位置。当前在哪个节点、下一步去哪完全由状态和路由函数决定和代码执行到第几行无关。这就意味着只要我把状态存下来下次读出来从当前节点继续跑就行不需要重放整个调用栈。我举个具体例子对比。手写 Loop 里要不要继续调模型是靠if last_message.tool_calls:判断的这个判断发生在循环体内部。LangGraph 里这个判断变成一个条件边def should_continue(state: AgentState) - str: last state[messages][-1] if getattr(last, tool_calls, None): return tools return end graph.add_conditional_edges(agent, should_continue, { tools: tools, end: END, })看起来只是换了个写法但本质变了这个判断现在是纯函数输入是状态输出是字符串。纯函数意味着可测试、可重放、可序列化。手写 Loop 里那个if依赖的是运行时的局部变量你没法在进程外复现它。2.3 状态可序列化是中断恢复的前置条件不是可选项很多人做 Agent 时觉得状态存不存无所谓反正跑得快。这个想法在单轮对话里没问题但只要涉及中断恢复状态必须能序列化。LangGraph 的状态默认要求是 JSON 可序列化的messages里的消息对象会被序列化成特定格式。这里有个坑我踩过如果你在状态里塞了不可序列化的东西比如数据库连接、文件句柄、自定义的不可 pickle 对象compile的时候不报错但一旦启用 Checkpointer 落盘就会炸。我当时的做法是在状态里只放数据把资源通过闭包或者依赖注入传给节点函数节点内部用完即弃。比如数据库连接我在节点函数里现开现关或者用一个全局的连接池绝不放进 State。注意State 里只放数据不放资源。连接、句柄、锁这类东西一旦进了 State序列化必炸而且报错信息往往很隐晦排查起来很痛苦。3. PostgreSQL Checkpoint状态落地的真实细节3.1 Checkpointer 到底存了什么为什么选 PostgreSQLLangGraph 的 Checkpointer 在每个超级步super-step结束后把当前状态快照写进存储。所谓超级步可以理解成一轮节点执行。比如 agent 节点跑完是一个超级步tools 节点跑完是另一个超级步。每个快照带一个checkpoint_id并且通过parent_checkpoint_id串成一条链形成完整的历史。存储后端有好几种内存版MemorySaver适合测试SQLite 适合单机PostgreSQL 适合生产。为什么生产选 PostgreSQL三个原因。第一它的事务和并发控制成熟多个请求同时读写不同 thread 的状态不会互相干扰。第二它支持 JSONB 类型状态快照直接存成 JSONB查询和索引都方便。第三运维体系成熟备份、监控、扩容都有现成方案。内存版的问题很直接进程一重启全没多实例部署时各存各的用户请求打到哪个实例就丢状态。SQLite 的问题是并发写会锁库Agent 场景下写频率不低容易成为瓶颈。所以只要上生产PostgreSQL 基本是默认选择。3.2 建表、连接池与序列化器的配置细节LangGraph 的 PostgreSQL Checkpointer 需要一张表来存 checkpoint。用官方提供的PostgresSaver时它会自带一个setup()方法帮你建表from langgraph.checkpoint.postgres import PostgresSaver DB_URI postgresql://user:passlocalhost:5432/agentdb with PostgresSaver.from_conn_string(DB_URI) as checkpointer: checkpointer.setup() graph builder.compile(checkpointercheckpointer)setup()会创建checkpoints、checkpoint_blobs、checkpoint_writes等表。我建议第一次跑的时候手动执行一次setup()之后在生产环境用迁移脚本管理不要每次启动都调setup()虽然它是幂等的但启动时多一次 DDL 检查没必要。连接池这块有个坑。from_conn_string每次会新建连接在高并发下会打爆数据库连接数。生产环境我一般自己建一个ConnectionPool然后传给PostgresSaverfrom psycopg_pool import ConnectionPool from langgraph.checkpoint.postgres import PostgresSaver pool ConnectionPool(conninfoDB_URI, max_size20, kwargs{autocommit: True}) checkpointer PostgresSaver(pool)注意autocommitTrue这个参数。Checkpointer 内部会自己管理事务边界如果你在外面又包一层事务容易出现事务里嵌套事务的问题表现为写入不生效或者死锁。我一开始没加这个参数调试了半天发现 checkpoint 写进去了但读不出来就是因为外层事务没提交。序列化器方面LangGraph 默认用JsonPlusSerializer它能处理大部分常见类型包括 LangChain 的消息对象、datetime、UUID 等。如果你有自定义类型需要注册对应的序列化逻辑否则落盘时会报TypeError。我的经验是尽量别在 State 里放自定义类型能转成 dict 就转成 dict省心。3.3 thread_id 是恢复的钥匙命名策略直接影响可维护性每个 checkpoint 都归属于一个thread_id。恢复的时候你传入同一个thread_idCheckpointer 就会把最新的状态读出来从那里继续。所以thread_id是整个恢复机制的钥匙。命名策略上我踩过坑。一开始我用自增 ID 当thread_id后来发现没法区分这是哪个用户的哪个会话。改成f{user_id}:{session_id}之后清晰多了。再后来做多租户又加上租户前缀f{tenant}:{user}:{session}。这个格式没有标准答案但原则是能从 thread_id 反推出归属方便排查和清理。清理也很重要。checkpoint 会无限增长一个活跃会话可能产生几百个快照。我一般写一个定时任务删除created_at超过 N 天且会话已结束的记录。注意别删正在用的判断依据是会话状态是否为终态。DELETE FROM checkpoints WHERE thread_id IN ( SELECT thread_id FROM checkpoints GROUP BY thread_id HAVING MAX(created_at) NOW() - INTERVAL 30 days );提示删 checkpoint 前一定要确认对应会话已经结束。如果误删了活跃会话的历史用户下次恢复时会从更早的快照开始表现为丢失了最近几轮对话这种 bug 很难查。4. 中断恢复的完整链路从 interrupt 到 resume4.1 interrupt 不是抛异常而是把暂停点写进状态LangGraph 的中断机制核心是interrupt()函数。在节点里调用它图的执行会在这里暂停当前状态被 checkpoint 保存控制权返回给调用方。等调用方决定继续时用Command(resume...)把值传回去图从暂停点继续执行。这里最容易误解的一点是interrupt()看起来像抛异常但它不是。它不会让节点函数退出后丢失局部变量而是把暂停在哪个节点、暂停时传了什么值记录进 checkpoint。恢复时节点函数会从头重新执行但interrupt()这次会直接返回你传入的 resume 值而不是再次暂停。这个重新执行的行为我第一次遇到时很困惑。假设节点里有副作用比如发了一封邮件然后才调用interrupt()那恢复时邮件会再发一次。所以副作用必须放在 interrupt 之后或者做成幂等的。我一般的做法是把 interrupt 放在节点最前面先拿到人工决策再执行有副作用的逻辑。def approval_node(state: AgentState): decision interrupt({ question: 是否执行删除操作, action: state[pending_action], }) if decision approve: return {messages: [AIMessage(已批准继续执行)]} return {messages: [AIMessage(已拒绝终止操作)]}4.2 resume 的两种触发方式同步 invoke 和异步事件恢复有两种典型场景。第一种是同步的调用方拿到 interrupt 信号后等用户操作然后调graph.invoke(Command(resumevalue), config)。第二种是异步的用户操作通过消息队列或者 HTTP 接口进来服务端收到后触发恢复。同步方式简单直接适合暂停后立刻等人的场景。异步方式适合暂停后可能等几分钟甚至几小时的场景比如审批流。异步方式的关键是恢复时要用同一个 thread_id 和 checkpointer否则找不到暂停点。config {configurable: {thread_id: user1:session1}} # 第一次执行会在 approval_node 暂停 result graph.invoke({messages: [HumanMessage(删除用户 123)]}, config) # 用户点了批准恢复执行 result graph.invoke(Command(resumeapprove), config)我实测下来异步恢复最容易出问题的地方是并发恢复。如果用户手快点了两次批准两个恢复请求同时进来可能都读到同一个暂停点导致节点执行两次。解决办法是在业务层加锁或者用数据库的行锁保证同一个 thread 同时只有一个恢复在执行。LangGraph 本身不负责这个得自己兜。4.3 恢复后状态从哪读、怎么验证恢复正确恢复时Checkpointer 会根据thread_id找到最新的 checkpoint把状态加载出来然后从暂停的节点继续。你可以通过graph.get_state(config)查看当前状态确认恢复点是否正确state graph.get_state(config) print(state.next) # 下一个要执行的节点 print(state.values[messages][-1]) # 最后一条消息state.next是个元组表示接下来要执行的节点。如果它是空的说明图已经执行到 END 了。如果它包含你暂停的那个节点说明恢复点正确。验证恢复正确性我一般做三件事。第一检查state.next是否符合预期。第二检查state.values里的关键字段有没有丢。第三跑一遍完整流程看最终结果和不中断直接跑是否一致。第三点最重要因为中断恢复最容易出的 bug 就是恢复后少执行了一步或者多执行了一步只有端到端对比才能发现。注意get_state返回的是快照不是实时状态。如果你在恢复过程中调用它拿到的可能是中间态。要拿最终态等invoke返回后再调。5. AG-UI 把 Runtime 暴露给前端协议层的取舍5.1 为什么不用裸 WebSocketAG-UI 解决了什么后端 Runtime 跑通了前端怎么接最直接的想法是开个 WebSocket后端把状态变化推过去。但真做起来你会发现一堆问题消息格式怎么定中断信号怎么表达恢复请求怎么发多个 Agent 事件怎么区分这些如果自己定每个项目都要重新设计一遍而且前后端容易对不齐。AG-UI 是一套面向 Agent 前端的协议它规定了事件类型和消息格式。核心事件包括TEXT_MESSAGE_CONTENT流式文本、TOOL_CALL_START/TOOL_CALL_END工具调用、STATE_SNAPSHOT状态快照、RUN_FINISHED执行结束等。前端按事件类型渲染后端按事件类型推送双方解耦。它相比裸 WebSocket 的价值在于标准化。你不用再纠结中断信号用什么字段协议里已经定义好了。前端也有现成的组件库可以对接省掉大量联调时间。当然代价是要接受它的抽象有些自定义需求得绕一下。5.2 把 LangGraph 的事件流翻译成 AG-UI 事件LangGraph 执行时可以通过astream_events拿到细粒度事件包括节点开始、节点结束、LLM token 流等。把这些事件映射到 AG-UI 事件就是适配层的核心工作。async for event in graph.astream_events(input, config, versionv2): kind event[event] if kind on_chat_model_stream: chunk event[data][chunk] if chunk.content: yield AGUIEvent.text_message_content(chunk.content) elif kind on_chain_start and event[name] approval_node: yield AGUIEvent.state_snapshot({status: waiting_approval})映射的时候有个细节LangGraph 的事件粒度比 AG-UI 细很多内部事件前端不关心需要过滤。我一般只透传三类LLM 的 token 流、工具调用的开始和结束、以及自定义的中断信号。其他内部事件要么忽略要么聚合成状态快照再推。中断信号的处理是重点。当interrupt()触发时适配层要推一个明确的事件告诉前端现在需要用户输入并带上 interrupt 的 payload。前端渲染成审批按钮或者输入框用户操作后通过另一个接口把 resume 值发回后端。5.3 前端拿到中断信号后恢复请求怎么发回后端前端收到中断事件后展示对应的 UI。用户操作完前端发一个恢复请求后端用Command(resumevalue)触发恢复。这个请求要带上thread_id后端才能找到对应的暂停点。// 前端伪代码 const resumeRun async (threadId, value) { await fetch(/api/resume, { method: POST, body: JSON.stringify({ thread_id: threadId, resume: value }), }); };后端收到后app.post(/api/resume) async def resume(req: ResumeRequest): config {configurable: {thread_id: req.thread_id}} async for event in graph.astream(Command(resumereq.resume), config): yield to_agui(event)这里有个坑恢复请求和原始执行请求可能打到不同的后端实例。如果 checkpointer 是 PostgreSQL这没问题任何实例都能读到状态。但如果用了内存版 checkpointer就会找不到暂停点。所以只要涉及中断恢复checkpointer 必须是共享存储这是硬性要求。6. 实测中那些文档不会写的坑6.1 checkpoint 写入频率过高导致的性能问题LangGraph 默认在每个超级步结束后写一次 checkpoint。一个复杂的 Agent 一轮对话可能产生十几个超级步每个都写一次数据库。单用户看不出问题几百个并发用户时数据库写入压力就上来了。我的优化思路有两个。第一评估是否真的需要每个超级步都持久化。如果只是为了防止进程崩溃丢状态可以接受稍微粗一点的粒度。LangGraph 支持配置 checkpoint 的写入策略具体参数看版本不同版本 API 有差异。第二把 checkpoint 表和其他业务表分到不同的数据库实例或者不同的表空间避免互相影响。还有一个更隐蔽的问题checkpoint 的 payload 可能很大。如果 State 里存了长文档、大 JSON每次快照都写一遍存储和 IO 都会爆。我的做法是把大对象存到对象存储State 里只存引用比如 URL 或 ID需要时再取。6.2 恢复时状态不一致的三种典型表现第一种表现是消息重复。恢复后模型收到了重复的历史消息导致回答错乱。原因通常是 reducer 配置不对或者恢复时手动往 State 里塞了已经存在的消息。排查方法是打印恢复前后的messages列表对比长度和内容。第二种表现是工具调用丢失。中断发生在工具调用之后、结果写回之前恢复时工具结果没被正确合并。这通常是因为工具节点返回的增量格式不对或者 reducer 没处理这种类型。检查工具节点的返回值确保它返回的是{messages: [ToolMessage(...)]}这种标准格式。第三种表现是路由走错分支。恢复后条件边判断出错走了不该走的分支。原因往往是 State 里某个字段在恢复后值不对导致路由函数返回了错误的结果。排查方法是把路由函数的输入输出打日志对比中断前后的差异。6.3 多实例部署下 thread 锁与并发恢复的冲突多实例部署时同一个 thread 的恢复请求可能同时打到两个实例。如果两个实例都读到同一个暂停点都执行恢复就会导致节点执行两次。对于有副作用的节点这是灾难性的。解决办法是在恢复入口加分布式锁锁的 key 是thread_id。用 Redis 的SET NX或者 PostgreSQL 的 advisory lock 都行。拿到锁的实例执行恢复没拿到的返回正在处理中。锁的过期时间要设得比正常恢复耗时长一些避免锁提前释放。# 用 PostgreSQL advisory lock 示例 with conn.cursor() as cur: cur.execute(SELECT pg_try_advisory_lock(%s), (hash(thread_id),)) locked cur.fetchone()[0] if not locked: return {error: concurrent resume in progress} try: # 执行恢复 ... finally: cur.execute(SELECT pg_advisory_unlock(%s), (hash(thread_id),))这个锁不是 LangGraph 提供的得自己在业务层实现。我踩过一次坑没加锁用户双击批准按钮结果删了两次数据。从那以后所有涉及副作用的恢复入口都加了锁。6.4 版本升级时 checkpoint 格式兼容性LangGraph 迭代很快checkpoint 的内部格式在不同版本间可能有变化。升级版本后旧版本写的 checkpoint 可能读不出来或者读出来字段对不上。我的做法是升级前先在测试环境用生产数据的副本验证确认旧 checkpoint 能正常读取和恢复。如果不行要么写迁移脚本转换格式要么接受旧会话无法恢复并在升级公告里说明。生产环境升级时我一般会保留旧版本的读路径一段时间新写入用新格式读取时根据版本号走不同逻辑等旧会话自然过期后再清理。提示checkpoint 格式兼容性是升级时最容易忽略的点。建议在 CI 里加一个测试用固定版本的 checkpoint 数据验证新版本能否正确恢复防止升级引入回归。7. 写在最后的一点个人体会这套东西我从手写 Loop 一路踩过来最大的感受是Agent 的可靠性不取决于模型多强而取决于执行过程多可控。模型再聪明如果中断后恢复不了、状态对不上、并发时互相踩产品就没法用。LangGraph 加 PostgreSQL Checkpoint 加 AG-UI 这套组合本质上是在给 Agent 补上操作系统该有的能力进程调度、状态持久化、中断处理、进程间通信。如果你现在还在用手写 Loop我的建议是别急着全量重构。先把手写 Loop 里最痛的那个点——通常是中断恢复——单独抽出来用 LangGraph 重写那一个流程跑通之后再逐步迁移。全量重写风险太大而且很多简单场景手写 Loop 反而更直接。工具是拿来解决问题的不是拿来炫技的。最后分享一个我调试中断恢复的小技巧在开发环境把 checkpointer 换成MemorySaver同时在每个节点入口打印state的摘要。这样你能清楚地看到每次恢复时状态是怎么变的比直接看数据库里的 JSONB 直观得多。等逻辑跑通了再切回 PostgreSQL 验证持久化能省掉大量排查时间。
返回列表