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