ARTICLE DETAIL

资讯详情

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

LangGraph+PostgreSQL:打造可中断恢复的Agent长任务运行时

LangGraph+PostgreSQL:打造可中断恢复的Agent长任务运行时 1. 手写 Loop 的三处硬伤为什么长任务一碰真实业务就崩1.1 状态只活在内存里进程一死就是失忆我最早做 Agent 应用的时候执行循环写得非常朴素一个while True套上模型调用、工具调用、结果判断所有上下文全部塞在内存的 dict 里。单机、单用户、跑 demo 的时候完全没问题速度还快。但一旦部署到真实业务里第一个被戳穿的就是状态无持久化这个问题。内存状态意味着只要进程重启、被 OOM Killer 杀掉、或者部署时滚动更新运行到一半的任务立刻变回白纸。普通对话丢了还能让用户重新说一遍但长任务不行——比如一个批处理流程已经处理完 30 份文档还差 20 份进程一挂就得全部重跑。用户不会关心你是 OOM 还是发布导致的他只知道又要重新来一次。1.2 没有断点概念人工介入等于重构第二个硬伤更隐蔽手写循环很难插入真正的人工确认。业务场景里经常需要机器跑一段暂停下来等人审批批完再继续。手写循环实现这种逻辑时你只能靠input()阻塞或者轮询数据库里的审批状态本质上把流程控制权拆得七零八落。我最初的做法是在循环里定期查一张审批表看到已通过才继续往下走。表面能用但问题在于任务运行到哪一步、状态长什么样、审批通过后从哪个节点续跑全部要靠自己用额外字段或者临时表来记录。这等于自己维护一套不完整的断点系统而且一旦记录和实际状态不一致恢复出来的根本不是你想要的现场。1.3 前端拿不到事件流只能靠轮询碰运气第三个问题来自前端。手写 Loop 跑起来之后网页端想知道现在到哪一步了最原始的办法就是定时轮询接口查状态。轮询的体验很差短了浪费请求、长了界面像卡住。更麻烦的是任务中断后要恢复前端根本不知道应该传什么参数、从哪儿续起。其实这整套痛点的本质就是缺一个可恢复的 Runtime既要能持久化状态又要能暂停等待外部输入还要能把运行过程以事件流的方式推给前端。后来我基于 LangGraph 做了一套带状态图语义的 Runtime用 PostgreSQL Checkpoint 解决状态存储再用 AG-UI 协议把中断恢复链路接到 UI才真正把这件事从玄学变成了可操作、可复现的系统。2. PostgreSQL Checkpoint 存储先搞明白它到底存了什么、为什么选 PG2.1 拆解 LangGraph Checkpoint 的存储模型LangGraph 的 Checkpoint 并不是简单地把整个对象 pickle 一下丢进数据库它的存储模型是围绕状态图执行设计的。一个 checkpoint 里至少包含四层信息当前执行到哪个节点、这个节点的输入输出是什么、整个共享状态State里的各个字段值、以及父子 checkpoint 之间的关系。有了父子关系LangGraph 才有能力做回放或者说续跑。恢复的时候框架拿到thread_id对应的最新 checkpoint找到上次执行被打断的位置然后从那个节点的入口重新进入。注意这里不是从图的最开始重新跑而是精确地回到断点这一点比很多手写的断点续传要严谨得多。在 PostgreSQL 存储实现里核心表大致是checkpoints、checkpoint_blobs、checkpoint_writes这三张。checkpoints记录每个 checkpoint 的元数据比如线程 ID、checkpoint ID、父 checkpoint IDcheckpoint_blobs存放序列化后的状态内容checkpoint_writes则记录节点执行过程中的中间写入这是恢复时重建节点输入的关键。初次接触的人容易只盯着状态存哪了实际上真正复杂的部分是执行到一半的现场怎么重建后者靠的就是这些中间写入记录。2.2 为什么选 PostgreSQL事务、并发、运维的三重兜底我评估过几种方案内存、SQLite、Redis、PostgreSQL各有各的适用场景但最终选了 PostgreSQL。先看内存存储快是真快但没有持久化进程重启就全没了等于白搭。SQLite 是单文件的部署简单适合本地开发或者单线程小流量场景但它的并发写入锁粒度和运维成熟度在真实多实例部署里不够看。Redis 可以做持久化性能很好但它本身不是为复杂条件查询 强事务设计的我需要检查某个线程的所有 checkpoint 历史、做数据订正、跨表查询的时候Redis 的模式就有点别扭了。PostgreSQL 的优势恰好补上前面几项的短板它天然支持事务一个 checkpoint 的写入要么完全成功要么完全不成功不会出现状态写一半的脏数据多个图实例并发执行时行级锁能保证同一个thread_id的写入不会互相覆盖运维层面pg_dump备份、流复制、监控告警都是现成生态。对于要上生产环境的可恢复 Runtime这些能力比性能数字好看重要得多。2.3 落地配置起一个 PG 容器并接入 Checkpointer开发环境我直接起一个 PostgreSQL 16 容器命名和端口按自己习惯来就行docker run -d --name pg-checkpoint \ -e POSTGRES_PASSWORDpostgres \ -e POSTGRES_DBlanggraph \ -p 5432:5432 postgres:16然后安装 Python 侧的适配包并连上。我用的主要是langgraph-checkpoint-postgres异步场景用AsyncPostgresSaverpip install langgraph-checkpoint-postgres psycopg[binary]from langgraph.checkpoint.postgres.aio import AsyncPostgresSaver checkpointer AsyncPostgresSaver.from_conn_string( postgresql://postgres:postgreslocalhost:5432/langgraph ) await checkpointer.setup()第一次运行时调用setup()会自动建好上面说的那几张表不需要手工写 DDL。这里有个小提醒setup()建议在启动流程里显式调用一次不要在每次请求里都执行免得启动时一堆实例同时对表结构做检查。3. StateGraph 如何变成可恢复 Runtime中断、恢复与线程隔离3.1 状态图和手写 While 循环的本质区别手写 Loop 是一段线性代码执行流程藏在控制流里状态藏在局部变量里。LangGraph 的 StateGraph 则把节点和边显式建模每个节点接收共享状态、返回状态增量边决定下一步走向哪里。这种显式结构带来的直接好处是——每一步的执行都是可以被记录、被回放、被恢复的转换。打个比方手写 Loop 像是你在黑板上一步步算题粉笔字就是状态下课一擦进程重启就没了StateGraph 像是每一步都拍了照片存档想接着算的时候把最新一张照片摆出来从那里继续算。Checkpoint 就是这些照片PostgreSQL 就是放照片的保险柜。3.2interrupt是真正的暂停不是报错退出LangGraph 提供interrupt机制它可以暂停图的执行并等待外部输入。含义上和抛异常终止完全不同异常终止会让整个执行栈销毁而interrupt是一种可控的挂起框架会把当前状态持久化到 Checkpoint然后安静地等待。恢复时只需要向同一个thread_id传入Command(resume...)框架就会从暂停的那个节点继续往下。别人第一次看到这套机制时最容易问那个被 interrupt 的节点会不会重新执行一遍答案是会从 interrupt 点的后半段继续前半段已经完成的副作用操作不会重复执行这就是_checkpoint_writes 中间写入在起作用。理解这一点对后面排查重复事件问题非常关键。3.3 通过 thread_id 构建可恢复会话LangGraph 里用config中的thread_id作为会话身份标识。同一个thread_id的所有执行共享一条状态历史不同thread_id之间完全隔离。这个设计让并发变得非常简单不需要自己加锁直接多开几个thread_id各跑各的。config {configurable: {thread_id: task-001}}我一开始不太适应这种把身份放在 config 里的用法总觉得应该显式传入函数。后来才明白这种方式的好处是任何一次invoke都是幂等的、可定位的只要thread_id一致无论在哪个进程执行都能命中同一份 Checkpoint。3.4 有 Checkpointer 的图才叫 Runtime否则只是函数这是我反复强调的一句话不带 Checkpointer 的图本质上还是一个普通函数调用输入进去吐输出出来跑完就没了只有编译时挂上 Checkpointer图才升级成运行时——有记忆、可暂停、可恢复、可审计。编译挂载的代码非常简单graph builder.compile(checkpointercheckpointer)但这一步背后的语义变化是巨大的。从此刻起图执行过程中的每次节点转换都会被持久化外部进程杀不掉它的记忆重启之后照样能续跑。很多教程只告诉你调用 compile(checkpointerx) 就行却没有强调这一步是把执行变成状态机的分水岭。4. 用 AG-UI 打通 Runtime 与前端事件流替代轮询4.1 AG-UI 补上的是Agent 与界面之间的协议空白之前的方案里后端跑 Agent 任务前端想看进度最常见的有两条路轮询 REST 接口或者自己定义一套 WebSocket 消息格式。轮询的体验前面说过了自建消息格式则会导致每次项目都从零造轮子事件命名、字段结构、错误语义全凭个人喜好。AG-UI 就好在它把这层通信标准化了。它定义了一组 Agent 运行过程中会向前端推送的事件类型比如文本输出、日志、状态变更通知、需要用户操作的通知等等。前端只需要实现 AG-UI 的事件解析就能复用到任何遵循该协议的 Runtime 上不用一个项目换一套协议。对于做平台型 Agent 产品的人来说这套标准等价于给前端不通后端这堵墙开了一扇统一的门。事件载荷的基本原则是轻量结构化核心字段大致围绕事件类型、事件 ID、载荷数据、关联线程。它不规定你必须怎么实现业务逻辑只是帮你把发生了什么、现在需要谁做什么讲清楚。4.2 把中断恢复流程映射成 AG-UI 事件接入 AG-UI 后我的 Runtime 服务端会做一层桥接把 LangGraph 节点的执行状态翻译成 AG-UI 事件推送出去。下面是一段我实际用的简化版桥接代码def runtime_to_agui(thread_id: str, event: dict) - dict: event_type event.get(type) if event_type interrupted: return { type: user_action, id: finterrupt-{thread_id}, payload: { thread_id: thread_id, need: approval, data: event.get(payload) } } if event_type resumed: return { type: agent_state_changed, id: fresume-{thread_id}, payload: { thread_id: thread_id, state: running } } ...这里最关键的是user_action事件它告诉前端现在任务暂停了需要人工处理这是暂停位置的数据。前端收到后直接弹出审批面板用户点通过/拒绝后后端带着Command(resumeresult)恢复执行再推送一个agent_state_changed事件通知界面刷新状态。4.3 前端断点续跑的交互闭环接上事件流之后用户侧的交互变得很顺页面加载时如果有未完成的thread_id先请求一次运行时快照拿到当前状态再建立事件流订阅。中断发生时界面从运行中切换到等待确认恢复后继续推流。整个过程前端不再需要关心数据库里 Checkpoint 长什么样也不需要知道 LangGraph 内部有多少节点只需消费标准事件。这也是 AG-UI 给我的最大价值业务逻辑和技术细节被事件协议隔开前端团队和后端团队各自面向协议编程而不是面向彼此的数据结构编程。5. 完整联调实录模拟进程被杀之后的中断恢复5.1 测试场景设计纸上谈兵没用我直接设计了一个贴近业务的测试场景一个长文档批处理 人工审批流程。流程包括三个节点process_docs模拟批量处理文档每个文档耗 1 秒human_review插入人工审批暂停点finalize在审批通过后执行汇总收尾。用 PostgreSQL 做 Checkpoint 存储在人工审批阶段把进程杀掉再重启恢复验证整个系统能把状态接上。5.2 图定义与执行代码import time from typing import TypedDict from langgraph.graph import StateGraph from langgraph.types import interrupt, Command class TaskState(TypedDict): doc_ids: list[str] processed: list[str] approved: bool def process_docs(state: TaskState): processed [] for doc_id in state[doc_ids]: time.sleep(1) processed.append(doc_id) return {processed: processed} def human_review(state: TaskState): decision interrupt({ question: 是否通过文档审核, processed: state[processed] }) return {approved: decision yes} def finalize(state: TaskState): return {result: fdone, approved{state[approved]}} builder StateGraph(TaskState) builder.add_node(process_docs, process_docs) builder.add_node(human_review, human_review) builder.add_node(finalize, finalize) builder.add_edge(process_docs, human_review) builder.add_edge(human_review, finalize) builder.set_entry_point(process_docs) builder.set_finish_point(finalize) graph builder.compile(checkpointercheckpointer)第一轮执行跑到人工审批就停下来config {configurable: {thread_id: task-001}} result graph.invoke({doc_ids: [a.txt, b.txt, c.txt]}, config)运行日志如下[13:01:02] node: process_docs start [13:01:05] node: process_docs done - 3 docs [13:01:05] node: human_review interrupted, waiting for approval此时任务挂起Checkpoint 已经落库。紧接着我手动模拟进程被杀直接kill -9掉 Python 进程然后重启服务。这个模拟很粗暴但恰好能验证任何进程意外死亡场景下的恢复能力。5.3 重启恢复同一个 thread_id 接着跑重启进程后用同一个thread_id传入人工审批结果config {configurable: {thread_id: task-001}} graph.invoke(Command(resumeyes), config)观察到的日志是[13:01:46] node: human_review resume, decisionyes [13:01:46] node: finalize start [13:01:46] node: finalize done - resultdone, approvedTrue [13:01:46] task-001 COMPLETED注意日志里没有重新出现process_docs start因为恢复是从human_review的暂停点继续的。三份文档的处理结果直接从 Checkpoint 重建出来没有被重复执行。这就是状态图 Checkpoint 相对手写 Loop 最直观的碾压式优势。5.4 并发隔离与中间态检查我同时开了另一个thread_id跑同样的文档流程它和task-001完全独立。看 PostgreSQL 表时两个线程的 checkpoint 历史各归各的SELECT thread_id, checkpoint_id, parent_checkpoint_id FROM checkpoints WHERE thread_id IN (task-001, task-002) ORDER BY checkpoint_id;thread_id就是天然的隔离边界不需要额外加锁或者排队。中间态也没必要只靠日志观察直接查库就能确认SELECT thread_id, checkpoint_id, parent_checkpoint_id, checkpoint FROM checkpoints WHERE thread_id task-001 ORDER BY checkpoint_id;这套数据库里能查到每一次执行现场的特性在排查问题和做审计时价值极高。线上如果有人说任务跑偏了我能直接拉出他的线程历史看到每一步状态和父节点关系而不是靠猜。6. 踩坑复盘文档里不会教你的恢复运行细节6.1 thread_id 没传或传错恢复会静默变成新会话这个坑我踩过不止一次。开发时偶尔忘了在 config 里写thread_id或者在不同环境下传了不同的thread_id结果每次invoke都从零开始。最可怕的是它不会报错看起来一切正常但你感觉不到恢复在发生。排查方法很简单检查checkpoints表里同一个thread_id的 checkpoint 曲线。如果每次执行都只有一条孤立记录、没有父子关系链那基本可以确定 thread_id 没有正确复用。建议在入口层做好配置的强制校验宁可报错也不要静默开新会话。6.2 interrupt 必须配合 Checkpointer否则直接报错interrupt不是随处可用的魔法函数。如果图没有在compile时挂载 Checkpointer执行到interrupt会直接报错因为它没法持久化暂停状态。这个规则的背后逻辑很简单暂停和恢复依赖 Checkpoint没有 Checkpoint 就没有可恢复这个概念。所以排查为什么 interrupt 挂了时先看compile(checkpointer...)是否真的生效而不是盯着 interrupt 的传参看半天。6.3 恢复后重复事件的幂等处理接入 AG-UI 之后我遇到一个实际问题恢复执行时桥接层可能会把同一个节点的事件推送给前端两次。原因在于 LangGraph 从 Checkpoint 重建现场时会重新走到中断节点附近如果桥接层不做去重前端就会看到等待审核弹窗闪一下又消失、然后又出现。我的经验是推送事件时把checkpoint_id或节点执行 ID 作为事件 ID 的一部分前端按事件 ID 去重。另一个更保险的做法是桥接层只推送增量状态变化而不是每次执行都全量广播当前节点状态。这两件事叠加之后恢复流程在前端看起来就非常干净了。6.4 PostgreSQL 连接池与长空闲断开可恢复 Runtime 跑了一段时间后会遇到一个经典问题闲置一段时间的连接报server closed the connection unexpectedly。这多半是 PostgreSQL 服务端把空闲连接断掉了而客户端连接池还认为它是活的。解决思路有几个数据库侧开启tcp_keepalives_idle等参数或者客户端用短连接池、定期清理空闲连接。我个人更倾向于控制连接池的max_size并且显式配置连接空闲回收尤其是大量thread_id并发的时候连接数很容易被子任务吃满。日志里如果看到connection pool exhausted不要急着加资源先看是不是有些查询把连接拽住不放。6.5 序列化格式可读性优先还是体积优先LangGraph Checkpoint 的序列化格式有得选核心取舍是 JSON 与 msgpack。JSON 可读性好出问题的时候能直接打开数据库看内容非常利于定位msgpack 体积小、序列化快适合存大量状态或高并发写入但排查问题时得先做反序列化才能看懂。我开发期和测试期全程用可读性好的格式只有临近上线前才评估是否需要切换成紧凑格式。日志观察和问题排查的成本在项目初期远高于那点存储空间的成本别为了显得很专业过早优化。6.6 版本兼容是隐藏的坑LangGraph 及其 Checkpoint 适配包迭代很快langgraph、langgraph-checkpoint-postgres之间的版本如果不匹配行为会很奇怪有时是建表结构对不上有时是恢复时报奇怪的校验错误。我踩过一次升级后旧 checkpoint 无法被新版本读取的坑最后只能手动迁移数据。我的建议是在项目里锁住主要版本升级 Checkpoint 存储相关依赖时,必须在测试环境先跑一遍中断 - 升级 - 恢复的完整链路确认旧数据能正常续跑再上生产。这条经验帮我避开过很多次线上事故也让我意识到可恢复这件事本身也需要被纳入版本兼容性的测试范围。最后再分享一个让我省心很多的做法Checkpoint 所在的 PostgreSQL 单独做一份流复制备份。以前我怕状态库挂了导致所有会话无法恢复后来把备份和恢复演练纳入例行检查真遇到磁盘故障时切到备库只需要改连接配置所有thread_id的断点历史都还在。可恢复 Runtime 的价值说到底就体现在这种即便基础设施出了意外业务会话依然能接上的底线上。
返回列表