Administrator
发布于 2026-08-24 / 1105 阅读
21

事件驱动架构与 Agent 的深度结合

四十秒的任务撞上三十秒的超时

8 月 11 号,我们的工单自动处理 Agent 上线一周后,监控上出现一条难看的曲线:任务失败率 4.7%,失败原因几乎全是 UpstreamTimeout

查了一下,原因很直白——这个 Agent 平均执行 11 秒,但 P99 是 41 秒。而我们网关的超时是 30 秒,那是三年前面向普通 API 定的值。

$ kubectl logs agent-orch-7d9f -n ai | grep TIMEOUT | tail -3
2026-08-11T14:22:07  WARN  task=AGT-8841201 timeout after 30012ms, killed
2026-08-11T14:22:11  WARN  task=AGT-8841207 timeout after 30004ms, killed
2026-08-11T14:22:19  WARN  task=AGT-8841213 timeout after 30011ms, killed

最气人的是第三单:被 kill 的时候它已经执行到 94%,再有两秒就出结果了。资源全浪费,用户还得重来一遍。

这篇记录我们把这类长任务从同步调用改成事件驱动的完整过程。不是什么新鲜架构,但和 Agent 结合时有几个地方挺特殊。

先确认:这不是调大超时能解决的

第一个冒出来的方案是把超时改成 120 秒。我否决了,理由有三条,都有数据支撑。

连接占用。我们的网关用 Spring Boot 4 + 虚拟线程,连接本身不贵,但整个链路上的 Nginx、ALB、客户端 SDK 都各自有超时。要改得全改,而且 ALB 的空闲超时上限是 60 分钟、默认 60 秒,改了会影响其他业务。

重试风暴。超时之后上游会重试。Agent 任务贵(单次 ¥0.2 左右)且慢,一次超时重试就是双倍成本。我们统计过,4.7% 的失败里有 1.8% 触发了上游重试,等于凭空多烧 38% 的钱。

用户体验。用户盯着一个转圈的界面看 40 秒,不知道进行到哪一步,这个体验本身就不合格。我们产品侧的要求是要能显示进度。

所以方向是异步化:请求返回一个任务 ID,用户在别处等结果。

改造后的架构

                  ┌──────────── Kafka: agent.task.submit ──────────┐
                  │                                                 │
工单系统 ──HTTP──> Agent 网关 ──> 立即返回 taskId ──> 用户看到"处理中"
                       │
                       └──> 写 agent_task 表(状态 INIT)
                                ▲
                                │
                  ┌─────────────┴──────────────┐
                  │      Agent Worker 集群      │
                  │  消费 → 执行 → 上报进度      │
                  └─────────────┬──────────────┘
                                │
                  ┌─────────────▼──────────────┐
                  │  Kafka: agent.task.progress │
                  │  Kafka: agent.task.done     │
                  └─────────────┬──────────────┘
                                │
              ┌─────────────────┼─────────────────┐
              ▼                 ▼                 ▼
         SSE 推送          Webhook 回调      站内消息

事件触发:不是所有事都该触发 Agent

第一个决策点是"什么事件触发 Agent"。我们最初的做法是"工单创建就触发",跑了三天发现 62% 的任务是白跑——用户只是创建了个工单还没填完内容,或者工单类型根本不适合自动处理。

改成两级过滤:

@Component
public class AgentTrigger {

    /** L1:廉价规则过滤,纳秒级,拦掉明显不合适的 */
    public boolean preFilter(TicketEvent e) {
        if (e.type() == SPAM || e.type() == INTERNAL_TEST) return false;
        if (e.contentLength() < 20) return false;             // 内容太短
        if (!AUTO_TYPES.contains(e.category())) return false;  // 类目不在白名单
        return dedupe.isFirstEvent(e.ticketId(), e.version()); // 同一版本只触发一次
    }

    /** L2:便宜模型判断"值不值得用贵的模型处理" */
    public boolean worthIt(TicketEvent e) {
        float score = smallModel.score("""
            判断这个工单是否适合 AI 自动处理。只输出 0 或 1。
            适合:有明确诉求、信息基本完整、不涉及投诉升级。
            工单内容:%s
            """.formatted(e.content()));
        return score > 0.5;
    }
}

L1 拦掉 44%,L2 再拦掉 19%,最终只有 37% 的工单会真正触发 Agent。按日均 8 万工单算,一天少跑 5 万个任务,省 ¥1 万。

L2 用小模型(¥0.00003/次),它的误判成本很低——把该处理的拦了,用户还能手工处理;把不该处理的放进去了,也就是浪费一次调用。这个不对称性决定了 L2 的阈值可以偏保守。

异步执行:状态机 + 进度上报

任务状态机是核心。我们定义了六个状态,每个状态的转移都有明确条件:

public enum TaskState {
    INIT,       // 已创建,未投递
    QUEUED,     // 已投递到 MQ
    RUNNING,    // Worker 已领取
    SUCCEEDED,  // 完成
    FAILED,     // 失败且不可重试
    UNKNOWN     // 状态不明,需要查证
}

UNKNOWN 这一态是必须的,理由我在 Agent 幂等那篇里写过:Worker 可能执行完了但在上报结果前挂了。没有这一态,就只能靠超时猜,猜错就是重复执行或者丢任务。

Worker 侧用虚拟线程 + 手动提交 offset:

@KafkaListener(
    topics = "agent.task.submit",
    groupId = "agent-worker",
    concurrency = "8",
    containerFactory = "virtualThreadFactory")
public void onTask(ConsumerRecord<String, TaskEvent> rec, Acknowledgment ack) {
    TaskEvent ev = rec.value();

    // 幂等:同一 taskId 只处理一次
    if (!taskDao.casState(ev.taskId(), QUEUED, RUNNING)) {
        ack.acknowledge();          // 已被别人领走或已完成,直接提交
        return;
    }

    try {
        AtomicInteger step = new AtomicInteger();
        AgentResult r = agent.run(ev,
            progress -> {           // 每次工具调用后回调
                progressPublisher.send(ev.taskId(), step.incrementAndGet(), progress);
                taskDao.touchHeartbeat(ev.taskId());
            });
        taskDao.finish(ev.taskId(), SUCCEEDED, r);
        donePublisher.send(ev.taskId(), r);
        ack.acknowledge();
    } catch (Exception ex) {
        taskDao.finish(ev.taskId(), FAILED, ex);
        ack.acknowledge();          // 业务失败也要提交,不能让消息重放
    }
}

两个细节值得说。

心跳。taskDao.touchHeartbeat() 每次工具调用后更新 last_heartbeat。另有一个定时任务扫 RUNNING 且心跳超过 3 分钟的任务,标记为 UNKNOWN,然后由 Worker 的 reconcile 逻辑去查真实状态。这个机制上线第二周就用上了——一个 Worker Pod 被 OOM kill,7 个任务进了 UNKNOWN,全部通过查证恢复,没有一个重复执行。

失败也要 ack。初版我写的是失败不提交 offset 等重投,这是错的。Agent 任务的失败大多是模型输出不合规、工具返回异常,重试大概率还是失败,白白浪费钱。真正的重试策略应该按错误类型分:网络错误、限流错误可以重投;业务错误直接进死信队列,人工看。

结果通知:三条通道各有用处

通知这块我们踩了点弯路,一开始只做了 SSE,后来发现不够。

通道适用方优点
SSE前端页面能推进度,实现简单断线重连要处理 Last-Event-ID
Webhook下游系统解耦,标准做法对端超时/不可用要重试
站内消息兜底,一定能到实时性差

SSE 的进度推送:

@GetMapping(value = "/tasks/{id}/stream", produces = MediaType.TEXT_EVENT_STREAM_VALUE)
public Flux<ServerSentEvent<Progress>> stream(@PathVariable String id,
                                               @RequestHeader(value = "Last-Event-ID", required = false) String lastId) {
    // 断线重连:先补发历史进度,再接实时流
    Flux<Progress> replay = lastId == null
        ? Flux.empty()
        : progressStore.since(id, Long.parseLong(lastId));

    return Flux.concat(replay, progressSink.asFlux().filter(p -> p.taskId().equals(id)))
               .map(p -> ServerSentEvent.builder(p).id(String.valueOf(p.seq())).build())
               .timeout(Duration.ofMinutes(5))
               .doOnCancel(() -> log.debug("client disconnected {}", id));
}

Last-Event-ID 这个头是浏览器自动带的,但我们的移动端 SDK 一开始没实现,导致 App 切后台再回来进度就从头开始了。这个 bug 修了一天。

Webhook 的重试策略要注意退避:1s, 5s, 30s, 5min, 30min,最多 5 次。而且回调体必须带 taskId 的幂等键,因为对端可能收到了但响应丢了,我们不重投,但对端可能重复收到(网络重试导致)。

可靠性:哪些地方会丢任务

改造完我们做了一次故障演练,专门找丢任务的场景。列出发现的四处,都修了。

一、投递后立刻宕机。网关写了 INIT 状态的事务提交了,但 Kafka 消息没发出去(进程被 kill)。修复:用本地消息表 + 定时补偿。

@Transactional
public String submit(TaskRequest req) {
    String taskId = idGen.next();
    taskDao.insert(taskId, INIT, req);
    outboxDao.insert("agent.task.submit", taskId, toJson(req));  // 同事务
    return taskId;
}

@Scheduled(fixedDelay = 1000)     // 补偿发送
public void relayOutbox() {
    outboxDao.pending(500).forEach(m -> {
        kafka.send(m.topic(), m.key(), m.body())
             .whenComplete((r, ex) -> {
                 if (ex == null) outboxDao.markSent(m.id());
                 else outboxDao.incrRetry(m.id());
             });
    });
}

二、消费者 rebalance 期间的重复。Kafka rebalance 时可能重复投递。上面那个 casState 挡住了绝大部分,但我们还是加了业务层的幂等检查。演练 200 次 rebalance,零重复执行。

三、进度消息乱序。SSE 推送的进度可能因为多分区而乱序(step 5 先于 step 3 到达)。前端加了序号排序 + 丢弃过期序号。这个 bug 在压测时发现,用户会看到进度条往回走。

四、结果通知全失败。SSE 断了、Webhook 对端挂了、站内消息服务不可用。这种情况我们靠轮询兜底:前端每 10 秒拉一次任务状态,最多拉 30 分钟。这是最后一道防线,实现最简单,也最可靠。

改造后的数据

8 月 20 日全量上线,两周的数据:

指标同步模式事件驱动
任务成功率91.2%99.4%
超时失败率4.7%0
重复执行1.8%(上游重试)0
提交接口 P9934ms
端到端 P50 / P9911s / 41s8.2s / 47s
日均任务数8.0 万3.0 万(过滤后)
日均成本¥1.6 万¥0.62 万
用户中途放弃率12.4%3.1%

端到端 P99 反而涨到了 47 秒,因为不再有 30 秒的硬截断,那些真正慢的任务(多轮工具调用、模型重试)能跑完了。但用户放弃率从 12.4% 降到 3.1%,因为能看到进度条,知道系统在动。这个数字比我预期的重要——用户能忍受慢,不能忍受不确定。

成本降幅主要来自前置过滤(少跑 62% 的任务),异步化本身的贡献是消除重复执行那部分。

还没解决的

三个问题。

任务优先级。现在所有任务一个队列,VIP 客户的高优工单和内部测试任务排在一起。压测时高优任务的排队延迟能到 4 分钟。我们打算做优先级队列,但 Kafka 原生不支持,现在的方案是按优先级分 topic、Worker 侧加权消费,还没实现完。

任务取消。用户点了"取消",但任务已经在 Worker 里跑了。我们能做的是在 Agent 的每一步检查取消标记,但模型调用本身无法中断——一次 8 秒的模型调用,点取消之后还得等它返回。目前只能在这 8 秒里不往下走。

跨任务的状态共享。同一个工单被拆成 3 个子任务并发执行,它们之间需要共享一些中间结果。我们现在靠数据库,但写冲突频发。考虑过引入事件溯源,评估下来复杂度太高,暂时用乐观锁扛着。

小结

这次改造技术上没什么新东西,全是消息队列、状态机、幂这些老套路。但和 Agent 结合有两个地方值得记下来。

一是Agent 任务的执行时间分布极长尾。传统 API 的 P99 可能是 P50 的 3 倍,我们的 Agent 任务是 5 倍以上。这意味着为同步调用设超时必然出错,异步不是优化,是必需。

二是Agent 任务贵到值得做前置过滤。传统消息消费的任务成本是微秒级的 CPU,随便跑;Agent 任务单次两毛钱,跑之前先花三厘钱判断一下值不值得跑,这笔账非常划算。这是事件驱动架构在 AI 时代的一个新变化。

最后一点体会:进度上报这件事,我一开始觉得是"锦上添花",做完发现它是改造里收益最大的部分。技术架构的事最后往往落在用户感受上,用户放弃率降了 9 个百分点,比成功率涨 8 个百分点更让我满意。

至于任务取消和优先级,我还在改,有结果再写。

参考