ARTICLE DETAIL

资讯详情

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

从手写Loop到生产级Agent:LangGraph+PostgreSQL Checkpoint+AG-UI中断恢复实战

从手写Loop到生产级Agent:LangGraph+PostgreSQL Checkpoint+AG-UI中断恢复实战 1. 为什么中断恢复是 Agent 工程绕不过去的一道坎做过 Agent 项目的人大概都有过这种体验一个跑了十几步的对话流程用户突然关掉页面或者服务端因为一次部署重启整个会话状态就没了。用户回来一看之前聊到哪儿、工具调到哪一步、中间产出了什么全部归零。这不是体验问题这是能不能上生产的问题。我最早做对话式 Agent 的时候用的是最朴素的手写 Loop一个while循环维护一个messages列表每轮把用户输入塞进去调一次模型判断有没有工具调用有就执行工具再把结果塞回去直到模型不再要求调工具为止。这套东西在 Demo 阶段非常好用几十行代码就能跑起来逻辑一目了然。但它有个致命缺陷——状态全在内存里。进程一挂messages列表就没了没有任何地方记录这个会话执行到哪一步了。后来我尝试自己往数据库里存messages每次循环结束就序列化一次。听起来简单真做起来一堆问题并发写怎么保证顺序工具执行到一半崩了重放的时候会不会重复执行副作用模型返回的中间态比如正在等待人工确认怎么表达这些问题一个个冒出来我才意识到可恢复不是加一个存盘动作那么简单它是一套运行时Runtime的设计问题。这篇内容就是把我从手写 Loop 迁移到 LangGraph PostgreSQL Checkpoint AG-UI 这套组合的完整过程拆开讲。核心要解决三件事状态怎么被结构化地持久化、中断之后怎么从正确的点恢复、前端怎么实时感知并驱动这个可恢复的运行时。适合已经写过基础 Agent Loop、准备往生产级方向走的人看。如果你还在纠结 LangChain 和 LangGraph 的区别简单说一句LangChain 更像工具箱LangGraph 更像状态机运行时前者给你零件后者给你一套能持久化、能中断、能恢复的执行框架。2. 手写 Loop 的天花板到底在哪里2.1 一个典型手写 Loop 的结构与它的隐含假设先把我早期那套 Loop 的骨架摆出来方便对照后面要讲的东西def run_agent(user_input, history): messages history [{role: user, content: user_input}] while True: response llm.invoke(messages) messages.append(response) if not response.tool_calls: return messages for call in response.tool_calls: result execute_tool(call.name, call.args) messages.append({role: tool, content: result})这段代码能跑但它建立在几个隐含假设上第一整个流程是单进程、单线程、短生命周期的第二工具执行是幂等且瞬时的第三用户不会在中途离开。只要这三个假设有一个不成立这套 Loop 就开始出问题。我踩得最狠的一次是工具执行超时。有个工具要调外部接口正常情况下两秒返回某次对方服务抖动卡了三十秒。用户等不及刷新了页面前端会话断了但后端那个while还在跑。等它终于拿到结果往一个已经没人听的会话里塞消息纯属浪费。更糟的是如果这个工具是有副作用的比如下单、发消息重试逻辑一旦没做好就会重复执行。2.2 状态、副作用、并发这三座山把问题抽象一下手写 Loop 撞墙的本质是三个东西没有被显式管理状态State。messages列表是一个隐式的、扁平的状态。它没有版本没有快照没有当前处于哪个阶段的概念。你没法回答这个会话现在是在等模型、等工具、还是等用户确认。副作用Side Effect。工具调用一旦发生就对外部世界产生了影响。手写 Loop 里工具执行和状态更新是耦合在一起的崩在中间就会出现工具执行了但状态没记或者状态记了但工具没执行的不一致。并发Concurrency。同一个会话如果被两个请求同时驱动比如用户开了两个标签页两个 Loop 会各自维护一份messages互相覆盖最后状态彻底乱掉。这三座山不解决Agent 就永远停在 Demo 阶段。而 LangGraph 的价值恰恰是把这三件事变成了框架层面的一等公民。2.3 从循环到图的思维转换手写 Loop 的思维是线性推进一轮接一轮直到结束。LangGraph 的思维是状态图你定义若干个节点Node节点之间用边Edge连接整个执行过程就是状态在图上流转。每个节点接收当前状态返回状态更新框架负责把更新合并进全局状态。这个转换一开始会让人觉得多此一举——明明一个循环能搞定的事为什么要画成图但当你需要在某个节点前暂停等用户输入某个节点失败后回到上一个节点重试多个分支根据条件走不同路径的时候图结构的表达力就体现出来了。更重要的是图结构天然适合做 Checkpoint每个节点执行完就是一个天然的存档点状态在这一刻是明确且完整的。3. LangGraph 的 StateGraph 与 Checkpoint 机制拆解3.1 StateGraph 的核心抽象State、Node、EdgeLangGraph 里最核心的三个概念我用一个实际例子说明。假设我们要做一个查资料 写摘要的 Agentfrom typing import TypedDict, Annotated from langgraph.graph import StateGraph, START, END from langgraph.graph.message import add_messages class AgentState(TypedDict): messages: Annotated[list, add_messages] research_done: bool summary: str def research_node(state: AgentState): # 调用搜索工具把结果塞进 messages result search_tool(state[messages][-1][content]) return {messages: [{role: tool, content: result}], research_done: True} def summarize_node(state: AgentState): summary llm.invoke(f总结以下内容{state[messages]}) return {summary: summary.content} graph StateGraph(AgentState) graph.add_node(research, research_node) graph.add_node(summarize, summarize_node) graph.add_edge(START, research) graph.add_edge(research, summarize) graph.add_edge(summarize, END) app graph.compile()这里有几个细节值得说。AgentState用TypedDict定义每个字段就是状态的一部分。Annotated[list, add_messages]这个写法是关键——它告诉 LangGraph这个字段的更新方式是追加而不是覆盖。add_messages是一个 reducer负责定义新值和旧值怎么合并。默认情况下字段更新是覆盖式的但messages这种列表显然需要追加所以要用 reducer 显式声明。节点函数接收完整状态返回一个字典表示我要更新哪些字段。注意是部分更新不需要返回整个状态。框架会拿你返回的字典按 reducer 规则合并进全局状态。这个设计让节点逻辑很干净不用关心状态的全貌。3.2 Checkpoint 到底存了什么Checkpoint 这个词容易让人误解以为只是存个messages。实际上 LangGraph 的 Checkpoint 存的是整个图状态在某个时刻的完整快照包括所有状态字段的当前值messages、research_done、summary等当前执行到哪个节点、下一步该走哪条边待处理的任务pending tasks比如某个节点还没执行完版本信息checkpoint_id、parent_checkpoint_id形成一条可追溯的链每次节点执行前后框架都会写一个 Checkpoint。这意味着你可以回到任意一个历史 Checkpoint从那里重新执行。这就是时间旅行能力的基础。Checkpoint 的存储由 Checkpointer 负责。LangGraph 提供了内存版MemorySaver只适合测试、SQLite 版、PostgreSQL 版等。生产环境我强烈建议用 PostgreSQL原因后面细讲。3.3 为什么是 PostgreSQL 而不是别的选 PostgreSQL 做 Checkpoint 存储不是因为它时髦而是几个硬需求它都能满足事务性。Checkpoint 的写入必须和业务操作在同一个事务语义下否则会出现状态存了但业务没做的不一致。PostgreSQL 的 ACID 是刚需。并发控制。同一个会话可能被多个请求驱动需要行级锁或者乐观锁来保证不会写冲突。PostgreSQL 的SELECT ... FOR UPDATE和 MVCC 直接可用。JSONB 支持。状态里经常有嵌套结构JSONB 既能存又能查还能建索引。LangGraph 的 PostgresSaver 内部就是用 JSONB 存状态快照的。成熟的运维生态。备份、监控、连接池这些不用自己造轮子。对比一下几种存储方案存储方案适用场景并发能力持久性我的建议MemorySaver本地调试无进程内只用于跑通逻辑SQLite单机小规模有限文件级个人项目可以PostgreSQL生产环境强强首选Redis高频读写缓存强可配置适合做二级缓存我自己的做法是 PostgreSQL 做主存储Redis 做热点会话的缓存层读的时候先查 Redis 再回源 PG。这个组合在会话量上来之后能明显降低 PG 的压力。4. 把 PostgreSQL Checkpoint 接进 LangGraph 的完整过程4.1 环境准备与依赖安装先把依赖装齐。这里要注意版本匹配LangGraph 迭代很快不同版本的 API 有差异pip install langgraph langgraph-checkpoint-postgres psycopg[binary] psycopg-poollanggraph-checkpoint-postgres是官方维护的 PostgreSQL Checkpointer 包底层用psycopg3。如果你还在用psycopg2需要换过来因为官方包是基于 psycopg3 写的。数据库这边建一个专用的库和用户别用超级用户跑应用CREATE DATABASE agent_runtime; CREATE USER agent_app WITH PASSWORD your_strong_password; GRANT ALL PRIVILEGES ON DATABASE agent_runtime TO agent_app;4.2 初始化 Checkpointer 与建表LangGraph 的 PostgresSaver 需要几张表来存 Checkpoint、写入记录和迁移版本。官方提供了setup()方法自动建表from langgraph.checkpoint.postgres import PostgresSaver from psycopg_pool import ConnectionPool DB_URI postgresql://agent_app:your_strong_passwordlocalhost:5432/agent_runtime pool ConnectionPool(conninfoDB_URI, max_size20, kwargs{autocommit: True}) checkpointer PostgresSaver(pool) checkpointer.setup() # 建表只需执行一次setup()会创建checkpoints、checkpoint_writes、checkpoint_migrations等表。生产环境我建议把setup()单独抽成一个初始化脚本在部署流程里跑一次而不是每次启动都调。虽然它内部有幂等判断但每次启动都连库建表不是好习惯。注意autocommitTrue这个参数在连接池配置里很关键。PostgresSaver 内部会自己管理事务边界如果连接不是 autocommit 模式会出现事务嵌套报错。这个坑我踩过报错信息很隐晦排查了半天。4.3 编译图时挂载 Checkpointer把 Checkpointer 挂到图上只需要在compile时传进去app graph.compile(checkpointercheckpointer)挂载之后每次执行图都要传一个config里面带上thread_id。thread_id就是会话的唯一标识同一个thread_id的所有执行共享一条 Checkpoint 链config {configurable: {thread_id: user-123-session-456}} result app.invoke({messages: [{role: user, content: 帮我查一下...}]}, config)这里有个设计要点thread_id的生成策略直接影响恢复能力。我一般用用户ID 会话ID组合会话ID由前端生成并持久化在本地。这样即使用户换了设备只要会话ID还在就能恢复。千万别用随机 UUID 每次请求都新生成那样等于没有恢复能力。4.4 验证 Checkpoint 真的写进去了跑完一次invoke去数据库里查一下SELECT thread_id, checkpoint_id, parent_checkpoint_id, jsonb_array_length(checkpoint-channel_values-messages) AS msg_count, created_at FROM checkpoints WHERE thread_id user-123-session-456 ORDER BY created_at DESC;如果能看到多条记录且parent_checkpoint_id形成链式关系说明 Checkpoint 正常工作了。msg_count能帮你快速确认消息有没有丢。我一般会写一个小的巡检脚本定期统计每个 thread 的 Checkpoint 数量异常增长比如某个会话几分钟内产生上千条往往意味着有死循环需要告警。5. 中断恢复的三种典型场景与实现5.1 场景一进程崩溃后的自动恢复这是最基础的场景。用户发了一条消息Agent 执行到第三步时服务重启了。重启后用户再发一条消息我们希望 Agent 能接着上次的状态继续而不是从头开始。LangGraph 的做法是用同一个thread_id再次invoke框架会自动从最新的 Checkpoint 加载状态然后继续执行。你不需要手动做什么只要保证thread_id一致。但这里有个细节如果上次崩溃时正好卡在某个节点执行到一半恢复时会怎么处理答案是——未完成的节点会被重新执行。这就是为什么工具必须幂等或者要有去重机制。我一般给每个工具调用生成一个call_id执行前先查一下这个call_id有没有执行记录有就直接返回缓存结果。5.2 场景二人工介入Human-in-the-loop有些操作需要人工确认比如是否真的要发送这封邮件。LangGraph 提供了interrupt机制可以在节点执行前暂停from langgraph.types import interrupt def send_email_node(state: AgentState): decision interrupt({question: 确认发送邮件, draft: state[draft]}) if decision approve: actually_send(state[draft]) return {status: sent} return {status: cancelled}调用interrupt时图会暂停状态被 Checkpoint 保存。前端拿到这个中断信号后展示确认界面。用户点了确认前端再调一次invoke传入恢复值from langgraph.types import Command app.invoke(Command(resumeapprove), config)框架会从暂停的 Checkpoint 恢复把approve作为interrupt的返回值继续执行。这个机制是可恢复 Runtime最有价值的体现之一——中断不是异常而是流程的正常组成部分。5.3 场景三跨设备、跨时间的会话续接用户在公司电脑上聊到一半回家用手机继续。这要求状态不仅持久化还要能被不同的客户端拉取。这时候 Checkpoint 就不只是恢复执行用了还要能读取当前状态给前端渲染。LangGraph 提供了get_state方法state app.get_state(config) print(state.values) # 当前所有状态字段 print(state.next) # 下一步要执行的节点前端拿到state.values里的messages就能把历史对话渲染出来。拿到state.next就知道当前是等待用户输入还是等待人工确认。这套东西配合 AG-UI 就能做出很自然的续接体验。6. AG-UI 如何把可恢复 Runtime 暴露给前端6.1 AG-UI 解决的是什么问题后端有了可恢复的 Runtime前端怎么知道现在该显示什么传统做法是前端自己维护一套状态和后端对不上就出 bug。AG-UI 的思路是后端把状态变化以事件流的形式推给前端前端只负责渲染。AG-UI 定义了一套标准的事件协议包括TEXT_MESSAGE_CONTENT文本增量、TOOL_CALL_START工具调用开始、STATE_SNAPSHOT状态快照、RUN_FINISHED执行结束等。后端把 LangGraph 的执行过程翻译成这些事件前端按事件更新 UI。6.2 把 LangGraph 的执行流翻译成 AG-UI 事件LangGraph 支持astream_events可以流式拿到执行过程中的各种事件。我们要做的是把这些事件映射到 AG-UI 协议async def run_agent_stream(thread_id, user_input): config {configurable: {thread_id: thread_id}} async for event in app.astream_events( {messages: [{role: user, content: user_input}]}, config, versionv2 ): kind event[event] if kind on_chat_model_stream: chunk event[data][chunk] if chunk.content: yield {type: TEXT_MESSAGE_CONTENT, delta: chunk.content} elif kind on_tool_start: yield {type: TOOL_CALL_START, tool: event[name], args: event[data].get(input)} elif kind on_chain_end and event[name] LangGraph: yield {type: RUN_FINISHED}这段代码是核心。astream_events会吐出模型流式输出、工具调用、链结束等各种事件我们按类型分发给前端。前端收到TEXT_MESSAGE_CONTENT就追加文字收到TOOL_CALL_START就显示正在调用 XX 工具。6.3 中断状态怎么传给前端当图因为interrupt暂停时astream_events会正常结束但状态里会有待处理的 interrupt。我们需要在流结束时检查一下state app.get_state(config) if state.next: # 有未完成的节点说明被中断了 interrupts state.tasks[0].interrupts if state.tasks else [] for itr in interrupts: yield {type: STATE_SNAPSHOT, state: {status: interrupted, question: itr.value}}前端收到status: interrupted就弹出确认框用户操作后再发起一次请求带上Command(resume...)。整个链路就闭环了。6.4 前端渲染的几个实操细节前端这边有几个坑我踩过。第一消息去重。流式输出时如果网络抖动导致重连可能会收到重复的增量。我的做法是给每条消息一个message_id前端按 ID 去重。第二工具调用的展示时机。工具调用可能耗时很长不能等它结束才显示。要在TOOL_CALL_START时就渲染一个进行中的状态TOOL_CALL_END时再更新为完成。第三恢复时的状态对齐。用户刷新页面后前端要先调一个get_state接口拿到完整状态再开始接收流。否则会出现流里只有新消息历史消息丢了的情况。7. 生产环境必须处理的几个硬骨头7.1 Checkpoint 表的膨胀与清理Checkpoint 是只增不减的一个活跃会话每天可能产生几百条。几个月下来checkpoints表能到千万级。我遇到过查询变慢、磁盘告警的情况。清理策略有两种。一种是按时间保留比如只保留最近 30 天的 Checkpoint更早的归档到冷存储。另一种是按会话保留每个 thread 只保留最近 N 条。LangGraph 提供了delete_thread方法可以删整个会话但细粒度的清理需要自己写 SQLDELETE FROM checkpoints WHERE thread_id IN ( SELECT thread_id FROM checkpoints GROUP BY thread_id HAVING MAX(created_at) NOW() - INTERVAL 30 days );我一般用定时任务每天凌晨跑一次配合VACUUM回收空间。注意别在业务高峰期跑大删除会锁表。7.2 并发写入的冲突处理同一个thread_id被两个请求同时驱动时会出现写冲突。LangGraph 的 PostgresSaver 内部用了乐观锁冲突时会抛异常。我的处理方式是在应用层加会话级锁同一个 thread 同时只允许一个执行流。简单做法是用 Redis 的SET NX做一个分布式锁lock_key fagent_lock:{thread_id} if not redis.set(lock_key, 1, nxTrue, ex300): raise Exception(该会话正在处理中请稍候) try: # 执行图 finally: redis.delete(lock_key)锁的过期时间要设得比最长执行时间略长防止死锁。同时要有兜底万一锁没释放得有机制能强制清理。7.3 状态迁移与版本兼容Agent 上线后状态结构可能会变——加字段、改字段名、调整 reducer。这时候老的 Checkpoint 用新代码加载会出问题。我的经验是状态结构变更必须做迁移。做法是在状态里加一个schema_version字段加载时检查版本不匹配就走迁移逻辑。迁移可以是懒加载读的时候转也可以是批量任务提前把所有老 Checkpoint 转掉。会话量大的话懒加载更现实。另外reducer 的变更要特别小心。比如messages从覆盖改成追加老数据加载后行为会变。这种变更我一般会新开一个字段而不是改老字段的语义。7.4 监控与可观测性可恢复 Runtime 的监控重点和普通服务不一样。我关注这几个指标指标含义告警阈值建议checkpoint_write_latencyCheckpoint 写入耗时P99 500msactive_threads活跃会话数突增 50%interrupted_threads处于中断状态的会话持续增长resume_failures恢复失败次数任意非零checkpoint_table_sizeCheckpoint 表大小超过磁盘 70%resume_failures这个指标最关键。恢复失败往往意味着状态损坏或者代码不兼容必须第一时间发现。我在恢复逻辑里加了 try-catch失败时上报指标并记录完整的thread_id和checkpoint_id方便事后排查。8. 我踩过的几个真实坑与排查过程8.1 恢复后消息重复一次 reducer 误用的排查有段时间用户反馈恢复会话后之前的消息出现两遍。我一开始怀疑是前端渲染问题查了半天前端代码没发现异常。后来去数据库里看 Checkpoint发现messages数组里确实有重复。根因是 reducer 配置错了。我有个自定义状态字段用了默认的覆盖式更新但节点返回时返回了完整的列表而不是增量导致每次更新都把整个列表又追加了一遍。改成add_messagesreducer 后正常。这个坑的教训是凡是列表类型的字段一定要显式声明 reducer。默认行为是覆盖但很多人包括当时的我会想当然以为列表是追加。8.2 工具重复执行幂等性缺失的代价前面提过工具幂等这里讲个具体的。有个创建工单的工具恢复时被重复执行结果创建了两个工单。排查发现是崩溃点正好在工具执行完但 Checkpoint 还没写的窗口期。解决方案是给工具加幂等键。每次工具调用生成一个基于thread_id 节点名 参数哈希的idempotency_key工具内部先查这个 key 有没有执行记录def create_ticket(args, idempotency_key): existing db.query(SELECT result FROM tool_executions WHERE key %s, idempotency_key) if existing: return existing[0][result] result do_create_ticket(args) db.execute(INSERT INTO tool_executions (key, result) VALUES (%s, %s), idempotency_key, result) return result这个模式对所有有副作用的工具都适用。别嫌麻烦重复下单、重复发消息的代价比多写几行代码大得多。8.3 连接池耗尽一个被忽视的配置上线初期遇到过一次服务假死所有请求都卡住。查日志发现是数据库连接池耗尽。原因是 PostgresSaver 的每个操作都要拿连接而我的连接池max_size设得太小只有 5并发一上来就不够用。调整到 20 之后缓解了但根本解法是控制并发。Agent 执行本身是长耗时操作如果每个请求都占着一个连接等模型返回连接池再大也不够。我的做法是把 Checkpoint 的读写和模型调用解耦——模型调用不占数据库连接只在需要读写状态时才拿连接。8.4 中断恢复后状态错乱一个 parent_checkpoint 的坑有次恢复后Agent 的行为完全不对像是回到了很早的状态。排查发现是parent_checkpoint_id链断了。原因是中间有一次手动清理数据删掉了某个 Checkpoint导致后续的链找不到父节点。LangGraph 恢复时是沿着parent_checkpoint_id往上找的链断了就会回退到能找到的最早节点。所以清理 Checkpoint 时不能随便删中间的要么删整个 thread要么从最早的开始删。这个约束我在清理脚本里加了检查删除前先确认不会破坏链的完整性。9. 从这套架构里我总结出的几条经验第一状态设计要面向恢复。别把状态当成临时变量要当成需要持久化的业务数据来设计。每个字段都要能回答恢复时它应该是什么值。第二副作用和状态更新要分离。工具执行归工具执行状态更新归状态更新两者通过幂等键关联。这样即使中间崩了重放也不会出问题。第三中断是常态不是异常。把人工确认、等待外部事件这些都建模成正常的中断点而不是用异常处理。LangGraph 的interrupt机制就是为这个设计的。第四监控要覆盖恢复路径。正常路径跑通不代表恢复路径没问题。恢复失败、状态不一致这些指标必须单独监控因为它们往往在出问题后才被发现。第五别过早优化。我一开始想自己实现一套 Checkpoint 存储折腾了两周发现还不如直接用官方的 PostgresSaver。框架已经解决的问题就别重复造轮子了。这套 LangGraph PostgreSQL Checkpoint AG-UI 的组合我从去年用到现在支撑了日均几万次会话恢复成功率稳定在 99.9% 以上。中间踩的坑基本都在上面了希望对正在做类似事情的人有点帮助。如果你的场景里工具副作用特别重建议在幂等这块多花点心思这是最容易出事的地方。
返回列表