百万级AI教案生成系统:批量生产、多线程并发与优雅暂停重启实践

一、背景与挑战

在教育领域,一个教材版本往往包含数十个章节,每个章节又有多达十多个课时。要为每个课时分别生成教案、学案或作业,一笔任务动辄包含几百甚至上千个子任务。每个子任务需要调用大模型 API(如豆包 Seed 1.6),单次调用耗时 10~60 秒不等。如果串行执行,一笔大任务跑完可能需要数小时。

因此,系统在设计上需要解决三个核心问题:

批量生产:一笔父任务要能自动拆分为 N 个子任务,并逐条追踪状态
多线程并发:多个子任务同时执行,充分利用 CPU 和网络 IO,将总耗时压缩到分钟级
暂停与重启:用户可以随时暂停正在执行的大任务,后续一键恢复,已完成的子任务不重复执行
下面分享这套系统的设计与实现。

二、整体架构设计

2.1 核心组件

组件 作用
RocketMQ 任务消息驱动,解耦前端请求与后端执行
Redis + Redisson 分布式锁,防止集群环境下子任务被重复执行
ThreadPoolTaskExecutor 自定义线程池,承载子任务的多线程并发
CountDownLatch 并发计数器,确保主线程等待所有子线程完成后统一收尾
MyBatis 数据持久化,子任务状态精确追踪

2.2 数据模型:三表联动

任务的状态流转依赖三张核心表:

1
2
3
4
5
6
7
8
t_xj_task (父任务)
├── status: 0=待生成, 1=生成中, 2=暂停中, 3=已完成, 5=已完成(全部)

├── t_xj_task_info (子任务详情,一个章节+课时 = 一条)
│ └── taskStatus: 0=待执行, 1=执行中, 2=已完成, 3=失败

└── t_xj_task_sub (子任务生成的AI结果)
└── subStatus: 0=待送审, 1=已送审

一条父任务被拆分为多条 t_xj_task_info,每条 info 对应一个 (章节 + 课时) 组合。每条 info 最终生成一条 t_xj_task_sub 存储 AI 返回的内容。

三、批量生产:任务的拆分与分发

3.1 创建任务

前端提交一个包含多个章节和课时的 JSON:

1
2
3
4
5
6
{
"sourceArray": [
{ "sourceId": 101, "keshi": "1,2,3" },
{ "sourceId": 102, "keshi": "1,2" }
]
}

后端在 /add 接口中做三件事:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
// 1. 保存父任务
tXjTask.setStatus("0"); // 待生成
tXjTaskService.save(tXjTask);

// 2. 遍历章节和课时,批量创建子任务 info
for (章节 : sourceArray) {
for (课时 : keshi.split(",")) {
TXjTaskInfo info = new TXjTaskInfo();
info.setTaskId(tXjTask.getId());
info.setTaskStatus("0"); // 待执行
list.add(info);
}
}
tXjTaskInfoService.saveBatch(list); // 批量插入

// 3. 发送MQ消息,触发异步执行
rocketProduer.sendMessage(Jiao_Xue_Task_Topic, String.valueOf(tXjTask.getId()));

设计要点:frontend 请求只是”登记任务 + 发一条 MQ 消息”即返回,真正的 AI 生成全部放在 MQ 消费端异步执行,前端不会阻塞等待。

四、多线程并发:CountDownLatch + 线程池 + Redisson 分布式锁

这是整个系统的核心亮点。MQ 消费端收到父任务 ID 后:

4.1 一次性拉取所有待执行子任务

1
2
3
4
5
List<TXjTaskInfo> allWaitList = tXjTaskInfoService.list(
new LambdaQueryWrapper<TXjTaskInfo>()
.eq(TXjTaskInfo::getTaskId, taskId)
.eq(TXjTaskInfo::getTaskStatus, 0) // 只取待执行的
);

关键设计:不逐条轮询,而是一次性把所有 status=0 的子任务全查出来,批量提交到线程池。

4.2 CountDownLatch 并发控制

1
2
3
4
5
6
7
8
9
10
11
12
13
14
CountDownLatch latch = new CountDownLatch(allWaitList.size());

for (TXjTaskInfo taskInfo : allWaitList) {
taskExecutor.execute(() -> {
try {
// ... 执行子任务 ...
} finally {
latch.countDown(); // 每个线程完成时计数减一
}
});
}

// 主线程阻塞,等待所有子线程执行完毕(最长等5分钟)
latch.await(5, TimeUnit.MINUTES);

这样,所有子任务并发执行,主线程等待所有子任务完成后,统一将父任务标记为”已完成”。

4.3 Redisson 分布式锁:防止重复执行

在集群多实例部署的情况下,同一个子任务可能被多个实例同时消费。这里用 Redis 分布式锁来保护:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
RLock lock = redissonClient.getLock("lock:item:" + infoItem.getId());
boolean lockOk = lock.tryLock(0, 30, TimeUnit.SECONDS); // 拿不到锁直接放弃
if (!lockOk) {
return; // 其他实例正在执行,直接跳过
}

// 加锁后,二次校验子任务状态
TXjTaskInfo check = tXjTaskInfoService.getById(infoItem.getId());
if (!"0".equals(check.getTaskStatus())) {
return; // 已被执行
}

// 标记执行中
infoItem.setTaskStatus("1");
tXjTaskInfoService.updateById(infoItem);

// 调用AI生成...
// 保存结果...

设计要点:

tryLock(0, …) — 不等待,拿不到锁立即返回,避免线程堆积
锁 30 秒自动过期 — 防止异常死锁
加锁后二次校验状态 — 典型的 Double Check 模式

4.4 线程池配置

1
2
3
4
5
6
7
8
9
10
@Bean("taskExecutor")
public ThreadPoolTaskExecutor taskExecutor() {
int cpuCores = Runtime.getRuntime().availableProcessors();
ThreadPoolTaskExecutor executor = new ThreadPoolTaskExecutor();
executor.setCorePoolSize(cpuCores); // 核心线程 = CPU核数
executor.setMaxPoolSize(cpuCores * 2); // 最大线程 = 2倍核数
executor.setQueueCapacity(200); // 队列容量
executor.setRejectedExecutionHandler(new CallerRunsPolicy()); // 拒绝策略:调用者运行
return executor;
}

采用 CPU 核数动态适配,队列满时由调用者线程兜底执行,不会丢任务。

五、暂停与重启:优雅的中断与恢复

5.1 暂停机制

用户可以随时调用 /stopTask 接口暂停任务:

1
2
3
4
5
// 前端只需传 id 和 status="2"
@RequestMapping("/stopTask")
public Msg stopTask(@RequestBody TXjTask tXjTask) {
tXjTaskService.updateById(tXjTask); // 将status更新为"2"(暂停中)
}

但仅仅改数据库是不够的——已经在执行中的线程怎么办?答案是在每个执行阶段都设置暂停检查点:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
// 检查点1:MQ消费入口
if ("2".equals(parentTask.getStatus())) {
return; // 整条任务已暂停,直接退出
}

// 检查点2:每个子线程开始前
TXjTask liveParent = tXjTaskService.getById(taskId);
if ("2".equals(liveParent.getStatus())) {
return; // 暂停了,跳过
}

// 检查点3:AI调用前(最关键)
TXjTask pauseCheck = tXjTaskService.getById(taskId);
if ("2".equals(pauseCheck.getStatus())) {
throw new RuntimeException("全局任务暂停,终止AI生成");
}

这样,暂停指令下发后:

未开始的子任务直接跳过(检查点1、2)
正在调用AI的通过异常中断,标记为失败(检查点3)
已完成的不受影响(状态已是2,不会重复执行)

5.2 重启机制

重启同样简单,调用 /startTask:

1
2
3
4
5
@RequestMapping("/startTask")
public Msg startTask(@RequestBody TXjTask tXjTask) {
tXjTaskService.updateById(tXjTask); // 恢复status(如改为"0"或"1")
rocketProduer.sendMessage(Jiao_Xue_Task_Topic, tXjTask.getId()); // 重新发MQ
}

重新发送 MQ 后,消费端只查询 taskStatus = 0(待执行)的子任务,已完成和失败的不受影响,实现了断点续传。

5.3 状态机总结

1
2
3
4
5
6
7
    创建任务

┌── 待生成(0) ──→ 生成中(1) ──→ 已完成(5)
│ ↑ │
└── 暂停(2) ←─────────┘
│ (重启后回到生成中)
└──→ 生成中(1)

六、关键设计总结

设计点 落地方案 核心收益
异步解耦 前端登记入库存任务 + MQ异步消费调度 前端无需等待耗时生成逻辑,实现秒级响应,核心生成流程异步隔离,不阻塞前端请求
批量任务精细化拆分 大任务按章节、课时维度拆分为多条独立子任务Info 任务粒度足够精细,支持单条任务独立追踪、暂停、重试,为断点续传、精准管控打下基础
多线程并发提速 自定义 ThreadPoolTaskExecutor + CountDownLatch 批量并行执行 海量子任务并行处理,将传统串行执行的小时级耗时,压缩至分钟级,大幅提升批量生成效率
分布式并发安全 Redisson 零等待抢占锁 + 双重校验防并发 集群多实例部署下,杜绝同一子任务重复执行、并发脏写,保障分布式任务执行唯一性
任务即时暂停机制 任务入口、线程启动、AI分片调用前三层状态检查点 暂停指令实时生效,无延迟、无残留执行逻辑,精准锁定单任务,不丢任务、不重复执行
服务重启断点续传 启动仅查询待执行状态(status=0)未完成子任务 服务宕机、重启后自动跳过已完成、已终止任务,接续未完成任务执行,无需全量重跑
服务优雅关闭 CountDownLatch 超时等待 + CallerRunsPolicy 饱和策略 服务停机时等待正在执行的任务收尾,未执行任务不丢失、不线程卡死,实现无损停机、优雅下线