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