ARTICLE DETAIL

资讯详情

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

从手写Loop到LangGraph Runtime:PostgreSQL Checkpoint与AG-UI中断恢复实战

从手写Loop到LangGraph Runtime:PostgreSQL Checkpoint与AG-UI中断恢复实战 1. 为什么我要从手写 Loop 切换到 LangGraph Runtime最早做 AI Agent 编排的时候我和大多数人一样直接上手写while循环。逻辑很直白调模型、解析输出、判断是否要调工具、执行工具、把结果塞回上下文、再调模型直到模型不再请求工具为止。这套东西在 demo 阶段跑得飞快几十行代码就能让一个 Agent 跑起来调试也简单打个断点就能看到每一步的状态。但问题很快就来了。第一个真实需求是用户中途关掉页面回来还能接着聊。手写 Loop 的状态全在内存里进程一重启就全没了。第二个需求是人工审核节点——某些敏感操作需要人工确认后才能继续这意味着执行流要能暂停、持久化、等外部信号再恢复。第三个需求是多轮长任务一个任务可能跑十几分钟中间涉及十几次工具调用任何一次网络抖动都可能导致整个流程崩掉重来。这三个需求叠加在一起手写 Loop 就彻底撑不住了。你当然可以自己实现状态序列化、自己搞断点续传、自己设计恢复协议但那就是在重新造一个 Runtime。LangGraph 的价值就在这里它把可中断、可恢复、可持久化的 Agent 执行流抽象成了图结构配合 PostgreSQL Checkpoint 做状态落盘再用 AG-UI 把执行过程实时推给前端。整套组合下来你得到的是一个真正意义上的 Agent Runtime而不是一段跑完就丢的脚本。这篇文章我会完整拆解这套方案为什么选 LangGraph 而不是裸写 LoopPostgreSQL Checkpoint 到底存了什么、怎么存AG-UI 在前后端之间扮演什么角色以及中断恢复这条链路从触发到恢复的完整实现。适合已经写过基础 Agent、正在被状态管理折磨、想上生产级编排的开发者。如果你还在纠结 LangChain 和 LangGraph 的区别我也会在对应章节里说清楚。2. 核心概念拆解LangGraph、Checkpoint、AG-UI 各自解决什么问题2.1 LangGraph 与 LangChain 的区别到底在哪很多人第一次接触这两个名字会懵。简单说LangChain 是组件库LangGraph 是编排引擎。LangChain 提供的是 LLM 封装、Prompt 模板、工具定义、检索器这些积木LangGraph 提供的是把这些积木按图结构串起来、并且管理执行状态的能力。打个比方LangChain 像是一箱乐高零件LangGraph 是那张带轨道的底板零件插在底板上才能跑起来。你完全可以用 LangChain 的组件配 LangGraph 的图两者不是替代关系而是互补关系。面试里常问的langchain 和 langgraph 区别标准答案就是LangChain 关注单个组件怎么用LangGraph 关注多个步骤怎么按状态流转、怎么中断、怎么恢复。LangGraph 的核心抽象只有三个State状态、Node节点、Edge边。State 是一个共享的数据结构所有节点读写它Node 是一个执行单元接收 State 返回 State 的增量更新Edge 决定下一个走哪个 Node可以是固定的也可以是条件分支。整个图编译之后变成一个可执行对象这个对象就是 Runtime 的载体。2.2 Checkpoint 不是日志是执行现场的快照Checkpoint 这个词在不同领域含义差别很大。数据库里的 checkpoint 是 WAL 落盘的检查点游戏里的 checkpoint 是存档点而 LangGraph 里的 Checkpoint 是图执行状态在某一时刻的完整快照。它存的东西包括当前执行到哪个节点、State 的完整值、待处理的下一步、以及这次执行的线程 IDthread_id。有了这些你就能在任意时刻把执行冻住之后用同样的 thread_id 把状态读回来从断点继续跑。这跟数据库 checkpoint 的思路其实是一致的——都是把易失的内存状态固化到持久化存储区别只是 LangGraph 固化的是 Agent 的对话与工具调用上下文。这里有个关键点Checkpoint 是按 thread 组织的。一个 thread 代表一条独立的执行线比如一个用户的会话。同一个 thread 下可以有多个 checkpoint按时间顺序排列最新的那个就是当前状态。这种设计让回到历史某一步重跑变得非常自然你只要指定 checkpoint_id 就能从任意历史点恢复。2.3 AG-UI 补上前后端之间缺失的那一环Agent 跑在后端用户在前端中间怎么通信传统做法是前端轮询或者后端推 SSE但 Agent 的输出是流式的、结构化的、带工具调用事件的普通 SSE 根本表达不了。AG-UI 就是为这个场景设计的协议它定义了一套标准的事件类型文本增量、工具调用开始、工具调用结束、状态更新、中断请求、恢复信号等等。AG-UI 的核心价值在于把 Agent 的执行过程变成前端可消费的事件流。前端不需要知道 LangGraph 内部怎么跑只需要按 AG-UI 协议接收事件、渲染 UI。当后端触发中断时AG-UI 会推一个中断事件给前端前端弹出确认框用户点确认后前端发一个恢复请求后端从 Checkpoint 恢复执行。整条链路是解耦的前端换框架、后端换模型都不影响协议层。3. 整体架构设计三层结构怎么搭3.1 分层思路与数据流向整套系统我分成三层编排层LangGraph、持久层PostgreSQL Checkpoint、交互层AG-UI 前端。数据流是这样的用户在前端发消息前端通过 AG-UI 协议把请求发给后端后端把消息塞进 LangGraph 的 State用 thread_id 启动或恢复图执行图执行过程中每个节点跑完都会写一次 Checkpoint 到 PostgreSQL同时节点产生的事件通过 AG-UI 推给前端如果遇到需要人工介入的节点图主动中断Checkpoint 记录中断位置AG-UI 推中断事件用户确认后后端用同一个 thread_id 恢复执行从 Checkpoint 读回状态继续跑。这个设计的精髓在于状态与执行分离。执行是瞬时的、可能崩的状态是持久的、可靠的。只要 Checkpoint 写成功了执行崩了也能恢复。这跟传统后端无状态服务 外部存储的思路一脉相承只不过这里的外部存储存的是 Agent 的执行现场。3.2 为什么持久层选 PostgreSQL 而不是 RedisLangGraph 官方支持多种 Checkpointer内存版、SQLite 版、PostgreSQL 版。内存版重启就丢只能做测试SQLite 版适合单机小规模但并发写会锁PostgreSQL 版是生产首选。选 PostgreSQL 的理由很实在。第一Checkpoint 本质是结构化数据用关系型数据库存天然合适查询、索引、事务都有保障。第二Agent 场景经常需要按用户、按会话、按时间查历史SQL 表达力足够。第三PostgreSQL 的 JSONB 类型能直接存 State 这种半结构化数据既保留了 schema 的严谨性又有文档数据库的灵活性。第四运维成熟备份、主从、监控一整套都是现成的。Redis 不是不行但它的持久化是尽力而为的RDB 有丢数据窗口AOF 性能又打折扣。Checkpoint 这种丢了就要重跑整个任务的数据还是交给 PostgreSQL 更稳妥。3.3 中断恢复的触发时机设计中断不是随便触发的得想清楚在哪些节点中断。我的经验是分三类强制中断涉及资金、删除、对外发送等不可逆操作前必须人工确认。条件中断当模型置信度低于阈值、或者工具返回异常时中断让人来判断。主动中断用户自己点暂停或者系统检测到长时间无响应。这三类中断在 LangGraph 里的实现方式略有不同。强制中断用interrupt_before在节点执行前打断条件中断在节点内部判断后主动抛中断主动中断则通过外部信号触发。不管哪种最终都会落到 Checkpoint 上恢复逻辑是统一的。4. 环境准备与依赖安装4.1 Python 环境与核心依赖我用的 Python 3.11太老的版本有些异步特性支持不好。核心依赖就三个pip install langgraph langgraph-checkpoint-postgres ag-ui-protocollanggraph是核心库langgraph-checkpoint-postgres是 PostgreSQL 的 Checkpointer 实现ag-ui-protocol是 AG-UI 的 Python SDK。如果你还要接具体模型再装对应的 SDK比如langchain-openai或langchain-anthropic。注意langgraph-checkpoint-postgres依赖psycopgpsycopg3不是老的psycopg2。如果你项目里已经有 psycopg2两个可以共存但别搞混连接字符串的格式。4.2 PostgreSQL 准备与表结构初始化PostgreSQL 建议 14 以上JSONB 的性能和索引支持更完善。建库建用户这些常规操作就不展开了重点说 Checkpointer 的表。LangGraph 的 PostgresSaver 提供了setup()方法会自动建表。你只需要在应用启动时调一次from langgraph.checkpoint.postgres import PostgresSaver DB_URI postgresql://user:passlocalhost:5432/agent_db with PostgresSaver.from_conn_string(DB_URI) as checkpointer: checkpointer.setup()它会建三张表checkpoints存主状态checkpoint_blobs存大的二进制数据比如消息历史checkpoint_writes存节点写入的中间结果。这三张表的分工要理解清楚后面排查问题全靠它们。4.3 连接池配置与生产注意事项生产环境千万别每次请求都新建连接。用ConnectionPool管理连接池大小根据并发量定。我的经验值是每个并发执行线程占一个连接池大小设为预期并发数 × 1.5留点余量。from psycopg_pool import ConnectionPool pool ConnectionPool( conninfoDB_URI, min_size5, max_size20, timeout30, ) checkpointer PostgresSaver(pool)提示checkpoint_blobs表会随着对话轮次增长而膨胀一定要配定期清理策略。我一般按 thread 的最后活跃时间超过 30 天没动的 thread 归档或删除。别等到表几个 G 了才想起来。5. 用 LangGraph 构建可中断的图5.1 State 设计存什么、不存什么State 是整个 Runtime 的核心设计得好后面省一半事。我的原则是存业务状态不存临时变量。from typing import Annotated, TypedDict from langgraph.graph.message import add_messages class AgentState(TypedDict): messages: Annotated[list, add_messages] user_id: str task_status: str pending_action: dict | None approval_result: str | Nonemessages用add_messages注解这是 LangGraph 提供的 reducer新消息会追加而不是覆盖。task_status记录任务阶段pending_action存待审批的操作approval_result存审批结果。不存什么不存数据库连接、不存 HTTP client、不存任何不可序列化的对象。State 最终要序列化进 PostgreSQL塞个连接对象进去直接报错。需要这些资源就在节点函数里现取或者用依赖注入。5.2 节点划分与职责边界节点划分我遵循单一职责 可独立恢复原则。一个节点只做一件事做完就写 Checkpoint。这样中断恢复的粒度最细恢复时重跑的成本最低。典型的节点划分agent_node调模型决定下一步是回复还是调工具tool_node执行工具调用approval_node人工审批检查点finalize_node收尾生成最终回复from langgraph.graph import StateGraph, START, END def agent_node(state: AgentState): response llm.invoke(state[messages]) return {messages: [response]} def tool_node(state: AgentState): last state[messages][-1] results execute_tools(last.tool_calls) return {messages: results} def approval_node(state: AgentState): action state.get(pending_action) if action and action[risk] high: decision interrupt({action: action, reason: 需要人工确认}) return {approval_result: decision} return {approval_result: auto_approved}interrupt()是 LangGraph 提供的中断原语调用它会立刻暂停图执行把当前 State 写进 Checkpoint然后等外部恢复信号。5.3 条件边与中断点的设置条件边决定流程走向中断点决定在哪停。这两个要配合设计。def route_after_agent(state: AgentState): last state[messages][-1] if hasattr(last, tool_calls) and last.tool_calls: return tool return end builder StateGraph(AgentState) builder.add_node(agent, agent_node) builder.add_node(tool, tool_node) builder.add_node(approval, approval_node) builder.add_edge(START, agent) builder.add_conditional_edges(agent, route_after_agent, { tool: approval, end: END, }) builder.add_edge(approval, tool) builder.add_edge(tool, agent) graph builder.compile( checkpointercheckpointer, interrupt_before[tool], )interrupt_before[tool]表示每次执行 tool 节点前都中断。这是最粗暴的做法实际项目里我会更精细只在特定条件下中断比如工具是删除或转账时才打断。实操心得interrupt_before是静态的编译时就定了。如果你需要动态决定是否中断用节点内的interrupt()函数更灵活。我踩过的坑是一开始全用interrupt_before结果所有工具调用都要人工确认用户体验极差。后来改成节点内判断只有高风险操作才中断。6. PostgreSQL Checkpoint 的落盘与恢复机制6.1 Checkpoint 的写入时机与内容LangGraph 的 Checkpointer 是每个超级步骤super-step写一次。超级步骤可以理解为一批可以并行执行的节点。每批节点跑完Checkpointer 就把当前 State 的完整快照写进checkpoints表把消息等大对象写进checkpoint_blobs。写入的内容包括thread_id、checkpoint_id本次快照的唯一 ID、parent_checkpoint_id上一个快照、state 的序列化值、以及元数据时间戳、来源节点等。这个链式结构让历史回溯成为可能——顺着 parent 指针能一路回到起点。6.2 恢复时怎么读回状态恢复的核心 API 是graph.invoke(None, config)。注意第一个参数传None表示不注入新输入从 Checkpoint 恢复。config 里带 thread_idconfig {configurable: {thread_id: user-123-session-1}} # 首次执行 result graph.invoke({messages: [user_msg]}, config) # 中断后恢复 result graph.invoke(None, config)LangGraph 会根据 thread_id 找到最新的 checkpoint把 State 读回来从中断点继续。如果中断是interrupt()触发的恢复时interrupt()会返回外部传入的值from langgraph.types import Command # 恢复并传入审批结果 result graph.invoke( Command(resumeapproved), config, )Command(resume...)是恢复中断的标准方式resume的值会成为interrupt()的返回值。6.3 多线程并发下的隔离thread_id 是隔离的关键。不同用户的会话用不同 thread_id状态天然隔离。但要注意同一个 thread_id 的并发执行会冲突。如果用户快速点两次发送两个执行流会同时读写同一个 thread导致状态错乱。我的处理方式是在应用层加锁同一个 thread_id 同时只允许一个执行流。用 Redis 分布式锁或者数据库行锁都行。锁的粒度是 thread 级别不同 thread 之间不影响。import redis lock_key fthread_lock:{thread_id} lock redis_client.lock(lock_key, timeout300) if not lock.acquire(blockingFalse): raise RuntimeError(该会话正在执行中请稍后) try: result graph.invoke(input_data, config) finally: lock.release()注意锁的超时时间要大于任务最长执行时间否则任务还没跑完锁就释放了并发问题照样出现。我一般设 5 分钟超长任务另做处理。7. AG-UI 打通前后端事件流7.1 AG-UI 的事件模型AG-UI 定义了一套事件类型核心的有这么几类事件类型触发时机前端处理TEXT_MESSAGE_CONTENT模型输出文本增量追加到消息气泡TOOL_CALL_START工具调用开始显示正在执行...TOOL_CALL_END工具调用结束显示结果STATE_UPDATEState 变化更新 UI 状态INTERRUPT图中断弹出确认框RUN_FINISHED执行结束结束 loading这套事件模型的好处是前端只依赖协议不依赖后端实现。后端从 LangGraph 换成别的编排引擎只要事件按 AG-UI 格式推前端一行不用改。7.2 后端事件推送实现LangGraph 支持astream_events能拿到执行过程中的细粒度事件。我把它转成 AG-UI 格式推给前端from ag_ui.core import EventType, TextMessageContentEvent async def stream_agent(input_data, config): async for event in graph.astream_events(input_data, config, versionv2): kind event[event] if kind on_chat_model_stream: chunk event[data][chunk] if chunk.content: yield TextMessageContentEvent( typeEventType.TEXT_MESSAGE_CONTENT, deltachunk.content, ) elif kind on_tool_start: yield ToolCallStartEvent( typeEventType.TOOL_CALL_START, tool_nameevent[name], ) elif kind on_tool_end: yield ToolCallEndEvent( typeEventType.TOOL_CALL_END, tool_nameevent[name], resultstr(event[data][output]), )中断事件要单独处理。当图因为interrupt()暂停时astream_events会结束你需要检查最终状态里有没有待处理的中断state graph.get_state(config) if state.next: yield InterruptEvent( typeEventType.INTERRUPT, reason需要人工确认, payloadstate.tasks[0].interrupts[0].value if state.tasks else None, )7.3 前端接收与恢复请求前端用 SSE 或 WebSocket 接收事件流按类型分发处理。中断事件到达时弹出确认 UI用户操作后发恢复请求async function resumeAgent(threadId, decision) { const response await fetch(/api/agent/resume, { method: POST, headers: { Content-Type: application/json }, body: JSON.stringify({ thread_id: threadId, resume: decision }), }); // 继续消费 SSE 流 consumeStream(response.body); }后端收到恢复请求后用Command(resumedecision)恢复图执行继续推事件流。整条链路闭环。实操心得SSE 连接容易断尤其是移动端切后台再回来。我的做法是前端记录最后收到的事件序号重连时带上序号后端从 Checkpoint 里找到对应位置重放。这样用户不会丢消息体验好很多。8. 完整实操从零跑通一次中断恢复8.1 初始化与首次执行先把所有组件串起来from langgraph.graph import StateGraph, START, END from langgraph.checkpoint.postgres import PostgresSaver from psycopg_pool import ConnectionPool pool ConnectionPool(DB_URI, min_size2, max_size10) checkpointer PostgresSaver(pool) checkpointer.setup() builder StateGraph(AgentState) builder.add_node(agent, agent_node) builder.add_node(approval, approval_node) builder.add_node(tool, tool_node) builder.add_edge(START, agent) builder.add_conditional_edges(agent, route_after_agent, {tool: approval, end: END}) builder.add_edge(approval, tool) builder.add_edge(tool, agent) graph builder.compile(checkpointercheckpointer) config {configurable: {thread_id: demo-thread-001}} result graph.invoke( {messages: [{role: user, content: 帮我删除订单 #12345}]}, config, )执行到approval_node时因为操作是删除interrupt()触发图暂停Checkpoint 落盘。8.2 中断状态检查state graph.get_state(config) print(下一步:, state.next) print(待处理中断:, state.tasks[0].interrupts if state.tasks else None)输出会显示next(approval,)和中断的 payload。这时候去数据库查checkpoints表能看到最新一条记录checkpoint_id就是当前快照。8.3 恢复执行与结果验证from langgraph.types import Command result graph.invoke(Command(resumeapproved), config) print(result[messages][-1].content)恢复后interrupt()返回approvedapproval_node继续执行返回approval_resultapproved然后走tool节点执行删除最后回到agent生成最终回复。验证恢复是否真的从断点续跑而不是从头重跑看agent_node有没有被重复调用。如果 Checkpoint 生效agent只会跑一次恢复后直接从approval继续。8.4 数据库侧验证SELECT thread_id, checkpoint_id, parent_checkpoint_id, metadata-step AS step, created_at FROM checkpoints WHERE thread_id demo-thread-001 ORDER BY created_at;你会看到多条记录parent_checkpoint_id串成一条链。中断点那条的metadata里会标记中断信息。恢复后新增的记录 parent 指向中断点证明是续跑而非重跑。9. 常见问题与排查技巧实录9.1 中断恢复典型问题速查问题现象可能原因排查方向恢复后从头重跑thread_id 不一致检查 config 里的 thread_id恢复报 State 反序列化失败State 里有不可序列化对象检查 State 字段类型中断事件前端收不到astream_events 提前结束检查中断后的状态读取逻辑Checkpoint 表暴涨没有清理策略加定期归档任务并发执行状态错乱同 thread 并发加 thread 级锁恢复后 interrupt 返回 Noneresume 值没传对检查 Command(resume...)9.2 我踩过的三个坑第一个坑State 里塞了 datetime 对象。本地测试没问题因为内存 Checkpointer 不序列化。换 PostgreSQL 后直接报错因为 JSONB 不认 datetime。解决办法是统一用 ISO 格式字符串或者自定义序列化器。第二个坑thread_id 用了随机 UUID。每次请求都生成新 thread_id结果永远恢复不了因为找不到历史 Checkpoint。thread_id 必须由业务逻辑生成比如user_id session_id保证同一会话用同一个。第三个坑中断后没检查 state.next。我以为astream_events结束就是执行完了结果中断时它也结束。前端一直显示 loading用户以为卡死了。后来加了状态检查发现state.next非空就推中断事件。9.3 性能优化建议Checkpoint 写入是同步的高频写入会成为瓶颈。我的优化手段批量写入如果一批节点可以并行让它们跑完一起写而不是每个节点写一次。异步 Checkpointer用AsyncPostgresSaver配合异步图执行写入不阻塞主流程。冷热分离活跃 thread 的 Checkpoint 放主库历史 thread 归档到冷存储。索引优化checkpoints表的thread_id和created_at建联合索引查询快很多。提示checkpoint_blobs表存的是消息历史增长最快。如果消息里带大文件或长文本考虑把大对象抽出来单独存Checkpoint 里只存引用。10. 生产部署的几点经验10.1 数据库连接与事务生产环境用 PgBouncer 做连接池前置应用侧的连接池大小可以调小。事务隔离级别用默认的 Read Committed 就够Checkpoint 写入本身是单条 INSERT不需要更严格的隔离。10.2 监控指标必须监控的指标Checkpoint 写入延迟、恢复成功率、中断到恢复的平均时长、活跃 thread 数、表增长速度。我用 Prometheus Grafana 搭的看板异常时告警。10.3 灰度与回滚新版本图结构上线前先用小流量验证。LangGraph 的图是编译时确定的改结构要重新编译。回滚时注意老版本的 Checkpoint 可能和新版本 State schema 不兼容要么做 schema 迁移要么让老 thread 用老版本图跑完。我在实际项目里跑这套组合已经大半年了最深的体会是中断恢复不是加个功能而是整个架构的地基。一旦你决定要支持中断恢复State 设计、节点划分、事件推送、前端交互全都要围绕它来。LangGraph PostgreSQL Checkpoint AG-UI 这套组合的好处是每一层职责清晰出问题能快速定位到是哪一层。最后分享一个小技巧调试恢复逻辑时直接在数据库里手动改 Checkpoint 的 state 值能模拟各种边界情况比写测试用例快得多。
返回列表