Agent 工程体系 · 第 69/98 篇。内容以 2026 年 9 月可验证的公开规范和稳定接口为基线;框架版本敏感能力会明确标注,不把实验行为写成通用保证。
Agent 持久化执行:事件日志、Checkpoint、租约、恢复和确定性
Agent 一旦需要运行数分钟、数小时,或跨越人工审批、外部系统回调和进程重启,就不能再把一次调用看成“函数从头执行到尾”。它更接近一个持续推进的状态机:
其中任何一步都可能遇到进程崩溃、网络超时、重复调度、模型结果变化、数据库不可用或两个 Worker 同时接管同一个任务。
持久化执行要解决的不是“把变量保存下来”这么简单,而是要回答五个问题:
- 系统如何记录 Agent 已经发生过什么?
- 系统如何快速找到可以恢复的位置?
- 多个 Worker 如何避免同时执行同一个任务?
- 崩溃后哪些步骤应该重做,哪些步骤绝不能重做?
- 恢复时如何保证控制流不会因为随机数、时间或模型调用变化而走向另一条路径?
事件日志、Checkpoint、租约、恢复和确定性分别回答这些问题,但它们不是互相独立的组件,而是一套有因果关系的执行协议。
一、先建立执行模型:Agent 不是函数,而是可暂停状态机
1. Agent 执行的基本状态
设一个 Agent 实例的持久化状态为:
其中:
- :输入,例如用户请求、任务参数和租户信息;
- :工作记忆,例如消息、计划、工具返回值;
- :控制状态,例如当前节点、下一步动作和重试次数;
- :持久化进度,例如已经完成的任务、Checkpoint 标识;
- :恢复所需的信息,例如幂等键、外部请求 ID 和租约代数。
Agent 的一个节点可以抽象为:
其中:
- 是新状态;
- 是执行事件;
- 是外部副作用,例如发邮件、扣库存、写数据库或调用第三方 API。
关键点在于:状态变化和外部副作用不一定发生在同一个存储系统中。Checkpoint 可能成功写入数据库,但邮件已经发送后进程才崩溃;也可能外部 API 已经完成,但写入结果的过程超时。这种“部分成功”正是持久化执行最难处理的地方。
2. 逻辑时间和物理时间
Agent 的执行通常同时存在两种时间:
- 逻辑时间:第几个节点、第几个事件、第几个 Checkpoint;
- 物理时间:真实时间、超时时间、租约到期时间和外部服务的响应时间。
恢复主要依赖逻辑时间。物理时间只能用于判断超时、重试和租约是否失效,不能直接决定 Agent 业务状态。
例如,不能因为“当前时间已经超过 10 分钟”就推断“支付一定失败”。正确做法是查询支付系统,或者使用带幂等键的状态查询接口。
二、事件日志:记录发生过什么,而不是只记录现在是什么
1. 事件日志的定义
事件日志是按顺序追加记录执行事实的持久化结构。一个事件通常包含:
{
"event_id": "evt_00042",
"run_id": "run_8f2",
"seq": 42,
"event_type": "tool.completed",
"node": "reserve_inventory",
"payload": {
"request_id": "order_123:reserve",
"reserved": true
},
"created_at": "2026-09-01T10:20:31.120Z"
}
事件描述的是已经发生的事实,例如:
run.startednode.startedtask.completedtool.requestedtool.completedinterrupt.waitingcheckpoint.committedrun.failedrun.resumed
事件日志的核心性质是追加写入。已提交事件不应被原地修改;纠正错误应追加一个新事件,例如 reservation.cancelled,而不是把原来的 reservation.created 改成别的内容。
2. 事件日志和普通日志的区别
普通应用日志主要服务于人类排障:
reserve inventory succeeded
事件日志服务于机器恢复:
{
"event_type": "task.completed",
"task_id": "reserve_inventory",
"idempotency_key": "order_123:reserve",
"result_ref": "inventory_reservation_abc"
}
普通日志可能缺少:
- 唯一事件 ID;
- 严格的顺序号;
- 任务输入;
- 任务结果;
- 状态版本;
- 幂等键;
- 关联的外部请求 ID。
因此,不能把应用日志文件直接当作恢复依据。日志可读,不代表可重放。
3. 事件溯源不是所有 Agent 都必须采用的方案
事件溯源要求把当前状态视为事件序列的折叠结果:
其中 reduce 是确定性的状态转换函数。
例如:
order.created
payment.authorized
inventory.reserved
shipment.requested
可以折叠为:
{
"order_id": "order_123",
"payment": "authorized",
"inventory": "reserved",
"shipment": "requested"
}
但长时间运行的 Agent 如果每次恢复都从第一条事件开始重放,成本会随着历史长度增长。因此生产系统通常采用:
也就是说:
- 事件日志适合保留完整事实和审计轨迹;
- Checkpoint 适合快速恢复;
- 两者结合后,恢复成本从 降为 ,其中 是最近 Checkpoint 之后的事件数量。
LangGraph 的持久化模型以 Checkpoint 为中心:图状态按线程保存,并可用于故障恢复、人工介入、时间旅行和状态检查;它同时保存节点级的中间写入,用于同一 super-step 内的失败恢复。这更接近“Checkpoint 加执行记录”,不应简单等同于完整的事件溯源系统。(docs.langchain.com)
三、Checkpoint:恢复游标,而不是数据库备份
1. Checkpoint 的定义
Checkpoint是执行过程中某个可恢复边界上的状态快照。它至少需要能回答:
从哪里恢复?
恢复时的状态是什么?
哪些任务已经完成?
当前图或工作流版本是什么?
可以形式化为:
其中:
run_id标识一次运行;version标识状态结构和执行逻辑版本;position_i表示逻辑执行位置;state_i是快照状态;completed_i是已经成功完成的任务集合;metadata_i用于恢复、调试和审计。
Checkpoint 的位置必须是一致性边界。如果只保存了半个 Python 对象,或者保存了“节点已经开始”但没有保存任务结果,恢复时就无法判断该节点应该继续、重做还是跳过。
2. 节点边界和 super-step 边界
在图执行模型中,Checkpoint 常见于两种边界:
- 节点边界:一个节点完成后保存;
- super-step 边界:同一轮中所有可并行节点完成后保存。
LangGraph 的 Graph API 在每个 super-step 后生成完整 Checkpoint;一个 super-step 内的节点结果还会以任务级写入保存。因此,如果同一轮中节点 A 成功、节点 B 失败,恢复时可以保留 A 的结果,而不必重新执行 A。(docs.langchain.com)
例如:
START
├── fetch_customer
└── fetch_order
↓
compose_answer
第一轮中:
fetch_customer 成功
fetch_order 失败
理想恢复状态不是:
两个节点都重新执行
而是:
复用 fetch_customer 的 pending write
重试 fetch_order
成功后进入 compose_answer
这一区别直接影响:
- 外部 API 调用次数;
- 模型费用;
- 数据库压力;
- 重复副作用风险;
- 恢复延迟。
3. 快照、增量和事件日志的关系
三种存储方式可以这样区分:
| 方式 | 保存内容 | 恢复成本 | 存储成本 | 主要风险 |
|---|---|---|---|---|
| 全量快照 | 完整状态 | 低 | 高 | 大状态频繁复制 |
| 增量快照 | 状态变化 | 中 | 中 | 增量链损坏或合并复杂 |
| 事件日志 | 执行事实 | 高 | 通常较高 | 重放逻辑必须稳定 |
全量快照适合状态较小、恢复要求高的 Agent。增量快照适合消息和列表持续追加的场景。LangGraph 文档目前将 DeltaChannel 描述为保存增量而非每次保存完整累积值的能力,并明确标注需要 langgraph>=1.2 且仍处于 beta,接口未来可能变化,因此不能把它当作跨版本稳定契约。(docs.langchain.com)
一个实用的混合布局是:
event_log:
保存所有业务事实和执行事实
checkpoint:
每 N 个事件或每个稳定边界保存状态快照
task_result:
保存已完成任务的结果和幂等键
outbox:
保存待发送的外部副作用
恢复时:
读取最近 Checkpoint
→ 读取其后的事件
→ 合并已完成任务
→ 检查未确认副作用
→ 继续执行
四、租约:避免两个 Worker 同时推进同一个 Agent
1. 为什么只靠队列去重不够
假设任务队列投递了:
run_123
Worker A 取到任务后开始执行,但在处理过程中发生网络抖动。队列认为消息未确认,于是把任务重新投递给 Worker B。
此时可能出现:
Worker A:仍在调用外部支付 API
Worker B:开始恢复同一个 run_123
如果没有并发控制,两个 Worker 可能:
- 同时写入不同的 Checkpoint;
- 重复发送邮件;
- 一个 Worker 覆盖另一个 Worker 的状态;
- 按不同路径修改同一个订单;
- 互相重试,形成执行风暴。
因此需要租约。
2. 租约的定义
租约是带过期时间的排他执行权:
其中:
run_id:被保护的 Agent 实例;owner:持有租约的 Worker;epoch:租约代数,也叫 fencing token;expires_at:租约过期时间。
Worker 只有在租约有效且代数匹配时,才有权提交状态:
UPDATE agent_runs
SET state = :state,
checkpoint_id = :checkpoint_id,
updated_at = CURRENT_TIMESTAMP
WHERE run_id = :run_id
AND lease_owner = :owner
AND lease_epoch = :epoch
AND lease_expires_at > CURRENT_TIMESTAMP;
如果影响行数为 0,说明 Worker 已经失去租约,必须停止提交结果。
3. 为什么需要 fencing token
仅使用过期时间仍然不够。考虑以下时序:
t1 Worker A 获取租约,epoch=7
t2 A 长时间暂停,租约过期
t3 Worker B 获取新租约,epoch=8
t4 A 恢复,尝试写入旧结果
如果存储层只检查 run_id,A 的旧写入可能覆盖 B 的新状态。
因此每次重新获取租约都必须递增 epoch:
A: epoch=7
B: epoch=8
所有状态写入和外部资源写入都应携带 epoch,旧 Worker 的操作由下游拒绝:
这就是 fencing。它解决的不是“谁以为自己持有锁”,而是“旧持有者即使苏醒,也无法继续产生有效写入”。
4. 租约续期并不等于无限执行
Worker 执行期间通常需要续租:
租约 TTL = 60 秒
每 20 秒续租一次
但续租失败时不能继续执行。否则 Worker 可能在失去所有权后仍然调用外部系统。
生产实现至少需要:
- 本地停止标志;
- 续租失败立即触发取消;
- 提交 Checkpoint 时再次校验租约;
- 外部请求携带 fencing token 或幂等键;
- 任务重新入队由新的 Worker 接管。
租约只能保证同一时刻尽量只有一个合法执行者,不能保证进程崩溃前已经完成的外部副作用自动撤销。
五、恢复:从最后一个安全边界继续,而不是盲目从头重跑
1. 恢复的分类
恢复至少分为三类:
进程恢复
Worker 崩溃,新的 Worker 继续同一个 run_id。
节点恢复
某个节点失败,保留其他节点已经完成的结果,只重试失败节点。
业务恢复
外部系统状态不确定,例如支付接口超时,需要查询最终状态,而不是简单再次支付。
三者的恢复策略不同。把它们都实现成“捕获异常后重新调用整个 Agent”,通常会导致重复执行和状态污染。
2. 一个完整的恢复算法
设最近提交的 Checkpoint 为 ,其之后存在任务记录:
task A: completed
task B: started
task C: absent
恢复步骤可以写成:
关键不是“从 Checkpoint 开始执行所有节点”,而是:
Checkpoint 状态
+ 已持久化的任务结果
+ 外部副作用确认状态
= 可安全继续的执行上下文
3. 失败路径示例
考虑订单 Agent:
flowchart TD
A[读取订单] --> B[检查库存]
B --> C[创建支付意图]
C --> D{人工审批}
D -->|批准| E[创建发货单]
D -->|拒绝| F[取消订单]
E --> G[完成]
执行到人工审批前:
Checkpoint K1:
order_status = pending_approval
inventory = reserved
payment_intent_id = pi_123
此时进程崩溃,恢复流程应当是:
1. 使用同一个 run_id 加载 K1
2. 重新获取租约
3. 发现 payment_intent_id 已存在
4. 不再创建新的支付意图
5. 继续等待审批,或处理新的审批输入
如果恢复逻辑只看到“创建支付意图节点尚未返回成功”,就再次调用创建接口,可能生成 pi_124,造成重复支付意图。正确做法是把外部资源 ID 写入持久状态,并以它作为幂等恢复依据。
六、确定性:恢复同一次运行时,控制流必须可解释
1. 确定性的定义
对 Agent 而言,确定性不是要求每次运行都得到相同答案,而是要求:
给定同一个已持久化的运行历史,恢复执行时应沿着同一条已记录的控制路径继续,不能因为重新读取了随机数、当前时间或未固定的外部结果而产生无法解释的分叉。
区分两个概念:
- 跨运行确定性:相同输入的两次运行结果相同;
- 同一运行恢复确定性:一次运行暂停后恢复,行为与暂停前的逻辑历史一致。
LLM Agent 通常难以保证跨运行确定性,但可以通过保存任务结果、模型响应和工具结果,尽量保证同一运行恢复确定性。LangGraph 的 Functional API 文档明确要求将随机性和非确定性操作封装在 task 中,以便恢复时读取已经保存的任务结果,而不是重新计算。(docs.langchain.com)
2. 非确定性来源
常见来源包括:
random.choice(...)
datetime.now()
uuid.uuid4()
os.environ["FEATURE_FLAG"]
llm.invoke(...)
external_api_call(...)
下面的代码不适合直接放在可恢复控制流中:
def workflow(state):
if random.random() > 0.5:
return call_a(state)
return call_b(state)
因为第一次执行可能走 A,恢复后重新计算随机数却走 B。此时 Checkpoint 中的状态也许是合法的,但控制流已经改变。
正确方式是把非确定性动作变成可持久化任务:
from langgraph.func import entrypoint, task
from langgraph.checkpoint.memory import InMemorySaver
import random
checkpointer = InMemorySaver()
@task
def choose_route() -> str:
return "a" if random.random() > 0.5 else "b"
@task
def call_a() -> str:
return "result from a"
@task
def call_b() -> str:
return "result from b"
@entrypoint(checkpointer=checkpointer)
def workflow(_: dict) -> str:
route = choose_route().result()
if route == "a":
return call_a().result()
return call_b().result()
这里的关键不是使用了 @task 这个装饰器,而是:
随机选择发生在可记录的任务中
任务结果进入持久化执行上下文
恢复时优先读取已完成任务的结果
LangGraph 的 task 输出需要可序列化,以支持 Checkpoint 和恢复;文档同时建议将 API 调用和副作用放入 task,并设计为可重复执行或幂等。(docs.langchain.com)
3. 模型调用也需要确定性边界
把模型调用直接放在复杂控制流中,会产生如下问题:
response = model.invoke(prompt)
if "approve" in response.content.lower():
...
恢复时如果再次调用模型:
- prompt 可能已经变化;
- 检索结果可能变化;
- 模型版本可能变化;
- temperature 或服务端配置可能变化;
- 上游工具结果可能变化。
更稳妥的设计是把模型调用视为一个有输入、有输出、有版本的任务:
{
"task_id": "classify_001",
"model": "model-name",
"prompt_hash": "sha256:...",
"input_hash": "sha256:...",
"output": {
"decision": "approve",
"reason": "..."
}
}
恢复时:
任务已完成 → 使用已保存输出
任务未完成 → 重新调用,但必须使用同一个幂等键
这不意味着模型结果永远正确,而是让系统能区分:
第一次模型判断是什么
恢复时是否重新判断
重新判断使用了哪一个版本
七、副作用:Checkpoint 成功不代表外部世界已经成功
1. 最危险的时序
下面的代码看似合理:
def send_email_and_update(state):
send_email(state["email"])
return {"status": "email_sent"}
但实际时序可能是:
1. send_email 成功
2. 进程崩溃
3. {"status": "email_sent"} 没有写入 Checkpoint
4. 恢复后再次 send_email
因此,Checkpoint 只能证明状态已保存,不能证明外部副作用没有重复发生。
这类问题通常属于“至少一次执行”语义:
只要任务可能在“外部操作成功、内部确认失败”后重试,就必须使用:
- 幂等键;
- 外部状态查询;
- 去重表;
- Outbox;
- 可补偿操作;
- 人工核对。
2. 幂等键的基本结构
幂等键应来自业务事实,而不是 Worker 随机生成:
order_123:send_confirmation_email
order_123:create_payment_intent
order_123:reserve_inventory:v2
接口端保存:
CREATE TABLE idempotency_records (
idempotency_key TEXT PRIMARY KEY,
operation TEXT NOT NULL,
request_hash TEXT NOT NULL,
status TEXT NOT NULL,
response_json TEXT,
created_at TIMESTAMP NOT NULL
);
执行逻辑:
1. 根据 idempotency_key 查询记录
2. 已成功:直接返回保存的 response
3. 处理中:查询外部系统或等待
4. 不存在:插入 processing 记录
5. 调用外部 API
6. 保存 succeeded 和 response
还必须校验 request_hash。如果同一个幂等键对应了不同请求参数,不能直接复用旧结果,否则会把两个业务操作错误地合并。
3. Outbox 解决的是“状态和待发送事件的一致提交”
如果 Agent 需要在数据库中更新状态并发送消息,可以使用 Outbox:
BEGIN;
UPDATE orders
SET status = 'approved'
WHERE order_id = 'order_123';
INSERT INTO outbox (
event_id,
aggregate_id,
event_type,
payload,
status
) VALUES (
'evt_order_123_approved',
'order_123',
'order.approved',
'{"order_id":"order_123"}',
'pending'
);
COMMIT;
后台发送器再读取 outbox.status = 'pending' 的记录并发送。发送成功后更新为 sent。
这仍然可能重复发送,因此消费者也要幂等。Outbox 的作用是避免:
数据库更新成功,但消息没有留下发送记录
它不是“消息只发送一次”的证明。
八、LangGraph 中的持久化执行
LangGraph 将图定义、运行状态、Checkpoint 和恢复机制分离:
- 图定义描述节点和边;
- 状态 schema 描述节点间共享的数据;
- Checkpointer 保存线程状态;
thread_id指向同一个可恢复执行上下文;- 节点或 task 执行外部操作;
- 恢复时使用相同的
thread_id读取已有状态。
官方文档将 LangGraph 定位为面向长时间运行、有状态 Agent 的低层编排运行时,核心能力包括持久化、可恢复执行、人工介入和确定性步骤与 Agent 步骤混合。(docs.langchain.com)
1. 最小可运行示例
前置条件:
python -m pip install -U langgraph
示例:
from typing import TypedDict
from langgraph.checkpoint.memory import InMemorySaver
from langgraph.graph import StateGraph, START, END
class State(TypedDict):
count: int
message: str
def increment(state: State):
return {
"count": state["count"] + 1,
}
def format_message(state: State):
return {
"message": f"count={state['count']}",
}
builder = StateGraph(State)
builder.add_node("increment", increment)
builder.add_node("format_message", format_message)
builder.add_edge(START, "increment")
builder.add_edge("increment", "format_message")
builder.add_edge("format_message", END)
checkpointer = InMemorySaver()
graph = builder.compile(checkpointer=checkpointer)
config = {
"configurable": {
"thread_id": "demo-run-001",
}
}
result = graph.invoke(
{
"count": 0,
"message": "",
},
config=config,
)
print(result)
预期结果类似:
{
"count": 1,
"message": "count=1"
}
这里有三个必须理解的点:
checkpointer负责保存图状态;thread_id是加载和继续同一条状态链的持久指针;- 如果换成新的
thread_id,就会开始一个新的执行上下文。
LangGraph 文档明确要求调用带 Checkpointer 的图时提供 thread_id,因为 Checkpointer 依赖它查找和恢复线程状态。(docs.langchain.com)
InMemorySaver 适合本地实验,因为进程退出后数据丢失。生产环境应选择持久化后端,例如 SQLite、PostgreSQL、MongoDB、Redis、DynamoDB 或其他官方集成;具体后端的可用性和包名应以对应版本文档为准。(docs.langchain.com)
2. 查看状态和历史
同一个 thread_id 可以读取当前状态:
snapshot = graph.get_state(config)
print(snapshot.values)
print(snapshot.next)
也可以查看历史:
for item in graph.get_state_history(config):
print(
item.config,
item.values,
item.next,
)
这些接口适合:
- 判断 Agent 卡在哪个节点;
- 查看最后一次成功状态;
- 分析某次模型决策前的上下文;
- 对人工审批后的状态进行审计;
- 选择某个历史 Checkpoint 进行分支实验。
但需要区分:查看历史不等于修改生产状态。从旧 Checkpoint 分支通常应生成新的运行标识或新的检查点分支,不能直接覆盖原运行,否则会破坏原始审计链。
3. 人工介入和恢复
动态中断的典型流程是:
节点调用 interrupt()
→ 保存当前图状态
→ 返回中断信息
→ 外部系统展示给人
→ 使用同一个 thread_id 和恢复命令继续
关键点是恢复时必须复用原来的 thread_id。如果改用新 ID,系统会把它视为新的线程,而不是原执行的继续。(docs.langchain.com)
在包含随机数、时间、模型调用或文件写入的节点中,不能假设“函数从中断位置的下一行继续”。恢复机制可能重新进入节点,因此副作用必须放进可持久化 task,并且任务应具备幂等性。LangGraph 文档特别指出,错误地把副作用放在中断前的普通流程代码中,恢复时可能再次执行该副作用。(docs.langchain.com)
九、多 Agent 场景:父图、子图和状态隔离
多 Agent 系统经常包含:
Supervisor
├── Researcher
├── Coder
└── Reviewer
此时要区分三种状态:
- 父图状态:任务目标、全局计划、子 Agent 结果;
- 子图调用状态:某次 Researcher 调用的中间过程;
- 跨会话记忆:用户偏好、长期事实和历史摘要。
子图是否持久化,决定了恢复边界:
- per-invocation:每次调用独立,但支持一次调用内部的中断和恢复;
- per-thread:同一线程跨调用累积状态;
- stateless:不保存状态,也不支持持久化恢复。
LangGraph 文档将 per-invocation 描述为多数多 Agent 场景的合适默认模式;如果子 Agent 需要跨轮次保留上下文,才使用 per-thread。无 Checkpoint 的 stateless 子图在进程崩溃后必须从头重跑。(docs.langchain.com)
一个常见错误是:同一个有状态子图实例在同一个父节点中被连续调用多次,导致多个调用写入相同的 Checkpoint 命名空间。若这些调用本应彼此独立,应使用每次调用独立的持久化边界,而不是共享同一份子图线程状态。(docs.langchain.com)
十、错误处理:区分可重试、可恢复和不可恢复
不是所有异常都应该重试。
1. 可重试错误
典型包括:
连接超时
HTTP 429
临时 DNS 失败
数据库暂时不可用
重试应满足:
并使用指数退避、最大次数和抖动:
2. 可恢复但不能盲目重试的错误
例如:
支付接口超时
订单创建请求超时
消息发送后确认丢失
这些错误的事实状态不确定。正确流程是:
查询外部系统状态
→ 已成功:记录结果并继续
→ 明确失败:按业务规则重试或终止
→ 仍未知:进入人工处理或延迟队列
3. 不可恢复错误
典型包括:
状态 schema 无法解析
Checkpoint 损坏
权限永久失败
业务参数违反不变量
代码版本不兼容
这类错误继续重试只会制造更多噪声。系统应保存:
{
"status": "blocked",
"failure_class": "non_recoverable",
"checkpoint_id": "ckpt_42",
"error_code": "SCHEMA_VERSION_UNSUPPORTED",
"requires_operator": true
}
LangGraph 提供错误处理和重试策略相关能力,但业务系统仍需定义哪些异常可重试、哪些异常必须进入补偿或人工处理;框架不会自动知道外部支付、库存或审批的业务语义。(docs.langchain.com)
十一、状态版本迁移:Checkpoint 不是永远兼容的对象
Agent 运行可能持续数天,而代码部署可能每天发生变化。因此,Checkpoint 必须带有 schema 版本:
{
"schema_version": 3,
"graph_version": "2026.09.01",
"state": {
"plan": "...",
"approval": {
"status": "pending"
}
}
}
假设旧版本状态为:
{
"approved": true
}
新版本改成:
{
"approval": {
"status": "approved"
}
}
恢复前必须执行迁移:
def migrate_state(state: dict, version: int) -> tuple[dict, int]:
if version == 1:
state = {
**state,
"approval": {
"status": "approved" if state.pop("approved", False)
else "pending"
},
}
version = 2
return state, version
迁移函数需要满足:
- 可重复执行;
- 有明确的源版本和目标版本;
- 迁移失败不会覆盖原 Checkpoint;
- 能在测试环境重放历史数据;
- 对不可迁移状态明确阻断,而不是静默丢字段。
图逻辑版本也需要记录。即使状态 schema 没变,节点名称、守卫条件、工具参数或模型版本变化,也可能使旧运行无法按原路径恢复。
十二、确定性不等于“没有重试”
一个重要反例是:
def charge_card(order_id):
return payment_api.charge(order_id, amount=100)
开发者可能认为“只要返回值相同,这个函数就是幂等的”。实际上,第一次调用可能已经扣款成功,但由于响应超时,第二次调用仍然会再次扣款。
真正需要的是业务幂等:
def charge_card(order_id, amount):
key = f"{order_id}:charge:v1"
existing = payment_api.lookup_by_idempotency_key(key)
if existing:
return existing
return payment_api.charge(
order_id=order_id,
amount=amount,
idempotency_key=key,
)
这里的确定性来自:
相同业务操作 → 相同幂等键 → 相同外部资源或结果
而不是来自 Python 函数每次返回同一个值。
另一个反例是:
def decide(state):
now = datetime.now()
if now.hour < 12:
return {"route": "morning"}
return {"route": "afternoon"}
如果 Checkpoint 保存于 11:59,恢复发生于 12:01,Agent 可能改变路径。正确方式是:
def decide(state):
execution_date = state["execution_date"]
if execution_date.hour < 12:
return {"route": "morning"}
return {"route": "afternoon"}
execution_date 应在运行开始时生成并持久化,而不是每次恢复重新读取当前时间。
十三、生产诊断:从“Agent 卡住了”还原到具体事实
排查持久化执行问题时,至少要围绕以下关联键建立查询:
tenant_id
run_id
thread_id
checkpoint_id
event_id
task_id
lease_epoch
idempotency_key
一个最小诊断视图应能回答:
最后一个已提交 Checkpoint 是什么?
当前 next 节点是什么?
最近一次 lease owner 是谁?
lease 是否已经过期?
哪些 task 已完成?
哪些 task 处于 started 但没有 completed?
哪些 outbox 事件仍在 pending?
最近一次异常属于哪一类?
建议将执行时序记录为:
10:00:00 lease acquired epoch=12
10:00:01 node.started reserve_inventory
10:00:03 external.requested key=order_123:reserve
10:00:04 external.completed reservation_id=res_9
10:00:05 task.completed reserve_inventory
10:00:06 checkpoint.committed id=ckpt_44
如果最后一条记录是:
external.completed
但没有:
task.completed
checkpoint.committed
那么恢复时必须把该操作视为“结果可能已经存在”,而不是直接认定为失败。
如果最后一条记录是:
lease.lost epoch=12
旧 Worker 即使后续打印了“执行成功”,也不应被视为合法提交。真正可信的是带 fencing 校验并成功提交到持久化存储的结果。
十四、可靠执行的边界:到底能保证什么
在没有分布式事务同时覆盖:
Checkpoint 存储
外部 API
消息系统
数据库
的情况下,Agent 通常无法凭空获得严格的 exactly-once 语义。
更现实的保证是:
对应到工程语义:
- Checkpoint 保证进度可保存;
- 事件日志保证事实可追踪;
- 租约降低并发接管造成的冲突;
- fencing 阻止旧 Worker 的延迟写入;
- 幂等键限制重复副作用;
- 恢复算法决定哪些任务复用、哪些任务重试;
- 确定性保证同一运行的控制流可解释;
- 补偿流程处理无法原子提交的外部操作。
因此,“持久化执行”不是把 Agent 变成永远不会失败的程序,而是把失败转换成可识别、可恢复、可审计的状态。
当一个 Agent 能够明确回答:
我最后确认完成了什么;
我当前拥有哪一代执行权;
我恢复时将从哪个 Checkpoint 继续;
哪些副作用可能已经发生;
哪些结果可以复用;
哪些操作需要查询或补偿;
它才真正具备长时间运行和跨进程恢复的工程基础。
系列导航与关联阅读
- 系列入口:Agent 工程完整路线:从运行循环、记忆与协议到安全、评测和生产交付
- 上一篇:Agent 并发与调度:会话隔离、任务池、优先级、公平和资源配额
- 下一篇:Agent Checkpoint 与恢复:快照、增量、版本迁移和副作用重放
- 延伸:Agent 状态机设计:节点、事件、守卫、转移和可恢复执行
官方资料
本文依据 Agent、模型、协议与框架官方资料重新梳理;正文、示例与生产清单由 WR BLOG 编写。

评论
0 条讨论