向量数据库实战:选型、调优与落地~系列文章18:百万级向量数据的写入优化:批量导入、增量更新的最佳实践

📅 2026/7/25 16:52:15
向量数据库实战:选型、调优与落地~系列文章18:百万级向量数据的写入优化:批量导入、增量更新的最佳实践
百万级向量数据的写入优化批量导入、增量更新的最佳实践 本文是《向量数据库实战选型、调优与落地》专栏第 18 篇⏱️阅读时间约 12 分钟 开篇写入性能为什么重要很多团队在搭建向量数据库时把精力都放在了查询优化上——却忽略了写入性能同样关键 首次数据导入百万条数据写入要多久增量更新每天新增 10 万条怎么不影响在线查询数据重建换了嵌入模型全量数据要重新向量化怎么高效导入写入慢 上线慢 用户体验差 写入性能基线测试测试环境Milvus 2.48C16GSSD100 万条 1024 维向量写入方式速度总耗时CPU内存逐条写入2k/s8.3 min15%0.5GBbatch10012k/s1.4 min35%1.2GBbatch100025k/s40s⚡70%4GBbatch500024k/s42s85%12GB并发 batch1000×480k/s12.5s95%8GB结论并发批量写入比逐条写入快 40 倍️ 全量导入最佳实践方案一标准批量导入importtimeimportnumpyasnpfrompymilvusimportCollectiondefbulk_insert(collection,texts,embeddings,batch_size1000):标准批量导入totallen(texts)start_timetime.time()foriinrange(0,total,batch_size):batch_textstexts[i:ibatch_size]batch_embeddingsembeddings[i:ibatch_size]collection.insert([batch_texts,batch_embeddings])if(i//batch_size)%100:progressmin(ibatch_size,total)/total*100elapsedtime.time()-start_time speedprogress*total/100/elapsedprint(f进度:{progress:.1f}%, 速度:{speed:.0f}条/秒)elapsedtime.time()-start_timeprint(f导入完成总耗时:{elapsed:.1f}s, 平均速度:{total/elapsed:.0f}条/秒)# 使用bulk_insert(collection,all_texts,all_embeddings,batch_size1000)方案二并发批量导入importconcurrent.futuresfrompymilvusimportCollectiondefconcurrent_insert(collection,data,batch_size1000,max_workers4):并发批量导入batches[]foriinrange(0,len(data),batch_size):batches.append(data[i:ibatch_size])total_written0start_timetime.time()withconcurrent.futures.ThreadPoolExecutor(max_workersmax_workers)asexecutor:futures[]forbatchinbatches:texts[item[text]foriteminbatch]embeddings[item[embedding]foriteminbatch]futureexecutor.submit(collection.insert,[texts,embeddings])futures.append(future)forfutureinconcurrent.futures.as_completed(futures):resultfuture.result()total_writtenresult.insert_count elapsedtime.time()-start_timeprint(f并发导入完成总耗时:{elapsed:.1f}s)returntotal_written方案三Bulk Insert文件导入# 对于超大数据集使用文件导入更高效importjson# Step 1: 准备 JSON 文件defprepare_bulk_files(data,output_dir,file_size100000):将数据分成多个 JSON 文件foriinrange(0,len(data),file_size):batchdata[i:ifile_size]file_pathf{output_dir}/batch_{i//file_size}.jsonrows[]foriteminbatch:rows.append({text:item[text],embedding:item[embedding]})withopen(file_path,w)asf:json.dump({rows:rows},f)# Step 2: 通过 Milvus Bulk Insert API 导入frompymilvusimportutility# 将文件上传到 MinIO或指定存储# ...# 触发 Bulk Inserttask_idsutility.do_bulk_insert(collection_nameknowledge_base,files[batch_0.json,batch_1.json,batch_2.json])# 监控进度fortask_idintask_ids:stateutility.get_bulk_insert_state(task_id)print(f任务{task_id}:{state.state})文件导入 vs API 导入对比项API 导入文件导入Bulk Insert适合数据量 1000 万 1000 万✅速度25-80k/s100-500k/s⚡内存占用较高较低灵活性高可实时低批量操作断点续传需自己实现自动支持 增量更新策略方案一直接 Upsert# Milvus 支持 Upsert有则更新无则插入collection.upsert([[新文档内容],[new_embedding]])方案二分区增量# 使用分区管理增量数据frompymilvusimportCollection# 按日期创建分区collection.create_partition(data_2025_01)collection.create_partition(data_2025_02)collection.create_partition(data_2025_03)# 新数据写入最新分区partitioncollection.partition(data_2025_03)partition.insert([new_texts,new_embeddings])# 查询时搜索所有分区collection.load()# 加载所有分区resultscollection.search(...)方案三双 Buffer 切换┌─────────────────────────────────────────────────────────┐ │ 双 Buffer 增量更新 │ ├─────────────────────────────────────────────────────────┤ │ │ │ 在线 Buffer A当前服务 │ │ ┌──────────────────────────────────────┐ │ │ │ 存量数据 已索引 │ │ │ │ → 正常提供查询服务 │ │ │ └──────────────────────────────────────┘ │ │ │ │ 离线 Buffer B构建中 │ │ ┌──────────────────────────────────────┐ │ │ │ 增量数据写入 索引构建 │ │ │ │ → 不影响 Buffer A 的查询 │ │ │ └──────────────────────────────────────┘ │ │ │ │ 切换B 构建完成后原子切换到 B 提供服务 │ │ A 变成离线 Buffer接收下一批增量 │ │ │ └─────────────────────────────────────────────────────────┘⚠️ 写入优化的常见坑坑 1写入时不建索引# ❌ 错误边写边建索引forbatchinbatches:collection.insert(batch)collection.create_index(...)# 每次写入都建索引极慢# ✅ 正确先写完最后建一次索引forbatchinbatches:collection.insert(batch)# 全部写完后建索引collection.create_index(embedding,index_params)坑 2忽略 Flush# Milvus 写入后需要 Flush 才能持久化collection.insert(data)collection.flush()# 确保数据持久化# 批量导入后统一 Flushforbatchinbatches:collection.insert(batch)collection.flush()# 最后统一 Flush 一次坑 3内存溢出# ❌ 错误一次性加载所有数据到内存all_dataload_all_data()# 100GB 数据直接 OOM# ✅ 正确流式读取defstream_data(file_path,batch_size1000):batch[]withopen(file_path)asf:forlineinf:batch.append(json.loads(line))iflen(batch)batch_size:yieldbatch batch[]ifbatch:yieldbatchforbatchinstream_data(large_file.json):collection.insert(batch) 写入性能优化 Checklist优化项推荐值效果batch_size1000比逐条快 12x并发数4单机/ 8集群比单线程快 3-4x先写后建索引全部写完再建避免重复建索引Bulk Insert 1000万条数据比 API 快 5-10x流式读取大文件避免 OOM统一 Flush最后一次性减少 IO 次数 本篇核心要点回顾要点说明批量写入batch_size1000 是最佳平衡点并发写入单机 4 线程集群 8 线程超大数据用 Bulk Insert文件导入增量更新Upsert / 分区 / 双 Buffer常见坑先写后建索引、统一 Flush、流式读取下篇预告《向量数据库 RAG 融合实战构建企业级知识库的完整链路 》有问题欢迎评论区讨论觉得有用请点赞收藏 作者高炉炼铁智能化技术研究者专注钢铁冶金与人工智能 交叉领域。 如果觉得有帮助请点赞、收藏、转发版权归作者所有未经许可请勿抄袭套用商用(或其它具有利益性行为)。 关注专栏不错过后续精彩内容