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 存储:先搞明白它到底存了什么、为什么选 PG
2.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 ID;checkpoint_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_PASSWORD=postgres \ -e POSTGRES_DB=langgraph \ -p 5432:5432 postgres:16然后安装 Python 侧的适配包并连上。我用的主要是langgraph-checkpoint-postgres,异步场景用AsyncPostgresSaver:
pip install langgraph-checkpoint-postgres psycopg[binary]from langgraph.checkpoint.postgres.aio import AsyncPostgresSaver checkpointer = AsyncPostgresSaver.from_conn_string( "postgresql://postgres:postgres@localhost: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(checkpointer=checkpointer)但这一步背后的语义变化是巨大的。从此刻起,图执行过程中的每次节点转换都会被持久化,外部进程杀不掉它的"记忆",重启之后照样能续跑。很多教程只告诉你"调用 compile(checkpointer=x) 就行",却没有强调这一步是把"执行"变成"状态机"的分水岭。
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": f"interrupt-{thread_id}", "payload": { "thread_id": thread_id, "need": "approval", "data": event.get("payload") } } if event_type == "resumed": return { "type": "agent_state_changed", "id": f"resume-{thread_id}", "payload": { "thread_id": thread_id, "state": "running" } } ...这里最关键的是user_action事件:它告诉前端"现在任务暂停了,需要人工处理,这是暂停位置的数据"。前端收到后直接弹出审批面板,用户点通过/拒绝后,后端带着Command(resume=result)恢复执行,再推送一个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": f"done, 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(checkpointer=checkpointer)第一轮执行,跑到人工审批就停下来:
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(resume="yes"), config)观察到的日志是:
[13:01:46] node: human_review resume, decision=yes [13:01:46] node: finalize start [13:01:46] node: finalize done -> result=done, approved=True [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 的价值,说到底就体现在这种"即便基础设施出了意外,业务会话依然能接上"的底线上。