Appearance
异步化与批处理:灌库不能拖垮在线服务
1. 本节产出
一个异步灌库管道:上传即返回(任务 ID)、后台队列消费、可查进度、失败可重试。并且在线检索的延迟不受灌库影响(这是本节的核心验收)。
2. 前置依赖
- 02-14 增量更新
- 02-05 Embedding 批处理
- Redis(做任务队列)或 Spring 的线程池
3. 为什么同步灌库一定会出事
同步灌库的三个问题,按严重程度排序:
| 问题 | 表现 |
|---|---|
| 拖垮在线服务 | 灌库占满 CPU/连接池,用户检索延迟从 50ms 涨到 3s |
| 请求超时 | 灌 100 篇文档要 10 分钟,HTTP 连接早就断了 |
| 失败即全丢 | 中途出错,已处理的 80 篇白做(或留下脏数据) |
第一个问题最严重且最隐蔽:它不会报错,只是「系统变慢了」。而灌库通常发生在业务低峰或运维手动触发,很容易被误判为「系统本来就这么慢」。
解决思路就一句话:离线任务和在线服务必须资源隔离。
4. 核心原理
4.1 生产者-消费者模型
上传接口(在线)
│ 立刻返回 taskId
▼
任务表 / 队列(持久化)
│
▼ 后台消费者(独立线程池)
解析 → 切分 → 向量化 → 写入
│
▼
更新任务状态(进度、成功/失败)任务必须持久化。只放内存队列的话,应用重启后任务全丢,而且你无法回答「上次灌到哪了」。
4.2 资源隔离的三种做法
| 做法 | 机制 | 适用 |
|---|---|---|
| 独立线程池 | 灌库用专用线程池,限制并发数 | 单体应用,推荐起步 |
| 独立进程/服务 | 灌库拆成独立微服务 | 规模较大时 |
| 限流 + 背压 | 灌库速率自适应,检测到在线延迟上升就降速 | 推荐叠加 |
背压是最容易被忽略的一层:即使做了线程池隔离,灌库仍然会和在线服务争抢数据库连接、CPU、Embedding 配额。做法:监控在线检索的 P95,超过阈值就自动降低灌库并发。
4.3 并发度怎么定
灌库吞吐 = 并发数 × 单文档速度
受限因素:
1. Embedding 厂商的 QPS/配额(最常见瓶颈)
2. 数据库连接池大小
3. CPU(解析 PDF 是 CPU 密集)经验起点:并发 2~4。别一上来设 16——Embedding 厂商的限流会让你全部失败重试,反而更慢。
4.4 进度与幂等
| 要素 | 做法 |
|---|---|
| 进度可查 | 任务表记录 total / done / failed |
| 单篇幂等 | 按 sourceId + 内容哈希跳过(02-14) |
| 失败重试 | 单篇失败不影响整批,记录失败原因 |
| 任务级重试 | 支持「重试失败项」,不重做成功的 |
5. 代码走查
5.1 任务表
sql
CREATE TABLE ingest_task (
task_id VARCHAR(64) PRIMARY KEY,
tenant_id VARCHAR(64) NOT NULL,
total INT DEFAULT 0,
done INT DEFAULT 0,
failed INT DEFAULT 0,
status VARCHAR(16), -- PENDING/RUNNING/DONE/PARTIAL/FAILED
error_msg TEXT,
created_at TIMESTAMPTZ DEFAULT now(),
updated_at TIMESTAMPTZ DEFAULT now()
);
CREATE TABLE ingest_task_item (
task_id VARCHAR(64),
source_id VARCHAR(200),
file_path TEXT,
status VARCHAR(16), -- PENDING/DONE/FAILED
error_msg TEXT,
PRIMARY KEY (task_id, source_id)
);5.2 独立线程池
java
// ch02-rag/src/main/java/com/aitech/rag/config/IngestExecutorConfig.java
@Configuration
public class IngestExecutorConfig {
@Bean("ingestExecutor")
public ThreadPoolTaskExecutor ingestExecutor(IngestProps props) {
ThreadPoolTaskExecutor ex = new ThreadPoolTaskExecutor();
ex.setCorePoolSize(props.concurrency()); // 建议 2~4
ex.setMaxPoolSize(props.concurrency() * 2);
ex.setQueueCapacity(500);
ex.setThreadNamePrefix("ingest-");
// 关键:队列满时由调用方线程执行,形成天然背压
ex.setRejectedExecutionHandler(new ThreadPoolExecutor.CallerRunsPolicy());
ex.initialize();
return ex;
}
}CallerRunsPolicy 是一个便宜好用的背压机制:队列满了就让提交任务的线程自己干,自然降低提交速度,而不是无限制堆积。
5.3 上传接口:立刻返回
java
// ch02-rag/src/main/java/com/aitech/rag/controller/IngestController.java
@RestController
@RequestMapping("/api/ingest")
public class IngestController {
@PostMapping(consumes = MediaType.MULTIPART_FORM_DATA_VALUE)
public IngestTaskResponse submit(@RequestParam("files") List<MultipartFile> files) {
String taskId = UUID.randomUUID().toString();
String tenantId = TenantContext.require();
List<IngestItem> items = files.stream()
.map(f -> new IngestItem(taskId, sourceIdOf(f), store(f)))
.toList();
taskRepo.create(taskId, tenantId, items);
ingestQueue.submit(taskId); // 异步触发,不等结果
return new IngestTaskResponse(taskId, items.size()); // 立刻返回
}
@GetMapping("/{taskId}")
public TaskStatus status(@PathVariable String taskId) {
return taskRepo.status(taskId);
}
@PostMapping("/{taskId}/retry")
public TaskStatus retryFailed(@PathVariable String taskId) {
ingestQueue.retryFailed(taskId); // 只重试失败项
return taskRepo.status(taskId);
}
}5.4 消费者
java
// ch02-rag/src/main/java/com/aitech/rag/ingest/IngestWorker.java
@Service
public class IngestWorker {
@Async("ingestExecutor")
public void run(String taskId) {
List<IngestItem> items = taskRepo.pendingItems(taskId);
Semaphore rate = new Semaphore(props.concurrency());
for (IngestItem item : items) {
try {
rate.acquire();
// 背压:在线延迟高时先等一等
backpressure.awaitIfOnlineSlow();
ingestService.ingest(toRequest(item)); // 幂等
taskRepo.markDone(item);
} catch (Exception e) {
log.error("灌库失败 sourceId={}", item.sourceId(), e);
taskRepo.markFailed(item, e.getMessage()); // 单篇失败不中断
} finally {
rate.release();
}
}
taskRepo.finish(taskId);
}
}「单篇失败不中断」是批处理的基本原则。一篇文档解析失败不应该让整个任务失败。
5.5 背压实现
java
// ch02-rag/src/main/java/com/aitech/rag/ingest/Backpressure.java
@Component
public class Backpressure {
private final MeterRegistry meters;
public void awaitIfOnlineSlow() throws InterruptedException {
double p95 = meters.timer("rag.query").takeSnapshot()
.percentileValue(0.95, TimeUnit.MILLISECONDS);
if (p95 > 300) { // 在线 P95 超过 300ms
log.info("在线延迟偏高({}ms),灌库降速", p95);
Thread.sleep(1000); // 简单退避
}
}
}6. 跑起来
bash
git checkout ch02-16-async-batch
docker compose up -d
mvn spring-boot:runbash
# 1. 提交 50 篇文档
curl -X POST http://localhost:8080/api/ingest -F "files=@d1.pdf" -F "files=@d2.pdf" ...
# 期望:立刻返回 {"taskId":"...","total":50}
# 2. 查进度
curl http://localhost:8080/api/ingest/{taskId}
# 期望:{"total":50,"done":23,"failed":1,"status":"RUNNING"}
# 3. 关键验收:灌库期间测在线检索延迟
while true; do
curl -s -o /dev/null -w "%{time_total}\n" -X POST http://localhost:8080/api/ask \
-d '{"q":"年假怎么算","tenantId":"acme"}'
done
# 期望:延迟稳定,不随灌库显著上升
# 4. 重试失败项
curl -X POST http://localhost:8080/api/ingest/{taskId}/retry| 检查项 | 通过标准 |
|---|---|
| 上传即返回 | 响应时间 < 500ms(与文档数无关) |
| 进度可查 | 能实时看到 done/failed |
| 在线延迟不受影响 | 灌库期间 P95 波动 < 20% |
| 单篇失败不中断 | 造一篇坏文档,其他仍完成 |
| 重试 | 只重做失败项,成功的被跳过 |
第三项是本节的核心验收,必须一边灌库一边测在线延迟。
7. 生产避坑
- 不要用在线请求的线程池做灌库。即使用了
@Async,如果没指定独立线程池,用的还是默认池——等于没隔离。做法:显式指定@Async("ingestExecutor"),并在配置里确认线程池参数。 - 任务状态必须持久化,不能只放内存。放内存的后果:应用重启后「正在灌的任务」状态丢失,你既不知道灌没灌完,也无法续跑——只能全量重来。
- 并发度不是越大越好。Embedding 厂商有限流,并发过大会导致大量 429,触发重试后反而更慢,还浪费配额。做法:从 2 开始压测,观察失败率与吞吐的关系,找到拐点。
8. 延伸与锚点
- 思考题:改了切分策略,怎么知道效果是变好还是变差?(02C 最核心的一课——答案在下一课时)
- 代码锚点:
git checkout ch02-16-async-batch - 下一课时:02-17 评测体系与回归
- 对应课件:L02-16 异步化与批处理