把几千个用例一次性串行提交,换来的是连接一挂状态全丢、内存差点 OOM。痛定思痛,我把批量流程重写成”批次 → 任务 → 子任务”三级结构,每一级都独立可重试、可恢复、可观测。这套策略在生产环境跑了半年,十万级用例的批次也能稳定收尾。下面把拆分和进度管理的完整思路写出来。
一、为什么一定要做任务拆分
把 5000 个测试一次性塞进去,看着简单,实际是灾难:
- 同步执行用户 HTTP 连接挂掉,状态丢失;
- 全部进内存,OOM 风险;
- 失败无法重试,只能重头跑;
- 进度不可见,用户不知道”还要等多久”。
拆分成”批次 → 任务 → 子任务”三级后,每级都可重试、可恢复、可观测。我们拆的时候遵循三个原则:
- 批次是用户可感知的单元,进度、取消、导出都挂在它上面;
- 任务是模型维度的统计单元,一个模型一条,方便横向对比;
- 子任务是最小执行单元,一个用例乘一个模型,粒度细到可以精确重试。
三句话说完,落地就是下一节的三级结构。
二、三级任务结构
三级结构里每一级承担不同的职责。Batch 是用户视角的入口,一个提交就是一个批次;Task 按模型拆,方便对比不同模型在同一批用例上的表现;SubTask 拆到最小,单条用例失败不影响整批。先看结构定义:
Batch (批次)
├─ Task (任务,对应一个 model × 一个 prompt 组合)
│ └─ SubTask (子任务,对应一个测试用例)
落到数据库就是三张表,外键逐级串联:
CREATE TABLE eval_batch (
id BIGINT PRIMARY KEY,
user_id BIGINT NOT NULL,
name VARCHAR(128),
total_count INT NOT NULL,
done_count INT DEFAULT 0,
fail_count INT DEFAULT 0,
status VARCHAR(16), -- pending/running/done/partial_error
created_at DATETIME
);
CREATE TABLE eval_task (
id BIGINT PRIMARY KEY,
batch_id BIGINT NOT NULL,
model_code VARCHAR(64) NOT NULL,
prompt_id BIGINT,
total_count INT,
done_count INT DEFAULT 0,
fail_count INT DEFAULT 0,
status VARCHAR(16),
INDEX idx_batch (batch_id)
);
CREATE TABLE eval_subtask (
id BIGINT PRIMARY KEY,
task_id BIGINT NOT NULL,
case_id BIGINT,
input TEXT,
expected TEXT,
output LONGTEXT,
score DECIMAL(4,2),
status VARCHAR(16), -- pending/running/success/fail
error_msg VARCHAR(512),
duration_ms INT,
retry_count INT DEFAULT 0,
INDEX idx_task (task_id)
);
三张表之间是严格的父子关系,进度和失败都能从子级向上汇聚,这给后面的并发控制和重试打好了地基。
三、任务拆分算法
结构定了,拆分算法就顺理成章。用户提交的是一次”用例集合 × 模型集合”的笛卡尔积,我们在创建批次时把它展开成任务和子任务,同时入队。核心逻辑在 createBatch 里,四步注释已经写清楚了:
public Long createBatch(BatchCreateRequest req) {
// 1) 创建 batch
EvalBatch batch = new EvalBatch();
batch.setTotalCount(req.getCaseIds().size() * req.getModelCodes().size());
batchService.save(batch);
// 2) 按模型拆 task
for (String model : req.getModelCodes()) {
EvalTask task = new EvalTask();
task.setBatchId(batch.getId());
task.setModelCode(model);
task.setTotalCount(req.getCaseIds().size());
taskService.save(task);
// 3) 用例拆 subtask
for (Long caseId : req.getCaseIds()) {
EvalSubtask sub = new EvalSubtask();
sub.setTaskId(task.getId());
sub.setCaseId(caseId);
sub.setStatus("pending");
subtaskService.save(sub);
// 4) 投到 MQ
mqSender.sendSubtask(sub.getId());
}
}
return batch.getId();
}
这段代码的关键是入队前先把子任务落库——先持久化再投递,消费者即使瞬间挂掉,任务也能从库里捞回来重投,一条都不丢。batch_id 是用户查询入口,task_id 是按模型维度的进度,subtask_id 是最小重试单元,三层各司其职。
四、并发控制
拆分做得再细,不加控制的并发照样会把平台打崩。我们踩过真坑:放开跑 200 个并发,模型供应商五分钟就把我们限流了。所以并发控制采用”全局信号量 + 单模型信号量”双层结构:
@Component
public class SubtaskExecutor {
private final Semaphore globalLimit = new Semaphore(50); // 全局 50 并发
private final Map<String, Semaphore> modelLimit = Map.of(
"gpt-4o", new Semaphore(20), // 单模型 20 并发
"claude", new Semaphore(15)
);
@RabbitListener(queues = "eval.subtask")
public void run(SubtaskMessage msg) {
Semaphore modelSem = modelLimit.get(msg.getModelCode());
try {
globalLimit.acquire();
modelSem.acquire();
processSubtask(msg);
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
} finally {
modelSem.release();
globalLimit.release();
}
}
}
全局 50 防止打爆平台,单模型 20/15 防止触发上游限流,两层都拿到才执行,执行完按逆序释放。这套参数是压测后调出来的,不同供应商的限流阈值不一样,上线前最好先实测一轮再定。
五、进度管理
并发跑起来之后,用户最关心的就是进度。进度管理分三层,查询路径各不同:批次进度给用户看总览,任务进度按模型拆分,子任务进度用于内部排查:
// 批次进度(用户最常看)
public BatchProgress getBatchProgress(long batchId) {
EvalBatch b = batchService.getById(batchId);
int total = b.getTotalCount();
int done = b.getDoneCount() + b.getFailCount();
int percent = total == 0 ? 100 : done * 100 / total;
return new BatchProgress(b.getId(), percent, done, total,
b.getFailCount(), b.getStatus());
}
// 任务进度(按模型)
public List<TaskProgress> getTaskProgress(long batchId) {
return taskService.listByBatch(batchId).stream()
.map(t -> new TaskProgress(t.getId(), t.getModelCode(),
t.getDoneCount() * 100 / t.getTotalCount(),
t.getStatus()))
.toList();
}
进度不能只靠查询,完成事件要主动推。我们用 MQ 解耦完成通知,消费者收到后原子累加父级计数,再经 WebSocket 推给前端:
@RabbitListener(queues = "eval.subtask.complete")
public void onComplete(SubtaskCompleteMessage msg) {
// 1) 写库
subtaskService.markDone(msg.getSubtaskId(), ...);
// 2) 累加父级计数(原子 SQL UPDATE ... SET done_count = done_count + 1)
taskService.incDone(msg.getTaskId());
batchService.incDone(msg.getBatchId());
// 3) WebSocket 推
websocketPush.pushToUser(msg.getUserId(),
new ProgressEvent(batchId, percent, ...));
}
这里必须强调:DB 是进度的唯一真理源,WebSocket 只负责实时展示——断线重连后重新拉一遍 DB 就能对齐,不会出现前端数字和后端对不上的情况。
六、失败重试
批量任务跑起来,失败是常态。子任务粒度的重试是代价更小的补救:一条用例失败只重跑一条,而不是整个批次推倒重来。核心是 retrySubtask,超过三次就转 fail_permanent:
// 单个子任务重试
public void retrySubtask(long subtaskId) {
EvalSubtask sub = subtaskService.getById(subtaskId);
if (sub.getRetryCount() >= 3) {
sub.setStatus("fail_permanent");
subtaskService.update(sub);
return;
}
sub.setRetryCount(sub.getRetryCount() + 1);
sub.setStatus("pending");
subtaskService.update(sub);
mqSender.sendSubtask(sub.getId());
}
// 整批失败任务重试
public void retryFailed(long batchId) {
List<EvalSubtask> failed = subtaskService.listByBatchAndStatus(
batchId, "fail");
for (EvalSubtask sub : failed) {
retrySubtask(sub.getId());
}
}
整批重试也简单,按状态捞失败的子任务逐个重投。重试时要留一手——原始错误信息不能覆盖:
ALTER TABLE eval_subtask ADD COLUMN first_error VARCHAR(512);
-- 重试时只覆盖 status,第一次错误永久保留
first_error 只写第一次的错误,重试只改 status,这样即使最终跑通了,也能复盘当初为什么失败。
七、关键设计取舍
方案跑顺之后,把关键决策回头看一遍,很多地方当时纠结过:
| 决策 | 方案 | 取舍 |
|---|---|---|
| 任务粒度 | 1 case × 1 model = 1 subtask | 粒度细,重试代价小 |
| 持久化时机 | 入队前 + 完成后 | 入队前持久化保证不丢,完成后更新 |
| 进度同步 | DB 累加 + WebSocket 推 | DB 是真理源,WebSocket 只做实时 |
| 失败策略 | 自动重试 3 次 | 指数退避 5s/30s/120s |
| 取消 | subtask.status = ‘cancelled’ | 消费者发现 status != ‘pending’ 跳过 |
这张表基本就是批量编排的答案——粒度、持久化时机、进度同步方式、失败策略、取消机制,五个决策点定了,剩下的都是实现细节。
八、踩过的坑
能跑通和跑得稳是两回事,踩坑记录比成功经验更有含金量:
- MQ 堆积:批量提交 10 万条全塞 MQ,消费者跟不上。要做”批量入库 + 批量投递”,不要每条一次 INSERT + 一次 send。
- 进度数累加竞争:多个消费者同时
done_count + 1不带条件会丢更新。用UPDATE ... SET done_count = done_count + 1 WHERE id = ?的原子 SQL。 - WebSocket 推送风暴:5000 个 subtask 全部完成,5000 次推送把 WS 打爆。批量聚合:攒 1 秒或 50 条再推一次。
- 失败重试雪崩:重试全在 5s 后,5000 条同时打模型触发限流。jitter 必须加(±50%)。
- 取消不生效:用户点”取消”,但 worker 正在跑 LLM 调用,停不下来。要在 LLM 调用的 Future 上做
cancel(true),配合 OkHttp 的 Call.cancel。
五个坑横跨 MQ、数据库、WebSocket、重试、取消五个环节,任何一个单点处理不好,批量任务都会变成线上事故。
九、可视化看板
最后给产品一个可视化的出口。批量任务页分成四块,用户不用猜进度:
- 顶部:批次总进度条 + 数字
- 中部:每个 model 一行,进度 + 成功率
- 底部:失败任务表格,可勾选重跑
- 右侧:实时日志流(WebSocket)
页面本身不复杂,但有了它,用户对”还要等多久”有了明确预期,咨询量下降得很明显。到这里,批量任务的拆分、并发、进度、重试、看板就串成了一条完整的链路。
常见问题(FAQ)
Q1:批次上限是多少?
单批 ≤ 5000 用例 × 5 模型 = 25000 子任务。再大要拆批,避免单批耗时长无法取消。
Q2:subtask 表会很大吗?
大。平台 6 个月 8000 万行,做按月分区 PARTITION BY RANGE (YEAR(created_at)*100 + MONTH(created_at))。
Q3:跑批时用户关页面会丢进度吗?
不会。进度在 DB 持久化,用户回来查询直接看最新状态。WebSocket 断线也不影响最终结果。