ARTICLE DETAIL

资讯详情

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

LangGraph实战:用断点恢复与幂等执行打造稳定的AI Agent工作流

LangGraph实战:用断点恢复与幂等执行打造稳定的AI Agent工作流 做 AI Agent 应用时最怕的不是模型回答错而是流程跑了一半突然进程崩溃、网络超时或者需要人工确认却不知道从哪个节点继续。前阵子我做一个订单退款审批的 Agent最初设计成“一跑到底”结果上线没几天就踩了大坑用户在最后一步触发人工审批等待期间服务重启所有上下文清空整个退款流程回到原点更离谱的是消息队列重试导致同一笔退款被提交了两次。后来把整个流程迁移到 LangGraph用断点恢复Checkpoint Interrupt把流程卡在“等待人工审批”的位置再用幂等执行Idempotent Execution挡住重复提交问题才彻底根治。这篇内容就是那次重构的记录适合已经会 LangChain 基础操作、正准备把 Agent 工程化落地的朋友。1. 为什么 Agent 流程需要“暂停”和“重试”1.1 断点恢复从“一跑到底”到“随时续跑”传统 LLM 应用的执行模型是线性的用户输入进来链式调用模型和工具最终输出结果。这套模型在“纯对话、不落库”的场景下没问题但一旦涉及多步骤业务逻辑问题就暴露了——中间任何一步需要人工介入整个链路就得拆开中间任何一步抛异常之前所有状态都会丢。LangGraph 的断点恢复本质上引入了“图执行引擎”的思维方式整个流程是一张有向图节点之间通过 State 传递数据。图可以随时在执行到某个节点时暂停把当前所有状态序列化到检查点Checkpoint然后退出进程。等外部条件满足比如人工审批完成再通过同一个线程 ID 从暂停的位置继续往下走。这里最关键的设计是线程 IDThread ID。你可以把 Thread ID 理解为这条业务流程的“病历本编号”所有状态快照、消息历史、每一步节点的输出都归档在这个编号下。恢复执行时LangGraph 会从检查点里拉出最新状态从断点处继续而不是从头重跑。1.2 幂等执行重复执行不等于重复副作用断点恢复解决的是“流程中断后能续跑”但续跑本身会带来一个新问题恢复操作可能被执行多次。比如用户等不及点了两次按钮或者消息队列把同一个任务投递了两次或者你的恢复接口被监控探活脚本重复调用。幂等执行要求的是同一个任务无论执行一次还是执行一百次对外部系统产生的最终效果必须等价于一次。听起来像数据库的“唯一约束”但落到 Agent 场景里要复杂得多。Agent 会调用 LLM、调用外部 API、写入数据库、发送通知这些动作都可能产生副作用。如果只靠“流程能恢复”而不做幂等保护重试一次就多扣一次款多发一封邮件多创建一条记录。在 LangGraph 里做幂等重点有两个位置一个是节点内部的状态判断另一个是外部工具调用的请求去重。前者让节点能识别“这个任务我已经处理过了”后者让外部系统能识别“这笔请求是重复的”。后面第 4 章我会给出具体写法。1.3 LangGraph 的底层设计State、Thread 与 Checkpoint要理解断点恢复的实现原理得先厘清 LangGraph 的三个核心概念。State 是一个可序列化的数据结构一般是 TypedDict 或 Pydantic 模型它就像一条流水线上的传送带每个节点从传送带上读取需要的信息处理完把结果放回传送带。Thread 是一组相关执行的会话标识通过配置里的 thread_id 指定。Checkpoint 则是某个时刻 State 的完整快照LangGraph 在每次节点执行完毕后自动保存。这三个概念加在一起形成了 LangGraph 的恢复模型Thread ID 定位到某条业务流程Checkpoint 保存该流程所有状态节点执行到断点处写入快照并暂停。恢复时不需要重放之前的 LLM 调用直接从快照里取上下文即可省时省力也不会偏离原状态。2. 核心 API 与工具选型2.1 安装与版本说明写作本文时我使用的版本是 langgraph 0.2.x 以上API 相比早期版本有一些调整。安装命令很简单pip install langgraph langchain-openai如果你还需要持久化数据库可以加装pip install langgraph-checkpoint-sqlite langgraph-checkpoint-postgres这里提醒一句LangGraph 的版本迭代挺快的早期文章里常见的interrupt_before/interrupt_after写法虽然还能用但官方已经把interrupt()函数作为主推方案。如果你看到网上旧教程代码跑不通先检查版本再对照官方 CHANGELOG。2.2 三种 Checkpoint 存储方案怎么选Checkpoint 的存储位置决定了断点恢复的可靠性。LangGraph 内置了多种 Checkpointer常用的有三种存储方案适用场景优点注意点MemorySaver本地调试、单元测试零配置速度最快进程重启后数据全丢不能用于生产SqliteSaver单机生产、小流量持久化到本地文件轻量可靠不支持跨实例共享PostgresSaver多实例部署、正式生产支持并发、跨实例共享状态需要维护数据库异步接口相对复杂我自己在本地调试阶段通常先用 MemorySaver逻辑验证完切到 SqliteSaver真正上线多实例时才用 PostgresSaver。很多新手直接在代码里用 MemorySaver结果服务一重启就“失忆”以为断点恢复失效其实是存储方案选错了。2.3 中断方式的演进从 interrupt_before 到 interrupt()早期版本中你需要在编译图的时候指定断点位置graph builder.compile(checkpointercheckpointer, interrupt_before[process_refund])这种方式的问题在于断点是编译期静态写死的运行时无法灵活决定要不要中断。比如你想“只有金额超过 1000 才人工审批”用interrupt_before写起来就很绕。现在推荐的做法是在节点内直接调用interrupt()函数def require_approval(state: RefundState): decision interrupt({ order_id: state[order_id], amount: state[amount], }) # decision 就是人工恢复时传入的值比如 approved 或 rejected return {approval_result: decision}interrupt()的特性是当图执行到这一行时会把传入的 payload 保存到检查点然后立即暂停。恢复时外部调用方通过Command(resume...)把值传回interrupt()函数则把这个值作为返回值继续执行下面的逻辑。这种动态中断方式让“人机协同”变得非常自然——需要审批就停不需要审批就继续。2.4 LangChain 和 LangGraph 的区别与协作这两者的关系经常被搞混。简单来说LangChain 提供的是“零件”模型封装、向量存储、工具调用、文档加载器、各种 Agent 的预置 Prompt。LangGraph 提供的是“装配线”一个状态化图执行框架让你精确控制流程的节点顺序、分支路由、循环条件、持久化和中断恢复。可以配合使用用 LangChain 封装 LLM 调用用 LangGraph 编排整个流程。比如在 LangGraph 节点里你可以直接调用ChatOpenAI(modelgpt-4o)也可以使用 LangChain 的 Tool 工具封装外部 API。它们不是替代关系而是互补关系。面试和实际工作里经常被问到的场景是“什么情况下只靠 LangChain 就够了”答案是流程固定、无分支、无人工介入、无状态恢复需求的链式调用。一旦需要循环、分支、人工审批、多轮工具调用后上下文还能跨会话恢复就直接上 LangGraph。3. 实战订单退款审批中的断点恢复3.1 业务场景与状态设计我拿一个具体的退款审批场景来做示例。用户发起退款请求Agent 先查询订单和用户信息然后进入人工审批环节审批通过后调用支付平台的退款接口最后发送通知。需求里有两个硬性约束金额超过 1000 元必须人工确认整个流程允许中断后恢复。State 设计如下from typing import TypedDict, List, Optional class RefundState(TypedDict): order_id: str user_id: str amount: float approval_result: Optional[str] # approved 或 rejected refund_response: Optional[dict] messages: List[str]字段里的approval_result就是人工审批环节写入的结果refund_response是外部退款接口的返回messages用于记录流程日志。所有节点都围绕这个 State 进行读写节点间通过 return 更新字段。3.2 图结构与节点定义整个图包含四个节点load_order负责加载订单和用户信息require_approval是断点所在节点调用interrupt()等待人工审批process_refund调用支付接口执行退款notify_user发送通知。节点定义如下from langgraph.graph import StateGraph, START, END from langgraph.graph.state import Command from langgraph.types import interrupt def load_order(state: RefundState): # 真实项目中这里会查询数据库 order_id state[order_id] return { amount: 1999.0, messages: [f订单 {order_id} 已加载金额 1999.0 元], } def require_approval(state: RefundState): if state[amount] 1000: return {approval_result: auto_approved} decision interrupt({ order_id: state[order_id], amount: state[amount], tip: 金额超过阈值需要人工审批, }) return {approval_result: decision} def process_refund(state: RefundState): if state[approval_result] ! approved: return {refund_response: {skipped: True}} # 这里调用支付平台退款接口 response {refund_id: RF state[order_id], status: success} return { refund_response: response, messages: [退款已发起 str(response)], } def notify_user(state: RefundState): if state.get(refund_response, {}).get(status) success: return {messages: [用户已收到退款通知]} return {messages: [流程结束未执行退款]}require_approval里的阈值判断是动态的低于 1000 自动通过不触发中断高于 1000 才停下等人审批。这正是interrupt()优于编译期静态断点的地方——同一套代码在不同金额下表现完全不同。3.3 编译、中断与首次执行接下来把节点连成图from langgraph.checkpoint.memory import MemorySaver checkpointer MemorySaver() builder StateGraph(RefundState) builder.add_node(load_order, load_order) builder.add_node(require_approval, require_approval) builder.add_node(process_refund, process_refund) builder.add_node(notify_user, notify_user) builder.add_edge(START, load_order) builder.add_edge(load_order, require_approval) builder.add_edge(require_approval, process_refund) builder.add_edge(process_refund, notify_user) builder.add_edge(notify_user, END) graph builder.compile(checkpointercheckpointer)执行首次调用时Agent 会一直运行到require_approval的interrupt()处暂停config {configurable: {thread_id: refund-order-001}} result graph.invoke( {order_id: A1001, user_id: u_42, messages: []}, config, ) print(result)此时程序并不会报错而是返回一个“暂停中”的状态。你可以通过graph.get_state(config)查看当前暂停位置和已保存的状态state graph.get_state(config) print(state.next) # 下一个要执行的节点例如 (require_approval,)state.next告诉你断点卡在哪里后续节点列表一目了然。这就是断点恢复的核心图没有结束只是被冻结了。3.4 人工审批与从断点恢复现在模拟人工审批环节。审批人打开后台看到interrupt()payload 里的订单信息和金额点击“批准”。代码侧通过Command(resume...)恢复执行from langgraph.types import Command resume_config {configurable: {thread_id: refund-order-001}} graph.invoke(Command(resumeapproved), resume_config)当Command(resumeapproved)传入后interrupt()函数会返回字符串approvedrequire_approval继续往下走把approval_result写入状态接着执行process_refund和notify_user。如果你想拒绝审批传入Command(resumerejected)即可process_refund会跳过退款流程正常结束。整个过程的上下文不会丢失也不需要重新调用前面的 LLM 节点。3.5 完整代码清单为了方便你直接跑通我把上面的代码合并成一个可执行脚本。测试时可以用MemorySaver但要注意进程重启即丢失from typing import TypedDict, List, Optional from langgraph.graph import StateGraph, START, END from langgraph.graph.state import Command from langgraph.types import interrupt from langgraph.checkpoint.memory import MemorySaver class RefundState(TypedDict): order_id: str user_id: str amount: float approval_result: Optional[str] refund_response: Optional[dict] messages: List[str] def load_order(state: RefundState): return {amount: 1999.0, messages: [f订单 {state[order_id]} 已加载]} def require_approval(state: RefundState): if state[amount] 1000: return {approval_result: auto_approved} decision interrupt({order_id: state[order_id], amount: state[amount]}) return {approval_result: decision} def process_refund(state: RefundState): if state[approval_result] ! approved: return {refund_response: {skipped: True}} return {refund_response: {refund_id: RF state[order_id], status: success}} def notify_user(state: RefundState): if state.get(refund_response, {}).get(status) success: return {messages: state[messages] [用户已收到退款通知]} return {messages: state[messages] [未退款流程结束]} checkpointer MemorySaver() builder StateGraph(RefundState) builder.add_node(load_order) builder.add_node(require_approval) builder.add_node(process_refund) builder.add_node(notify_user) builder.add_edge(START, load_order) builder.add_edge(load_order, require_approval) builder.add_edge(require_approval, process_refund) builder.add_edge(process_refund, notify_user) builder.add_edge(notify_user, END) graph builder.compile(checkpointercheckpointer) config {configurable: {thread_id: refund-order-001}} graph.invoke({order_id: A1001, user_id: u_42, messages: []}, config) # 模拟人工审批通过 graph.invoke(Command(resumeapproved), config) print(graph.get_state(config).values)如果你把金额改成 500 元流程会直接自动通过不会触发中断。这个差异正好解释了动态断点的价值。4. 幂等执行的落地策略4.1 状态标记法最简单可靠的幂等断点恢复解决了“流程能续跑”但没解决“重复执行不产生副作用”。最朴素的幂等策略是在 State 里加一个标志位节点执行前先检查处理过就直接跳过。以退款节点为例def process_refund(state: RefundState): if state.get(refund_processed): return {refund_response: state[refund_response]} # 执行外部退款调用 response call_payment_refund_api(state[order_id]) return { refund_processed: True, refund_response: response, }第二次执行这个节点时看到refund_processedTrue就不会再调用外部接口。这个思路跟数据库里的“幂等键”类似胜在简单直接。但要注意一点refund_processed只对“同一个线程 ID、同一个状态”生效如果是两个完全不同的线程各自执行同一笔订单退款标志位就没法拦截了。4.2 外部系统调用的幂等request_id 去重真正能拦住“跨线程重复执行”的手段是在外部系统层面做去重。调用支付、短信、邮件这些接口时生成一个全局唯一的 request_id作为参数传过去外部系统根据 request_id 判断是否已经处理过如果处理过直接返回原有结果不再执行扣款或发信。LangGraph 节点里可以这样设计import uuid def process_refund_with_idempotency(state: RefundState): request_id state.get(request_id) if not request_id: request_id freq_{uuid.uuid4().hex} # 外部接口通过 request_id 去重 response call_payment_refund_api( order_idstate[order_id], request_idrequest_id, ) return {request_id: request_id, refund_response: response}这样即使两个不同线程同时发起退款只要外部系统认 request_id重复请求会被自动丢弃或直接返回原结果。这是支撑生产环境幂等最可靠的一层因为 State 标志位只能覆盖单线程内部外部接口去重才能真正兜底跨线程场景。4.3 update_state 与 as_node人工修正状态生产环境还有个很有意思的场景外部退款接口因为第三方系统故障一直没返回但人工已经通过线下渠道完成了退款操作。这时候再让 Agent 去调用退款接口就变成重复执行了。LangGraph 允许你直接修改状态假装某个节点已经执行过from langgraph.graph.state import Command graph.update_state( config, values{ refund_processed: True, refund_response: {refund_id: manual_001, status: success}, messages: [人工线下退款已完成], }, as_nodeprocess_refund, )as_nodeprocess_refund的意思是这组状态更新被视为process_refund节点的输出。后续节点会认为退款已经执行完毕直接进入通知环节不再重跑支付接口。这招在人工兜底、故障降级时非常实用相当于给人留了一个“手动改写流程记录”的后门。4.4 重放与幂等的组合场景断点恢复加上幂等执行最常见的组合是消息队列重放。假设 Agent 服务收到一个退款任务执行到process_refund前突然进程崩溃。消息队列基于“至少一次投递”原则过几秒会重新投递同一条消息。如果新投递的任务使用了全新的线程 ID那么状态标志位就失效了必须靠外部接口的 request_id 兜底。相反如果任务里携带了固定的任务 ID你可以在入口处把它映射成 LangGraph 的 thread_id那么即便重启、重放状态也会落到同一个线程节点内部的refund_processed就能生效。我见过不少团队在“至少一次投递”场景下没做线程 ID 映射导致同一笔业务被重复处理。最后总结出的实践是消息里的业务主键比如订单号直接作为 thread_id而不是每次发消息都生成新 UUID这样断点恢复和幂等标志位才能在同一套状态上下文里发挥作用。5. 常见问题与排查技巧5.1 恢复后状态丢失症状流程中断后再次调用graph.invoke(Command(resume...))报错或者从头执行断点位置完全丢失。原因九成是用 MemorySaver 部署到了生产环境进程重启后所有检查点烟消云散。解决方法是换成 SqliteSaver 或 PostgresSaver。代码改动很小from langgraph.checkpoint.sqlite import SqliteSaver with SqliteSaver.from_conn_string(checkpoints.db) as checkpointer: graph builder.compile(checkpointercheckpointer) # 在这里调用 graphSqliteSaver 会把检查点写入本地文件进程重启后依然能恢复。多实例部署则需要 PostgresSaver注意配置连接串和异步接口。5.2 同一个 thread 的脏状态症状同一个线程 ID 重复执行后State 里残留了上一次的approval_result、refund_response导致新任务直接跳过某些节点。原因通常是业务逻辑里复用了同一个 thread_id而任务本身应该是全新的。解决方式是每个新任务用新的 thread_id或者在入口处对 State 做一次清理。调试时可以先graph.get_state(config)检查当前值确认是不是上个任务留下的“脏数据”。5.3 并发与检查点竞争症状两个并发的请求同时操作同一个 thread_id后写入的检查点覆盖前者导致状态丢失或任务重复。LangGraph 的 Checkpoint 是基于“版本号”的底层存储为 Postgres 时可以配置条件更新但多实例并发更新同一个 thread 仍然建议做串行化处理比如在业务层对 thread_id 加分布式锁。大多数情况下一条业务线程同时被两处操作的频率不高但一旦出现排查起来非常头疼。5.4 调试技巧get_state_history 是神器如果你想知道这个线程之前每一步执行了什么不要只盯着get_state用get_state_history查看完整历史快照config {configurable: {thread_id: refund-order-001}} for snapshot in graph.get_state_history(config): print(next:, snapshot.next) print(values:, snapshot.values) print(---)它能列出该线程从开始到现在每一步的节点执行顺序、节点输出、当前状态。这在排查“为什么恢复后走向不对”时非常高效相当于给 Agent 流程装了行车记录仪。5.5 版本差异坑老接口 vs 新接口网上很多教程还是老写法interrupt_before[process_refund]然后在恢复时直接graph.invoke(None, config)。这个写法在 0.2.x 版本还能用但如果在新版本中混用interrupt()会出现混淆。我的建议是统一使用新 API节点内用interrupt()暂停恢复用Command(resume...)编译时不再传interrupt_before或interrupt_after。旧接口的语义是“在某个节点前后强制暂停”新接口的语义是“在流程中间某个点动态等待输入”后者明显更适合人机协同场景。最后再分享一个我实操中摸索出来的组合套路业务入口用消息主键作为 thread_id关键外部调用节点一律带 request_id数据库存储直接上 Persistent Checkpointer调试阶段善用get_state_history回溯每一步。这套组合下来无论是断点恢复还是幂等执行都稳得很。尤其是线程 ID 映射这个细节很多人一开始没注意等到重放、重试、人工审批混合在一起时才意识到它的重要性——早设计早省心。
返回列表