ARTICLE DETAIL

资讯详情

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

从手写Loop到可恢复Runtime:LangGraph与PostgreSQL Checkpoint实战

从手写Loop到可恢复Runtime:LangGraph与PostgreSQL Checkpoint实战 做 Agent 服务这大半年我踩得最深的一个坑就是任务跑到一半状态全没了。早先我图省事用 while 循环加全局字典手写了一套 Agent 调度逻辑当时觉得挺酷一上线就现原形——进程一重启、用户一刷新、网络抖一下就从头再来用户骂声一片。后来我把整套调度重构成基于 LangGraph 的图执行 Runtime状态快照落进 PostgreSQL Checkpoint前端通信统一走 AG-UI 风格的事件流中断恢复才算真正跑通。这篇文章会把从手写 Loop 到可恢复 Runtime 的完整改造过程、核心代码和踩坑记录都整理出来适合正在做 Agent 落地、被 human-in-the-loop 和状态持久化折磨的开发者。就算你还没用过 LangGraph按照文中的思路也能搭出一套能扛住故障的执行底座。1. 为什么最终放弃了手写 Loop调度逻辑的三宗罪1.1 状态全靠内存变量重启就失忆先说我最早期的手写版本。结构很简单每个用户一个 session 字典里面塞着对话历史、当前步骤、待处理的工具调用然后一个 while 循环不停跑直到模型不再要求调用工具为止。核心逻辑大概长这样# 早期手写版本伪代码 def run_agent(user_id, user_input): session sessions.get(user_id) # sessions 是全局字典 if session is None: session {history: [], step: 0, pending_tool_calls: None} while True: messages build_messages(session) resp llm.invoke(messages) if resp.tool_calls: session[pending_tool_calls] resp.tool_calls tool_results execute_tool(resp.tool_calls) session[history].append(tool_results) continue return resp.content表面问题谁都知道全局字典放在进程内存里进程一挂全部归零多实例部署更是各说各话。但更麻烦的是隐藏问题——为了让循环能接续执行我不得不在 session 里塞越来越多字段执行到哪一步了、哪个工具调用已经返回、哪些消息已经发给过模型、用户答没答完问题……这些字段之间还有隐式依赖业务复杂一点就互相踩。所有状态都存在内存变量里意味着每一次进程重启、每一次部署发布对用户来说都是一次失忆。用户正在填的订票表单、Agent 已经查完的航班列表、模型已经推理到一半的结论全部清零。这在 C 端产品里几乎是不可接受的。1.2 断点续跑要自己存整个执行栈轮子越造越沉有人会说用 Redis 把 session 存下来不就行了我也这么试过。序列化一份 JSON把 pending_tool_calls、history、step 全放进去恢复的时候读出来继续 while 循环。但问题在于一个真实的 Agent 任务不只是对话历史那么简单。它包含当前图走到哪个节点、某个工具调用的中间结果、等待用户确认的挂起请求、以及模型调用所依赖的完整消息上下文。这些状态分布在整个执行栈里。手写实现恢复逻辑等于自己造一套断点续跑系统而且这套系统还没有统一接口散落在业务代码里加一个交互节点就要改一遍恢复代码。LangGraph 这类图执行引擎出现本质就是把流程编排和状态持久化从业务代码里抽出来。流程写成一张图状态由框架统一做快照我只要关心节点里干了什么不用再操心怎么恢复、怎么续跑。1.3 前端通信没有协议每次都要改接口手写 Loop 还有一个很隐蔽的成本前端交互。Agent 执行过程中经常需要用户介入——选一个航班、确认一个订单、补充一个参数。后端为了配合前端每一类交互都自定义一套 JSON 结构前端再为每一类写一个渲染分支每次新增节点前后端都要同步联调。时间一长光维护交互消息协议就占掉了大量开发量。AG-UI 解决的就是这个痛点后端不暴露 LangGraph 内部结构而是把所有运行状态翻译成标准事件流前端只认事件、不认框架。你一换引擎前端代码一行都不用动。2. LangGraph 图执行 Runtime把流程和状态交给框架2.1 节点、边、状态从循环思维切换到图思维LangGraph 的核心抽象只有三个节点、边、状态。节点是执行单元可以是 LLM 调用、工具函数、任意 Python 函数边定义节点之间的流转关系状态则是贯穿整张图的共享数据通常用 TypedDict 声明。用生活化的类比以前手写 Loop 像一个人在一张纸上不断往下写流程写到哪、写到什么程度全凭自己记LangGraph 则像是画了一张地铁线路图每一站是节点路线是边车厢里运的乘客和货物就是状态。地铁开到哪里、车厢里装了什么都是系统统一记录的你不需要亲自盯着。这张图的运行由框架负责。节点函数输入当前状态、返回状态更新框架负责把更新合并回去并决定下一个该执行哪个节点。无用的 while 循环、标志位、临时变量全部消失了剩下的只有清晰的图结构。2.2 Checkpoint每一次执行都是一条可回放的状态链LangGraph 的 Checkpoint 机制是整套架构的灵魂。每执行完一个 SuperStep框架会把当前状态、节点写入位置、时间戳打包成一个 Checkpoint交给 Saver 持久化。多个 Checkpoint 会形成一条状态链——每个新快照都保存了父快照的引用。这意味着你不仅可以续跑还能回放和分叉。某一个时间点出了问题可以回到那个 Checkpoint 重新走一遍。这在人工审核、异常定位、任务重试的场景里价值极大。Saver 是一个接口层。LangGraph 自带 MemorySaver内存版重启即丢和 PostgresSaverPostgreSQL 持久化版也支持自定义 Saver。内存版适合开发调试生产环境必须上持久化版本。我的经验是只要涉及真实业务直接上 Postgres不要先内存后迁移迁移的坑比想象中多。2.3 interrupt让机器等人类并且能优雅地接上在 Agent 场景里最常遇到的情况是机器需要等用户输入。查完航班等用户选生成订单等用户确认缺参数等用户补充。手写 Loop 处理这种需求非常别扭而 LangGraph 原生提供了 interrupt 机制。from langgraph.types import interrupt, Command def confirm_selection(state): # 执行到这里会暂停等待外部输入 user_choice interrupt({ type: flight_selection_required, options: state[flight_options] }) # 恢复执行时user_choice 就是外部传入的值 return {selected_flight: user_choice}interrupt 触发后图执行会暂停Runtime 向调用方抛出一个__interrupt__事件同时当前状态已经由 Checkpoint 持久化。之后外部拿到用户的输入通过Command(resume...)把值传回来图会从中断节点继续往后执行。挂起的任务不会丢节点内部临时变量也都能接上这是手写 Loop 很难做到的。2.4 LangChain 与 LangGraph 的边界很多新手容易把 LangChain 和 LangGraph 混在一起。简单区分LangChain 负责提供组件比如模型封装、Prompt 模板、工具接入、输出解析LangGraph 负责编排和执行管流程、状态、持久化和人工介入。两者可以配合用也可以只用 LangGraph。我做这套 Runtime 时只用了 LangGraph 做图执行模型调用自己封装工具链还是 LangChain 的那套 load_tools。边界清晰反而更好维护。3. PostgreSQL Checkpoint 落地从连接串到生产可用的 Saver3.1 为什么是 PostgreSQL而不是 Redis 或内存选型的时候我做过一轮对比方案持久化事务能力运维成本适用场景MemorySaver无无极低本地调试、单进程测试Redis 自建代码有无中已有 Redis 基础设施PostgreSQL Checkpoint有有低可复用业务库生产环境、需要审计和回放自研存储有取决于实现高不推荐除非需求极其特殊PostgreSQL 在我这边的核心优势有两个。第一是事务一次状态写入要么成功要么失败不会有半截快照这对恢复可靠性是底线要求。第二是可审计每一份 Checkpoint 都是数据库里的一行记录天然可以和业务数据放在一起查、一起备份。Redis 虽然快但快照语义弱掉电丢数据、AOF 重放复杂反而增加心智负担。3.2 接入步骤与底层表结构接入非常直接。我用的依赖是psycopgpsycopg3注意不是 psycopg2LangGraph 官方适配的是 3.x 版本然后创建连接、实例化 PostgresSaver、调用 setup 建表、最后 compile 的时候把 saver 挂进去。import psycopg from langgraph.checkpoint.postgres import PostgresSaver conn psycopg.connect( postgresql://agent:agent_pass127.0.0.1:5432/agent_runtime ) checkpointer PostgresSaver(conn) checkpointer.setup() # 创建 checkpoints 表和 writes 表 app workflow.compile(checkpointercheckpointer)setup 会自动建两张表checkpoints用于存储每次执行完成后的快照writes用于记录每个节点对 channel 的写入明细。checkpoints 表里有 thread_id、checkpoint_ns、checkpoint_id、parent_checkpoint_id、checkpoint、metadata 等字段。thread_id 标识一条独立的会话流checkpoint_id 是快照的唯一标识parent_checkpoint_id 把快照串成可回放的链。根据官方文档的说明表会按照 checkpoint_ns 分区实际查询时可以按 namespace 过滤避免单表数据量过大影响性能。代码审查时只需保证 setup 执行一次不要在每次请求里反复调用。3.3 序列化别让一个 datetime 毁掉整个快照接入最开始时我没注意序列化问题直接往 state 里塞了一个 Pydantic 对象结果跑起来报错提示无法序列化。LangGraph 默认对 checkpoint 做 JSON 序列化state 字段必须是 JSON 友好的类型。如果你有自定义类型、dataclass、datetime要么在进入 state 前转成字符串或字典要么给 PostgresSaver 配置 pickler。我的建议是前者让 state 保持简单结构。复杂对象放到独立存储里state 里只存 ID 或者摘要。这样恢复速度快而且从前端做状态展示也更方便。检查一个 state 是否能正常快照很简单跑一次完整流程然后调用app.get_state(config)如果返回值里没有异常说明快照链路是通的。3.4 运维要点连接池、清理与多实例并发生产环境下有几个必须处理的问题。第一连接管理。PostgresSaver 构造时传入的 Connection 会在整个图生命周期内被复用。不要把连接对象放在请求作用域里随手关闭建议用连接池或者至少应用级单例。多实例部署时每个实例维护自己的连接池通过同一个数据库表实现状态共享。第二快照清理。checkpoints 表会随着任务执行不断膨胀尤其高频会话特别明显。我一般做一个定时任务清理超过保留周期的旧快照DELETE FROM checkpoints WHERE thread_id %s AND checkpoint_id NOT IN ( SELECT checkpoint_id FROM checkpoints WHERE thread_id %s ORDER BY checkpoint_id DESC LIMIT 10 );保留最近 N 个快照既避免表无限增长也保留了回放和问题排查的能力。第三并发隔离。不同 thread_id 天然隔离同一个 thread_id 同时被两个请求 resume 理论上可能互相覆盖。所以对同一会话的恢复操作一定要串行化建议在应用层对 thread_id 做分布式锁或者用任务队列保证恢复请求逐个执行。4. AG-UI 协议层让前端只知道事件流不感知 Runtime4.1 Agent 页面交互为什么需要一个统一协议Agent 的前端交互和传统接口完全不同。传统接口是请求-响应Agent 则是一个持续一段时长、中间多次停顿、需要用户介入的异步过程。前端要展示的内容也很杂模型生成的消息、工具调用的状态、当前执行到哪个节点、等待用户选择的下拉框。如果没有统一协议每一个 Agent 项目都在重新发明一套前后端通信格式而且格式之间根本不兼容。AG-UI 的核心思想是把 Runtime 的执行过程抽象成统一的事件流。后端不暴露内部实现只向事件流里推标准结构的事件前端只消费事件根据事件类型渲染对应 UI。这样后端可以换引擎前端不用改后端升级交互节点前端依然复用同一套渲染器。4.2 用 FastAPI SSE 实现一个类 AG-UI 的接口我实际落地用的是 FastAPI 加 SSEServer-Sent Events。选择 SSE 而不是 WebSocket是因为 Agent 执行天然是单向流式输出SSE 足够用而且自带断线重连机制。一个核心注意点如果使用的是同步版 PostgresSaverFastAPI 端点要写成 def不要写成 async def。psycopg 的同步连接不能塞进事件循环里跑否则会把整个服务阻塞。Go 代码看起来就是路由层同步、流式输出实际体验没有问题。from fastapi import FastAPI, Request from fastapi.responses import StreamingResponse import json from booking_agent import get_graph app FastAPI() def to_agui_event(event): # 把 LangGraph 原生事件翻译成类 AG-UI 事件 if __interrupt__ in event: interrupt_value event[__interrupt__][0] return { type: user_action, payload: interrupt_value, } for node_name, update in event.items(): if node_name find_flights: return {type: agent.progress, payload: {node: node_name}} if node_name book_flight: return {type: agent.done, payload: update} return {type: agent.message, payload: event} app.post(/v1/agent/run) def run_agent(payload: dict): thread_id payload[thread_id] message payload.get(message) config {configurable: {thread_id: thread_id}} def event_generator(): inputs {user_query: message} for event in get_graph().stream(inputs, config, stream_modeupdates): agui_event to_agui_event(event) yield fdata: {json.dumps(agui_event, ensure_asciiFalse)}\n\n return StreamingResponse(event_generator(), media_typetext/event-stream)前端接到user_action事件后把 options 渲染成选择框用户点击确认后带着选项回传给同一个接口。前后端之间只存在一种协议——事件流。LangGraph 内部是安全具体的还是换成了别的引擎对于接入层完全透明。4.3 恢复路径的实现同一个接口resume 参数中断恢复不需要单独的接口。我在/v1/agent/run里增加了 resume 字段如果请求带了 resume就调用graph.invoke(Command(resumeresume), config)如果没带就作为新消息启动。这样前端只需要用一个端点处理发消息和恢复两类操作。from langgraph.types import Command app.post(/v1/agent/run) def run_agent(payload: dict): thread_id payload[thread_id] config {configurable: {thread_id: thread_id}} resume_value payload.get(resume) def event_generator(): if resume_value is not None: events get_graph().stream(Command(resumeresume_value), config, stream_modeupdates) else: events get_graph().stream( {user_query: payload.get(message)}, config, stream_modeupdates, ) for event in events: agui_event to_agui_event(event) yield fdata: {json.dumps(agui_event, ensure_asciiFalse)}\n\n return StreamingResponse(event_generator(), media_typetext/event-stream)这个设计让前端几乎不需要理解恢复概念。页面上一个按钮、一个输入框提交时带上 thread_id 和对应的 resume 值就完成了从机器等用户到用户反馈后继续执行的闭环。4.4 get_state 与历史查询前端刷新后不丢局面SSE 断线或者用户刷新页面后前端需要重新拉取当前状态。LangGraph 提供了graph.get_state(config)返回当前快照、下一个待执行的节点、挂起的中断信息等。我把它包成一个 HTTP 接口app.get(/v1/agent/state) def get_state(thread_id: str): config {configurable: {thread_id: thread_id}} state get_graph().get_state(config) return { values: state.values, next: state.next, interrupts: [i.value for i in state.interrupts], }前端只需要把 thread_id 存到本地刷新后调用这个接口就能恢复整个页面的状态已经查到的航班还在、当前停在哪个确认节点、需要展示哪些选项一目了然。这也是整个方案最值钱的地方前端无状态状态全在底盘里。5. 实操记录一个订票 Agent 从崩溃到恢复的全过程5.1 场景设计与中断点选择我统一用一个航班订票助手作为示例因为它具备 Agent 应用最典型的形态第一步查航班是工具调用第二步需要用户选择是人工介入点第三步执行订票是收尾动作。中途任何一步崩溃都应该能从数据库快照恢复。图结构很简单START - find_flights - confirm_selection - book_flight - END。唯一的中断点放在 confirm_selection 节点。这个节点里调用 interrupt等用户选好航班后再拿出选择结果继续往下走。5.2 完整代码与依赖清单依赖文件如下注意 LangGraph 版本建议使用 0.3.x 以上版本interrupt 相关的 API 在这个版本里已经稳定。langgraph0.3.0 psycopg[binary]3.2 langchain-core0.3 fastapi0.115 uvicorn[standard]0.30图定义和运行逻辑的完整代码import json from typing import TypedDict, Optional import psycopg from langgraph.graph import StateGraph, START, END from langgraph.checkpoint.postgres import PostgresSaver from langgraph.types import interrupt, Command class BookingState(TypedDict): user_query: str flight_options: list selected_flight: Optional[dict] booking_result: Optional[str] def find_flights(state: BookingState): return {flight_options: [ {id: CA1801, price: 1280, time: 08:00, route: 上海-北京}, {id: MU5101, price: 1560, time: 10:30, route: 上海-北京}, ]} def confirm_selection(state: BookingState): choice interrupt({ type: flight_selection_required, options: state[flight_options], }) return {selected_flight: choice} def book_flight(state: BookingState): f state[selected_flight] return {booking_result: f预订成功{f[id]} {f[route]} {f[time]}票价 {f[price]} 元} workflow StateGraph(BookingState) workflow.add_node(find_flights, find_flights) workflow.add_node(confirm_selection, confirm_selection) workflow.add_node(book_flight, book_flight) workflow.add_edge(START, find_flights) workflow.add_edge(find_flights, confirm_selection) workflow.add_edge(confirm_selection, book_flight) workflow.add_edge(book_flight, END) def get_graph(): conn psycopg.connect( postgresql://agent:agent_pass127.0.0.1:5432/agent_runtime ) saver PostgresSaver(conn) saver.setup() graph workflow.compile(checkpointersaver) return graph if __name__ __main__: graph get_graph() config {configurable: {thread_id: booking-demo-001}} print( 第一次运行执行到中断点 ) for event in graph.stream( {user_query: 帮我订一张明天上午去北京的航班}, config, stream_modeupdates, ): print(json.dumps(event, ensure_asciiFalse, defaultstr)) print( 模拟恢复用户选择 MU5101 ) for event in graph.stream( Command(resume{id: MU5101, price: 1560, time: 10:30}), config, stream_modeupdates, ): print(json.dumps(event, ensure_asciiFalse, defaultstr))这里有个细节值得说Command(resume...)传进去的字典会被当作interrupt()的返回值。所以 resume 的数据结构要和中断事件里展示给用户的选项保持一致前端再把用户选中的那项完整传回来。5.3 第一次运行看到中断事件运行脚本后的输出大致分成两段。第一段会看到 find_flights 节点产生了航班列表然后 confirm_selection 节点触发了中断。LangGraph 的事件流里会出现__interrupt__标识里面带有我们在 interrupt 调用时传入的事件对象。此时整个图处于挂起等待用户反馈状态没有继续往下跑也不会自动过期。此时打开数据库查一下 checkpoints 表会看到这个 thread_id 已经存在快照记录。这意味着即使此刻整个进程被 kill 掉任务现场也已经被完整保存下来。5.4 模拟崩溃与用 thread_id 恢复现在模拟真实的生产事故运行到中断点之后直接 CtrlC 杀掉进程。然后重新启动一个新的 Python 进程只要业务代码不变使用同一个 thread_id 就可以从断点继续。# 新进程里执行 graph get_graph() config {configurable: {thread_id: booking-demo-001}} print(graph.get_state(config)) # 可以看到 next 是 confirm_selection状态里还有 flight_options for event in graph.stream( Command(resume{id: MU5101, price: 1560, time: 10:30}), config, stream_modeupdates, ): print(json.dumps(event, ensure_asciiFalse, defaultstr))新进程里没有重新输入用户查询也没有重新执行 find_flights航班列表直接来自数据库快照马上就能进入确认选择环节。整个续跑时间和完整执行相比几乎可以忽略这才是可恢复 Runtime 的真实价值。5.5 通过 HTTP 接口触发完整的中断-恢复闭环如果按 4.2 起一个 FastAPI 服务整个交互就是纯粹的 HTTP 事件流。先用 curl 发起第一次运行线程 ID 固定为booking-demo-002curl -N -X POST http://localhost:8000/v1/agent/run \ -H Content-Type: application/json \ -d {thread_id:booking-demo-002,message:帮我订航班}服务端返回 SSE 数据流前端解析出user_action事件后展示航班选项。用户点击航班后再把选中的航班对象作为 resume 值发回同一个接口服务端从断点恢复最后返回agent.done事件。整个链路不需要额外约定任何私有协议。6. 常见问题与排查实录环境、序列化、并发6.1 PostgreSQL 环境问题先谈环境。很多人在安装和使用 PostgreSQL 阶段就卡住了我整理一个速查表都是实际碰到过的高频问题。现象根因处理方式password authentication failed for user agent密码错误或认证方式不匹配确认密码检查 pg_hba.conf 中对应 host 的认证方式本地开发可用 md5 或 scram-sha-256database agent_runtime does not exist数据库未创建先用 postgres 超级用户建库建议一并指定 owner端口 5432 连不上服务未启动、防火墙拦截或装了多实例查看 pg_ctl status 或服务状态Windows 下注意服务名比如 postgresql-x64-17查询没有返回任何行但代码正常连接到了错误的数据库查看连接串的 dbname 和 host 是否匹配创建用户和数据库的 SQL 示例CREATE USER agent WITH PASSWORD agent_pass; CREATE DATABASE agent_runtime OWNER agent;6.2 LangGraph 接入报错psycopg.errors.UndefinedTable: relation checkpoints does not exist基本可以确定是没执行 setup。每次实例化 saver 之后先调saver.setup()它会执行建表语句。另一个常见问题是 psycopg 版本不匹配——接入langgraph.checkpoint.postgres一定要用 psycopg 3.x不要用 psycopg2。官方文档明确说明适配的是 psycopg3装错库会出现接口签名对不上的怪问题。6.3 序列化报错往 state 里塞非 JSON 类型时快照阶段会抛出序列化相关异常。做法是在节点内部就把复杂对象转换成基本类型或者查一下是否绕过了 state 类型声明直接放了对象。我踩过的一次具体事故是把一个 SQLAlchemy 模型对象直接放进了 state结果 checkpoint 写入时全链路报错排查了很久才发现是模型对象里有 bytes 字段。6.4 并发与数据膨胀同一 thread_id 被多个请求同时 resume会出现一方覆盖另一方写回的问题。这种不是数据库层面的死锁而是逻辑层面的互相覆盖。解决方案就是给同一个 thread_id 加锁或者让前端保证同一会话同一时间只有一个提交请求。快照表膨胀问题前面说过定时保留最近 N 条即可不需要全量手工清空。还有一个容易忽略的点每次执行完成的 Checkpoint 不是只存一份。中断前、中断后、恢复后都会产生新的 Checkpoint形成一条历史链。这也是为什么 get_state_history 可以做回到过去的核心但也意味着如果不做清理数据量增长会比想象快。坦白讲这套改造我前后折腾了两个版本。第一版把图定义、数据库连接、HTTP 路由全耦合在一个文件里虽然能跑但每次改业务都要碰底层。后来才意识到可恢复 Runtime 最大的价值在于分层图只关心流程Checkpoint 只关心状态AG-UI 协议层只关心交互。如果你想在自己的项目里引入这套思路我建议直接按这个边界切模块不要图省事揉在一起。最后再分享一个省心技巧把 thread_id 的生成规则和业务主键对齐。比如用户 ID 加任务类型加时间戳这样排查问题时直接从 Postgres 里按用户维度就能捞出一整条会话链不用翻日志猜 ID。状态可恢复是地基可排查才是日常真正救命的东西。
返回列表