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

选型建议:

  • 简单多步骤任务 → 步骤检查点(实现成本最低)
  • 需要审计/回溯 → 事件溯源(可复现任意时刻状态)
  • 高频短任务 → 状态快照(恢复速度最快)

五、生产环境注意事项

  1. 幂等性设计:每个步骤的执行函数必须支持重复调用不产生副作用。例如发送邮件前先查是否已发送。

  2. 超时与心跳:长时间运行的任务需要心跳机制,让调用方知道任务还活着:
    # 心跳更新
    self.conn.execute(
        "UPDATE tasks SET heartbeat_at=? WHERE task_id=?",
        (datetime.now().isoformat(), task_id)
    )
    
  3. 失败重试策略:建议指数退避 + 最大重试次数:
    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, "超过最大重试次数")
    
  4. 存储选择:轻量场景用 SQLite/PostgreSQL;高并发场景考虑 Redis + 持久化队列(如 BullMQ、Celery)。

总结

  • 持久化执行让 Agent 从”一次性工具”进化为”可靠服务”
  • 核心思想:状态外置 + 步骤检查点 + 断点续传
  • MCP Tasks 提供了标准化的异步任务模型
  • 实现上优先选择步骤检查点模式,简单有效

延伸阅读建议: