四十秒的任务撞上三十秒的超时
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 |
| 提交接口 P99 | — | 34ms |
| 端到端 P50 / P99 | 11s / 41s | 8.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 个百分点更让我满意。
至于任务取消和优先级,我还在改,有结果再写。