构建索引
我们接着上一篇“策略推荐与方案持久化”往下走。
上一篇的终点是:
plan 主记录已落库,状态 WAIT_CONFIRM
父子步骤已批量入库,executeStatus 都是 WAIT_EXECUTE
文档主表 parseStatus = PARSE_SUCCESS,strategyStatus = RECOMMENDED
解析任务按 SUCCESS 收尾
接下来用户在前端确认了方案(strategyStatus 切到 CONFIRMED),然后点击“构建索引”按钮——这一篇就是从这一刻开始的。
学完这一篇,你会理解:
1. 一次"构建索引"请求,后端同步链路到底做了哪些事
2. 为什么要做这么多前置校验
3. 为什么策略快照要固化到任务记录上
4. 为什么 kafkaTemplate.send 后面要跟 .get()
5. 同步链路和异步链路的边界在哪里
下一篇才进入真正的重头戏:异步消费端的 handleIndexBuild() 处理逻辑。
一、先建立整体认知:这一节在做什么?
可以一句话概括:
用户在前端点"构建索引"后,Controller 接收请求,Service 做四道前置校验、创建 BUILD_INDEX 类型的任务记录、把策略快照固化到任务上、把文档状态切到 BUILDING、记录启动日志,最后通过 Kafka 同步发送一条只包含三个 ID 的消息把后续繁重工作交给消费端处理。
整体流程图:
flowchart TD
A[前端 POST /index/build] --> B[Controller 接收 DocumentIndexBuildDto]
B --> C[Service.buildIndex 入口]
C --> D[校验1: 文档存在 + parseStatus + strategyStatus]
D --> E[校验2: planId 与 currentPlanId 一致]
E --> F[校验3: 无运行中的索引任务]
F --> G[校验4: plan 真实存在且有效]
G --> H[创建 BUILD_INDEX 任务记录]
H --> I[固化 strategySnapshot 到任务]
I --> J[文档 indexStatus 切到 BUILDING]
J --> K[写 START 任务日志]
K --> L[Kafka 同步发送 IndexBuildMessage]
L --> M[返回 VO 给前端]
L --> N[Consumer 异步消费]
N --> O[handleIndexBuild 真正处理]
这一节的定位
它是“策略推荐”和“异步索引构建”之间的桥梁。这条同步链路不做任何重活:
不做切块
不做向量化
不写索引
甚至不读解析后的文本
它只做两件事:校验 和 投递。这种“轻同步 + 重异步”的拆分,是后端设计中非常重要的工程模式。
二、为什么要把"构建索引"做成异步任务?
很多新手一上来会想:直接在 Controller 里把切块和向量化做完不行吗?
不行。原因有四个层面。
1. 单次操作耗时太长
索引构建的实际耗时分布:
切块:数百毫秒 ~ 几秒
向量化:每个 chunk 调一次模型,可能几百次,十几秒~分钟级
写 PG/ES/Neo4j:几秒
图谱关系构建:可能十几秒
加起来一份大文档可能需要 30 秒到几分钟。
如果同步执行:
HTTP 连接挂着不断
网关超时(通常 60 秒)
浏览器看到的是请求 pending
用户体验极差
2. 接口性能稳定性
同步链路里调用 LLM 或 embedding 服务,受外部服务抖动影响极大:
LLM 接口偶尔慢一些 → 整个 HTTP 请求慢
模型限流 → 请求堆积
向量库慢查询 → 接口变慢
把重活推到后台,接口的 P99 永远稳定在毫秒级。
3. 失败重试与幂等性
异步任务可以基于 Kafka 的消费失败机制天然支持重试:
消费者抛异常 → broker 自动重投
重试次数可配
死信队列兜底
而 HTTP 同步接口的重试只能依赖客户端,复杂度高。
4. 资源削峰
如果同时有 100 个用户点构建索引:
同步执行 → 100 个连接 + 100 个线程同时调 LLM → 服务被打挂
异步队列 → 消息排队,Consumer 按消费能力消化
Kafka 在这里既是异步通道,也是流量缓冲。
结论
“构建索引”典型符合长耗时、高风险、可异步的特征,必须做成异步任务。同步链路只做请求受理。
三、Controller 层:极简入口
@PostMapping("/index/build")
public ApiResponse<DocumentIndexBuildVo> buildIndex(@Valid @RequestBody DocumentIndexBuildDto dto) {
return ApiResponse.ok(documentManageService.buildIndex(dto));
}
Controller 这一层非常薄,几乎没有逻辑。这是合理的设计——Controller 的职责就是路由和参数校验,业务逻辑在 Service。
请求 DTO 只有三个字段:
@NotNull
private Long documentId;
@NotNull
private Long planId;
private Long operatorId; // 允许为空
为什么 documentId 和 planId 都要传?
只传 documentId 不行吗?后端不是有 currentPlanId 吗?
不行,原因是前端可能基于过期数据发请求:
用户在 t1 时刻打开方案确认页,看到方案 A
用户在 t2 时刻管理员把方案改成了 B(currentPlanId 变了)
用户在 t3 时刻点了"构建索引"按钮
如果只传 documentId,后端会用 B 去构建。但用户期望的是 A。
让前端显式传 planId,后端校验“传入的 planId == currentPlanId”,发现不一致就拒绝。这样前端会刷新页面让用户重新确认,避免出现“用户以为构建的是 A 实际却是 B”的语义错误。
这是一种 乐观锁思想 的体现。
operatorId 为什么允许为空?
因为索引构建有两种触发方式:
用户主动触发:operatorId 不为空
系统自动触发(比如批量重建、定时任务):operatorId 为空
后续会看到,operatorId 决定了任务的 triggerSource(USER 还是 SYSTEM)和日志的 operatorType,便于审计和统计。
四、Service 层:四道校验
buildIndex() 进来第一件事就是校验,一共四道,层层递进。
校验 1:文档存在 + 状态合法
SuperAgentDocument document = getDocumentOrThrow(dto.getDocumentId());
if (!Objects.equals(document.getParseStatus(), PARSE_SUCCESS)
|| !Objects.equals(document.getStrategyStatus(), CONFIRMED)) {
throw new SuperAgentFrameException(...,
"当前文档尚未完成"解析成功 + 策略确认",不能构建索引。");
}
这一步检查两个前置条件:
parseStatus = PARSE_SUCCESS:解析必须成功
strategyStatus = CONFIRMED:策略必须用户确认过
为什么必须 parseStatus = PARSE_SUCCESS?
因为索引构建的输入是解析后的纯文本(也就是 MinIO 里那份 parseTextPath 指向的 .txt)。如果解析没成功:
没文本 → 切块切空 → 索引为空
解析失败 → 文本可能不完整 → 索引质量差
为什么必须 strategyStatus = CONFIRMED?
回顾上一篇,策略推荐完成后状态是 RECOMMENDED,意味着:
方案是系统推荐的草案
还没经过用户审核
不能直接拿来执行
只有用户在前端点了“确认”按钮,状态切到 CONFIRMED,才表示:
用户审核了这套方案
用户接受了系统的推荐(或自己改过)
可以放心执行
这是上一篇“人机协同”设计的延续:方案必须经过人确认,才能进入执行阶段。
为什么用 ! 或 ! 而不是 && 取反?
if (! parseStatus == PARSE_SUCCESS || ! strategyStatus == CONFIRMED)
这种写法的语义是:任何一个条件不满足就拒绝。也就是说:
解析未成功 → 拒绝
解析成功但策略未确认 → 拒绝
解析未成功且策略未确认 → 拒绝
两个都满足 → 放行
逻辑等价于德摩根律:!(A && B) = !A || !B。
校验 2:planId 一致性
if (!Objects.equals(document.getCurrentPlanId(), dto.getPlanId())) {
throw new SuperAgentFrameException(...,
"当前文档的生效方案与请求方案不一致。");
}
这一步前面解释过了,是为了防止前端基于过期数据发请求。
为什么不直接以 currentPlanId 为准?
也就是“忽略前端传的 planId,直接用 document.currentPlanId”。
不行。因为这样会把语义错误变成静默执行:
用户期望构建方案 A
后端默默用方案 B 构建了
用户看到结果跟预期不一致也不知道为什么
显式校验 + 显式拒绝,能让前端弹出“方案已变更,请刷新”的提示,用户可以主动刷新重新确认,避免认知偏差。
校验 3:防重复提交
long runningTaskCount = taskMapper.selectCount(new LambdaQueryWrapper<>()
.eq(SuperAgentDocumentTask::getDocumentId, dto.getDocumentId())
.eq(SuperAgentDocumentTask::getTaskType, BUILD_INDEX)
.in(SuperAgentDocumentTask::getTaskStatus, NEW, RUNNING)
.eq(SuperAgentDocumentTask::getStatus, BusinessStatus.YES));
if (runningTaskCount > 0) {
throw new SuperAgentFrameException(...);
}
直接查数据库:这个文档当前有没有正在跑的索引任务?
场景模拟
不做这个校验会怎样?
用户点了一次"构建索引",任务 T1 创建,Kafka 消息发出去
1 秒后用户没看到反馈,以为没点上,又点一次
任务 T2 创建,Kafka 消息再发一次
Consumer 同时消费两条消息:
T1 在切块、向量化、写 PG
T2 也在切块、向量化、写 PG
两个任务并发写同一份文档的索引数据
后果非常严重:
(同一段被写两次)
索引数据相互覆盖或交错
切块结果不一致(LLM 切块本身有随机性)
检索召回会捞到双倍内容,影响 RAG 质量
向量库存储成本浪费
查询条件细节
documentId 限定:只看当前文档
taskType 限定:只看 BUILD_INDEX 类型(解析任务不算)
taskStatus IN (NEW, RUNNING):新建中或执行中都算"在跑"
status = YES:逻辑未删除的
四个条件缺一不可。特别是 taskType 的限定很关键——同一文档可能同时有“解析任务”和“索引任务”,不能混淆。
为什么用 NEW + RUNNING 两个状态?
因为任务从创建到被消费有时间差:
NEW:刚 insert,Kafka 消息还没被消费
RUNNING:Consumer 已经接到消息,正在处理
两个都要拦,否则 NEW → RUNNING 的瞬间漏掉。
这是不是绝对防并发?
不是。这只是“业务级幂等”,不是“数据库级幂等”。
考虑高并发:
请求 A 进来,查询发现没有运行中的任务
请求 B 进来,查询也发现没有运行中的任务(此时 A 还没 insert)
A insert 任务 T1
B insert 任务 T2
两个任务都创建了
这是典型的TOCTOU 竞态(Time of Check to Time of Use)。
要绝对防并发,需要:
方案一:数据库唯一索引(documentId + 状态)
方案二:分布式锁(Redis SETNX)
方案三:乐观锁(更新 document 的 indexVersion)
但实际产品中这种并发很少(一个用户不太可能毫秒级双击),加上 Kafka 消费端的去重机制,业务级校验已经足够。在工程上选择"足够好"而不是"完美",这是务实的判断。
校验 4:plan 真实存在
SuperAgentDocumentStrategyPlan plan = planMapper.selectById(dto.getPlanId());
if (plan == null || !Objects.equals(plan.getStatus(), BusinessStatus.YES)) {
throw new SuperAgentFrameException(...);
}
虽然校验 2 已经确认了 planId == currentPlanId,但这里还要再查一次 plan 表。
为什么要重复查?
因为校验 2 和校验 4 防的是不同问题:
校验 2:防方案不一致(请求和当前生效方案对不上)
校验 4:防方案被删除或失效(plan 表里 status=NO 或记录被物理删除)
理论上 currentPlanId 指向的方案不应该被删,但工程上不能假设上游数据完美。多查一次是廉价的保险。
更重要的目的:拿到 strategySnapshot
后面要把 strategySnapshot 固化到任务上。这个字段在 plan 表里,不重新查就拿不到。
也就是说,校验 4 既是校验,也是数据准备。
四道校验的顺序很讲究
1. 文档存在 + 状态对(最基础,大概率失败)
2. planId 一致(中频失败,前端缓存场景)
3. 无运行中任务(中低频)
4. plan 有效(几乎不会失败,但顺手拿数据)
按失败概率从高到低排,越早失败越早返回,避免做无用功。这是后端校验的常见排序原则。
五、创建任务记录:把"承诺"落库
校验全部通过后,创建任务:
Long taskId = uidGenerator.getUid();
SuperAgentDocumentTask task = new SuperAgentDocumentTask();
task.setId(taskId);
task.setDocumentId(document.getId());
task.setPlanId(dto.getPlanId());
task.setTaskType(BUILD_INDEX);
task.setTaskStatus(NEW);
task.setCurrentStage(CHUNK_EXECUTE);
task.setTriggerSource(resolveTriggerSource(operatorId));
task.setStrategySnapshot(plan.getStrategySnapshot());
task.setRetryCount(0);
taskMapper.insert(task);
每个字段都有明确目的,逐个分析。
1. taskId:UID 生成器
Long taskId = uidGenerator.getUid();
不用数据库自增 ID,而是用分布式 UID 生成器。原因:
分布式部署多实例时不冲突
插入前就能拿到 ID,避免 insert 后再回查
ID 本身带时间信息,方便排序和分析
2. taskType = BUILD_INDEX
任务类型枚举区分:
PARSE_DOCUMENT:解析任务(上一篇讲过)
BUILD_INDEX:索引构建任务(这一篇)
DELETE_INDEX:删除索引(可能存在)
REBUILD_INDEX:重建索引(可能存在)
后续查任务、统计、监控都按 taskType 分类。
3. taskStatus = NEW
任务初始状态。后续状态机:
NEW → RUNNING(Consumer 开始消费)
RUNNING → SUCCESS(全部阶段完成)
RUNNING → FAILED(任意阶段失败)
为什么不是 RUNNING?因为 Kafka 还没投递成功,任务还没被消费。这一步只是“预约”。
4. currentStage = CHUNK_EXECUTE
这里非常关键。索引构建链路的第一个执行阶段是“切块执行”。
回想上一篇任务阶段:
FILE_UPLOAD → CONTENT_PARSE → STRATEGY_ROUTE (解析任务)
CHUNK_EXECUTE → VECTORIZE → INDEX_WRITE (索引任务)
把 currentStage 初始化为 CHUNK_EXECUTE,是给 Consumer 的“起点提示”。Consumer 不需要自己判断从哪个阶段开始,直接读 currentStage 就行。
这样设计的好处
Consumer 失败重试时,可以从 currentStage 继续(而不是从头来过)
监控可以基于 currentStage 看任务推进到哪一步
卡住的任务一眼能看出卡在哪个阶段
5. triggerSource:触发来源
task.setTriggerSource(resolveTriggerSource(operatorId));
resolveTriggerSource 大概率是这样的逻辑:
return operatorId == null ? SYSTEM : USER;
为什么要区分?
USER 触发:用户主动构建,失败要给用户反馈
SYSTEM 触发:批量任务或定时任务,失败发告警给运维
统计上分开看:用户行为 vs 系统自动化
这是审计字段的常见设计。
6. strategySnapshot 固化(核心设计)
task.setStrategySnapshot(plan.getStrategySnapshot());
这是整段代码最值得讲的设计点。
为什么要在任务上冗余存一份策略快照?
明明 task.planId 已经指向了 plan 表,为什么不直接 join?
考虑下面这个时间线:
t0:任务创建,planId=100,plan 100 的 snapshot 是 "PARENT:1;CHILD:4,2"
t1:Kafka 消息投递成功
t2:用户在前端"反悔"了,重新调整了方案,plan 100 的 snapshot 变成 "PARENT:2;CHILD:2"
t3:Consumer 开始消费 t0 的消息
如果 Consumer 实时查 plan 表:
拿到的是 t3 时刻的方案 "PARENT:2;CHILD:2"
但用户在 t0 确认的是 "PARENT:1;CHILD:4,2"
执行结果和用户预期不一致
把快照固化到任务上后:
Consumer 直接读 task.strategySnapshot
拿到的是 t0 时刻的快照
按用户当时确认的方案执行
这是"事件溯源"思想
任务记录本质上是一个“已发生的事件”,事件应该是不可变的。事件的语义是:
"在 t0 时刻,用户基于方案快照 X 触发了索引构建"
事件一旦发生,用什么方案是确定的,不应该被后续操作改变。
这种思想在金融、订单、审计系统中非常常见:
订单创建时商品价格冗余存(防止商品涨价后历史订单变价)
合同签约时条款冗余存(防止条款修订影响已签合同)
任务触发时配置冗余存(防止配置变更影响已触发任务)
冗余存储的代价
存储多一份字符串
plan 改了之后任务里的快照不会自动更新
这些代价是值得的,因为换来的是确定性和可追溯性。
7. retryCount = 0
预留字段,给 Consumer 失败重试时累加。
每次失败重试 → retryCount++
超过最大重试次数 → 进死信队列
8. 为什么这一阶段不写策略步骤表?
回想上一篇,策略步骤表(strategy_step)已经在“策略推荐”阶段写好了,每一步都是 WAIT_EXECUTE 状态。
索引构建任务不重写步骤,而是:
通过 task.planId 找到 plan
通过 plan.id 找到 step 列表
按 stepNo 顺序执行
执行时把 step.executeStatus 从 WAIT_EXECUTE 切到 RUNNING/SUCCESS/FAILED
这是一个很自然的“配置 vs 执行”分离:
plan + step:配置层(描述"应该怎么做")
task:执行层(描述"什么时候做、做得怎么样")
六、文档状态切换:BUILDING
document.setIndexStatus(BUILDING);
documentMapper.updateById(document);
把文档主表的 indexStatus 切到 BUILDING。
indexStatus 状态机
回想完整状态:
WAIT_BUILD → BUILDING → BUILT(成功)
↘ BUILD_FAILED(失败)
为什么要立刻切?
因为前端需要实时反馈:
用户点了"构建索引"
列表页应该立刻显示"构建中"
详情页应该立刻显示进度状态
如果等 Consumer 消费时再切,会有几秒~几十秒延迟,用户会以为系统没反应。
状态字段的颗粒度
注意文档主表上有三个独立的状态:
parseStatus:解析阶段状态
strategyStatus:策略阶段状态
indexStatus:索引阶段状态
每个阶段一个独立字段,互不干扰。这样:
parseStatus=PARSE_SUCCESS + strategyStatus=CONFIRMED + indexStatus=BUILDING
解析成功 + 策略确认 + 正在建索引
parseStatus=PARSE_SUCCESS + strategyStatus=CONFIRMED + indexStatus=BUILT
全部完成,可用于 RAG
如果合并成单一状态字段,组合数会爆炸。多状态字段是工程上的合理设计。
七、记录启动日志
taskLogService.saveLog(taskId, document.getId(),
CHUNK_EXECUTE,
START,
INFO,
resolveOperatorType(operatorId),
operatorId,
"索引构建任务已创建,等待异步执行。",
Map.of("planId", dto.getPlanId(), "strategySnapshot", plan.getStrategySnapshot()));
这条日志和上一篇的解析日志结构一致,但语义不同:
stage = CHUNK_EXECUTE:对应任务的 currentStage
eventType = START:任务启动事件
level = INFO:正常流程
operatorType = USER 或 SYSTEM:看 operatorId 是否为空
metadata:planId + strategySnapshot
为什么 metadata 里要塞 strategySnapshot?
因为日志是“那一刻的快照”。
设想运维场景:
某用户反馈"我建的索引结果跟预期不一样"
排查时打开任务日志
立刻看到当时的 strategySnapshot
不需要回去查 plan 表(plan 可能后来又改过)
任务日志的 metadata 应该包含重现现场所需的所有关键信息。
日志和任务记录的关系
任务记录(task):描述"任务现在的状态"
任务日志(task_log):描述"任务发生过哪些事件"
任务表的字段会被覆盖(比如 taskStatus 从 NEW 变成 SUCCESS),但日志是追加写:
START 事件
RUNNING 事件
CHUNK_DONE 事件
VECTORIZE_DONE 事件
SUCCESS 事件
任意时刻把所有日志连起来,就能还原完整的任务执行历史。
八、Kafka 消息投递:同步等待
kafkaProducer.sendIndexBuild(
new DocumentIndexBuildMessage(document.getId(), taskId, dto.getPlanId()));
消息体设计:只有三个 ID
public class DocumentIndexBuildMessage {
private Long documentId;
private Long taskId;
private Long planId;
}
为什么消息体这么瘦?为什么不把 strategySnapshot、parseTextPath 这些都塞进去?
原因 1:消息只是"通知",不是"数据"
Kafka 消息的角色是“告诉 Consumer 该干活了”,具体要干什么、用什么数据,Consumer 自己去数据库查。
胖消息:把所有上下文塞到消息里
优点:Consumer 不用查库
缺点:消息大、Kafka 压力大、信息冗余、易过期
瘦消息:只塞 ID
优点:消息小、Kafka 高效
缺点:Consumer 多几次查库
工程上几乎都选瘦消息:
Kafka 不是数据库,不该承担数据存储职责
ID 永远准确,具体数据查表保证最新
消息无状态,容易重试和补偿
原因 2:避免数据不一致
如果消息里塞 strategySnapshot,但任务表里也存了一份,两份可能不一致(比如发送中网络抖动消息丢了一部分字段)。瘦消息天然避免这个问题——所有数据以数据库为准。
三个 ID 各自的作用
documentId:Consumer 查文档主表(parseTextPath、indexStatus 等)
taskId:Consumer 更新任务状态、记录日志
planId:Consumer 查 strategy_step 列表(也可以从 task 反查)
planId 其实有点冗余(通过 taskId 也能查到),但显式带上让 Consumer 不用再查一层。轻微冗余换性能,划算。
生产者发送:.get() 阻塞等待
kafkaTemplate.send(topic, key, payload).get();
这是这一段最值得讲的细节。
kafkaTemplate.send() 默认是异步的
返回的是 ListenableFuture<SendResult>,调用线程不等待发送完成就返回。这种异步发送性能更高,但有风险:
应用收到了 send() 调用,以为发出去了
实际 Kafka 可能因为网络问题、broker 不可达等原因没发出去
任务记录已经插入,但消息没投递成功
Consumer 永远不会消费 → 任务永远卡在 NEW
这就是**"任务存在但永远不会执行"**的不一致状态。
加 .get() 变成同步发送
.get() 会阻塞直到收到 broker 的 ack
broker ack 表示消息已写入 Kafka
此时再返回,后面的代码才执行
如果发送失败
.get() 抛异常
被外层 try-catch 捕获
转换成 SuperAgentFrameException 抛出
Spring 事务回滚 → task insert、document update 全部回滚
返回给前端"Kafka 发送失败"
这就保证了事务性:
要么(任务创建 + 文档状态切换 + 消息投递)全部成功
要么全部回滚
不会出现部分成功的中间态
代价
同步等待会阻塞接口
延迟变大(通常几毫秒~几十毫秒)
吞吐量下降
但对“构建索引”这种低频操作来说,正确性远比性能重要。如果是高 QPS 场景(比如埋点上报),才会考虑 fire-and-forget。
更严格的事务保证
注意:即使加了 .get(),也不是 100% 严格的分布式事务。考虑:
.get() 返回成功 → Kafka 已收到
但事务还没提交
事务在 commit 时数据库挂了 → 事务回滚
但 Kafka 消息已经发出去了
Consumer 消费时发现 task 不存在 → 异常
这是经典的**"本地事务 + 消息队列"分布式一致性问题**。
完美方案有:
1. 事务消息(RocketMQ Transactional Message)
2. 本地消息表 + 定时扫描补偿
3. CDC(监听数据库 binlog,自动发消息)
这套代码用最朴素的“先 insert 再发消息 + 同步 .get()”,在大部分场景够用,但极端情况下可能有问题。Consumer 端需要做幂等保护:
查 task 不存在 → 直接丢弃消息(不报错)
查 task 状态已经是 SUCCESS → 直接丢弃(防重复消费)
工程上接受**"小概率不一致 + Consumer 幂等"**的组合,性价比最高。
Topic 命名:环境隔离
SpringUtil.getPrefixDistinctionName() + "-" + properties.getKafka().getIndexTopic()
prefixDistinctionName 通常是环境名(dev/test/prod)或租户标识。最终 topic 可能是:
prod-document-index-build
test-document-index-build
dev-document-index-build
为什么要这样?
多环境共用一个 Kafka 集群时不会串数据
开发环境的消息不会被生产 Consumer 消费
线上故障时单环境隔离影响
这是多环境部署的标准做法。
Key 用 documentId:保证顺序
send(topic, String.valueOf(message.getDocumentId()), message);
Kafka 中相同 key 的消息会进入同一分区,同一分区按写入顺序消费。用 documentId 做 key 的好处:
同一文档的多个任务(比如重建索引)会被同一 Consumer 顺序消费
不会出现"重建任务比初次构建任务先执行"的乱序
如果 key 用 taskId 或随机值,就没有这个保证。
九、返回 VO:即时反馈
return new DocumentIndexBuildVo(
document.getId(), taskId,
task.getTaskType(), enumMsg(...),
task.getTaskStatus(), enumMsg(...),
document.getIndexStatus(), enumMsg(...)
);
返回给前端的 VO 包含:
documentId:文档 ID
taskId:任务 ID(前端可以用它轮询进度)
taskType / taskStatus / indexStatus:三个枚举值
对应的中文描述(enumMsg)
为什么要返回 taskId?
前端拿到 taskId 后可以:
1. 轮询任务进度接口:基于 taskId 查询当前阶段
2. 订阅 WebSocket:推送任务进度更新
3. 失败时让用户点"查看详情":跳到任务详情页(URL 带 taskId)
为什么枚举值要带"中文描述"?
taskStatus = 1(原始码)
enumMsg(...) = "新建中"
返回原始码是为了前端逻辑判断,返回中文描述是为了前端可以直接渲染到 UI,不需要前端再维护一份枚举映射表。
这种“码 + 描述”的双字段设计在前后端协作中很常见。
十、消费端入口:极简反序列化
@KafkaListener(topics = ..., groupId = "...-index")
public void consumeIndexBuild(String payload) {
try {
DocumentIndexBuildMessage message =
objectMapper.readValue(payload, DocumentIndexBuildMessage.class);
asyncProcessService.handleIndexBuild(
message.getDocumentId(), message.getTaskId(), message.getPlanId());
} catch (Exception exception) {
log.error("消费索引构建消息失败,payload={}", payload, exception);
}
}
消费端非常简洁,三件事:
1. 反序列化消息
2. 调用 asyncProcessService.handleIndexBuild()
3. 异常打日志
为什么 Consumer 这么薄?
因为它只是“入口”:
解析消息格式
路由到业务方法
不做任何业务逻辑
业务逻辑在 asyncProcessService.handleIndexBuild() 里。这种分层让:
Kafka 框架细节(序列化、消费者配置)集中在 Consumer
业务逻辑在 Service,可以单独测试(不依赖 Kafka)
groupId 的作用
groupId = "...-index"
groupId 决定了消费分组:
同一 groupId 的多个 Consumer 实例:互相分摊消息(每条只被一个实例消费)
不同 groupId:各自独立消费(每条被多个组消费)
这里 -index 后缀让索引构建有自己独立的消费组,与其他任务消费组(比如 -parse、-delete)隔离:
索引消费阻塞 → 不影响解析消费
索引消费扩容 → 只扩 -index 组
catch Exception 的设计
catch (Exception exception) {
log.error("消费索引构建消息失败,payload={}", payload, exception);
}
注意这里只打日志,不向上抛。这意味着:
消息会被认为消费成功(commit offset)
不会自动重试
为什么这样设计?
1. 异步链路有自己的失败处理(failTask 标记 FAILED)
2. 反序列化失败的消息再消费多少次都还是失败,重试无意义
3. 业务异常已经在 handleIndexBuild 内部 try-catch 处理过
这是“消息消费成功 ≠ 业务执行成功”的设计。Kafka 只保证消息不丢,业务成败由数据库的任务状态记录。
如果要更严格的重试,可以:
不 catch,让 Kafka 自动重试(配合 SeekToCurrentErrorHandler)
或者主动发到 DLQ(死信队列)
工程上看团队策略,这套代码选了“Consumer 幂等 + 任务状态托管”的方案。
十一、整段同步链路的状态视图
把所有状态变化串起来:
sequenceDiagram
participant FE as 前端
participant C as Controller
participant S as Service
participant DB as 数据库
participant K as Kafka
participant CON as Consumer
FE->>C: POST /index/build {documentId, planId}
C->>S: buildIndex(dto)
S->>DB: 查 document
S->>S: 校验 parseStatus + strategyStatus
S->>S: 校验 planId 一致
S->>DB: 查运行中任务数
S->>S: 校验无并发任务
S->>DB: 查 plan
S->>S: 校验 plan 有效
S->>DB: insert task (NEW, CHUNK_EXECUTE)
S->>DB: update document (indexStatus=BUILDING)
S->>DB: insert task_log (START)
S->>K: send(topic, documentId, message).get()
K-->>S: ack
S-->>C: VO
C-->>FE: 200 OK
K-->>CON: deliver message
CON->>CON: handleIndexBuild (异步,下一篇)
文档/任务状态变化:
入口前:
document: parseStatus=PARSE_SUCCESS, strategyStatus=CONFIRMED, indexStatus=WAIT_BUILD
task:无运行中
入口后:
document: indexStatus=BUILDING
task(新): status=NEW, currentStage=CHUNK_EXECUTE, 携带 strategySnapshot
task_log:START 事件
待 Consumer 消费后:
task: status=RUNNING
后续状态由异步链路推进
十二、把这条链路串成一句话
前端 POST /index/build 携带 documentId 和 planId 进入 buildIndex,先依次做四道校验:文档存在且 parseStatus=PARSE_SUCCESS 且 strategyStatus=CONFIRMED、请求 planId 与 document.currentPlanId 一致、当前文档没有 NEW/RUNNING 状态的 BUILD_INDEX 任务、plan 真实存在且有效;校验通过后基于 UidGenerator 生成 taskId,创建一条 BUILD_INDEX 类型、状态 NEW、currentStage=CHUNK_EXECUTE 的任务记录,把 plan.strategySnapshot 固化到任务上以保证异步执行时不受方案后续修改影响;接着把文档 indexStatus 切到 BUILDING 让前端实时看到"构建中",写一条 START 任务日志带上 planId 和 strategySnapshot 作为现场快照;最后通过 kafkaTemplate.send(topic, documentId, message).get() 同步发送只含 documentId/taskId/planId 三个 ID 的瘦消息到环境前缀 + 索引 topic,发送失败抛异常让整个事务回滚保证不出现"任务存在但消息未投递"的不一致状态;返回 VO 包含 taskId 让前端可以轮询进度。Consumer 在 -index 消费组接到消息后反序列化并调用 asyncProcessService.handleIndexBuild(),真正的切块/向量化/索引写入由异步链路负责。
十三、核心技术点提炼
1. 轻同步 + 重异步 的接口模式
同步链路只做校验和投递,不做任何重活。索引构建这种长耗时、高风险、可异步的操作必须解耦。
2. 四道校验的递进设计
文档状态 → 方案一致 → 防并发 → plan 有效,按失败概率排序,越早失败越早返回。
3. 双状态前置依赖
parseStatus=PARSE_SUCCESS 保证有文本可切,strategyStatus=CONFIRMED 保证策略经用户确认。两者缺一不可。
4. planId 显式校验防过期数据
不直接以 currentPlanId 为准,避免"用户期望 A 实际执行 B"的语义偏差。这是乐观锁思想。
5. 业务级防并发
通过查询运行中任务数防止用户重复点击。虽然有 TOCTOU 竞态但低频场景够用。
6. strategySnapshot 固化到任务(核心设计)
事件溯源思想,任务一旦创建用什么方案就确定,不受 plan 后续修改影响。换来确定性和可追溯性。
7. Kafka 瘦消息:只传 ID
消息只是通知不是数据。所有数据以数据库为准,避免不一致,保持消息无状态。
8. Kafka 同步发送 .get() 保证事务一致
防止"任务创建了但消息没投递成功"的不一致状态。失败抛异常让 Spring 事务回滚。
9. 文档 indexStatus 立即切到 BUILDING
前端实时反馈,避免用户以为系统没响应。状态机字段独立,多阶段并行追踪。
10. 任务日志 metadata 带 strategySnapshot
日志是事件快照,保留重现现场所需的所有关键信息。
11. Topic 环境前缀 + Key 用 documentId
环境隔离防止串数据,同 documentId 进同一分区保证顺序消费。
12. Consumer 极简:反序列化 + 路由
业务逻辑在 Service,Consumer 只负责消息层入口。便于测试和扩展。
13. Consumer catch Exception 不重抛
消息消费成功 ≠ 业务执行成功。失败由任务状态记录,避免无意义重试。
十四、面试官可能会问的问题
问题 1:为什么"构建索引"要做成异步?同步不行吗?
可以回答:
不行。索引构建涉及切块、向量化(每个 chunk 调一次 LLM/embedding 模型)、写多个存储(PG/ES/Neo4j),单次耗时可能从几秒到几分钟。同步执行会让 HTTP 连接长期挂着、网关超时、用户体验极差。而且受 LLM 服务抖动影响接口性能不稳定,并发场景下还会把外部服务打挂。异步化之后接口永远毫秒级返回,Kafka 既是异步通道也是流量缓冲,失败可基于消费机制自然重试。这是典型的"轻同步 + 重异步"模式。
问题 2:为什么要把 strategySnapshot 冗余存一份到任务表?
可以回答:
这是事件溯源思想。任务记录代表"已发生的事件",事件应该不可变。从任务创建到 Consumer 消费有时间差,这段时间方案可能被改。如果 Consumer 实时查 plan 表,会拿到改后的方案,但用户当时确认的是旧方案,执行结果跟用户预期不一致。把 snapshot 固化到任务上,Consumer 直接读 task.strategySnapshot 就能用当时确认的版本执行。这种思想在订单(商品价格冗余)、合同(条款冗余)系统中也很常见,用冗余存储换确定性和可追溯性。
问题 3:kafkaTemplate.send() 后面为什么要加 .get()?
可以回答:
默认 send() 是异步的,调用线程不等 broker 确认就返回。这有风险:应用以为发出去了,实际可能因为网络问题或 broker 不可达没发出去;任务记录已经 insert 但消息没投递,Consumer 永远不会消费,任务卡在 NEW 状态。加 .get() 变成同步等待,发送失败会抛异常被 try-catch 捕获,外层事务回滚,task insert 和 document update 全部撤销,避免"任务存在但消息未投递"的不一致状态。代价是接口延迟变大,但对低频操作来说正确性远比性能重要。
问题 4:为什么 Kafka 消息只传三个 ID 不传完整数据?
可以回答:
因为 Kafka 消息的角色是"通知 Consumer 该干活了",不是"数据存储"。瘦消息的好处:消息小、Kafka 高效、Consumer 总能从数据库拿到最新数据避免不一致、消息无状态便于重试。如果塞胖数据:Kafka 压力变大、消息可能过期、和数据库可能不同步。所有数据以数据库为准,消息只做触发,这是工程上几乎一致的选择。
问题 5:为什么要校验 planId 一致而不是直接用 currentPlanId?
可以回答:
防止前端基于过期数据发请求。用户可能在 t1 时刻看到方案 A 并点了构建按钮,但 t1.5 时刻管理员把方案改成了 B。如果后端直接用 currentPlanId(也就是 B),用户期望构建 A 实际却用 B,结果跟预期不一致用户也不知道为什么。让前端显式传 planId 后端校验"传入的 == 当前的",发现不一致就拒绝并让前端刷新页面让用户重新确认,避免静默执行。这是乐观锁思想,把语义错误显式暴露而不是默默执行。
问题 6:防重复提交的查询有竞态条件吗?怎么处理?
可以回答:
有 TOCTOU 竞态。两个并发请求同时查询都发现没有运行中任务,然后都 insert,最终创建两个任务。要绝对防并发可以用:数据库唯一索引(documentId + 状态字段)、分布式锁(Redis SETNX)、乐观锁(document.indexVersion)。但实际场景中用户毫秒级双击概率低,而且 Consumer 端可以做幂等保护——查 task 状态如果已经 SUCCESS 或 RUNNING 就跳过。工程上选了"业务级校验 + Consumer 幂等"的组合,足够好且实现简单。完美防并发不是必须的,够用且成本低比绝对正确更重要。
问题 7:Consumer 里 catch Exception 不重抛有什么后果?
可以回答:
后果是消息会被 commit offset,不会自动重试。这是有意为之的设计:业务异常已经在 handleIndexBuild 内部 try-catch 里把任务标记为 FAILED 了,Kafka 层再重试一次也是失败;反序列化失败的消息再消费多少次都还是错的。所以这里采用"消息消费成功 ≠ 业务执行成功"的设计,Kafka 只保证消息不丢,业务成败由数据库 taskStatus 字段记录。如果要更严格的重试可以让 Kafka 自动重投或发到 DLQ,看团队策略。
企业级项目导航:⬅️ 06-异步索引构建:落库、向量化与收尾 | 07-构建索引 | ➡️ 08-索引构建链路白话讲解
💬 评论