Administrator
发布于 2025-06-14 / 1017 阅读
19

Agent 任务编排与分布式调度的结合

一次发版,37 个跑了两小时的任务全没了

五月底的一次例行发版,我们的合同审核 Agent 服务重启。重启完没多久,业务部门来问:"我昨天下午提交的那批合同怎么一直显示'审核中'?" 查了一下,那次重启时正在执行的 37 个长任务全部丢失,其中最长的已经跑了 2 小时 17 分钟,完成了 60 多个步骤。

更要命的是,这些任务的前端状态永远停在"审核中",用户既看不到结果也无法重试,最后是运维手动改数据库才解决的。这次事故让我意识到:Agent 长任务是分布式系统里的经典难题,不是 AI 的新问题,但我们当时完全没按分布式系统的思路设计它。

根因:把长任务当成了 HTTP 请求

我们第一版的执行逻辑长这样:

@PostMapping("/agent/contract-review")
public Result submit(@RequestBody ContractReviewRequest req) {
    String taskId = idGen.next();
    // 异步执行,但状态全在内存里
    CompletableFuture.runAsync(() -> {
        AgentContext ctx = new AgentContext(taskId);
        runningTasks.put(taskId, ctx);          // ← 问题在这
        while (!ctx.isFinished()) {
            Step step = planner.nextStep(ctx);
            StepResult r = toolExecutor.execute(step);
            ctx.record(step, r);
        }
        resultStore.save(taskId, ctx.summary());
    });
    return Result.ok(taskId);
}

问题一目了然:runningTasks 是个内存 Map。进程一重启,所有执行中的上下文全部消失。而且没有持久化的执行状态,重启后连"哪些任务在跑"都不知道,更别说恢复。

还有个隐藏问题:这个任务跑在同一个 JVM 里,如果某个工具调用阻塞(我们遇到过一个大文件解析卡了 20 分钟),会一直占着一个线程。我们的线程池只有 32 个,几个长任务就能占满。

重新设计:把执行状态持久化

核心思路是把 Agent 执行建模成一个持久化的状态机,每一步都落库,进程崩溃后可以从最后一步恢复。表结构:

CREATE TABLE agent_task (
    id            BIGINT PRIMARY KEY,
    task_type     VARCHAR(64) NOT NULL,
    status        VARCHAR(32) NOT NULL,   -- PENDING/RUNNING/PAUSED/SUCCESS/FAILED
    current_step  INT NOT NULL DEFAULT 0,
    input         JSON NOT NULL,
    result        JSON,
    error_msg     TEXT,
    retry_count   INT NOT NULL DEFAULT 0,
    version       INT NOT NULL DEFAULT 0,  -- 乐观锁
    owner         VARCHAR(64),             -- 当前处理该任务的实例
    lease_expire  DATETIME,                -- 租约过期时间
    created_at    DATETIME NOT NULL,
    updated_at    DATETIME NOT NULL,
    INDEX idx_status_lease (status, lease_expire)
);

CREATE TABLE agent_step (
    id          BIGINT PRIMARY KEY,
    task_id     BIGINT NOT NULL,
    step_no     INT NOT NULL,
    step_type   VARCHAR(32) NOT NULL,
    input       JSON NOT NULL,
    output      JSON,
    compressed  TINYINT NOT NULL DEFAULT 0,
    tokens      INT,
    duration_ms INT,
    status      VARCHAR(16) NOT NULL,
    UNIQUE KEY uk_task_step (task_id, step_no)
);

两个表的关系是关键:agent_task 存任务的当前位置和整体状态,agent_step 存每一步的完整记录。恢复上下文时不是从内存读,而是从 agent_step 重建。

调度:数据库轮询 + 租约

我们评估过几种方案:

  • Kafka 驱动:每步完成后发消息触发下一步。优点是实时性好,缺点是任务状态和消息状态要保持一致,且失败补偿复杂;
  • XXL-Job / PowerJob 这类调度框架:成熟,但它们的模型是"定时触发的任务",跟 Agent 的"前一步决定后一步"不太匹配;
  • 数据库轮询 + 租约:最简单,也最可控。

最后选了第三种。核心逻辑是:每个实例定期扫描 PENDINGlease_expire 已过期的任务,用乐观锁抢占,抢到了就执行一步。

@Scheduled(fixedDelay = 500)
public void pollAndExecute() {
    List<AgentTask> tasks = taskRepo.findRunnable(
            LocalDateTime.now(), INSTANCE_ID, 20);
    tasks.forEach(this::executeOneStep);
}

@Transactional
public boolean acquire(AgentTask task) {
    int updated = taskRepo.leaseTask(
            task.getId(),
            INSTANCE_ID,
            LocalDateTime.now().plusSeconds(120),   // 租约 2 分钟
            task.getVersion());
    return updated > 0;   // 乐观锁,抢不到就跳过
}

对应 SQL:

UPDATE agent_task
   SET owner = #{instanceId},
       lease_expire = #{leaseExpire},
       status = 'RUNNING',
       version = version + 1
 WHERE id = #{id}
   AND version = #{version}
   AND (status = 'PENDING' OR lease_expire < NOW())

租约机制解决的是"实例崩溃后任务卡死"的问题:如果一个实例执行到一半挂了,租约到期后任务会被其他实例接管。租约时长要大于单步最长执行时间,我们设的 120 秒(单步 P99 是 43 秒)。

这里有个细节:执行前必须再检查一次租约是否仍然有效。因为从"抢占"到"执行"之间可能有延迟,如果某一步特别慢,租约可能已经过期被别人接管了。我们在每一步工具调用前都校验一次 owner = INSTANCE_ID

断点续跑:上下文重建

恢复执行的关键是从 agent_step 重建出 Agent 上下文。这里有个权衡:完整重建会把所有历史步骤读出来,长任务的上下文会很大(我们最极端的有 143 步,历史记录 2.3MB)。

我们的做法是分级重建

public AgentContext rebuild(Long taskId) {
    AgentTask task = taskRepo.findById(taskId);
    int total = task.getCurrentStep();

    List<StepRecord> steps = stepRepo.findByTaskId(taskId);

    // 最近 4 步保留完整内容,更早的用摘要
    List<StepRecord> recent = steps.stream()
            .filter(s -> s.getStepNo() > total - 4)
            .toList();

    List<StepRecord> old = steps.stream()
            .filter(s -> s.getStepNo() <= total - 4)
            .toList();

    String summary = old.isEmpty() ? ""
            : summaryStore.loadLatest(taskId);   // 复用之前生成的摘要

    return AgentContext.of(task, summary, recent);
}

摘要不是重建时临时算的(那样恢复太慢),而是在执行过程中每 4 步就生成一次并存到 agent_task 的一个字段里。这样恢复只需要读一份摘要加 4 步记录,实测重建耗时从 1.8 秒降到 90 毫秒。

这里踩过一个坑:早期我们把摘要存在 result 字段里,结果把真正的任务结果覆盖了。后来单独加了 context_summary 字段。

失败补偿:分三类处理

Agent 的每一步失败,处理方式取决于这一步的性质:

步骤类型示例失败处理
只读操作查询监控、检索文档重试 3 次,仍失败则把错误信息交给模型让它换思路
可补偿的写创建工单、发通知执行前记录补偿动作,失败时执行补偿
不可补偿的写调用支付、审批通过不自动重试,转人工确认

第二类需要工具自己声明补偿逻辑:

public interface CompensableTool extends Tool {
    /** 返回补偿所需的快照数据 */
    default Map<String, Object> snapshot(StepInput input) { return Map.of(); }

    /** 执行补偿 */
    void compensate(Map<String, Object> snapshot);
}

// 示例
@Component
public class CreateTicketTool implements CompensableTool {
    public Map<String, Object> snapshot(StepInput in) {
        return Map.of("draftId", in.get("draftId"));
    }
    public void compensate(Map<String, Object> snap) {
        ticketClient.deleteDraft((String) snap.get("draftId"));
    }
}

第三类是最需要谨慎的。我们的规则是:涉及资金、对外承诺、权限变更的操作,Agent 只能生成"待执行计划",必须人工点击确认后才真正执行。这一步在架构上体现为:这类工具的执行结果不是直接生效,而是写入一个 pending_action 表,等待人工审批。

这个设计在上线第二周就派上用场了:Agent 因为误解了合同条款,生成了一条"向供应商发送催款函"的待执行操作,被审批人拦下了。如果当时是自动执行,就是一次对外事故。

幂等:每一步都要有唯一键

因为有了重试和任务接管,同一步可能被重复执行。我们的要求是所有工具调用必须幂等,实现方式是给每步生成幂等键:

String idempotentKey = taskId + ":" + stepNo + ":" + toolName;

工具执行前先查一下这个 key 的结果是否已存在(存在 agent_step 表里),存在就直接返回,不重复执行:

Optional<StepRecord> existing =
        stepRepo.findByIdemKey(taskId, stepNo, toolName);
if (existing.isPresent() && existing.get().isSuccess()) {
    log.info("step already executed, skip: {}", idempotentKey);
    return existing.get().getOutput();
}

对于无法保证幂等的外部调用(比如发送通知),我们在调用前记录一条 pending 状态的 step 记录,调用成功后更新为 success。如果进程在中间崩溃,恢复时会看到一条 pending 记录,此时不盲目重试(因为可能已经发出去了),而是标记为 unknown 转人工确认。

这个 unknown 状态一开始没设计,是上线后发现有几条通知重复发送了两遍才补上的。教训是:外部副作用的状态机里,"不知道成功没有"是一个必须存在的状态,不能只有成功和失败。

孤儿任务回收

还有个必须要处理的:任务被标记为 RUNNING,但处理它的实例已经不在了(进程被 kill -9、机器宕机),而租约因为某种原因没有正常过期。我们加了一个定时任务:

@Scheduled(fixedDelay = 60_000)
public void reclaimOrphans() {
    // 租约过期超过 5 分钟仍未推进的任务,强制回收
    List<AgentTask> orphans = taskRepo.findStale(
            LocalDateTime.now().minusMinutes(5));
    orphans.forEach(t -> {
        log.warn("reclaim orphan task: id={}, step={}, lastOwner={}",
                 t.getId(), t.getCurrentStep(), t.getOwner());
        taskRepo.resetToPending(t.getId());
    });
}

这个任务的告警我们接到了值班群,因为孤儿任务通常意味着有实例非正常退出,值得关注。

改造后的效果

指标改造前改造后
发版重启导致的任务丢失37 个/次0
单任务最长执行时间2h17m(记录)4h06m
任务最终成功率82.3%94.7%
崩溃恢复时间不可恢复平均 3 秒
重复执行导致的事故3 次0
单实例可承载并发任务32400+

并发能力的大幅提升是因为改成了"轮询 + 执行一步"的模型,不再有线程被长任务长期占用。配合虚拟线程,单实例轻松跑几百个并发任务。

成功率提升的 12 个点主要来自两点:崩溃后可恢复(原来直接失败),以及单步失败可以重试(原来整个任务失败)。

几条经验

  • Agent 长任务本质上就是分布式任务调度,别因为它是 AI 就忘了基本的工程套路:状态持久化、租约、幂等、补偿;
  • 每一步都落库,这是所有恢复能力的地基。代价是数据库写入变多(我们峰值每秒 2000 次 step 写入),但值得;
  • 上下文重建要提前算摘要,别在恢复时临时生成;
  • 状态机里要有 unknown 状态,处理"不确定有没有执行成功"的情况;
  • 不可逆操作必须人工确认,这条救过我们一次;
  • 租约时长要大于单步 P99 的 2 倍以上,否则会频繁发生任务被误接管。

小结

这次改造让我重新认识到一件事:Agent 的"智能"部分(规划、工具选择)反而不是工程难点,难的是让一个可能跑几小时、涉及外部副作用、随时可能崩溃的流程可靠地完成。这些问题在 SOA 时代就有成熟的答案——状态机、Saga、幂等、补偿——只是套了一层 AI 的外衣。

如果你的 Agent 要跑超过几分钟,或者涉及任何不可逆操作,建议在设计阶段就把这套东西考虑进去。事后补的代价,我们这次是两周加一次线上事故。

参考