30 秒快速回答
Agent 持久化执行(Durable Execution) 是一种让 AI Agent 的任务状态在进程崩溃、客户端断连、服务器重启后仍能恢复并继续执行的技术。它的核心价值:让 Agent 从”一次性调用”升级为”可靠的异步任务引擎”,告别”跑了 30 分钟,断网就白干”的尴尬。
一、为什么 Agent 需要持久化执行?
传统的 Agent 调用模型是一个同步请求-响应链路:
用户 → Agent Server → LLM API → 工具调用 → 返回结果
这个模型在以下场景会直接崩溃:
| 场景 | 后果 |
|---|---|
| 执行到第 4 步工具调用时客户端断网 | 前 3 步结果全部丢失 |
| Agent 服务 OOM 重启 | 正在跑的长流程任务直接消失 |
| 大文件上传中浏览器标签被关闭 | 上传中断,无恢复机制 |
| 多 Agent 协作中某个节点超时 | 整个编排链断裂,无法从断点继续 |
持久化执行解决的核心问题:将 Agent 的运行时状态与进程生命周期解耦。
二、持久化执行的核心机制
2.1 工作流状态外置
持久化 Agent 不再把状态放在内存里,而是写入外部存储(数据库/消息队列):
# ❌ 传统方式:状态在内存中
class SimpleAgent:
def execute_task(self, task):
step1 = self.call_tool("search", task.query) # 内存中
step2 = self.call_tool("analyze", step1.result) # 内存中
step3 = self.call_tool("write", step2.result) # 内存中
return step3
# 任何一步崩溃 → 全部丢失
# ✅ 持久化方式:每一步都落盘
class DurableAgent:
def execute_task(self, task_id):
state = self.db.load_state(task_id) # 从 DB 恢复
for step in state.pending_steps:
try:
result = self.call_tool(step.tool, step.input)
state.mark_completed(step, result) # 每步落盘
self.db.save_state(state)
except Exception as e:
state.mark_failed(step, e) # 失败也记录
self.db.save_state(state)
# 下次重启或重试时从失败步骤继续
2.2 MCP Tasks:官方标准化方案
MCP 协议在 2026 年引入 MCP Tasks 规范,为持久化执行提供了标准化的异步调用模型:
// 客户端发起异步任务
{
"method": "tasks/run",
"params": {
"task": {
"name": "generate_report",
"input": { "query": "Q3 销售数据" }
}
}
}
// 服务端返回 task_id,客户端可随时查询状态
{
"task_id": "task_abc123",
"status": "running"
}
// 客户端断开重连后,查状态继续
{
"method": "tasks/status",
"params": { "task_id": "task_abc123" }
}
// → { "status": "completed", "result": {...} }
MCP Tasks 的三个关键特性:
| 特性 | 说明 |
|---|---|
| 异步提交 | 客户端提交任务后立即返回 task_id,不阻塞 |
| 状态可查询 | 随时通过 task_id 查询进度(pending → running → completed/failed) |
| 断连不丢 | 任务在服务端独立运行,客户端断开不影响执行 |
三、实现一个最小可用的持久化 Agent
以下用 Python + SQLite 展示核心思路:
import sqlite3
import json
import uuid
from datetime import datetime
class DurableTaskEngine:
def __init__(self, db_path="agent_tasks.db"):
self.conn = sqlite3.connect(db_path)
self._init_db()
def _init_db(self):
self.conn.execute("""
CREATE TABLE IF NOT EXISTS tasks (
task_id TEXT PRIMARY KEY,
status TEXT DEFAULT 'pending',
current_step INTEGER DEFAULT 0,
total_steps INTEGER,
step_results TEXT DEFAULT '[]',
created_at TEXT,
updated_at TEXT
)
""")
def submit_task(self, steps: list) -> str:
"""提交一个多步骤任务,返回 task_id"""
task_id = str(uuid.uuid4())[:8]
now = datetime.now().isoformat()
self.conn.execute(
"INSERT INTO tasks (task_id, total_steps, created_at, updated_at) VALUES (?, ?, ?, ?)",
(task_id, len(steps), now, now)
)
self.conn.commit()
return task_id
def resume_task(self, task_id: str) -> dict:
"""恢复任务:从上次中断的地方继续"""
row = self.conn.execute(
"SELECT * FROM tasks WHERE task_id = ?", (task_id,)
).fetchone()
if not row:
return {"error": "task not found"}
status = row[1]
current_step = row[2]
if status == "completed":
return {"status": "completed", "result": json.loads(row[4])}
return {
"task_id": task_id,
"status": status,
"current_step": current_step,
"total_steps": row[3],
"hint": f"从第 {current_step + 1} 步继续执行"
}
def mark_step_done(self, task_id: str, step_result: dict):
"""标记一步完成,持久化中间结果"""
row = self.conn.execute(
"SELECT current_step, total_steps, step_results FROM tasks WHERE task_id = ?",
(task_id,)
).fetchone()
new_step = row[0] + 1
results = json.loads(row[2])
results.append(step_result)
new_status = "completed" if new_step >= row[1] else "running"
self.conn.execute(
"UPDATE tasks SET current_step=?, step_results=?, status=?, updated_at=? WHERE task_id=?",
(new_step, json.dumps(results), new_status, datetime.now().isoformat(), task_id)
)
self.conn.commit()
使用示例:
engine = DurableTaskEngine()
# 提交一个 3 步任务
task_id = engine.submit_task([
{"tool": "search", "input": "AI 行业 2026 趋势"},
{"tool": "summarize", "input": ""},
{"tool": "write_report", "input": ""}
])
# → task_id: "a1b2c3d4"
# 模拟第 1 步执行后崩溃
engine.mark_step_done(task_id, {"output": "5大趋势..."})
# 💥 崩溃!
# 重启后恢复
state = engine.resume_task(task_id)
print(state)
# → current_step: 1, hint: "从第 2 步继续执行"
四、持久化执行的三种架构模式
| 模式 | 原理 | 适用场景 | 典型实现 |
|---|---|---|---|
| 事件溯源(Event Sourcing) | 记录所有事件日志,状态 = 事件重放 | 需完整审计轨迹的任务 | Temporal、Restate |
| 状态快照(State Snapshot) | 定期将完整状态序列化存储 | 状态体积小、恢复频率高的任务 | Redis RDB、自研方案 |
| 步骤检查点(Step Checkpoint) | 每完成一步就落盘中间结果 | 多步骤、步骤间有依赖的工作流 | MCP Tasks、LangGraph Checkpointer |
选型建议:
- 简单多步骤任务 → 步骤检查点(实现成本最低)
- 需要审计/回溯 → 事件溯源(可复现任意时刻状态)
- 高频短任务 → 状态快照(恢复速度最快)
五、生产环境注意事项
-
幂等性设计:每个步骤的执行函数必须支持重复调用不产生副作用。例如发送邮件前先查是否已发送。
- 超时与心跳:长时间运行的任务需要心跳机制,让调用方知道任务还活着:
# 心跳更新 self.conn.execute( "UPDATE tasks SET heartbeat_at=? WHERE task_id=?", (datetime.now().isoformat(), task_id) ) - 失败重试策略:建议指数退避 + 最大重试次数:
max_retries = 3 for attempt in range(max_retries): try: result = execute_step(step) break except Exception: wait = 2 ** attempt # 1s → 2s → 4s time.sleep(wait) else: mark_task_failed(task_id, "超过最大重试次数") - 存储选择:轻量场景用 SQLite/PostgreSQL;高并发场景考虑 Redis + 持久化队列(如 BullMQ、Celery)。
总结
- 持久化执行让 Agent 从”一次性工具”进化为”可靠服务”
- 核心思想:状态外置 + 步骤检查点 + 断点续传
- MCP Tasks 提供了标准化的异步任务模型
- 实现上优先选择步骤检查点模式,简单有效
延伸阅读建议: