--- title: "09-落库向量化收尾" created: 2026-05-20 aliases: - 落库向量化收尾 tags: - 项目 --- # 落库向量化收尾 我们接着上一篇“四种切块策略详解”继续往下走。 上一篇的终点是: ```text executePipeline 执行完毕 四种策略各自产出候选 chunk 所有候选都还在内存里(ParentBlockCandidate 列表) 还没有任何数据落库,没有任何向量被写入 ``` 这一篇就是把这批候选结果**真正变成可用的索引**: ```text 阶段三:切块后处理(候选 → 数据库实体) 阶段四:向量化(chunk → 向量,双索引并行) 阶段五:收尾(状态机最终推进) 异常处理:统一失败兜底 ``` 学完这一篇,整条文档处理流水线**就彻底打通了**:从用户上传到 RAG 可用,每一步的代码细节都清楚了。 下一篇可以进入更上层的话题——检索链路(用户提问到答案返回)的设计。 --- ### **一、整体认知:这一节在做什么** 可以一句话概括: > 切块产出候选结果后,先做最后一轮过滤淘汰无效父块,再用 buildParentChildEntities 给每个父子块分配全局 ID + 全局递增 chunkNo + token 估算等衍生字段把内存对象转成数据库实体,先写父块再写子块保证外键引用完整;然后调用 vectorGateway.vectorize 按 10 个一批调 EmbeddingModel 生成向量,通过 PGVector 的 INSERT ON CONFLICT 批量 upsert,并在 chunk 实体上原地回填 vectorStatus / vectorId 等状态;紧接着可选地把同一批 chunk 通过 ElasticsearchKeywordSearchGateway 写入 ES 关键词索引,文档元数据冗余进每条记录避免回表,标签字符串拆成数组方便 multi\_match;最后把 plan 标 EXECUTED、document 标 BUILD\_SUCCESS、task 标 SUCCESS,所有状态机推进到终态。任何环节失败由外层 try-catch 统一接管,把 chunk 中间态 VECTOR\_FAILED 化、step 标 EXECUTE\_FAILED、document 标 BUILD\_FAILED、task 标 FAILED,避免库里残留半成品状态。 #### **整体流程图** ```mermaid flowchart TD A[buildParentBlocks 返回候选] --> B[阶段三:后处理] B --> C[过滤无效父块] C --> D[buildParentChildEntities] D --> E[父块 insert] E --> F[chunk insert] F --> G[阶段四:向量化] G --> H[vectorGateway.vectorize] H --> I[EmbeddingModel 批量 embed] I --> J[PGVector batchUpsert] J --> K[markSuccess 回填状态] K --> L{keywordSearchGateway 可用?} L -->|是| M[indexChunks 写 ES] L -->|否| N[跳过] M --> O[chunk updateById 持久化向量状态] N --> O O --> P[阶段五:收尾] P --> Q[plan EXECUTED] Q --> R[document BUILD_SUCCESS] R --> S[task SUCCESS] G -.异常.-> X[统一失败处理] H -.异常.-> X M -.异常.-> X X --> Y[document BUILD_FAILED] Y --> Z[chunk VECTOR_FAILED] Z --> AA[step EXECUTE_FAILED] AA --> BB[task FAILED] ``` --- ### **二、阶段三:切块后处理** #### **1. 最后一轮过滤** ```java List finalParentBlockList = parentBlockCandidateList.stream() .filter(item -> item != null && StrUtil.isNotBlank(item.getText()) && item.getChildChunks() != null && item.getChildChunks().stream() .anyMatch(child -> StrUtil.isNotBlank(child.getText()))) .toList(); ``` ##### **三个过滤条件** ```text 1. item 非空且 text 非空白:父块本身要有内容 2. childChunks 不为 null:子块列表存在 3. anyMatch 至少一个有效子块:不允许全空 ``` ##### **为什么 buildParentBlocks 已经清洗过还要再过滤?** 虽然上一篇 buildParentBlocks 内部已经做了三轮 cleanupChunkList,但这里再过滤一次的原因: ```text 1. cleanupChunkList 是"层内清洗"——清理 chunk 列表 2. 这里是"层间清洗"——验证父子关系完整性 3. anyMatch 检查特别重要:即使父块有内容,子块全空也不能落库 ``` ##### **anyMatch vs allMatch 的选择** ```text .anyMatch(child -> StrUtil.isNotBlank(child.getText())) ``` ```text anyMatch:至少一个子块有内容就 OK allMatch:所有子块都必须有内容 ``` 为什么选 anyMatch? ```text 子块列表里偶尔有空块是正常的(被清洗后混进来) 只要有一个有效子块就足以支撑检索 allMatch 太严格,会误杀很多有效父块 ``` 但配合后面的 `if (StrUtil.isBlank(childCandidate.getText())) continue;`,最终落库的 chunk 一定都是有效的——**先放过父块,再在循环里过滤无效子块**。这种"宽进严出"的设计避免了过度严格导致的数据丢失。 #### **2. buildParentChildEntities:候选 → 实体** 这个方法是阶段三的核心,做四件事: ```text 1. 给父块和子块分配全局唯一 ID 2. 建立父子关系 3. 生成全局递增的 chunkNo 4. 计算字符数、token 估算等衍生字段 ``` ##### **关键设计 1:globalChunkNo 全局递增** ```java int globalChunkNo = 1; for (... parentIndex < parentBlockCandidateList.size(); parentIndex++) { ... for (ChunkCandidate childCandidate : parentCandidate.getChildChunks()) { ... chunk.setChunkNo(globalChunkNo++); } } ``` 设想两种编号方式: ```text 方式 A:每个父块内从 1 开始 父块 1 的子块:1, 2, 3 父块 2 的子块:1, 2 父块 3 的子块:1, 2, 3, 4 方式 B:整篇文档全局递增(当前实现) 父块 1 的子块:1, 2, 3 父块 2 的子块:4, 5 父块 3 的子块:6, 7, 8, 9 ``` 为什么选方式 B? ```text 1. 检索结果按 chunkNo 排序就是原文顺序 2. 用户看引用时显示 "chunk #5",一目了然在文档中的位置 3. 调试和定位时不需要 (parentNo, chunkNo) 二元组 4. 全局唯一,可以独立做主键 ``` 方式 A 在概念上更"对齐父子关系",但工程上方式 B 更好用。 ##### **关键设计 2:startChunkNo / endChunkNo 区间回填** ```java int startChunkNo = globalChunkNo; int childCount = 0; for (... childCandidate : parentCandidate.getChildChunks()) { ... chunk.setChunkNo(globalChunkNo++); childCount++; } parentBlock.setChildCount(childCount); parentBlock.setStartChunkNo(childCount == 0 ? null : startChunkNo); parentBlock.setEndChunkNo(childCount == 0 ? null : globalChunkNo - 1); ``` 这是个巧妙的设计——**父块上冗余存储了它管辖的 chunk 编号区间**。 为什么需要这个? ```sql 父块详情页要展示"这个父块包含 chunk 5-9" 不冗余 → 每次都要 SELECT MIN/MAX FROM chunk WHERE parent_block_id = ? 冗余 → 直接读 startChunkNo / endChunkNo ``` 类似上一篇 strategySnapshot 的反范式优化,**用冗余字段换查询性能**。 ##### **关键设计 3:childCount == 0 时填 null** ```java parentBlock.setStartChunkNo(childCount == 0 ? null : startChunkNo); ``` 为什么要区分 0 和 非0? ```sql 如果直接填 0 或者 globalChunkNo "0" 看起来像有效的 chunk 编号 可能误导查询(SELECT WHERE chunkNo BETWEEN 0 AND 0) 填 null 明确表达"这个父块没有任何 chunk" SQL 查询条件可以用 IS NOT NULL 排除 ``` ##### **关键设计 4:子块继承父块 sectionPath** ```text chunk.setSectionPath(StrUtil.blankToDefault( childCandidate.getSectionPath(), parentCandidate.getSectionPath() )); ``` ```text 子块有自己的 sectionPath → 用子块的(更精细) 子块没有 sectionPath → 继承父块的(至少有章节信息) ``` 什么场景下子块没有 sectionPath? ```text 父块走结构切块产生(带 sectionPath="第一章 > 1.1节") 子块走递归切块产生(只是按字符切,没有路径) 此时继承父块路径,保证检索时仍能展示"这段来自第一章 > 1.1节" ``` 这是**信息向下传递**的常见手法,让子块继承父块的上下文。 ##### **关键设计 5:vectorStatus 初始化为 WAIT\_VECTOR** ```java chunk.setVectorStatus(DocumentVectorStatusEnum.WAIT_VECTOR.getCode()); chunk.setVectorStoreType(DocumentVectorStoreTypeEnum.PG_VECTOR.getCode()); ``` ```text WAIT_VECTOR:等待向量化 后续 vectorize 会切到 VECTORIZING → VECTOR_SUCCESS 失败时切到 VECTOR_FAILED ``` 为什么 chunk 落库时已经预设了 vectorStoreType? ```text 即使向量化还没开始,目标存储已经确定(PG_VECTOR) 后续 vectorize 失败时,这个字段还在,运维能立刻知道"应该写到哪" 不会出现"chunk 在但不知道目标向量库"的歧义 ``` #### **3. estimateTokenCount:轻量 token 估算** ```java int chineseCount = 0; for (char current : text.toCharArray()) { if (String.valueOf(current).matches("[\\u4e00-\\u9fa5]")) { chineseCount++; } } int englishCount = 0; for (String word : text.split("\\s+")) { if (word.matches(".*[A-Za-z].*")) { englishCount++; } } return chineseCount + englishCount + Math.max(1, (text.length() - chineseCount) / 4); ``` ##### **为什么不用真正的 tokenizer?** OpenAI / Anthropic 等模型都有自己的 tokenizer(如 tiktoken),能精确算出 token 数。但这里选了**简化估算**。 理由: ```text 1. 真 tokenizer 慢:每个 chunk 算一次,大文档累积可观 2. 真 tokenizer 重:加依赖、加内存 3. 真 tokenizer 不通用:不同模型 tokenizer 不同 4. 这里 token 用途仅限"统计展示"——不参与计费、不参与限流 ``` 精确度要求不高的场景,**轻量估算就够**。 ##### **估算公式的意图** ```text 中文 1 字 ≈ 1 token (中文模型常见比例) 英文 1 词 ≈ 1 token (粗略估计) 其他字符:每 4 个字符算 1 token(标点、数字、符号等) ``` 误差在 ±20% 以内,对统计展示足够。 ##### **Math.max(1, ...) 的边界保护** ```text Math.max(1, (text.length() - chineseCount) / 4) ``` 如果文本全是中文: ```text text.length() - chineseCount = 0 0 / 4 = 0 chineseCount + englishCount + 0 ``` 加 `Math.max(1, ...)` 后: ```text 全中文文本至少加 1 保证短文本 token 不会为 0(看起来像异常数据) ``` 但这个保护其实有点过度——纯中文 1 字 1 token 时其他字符部分加 0 也合理。这是**经验主义微调**,不影响正确性。 #### **4. 父子先后入库:外键约束的暗示** ```java for (... parentBlock : parentBlockEntityList) { parentBlockMapper.insert(parentBlock); } for (... chunk : chunkEntityList) { chunkMapper.insert(chunk); } ``` 为什么父块先写、子块后写? ```text chunk.parentBlockId 引用 parent_block.id 即使数据库没建外键,业务上也要保证"被引用的先存在" ``` 这种**插入顺序约束**让数据库层面的关联在任何时刻都是闭合的: ```text 万一中途崩溃 父块 5 个写完,子块写到一半 库里:5 个父块 + N 个子块(都引用已存在的父块) 数据完整性仍然保持 ``` 如果反过来子块先写: ```text 子块写入时父块还不存在 chunk.parentBlockId 引用了"虚空" ID 此时崩溃,留下"孤儿子块" ``` ##### **为什么不用一次 batchInsert?** 代码用的是循环逐条 insert 而不是批量插入。这其实是**性能上可优化**的点: ```text 循环 insert:每条一次 SQL,N 次网络往返 batchInsert:一次 SQL,1 次网络往返 ``` 但当前实现选了简单: ```text 事务内逐条 insert,网络往返开销在事务内被摊薄 代码更直观 父块通常 < 50 个,子块通常 < 500 个,性能可接受 ``` 如果后续遇到大文档性能问题,可以改成 `parentBlockMapper.insertBatchSomeColumn(...)`。 --- ### **三、阶段四:向量化** 向量化是整条链路**最重的一步**——既要调外部模型(embedding API),又要写大量数据(向量比文本大得多)。 #### **1. 主流程:三个独立调用** ```java vectorGateway.vectorize(chunkEntityList); DocumentKeywordSearchGateway keywordSearchGateway = keywordSearchGatewayProvider.getIfAvailable(); if (keywordSearchGateway != null) { keywordSearchGateway.indexChunks(chunkEntityList); } for (SuperAgentDocumentChunk chunk : chunkEntityList) { chunkMapper.updateById(chunk); } ``` 三个调用各自独立: ```text 1. vectorize:语义召回索引(必选) 2. indexChunks:关键词召回索引(可选) 3. updateById:把内存中的状态变更持久化回业务表 ``` ##### **为什么 vectorize 是必选,indexChunks 是可选?** ```text vectorize:RAG 的核心能力,没向量就检索不了语义 indexChunks:增强能力,关键词检索是锦上添花 ``` ##### **getIfAvailable 的可插拔设计** ```java DocumentKeywordSearchGateway keywordSearchGateway = keywordSearchGatewayProvider.getIfAvailable(); ``` `ObjectProvider.getIfAvailable()` 是 Spring 的优雅注入: ```text Bean 存在 → 返回实例 Bean 不存在 → 返回 null,不抛异常 ``` 为什么用这种方式而不是 @Autowired(required=false)? ```java @Autowired(required=false):字段级,启动时绑定 ObjectProvider.getIfAvailable():运行时按需获取 ``` 后者的好处: ```text 1. Bean 配置可以动态变化(虽然 Spring 单例容器下意义不大) 2. 可选依赖的语义更明确 3. 不强依赖在字段上,代码组织更灵活 ``` 这是 Spring 中表达"**可选依赖**"的标准方式。 ##### **可选 ≠ 失败可忽略** ```java if (keywordSearchGateway != null) { keywordSearchGateway.indexChunks(chunkEntityList); } ``` 注意这里 **没有 try-catch**。意思是: ```text 关键词检索网关不存在 → 跳过(可选) 关键词检索网关存在但写入失败 → 抛异常,任务失败 ``` 这种设计的逻辑: ```text "我不提供这个能力" → 跳过没问题 "我提供这个能力但搞砸了" → 必须失败,不能让用户以为关键词索引建好了实际却没有 ``` 这是**契约一致性**的保证——一旦承诺提供能力,就要保证能力可用。 #### **2. vectorize 主方法** ```java public static final int EMBEDDING_BATCH_SIZE_LIMIT = 10; ``` ##### **为什么批大小是 10?** ```text 太小(比如 1):每个 chunk 一次 API 调用,网络开销爆炸 太大(比如 100): 单次 prompt 累积太长可能超模型限制 单次失败影响范围太大 重试代价高 ``` 10 是经验值,平衡了吞吐和失败粒度。不同 embedding 模型的最佳批大小不同,可以根据具体模型调。 ##### **过滤空 chunk** ```java List<...> validChunkList = chunkList.stream() .filter(chunk -> chunk != null && StrUtil.isNotBlank(chunk.getChunkText())) .toList(); if (validChunkList.isEmpty()) return; ``` 虽然前面 buildParentChildEntities 已经过滤过空 chunk,这里再过滤一遍: ```text 深度防御:即使上游漏过空块,这里也能拦 EmbeddingModel 输入空字符串可能抛异常或返回零向量,污染数据 ``` ##### **批次循环** ```java int totalBatchCount = (validChunkList.size() + batchSize - 1) / batchSize; for (int startIndex = 0; startIndex < validChunkList.size(); startIndex += batchSize) { int endIndex = Math.min(startIndex + batchSize, validChunkList.size()); List<...> currentBatch = validChunkList.subList(startIndex, endIndex); int currentBatchIndex = (startIndex / batchSize) + 1; log.info("...batchIndex={}/{}, ...", currentBatchIndex, totalBatchCount, ...); List embeddingList = embeddingModel.embed(...); if (embeddingList.size() != currentBatch.size()) { throw new IllegalStateException("EmbeddingModel 返回的向量数量与 chunk 数量不一致。"); } batchUpsert(...); markSuccess(currentBatch); } ``` ##### **subList 的零拷贝** ```java List<...> currentBatch = validChunkList.subList(startIndex, endIndex); ``` `subList` 返回的是**原 list 的视图**,不是新 list: ```text 零拷贝,不占额外内存 修改 subList 会反映到原 list(这里只读不修改,无问题) 但要避免在 subList 期间修改原 list,会导致 ConcurrentModificationException ``` ##### **(size + batchSize - 1) / batchSize 上取整** ```java int totalBatchCount = (validChunkList.size() + batchSize - 1) / batchSize; ``` ```text size = 25, batchSize = 10 (25 + 9) / 10 = 3 3 个批次:[1-10], [11-20], [21-25] ``` 这是上取整的标准技巧,比 `Math.ceil((double)size / batchSize)` 更高效(不涉及浮点)。 ##### **数量校验** ```java if (embeddingList.size() != currentBatch.size()) { throw new IllegalStateException("..."); } ``` 为什么这么严格? ```text EmbeddingModel 返回向量必须与输入 chunk 一一对应 如果数量不一致,后续 batchUpsert 时 chunk[i] 和 embedding[i] 错位 错位的向量会让检索完全失效(用 chunk 5 的文本搜出 chunk 7 的内容) ``` 这种**契约违反**只能立刻抛异常,绝不能继续。 ##### **批次粒度的失败隔离** 注意这里**没有按批次 try-catch**。一批失败整个 vectorize 失败: ```text 为什么不允许部分成功? chunk 5-10 向量化成功,chunk 11-20 失败 库里有"已向量化的"和"未向量化的"混合状态 后续要么补偿要么回滚,极复杂 所以:要么全部成功,要么全部失败 失败由外层 catch 接管,统一标记 VECTOR_FAILED ``` #### **3. batchUpsert:PGVector 批量写入** ```java pgVectorJdbcTemplate.batchUpdate(upsertSql, new BatchPreparedStatementSetter() { @Override public void setValues(PreparedStatement ps, int index) throws SQLException { SuperAgentDocumentChunk chunk = chunkBatch.get(index); float[] embedding = embeddingBatch.get(index); chunk.setVectorStatus(DocumentVectorStatusEnum.VECTORIZING.getCode()); String metadataJson = buildMetadataJson(chunk, embeddingModelName); ps.setLong(1, chunk.getId()); ... ps.setString(18, toVectorLiteral(embedding)); } @Override public int getBatchSize() { return chunkBatch.size(); } }); ``` ##### **为什么用 JdbcTemplate 而不是 MyBatis?** ```text PGVector 写入需要 JDBC 层处理向量字面量 MyBatis 的 ORM 模型不擅长处理 vector 类型 JdbcTemplate + 原生 SQL 更直接 ``` 这是混合 ORM/JDBC 的常见做法——**简单查询用 ORM,特殊类型用 JDBC**。 ##### **专用 jdbcTemplate** ```text pgVectorJdbcTemplate ``` 注意这里是 `pgVectorJdbcTemplate`,不是默认的 jdbcTemplate。说明: ```text PGVector 可能是独立的数据库实例 配置了独立的 DataSource 和 JdbcTemplate 业务库和向量库分离,各自独立伸缩 ``` 这种**数据库分离**架构在大型系统里很常见: ```text 业务库:OLTP 优化,SSD,主从复制 向量库:特殊扩展(pgvector),HNSW 索引,可能不同实例 ``` ##### **INSERT ON CONFLICT DO UPDATE** ```sql UPSERT_SQL_TEMPLATE // INSERT INTO ... ON CONFLICT (chunk_id) DO UPDATE SET ... ``` 这是 PostgreSQL 的 upsert 语法: ```text chunk_id 不存在 → INSERT chunk_id 已存在 → UPDATE ``` 为什么用 upsert 而不是纯 INSERT? ```text 任务可能重试 失败后部分 chunk 已经写入向量库 重试时 INSERT 会冲突 upsert 让重试天然幂等 ``` 这是**幂等设计**的经典手法——让重复执行的结果和单次执行一致。 ##### **状态切到 VECTORIZING 的时机** ```java chunk.setVectorStatus(DocumentVectorStatusEnum.VECTORIZING.getCode()); ``` 注意这是在 `setValues` 里改内存状态,而不是落库前改。逻辑: ```text PreparedStatement 设值时,内存对象从 WAIT_VECTOR 切到 VECTORIZING batchUpdate 执行后,markSuccess 再把内存对象切到 VECTOR_SUCCESS 最终 chunkMapper.updateById 把 VECTOR_SUCCESS 持久化到 chunk 表 ``` 中间态 VECTORIZING 在内存里短暂存在,但不一定持久化。这是**内存状态转换**和**库状态转换**的差异。 ##### **toVectorLiteral:向量字符串字面量** ```java ps.setString(18, toVectorLiteral(embedding)); ``` PGVector 接收的向量格式是 `'[0.1, 0.2, 0.3]'` 字符串: ```text float[] {0.1f, 0.2f, 0.3f} → "[0.1,0.2,0.3]" → ps.setString(...) PostgreSQL 解析时再转回向量类型 ``` 这种**字符串字面量**接口是 PGVector JDBC 层的特殊设计,因为标准 JDBC 没有 vector 类型。 #### **4. markSuccess:回填状态** ```java private void markSuccess(List<...> chunkBatch) { for (... chunk : chunkBatch) { chunk.setVectorId(String.valueOf(chunk.getId())); chunk.setVectorStoreType(DocumentVectorStoreTypeEnum.PG_VECTOR.getCode()); chunk.setVectorStatus(DocumentVectorStatusEnum.VECTOR_SUCCESS.getCode()); } } ``` ##### **vectorId = chunk.id 的设计** ```java chunk.setVectorId(String.valueOf(chunk.getId())); ``` 为什么 vectorId 直接用 chunk 主键? ```text 1. 简化:不需要单独维护向量库的 ID 空间 2. 一致:业务库和向量库的 ID 对应关系一目了然 3. 幂等:重写时 vectorId 不变,upsert 自然覆盖 ``` 如果用独立 vectorId(比如 UUID): ```text 要存"业务 chunk_id ↔ 向量库 vector_id"的映射 查找复杂度增加 看起来更"专业"但徒增复杂度 ``` 简单设计胜出。 ##### **为什么是 String 类型?** ```java chunk.setVectorId(String.valueOf(chunk.getId())); ``` 明明 chunk.id 是 Long,为什么 vectorId 用 String? ```text 不同向量库 ID 类型不同 PGVector:数字主键 Milvus:可能 String UUID Pinecone:String Weaviate:UUID ``` vectorId 用 String 是**通用接口**: ```text 如果以后切换到其他向量库,ID 类型变了不需要改 schema 牺牲一点存储(数字转字符串多几个字节),换接口稳定 ``` 这是为**多向量库支持**留的扩展位。 #### **5. buildMetadataJson:metadata 设计** ```java Map metadata = new LinkedHashMap<>(); metadata.put("documentId", chunk.getDocumentId()); metadata.put("taskId", chunk.getTaskId()); metadata.put("planId", chunk.getPlanId()); metadata.put("parentBlockId", chunk.getParentBlockId()); metadata.put("chunkNo", chunk.getChunkNo()); metadata.put("sourceType", chunk.getSourceType()); metadata.put("sectionPath", chunk.getSectionPath()); metadata.put("charCount", chunk.getCharCount()); metadata.put("tokenCount", chunk.getTokenCount()); metadata.put("embeddingModel", embeddingModelName); return objectMapper.writeValueAsString(metadata); ``` ##### **metadata 的两个用途** ```text 1. 检索过滤:WHERE metadata->>'documentId' = '123' 限定只搜某个文档的 chunk 限定只搜某个章节路径下的 chunk 2. 调试追溯:看到一条向量记录就知道 属于哪个文档、哪个任务、哪个父块 用什么 embedding 模型生成的 ``` ##### **为什么字段都这么"基础"?** metadata 里都是基础维度,没有业务字段(比如文档名、标签)。原因: ```text 1. metadata 跟 chunk 一比一关系,信息冗余多了 PGVector 表会膨胀 2. 业务字段(文档名、标签)更适合放 ES 关键词索引 3. metadata 只保留"过滤检索时常用的轻量字段" ``` 这是**职责分离**: ```text PGVector metadata:语义检索过滤 ES 索引文档:关键词检索字段 两者各有侧重,不互相覆盖 ``` ##### **embeddingModel 的版本追溯** ```java metadata.put("embeddingModel", embeddingModelName); ``` 为什么要存 embedding 模型名? ```text embedding 模型升级时(比如 text-embedding-ada-002 → text-embedding-3-small) 新生成的向量维度可能不同 新旧向量混存检索质量会出问题 通过 metadata 能快速找出"用旧模型生成的需要重建" ``` 这是**向量库版本管理**的关键字段。 ##### **LinkedHashMap 保序** ```text new LinkedHashMap() ``` JSON key 顺序对功能没影响,但对**调试和对比**有用: ```text Linked HashMap:每次输出 key 顺序一致 HashMap:key 顺序随机 看 metadata 时知道在哪个位置看哪个字段 diff 两条记录时不会因 key 顺序不同造成假差异 ``` ##### **JsonProcessingException 包装** ```java catch (JsonProcessingException exception) { throw new IllegalStateException("序列化 PGVector metadata 失败。", exception); } ``` JSON 序列化几乎不会失败(key 都是 String,value 都是基础类型)。这里 catch 是为了**契约清晰**: ```text 方法签名不抛 checked exception(JsonProcessingException 是 checked) 转成 IllegalStateException(unchecked) 让调用方不需要处理 "内部状态错误" 比 "JSON 处理异常" 更接近业务语义 ``` #### **6. indexChunks:Elasticsearch 关键词索引** ##### **整体流程** ```java public void indexChunks(List chunkList) { if (CollUtil.isEmpty(chunkList)) return; Map documentMap = loadDocumentMap(chunkList); BulkRequest.Builder bulkBuilder = new BulkRequest.Builder() .index(properties.getElasticsearch().getIndexName()) .refresh(Refresh.WaitFor); for (... chunk : chunkList) { SuperAgentDocument document = documentMap.get(chunk.getDocumentId()); DocumentKeywordIndexRecord indexRecord = toIndexRecord(chunk, document); bulkBuilder.operations(operation -> operation .index(index -> index .id(indexRecord.getChunkId()) .document(indexRecord))); } BulkResponse response = elasticsearchClient.bulk(bulkBuilder.build()); if (response.errors()) { String errorMessage = response.items().stream()... throw new IllegalStateException("..."); } } ``` ##### **关键词索引和向量索引的对比** ```text 向量索引(PGVector): 解决"语义相似"召回 用户搜 "怎么报销" 也能搜到 "费用报销流程" 缺点:对精确字面要求差 关键词索引(ES): 解决"字面匹配"召回 用户搜 "FORM-2024-001" 这种唯一编号 用户按"分类=财务" 过滤 缺点:对语义变化敏感 ``` 两者形成**互补检索**,后续检索链路会做混合召回。 ##### **Refresh.WaitFor 的权衡** ```java .refresh(Refresh.WaitFor); ``` ES 的 refresh 策略: ```text NONE:写入后异步刷新(可能 1 秒后才可见) WaitFor:写入后等待下次 refresh 完成 Immediate:写入后立即 refresh(强制刷新) ``` 为什么选 WaitFor 而不是 NONE? ```text NONE: 速度最快 但任务返回后用户立刻搜索可能搜不到刚建的索引 WaitFor: 略慢(等下次 refresh) 但保证返回后立即可搜 适合"建完索引立刻试用"的场景 ``` 为什么不选 Immediate? ```text Immediate 强制刷新对 ES 集群压力大 高并发下会拖垮集群 WaitFor 是更友好的折中 ``` ##### **chunk\_id 作为 ES 文档 ID 的幂等设计** ```text .id(indexRecord.getChunkId()) ``` 和 PGVector 的 vectorId = chunk.id 是同一思路: ```text 重写时 ES 自动覆盖原文档 任务重试天然幂等 不会留下重复索引 ``` ##### **bulk API 的细节:errors() 检查** ```java BulkResponse response = elasticsearchClient.bulk(bulkBuilder.build()); if (response.errors()) { String errorMessage = response.items().stream() .filter(item -> item.error() != null) .map(item -> item.id() + ":" + item.error().reason()) .collect(Collectors.joining("; ")); throw new IllegalStateException("批量写入 Elasticsearch 失败: " + errorMessage); } ``` 这是 ES Bulk API 的**陷阱**: ```text HTTP 200 OK 不代表所有 item 成功 response 里每个 item 有自己的成功/失败状态 必须显式检查 response.errors() ``` 很多开发者会漏掉这个检查,导致部分写入失败但任务标记成功。这套代码做了完整检查: ```text 1. response.errors() 判断是否有任何 item 失败 2. stream 过滤出失败 item 3. 拼接 "id:reason" 让错误信息可定位 4. 抛 IllegalStateException 让外层接管 ``` ##### **错误信息格式 "id:reason"** ```text .map(item -> item.id() + ":" + item.error().reason()) ``` 为什么这么拼? ```text id:具体哪个 chunk 失败 reason:为什么失败 分号分隔:多个失败合并展示 最终信息:"chunk_5:mapping conflict; chunk_8:disk full" 运维一眼就能定位问题 ``` ##### **IOException 包装** ```java catch (IOException exception) { throw new IllegalStateException("写入 Elasticsearch 失败", exception); } ``` ES 客户端的方法签名抛 IOException(checked)。这里包装成 IllegalStateException: ```text 让 indexChunks 方法签名不带 checked exception 调用方不需要 try-catch IOException IllegalStateException 在外层 handleIndexBuild 的 Exception catch 里被统一接管 ``` #### **7. loadDocumentMap:N+1 优化** ```java List documentIds = chunkList.stream() .map(SuperAgentDocumentChunk::getDocumentId) .filter(Objects::nonNull) .distinct() .toList(); if (documentIds.isEmpty()) return Map.of(); List<...> documents = documentMapper.selectBatchIds(documentIds); Map documentMap = new LinkedHashMap<>(); for (... document : documents) { documentMap.put(document.getId(), document); } return documentMap; ``` ##### **N+1 问题的典型场景** ```text 错误写法: for (chunk in chunks): document = documentMapper.selectById(chunk.documentId) ... 1 次查询 + N 次查询 = N+1 正确写法: documentIds = chunks.map(documentId).distinct() documents = documentMapper.selectBatchIds(documentIds) 1 次查询 ``` 这套代码用 distinct 去重的细节也很重要: ```text 500 个 chunk 可能都属于同一个 document 不去重:selectBatchIds 传 500 个相同的 ID 去重:只传 1 个 ID 减少 SQL 参数和数据库压力 ``` ##### **Map 而非 List 的访问优化** ```text Map documentMap ``` 把 List 转成 Map 后,主循环可以 O(1) 访问: ```text List 查找:for + equals 比较,O(N) Map 查找:hash 直接索引,O(1) 500 个 chunk × 50 个 document List:500 * 50 = 25000 次比较 Map:500 次 hash = 500 次比较 ``` ##### **空文档兜底** ```java SuperAgentDocument document = documentMap.get(chunk.getDocumentId()); // 后续 toIndexRecord 里: .documentName(document == null ? "" : safeText(document.getDocumentName())) ``` 如果 chunk 关联的 document 已被删除或查询不到: ```text documentMap.get(...) 返回 null toIndexRecord 用空字符串兜底 chunk 仍然能写入 ES,只是文档维度字段为空 ``` 这种**容错降级**让 chunk 不会因为 document 缺失而整体丢失。 #### **8. toIndexRecord:领域对象映射** ##### **字段来源的两类** ```text // 来自 chunk 的字段 .chunkId(String.valueOf(chunk.getId())) .documentId(chunk.getDocumentId()) ... .chunkText(safeText(chunk.getChunkText())) // 来自 document 的冗余字段 .documentName(document == null ? "" : safeText(document.getDocumentName())) .knowledgeScopeCode(...) .knowledgeScopeName(...) .businessCategory(...) .documentTags(splitTags(...)) ``` ##### **为什么把 document 字段冗余进每条 chunk 索引?** 考虑检索场景: ```text 用户搜 "财务报销" 要求 1:文本匹配 "财务报销"(在 chunk_text) 要求 2:只搜 "财务" 分类的文档 如果不冗余 documentMap: ES 查 "财务报销" 命中 chunk_5 回数据库查 chunk_5 的 document.businessCategory 判断是否 "财务" 不是 → 过滤掉 如果冗余: ES 查 "财务报销" AND businessCategory="财务" 一次搜索完成 ``` 冗余的代价(存储)远小于回报(检索性能 + 简化)。 ##### **chunkId 用 String** ```text .chunkId(String.valueOf(chunk.getId())) ``` ES 文档 ID 是 String 类型。明明 Long 也能存,为什么转 String? ```text ES API 接口 id 字段就是 String Long → String 转换在写入时只做一次 检索回来还是 String,直接用,不需要再转 ``` ##### **chunkId 作为索引文档 ID 而非 chunk\_id 字段** 注意这里 `.id(indexRecord.getChunkId())` 把 chunkId 作为 ES 文档的 \_id(主键),而 record 里又有 chunk\_id 字段: ```text ES _id:文档主键,用于幂等覆盖 record.chunk_id:可检索字段,用于 term 查询 两者数据相同但语义不同 _id 在 ES 元数据,chunk_id 在文档体内 ``` 这种"**主键 + 字段双存**"在 ES 设计中很常见,因为 \_id 不能直接 term 查询,必须有同步的字段。 #### **9. splitTags:字符串数组化** ```text return java.util.Arrays.stream(documentTags.split(",")) .map(String::trim) .filter(StrUtil::isNotBlank) .distinct() .toList(); ``` ##### **三步清洗** ```text trim:去掉用户录入或数据库存储产生的空格 filter isNotBlank:过滤连续逗号、尾部逗号产生的空标签 distinct:去重 ``` ##### **为什么 ES 用数组而不是字符串?** ```text 数据库存:"报销,财务,发票"(字符串) ES 存:["报销", "财务", "发票"](数组) ``` 数组在 ES 中的优势: ```text 1. multi_match 可以对单个标签命中 搜 "财务" → 直接命中 "财务" 数组项 字符串方案需要用 LIKE %财务% 模糊匹配,效率低 2. terms 聚合能统计标签分布 "总共有多少 chunk 带 '财务' 标签?" 一个聚合查询搞定 3. ES 倒排索引天然适合数组 ``` #### **10. safeText:null → 空字符串** ```java private String safeText(String text) { return text == null ? "" : text; } ``` ##### **为什么不 trim?** ```java // 不是 return text == null ? "" : text.trim(); return text == null ? "" : text; ``` ```text 1. 这些字段(章节路径、文档名、正文)有自己的语义 2. trim 会改变用户原文 3. 索引层不应该擅自改写文本 ``` 只把 null 变成空字符串,**不改写内容**——这是**最小干预**原则。 ##### **为什么需要这个方法?** ```text ES 对 null 字段的处理: 存为 null 检索 DSL 要写 exists 判断 统一用空字符串: 检索 DSL 简单(不用 exists) 字段类型一致(都是 string) 避免 NPE 风险 ``` 工程上**统一用 "无值字符串"代替 null**,让下游代码假设永远不会 null,简化逻辑。 --- ### **四、阶段五:收尾** ```java plan.setPlanStatus(DocumentPlanStatusEnum.EXECUTED.getCode()); planMapper.updateById(plan); document.setIndexStatus(DocumentIndexStatusEnum.BUILD_SUCCESS.getCode()); document.setLastIndexTaskId(taskId); documentMapper.updateById(document); finishTaskSuccess(task, DocumentTaskStageEnum.STORE_COMPLETE.getCode(), startTime); ``` #### **三层状态机的最终态** ```yaml plan: WAIT_CONFIRM → CONFIRMED → EXECUTED document: WAIT_BUILD → BUILDING → BUILD_SUCCESS task: NEW → RUNNING → SUCCESS step: WAIT_EXECUTE → EXECUTING → EXECUTE_SUCCESS ``` 四条状态机各自推进到终态。每条状态机独立追踪不同维度: ```text plan:这个方案被执行过了吗 document:这份文档可用于 RAG 了吗 task:这次任务跑完了吗 step:每个策略步骤跑完了吗 ``` #### **lastIndexTaskId 的作用** ```java document.setLastIndexTaskId(taskId); ``` 记录"最近一次成功索引的任务 ID": ```text 重建索引时:对比 lastIndexTaskId,知道是否需要清理旧索引 排查问题时:从 document 直接跳到任务详情 监控统计时:计算"索引时长" "索引频率" ``` 这种**反向引用**让 document 表自带溯源能力。 #### **finishTaskSuccess 的复用** ```java private void finishTaskSuccess(SuperAgentDocumentTask task, Integer stage, Date startTime) { Date finishTime = new Date(); task.setTaskStatus(SUCCESS); task.setCurrentStage(stage); task.setFinishTime(finishTime); task.setCostMillis(finishTime.getTime() - startTime.getTime()); task.setErrorCode(null); task.setErrorMsg(null); taskMapper.updateById(task); } ``` 这个方法在解析任务和索引任务都用,逻辑完全一样: ```text 切到 SUCCESS 记录完成时间 计算耗时 清空错误字段(防止上次失败的 errorMsg 残留) ``` **清空 errorMsg / errorCode** 是关键细节: ```text 任务失败 → 写入 errorMsg 任务后来重试成功 → 不清空就会留着错误信息 查任务表会看到"成功的任务但有错误信息" → 误导 ``` 这是**防御性收尾**。 #### **STORE\_COMPLETE 阶段命名** ```java finishTaskSuccess(task, DocumentTaskStageEnum.STORE_COMPLETE.getCode(), startTime); ``` 阶段名是 `STORE_COMPLETE`(存储完成),不是 `INDEX_COMPLETE` 或 `VECTORIZE_COMPLETE`。 这个命名暗示了任务的**终点语义**: ```text 不只是向量化完成 而是所有存储(业务库 + 向量库 + ES)都已写入 站在"存储层"视角的统一收尾 ``` 阶段命名映射的是**业务语义**而不是**实现细节**,让后续即使实现变化(比如加了 Neo4j 图谱)也不需要改阶段名。 --- ### **五、统一错误处理** ```typescript catch (Exception exception) { log.error("异步构建索引失败,documentId={}, taskId={}, planId={}", documentId, taskId, planId, exception); document.setIndexStatus(DocumentIndexStatusEnum.BUILD_FAILED.getCode()); documentMapper.updateById(document); chunkMapper.update(null, new LambdaUpdateWrapper() .eq(SuperAgentDocumentChunk::getTaskId, taskId) .eq(SuperAgentDocumentChunk::getStatus, BusinessStatus.YES.getCode()) .set(SuperAgentDocumentChunk::getVectorStatus, VECTOR_FAILED) .set(SuperAgentDocumentChunk::getVectorStoreType, PG_VECTOR)); updateStepExecuteStatus(planId, EXECUTE_FAILED); failTask(task, startTime, exception, task.getCurrentStage()); } ``` #### **1. 统一 catch 的设计哲学** 整个 `handleIndexBuild()` 方法被一个大 try-catch 包裹,任意阶段失败都进这里: ```text 切块阶段失败 → catch 落库阶段失败 → catch 向量化失败 → catch ES 写入失败 → catch 状态更新失败 → catch ``` 为什么不分阶段 try-catch? ```text 分阶段 catch: 每个阶段有自己的 catch 逻辑 代码冗余,容易遗漏 部分阶段失败时状态机推进逻辑复杂 统一 catch: 一处兜底,逻辑集中 任意阶段失败都按同一套规则处理 代码简洁 ``` 这是**异常处理集中化**的常见做法。代价是失败时无法知道"具体哪个细节阶段挂了",但 task.currentStage 已经记录了进入 catch 时的阶段,足够定位。 #### **2. 失败处理的四件事** 按顺序执行四个收尾动作: ```text 1. 文档主表标记 BUILD_FAILED 2. 当前任务的所有 chunk 标记 VECTOR_FAILED 3. 所有 step 标记 EXECUTE_FAILED 4. 任务标记 FAILED ``` ##### **顺序为什么重要?** ```text 最先改"文档可见状态"(indexStatus):前端立刻能看到"构建失败" 然后处理"中间数据"(chunk vectorStatus):清理半成品状态 再处理"步骤状态"(step):标记策略执行失败 最后处理"任务本身"(task):整体收尾 ``` 这种**从外向内**的失败传播让用户感知最快——前端最先看到 document 的状态变化。 ##### **如果某一步失败了怎么办?** 注意这四个动作之间**没有 try-catch**。如果 documentMapper.updateById 在 catch 块里又抛异常: ```text 后面三个动作不会执行 chunk / step / task 留在中间态 需要外部补偿任务清理 ``` 为什么不嵌套 try-catch 保护? ```text 1. 已经在异常处理路径上,再嵌套异常会让逻辑非常复杂 2. 在 catch 里再失败概率极低(数据库连接还在) 3. 真出现这种情况通常是数据库整体不可用,补偿也没用 ``` 工程上选了简单——**接受罕见极端情况下的状态不一致**,由外部补偿任务兜底。 #### **3. chunk vectorStatus 全量切 VECTOR\_FAILED** ```text chunkMapper.update(null, new LambdaUpdateWrapper<...>() .eq(...taskId, taskId) .set(...vectorStatus, VECTOR_FAILED) .set(...vectorStoreType, PG_VECTOR)); ``` ##### **按 taskId 范围更新** ```text 只更新当前 taskId 的 chunk 不影响其他历史任务的 chunk ``` 这个范围限定很重要: ```text 同一文档可能有多次索引任务 旧任务可能已经成功 新任务失败时不应该影响旧 chunk 状态 ``` ##### **vectorStatus 的所有可能终态** ```text 正常路径:WAIT_VECTOR → VECTORIZING → VECTOR_SUCCESS 失败路径:WAIT_VECTOR / VECTORIZING → VECTOR_FAILED ``` 为什么要主动切到 VECTOR\_FAILED 而不是留 WAIT\_VECTOR? ```text WAIT_VECTOR:看起来像"还没开始" VECTOR_FAILED:明确"已尝试但失败" 留 WAIT_VECTOR: 定时扫描"待向量化 chunk"的任务可能误捡这批 重新尝试可能继续失败 污染监控数据 切 VECTOR_FAILED: 明确终态 监控能正确统计失败率 重试逻辑可以基于 VECTOR_FAILED 显式触发 ``` 这呼应了文档里那段话:"避免数据库里残留 WAIT\_VECTOR 这种中间态"。 ##### **同时设置 vectorStoreType** ```text .set(...vectorStoreType, PG_VECTOR) ``` 为什么失败时还要设置目标存储? ```text chunk 表里 vectorStoreType 可能是 null(失败发生在切块阶段,还没到向量化) 即使失败也明确告诉运维"这批 chunk 本来要写到 PGVector" 后续重试时知道目标库 ``` 这是**信息完整性**的小细节——失败状态也要保持元数据完整。 #### **4. failTask:任务失败收尾** ```java private void failTask(SuperAgentDocumentTask task, Date startTime, Exception exception, Integer currentStage) { Date finishTime = new Date(); task.setTaskStatus(FAILED); task.setCurrentStage(currentStage); task.setFinishTime(finishTime); task.setCostMillis(finishTime.getTime() - startTime.getTime()); task.setErrorCode("TASK_FAILED"); task.setErrorMsg(exception.getMessage()); taskMapper.updateById(task); } ``` ##### **currentStage 保留进入 catch 时的值** ```java failTask(task, startTime, exception, task.getCurrentStage()); ``` 注意这里传的是 `task.getCurrentStage()`——**也就是异常发生那一刻的阶段**。 为什么保留这个? ```text 切块阶段失败 → currentStage = CHUNK_EXECUTE 向量化失败 → currentStage = VECTORIZE(如果切到这个阶段) 存储失败 → currentStage = STORE_COMPLETE 排查时直接看 currentStage 就知道在哪个阶段挂的 ``` 如果固定填 STORE\_COMPLETE,所有失败任务看起来都"卡在最后一步",定位困难。**保留进入 catch 时的真实阶段**让失败有可追溯性。 ##### **errorCode 用粗粒度 "TASK\_FAILED"** ```java task.setErrorCode("TASK_FAILED"); task.setErrorMsg(exception.getMessage()); ``` errorCode 是粗粒度的固定值,errorMsg 才是异常的详细信息。 为什么不细化 errorCode? ```text 异常种类太多:NPE / SQLException / IOException / IllegalState ... 统一编码成有限分类工作量大 errorMsg 已经携带详细信息 errorCode 主要用于"是否失败"的快速判断 ``` 如果将来运营需要按错误类型统计,可以再加细分逻辑: ```text 判断 exception 类型 → 映射到细分 errorCode 比如 EmbeddingException → errorCode="EMBEDDING_FAILED" ``` 但当前阶段,**粗粒度足够**。 ##### **costMillis 也要算** ```java task.setCostMillis(finishTime.getTime() - startTime.getTime()); ``` 失败任务也要记耗时: ```text 失败前跑了多久 帮助分析"是早期失败还是晚期失败" 失败时间分布可能揭示瓶颈 ``` #### **5. 失败处理 vs 成功处理的对称性** 把 finishTaskSuccess 和 failTask 放一起对比: | 字段 | 成功 | 失败 | | --- | --- | --- | | taskStatus | SUCCESS | FAILED | | currentStage | 显式传入(STORE\_COMPLETE) | 保留发生时阶段 | | finishTime | now | now | | costMillis | 算 | 算 | | errorCode | null(清空) | "TASK\_FAILED" | | errorMsg | null(清空) | exception.getMessage() | 两个方法**结构对称**:都是设置 6 个字段,差别只在值。这种对称设计让代码读起来一目了然——成功和失败是两条对偶的路径。 --- ### **六、状态机全景** 到这里,整条文档处理流水线的状态机就完整了。 #### **1. 四条状态机** ```mermaid stateDiagram-v2 [*] --> WAIT_PARSE: 上传完成 WAIT_PARSE --> PARSING: Consumer 消费 PARSING --> PARSE_SUCCESS: 解析成功 PARSING --> PARSE_FAILED: 解析失败 state PARSE_SUCCESS { [*] --> WAIT_RECOMMEND WAIT_RECOMMEND --> RECOMMENDED: 策略推荐完成 RECOMMENDED --> CONFIRMED: 用户确认 CONFIRMED --> [*]: 触发索引构建 } PARSE_SUCCESS --> WAIT_BUILD: 索引未建 WAIT_BUILD --> BUILDING: 任务启动 BUILDING --> BUILD_SUCCESS: 索引完成 BUILDING --> BUILD_FAILED: 索引失败 BUILD_SUCCESS --> [*]: RAG 可用 ``` #### **2. 状态机彼此联动** 四条状态机不是孤立的,而是**联动推进**: ```text document.parseStatus = PARSE_SUCCESS 是前置: 必须解析成功才能推荐策略 document.strategyStatus = CONFIRMED 是前置: 必须用户确认才能构建索引 document.indexStatus = BUILD_SUCCESS 是终点: 文档可被 RAG 检索 ``` 每个状态都有**入口条件**和**推进条件**,整体形成有向无环图(DAG)。 #### **3. 失败状态的可恢复性** 注意所有失败状态都不是终态: ```text PARSE_FAILED → 用户重新上传 / 后台重新解析 → PARSING → PARSE_SUCCESS BUILD_FAILED → 用户重试构建 → BUILDING → BUILD_SUCCESS ``` 这种**失败可恢复**的设计让系统天然鲁棒——任何一步失败都不是死局。 --- ### **七、阶段四的关键技术点提炼** #### **1. 父子先后入库保证外键完整性** 即使没建数据库外键,业务上严格保证"被引用的先存在",避免崩溃留下孤儿数据。 #### **2. globalChunkNo 全局递增** 整篇文档统一编号而非父块内重置,让 chunkNo 自然映射原文顺序。 #### **3. startChunkNo / endChunkNo 区间冗余** 父块上冗余存子块编号区间,避免详情页查询时回表 MIN/MAX。 #### **4. childCount==0 填 null 而非 0** 明确区分"没有子块"和"有 0 号子块",避免误导性查询。 #### **5. 子块继承父块 sectionPath** 子块没有自己路径时继承父块,保证检索结果总能展示章节信息。 #### **6. token 轻量估算** 不依赖 tokenizer 库,中文按字英文按词其他按 4 字符,统计展示够用。 #### **7. ObjectProvider.getIfAvailable 表达可选依赖** ES 网关不存在时跳过,存在时失败必须传播——可选 ≠ 失败可忽略。 #### **8. EMBEDDING\_BATCH\_SIZE\_LIMIT = 10** 平衡吞吐和失败粒度的经验值,太小开销大太大失败影响范围大。 #### **9. INSERT ON CONFLICT 实现幂等** PGVector upsert 让重试天然幂等,不会因重复执行留下重复向量。 #### **10. vectorId = chunk.id 简化映射** 业务库和向量库共用 ID 空间,不需要额外维护映射表。 #### **11. metadata 只存"轻量过滤字段"** PGVector metadata 用于检索过滤,业务字段(文档名、标签)放 ES,职责分离。 #### **12. embeddingModel 名字写进 metadata** 为模型版本升级时的向量重建留追溯线索。 #### **13. ES Refresh.WaitFor 平衡可见性和性能** 写入后等下次刷新可见,比 NONE 更友好比 Immediate 更稳。 #### **14. ES Bulk API 必须检查 response.errors()** HTTP 200 ≠ 所有 item 成功,漏检会导致部分写入失败但任务标成功。 #### **15. loadDocumentMap 解决 N+1** distinct + selectBatchIds 一次查完,O(1) Map 访问替代 O(N) 列表查找。 #### **16. document 字段冗余进 ES 索引** 避免回表关联,让单次 ES 查询完成"内容匹配 + 文档维度过滤"。 #### **17. chunkId 既作 ES \_id 又作字段** \_id 用于幂等覆盖,同名字段用于 term 查询,两者数据相同语义不同。 #### **18. 标签字符串 → ES 数组** 数据库存 CSV,ES 存数组,让 multi\_match 能命中单个标签。 #### **19. safeText 只 null → 空,不 trim** 最小干预,索引层不擅自改写用户原文语义。 #### **20. 失败时 vectorStatus 主动切 VECTOR\_FAILED** 避免库里残留 WAIT\_VECTOR 中间态污染监控和重试逻辑。 #### **21. failTask 保留 currentStage 真实值** 记录异常发生时的阶段而非固定值,让失败可追溯到具体环节。 #### **22. 成功失败收尾对称设计** finishTaskSuccess / failTask 结构对称,让代码可读性极高。 --- ### **八、面试官可能会问的问题** #### **问题 1:为什么 chunkNo 用全局递增而不是父块内从 1 开始?** 可以回答: > 因为 chunkNo 的核心用途是"在原文中的位置标记"。全局递增让 chunkNo 直接映射原文顺序,按 chunkNo 排序就是阅读顺序,用户看到检索引用 "chunk #5" 就能知道大致位置。父块内重置虽然概念上对齐父子关系,但工程上要二元组 (parentNo, chunkNo) 才能定位,调试麻烦。而且全局编号配合 startChunkNo / endChunkNo 在父块上冗余,反而能用一个区间字段表达父子覆盖关系,比父块内重置更优雅。 #### **问题 2:为什么父块要冗余 startChunkNo / endChunkNo?** 可以回答: > 反范式优化。如果不冗余,每次详情页要展示"这个父块包含哪些 chunk"都要 SELECT MIN(chunkNo), MAX(chunkNo) FROM chunk WHERE parentBlockId=?。落库时一次性算好存到父块,后续读取免聚合查询。代价是 chunk 表更新时要同步父块上的区间(但 chunk 一旦落库不会变,所以这个代价几乎为 0),回报是查询性能。这是典型的"读多写少场景下用冗余换性能"的反范式设计。childCount==0 时填 null 而非 0,是为了明确区分"没有子块"和"有第 0 号子块",避免 SQL 查询出现误导。 #### **问题 3:为什么 PGVector 用 INSERT ON CONFLICT 而不是直接 INSERT?** 可以回答: > 为了让任务重试天然幂等。索引构建任务可能因为各种原因失败重试,第一次执行可能写入了部分向量,第二次执行时这些向量已经存在。如果用纯 INSERT,遇到主键冲突会直接报错,重试逻辑要么先 DELETE 再 INSERT(额外开销且不原子),要么手动检查存在再决定 INSERT 或 UPDATE(两次查询)。INSERT ON CONFLICT DO UPDATE 让 PostgreSQL 一条语句搞定:不存在就插入,存在就覆盖。重试时无需特殊处理,直接重新跑一遍 vectorize 就行。这是分布式系统中实现幂等的经典手法,把幂等性下沉到存储层。 #### **问题 4:为什么向量化要批量调用 EmbeddingModel 而不是逐个?** 可以回答: > 三个理由。第一性能:每次 API 调用都有网络往返开销,500 个 chunk 逐个调要 500 次往返,批量调只要 50 次(按 10 个一批)。第二成本:很多 embedding 模型按调用次数收费,批量能减少计费次数。第三吞吐:embedding 模型通常支持批量推理且批处理效率更高(GPU 并行)。批大小选 10 是经验值,太小开销摊不薄,太大又有问题:单次 prompt 累积太长可能超模型限制,单次失败影响范围大。10 是平衡点。最关键的是要校验返回的向量数量与输入 chunk 数量一致,否则一一对应关系破坏会让检索完全错位。 #### **问题 5:为什么 PGVector 和 ES 都要写?只写一个不行吗?** 可以回答: > 因为它们解决不同的检索问题。PGVector 走向量相似度,擅长"语义相近"召回——用户搜"怎么报销"也能搜到"费用报销流程"这种同义不同字的内容。ES 走倒排索引,擅长"字面匹配"召回——用户搜"FORM-2024-001"这种唯一编号必须精确匹配,或按"分类=财务"过滤特定子集。两者各有短板:向量召回对精确字面差,关键词召回对同义不敏感。后续检索链路会做混合召回——两路都搜然后重排,让用户既能搜到语义相关也能搜到字面命中。这就是 RAG 系统中常见的 hybrid search 思路。所以两个都要写,缺一个体验都打折。 #### **问题 6:为什么 ES 写入用 Refresh.WaitFor 而不是 NONE 或 Immediate?** 可以回答: > 三种 refresh 策略各有取舍。NONE 写入后异步刷新,速度最快,但任务返回后用户立刻搜索可能搜不到刚建的索引——对"建完立刻试用"的场景体验差。Immediate 强制刷新立即可见,但每次写入都触发 refresh 对 ES 集群压力大,高并发下可能拖垮集群。WaitFor 是折中——写入后等待下一次自然 refresh 完成才返回,既保证返回后立即可搜,又不强制额外 refresh。代价是单次写入延迟略高(最多等一个 refresh 间隔,通常 1 秒)。对索引构建这种低频长耗时操作,多等 1 秒可以接受,换来的是用户体验的连续性。 #### **问题 7:ES Bulk API 返回 200 就代表成功了吗?** 可以回答: > 不是,这是 ES Bulk API 的经典陷阱。HTTP 200 OK 只表示请求格式正确、ES 集群接收了批量请求,但单个 item 的成功失败要看 response 里每个 BulkResponseItem 的 error 字段。可能整批 100 条成功 99 条失败 1 条,HTTP 仍然返回 200。所以必须显式调 response.errors() 判断有没有任何失败,再用 response.items() 遍历找出失败的。这套代码做了完整检查:errors() 判断 + stream filter 出失败项 + 拼接 "id:reason" 让错误可定位 + 抛异常让外层接管。漏掉这个检查很常见,会导致部分 chunk 没建到关键词索引但任务标成功,用户搜不到这部分内容也不知道为什么。 #### **问题 8:为什么 ES 索引要冗余文档维度字段(documentName / businessCategory 等)?** 可以回答: > 因为 ES 检索时要在单条索引文档内完成所有过滤和打分。设想用户搜"财务报销" + "只看财务分类"。如果不冗余 businessCategory:ES 先搜文本命中得到一批 chunk\_id,然后回数据库查每个 chunk 对应文档的 businessCategory,再过滤"财务"分类——链路长且打分时无法考虑分类匹配权重。如果冗余 businessCategory 进 chunk 索引:ES 一次查询就能 "match chunkText AND term businessCategory",单跳完成检索。代价是存储冗余(同一文档的所有 chunk 都重复存 documentName 等字段),但 ES 本身是为搜索而生,存储成本远小于查询性能收益。这是 ES 设计中"宽表化"的常见做法。 #### **问题 9:为什么 ES 文档 \_id 用 chunkId?** 可以回答: > 为了幂等覆盖。ES 中 \_id 是文档主键,同 \_id 的写入会直接覆盖原文档。任务重试时同一批 chunk 重新走 indexChunks,因为 \_id 都是 chunkId,新写入会覆盖旧文档而不是新增重复——天然幂等。如果用 ES 自动生成 \_id 或者随机 UUID,重试会写入新文档,留下重复索引污染检索结果。注意这里有个细节:indexRecord 内部还有 chunk\_id 字段,因为 ES \_id 不能直接 term 查询,必须有同步的字段才能用作过滤条件。所以是"主键 + 字段双存",主键负责幂等,字段负责检索。这种设计在 ES 中很常见。 #### **问题 10:失败处理为什么要主动把 chunk 切到 VECTOR\_FAILED 而不是留 WAIT\_VECTOR?** 可以回答: > 为了避免中间态残留污染系统。WAIT\_VECTOR 和 VECTORIZING 都是中间态,VECTOR\_FAILED 才是失败的终态。如果失败时不主动切:定时扫描"待向量化 chunk"的补偿任务可能误捡这批已经失败过的 chunk 反复重试;监控指标统计不准(看到一堆"待向量化"实际是失败的);运维排查时分不清"还没开始"和"已经失败"。主动切 VECTOR\_FAILED 让中间态明确终态化,监控、重试、排查都基于明确状态决策。这呼应了一个原则:永远不让中间态在系统中长期存在,要么推进到成功终态,要么明确标记失败终态,避免"既不是开始也不是结束"的灰色地带。 #### **问题 11:catch 块里如果再抛异常怎么办?** 可以回答: > 当前代码没有保护——catch 块里的 4 个 update 操作如果中途某一步失败,后续操作不会执行,留下部分中间态。比如 document 切到 BUILD\_FAILED 成功,chunk 还没切就再次抛异常,chunk 表里会留 WAIT\_VECTOR 和 BUILD\_FAILED 不一致的状态。这是个潜在缺陷,但工程上接受了:第一概率极低(catch 里都是简单 update,数据库连接还在),第二真出现这种情况通常是数据库整体不可用,再怎么保护也救不回来,第三嵌套 try-catch 会让代码非常复杂。更好的做法是写一个外部补偿任务定期扫描"卡在中间态超过 N 分钟"的 chunk/task,自动修复或告警。这是"主流程简单 + 外部补偿兜底"的常见组合。 #### **问题 12:从用户上传到 RAG 可用,整个流水线最关键的几个设计点是什么?** 可以回答: > 要点很多,挑核心的几个。第一,**异步链路解耦**:上传/解析/索引各自异步,Kafka 串联,避免同步链路超时和并发瓶颈。第二,**人机协同确认**:策略推荐后等用户确认才执行索引,让自动化有人工把关余地。第三,**Parent-Child 双层切块**:父块负责回答上下文,子块负责检索精度,解决"小块召回 vs 大块完整"的核心矛盾。第四,**多策略组合 + 多层降级**:四种切块策略按方案组合,每种内部都有降级路径,保证任何文档都能跑通。第五,**双索引并行**:PGVector 语义召回 + ES 关键词召回,混合检索覆盖所有查询场景。第六,**幂等设计**:UPSERT、\_id 覆盖、状态机,让任务重试天然安全。第七,**状态机驱动**:document/task/plan/step 四条独立状态机相互联动,让链路任意时刻状态可观测可追溯。这七个设计点撑起了整套系统的健壮性和可演进性。 --- 到这里,整个文档处理流水线就**彻底讲完了**。 回顾从第一篇到这一篇走过的路: ```text 1. 上传接口与文档主记录创建 2. Kafka 异步消费 + Tika 文本解析 3. 结构节点提取四阶段流水线 4. 解析结果统计与异步收尾 5. 策略推荐与方案持久化 6. 索引构建入口与 Kafka 消息投递 7. 异步索引构建:初始化与切块执行 8. 四种切块策略详解 9. 异步索引构建:落库、向量化与收尾(本篇) ``` 九篇连起来覆盖了从用户点击"上传"到文档"可用于 RAG 检索"的完整链路。每一篇都不只是讲代码,更讲背后的设计权衡——为什么这么做、不这么做会怎样、有哪些备选方案、工程上为什么这么取舍。 下一阶段可以进入更上层的话题: ```text 检索链路:用户提问 → 召回 → 重排 → 答案生成 混合检索:PGVector + ES 双路召回怎么合并打分 重排模型:rerank 模型的接入和效果优化 答案生成:prompt 设计、引用展示、流式输出 反馈闭环:用户反馈如何改进检索质量 ``` 这些话题会让你看到**RAG 从存储到使用**的完整闭环——索引建好只是基础,怎么用好这些索引才是产品价值的来源。 --- **企业级项目导航**:⬅️ [[08-索引构建链路白话讲解|08-索引构建链路白话讲解]] | 09-落库向量化收尾 | ➡️ [[01-从用户提问到答案返回的总流程|01-从用户提问到答案返回的总流程]]