Agent 工程体系 · 第 70/98 篇。内容以 2026 年 9 月可验证的公开规范和稳定接口为基线;框架版本敏感能力会明确标注,不把实验行为写成通用保证。

Agent Checkpoint 与恢复:快照、增量、版本迁移和副作用重放

Agent 一旦需要运行数分钟、数小时甚至数天,失败就不再是“请求返回错误”这么简单。进程可能在模型调用后崩溃,工具请求可能已经发出但响应尚未收到,多个 Agent 可能并行修改同一份状态,代码部署后旧任务又可能在新版本上继续执行。

这类系统需要回答四个问题:

  1. 已经完成的工作保存在哪里?
  2. 恢复时从哪里继续,而不是从头开始?
  3. 状态结构和流程代码升级后,旧任务还能否继续?
  4. 一个已经执行过、但恢复时可能再次执行的副作用,如何避免重复?

Checkpoint 是这四个问题的共同基础,但它不是简单的“把 Python 对象存进数据库”。一个可靠的恢复系统,至少要同时处理:

  • 可恢复状态;
  • 执行位置;
  • 并发任务;
  • 版本兼容;
  • 外部副作用;
  • 重放时的确定性;
  • 持久化提交本身的失败。

LangGraph 的持久化模型可以作为一个具体实现来理解:它把线程状态保存为 checkpoint,在图的 super-step 边界持久化完整状态,并额外保存 super-step 内部节点的 pending writes,以支持部分完成后的恢复。生产系统还需要将这些机制与事件日志、请求键、工具调用幂等和数据库去重结合起来。(docs.langchain.com)

一、先区分四种数据:状态、快照、事件和副作用

1. Agent 状态不是全部业务事实

设 Agent 在时刻 tt 的可恢复状态为:

St=(Ct,Pt,Mt,Rt,Vt)S_t = (C_t, P_t, M_t, R_t, V_t)

其中:

  • CtC_t:对话上下文,例如消息、当前任务;
  • PtP_t:流程位置,例如下一步节点、等待中的人工审批;
  • MtM_t:模型或工具产生的中间结果;
  • RtR_t:恢复所需的运行元数据,例如重试次数、租约、请求键;
  • VtV_t:状态版本和业务流程版本。

但不是所有运行时对象都应该放入 StS_t。以下对象通常不适合作为 checkpoint 内容:

  • 数据库连接;
  • HTTP 客户端;
  • Python 协程和线程;
  • 文件句柄;
  • 未序列化的模型客户端;
  • 包含密钥的临时对象;
  • 只在当前进程有效的内存缓存。

Checkpoint 应保存“恢复所需的数据”,而不是保存“当前进程的全部内存”。

在 LangGraph 中,图状态由 state schema 和 reducer 组成,节点返回的是对状态的更新,而不是任意地修改共享对象。状态可以使用 TypedDict、dataclass 或 Pydantic 模型定义;多个节点并行写入同一 channel 时,reducer 决定这些更新如何合并。(docs.langchain.com)

2. 快照是某个时刻的完整状态

快照是某个一致性边界上的完整状态:

Ki=(Si,Pi,Mi,metai)K_i = (S_i, P_i, M_i, \text{meta}_i)

它至少需要包含:

  • 状态值;
  • 下一步要执行的节点或任务;
  • 当前 checkpoint 标识;
  • 父 checkpoint 标识;
  • 创建时间;
  • 产生该状态的节点写入;
  • 任务错误和中断信息。

快照的核心优势是恢复简单:

Recover(Ki)从 Ki 指定的位置继续执行\text{Recover}(K_i) \Rightarrow \text{从 } K_i \text{ 指定的位置继续执行}

代价是存储量可能随状态大小和运行步数增长。假设第 ii 个快照大小为 sis_i,保存 nn 个快照的空间复杂度是:

O(i=1nsi)O\left(\sum_{i=1}^{n} s_i\right)

如果消息列表不断追加,每个快照都复制完整消息列表,那么单个线程的总存储量可能接近二次增长。

3. 增量只保存状态变化

增量保存的是从上一个状态到当前状态的变化:

Δi=SiSi1\Delta_i = S_i \ominus S_{i-1}

恢复时需要从基准快照开始折叠增量:

Sn=K0Δ1Δ2ΔnS_n = K_0 \oplus \Delta_1 \oplus \Delta_2 \oplus \cdots \oplus \Delta_n

增量的优势是写入小,尤其适合追加型 channel,例如消息、事件、工具调用记录。代价是读取和恢复需要重放多个增量;如果增量链过长,恢复延迟和损坏传播范围都会增加。

因此实际系统通常采用:

基准快照 K0
  + 增量 Δ1
  + 增量 Δ2
  + ...
  + 增量 Δm

mm 超过阈值时重新生成快照:

K0 + Δ1 + ... + Δm
              ↓ compaction
             K1

LangGraph 默认会在每个 super-step 保存各 state channel 的完整值;文档同时提供了 DeltaChannel,用于只保存增量,当前文档标记其需要 langgraph>=1.2 且处于 beta 状态,因此不能把它当作跨版本稳定契约使用。(docs.langchain.com)

4. 事件日志记录“发生过什么”

事件日志不是状态快照的替代品,而是另一种记录:

RunStarted
NodeStarted
ModelCalled
ToolRequested
ToolSucceeded
StatePatched
HumanApproved
RunFailed

事件通常是追加写入的:

E=(e1,e2,,en)E = (e_1, e_2, \ldots, e_n)

通过事件重建状态:

Sn=fold(S0,E)S_n = \operatorname{fold}(S_0, E)

快照和事件日志的差异如下:

维度 快照 事件日志
记录内容 当前结果 状态变化和事实
恢复方式 直接加载 从某个位置折叠事件
读取速度 通常较快 取决于日志长度
审计能力 有限 较强
修改历史 通常不可变 通常追加不可变
外部副作用证明 不充分 可记录请求、回执和结果

最稳妥的工程组合通常是:

事件日志:记录事实和外部交互
Checkpoint:保存最近可恢复状态
增量:降低频繁保存的成本

但要注意:事件日志只能证明“系统记录了某件事”,不能自动证明外部系统确实执行了这件事。外部系统仍需要回执、查询接口或幂等协议。

二、Checkpoint 的一致性边界:为什么节点边界很重要

1. Checkpoint 不是任意代码行的快照

假设一个节点执行三步:

读取数据库
调用模型
写入支付系统

如果进程在第三步之后、节点返回之前崩溃,运行时可能只有两种选择:

  • 从节点开始重跑;
  • 依赖节点内部更细粒度的任务结果。

它通常不能像调试器一样从 Python 函数的某一行继续。LangGraph 的 Graph API 以 super-step 为 checkpoint 边界;恢复时从发生故障的节点开始。节点越大,失败时重复执行的工作越多。(docs.langchain.com)

因此,节点粒度实际上决定了恢复粒度:

重复工作量故障节点中已完成但未持久化的工作\text{重复工作量} \approx \text{故障节点中已完成但未持久化的工作}

2. Super-step 中的并行写入

考虑如下图:

flowchart LR
    START --> A[解析任务]
    A --> B[查询库存]
    A --> C[计算优惠]
    B --> D[汇总]
    C --> D
    D --> END

BC 并行执行的 super-step 中,可能出现:

B 成功并写入
C 执行失败

如果只在整个 super-step 结束后保存完整快照,那么 B 的结果会丢失,恢复时需要重新执行 B

如果运行时先持久化节点级写入:

checkpoint_writes:
  B -> {"stock": ...}

恢复时就可以复用 B 的结果,只重试 C。LangGraph 的持久化模型同时保存 super-step 边界的完整 checkpoint 和 super-step 内部的 task writes,后者用于 pending writes recovery。(docs.langchain.com)

这里有一个重要边界:

节点级写入不是完整的 StateSnapshot,不能把它当作任意时间点的用户可恢复快照。

它解决的是“同一 super-step 中已经成功的任务不要重复计算”,而不是“允许用户从任意 Python 指令位置恢复”。

3. 三种持久化耐久性

系统通常有三种 checkpoint 写入策略:

async:后台写入,吞吐和延迟较好,但进程立即崩溃时可能丢最近写入
sync:当前步骤等待 checkpoint 写入完成,再继续
exit:只在流程结束时写入

LangGraph 文档将默认的异步持久化描述为在后台写 checkpoint;同时提供同步和仅退出时写入的模式。(docs.langchain.com)

可以用一个简单的故障窗口理解:

W=[tstate changed,tcheckpoint committed]W = [t_{\text{state changed}}, t_{\text{checkpoint committed}}]

如果进程在窗口 WW 内崩溃,那么最近一次状态变化是否可恢复,取决于持久化模式。

  • sync:尽量使 WW 接近零,但会增加步骤延迟;
  • asyncWW 非零,但吞吐较好;
  • exit:整个运行期间都可能没有中间恢复点,不适合长流程。

这不是“哪个模式最好”,而是丢失工作量与每步延迟之间的取舍。

三、一个可运行的最小恢复示例

下面的例子使用 LangGraph 的 Graph API 和内存 Checkpointer,演示:

  • 通过 thread_id 关联同一个执行线程;
  • 在节点边界保存状态;
  • 发生异常后使用同一线程恢复;
  • 用请求键避免工具副作用重复提交。

内存 saver 只适合演示,因为进程重启后数据会丢失;生产环境应使用持久化数据库实现。(docs.langchain.com)

from typing import TypedDict

from langgraph.checkpoint.memory import InMemorySaver
from langgraph.graph import END, START, StateGraph


class State(TypedDict, total=False):
    order_id: str
    amount: int
    payment_request_id: str
    payment_status: str
    charged_amount: int
    fail_once: bool


class FakePaymentGateway:
    def __init__(self):
        self.charges: dict[str, int] = {}

    def charge(self, request_id: str, amount: int) -> dict:
        # 模拟幂等支付接口:
        # 相同 request_id 再次调用时,返回原结果而不重复扣款。
        if request_id in self.charges:
            return {
                "status": "already_charged",
                "amount": self.charges[request_id],
            }

        self.charges[request_id] = amount
        return {"status": "charged", "amount": amount}


gateway = FakePaymentGateway()


def prepare_payment(state: State):
    request_id = f"charge:{state['order_id']}"
    return {
        "payment_request_id": request_id,
        "payment_status": "prepared",
    }


def charge_payment(state: State):
    # 模拟“第一次执行到这里时进程崩溃”。
    if state.get("fail_once"):
        raise RuntimeError("worker crashed before payment result was saved")

    result = gateway.charge(
        request_id=state["payment_request_id"],
        amount=state["amount"],
    )

    return {
        "payment_status": result["status"],
        "charged_amount": result["amount"],
    }


builder = StateGraph(State)
builder.add_node("prepare_payment", prepare_payment)
builder.add_node("charge_payment", charge_payment)
builder.add_edge(START, "prepare_payment")
builder.add_edge("prepare_payment", "charge_payment")
builder.add_edge("charge_payment", END)

checkpointer = InMemorySaver()
graph = builder.compile(checkpointer=checkpointer)

config = {
    "configurable": {
        "thread_id": "order-thread-1001",
    }
}

initial_state = {
    "order_id": "order-1001",
    "amount": 500,
    "fail_once": True,
}

try:
    graph.invoke(initial_state, config)
except RuntimeError as exc:
    print(f"first run failed: {exc}")

snapshot = graph.get_state(config)
print("after failure:", snapshot.values)
print("next nodes:", snapshot.next)

# 模拟修复故障:后续恢复不再触发 fail_once。
graph.update_state(config, {"fail_once": False})

result = graph.invoke(None, config)
print("after recovery:", result)
print("gateway charges:", gateway.charges)

预期输出的关键部分类似:

first run failed: worker crashed before payment result was saved
after failure: {'order_id': 'order-1001', ..., 'payment_status': 'prepared'}
next nodes: ('charge_payment',)
after recovery: {..., 'payment_status': 'charged', 'charged_amount': 500}
gateway charges: {'charge:order-1001': 500}

这个例子中有三个关键事实。

第一,prepare_payment 已经成功并形成 checkpoint,所以恢复时不需要重新准备请求键。thread_id 是加载和恢复线程状态的索引;没有它,checkpointer 无法知道应该读取哪条执行历史。(docs.langchain.com)

第二,故障发生在 charge_payment 节点内,因此恢复从该节点开始,而不是从函数中间开始。

第三,即使支付网关已经扣款、但 Agent 在保存结果前崩溃,恢复时仍可能再次调用网关。payment_request_id 让网关可以识别重复请求,把“至少一次调用”转化为“至多一次业务效果”。

示例中的 FakePaymentGateway 使用内存字典实现去重。真实系统应把去重表放在数据库中,例如:

CREATE TABLE tool_call_receipts (
    request_id      VARCHAR(255) PRIMARY KEY,
    tool_name       VARCHAR(100) NOT NULL,
    request_hash    CHAR(64) NOT NULL,
    status          VARCHAR(32) NOT NULL,
    response_json   JSONB,
    created_at      TIMESTAMPTZ NOT NULL DEFAULT now()
);

工具调用逻辑应满足:

1. 根据业务主键生成 request_id;
2. 计算请求参数哈希;
3. 查询 request_id;
4. 已成功:直接返回历史回执;
5. 已存在但参数哈希不同:拒绝,防止键冲突;
6. 不存在:执行外部操作;
7. 保存结果和状态;
8. 后续重复调用返回同一回执。

注意第 6 步和第 7 步之间仍然存在经典的双写窗口:

外部支付成功
进程崩溃
本地回执尚未写入

所以真正可靠的外部工具不能只依赖本地 checkpoint。它至少需要以下一种能力:

  • 外部 API 原生支持幂等键;
  • 可以通过业务查询接口确认结果;
  • 采用带状态机的“创建请求—轮询结果”协议;
  • 使用事务性消息或 outbox,避免业务状态和待发送事件分离。

四、恢复不是回到原始代码行,而是重放到下一个边界

1. 恢复和重放的区别

恢复通常指从最近一次失败的 checkpoint 继续当前运行。

重放指选择一个历史 checkpoint,从该位置重新执行后续步骤。它用于:

  • 诊断错误;
  • 回归测试;
  • 比较新旧提示词或模型;
  • 从某个历史状态创建分支;
  • 人工修改状态后继续执行。

LangGraph 支持通过历史 checkpoint_id 进行 time travel 和 replay;checkpoint 之前的节点使用已保存结果,之后的节点重新执行。文档特别指出,重放会再次触发后续的模型调用、API 请求和中断,因此副作用必须具备幂等性或被隔离。(docs.langchain.com)

设历史状态序列为:

K0 -> K1 -> K2 -> K3 -> K4

K2 重放时:

K0、K1、K2:直接读取历史结果
K3、K4:重新执行

如果 K3 是“发送邮件”,那么重放可能再次发送邮件。Checkpoint 只保存 Agent 的状态,不会自动撤销已经发送到 SMTP 服务的邮件。

2. 为什么副作用必须被包装成任务

在 Functional API 中,@task 将一个离散工作单元的结果持久化,使恢复时可以复用已经完成的任务结果。任务结果需要可序列化;当任务已经完成并写入 checkpoint,恢复可以读取结果,而不是重新计算。(docs.langchain.com)

错误写法:

@entrypoint(checkpointer=checkpointer)
def workflow(inputs: dict):
    send_email(inputs["to"])  # 非幂等副作用直接放在流程函数中
    answer = interrupt("请审批")
    return answer

恢复时,流程可能从 entrypoint 开头重新运行,于是 send_email 再次执行。

更安全的结构是:

from langgraph.checkpoint.memory import InMemorySaver
from langgraph.func import entrypoint, task
from langgraph.types import interrupt


@task
def send_email_once(message_id: str, to: str, body: str) -> str:
    # 真实实现应调用支持 message_id 幂等的邮件服务,
    # 或在本地 outbox / receipt 表中做去重。
    print(f"send email: {message_id} -> {to}")
    return message_id


@entrypoint(checkpointer=InMemorySaver())
def workflow(inputs: dict):
    send_email_once(
        message_id=inputs["message_id"],
        to=inputs["to"],
        body=inputs["body"],
    ).result()

    approved = interrupt({
        "type": "approval",
        "message": "是否继续?",
    })

    return {
        "approved": approved,
    }

但“包装成 task”不等于自动获得 exactly-once。它主要解决的是:

  • 任务结果可以被 checkpoint;
  • 已完成任务在恢复时可以复用;
  • 工作流中的非确定性结果可以固定下来。

如果任务已经开始执行、但在任务结果持久化前失败,任务仍可能再次运行。因此外部副作用仍需要请求键、查询确认或幂等写入。LangGraph 文档也明确要求将 API 调用和文件写入等副作用放入任务,并设计幂等操作或使用幂等键。(docs.langchain.com)

3. 中断前后的副作用边界

如果节点逻辑是:

def approve_order(state):
    create_audit_record()
    approved = interrupt("请审批")
    charge_card()
    return {"approved": approved}

恢复时,包含 interrupt 的节点可能从节点开头重新执行。因此:

  • create_audit_record() 会重复;
  • charge_card() 在审批返回后执行,通常不会因为同一个中断再次触发,但仍需考虑节点重试;
  • 审批前的写入必须幂等;
  • 更好的设计是把审计写入、审批等待、扣款拆为不同节点。

LangGraph 的中断文档明确指出:中断前的副作用会在恢复时重新执行,应使用幂等操作,或者将副作用放到中断之后、拆到独立节点中。(docs.langchain.com)

五、增量 Checkpoint 的正确使用方式

增量保存最适合以下状态:

class State(TypedDict):
    messages: list
    tool_events: list
    current_step: str

如果每次只追加一条消息,可以记录:

{
  "messages": {
    "append": [
      {
        "role": "tool",
        "name": "inventory",
        "content": "in_stock"
      }
    ]
  }
}

但并不是所有字段都适合增量。

1. 适合增量的字段

  • 追加型事件;
  • 只增不减的审计记录;
  • 工具调用历史;
  • 大型消息序列;
  • 可通过顺序合并的日志。

这些字段通常有结合律:

(ab)c=a(bc)(a \oplus b) \oplus c = a \oplus (b \oplus c)

如果还满足交换律:

ab=baa \oplus b = b \oplus a

那么并行合并更容易处理。

2. 不适合直接增量的字段

  • 当前余额;
  • 当前租约持有者;
  • 任务状态;
  • 订单状态;
  • “最后一次模型答案”;
  • 需要删除或覆盖的配置;
  • 依赖严格版本顺序的状态。

例如余额从 100 改为 80,增量不能简单记录“减 20”,因为重试、重复应用或乱序都会产生错误。更安全的方式是记录带版本的写入:

{
  "account_id": "a-1",
  "expected_version": 7,
  "new_balance": 80
}

数据库使用乐观锁:

UPDATE accounts
SET balance = 80,
    version = version + 1
WHERE account_id = 'a-1'
  AND version = 7;

如果影响行数为 0,说明状态已被其他执行者修改,不能静默覆盖。

3. 增量链必须可压缩

增量系统至少需要:

  • 基准快照;
  • 增量序号;
  • 父版本或父哈希;
  • 校验和;
  • 压缩策略;
  • 并发写入检测;
  • 损坏后的回退点。

一个简单的增量记录可以是:

{
  "thread_id": "t-100",
  "checkpoint_id": "cp-42",
  "parent_id": "cp-41",
  "sequence": 42,
  "state_version": 3,
  "delta": {
    "current_step": "charge_payment",
    "tool_events_append": [
      {"request_id": "charge:order-1001", "status": "started"}
    ]
  },
  "sha256": "..."
}

恢复过程先检查:

sequence 是否连续
parent_id 是否匹配
sha256 是否正确
state_version 是否可迁移

任一步失败,都不能继续盲目折叠后续增量。应回退到最近的完整快照或进入人工修复流程。

六、版本迁移:Checkpoint 是持久化 API

Checkpoint 的结构一旦落盘,就成为代码和数据之间的 API。代码升级不能只考虑“新请求能否运行”,还必须考虑:

旧状态 -> 新代码
旧节点名 -> 新节点名
旧字段 -> 新字段
旧业务流程 -> 新业务流程

LangGraph 默认让正在运行和恢复中的线程使用当前部署的图代码,而不是把线程永久绑定到启动时的代码版本。这使得修复可以立即作用于旧线程,但也意味着每次部署都可能影响已有 checkpoint。(docs.langchain.com)

1. 技术兼容和业务兼容

技术兼容关心旧数据能不能被新代码加载和执行:

  • 节点是否还存在;
  • state key 是否还存在;
  • 类型是否仍可读取;
  • 新增字段是否有默认值;
  • 条件路由是否仍能处理旧状态。

业务兼容关心旧任务是否应该继续使用旧业务逻辑:

  • 旧订单是否应该继续旧计费规则;
  • 已经审批一半的流程是否应该走新风控;
  • 新字段是否意味着流程分支变化;
  • 新模型是否会改变旧任务的决策。

二者不能混为一谈。代码可以技术上成功加载旧 checkpoint,但业务语义已经变化。

2. 安全的字段迁移

假设旧版本:

class State(TypedDict):
    order_id: str
    amount: int

新版本需要增加 currency。不要直接写成必填字段:

class State(TypedDict):
    order_id: str
    amount: int
    currency: str

旧 checkpoint 没有 currency,恢复时可能无法满足新 schema。更安全的是:

from typing import NotRequired
from typing_extensions import TypedDict


class State(TypedDict):
    order_id: str
    amount: int
    currency: NotRequired[str]

节点中提供兼容默认值:

def normalize_order(state: State):
    currency = state.get("currency", "CNY")
    return {"currency": currency}

LangGraph 的兼容性文档建议新增字段使用 NotRequired 或可选默认值;删除字段应先保留一段排空周期;重命名则采用“新增—双写或双读—确认旧线程耗尽—删除”的过程。(docs.langchain.com)

3. 节点重命名不能直接做

旧状态可能保存:

{
  "next": ["charge_payment"]
}

如果新版本把节点改名为 capture_payment,恢复时运行时仍可能尝试查找 charge_payment,结果找不到节点。

安全做法不是直接删除旧节点,而是保留兼容入口:

builder.add_node("charge_payment", charge_payment_compat)
builder.add_node("capture_payment", capture_payment)

兼容节点可以:

  • 调用新实现;
  • 将旧状态转换为新状态;
  • 写入新的流程版本;
  • 将后续路由切换到新节点。

等确认所有可能停在旧节点的线程都完成后,再删除旧节点。

4. 用流程版本控制业务分支

如果新旧业务逻辑必须并存,应在流程开始时写入版本:

def intake(state: State):
    if "flow_version" in state:
        return {}

    return {
        "flow_version": 2,
    }

之后使用显式路由:

def route_after_triage(state: State):
    if state.get("flow_version", 1) == 1:
        return "legacy_path"
    return "new_path"

重要的是,版本必须在需要分支之前写入 checkpoint。否则一个尚未写入版本的旧线程可能在升级后被错误地解释为新流程。

LangGraph 的兼容性文档给出的迁移模式也是:旧线程从状态中的流程版本读取并继续旧路径,新线程写入新版本并走新路径;待旧线程排空后再删除版本分支。(docs.langchain.com)

七、多 Agent 恢复:避免共享状态的隐式覆盖

多 Agent 系统通常有三种状态:

主编排器状态
子 Agent 私有状态
跨 Agent 共享业务状态

它们不应全部放在同一个可变字典中。

1. 子 Agent 状态隔离

例如:

Supervisor
├── Researcher
├── Coder
└── Reviewer

Researcher 的中间搜索结果通常属于 Researcher 的 checkpoint namespace;Coder 不应该通过读取同一个可变字段来猜测 Researcher 当前执行到了哪里。

LangGraph 为子图 checkpoint 提供 namespace,用来区分父图和子图的执行状态;跨线程、跨 Agent 的长期共享数据则更适合使用 Store,而不是依赖某个线程的短期 checkpoint。(docs.langchain.com)

2. 共享写入必须有并发条件

错误模型:

state["balance"] -= 10

两个 Agent 同时读取 100:

Agent A 读取 100,计算 90
Agent B 读取 100,计算 80
A 写入 90
B 写入 80

最终结果是 80,但实际应该是 70。

安全模型需要将更新表达为原子操作或带版本条件的命令:

{
  "operation": "debit",
  "account_id": "a-1",
  "amount": 10,
  "request_id": "run-1:debit:1"
}

由拥有该业务状态的服务执行:

UPDATE accounts
SET balance = balance - 10,
    version = version + 1
WHERE account_id = 'a-1'
  AND balance >= 10;

然后用 request_id 做去重。Checkpoint 负责恢复 Agent 的执行位置,不能替代业务数据库的并发控制。

3. 租约和恢复的关系

当一个 worker 处理线程时,需要租约:

线程 t-100
租约持有者:worker-A
租约过期时间:10:30:00

worker 崩溃后,另一个 worker 可以在租约过期后接管。但接管不代表旧 worker 一定停止;网络分区时,旧 worker 可能仍在运行。因此必须防止“失去租约的 worker 继续写入”。

常见办法是将租约版本作为 fencing token:

worker-A 获得 token=41
worker-B 接管后获得 token=42

所有 checkpoint 和业务写入都要求 token 不小于当前有效 token。token=41 的旧 worker 即使恢复,也不能覆盖 token=42 的新写入。

八、常见错误和诊断方法

错误一:把 InMemorySaver 当生产存储

现象:

开发环境恢复正常
进程重启后 thread_id 仍存在
但历史状态完全消失

原因是内存 saver 只保存于当前进程。生产环境应使用数据库支持的 checkpointer,例如 SQLite 用于本地流程,Postgres 用于生产持久化;具体实现和安装包应以目标版本文档为准。(docs.langchain.com)

错误二:只保存状态,不保存执行位置

如果只保存:

{"messages": [...], "order_id": "..."}

却没有保存:

{"next": ["charge_payment"]}

恢复时无法判断:

  • 哪些节点已经完成;
  • 哪些节点还未开始;
  • 哪些并行分支已经写入;
  • 当前是否在等待人工输入。

状态和程序计数器必须作为一个一致性单元保存。

错误三:把“调用过工具”当作“工具成功”

工具调用应至少区分:

requested
started
succeeded
failed
unknown

其中 unknown 很重要:

请求已发出
客户端没有收到响应

此时不能直接重试非幂等操作,也不能直接标记失败。应通过 request_id 查询外部系统:

查询成功 -> 记录 succeeded
查询不到 -> 根据协议决定重试
无法确认 -> 转人工或进入补偿流程

错误四:恢复时随机数和当前时间重新生成

如果控制流依赖:

if random.random() < 0.5:
    ...

或者:

deadline = datetime.now() + timedelta(minutes=5)

恢复时重新计算可能得到不同分支。非确定性值应在任务中生成并持久化,后续重放读取原结果,而不是再次生成。LangGraph Functional API 文档将随机数、外部 API 结果等非确定性操作列为需要任务封装的场景。(docs.langchain.com)

错误五:把 checkpoint 当作回滚

Checkpoint 可以让 Agent 回到过去的状态,但不能自动撤销:

  • 已发出的邮件;
  • 已创建的订单;
  • 已扣除的余额;
  • 已提交的代码;
  • 已发送到第三方的 HTTP 请求。

回滚 Agent 状态和补偿外部副作用是两个不同动作:

状态回退:让流程从旧 checkpoint 重新执行
业务补偿:调用退款、撤销、删除或反向记账接口

如果外部系统没有补偿接口,就不能假设 time travel 能提供事务回滚。

九、生产设计的验证闭环

一个可上线的恢复机制,不能只测试“正常运行后读取 checkpoint”。至少应验证以下故障路径:

1. 节点执行前进程崩溃
2. 节点执行中进程崩溃
3. 外部请求成功、回执保存前崩溃
4. checkpoint 写入前进程崩溃
5. 两个 worker 同时获取同一线程
6. 租约过期后旧 worker 继续运行
7. 代码升级后旧线程恢复
8. 从历史 checkpoint 重放
9. checkpoint 增量链中间损坏
10. 人工中断后重复恢复

每条故障路径都应定义可观测结果:

thread_id
run_id
checkpoint_id
parent_checkpoint_id
state_version
flow_version
node_name
task_id
request_id
attempt
lease_token
tool_receipt_status

诊断时不要只看最终回答。应沿着下面的链路查询:

run
 -> checkpoint history
   -> node/task writes
     -> tool request_id
       -> external receipt
         -> business record

例如用户报告“订单被扣了两次”,需要同时回答:

  • Agent 是否重试了工具节点?
  • 两次请求是否使用了同一个 request_id?
  • 外部支付系统是否支持幂等?
  • 本地回执表是否有唯一约束?
  • 两次扣款是否真的对应两个业务请求?
  • 是否存在两个 worker 同时持有租约?

只有 checkpoint 记录,通常不足以回答最后几个问题。

十、核心边界:Checkpoint 提供的是可恢复执行,不是魔法级 exactly-once

可以把 Agent 执行抽象成:

Step(Si)(Si+1,Oi)\text{Step}(S_i) \rightarrow (S_{i+1}, O_i)

其中:

  • SiS_i 是可恢复状态;
  • Si+1S_{i+1} 是下一状态;
  • OiO_i 是对外部系统产生的副作用。

Checkpoint 能可靠保存的是 SiS_i 或其增量。对于 OiO_i,如果外部系统与 checkpoint 存储不在同一个事务中,就无法仅靠 checkpoint 保证:

外部副作用恰好执行一次\text{外部副作用恰好执行一次}

工程上更现实的目标是:

至少一次尝试+幂等业务效果\text{至少一次尝试} + \text{幂等业务效果}

或者:

可重试调用+可查询回执+可补偿状态机\text{可重试调用} + \text{可查询回执} + \text{可补偿状态机}

这也是快照、增量、版本迁移和副作用重放必须放在一起讨论的原因:

  • 快照决定从哪里恢复;
  • 增量决定如何高效保存变化;
  • 版本迁移决定旧状态能否被新代码解释;
  • 副作用幂等决定重放是否会造成真实业务损害;
  • 事件日志和回执决定系统能否解释“到底发生了什么”。

当这几个层次被分别建模后,Agent 才不只是“失败后再调用一次模型”,而是一个能够检查点化、恢复、迁移、审计和安全重放的持久化执行系统。


系列导航与关联阅读

官方资料

本文依据 Agent、模型、协议与框架官方资料重新梳理;正文、示例与生产清单由 WR BLOG 编写。