用Spring AI+Qdrant实现可恢复的RAG增量索引服务

📅 2026/7/24 21:22:22
用Spring AI+Qdrant实现可恢复的RAG增量索引服务
文章摘要本文实现一个面向企业知识库的增量索引服务通过文档Checksum和Chunk Hash识别变化只对新增或修改片段生成Embedding使用稳定Point ID写入Qdrant通过版本状态实现新旧索引原子切换失败任务可重试删除操作可幂等执行。示例使用Spring Boot、Spring AI和Qdrant重点展示工程结构、数据模型、任务状态和核心代码。一、目标架构上传或数据变更 → 文档版本记录 → 解析与分块 → Chunk差异计算 → Embedding → Qdrant写入READY → 完整性校验 → 版本切换 → 清理旧版本与缓存目标未变化Chunk不重复Embedding新版本未完成前不影响生产任务失败可以恢复删除操作幂等所有Point可以追踪来源支持多租户和版本过滤。二、项目结构rag-indexing-service ├── domain │ ├── KnowledgeDocument.java │ ├── DocumentVersion.java │ ├── KnowledgeChunk.java │ └── IndexTask.java ├── application │ ├── IncrementalIndexService.java │ ├── DocumentParser.java │ └── VersionSwitchService.java ├── infrastructure │ ├── QdrantChunkRepository.java │ ├── SpringAiEmbeddingGateway.java │ └── JdbcDocumentRepository.java └── web └── KnowledgeIndexController.java三、核心数据模型publicenumVersionStatus{UPLOADED,PARSING,INDEXING,READY,EFFECTIVE,EXPIRED,FAILED}publicrecordDocumentVersion(StringversionId,StringlogicalDocumentId,longversion,StringfileChecksum,VersionStatusstatus,InstantcreatedAt){}ChunkpublicrecordKnowledgeChunk(StringchunkId,StringlogicalDocumentId,StringversionId,StringsectionPath,intchunkIndex,Stringcontent,StringcontentHash,MapString,Objectmetadata){}四、稳定Point ID不要每次生成随机UUID。publicStringbuildPointId(KnowledgeChunkchunk){StringrawString.join(:,chunk.logicalDocumentId(),chunk.versionId(),chunk.sectionPath(),String.valueOf(chunk.chunkIndex()));returnsha256(raw);}稳定ID便于Upsert删除重试引用一致性校验。五、计算文件和Chunk HashpublicStringsha256(Stringcontent){try{MessageDigestdigestMessageDigest.getInstance(SHA-256);byte[]bytesdigest.digest(content.getBytes(StandardCharsets.UTF_8));returnHexFormat.of().formatHex(bytes);}catch(NoSuchAlgorithmExceptionexception){thrownewIllegalStateException(exception);}}解析前先比较文件Hash相同 → 跳过整个文档 不同 → 解析并比较Chunk Hash六、差异模型publicrecordChunkDiff(ListKnowledgeChunkadded,ListKnowledgeChunkmodified,ListKnowledgeChunkunchanged,ListKnowledgeChunkremoved){}比较逻辑publicChunkDiffdiff(ListKnowledgeChunkoldChunks,ListKnowledgeChunknewChunks){MapString,KnowledgeChunkoldByPatholdChunks.stream().collect(Collectors.toMap(this::logicalChunkKey,Function.identity()));MapString,KnowledgeChunknewByPathnewChunks.stream().collect(Collectors.toMap(this::logicalChunkKey,Function.identity()));// 根据逻辑位置与contentHash分类// 示例中省略集合拼装细节returncalculateDiff(oldByPath,newByPath);}逻辑Chunk Keysection_path chunk_index如果分块策略变化应提升chunking_version并执行全量重建而不是错误复用旧向量。七、Embedding网关publicinterfaceEmbeddingGateway{Listfloat[]embed(ListStringtexts);}Spring AI实现ComponentpublicclassSpringAiEmbeddingGatewayimplementsEmbeddingGateway{privatefinalEmbeddingModelembeddingModel;publicSpringAiEmbeddingGateway(EmbeddingModelembeddingModel){this.embeddingModelembeddingModel;}OverridepublicListfloat[]embed(ListStringtexts){returntexts.stream().map(embeddingModel::embed).toList();}}生产环境应批量调用并限制批大小并发超时重试Token单任务成本。八、Qdrant Payload设计{tenant_id:T001,logical_document_id:TRAVEL-POLICY,version_id:V4,source_version:4,status:READY,section_path:4.2,content_hash:...,embedding_model:...,chunking_version:v2}生产检索只允许status EFFECTIVE九、增量索引核心服务ServicepublicclassIncrementalIndexService{privatefinalDocumentRepositorydocuments;privatefinalDocumentParserparser;privatefinalEmbeddingGatewayembeddings;privatefinalChunkVectorRepositoryvectors;privatefinalVersionSwitchServiceswitchService;TransactionalpublicIndexResultindex(IndexCommandcommand){DocumentVersionversiondocuments.createVersion(command);try{documents.updateStatus(version.versionId(),VersionStatus.PARSING);ListKnowledgeChunknewChunksparser.parse(command.file());ListKnowledgeChunkoldChunksdocuments.findEffectiveChunks(command.logicalDocumentId());ChunkDiffdiffdiff(oldChunks,newChunks);documents.updateStatus(version.versionId(),VersionStatus.INDEXING);indexChangedChunks(version,diff);verify(version,newChunks.size());documents.updateStatus(version.versionId(),VersionStatus.READY);switchService.activate(version);returnIndexResult.success(version.versionId(),diff);}catch(RuntimeExceptionexception){documents.markFailed(version.versionId(),exception.getMessage());throwexception;}}}真正项目中不要将长时间Embedding放在单个数据库事务内。上面主要展示流程生产实现应使用任务状态和短事务。十、批量写入QdrantprivatevoidindexChangedChunks(DocumentVersionversion,ChunkDiffdiff){ListKnowledgeChunkchangedStream.concat(diff.added().stream(),diff.modified().stream()).toList();for(ListKnowledgeChunkbatch:BatchUtils.partition(changed,64)){Listfloat[]vectorBatchembeddings.embed(batch.stream().map(KnowledgeChunk::content).toList());vectors.upsertReady(version,batch,vectorBatch);}}写入状态先使用READY不要直接设置EFFECTIVE。十一、版本切换ServicepublicclassVersionSwitchService{Transactionalpublicvoidactivate(DocumentVersionversion){repository.expireCurrentVersion(version.logicalDocumentId());repository.activateVersion(version.versionId());vectorRepository.updateStatus(version.logicalDocumentId(),version.versionId(),EFFECTIVE);cacheVersionRepository.increment(version.logicalDocumentId());}}理想情况下数据库状态和向量Payload更新需要设计补偿机制因为它们不能参与同一个本地事务。可以采用状态机 Outbox 幂等补偿任务十二、删除旧版本版本切换后先逻辑失效旧版本 → EXPIRED再异步清理publicvoidcleanupExpired(StringlogicalDocumentId,Durationretention){ListStringexpiredVersionsrepository.findExpiredBefore(logicalDocumentId,Instant.now().minus(retention));for(StringversionId:expiredVersions){vectorRepository.deleteVersion(versionId);}}十三、失败恢复索引任务记录task_id version_id stage batch_no retry_count error_message updated_at恢复时从失败批次继续而不是重跑全部文档。重复写入依赖稳定Point ID所以Upsert是幂等的。十四、完整性校验privatevoidverify(DocumentVersionversion,intexpectedCount){longactualvectors.countByVersion(version.versionId());if(actual!expectedCount){thrownewIllegalStateException(Chunk数量不一致expectedexpectedCount, actualactual);}}还可以抽样验证文本Hash向量维度Payload字段检索结果引用映射。十五、查询过滤must: tenant_id 当前租户 status EFFECTIVE不要只按logical_document_id查询也不要让前端直接传任意tenantId。租户信息应来自认证上下文。十六、接口设计RestControllerRequestMapping(/api/knowledge-index)publicclassKnowledgeIndexController{privatefinalIndexTaskServicetasks;PostMappingpublicIndexTaskResponsesubmit(RequestBodyIndexRequestrequest){returntasks.submit(request);}GetMapping(/{taskId})publicIndexTaskStatusstatus(PathVariableStringtaskId){returntasks.status(taskId);}}大文档使用异步任务不要让HTTP连接等待整个Embedding过程。十七、监控指标index_task_success_rate index_task_duration changed_chunk_count embedding_reuse_rate embedding_batch_failure_count index_version_switch_failure_count ready_version_age stale_effective_version_count source_to_searchable_latency十八、生产环境继续补齐分布式锁同一文档并发更新任务队列断点续跑限流成本预算回归评测删除传播多Collection切换Embedding模型迁移。总结可恢复的RAG增量索引服务需要同时具备稳定ID 文件与Chunk Hash 差异计算 批量Embedding READY/EFFECTIVE状态 原子版本切换 幂等重试 完整性验证只做“新文件重新Embedding并Upsert”还不足以支撑企业知识库长期运行。