8.2. 用 LangGraph 实现带记忆断点恢复的工作流

你需要什么

  • 环境:Python 3.10+,已安装 langgraphlangchainletta-clientlanggraph-checkpoint-postgres(本地开发可用 langgraph-checkpoint-sqlite 作为轻量替代)。
  • 工具:一个可用的 LLM API 密钥(如 OpenAI),以及本地或远程的 Postgres 数据库(若使用 Sqlite 可跳过)。
  • 前置知识:了解 LangGraph 的基本概念——状态、节点、边,并阅读过上一章《从零搭建一个能记住一切的客服智能体》。
  • 预计时间:约 45–60 分钟。

最终成果

我们将构建一个多步请假审批智能体,它能记住每一次审批过程,并在任何节点暂停、等待人工干预,甚至宕机后仍能从中断处精确恢复。最终效果相当于一个“工作流 + 长期记忆”的系统:

  • 用户的每一次申请都会被永久记录在 Letta 记忆块中,跨会话可用。
  • 审批流程在每一个关键节点前自动暂停(human‑in‑the‑loop),等待人工输入“同意”或“驳回”。
  • 如果系统崩溃或主动挂起,重启后可以通过线程 ID 恢复至上一检查点,不会重复执行已通过的步骤。
  • 当某步发生错误(例如调用外部服务失败),系统能自动回退到上一个安全检查点,并利用记忆补偿丢失的上下文。

这样做是为了解决两个实际痛点:复杂业务流程不能全是自动决策,需要人类把关;而任何中断(服务器重启、用户掉线)都不应导致业务状态丢失。

步骤一:设计状态图与检查点

审批流程定义如下:

  1. 用户提交请假申请(开始)
  2. 经理审批(若金额 < 5000 元,直接通过,否则必须等待经理人工确认)
  3. 总监审批(当请假时间超过 3 天,自动触发,同样暂停等待)
  4. HR 备案(自动节点,但会记录到记忆)
  5. 完成

我们将这个流程建模为 LangGraph 的有向图 StateGraph,并用一个 TypedDict 定义共享状态:

from typing import TypedDict, Optional
from langgraph.graph import StateGraph, END
from langgraph.checkpoint.postgres import PostgresSaver
# 开发环境可替换为 from langgraph.checkpoint.memory import MemorySaver

class ApprovalState(TypedDict):
    user_id: str
    leave_days: int
    amount: float          # 请假涉及的费用(如有)
    manager_approved: Optional[bool]
    director_approved: Optional[bool]
    hr_recorded: bool

checkpointer 是断点恢复的核心。每执行完一个节点,LangGraph 都会将完整状态序列化并持久化到检查点存储中。生产环境必须使用 PostgresSaver,因为它支持并发安全和高可用(注意框:MemorySaver 仅在单次进程内有效,重启后丢失所有状态,绝不用于生产)。实例化时将 PostgresSaver 注入图:

with PostgresSaver.from_conn_string("postgresql://user:pass@localhost/db") as checkpointer:
    graph = StateGraph(ApprovalState)
    # ... 添加节点与边 ...
    compiled_graph = graph.compile(checkpointer=checkpointer)

interrupt_beforeinterrupt_aftercompile() 方法的参数,而非 add_node 的参数(重要修正)。我们希望在“经理审批”和“总监审批”两个节点执行前暂停,因此设置:

compiled_graph = graph.compile(
    checkpointer=checkpointer,
    interrupt_before=["manager_approval", "director_approval"]
)

这样,当流程运行到这些节点前,会自动抛出一个 GraphInterrupt 异常并保存当前状态。我们可以捕获该异常,向用户请求人工决策,然后通过 Command(resume=...) 恢复。

预期结果:运行第一步,状态应停留在经理审批前,数据库的 checkpoints 表会新增一条记录,包含序列化的 ApprovalState

步骤二:实现挂起与恢复逻辑

挂起点(interrupt)发生时,工作流从 invokeastream 中抛出异常。我们需要捕获它,并给出一个外部接口(如命令行、API 端点)供输入:

from langgraph.errors import GraphInterrupt
from langgraph.types import Command

config = {"configurable": {"thread_id": "user_req_123"}}

try:
    # 第一次调用:会停在第一个 interrupt 前
    result = compiled_graph.invoke(
        {"user_id": "u1", "leave_days": 7, "amount": 12000, "hr_recorded": False},
        config
    )
except GraphInterrupt:
    print("⏸️ 工作流已暂停,等待人工输入。")

# 模拟人类回答(通常来自队列、UI 回调等)
human_decision = "approve"  # 或者 "reject"

# 恢复执行:注意要传入同样的 config,保证 thread_id 一致
compiled_graph.invoke(Command(resume={"decision": human_decision}), config)

当执行恢复命令时,LangGraph 会从最近的检查点加载状态,并继续执行 manager_approval 节点。节点内部可以访问 Command.resume 传入的数据:

def manager_approval(state: ApprovalState, config: RunnableConfig):
    # 获取中断时传递的恢复数据(需要结合 LangGraph 的状态写入)
    decision = config.get("configurable", {}).get("__pregel_resume", {}).get("decision")
    # 实际实现中,我们会在 interrupt 后通过额外的状态 key 接收
    state["manager_approved"] = (decision == "approve")
    # 如果驳回,直接结束流程
    if not state["manager_approved"]:
        return {"manager_approved": False, "next": "END"}
    return state

注意:调研素材指出,子图命令式调用需要显式传递检查点配置才能恢复。但本章不涉及子图,我们的流程都是一个主图。若日后扩展,应优先将子图注册为父图节点,以自动继承检查点传播。

恢复执行后会继续往下走,如果请假天数超过 3 天,自动进入总监审批节点。因为我们也设置了 interrupt_before=["director_approval"],所以它会再次暂停。此时,用户可以用同样的 thread_id 再次发起恢复调用,工作流不会从零开始,而是精准地停在总监审批这一步。

预期结果:通过 thread_id 多次调用 invoke,每次只推进一个人类等待的节点,其余自动节点照常运行。

步骤三:错误恢复与记忆补偿

审批流程中,HR 备案节点若调用外部系统失败(如 API 超时),我们不能让状态丢失。LangGraph 本身的重试机制(retry 参数)可以处理暂时故障,但对于不可恢复的错误,我们可以回退到上一个检查点,同时利用 Letta 记忆补偿丢失的上下文。

我们集成 Letta 客户端,将长期记忆挂载到 config 中:

from letta import Letta

letta_client = Letta(base_url="http://localhost:8283")
memory_block = letta_client.create_block(label="approval_history", value="")

# 在节点中写入记忆
def hr_record(state, config):
    # 尝试调用外部系统
    try:
        send_to_hr_system(state)
    except Exception as e:
        # 将失败信息写入长期记忆,供重试或人工处理
        memory_block.value += f"\n[失败] 用户 {state['user_id']} 的备案失败:{e}"
        raise  # 仍抛出异常,让 LangGraph 捕获并回退

    # 成功则记录
    memory_block.value += f"\n[成功] 已备案用户 {state['user_id']} 请假 {state['leave_days']} 天"
    return {"hr_recorded": True}

当节点抛出异常后,LangGraph 会标记该执行失败,但不会丢弃检查点。我们可以根据错误类型选择恢复策略。如果是永久错误,我们可能希望跳过该节点,手动将状态修改为 hr_recorded=False 并继续。此时记忆块里保留了详细历史,人工可通过 Letta 查看并手动补偿。

若要回退到上一个安全点,我们可以手动重置 thread 的状态:通过检查点接口读取历史快照,将当前状态替换为最后一个成功的检查点,然后从那个位置恢复。LangGraph 的 checkpointer 提供 listget_tuple 方法:

# 获取 thread 的所有检查点
checkpoint_tuples = list(checkpointer.list(config))
last_success = None
for t in checkpoint_tuples:
    if t.metadata.get("step") != -1:  # -1 表示错误中断
        last_success = t
# 如果找到了安全点,则用其状态重新 invoke
if last_success:
    compiled_graph.invoke(last_success.checkpoint["channel_values"], config)

这种方式相当于手动补偿,适合灾难性故障后的恢复。结合 Letta 记忆,即使检查点里的部分数据丢失,也能从记忆块中找回历史决策依据。

预期结果:模拟 HR 系统故障后,可从数据库检查点列表中找到前一个正常节点状态,重放执行,并且 Letta 记忆块中始终保留着完整的审计线索。

回顾

  • 我们用 StateGraph + PostgresSaver 构建了带持久检查点的多步审批流程。
  • 通过 compile(interrupt_before=...) 在关键节点前自动插入等待点,实现人机协同。
  • 恢复时只需提供相同的 thread_idCommand(resume=...),工作流就能精准继续。
  • 遇到节点执行失败,我们借助检查点历史与 Letta 长期记忆,实现状态回退和上下文补偿,保证业务流程不丢失。
  • 总共花费约 50 分钟完成从零搭建到错误恢复验证。

行动清单

  1. 在 Postgres 中创建检查点表,并实例化 PostgresSaver
  2. 定义状态字典,画出你的业务流程图,映射为 LangGraph 节点。
  3. 使用 compile(checkpointer=..., interrupt_before=[...]) 设置暂停点。
  4. 捕获 GraphInterrupt 并设计用户输入接口,通过 Command 恢复。
  5. 集成 Letta 记忆,在关键节点写入日志,在失败时利用检查点历史回退。

下一章《构建自我进化的编码智能体》,我们将把记忆能力再推进一步:让智能体在编写代码时自动借鉴历史项目中的最佳实践,并像人类开发者一样从错误中学习。你刚才实现的断点恢复机制,将在编码智能体的长时间代码生成任务中大放异彩——无论是意外的语法错误还是模型输出超时,都能像本章的审批流程一样优雅重试。

本文章首发在 LearnKu.com 网站上。

上一篇 下一篇
讨论数量: 0
发起讨论 只看当前版本


暂无话题~