ARTICLE DETAIL

资讯详情

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

从手写Loop到可恢复Runtime:LangGraph+PostgreSQL Checkpoint+AG-UI实战

从手写Loop到可恢复Runtime:LangGraph+PostgreSQL Checkpoint+AG-UI实战 上个月我把负责的一个 agent 服务整层重写了。这事得从一次客户反馈说起他们要求服务半夜崩溃了第二天能不能自己接着跑不用重头再来。听起来很简单但真上手做的时候我那套手写 while 循环完全扛不住。原因很直白——没有持久化、没有状态机、没有协议层。后来我用 LangGraph 做编排用 PostgreSQL Checkpoint 做状态持久化再用 AG-UI 统一了 Runtime 和前端之间的交互协议算是把中断恢复这件事从口号变成了能落地的能力。这篇文章就把这次改造的完整思路、关键代码和踩过的坑记录下来给同样在手写 Loop 和可恢复 Runtime 之间纠结的人一个参考。1. 为什么我放弃了手写 Loop1.1 手写循环的三宗罪先说说我之前是怎么写的。几乎所有 agent 的雏形都是同一个套路把消息列表塞给大模型如果返回了工具调用就执行工具再把结果追加回去循环直到模型不再调用工具为止。核心代码大概长这样messages [{role: user, content: 帮我查下这个月的订单量}] while True: response llm.invoke(messages) messages.append(response) if not response.tool_calls: break for tc in response.tool_calls: tool_result call_tool(tc[name], tc[args]) messages.append({ role: tool, tool_call_id: tc[id], content: tool_result })这段代码跑通 demo 只要一个下午但真正做产品的时候它会从三个方向同时恶心你。第一宗罪进程一崩上下文归零。整个对话状态存在内存里挂在 Python 进程的局部变量上。进程挂掉、服务器重启、发布时滚动更新所有东西瞬间蒸发。客户半夜问我上周让 agent 查的那批数据后来怎么了我什么都答不上来。第二宗罪没法优雅地暂停。工具执行经常需要人工审批比如转账、发邮件、删除数据。手写循环碰到这种情况要么让模型自己决定要么我在循环里塞一个input()等用户输入。前者不可控后者直接卡死整个服务。真实业务里人工确认可能要等十分钟甚至隔天一个循环不可能干等那么久。第三宗罪没有外界可以观测和干预的入口。前端想要展示agent 正在调数据库agent 正在等审批我就得在循环里到处埋点、写回调、维护一套自己发明的轮询接口。前端要改一次交互我要改一整轮代码。这三宗罪本质上是同一个问题我把一次会话当成了一个驻留在内存里的活对象但它实际上应该是一个可以随时暂停、保存、恢复、检查的独立状态机。1.2 LangGraph 到底改变了什么LangGraph 解决的就是这个把 agent 变成状态机的问题。它的核心概念不多一个 StateGraph你可以往里面注册节点node和边edge节点就是一段处理逻辑边决定下一步走哪个节点中间流转的数据是一个结构化的 State。我打个比方手写 Loop 像是你一个人拿着流程图在脑子里执行每一步走到哪完全靠你自己记。LangGraph 则是把流程图本身变成可执行文件每一步走到哪、当前数据是什么、已经执行了哪些步骤框架都会显式地记录和托管。实际写起来是这样的from 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] pending_tool_call: dict | None builder StateGraph(AgentState) builder.add_node(agent, call_llm) builder.add_node(tools, execute_tools) builder.add_node(human_review, review_node) builder.add_edge(START, agent) builder.add_conditional_edges( agent, decide_next, {tools: tools, ask_human: human_review, END: END} ) builder.add_edge(tools, agent) builder.add_edge(human_review, agent)State 的messages字段用了add_messages这个 reducer它的作用是每次节点返回新的消息时不是覆盖而是自动按 message id 去重合并。这玩意看着简单实际是状态管理里非常关键的设计——它保证了重试、恢复、多处修改消息时不会冲突。变更的核心在于LangGraph 把一个不可控的循环变成了一个有向图。图一旦编译出来每一轮节点执行就是一个超级步super-step每一步执行完框架都有机会做一次状态快照。这个快照就是后面要说的 Checkpoint。1.3 LangGraph 与 LangChain 的关系别搞混搜资料的时候一定会碰到LangChain 和 LangGraph 有什么区别这种高频问题这里顺手说清楚。LangChain 是模型调用、工具封装、检索这些组件的工具库LangGraph 是专门做 agent 编排和状态管理的运行时框架。它俩不是替代关系LangGraph 内部大量复用了 LangChain 的模型接口、tool 定义和消息结构。用生活化一点的类比LangChain 是工具箱里面有扳手、螺丝刀、电钻LangGraph 是流水线的传送带和控制台它决定你什么时候用哪把工具、用完下一步干嘛、哪一步要停下来质检。你完全可以用 LangGraph 而不碰 LangChain 的组件自己封装模型调用但那样会失去消息标准化、工具调用解析这些顺手的好东西。实际项目里我通常两个一起用LangGraph 负责编舞LangChain 负责提供演员。2. Checkpoint 才是恢复能力的根基2.1 没有持久化谈恢复就是空话光有状态机还不够。如果状态只存在内存里那和手写 Loop 没有本质区别只是把循环换成了图。要让中断恢复成立必须在每一步执行完之后把整个 State 落盘。LangGraph 里的 Checkpoint 干的就是这件事每一次超级步结束后把当前节点、状态数据、下一步该走哪条边全部保存到 Checkpoint 里。恢复的时候就很简单了给定一个 thread_id线程 ID你可以理解成一个会话的全局唯一编号框架从 Checkpoint 里读回上次的状态直接从断点继续执行。整个过程不需要重新调用之前的模型也不需要重放全部历史消息。Checkpoint 还有一个容易被忽略的价值它让 agent 的执行变成了可回放的。线上出了诡异问题我可以把某个 thread_id 的完整状态导出来本地恢复、逐步调试而不是靠日志猜。这对手写 Loop 来说几乎不可能。2.2 为什么选 PostgreSQL 而不是 SQLiteLangGraph 官方提供了几种 Checkpointer内存版 MemorySaver、SQLite 版 SqliteSaver、PostgreSQL 版 PostgresSaver还有 Redis 版。我最终选 PostgreSQL是因为它有几个不可替代的优势。第一事务性。Checkpoint 写入涉及多张表的协同更新PostgreSQL 的事务保证让部分写入不会发生。SQLite 也能开事务但在高并发写入下锁竞争明显容易报database is locked。第二多进程能力。agent 服务一般会横向扩容成多个进程SQLite 的文件锁机制在多进程下很痛苦PostgreSQL 天生支持并发连接多个 worker 可以同时读写不同 thread_id 的 Checkpoint互不干扰。第三复用现有基础设施。绝大多数公司已经有 PostgreSQL 实例不需要为 agent 单独引入一套存储运维成本和心智负担都低。从性能上说PostgreSQL 对 Checkpoint 这种单会话低频小写入的场景绰绰有余。每个超级步写入几个 KB 到几十 KB 的数据即使每秒几十个活跃会话PostgreSQL 也毫无压力。真正需要 Redis 级的吞吐是早期大促流量那种量级agent 应用现阶段远够不着。PostgresSaver 的连接和建表代码非常简洁from langgraph.checkpoint.postgres import PostgresSaver DB_DSN postgresql://agent_user:your_password127.0.0.1:5432/agent_db checkpointer PostgresSaver.from_conn_string(DB_DSN) # 首次运行必须执行创建 checkpoint 相关表 checkpointer.setup()setup()会在数据库里创建checkpoints、checkpoint_blobs、checkpoint_writes这几张表。checkpoints存的是每次超级步的主记录比如 thread_id、图的执行状态、节点位置checkpoint_blobs存序列化之后的完整状态数据checkpoint_writes存的是某个节点写入的增量数据。2.3 Checkpoint 写入的时机与代价值得专门说清楚的是 Checkpoint 的写入时机。它不是每执行一句模型输出就写一次而是在每个超级步完成后、图准备进入下一个节点之前统一落盘。也就是说一个节点内部循环处理 100 条消息不会产生 100 次写入只会在节点返回后写一次。这个设计把写入频率压得比较低避免了频繁磁盘 IO。不过要注意一个代价问题checkpoint 是每次完整保存整个状态。如果 State 里放了大的中间数据比如一段完整的文档内容、一串很长的检索结果那么每次快照都会把这堆东西重新序列化一遍。所以我在设计 State 时有条铁律只往 State 里放需要跨节点传递的核心数据能放引用 ID 的绝不放大对象能在恢复时重新拉取的数据绝不塞进 State。这跟你写 React 组件时控制 props 大小的道理一模一样状态越瘦快照越便宜。3. 用 interrupt 把中断变成一等公民3.1 interrupt 机制到底怎么工作Checkpoint 解决了崩溃后恢复但还有一个场景是主动中断你想让 agent 停下来等用户确认后再继续。LangGraph 提供了一个专门的原语叫interrupt()。interrupt()的用法简单得让人意外在节点函数里调用它传入一个值图的执行就会在那一刻暂停并把那个值作为可恢复的关键信息写入 Checkpoint。暂停后graph.invoke()调用会正常返回不是报错返回的是中断时的状态和__interrupt__数据。之后你想让它继续只需要对同一个 thread_id 再次调用invoke(None)或传入Command(resume...)。为什么这个东西比手写 Loop 里塞一个input()优雅因为 interrupt 是可持久化的暂停。进程即使在中途被杀掉Checkpoint 里已经记录了中断位置。重启进程、恢复同一个 thread_id它依然停在那里等你不会丢失我正在等用户审批这个事实。一个典型的审批节点长这样from langgraph.types import interrupt def human_review_node(state: AgentState): decision interrupt({ type: tool_approval, tool_name: state[pending_tool_call][name], args: state[pending_tool_call][args], message: 是否允许执行该工具 }) if decision ! approve: return {messages: [{role: assistant, content: 用户拒绝了工具调用}], pending_tool_call: None} return {pending_tool_call: state[pending_tool_call]}3.2 Human-in-the-loop 的两种常见模式实际项目里interrupt 最常见的两种用法我分开讲。第一种是执行前审批。比如 agent 要调用一个发送邮件的工具这是有副作用的操作不能让它自己拍板。图的流程走到human_review节点时 interrupt把将要执行的工具和参数发给前端。用户点批准Command(resumeapprove)让流程继续点拒绝Command(resumereject)走另一条分支。这一模式的要点是审批信息必须足够完整前端才能做出决策。我踩过坑是只传了工具名没传参数用户根本不知道 agent 要拿这个工具干什么。第二种是索取补充信息。agent 做数据分析时发现信息不够比如用户说帮我安排一趟行程但没给出日期和预算。此时 interrupt 向用户抛出问题用户回复具体信息后恢复执行。这比让模型瞎猜或者不停追问要可控得多因为每次中断都是显式的用户明确知道自己被问了什么。3.3 和手写 Loop 的直观对比手写 Loop 里做人工审批最粗暴的写法是user_input input(批准吗? y/n)这就堵死了异步交互。如果你改成用消息队列等审批结果那又是一套状态管理噩梦你得自己记录当前循环到哪个工具了审批单号和工具调用的对应关系。而 LangGraph 的 Checkpoint interrupt 组合把当前在哪、等什么、恢复后去哪全部自动记录。所以我的判断是只要你的 agent 会涉及人工确认、需要跨进程或跨天数的会话、或者有高可用要求就应该放弃手写 Loop 转向带 Checkpoint 的图运行时。这个决定早做早省心。4. AG-UI让 Runtime 和前端有共同语言4.1 AG-UI 是什么状态机和 Checkpoint 解决了服务端的问题但前端怎么和 Runtime 通信是另一个隐藏的坑。我之前是自己定义一套 JSON 轮询协议前端每隔几秒问一次agent 跑完了吗不仅浪费资源而且语义混乱——等待运行中需要审批这些状态全靠我拍脑袋命名。这次改造我用 AG-UI 统一了交互层。AG-UIAgent-User Interface是一个面向 agent 与用户之间交互的开放协议它把 agent 运行过程中的各种事件标准化定义了一组统一的 JSON 事件格式。它的定位类似 WebSocket 之于聊天、SSE 之于流式输出目标是让任何前端Web、桌面、移动端都能以一种通用的方式订阅 agent 的运行时状态和交互节点而不是每个项目各搞一套。有人会问这东西和 LangGraph 是什么关系我的理解是LangGraph 是 agent 的执行引擎AG-UI 是执行引擎与 UI 之间的通信契约两者互补。LangGraph 的interrupt是状态机层面的暂停原语AG-UI 解决的是这个暂停如何呈现给用户、用户如何把决定传回运行时。4.2 事件映射与交互闭环AG-UI 的事件类型覆盖了 agent 交互的完整生命周期。我这里列一下我在项目里实际用到的核心事件类型事件类型作用对应 LangGraph 场景agent_message传输 agent 的流式文本输出普通模型生成tool_call通知前端某个工具被调用节点进入 tools 前tool_result工具执行结果回传工具节点完成input_required告知前端需要用户输入/批准interrupt 触发时agent_state同步当前状态如 waiting/workingCheckpoint 状态查询lifecycle_event会话开始、结束、恢复invoke/resume这些事件的妙处在于前端不需要知道 LangGraph 的存在。它只负责渲染 AG-UI 事件收到input_required就弹出审批卡片收到agent_message就流式渲染文本收到tool_call就展示正在调用工具的状态。整个前端逻辑变成了一张事件驱动的关系表和具体后端框架解耦。之后就算我把 LangGraph 换成别的运行时只要还遵守 AG-UI前端一行都不用改。4.3 会话恢复时的协议配合中断恢复不只是后端状态恢复前端也需要重连。比如用户在手机上打开一个历史会话这个会话可能三天前被打断过。前端的做法是先通过 thread_id 从 Runtime 拿到当前状态agent_state如果发现存在未完成的中断就立刻渲染上次进行到一半需要您确认的界面同时注入一条input_required事件让用户继续操作。这块我自己的经验是恢复场景下事件重放要克制。不能把三天前所有的agent_message全部重放一遍前端只关心两件事——当前卡在哪一步、需要用户做什么。AG-UI 的agent_state和input_required组合可以精确回答这两件事其他历史事件按需从 Checkpoint 里查不作为重放内容。5. 完整迁移实操从零跑通中断恢复5.1 前置准备PostgreSQL 安装与连接先说数据库准备。PostgreSQL 的安装根据系统不同有些差异我在 macOS 上用 Homebrew 一条命令解决Linux 上则用包管理器Windows 上直接装官方安装包主要注意端口和密码配置# macOS brew install postgresql16 brew services start postgresql16 # Debian/Ubuntu sudo apt update sudo apt install postgresql postgresql-contrib sudo systemctl enable --now postgresql装完创建 agent 专用的数据库和账号不建议直接用 postgres 超级用户跑业务CREATE USER agent_user WITH PASSWORD your_password; CREATE DATABASE agent_db OWNER agent_user; GRANT ALL PRIVILEGES ON DATABASE agent_db TO agent_user;然后在本地的 Python 环境里装依赖pip install langgraph langgraph-checkpoint-postgres langchain-openai这里有个坑值得提醒langgraph-checkpoint-postgres依赖psycopg新的 Psycopg 3安装的时候会自动带上来但如果你的环境里有老版本的psycopg2可能出现驱动冲突。我遇到过连接串写对了、密码也对但一直报module psycopg has no attribute connect的情况最后检查发现是环境里同时存在两个 psycopg 版本导致的。建议在干净的虚拟环境里安装或者显式指定psycopg[binary]。5.2 编译一个带 Checkpoint 和中断的 Graph下面是我改造后的一个精简但完整的例子。它包含一个会调用工具的 agent 节点、一个工具节点、一个人工审批节点并且用 PostgresSaver 做 Checkpoint。import json from typing import TypedDict, Annotated from langgraph.graph import StateGraph, START, END from langgraph.graph.message import add_messages from langgraph.checkpoint.postgres import PostgresSaver from langgraph.types import interrupt, Command from langchain_openai import ChatOpenAI from langchain_core.tools import tool DB_DSN postgresql://agent_user:your_password127.0.0.1:5432/agent_db class AgentState(TypedDict): messages: Annotated[list, add_messages] pending_tool_call: dict | None tool def query_order(month: str) - str: 查询指定月份的订单量 return json.dumps({month: month, order_count: 128}) def call_llm(state: AgentState): model ChatOpenAI(modelgpt-4o-mini).bind_tools([query_order]) response model.invoke(state[messages]) return {messages: [response]} def execute_tools(state: AgentState): tool_call state[pending_tool_call] result query_order.invoke(tool_call[args]) return { messages: [{ role: tool, tool_call_id: tool_call[id], content: result }], pending_tool_call: None } def human_review(state: AgentState): tc state[pending_tool_call] decision interrupt({ type: tool_approval, tool_name: tc[name], args: tc[args] }) if decision ! approve: return {pending_tool_call: None} return {pending_tool_call: tc} def decide_next(state: AgentState): last_msg state[messages][-1] if not last_msg.tool_calls: return END if last_msg.tool_calls[0][name] query_order: # query_order 需要人工审批 return ask_human return tools builder StateGraph(AgentState) builder.add_node(agent, call_llm) builder.add_node(tools, execute_tools) builder.add_node(human_review, human_review) builder.add_edge(START, agent) builder.add_conditional_edges(agent, decide_next, {tools: tools, ask_human: human_review, END: END}) builder.add_edge(tools, agent) builder.add_edge(human_review, agent) checkpointer PostgresSaver.from_conn_string(DB_DSN) checkpointer.setup() app builder.compile(checkpointercheckpointer)注意我故意把query_order这个工具设置成必须经过人工审批目的是演示 interrupt 的真实效果。decide_next里根据工具名分流走到human_review节点就暂停。5.3 模拟崩溃与恢复现场编译完成后先发起第一次调用config {configurable: {thread_id: demo-001}} result app.invoke( {messages: [{role: user, content: 帮我查一下3月份的订单量}]}, configconfig ) print(中断信息:, result[__interrupt__])程序会打印类似这样的中断信息中断信息: [Interrupt(value{type: tool_approval, tool_name: query_order, args: {month: 3月}})]此刻图已经暂停在human_review节点等用户审批。注意这个暂停是持久化的——即使我现在把 Python 进程杀掉数据也已经写入 PostgreSQL。模拟崩溃强制结束进程然后再写一个全新的脚本用同一个 thread_id 去查询当前状态from langgraph.checkpoint.postgres import PostgresSaver checkpointer PostgresSaver.from_conn_string(DB_DSN) config {configurable: {thread_id: demo-001}} state checkpointer.get(config) print(当前状态:, state.values) print(下一个节点:, state.next)输出会显示pending_tool_call还在next指向human_review证明状态没有丢。接着模拟用户点了批准恢复执行result app.invoke(Command(resumeapprove), configconfig) print(result[messages][-1].content)这一步会继续执行human_review之后的流程节点返回了审批通过的 pending_tool_call图进入tools执行真实工具然后回到agent生成最终回答。整个恢复过程完全自动不需要重新调用模型生成前面的历史。我实测的恢复时间在毫秒级主要耗时是查询本节点状态的耗时和生产空闲期调用模型的时间。6. 常见问题与排查实录6.1 高频问题速查表改造过程中我整理了这张问题对照表基本都是网上社区里被反复问过的类型也覆盖了我自己真实遇到的情况现象根本原因解法执行时报表不存在、KeyError 指向 checkpoint忘记调用checkpointer.setup()首次运行先建表连接数据库超时 / connection refused端口、密码、监听地址配置问题检查pg_hba.conf和postgresql.conf监听设置恢复时找不到 thread_id配置里的thread_id拼错或大小写不一致统一用get_state验证 thread_idinterrupt 后invoke(None)不生效中断点有多个恢复需要显式Command(resume...)用Command(resumevalue)显式传值状态里存自定义对象导致序列化失败LangGraph 默认用 JSON 序列化 State只存基本类型或自行实现序列化逻辑多个进程同时写同一 thread_id 导致状态错乱同一会话并发更新应用层按 thread_id 加锁或避免同一会话并发触发AG-UI 事件收不到输出错误混用了invoke阻塞调用和流式事件流按场景区分原语调用与事件订阅模式6.2 几个值得注意的坑第一个坑是 State 的序列化边界。LangGraph 的 Checkpoint 默认会把 State 序列化存储虽然它内部用了高效的序列化方案但如果你在 State 里塞了自定义的 Python 类实例恢复时大概率反序列化失败。我的规避方法是State 里只放 dict、list、str、int 这类基础结构需要复杂对象时用ID 存储层替代需要时再重新加载。这个原则我第一次重构时没遵守结果线上一次崩溃恢复直接报错被迫回滚。第二个坑是messages的 reducer。刚开始我图省事给messages直接写list类型而不加add_messages结果每次节点返回消息都是整体覆盖历史消息被丢得乱七八糟。如果你也想保住完整的对话上下文务必使用Annotated[list, add_messages]。这玩意还有一个额外好处重复恢复调用时相同 id 的消息不会重复插入天然幂等。第三个坑和 AG-UI 有关别把 interrupt 的数据格式当成协议格式直接透传。我早期是interrupt({type:tool_approval, ...})然后把整个 dict 塞进自定义事件发给前端后来发现前端要做各种 if 判断代码很脆弱。正确做法是interrupt 的数据是给后端逻辑复用的前端展示和交互应该走 AG-UI 的标准事件类型在后端做一次映射。业务数据和协议数据分离前端才不会被你的内部实现绑架。第四个坑是恢复后的next节点判断。get_state返回的next字段表示接下来要进入的节点有时候它是(human_review,)这种元组形式有时候又是空元组表示已结束。判断是否还有未完成的工作时不能简单用if state.next:因为状态里可能还有 pending 的数据但next为空。这是我调试恢复逻辑时浪费过时间的地方建议仔细读一下官方对StateSnapshot.next语义的说明。最后再分享一个小技巧我个人实际操作中最受益的一个习惯是把恢复和新建统一封装成一个入口函数而不是在业务代码里到处写invoke。比如def run_or_resume(thread_id, new_inputNone): config {configurable: {thread_id: thread_id}} snapshot app.get_state(config) if snapshot.next: # 存在未完成任务执行恢复 return app.invoke(Command(resumenew_input), configconfig) # 新会话或已结束走正常输入 return app.invoke({messages: [{role: user, content: new_input}]}, configconfig)这个函数的思路很简单先看这个 thread_id 有没有未完成的状态有就恢复没有就当新会话处理。它解决了我当时最头疼的前端传同一个 thread_id 过来我怎么区分是要继续还是新问题的需求。配合 AG-UI 的agent_state事件前端每次打开页面先拉一次状态再决定显示继续上次任务还是发起新提问整个交互就顺了。这次重构下来我最大的体会是可恢复 Runtime 不是靠某个单一技术实现的它是状态机LangGraph、持久化PostgreSQL Checkpoint和通信协议AG-UI三个层次配合的结果。每一层解决一个独立的问题缺一个都不完整。如果你也正在被agent 跑一半就丢状态折磨建议按这个顺序去拆先把图跑通再挂上 Checkpoint最后再接入 AG-UI。走完这三步你的 agent 服务才真正配得上Runtime这个称呼。
返回列表