Skip to content

状态机与持久化:Agent 跑一半重启了怎么办 ​

1. 本节产出 ​

一个可恢复的 Agent 运行时:每轮迭代落检查点(checkpoint)、服务重启后能从断点继续、长时间任务可查询进度。并且有明确的上下文压缩策略防止状态无限膨胀。

2. 前置依赖 ​

3. 为什么状态持久化是「能不能上生产」的分水岭 ​

Agent 任务和普通 HTTP 请求有三个本质区别:

区别后果
耗时长(几十秒到几分钟)重启/发布必然打断进行中的任务
多轮迭代(每轮一次模型调用)中断损失大,重跑要重付钱
状态复杂(目标、历史、工具结果)只存内存无法恢复

没有持久化时,一次滚动更新就会让所有进行中的 Agent 任务全部丢失,用户看到的是「转了半天然后没了」。

这在企业环境里更严重:K8s 会周期性重启 Pod、节点会被调度、HPA 会扩缩容——重启不是意外,是常态。

4. 核心原理 ​

4.1 Agent 状态由什么组成 ​

AgentState = {
  runId,              运行标识
  goal,               用户目标(不变)
  status,             RUNNING / WAITING_APPROVAL / DONE / FAILED
  iteration,          当前轮次
  transcript,         累积的 Thought/Action/Observation
  contextSummary,     早期内容的压缩摘要(关键)
  tokenUsed,          已消耗
  pendingApproval,    等待审批的工具调用(如有)
  updatedAt
}

contextSummary 是最容易被忽略但最重要的字段。它是对早期 transcript 的压缩,用来对抗上下文膨胀。

4.2 检查点策略 ​

策略做法权衡
每轮落盘每次迭代后写一次最安全,但写入频繁
每 N 轮落盘批量写折中,可能丢 N 轮
关键节点落盘只在工具调用后、等待审批时写写入少,但普通迭代可能丢

推荐:每轮落盘 + 异步写入。Agent 每轮本来就有几百毫秒的模型调用延迟,异步写一次数据库(几毫秒)完全可以接受。

4.3 恢复流程 ​

服务启动
   │
   ▼ 扫描 status IN (RUNNING, WAITING_APPROVAL) 的运行
   │
   ├─ 未超过总超时 → 从 checkpoint 恢复,继续循环
   └─ 已超过总超时 → 标记 FAILED,通知用户

必须检查总超时。否则一个挂了三天的任务被恢复后继续烧钱。

4.4 上下文压缩(关键) ​

问题:10 轮之后 transcript 可能有 5 万 Token
方案:保留最近 K 轮的完整内容 + 早期内容的摘要

transcript = [
  摘要:「第1-5轮:查询了订单 SO123(已发货)、
         查询了物流(已到杭州)、尝试改状态被拒」,
  第6轮完整内容,
  第7轮完整内容,
  ...
  第10轮完整内容
]

压缩时机:估算 Token 超过阈值(比如窗口的 60%)时触发,把最早的几轮交给模型生成摘要。

注意:摘要本身要花钱(一次模型调用),但比起上下文爆炸导致任务失败或成本翻倍,这笔钱值得。

5. 代码走查 ​

5.1 状态定义与存储 ​

sql
CREATE TABLE agent_run (
  run_id          VARCHAR(64) PRIMARY KEY,
  tenant_id       VARCHAR(64) NOT NULL,
  user_id         VARCHAR(64),
  goal            TEXT NOT NULL,
  status          VARCHAR(24) NOT NULL,
  iteration       INT DEFAULT 0,
  max_iterations  INT DEFAULT 10,
  transcript      JSONB,          -- 最近若干轮
  context_summary TEXT,           -- 早期压缩摘要
  token_used      INT DEFAULT 0,
  pending_tool    JSONB,          -- 等待审批的工具调用
  created_at      TIMESTAMPTZ DEFAULT now(),
  updated_at      TIMESTAMPTZ DEFAULT now()
);

CREATE INDEX idx_run_status ON agent_run (status, updated_at);

5.2 检查点写入 ​

java
// src/main/java/com/example/harness/state/AgentStateRepository.java
@Repository
public class AgentStateRepository {

    /** 异步落盘,不阻塞主循环 */
    @Async("stateExecutor")
    public void checkpoint(AgentState s) {
        jdbc.update("""
                UPDATE agent_run
                   SET status=?, iteration=?, transcript=?,
                       context_summary=?, token_used=?, updated_at=now()
                 WHERE run_id=?
                """, s.status().name(), s.iteration(),
                Json.encode(s.recentTranscript()), s.contextSummary(),
                s.tokenUsed(), s.runId());
    }

    public List<AgentState> resumable() {
        return jdbc.query("""
                SELECT * FROM agent_run
                 WHERE status IN ('RUNNING','WAITING_APPROVAL')
                   AND updated_at > now() - INTERVAL '30 minutes'
                """, mapper);
    }
}

5.3 可恢复的循环 ​

java
// src/main/java/com/example/harness/agent/StatefulAgentRunner.java
@Service
public class StatefulAgentRunner {

    public AgentResult run(String runId, String goal, int maxIterations) {
        AgentState state = repo.load(runId)
                .orElseGet(() -> AgentState.start(runId, goal, maxIterations));

        while (state.iteration() < maxIterations) {
            // 1. 上下文压缩检查
            if (needsCompress(state)) {
                state = compress(state);
            }

            // 2. 模型决策
            Action action = decide(state);
            state = state.withIteration(state.iteration() + 1);

            if (action.step() == FINISH) {
                return finish(state, action.finalAnswer());
            }

            // 3. 工具执行(可能触发审批 → 保存状态并暂停)
            if (requiresApproval(action)) {
                state = state.waitingApproval(action);
                repo.checkpoint(state);
                return AgentResult.needApproval(runId);
            }

            String obs = executeTool(action);
            state = state.append(action, obs);

            // 4. 落检查点
            repo.checkpoint(state);
        }
        return AgentResult.maxIterations(state);
    }
}

5.4 上下文压缩 ​

java
// src/main/java/com/example/harness/state/ContextCompressor.java
@Component
public class ContextCompressor {

    public AgentState compress(AgentState state) {
        List<String> all = state.fullTranscript();
        int keep = 4;                                  // 保留最近 4 轮
        List<String> older = all.subList(0, all.size() - keep);

        String summary = chatClient.prompt()
                .system("""
                        把下面的 Agent 执行记录压缩成一段摘要。
                        保留:已完成的事实、已获得的关键数据、失败的尝试及原因。
                        丢弃:重复的尝试、工具返回的冗余细节。
                        控制在 300 字以内。
                        """)
                .user(String.join("\n", older))
                .options(OpenAiChatOptions.builder().temperature(0.0).maxTokens(500).build())
                .call().content();

        return state.withSummary(summary).withRecent(all.subList(all.size() - keep, all.size()));
    }
}

5.5 启动时恢复 ​

java
// src/main/java/com/example/harness/state/ResumeOnBoot.java
@Component
public class ResumeOnBoot implements ApplicationRunner {

    @Override
    public void run(ApplicationArguments args) {
        List<AgentState> pending = repo.resumable();
        log.info("发现 {} 个待恢复的 Agent 运行", pending.size());

        for (AgentState s : pending) {
            if (expired(s)) {
                repo.markFailed(s.runId(), "恢复时已超时");
            } else {
                runner.resume(s.runId());      // 从断点继续
            }
        }
    }
}

6. 跑起来 ​

bash
git checkout ch03-05-state-machine
docker compose up -d
mvn spring-boot:run
bash
# 1. 启动一个长任务(设 maxIterations=10)
TASK=$(curl -s -X POST http://localhost:8080/api/agent/run \
  -d '{"goal":"调研三个云厂商的向量数据库方案并对比"}' | jq -r .runId)

# 2. 跑到第 3 轮时 Ctrl+C 停掉应用
sleep 8 && pkill -f spring-boot

# 3. 重启
mvn spring-boot:run
# 期望日志:发现 1 个待恢复的 Agent 运行 → 从第 3 轮继续

# 4. 查进度
curl http://localhost:8080/api/agent/$TASK
# 期望:{"iteration":7,"status":"RUNNING","tokenUsed":18420}
检查项通过标准
检查点每轮后数据库有更新
恢复重启后从断点继续(日志显示起始轮次 > 1)
超时判断超时的运行被标记失败而非恢复
压缩触发长任务中日志出现压缩记录
幂等重复恢复不会产生重复执行

7. 生产避坑 ​

  1. 恢复时必须检查总超时。只看「状态是 RUNNING」就恢复,会让几天前挂掉的任务在今天重新启动并继续烧钱。做法是恢复前判断 updated_at 与总超时,超时直接标失败。
  2. 上下文压缩的触发阈值要留余量。等上下文已经到 90% 才压缩,往往来不及(压缩本身要调用模型,返回后可能已经超了)。经验值:60% 触发。
  3. 检查点写入要异步,且失败不能中断主循环。写数据库失败就中断 Agent 是过度反应——应该记日志告警,让任务继续跑(状态丢了可以重跑,但不该因为一次写失败就放弃已经花了钱的任务)。

8. 延伸与锚点 ​

  • 思考题:Agent 要执行一个「删除生产数据」的工具,怎么确保它不会误操作?(答案在下一课时:Human-in-the-loop)
  • 代码锚点:git checkout ch03-05-state-machine
  • 下一课时:03-06 Human-in-the-loop
  • 对应课件:L03-05 状态机与持久化