ARTICLE DETAIL

资讯详情

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

从手写Loop到LangGraph+PostgreSQL Checkpoint:Agent中断恢复实战

从手写Loop到LangGraph+PostgreSQL Checkpoint:Agent中断恢复实战 1. 为什么手写 Loop 撑不过第三次线上事故大部分做 Agent 的团队起步阶段都是手写一个while循环调模型、解析工具调用、执行工具、把结果塞回消息列表、再调模型直到模型不再请求工具为止。这套东西在 Demo 阶段跑得飞快代码不到一百行逻辑一眼看得懂谁看了都觉得自己掌控全局。但只要把它放到真实业务里问题就会一个接一个冒出来而且往往是在最不该出问题的时候。我第一次被这套手写 Loop 教育是在一个需要调用外部接口做数据核对的场景里。整个流程要串七八个工具调用中间还要等人工审批。结果某天上游接口超时进程直接挂掉重启之后所有上下文全丢用户那边看到的是任务执行到一半突然消失。更麻烦的是有些工具是非幂等的——比如已经发出去的写操作重启后如果从头再跑一遍就会重复执行。那一次我们花了整整两天做数据回滚从那以后我就明白Agent 的运行时状态不能只活在内存里。手写 Loop 的第二个致命伤是中断恢复。真实场景里Agent 执行到一半需要暂停的情况太多了等人工确认、等外部回调、等定时任务触发、等一个长耗时任务完成。手写 Loop 要么阻塞在那里占着进程要么直接退出丢掉状态。你可能会说我把状态存数据库不就行了但问题是你要存的不只是消息列表还有当前执行到哪个节点、哪些分支已经走过、哪些工具的结果还没消费、下一步该往哪走。这些信息散落在你的业务代码里没有统一的抽象存起来容易恢复起来就是一场灾难。第三个问题是可观测性。手写 Loop 的执行路径是一条黑盒出了问题你只能靠打日志。但 Agent 的执行是树状甚至图状的有分支、有循环、有并行日志打出来是一团乱麻。你想知道为什么这次它选了工具 A 而不是工具 B光看日志根本还原不出来。所以当我第一次接触 LangGraph 的时候最打动我的不是它的图这个概念有多优雅而是它把状态管理和**检查点Checkpoint**做成了框架的一等公民。这意味着中断恢复不再是我自己拼凑的补丁而是运行时天然具备的能力。这篇文章我就把从手写 Loop 迁移到 LangGraph PostgreSQL Checkpoint AG-UI 的完整过程拆开讲包括我踩过的坑、参数怎么配、恢复逻辑怎么写以及前端怎么通过 AG-UI 把中断和恢复的状态实时呈现出来。提示这篇文章假设你已经对 LangGraph 的基本概念StateGraph、Node、Edge、State有初步了解。如果你完全没接触过建议先跑通一个最简单的两节点图再回来看否则后面的 Checkpoint 细节会有点吃力。2. LangGraph 的图模型到底比手写 Loop 强在哪2.1 状态是显式声明的不是散落在闭包里的手写 Loop 最隐蔽的问题是状态藏在各种变量和闭包里。你有一个messages列表、一个step_count计数器、一个pending_tool_calls队列可能还有一个retry_times字典。这些东西分散在函数作用域里你想序列化它们得一个个手动抠出来。LangGraph 的做法是让你先定义一个State通常是一个TypedDict或者 Pydantic 模型把所有需要跨节点传递的数据都声明在里面。比如from typing import Annotated, TypedDict from langgraph.graph.message import add_messages class AgentState(TypedDict): messages: Annotated[list, add_messages] task_id: str approval_status: str retry_count: int intermediate_results: dict这里有个关键细节Annotated[list, add_messages]里的add_messages是一个reducer。它的作用是定义当多个节点都想更新messages时怎么合并。默认情况下LangGraph 对同一个 key 的更新是覆盖式的但消息列表显然不能覆盖得追加。这个 reducer 机制是 LangGraph 状态管理的核心也是它比手写 Loop 强的地方——合并逻辑是声明式的不是命令式的。我一开始没重视 reducer结果在一个并行分支的场景里踩了坑两个节点同时往intermediate_results这个 dict 里写数据后写的把先写的覆盖了。后来改成用自定义 reducer 做 dict 合并才解决。这个坑在手写 Loop 里其实也存在只是你写的是dict.update()出了问题更容易定位而在图模型里你得理解 reducer 的执行时机才能排查。2.2 节点是纯函数副作用被隔离LangGraph 的节点本质上是一个接收 State、返回 State 更新的函数。这个约束看起来简单但它强制你把副作用调 API、写数据库、发消息和状态流转分开。手写 Loop 里这两者是混在一起的你调完工具直接把结果 append 到列表里中间没有任何边界。分开之后的好处是可测试。你可以单独测一个节点给它一个构造好的 State断言它返回的更新是否符合预期完全不用启动整个流程。我们团队现在的做法是每个节点都有对应的单元测试覆盖率要求 80% 以上。这在手写 Loop 时代是不可想象的因为一个函数里既有网络调用又有状态变更你只能做集成测试。另一个好处是可替换。比如某个工具节点一开始是直接调外部 API后来要改成走消息队列异步执行你只需要换掉这个节点的实现图的拓扑结构完全不用动。手写 Loop 里这种改动往往要动到主循环的逻辑。2.3 条件边让分支逻辑从 if-else 里解放出来手写 Loop 的分支逻辑就是一堆if-elif-else嵌套深了之后可读性极差。LangGraph 用**条件边conditional edge**把分支抽出来def should_continue(state: AgentState) - str: last_message state[messages][-1] if last_message.tool_calls: return tools if state[approval_status] pending: return wait_approval return end graph.add_conditional_edges( agent, should_continue, { tools: tool_node, wait_approval: approval_node, end: END, } )这个should_continue函数是纯的输入 State 输出字符串测试起来非常方便。而且路由表是显式声明的你一眼就能看出这个节点可能流向哪些地方。手写 Loop 里你想知道这个分支会走到哪得顺着代码往下读遇到函数调用还得跳进去看。2.4 循环是图的一部分不是 while 的副作用Agent 的典型执行模式是调模型 → 调工具 → 再调模型这本质上是一个循环。手写 Loop 用while True实现循环的终止条件藏在break语句里。LangGraph 里这个循环是图上的边agent → tools → agent终止条件由条件边决定。这个区别在中断恢复场景下是决定性的。手写 Loop 的while循环一旦中断你根本不知道它执行到第几轮、下一轮该从哪开始。而图模型里每一轮执行都会产生一个 Checkpoint记录当前在哪个节点、State 是什么。恢复的时候直接从最后一个 Checkpoint 继续循环的轮次这个概念被 Checkpoint 的序号天然替代了。我实测下来一个典型的 ReAct 风格 Agent用 LangGraph 重写之后代码量大概增加了 30%但换来的是可测试性、可观测性和中断恢复能力。这笔账怎么算都划算。3. PostgreSQL Checkpoint 的配置细节与那些文档没写的坑3.1 为什么选 PostgreSQL 而不是内存或 SQLiteLangGraph 内置了几种 CheckpointerMemorySaver内存、SqliteSaver本地文件、PostgresSaverPostgreSQL。Demo 阶段用MemorySaver没问题但生产环境必须换掉原因有三个。第一内存 Checkpointer 进程重启就丢。这跟手写 Loop 的问题一模一样等于白折腾。第二SQLite 不适合多实例部署。如果你的服务是水平扩展的多个实例同时读写同一个 SQLite 文件会锁竞争严重时直接报database is locked。第三PostgreSQL 支持并发写入和事务这是多实例场景的刚需。我们线上用的是 PostgreSQL 14配合psycopg驱动。LangGraph 的PostgresSaver会自动建表但有几个配置项必须手动调否则性能会很差。3.2 建表与连接池配置初始化 Checkpointer 的代码大概长这样from langgraph.checkpoint.postgres import PostgresSaver from psycopg_pool import ConnectionPool DB_URI postgresql://user:passhost:5432/agent_db?sslmodedisable pool ConnectionPool( conninfoDB_URI, min_size5, max_size20, timeout30, kwargs{autocommit: True, prepare_threshold: 0}, ) checkpointer PostgresSaver(pool) checkpointer.setup() # 首次运行建表这里有几个坑我要重点说。第一个坑autocommit必须开。LangGraph 的 Checkpointer 内部会自己管理事务边界如果你用默认的autocommitFalse会出现写入没提交的情况表现为 Checkpoint 查不到。我一开始没注意调试了半天以为是表没建对。第二个坑prepare_threshold0。这是psycopg3的一个优化参数设为 0 表示禁用 prepared statement 缓存。为什么禁用因为连接池里的连接会被复用如果开了 prepared statement 缓存某些情况下会出现prepared statement already exists的错误。这个坑在官方文档里没写是我在 GitHub issue 里翻到的。第三个坑连接池大小。max_size不要设太大PostgreSQL 默认最大连接数是 100你如果有 5 个服务实例每个池子 20 个连接加起来就 100 了再有点别的服务直接打满。我们的经验是每个实例max_size20足够因为 Checkpoint 的写入是短事务连接周转很快。3.3 Checkpoint 的存储结构长什么样PostgresSaver建的表主要有四张checkpoints、checkpoint_blobs、checkpoint_writes、checkpoint_migrations。核心是前两张。checkpoints表存的是每个 Checkpoint 的元信息thread_id、checkpoint_id、parent_checkpoint_id、checkpointJSONB 格式的 State 快照、metadata。checkpoint_blobs存的是大的二进制数据比如消息历史里的大文本。理解这个结构对排查问题很重要。有一次我们遇到恢复后 State 不对的问题直接查checkpoints表SELECT thread_id, checkpoint_id, parent_checkpoint_id, metadata FROM checkpoints WHERE thread_id your-thread-id ORDER BY checkpoint_id DESC LIMIT 10;发现parent_checkpoint_id链断了说明中间有 Checkpoint 没写成功。进一步查日志发现是某次写入时连接池超时事务回滚了。这个排查过程如果不懂表结构根本无从下手。3.4 thread_id 的设计直接决定恢复粒度thread_id是 Checkpoint 的隔离维度。同一个thread_id下的 Checkpoint 构成一条链恢复时从最新的开始。不同thread_id之间完全隔离。这里的设计决策很关键一个用户会话对应一个 thread_id还是一个任务对应一个 thread_id我们的做法是一个任务一个 thread_id格式是task_{task_id}。原因是用户会话可能包含多个任务如果共用一个 thread_id恢复的时候会串。而任务级别的隔离让每个任务的恢复互不影响。但这里有个反直觉的点同一个 thread_id 下你可以有多个并发的执行。LangGraph 用checkpoint_nsnamespace来区分。默认情况下checkpoint_ns是空字符串如果你用了子图subgraph子图会有自己的 namespace。这个机制在嵌套图场景下很有用但如果你不理解会出现明明写了 Checkpoint 却查不到的情况——因为查的时候没指定 namespace。4. 中断恢复的完整实现链路4.1 中断的两种类型主动中断与被动中断在讲实现之前得先把中断这个概念拆清楚。LangGraph 里的中断分两种。主动中断是用interrupt()函数显式触发的。比如人工审批场景from langgraph.types import interrupt def approval_node(state: AgentState): decision interrupt({ question: 是否批准执行该操作, details: state[intermediate_results], }) return {approval_status: decision}这个interrupt()一调用图的执行就暂停了当前 State 被写入 Checkpoint控制权返回给调用方。调用方拿到的是一个包含中断信息的对象可以展示给用户。用户做出决策后通过Command(resume...)恢复执行。被动中断是进程崩溃、网络断开、超时等外部原因导致的。这种中断没有显式的interrupt()调用但 Checkpoint 机制保证了最后一个成功的节点状态被持久化。恢复的时候从最后一个 Checkpoint 继续。这两种中断的处理方式不同主动中断需要业务层处理等待用户输入的逻辑被动中断只需要重新调用invoke或stream并传入相同的thread_id。4.2 主动中断的恢复Command 与 resume 值主动中断的恢复代码大概是这样from langgraph.types import Command config {configurable: {thread_id: task_123}} # 第一次执行会在 approval_node 中断 result graph.invoke({messages: [...]}, config) # 用户审批后恢复 result graph.invoke( Command(resumeapproved), config, )这里有个极易踩的坑Command(resume...)的值会作为interrupt()函数的返回值注入到approval_node里。也就是说decision变量的值就是approved。如果你在interrupt()之后还有代码那些代码会继续执行。我第一次用的时候以为恢复是重新执行整个节点结果发现节点前半部分的代码不会重跑只有interrupt()之后的代码会执行。这个行为其实很合理——因为interrupt()之前的副作用不应该重复执行——但文档里没强调我调试了好一会儿才想明白。注意interrupt()之前的代码如果有副作用比如写数据库恢复时不会重跑这是好事。但如果你在interrupt()之前做了状态更新那些更新在恢复时也不会重放因为 State 是从 Checkpoint 恢复的。所以所有状态更新都应该在interrupt()之后做或者确保更新是幂等的。4.3 被动中断的恢复幂等性是前提被动中断的恢复看起来简单——重新invoke就行——但前提是所有节点都是幂等的。因为恢复的时候最后一个 Checkpoint 之后的节点会重新执行如果那个节点有非幂等的副作用就会重复执行。我们的做法是给每个有副作用的节点加幂等键。比如调用外部写接口的节点def write_node(state: AgentState): idempotency_key f{state[task_id]}_{state[step]} # 先查是否已执行 if redis.get(fexecuted:{idempotency_key}): return {intermediate_results: {write: already_done}} # 执行写操作 result external_api.write(..., idempotency_keyidempotency_key) redis.setex(fexecuted:{idempotency_key}, 86400, 1) return {intermediate_results: {write: result}}这个幂等键的设计有个讲究它必须能唯一标识这一次执行。用task_id step是个常见方案但如果同一个 step 可能因为重试而执行多次就得再加一个重试计数。我们踩过的坑是某个节点因为超时被重试了三次幂等键没变结果第二次重试直接返回了已执行但实际上第一次执行是失败的。后来改成幂等键里带上执行尝试的序号才解决。4.4 恢复时的 State 合并逻辑恢复的时候LangGraph 会从 Checkpoint 加载 State然后从最后一个成功的节点继续。这里有个细节Checkpoint 里存的是节点的输入 State不是输出。也就是说如果节点 A 执行到一半崩溃了Checkpoint 里存的是 A 的输入恢复时 A 会重新执行。这个设计的原因是为了保证一致性如果存的是输出那 A 的副作用可能已经发生了但输出没写成功恢复时就会丢失。存输入的话A 重新执行副作用通过幂等键去重。理解这一点对调试很重要。有一次我们发现恢复后某个工具被调用了两次查 Checkpoint 才发现那个工具节点执行到一半崩溃了Checkpoint 存的是它的输入恢复时它重新执行了一遍。加上幂等键之后就正常了。5. AG-UI 如何把中断状态实时推到前端5.1 AG-UI 解决的是什么问题后端的中断恢复逻辑再完善如果前端不知道现在中断了、在等什么、恢复后是什么状态用户体验依然是一团糟。AG-UI 是一个专门为 Agent 前端交互设计的协议它定义了一套事件流把 Agent 的执行过程包括中断、恢复、工具调用、流式输出标准化地推给前端。在没有 AG-UI 之前我们的做法是自己定义 WebSocket 消息格式后端每执行一步就推一条消息。问题是这套格式是私有的前端每接一个新 Agent 就要重新适配。AG-UI 的价值在于协议标准化它定义了RUN_STARTED、TEXT_MESSAGE_CONTENT、TOOL_CALL_START、STATE_SNAPSHOT、RUN_FINISHED等事件类型前端只要实现一次就能对接所有兼容 AG-UI 的后端。5.2 中断事件在 AG-UI 里怎么表达AG-UI 里跟中断相关的事件主要是两类STATE_SNAPSHOT和自定义事件。STATE_SNAPSHOT是 State 的完整快照在关键节点推送。当 Agent 中断时后端会推一个STATE_SNAPSHOT里面包含当前 State 和中断信息。前端收到后根据中断类型展示不同的 UI如果是审批中断弹审批框如果是等待外部回调展示处理中。自定义事件用来传递业务特定的信息。比如我们的审批中断会推一个CUSTOM事件name是approval_requiredvalue里包含审批的问题和选项。前端监听这个事件渲染审批组件。这里有个实践中的坑STATE_SNAPSHOT是全量推送如果 State 很大比如消息历史很长推送的数据量会很大。我们的做法是在 State 里只放必要的字段大的中间结果存到对象存储State 里只放引用。这样STATE_SNAPSHOT的大小可控。5.3 恢复时前端怎么触发前端触发恢复的方式是发一个RUN请求带上thread_id和resume值。AG-UI 的协议里这个请求的格式大概是{ threadId: task_123, runId: run_456, resume: { interruptId: int_789, value: approved } }后端收到后调用graph.invoke(Command(resumevalue), config)然后把新的事件流推给前端。这里的关键是interruptId的传递。LangGraph 的interrupt()会返回一个包含interrupt_id的对象这个 ID 必须原样传回给后端否则恢复的时候找不到对应的中断点。我们一开始没传这个 ID结果恢复时报找不到中断点查了半天才发现是这个问题。5.4 流式输出与中断的配合AG-UI 的流式输出TEXT_MESSAGE_CONTENT和中断是配合使用的。Agent 在生成文本的时候前端实时展示如果生成到一半触发了中断流式输出暂停前端展示中断 UI恢复后流式输出继续。这个继续的实现有个细节LangGraph 的流式输出是基于 Checkpoint 的恢复后会从最后一个 token 继续不会重复输出已经推过的内容。但前提是前端要记录已经收到的 token 位置。我们的做法是前端维护一个messageId到lastTokenIndex的映射恢复时带上这个索引后端从对应位置继续推。实测下来这套机制在长文本生成的场景下体验很好用户几乎感觉不到中断和恢复的切换。但如果中断发生在工具调用期间恢复后工具会重新执行幂等键保证不重复副作用前端需要展示工具重新执行中的状态这个状态切换要处理好否则用户会困惑。6. 从手写 Loop 迁移的实操路线与踩坑记录6.1 迁移不是重写是逐步替换我的建议是不要一次性把手写 Loop 全部重写成 LangGraph风险太大。正确的做法是逐步替换先把状态管理抽出来用 LangGraph 的 State 定义替代散落的变量然后把主循环改成一个两节点的图agent tools最后逐步把分支逻辑改成条件边。我们当时的迁移分了三周。第一周只做 State 定义和 Checkpointer 接入主循环还是手写的但状态已经存到 PostgreSQL 了。第二周把主循环改成图跑通基本的 ReAct 流程。第三周加中断恢复和 AG-UI。每周都有可运行的版本出问题可以快速回滚。6.2 迁移中最容易出问题的三个地方第一个是 reducer 的语义。手写 Loop 里你直接list.append语义很明确。LangGraph 里 reducer 是隐式的你得理解多个节点更新同一个 key 时怎么合并。我们踩的坑是并行分支同时更新一个 dict后写的覆盖先写的。解决方案是用自定义 reducerdef merge_dict(left: dict, right: dict) - dict: return {**left, **right} class AgentState(TypedDict): intermediate_results: Annotated[dict, merge_dict]第二个是 Checkpoint 的写入时机。LangGraph 默认在每个节点执行后写 Checkpoint但这个行为可以通过checkpoint_during参数控制。如果你有节点执行很频繁但状态变化很小可以关掉它的 Checkpoint减少数据库压力。我们有个心跳节点每秒执行一次一开始没关 Checkpoint数据库写入量直接爆了。后来改成checkpoint_duringFalse问题解决。第三个是 thread_id 的生命周期管理。手写 Loop 里没有 thread_id 这个概念迁移后需要设计一套 thread_id 的生成和清理机制。我们的做法是任务创建时生成task_{uuid}任务完成后保留 7 天用于排查问题然后定时任务清理。如果不清理checkpoints表会无限增长查询性能会下降。6.3 性能实测数据我们做了一个对比测试同样的 ReAct 流程平均 5 轮工具调用手写 Loop 和 LangGraph PostgreSQL Checkpoint 的性能差异指标手写 LoopLangGraph PG Checkpoint单次执行耗时无中断3.2s3.8s单次执行耗时含 1 次中断恢复不支持4.5s进程崩溃后恢复状态丢失从最后 Checkpoint 恢复并发 100 任务内存占用 2.1GB内存占用 1.4GB可观测性仅日志完整执行链路可以看到无中断场景下 LangGraph 有约 18% 的性能开销主要来自 Checkpoint 的序列化和数据库写入。但这个开销换来的是中断恢复能力在真实业务里完全值得。而且并发场景下 LangGraph 的内存占用反而更低因为状态不在内存里堆积。6.4 几个值得分享的小技巧技巧一Checkpoint 的序列化用ormsgpack而不是pickle。LangGraph 默认用ormsgpack它比pickle快很多而且更安全。如果你有自定义的 State 类型确保它能被ormsgpack序列化否则会报错。我们有个 State 里放了datetime对象一开始序列化失败后来改成 ISO 格式的字符串才解决。技巧二用graph.get_state(config)调试。这个方法能拿到当前 thread 的最新 State 和下一个要执行的节点调试的时候非常有用。我经常在恢复逻辑里加一行日志打印get_state的结果确认恢复的起点对不对。技巧三Checkpoint 的metadata字段可以放业务信息。比如任务的用户 ID、来源渠道等。这些信息不参与图的执行但在排查问题时很有用。我们会在每次invoke时通过config传入 metadataCheckpoint 会自动带上。技巧四AG-UI 的事件要加序号。网络传输可能乱序前端需要根据序号重排。AG-UI 协议里事件本身没有序号字段我们在自定义事件里加了一个seq字段前端按这个排序。这个细节在弱网环境下很重要。7. 一些关于 Runtime 设计的个人体会从手写 Loop 到可恢复 Runtime最大的认知转变是Agent 的执行不是一个函数调用而是一个可以暂停、恢复、观测的长期过程。这个转变带来的不只是技术方案的改变还有思维方式的改变。手写 Loop 时代我关注的是怎么让这一次执行跑通。有了 Checkpoint 之后我关注的是这一次执行在任何时刻中断能不能恢复到一致的状态。这个视角的切换让我在设计每个节点的时候都会问自己这个节点的副作用幂等吗它的状态更新会不会在恢复时丢失它的执行时间会不会太长导致 Checkpoint 频繁写入另一个体会是中断恢复不是加一个功能而是改一种架构。你不能在现有系统上打补丁式地加中断恢复因为状态管理、幂等性、可观测性这些是相互关联的。LangGraph 的价值在于它把这些能力做成了框架的一部分你只要按它的模式写就自然获得了这些能力。最后说一个我踩过的坑不要过度依赖 Checkpoint 做业务状态存储。Checkpoint 是给图执行用的它的生命周期跟 thread 绑定。如果你把业务状态也塞进去会导致 State 膨胀、序列化变慢、恢复逻辑复杂。我们的做法是业务状态单独存一张表Checkpoint 里只放执行相关的状态两者通过task_id关联。这样职责清晰各自的生命周期也好管理。这套方案我们线上跑了半年多处理了大概几十万个任务中断恢复的成功率在 99.9% 以上。剩下的 0.1% 失败基本都是因为外部依赖不可用导致的跟 Checkpoint 机制本身无关。如果你也在做 Agent 的运行时我强烈建议尽早把状态管理和中断恢复做起来越晚做迁移成本越高。
返回列表