跳到正文
Elaine Blog
返回

用 LangGraph 实现可恢复任务:完整 Python 实战与源码解读

更新于:
Agent Runtime

前面的命令行 Runtime 把状态和支付回执保存到 SQLite,但下一步执行哪个函数仍由我们手动决定。LangGraph 可以把这种控制流写成图,并保存图执行状态。

本篇做一件具体的事:第一次命令创建退款任务并停在审批,进程退出;第二次命令加载同一份数据库,提交决定,再继续执行。示例不调用模型,因此读者可以先理解恢复,再把模型节点接进来。

图中保存什么,业务库保存什么

图状态只保存 task ID 与流程状态。任务库保存金额、订单版本、审批与回执,模拟支付库保存支付事实。它们之间通过稳定身份关联。

如果把所有数据都复制进图状态,恢复后可能拿着过期订单继续执行。如果什么都不保存,只存一段模型总结,又无法核对审批。因此要按事实来源决定存储位置。

class RefundState(TypedDict):
    # 图通过任务 ID 读取权威业务记录,避免在多个节点各自创造申请身份。
    task_id: str
    # 此字段是图的执行投影;最终支付结果仍读取业务库与支付回执。
    status: str

完整文件包含导入、图构建与命令行:durable_refund_graph.py下载。正文片段用于解释,不是另一套独立实现。

interrupt 为什么会再次经过节点开头

审批节点读取任务,将订单和金额传给 interrupt。第一次调用在这里暂停;后续用 Command 的 resume 值继续时,节点会从头重新进入,interrupt 才返回恢复值。

因此 interrupt 前不能无保护地创建支付或发送通知。即使暂停只发生一次,恢复也可能再次经过这段代码。示例这里只读任务,审批落库放在恢复值之后,并用 decision ID 去重。

# 这段数据用于展示审批对象;真实身份必须由外部登录和权限系统验证。
decision = interrupt({"task_id": task["id"], "amount_minor": task["amount_minor"]})
# 字符串 "false" 不是布尔 False,不能使用 bool(decision) 做审批判断。
if not isinstance(decision, dict) or type(decision.get("approved")) is not bool:
    raise ValueError("审批必须携带布尔 approved")

如果人工拒绝,业务状态明确变为 rejected,图结束。不能只把控制流跳到 END,却留下 awaiting_approval,让界面永远显示等待。

文件型 Saver 才能跨进程找回暂停点

InMemorySaver 适合进程内演示,进程退出后无法依靠它恢复。这里使用 SqliteSaver 的文件连接,并在每次命令中重新打开同一个数据库。

# with 保持连接在图调用期间有效;离开时关闭连接,而数据库文件继续存在。
with SqliteSaver.from_conn_string(".runtime-demo/checkpoints.sqlite") as saver:
    graph = build_graph(store, saver)
    # 恢复必须同时复用 thread_id、数据库和兼容的图定义。
    config = {"configurable": {"thread_id": "T17"}}

完整示例使用 sync durability,使图的持久化与执行顺序采取同步策略。它增加等待开销,也不能让图数据库与支付服务成为同一个事务。官方持久化文档解释了 Checkpoint 与线程的关系。

按三个独立命令观察

在配套 Python 目录中,使用三个独立命令观察同一任务:

# 第一次进程:创建任务并在审批节点暂停。
python durable_refund_graph.py submit T17
# 第二次进程:读取同一暂停点,提交虚构审批并执行模拟退款。
python durable_refund_graph.py approve T17
# 第三次进程:同时读取业务记录与图快照,观察二者的责任范围。
python durable_refund_graph.py status T17

默认目录相同,不要在不同工作目录误建两份数据库。想研究支付响应丢失,先用配套 refund_cli 的 lose-response 选项;图示例正常路径不默认注入故障。

如果某节点执行异常而留下待执行任务,可以使用 retry 命令以 None 继续已有图;如果停在人工 interrupt,仍必须明确提供审批值。两者不是同一种恢复输入。

从 StateGraph 读到实际执行器

源码固定为提交 e539ac1,不根据未来 main 分支猜实现。graph/state.py定义状态图的节点、边、状态通道和编译过程。StateGraph 是构建描述,compile 后的对象才负责运行。

随后看 pregel/_loop.py。它协调步骤、已有 Checkpoint 和待处理写入;Runner 负责具体任务执行。这样能理解为什么画一条边还不等于一切都只是普通函数递归。

图按 superstep 推进。并行节点在一个阶段产生写入,再由状态合并规则形成后续可见状态。若多个节点写同一字段,就需要前文解释的 Reducer,不能指望框架猜测业务合并规则。

为什么还要保存 pending writes

同一阶段两个节点并行,其中一个成功、另一个失败。如果只保存整阶段最终快照,恢复可能重复已经成功的工作。实现中保留与任务相关的中间写入,供恢复时识别已经完成的部分。

继续看 SqliteSaver:get_tuple 读取 Checkpoint 及相关信息,put 保存快照,put_writes 保存任务写入。不同接口对应不同恢复需求。

但若外部工具成功、写入尚未保存就崩溃,仍存在未知窗口。pending writes 缩小某些重复执行范围,不能替远端服务提供幂等。

到达 END 是否代表业务成功

不一定。图可能因拒绝、证据不足或结果未知而结束。示例同时打印 graph_values 与 business,明确区分图投影和权威业务状态。

如果支付结果是 reconciling,应使用业务对账流程查询原回执。图没有无限循环重试,也没有为了显示成功把未知改成 completed。生产系统可以把对账设计为定时节点或独立工作流。

SqliteSaver 的源码说明它面向轻量同步场景。真实并发部署需要选择合适持久后端、处理会话并发与版本升级。学习本篇的重点是理解“什么被保存、哪里会重执行、外部事实由谁确认”,然后再选择规模化方案。


分享这篇文章:

上一篇
生产中的 Runtime 怎样设计:状态、执行、存储与版本升级
下一篇
读懂 Temporal:历史重放怎样驱动长期业务流程