这一步我们是把前面分段后的一个个segment保存在向量数据库中,那么我们就在KnowledgeDocumentService中定义一个embeddingAndStore:
@Override
public boolean embeddingAndStore(Long docId) {
//todo 增加分布式锁
KnowledgeDocument knowledgeDocument = getById(docId);
if (knowledgeDocument == null) {
return false;
}
if (knowledgeDocument.getStatus() == DocumentStatus.VECTOR_STORED) {
return true;
}
//todo 状态的校验
//分页扫描全部document_id为docId且status为INIT的文档片段
LambdaQueryWrapper<KnowledgeSegment> queryWrapper = Wrappers.<KnowledgeSegment>lambdaQuery()
.eq(KnowledgeSegment::getDocumentId, docId)
.eq(KnowledgeSegment::getStatus, SegmentStatus.INIT)
.isNull(KnowledgeSegment::getEmbeddingId)
.eq(KnowledgeSegment::getSkipEmbedding, 0);
Page<KnowledgeSegment> page = knowledgeSegmentService.page(new Page<>(1, 100), queryWrapper);
while (page.getCurrent() == 1 || page.hasNext()) {
List<KnowledgeSegment> textSegmentsToEmbed = page.getRecords();
List<TextSegment> textSegments = textSegmentsToEmbed.stream()
.map(segment -> TextSegment.from(segment.getText(), Metadata.from(segment.getMetadataMap())))
.toList();
// 获取嵌入向量
Response<List<Embedding>> embeddingResponse = openAiEmbeddingModel.embedAll(textSegments);
// 存储嵌入向量
List<String> embeddingIds = elasticsearchEmbeddingStore.addAll(embeddingResponse.content(), textSegments);
// 更新文档片段状态
for (int i = 0; i < textSegmentsToEmbed.size(); i++) {
String embeddingId = embeddingIds.get(i);
KnowledgeSegment knowledgeSegment = textSegmentsToEmbed.get(i);
knowledgeSegment.setEmbeddingId(embeddingId);
knowledgeSegment.setStatus(SegmentStatus.VECTOR_STORED);
knowledgeSegmentService.updateById(knowledgeSegment);
}
// 继续扫描下一页
page = knowledgeSegmentService.page(new Page<>(page.getCurrent() + 1, 100), queryWrapper);
}
//todo 需要对所有的segment做检查,确保所有的segment都已转换为vector
// 更新文档状态
knowledgeDocument.setStatus(DocumentStatus.VECTOR_STORED);
return updateById(knowledgeDocument);
}
这里面的主要步骤: 这段代码实现了一个将文档切片并进行向量化存储的核心流程。我们可以把它拆解为以下几个关键步骤: 前置校验与状态检查 - 获取文档:根据 docId 查询 KnowledgeDocument 实体。 - 空值与状态判断:如果文档不存在直接返回 false;如果文档已经是向量化完成状态 (VECTOR_STORED),则直接返回 true(避免重复处理)。 分页扫描待处理切片 - 构建查询条件:使用 LambdaQueryWrapper 筛选属于该文档 (docId) 且满足以下条件的切片 (KnowledgeSegment): - 状态为初始化 (INIT)。 - 没有关联的向量ID (embeddingId 为空)。 - 未被标记为跳过向量化 (skipEmbedding == 0)。 - 分页处理:使用分页查询(每页 100 条)来防止一次性加载过多数据导致内存溢出。 向量化循环处理 代码进入一个 while 循环,直到处理完所有符合条件的切片页: - 数据转换:将数据库实体 KnowledgeSegment 转换为向量化模型所需的 TextSegment 对象(包含文本内容和元数据)。 - 调用模型生成向量:调用 openAiEmbeddingModel.embedAll 批量获取文本的嵌入向量。 - 存储向量:调用 elasticsearchEmbeddingStore.addAll 将生成的向量存储到 Elasticsearch 中,并获取返回的 embeddingId。 - 更新切片状态:遍历处理过的切片,将获取到的 embeddingId 回填,并将状态更新为 VECTOR_STORED,然后更新数据库。 - 翻页:继续查询下一页数据,直到处理完毕。 最终状态更新 - 更新文档状态:将主文档 KnowledgeDocument 的状态更新为 VECTOR_STORED,标志着整个文档的向量化流程结束。