☰
向量构建延迟调优:批量(Batch)计算 Embedding 与异步索引构建实战
2026/10/11 1:57:55 网站建设 项目流程

做知识库或者商品语义搜索(RAG)的项目,上线初期往往都挺平稳,但随着业务数据量上来,痛点就会集中爆发在“离线初始化”和“增量同步”两个环节。

上个月我们团队接了一个需求,要把商户后台积累的近 50 万份专业售后工单切片后存入向量数据库,作为智能客服的知识底座。按每篇文档切分出 8~12 个 Chunk(文本块)来算,总共涉及近 500 万条向量数据的生成与索引构建。

一开始负责这个模块的同事写得很简单粗暴:遍历文本列表,一条一条调用远程的 Embedding 模型接口,拿到 1536 维的浮点数组,再调用 Milvus 的单条插入接口。跑了一晚上,第二天早上过来看进度条,才完成了不到 3%,期间还被云厂商的 API 报了一堆429 Too Many Requests限流错误。

这篇文章聊聊我们后来是如何通过批量计算(Batch Embedding)、异步流水线以及向量索引延迟构建,将原本需要跑两天的初始化任务压缩到两小时以内的。


为什么单条同步构建会慢如蜗牛?

很多同学觉得慢是因为网络延迟,这只说对了一半。真正拖垮整体吞吐的,是以下三个层面的资源浪费:

  1. HTTP 连接与协议开销:单条文本通常只有 200~500 个字,而一次完整的 HTTP/1.1 或 HTTP/2 请求头、TLS 加密开销甚至比载荷还大。500 万次往返(RTT),哪怕单次网络延迟只有 15ms,光网络消耗的时间就高达 20 多个小时。
  2. GPU 推理没有吃满并行算力:大模型无论是做文本生成还是向量特征提取,GPU 的张量核心最擅长做矩阵并行运算。单条文本送进去,GPU 显存带宽和计算核心利用率往往不足 10%;而一次性送入 64 或 128 条文本,GPU 耗时可能只增加 20%~30%,吞吐量却能翻 10 倍以上。
  3. 向量数据库的 HNSW 动态插入锁竞争:像 Milvus、pgvector 等数据库,在建立图索引(如 HNSW)时,每插入一个新向量,都需要在现有的多层图中做局部邻域搜索并更新连边。一条一条插入,索引反复触发加锁与拓扑重平衡;而批量插入时,引擎可以利用内存缓冲区批量建图,效率完全是两个量级。

批量化与流水线架构设计

为了把吞吐拉到极限,我们把整个构建链路拆分为三个解耦的阶段:

  1. 分批打包器(Batch Packer):按文本条数和最大 Token 双重水位线打包,避免单批次超载。
  2. 异步 Embedding 执行器:结合虚拟线程与信号量进行并发限流,批量请求模型。
  3. 分批落盘与索引延后:先写入扁平结构(Flat)或分批写入暂存表,全部灌完后再统一触发索引构建(Build Index)。

1. 动态 Batching 打包器

不同厂商的 Embedding 接口对单次 Batch 的限制不同,既有条数限制(例如最多 64 条),也有 Token 数量限制(例如一次请求所有文本总 Token 不能超过 8192)。单纯按条数切片很容易因为遇到长文本而直接触发服务端报错。

我们写了一个简单的滑动聚合打包器:

public class TextBatchPacker { private final int maxBatchSize; private final int maxBatchTokens; public TextBatchPacker(int maxBatchSize, int maxBatchTokens) { this.maxBatchSize = maxBatchSize; this.maxBatchTokens = maxBatchTokens; } public List<List<DocumentChunk>> pack(List<DocumentChunk> chunks) { List<List<DocumentChunk>> batches = new ArrayList<>(); List<DocumentChunk> currentBatch = new ArrayList<>(); int currentTokens = 0; for (DocumentChunk chunk : chunks) { int chunkToken = estimateTokenCount(chunk.getContent()); // 超过条数上限或超过 Token 上限,立即封箱 if (currentBatch.size() >= maxBatchSize || (currentTokens + chunkToken) > maxBatchTokens) { if (!currentBatch.isEmpty()) { batches.add(new ArrayList<>(currentBatch)); currentBatch.clear(); currentTokens = 0; } } currentBatch.add(chunk); currentTokens += chunkToken; } if (!currentBatch.isEmpty()) { batches.add(currentBatch); } return batches; } private int estimateTokenCount(String text) { // 中文粗略预估:字符数 * 1.3,或使用分词器精确计数 return (int) (text.length() * 1.3); } }

2. 基于虚拟线程的并发 Embedding 调度

在 Java 21+ 时代,处理高并发网络 I/O 任务不再需要维护复杂的响应式链条或重型线程池。直接利用虚拟线程(Virtual Thread)配合Semaphore,既能保证代码的直观可读,又能轻松压满网络带宽。

@Service public class AsyncEmbeddingService { private final EmbeddingModel embeddingModel; private final Semaphore rateLimiter = new Semaphore(16); // 控制最大并发批次数 private final VectorDatabaseClient vectorDbClient; public AsyncEmbeddingService(EmbeddingModel embeddingModel, VectorDatabaseClient vectorDbClient) { this.embeddingModel = embeddingModel; this.vectorDbClient = vectorDbClient; } public void processAllChunks(List<DocumentChunk> allChunks) { TextBatchPacker packer = new TextBatchPacker(64, 8000); List<List<DocumentChunk>> batches = packer.pack(allChunks); try (var executor = Executors.newVirtualThreadPerTaskExecutor()) { List<CompletableFuture<Void>> futures = batches.stream() .map(batch -> CompletableFuture.runAsync(() -> processSingleBatch(batch), executor)) .toList(); CompletableFuture.allOf(futures.toArray(new CompletableFuture[0])).join(); } } private void processSingleBatch(List<DocumentChunk> batch) { try { rateLimiter.acquire(); List<String> texts = batch.stream().map(DocumentChunk::getContent).toList(); // 批量请求 Embedding 模型 EmbeddingResponse response = embeddingModel.embedForResponse(texts); List<Embedding> embeddings = response.getResults(); // 组装带向量的实体 List<VectorRecord> records = new ArrayList<>(batch.size()); for (int i = 0; i < batch.size(); i++) { records.add(new VectorRecord( batch.get(i).getId(), batch.get(i).getDocumentId(), embeddings.get(i).getOutput(), batch.get(i).getContent() )); } // 批量写入向量库(此时底层只做 Append,不建立强一致索引) vectorDbClient.insertBatch(records); } catch (InterruptedException e) { Thread.currentThread().interrupt(); throw new RuntimeException("处理被中断", e); } finally { rateLimiter.release(); } } }

生产落地的硬核细节与避坑指南

1. 避免“毒药文本”击垮整批任务

在批量调用 Embedding 时,只要批次里有一条文本格式异常(例如包含超长特殊字符、不可见控制符,导致模型返回 400 Bad Request),整个包含 64 条文本的 Batch 都会直接报错失败。
如果简单粗暴地做整批重试,不仅徒劳无功,还会陷入死循环。

我们的生产处理策略是:

  • 批量请求如果捕获到非 429、非 5xx 的客户端错误(400~499),立即对当前批次降级为“二分法重试”或“逐条单发”。
  • 把无法通过模型推理的异常文本记录到死信表(Dead Letter),记录其原始 doc_id 并跳过,绝不卡死主数据流。

2. 内存反压:小心堆内存被 List 撑爆

浮点向量是非常占内存的。一个 1536 维的float数组在 JVM 堆里占约 6KB 内存。如果有 100 万条向量暂存在内存队列中,光向量本身就需要近 6GB 内存,还没算对象的元数据和包裹对象开销。

在做大规模数据迁移时,绝对不能一次性将全量 Chunk 加载进内存。必须使用分页游标(如 MyBatis 的Cursor或 Spring Data 的分块滚动拉取),或者利用阻塞队列维护固定容量的滑动窗口(比如最多在堆中保留 500 个 Batch)。

3. 向量库索引构建时机:先插数据,后建索引

不管是 Milvus 还是自建的 pgvector,在写入几百万条数据时,如果先建好了 HNSW 索引再去逐批写入,写入性能会呈现断崖式下跌。因为每写一批,CPU 都在拼命计算图的连通度和重平衡。

黄金法则是:在大量数据导入之前,先把 Collection 或表配置为未建立索引状态(或者先 drop 索引),只保留 raw 向量数据。当全部批量数据落盘完成后,再发起一次统一的createIndex指令。
这时候,向量数据库会利用多核 CPU 的并行计算能力,一口气构建整棵 HNSW 图。实测对比表明,先插后建相比边插边建,总耗时能直接缩短 60% 以上。

工程落地没有银弹,把批量吞吐压榨到位,把异常隔离做好,系统自然既快又稳。

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询