Skip to content

异步化与批处理:灌库不能拖垮在线服务 ​

1. 本节产出 ​

一个异步灌库管道:上传即返回(任务 ID)、后台队列消费、可查进度、失败可重试。并且在线检索的延迟不受灌库影响(这是本节的核心验收)。

2. 前置依赖 ​

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:run
bash
# 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. 生产避坑 ​

  1. 不要用在线请求的线程池做灌库。即使用了 @Async,如果没指定独立线程池,用的还是默认池——等于没隔离。做法:显式指定 @Async("ingestExecutor"),并在配置里确认线程池参数。
  2. 任务状态必须持久化,不能只放内存。放内存的后果:应用重启后「正在灌的任务」状态丢失,你既不知道灌没灌完,也无法续跑——只能全量重来。
  3. 并发度不是越大越好。Embedding 厂商有限流,并发过大会导致大量 429,触发重试后反而更慢,还浪费配额。做法:从 2 开始压测,观察失败率与吞吐的关系,找到拐点。

8. 延伸与锚点 ​