--- title: "01-索引构建入口与Kafka消息投递" created: 2026-05-18 aliases: - 索引构建入口与Kafka消息投递 tags: - 项目 --- # 索引构建入口与Kafka消息投递 这篇文档讲的是:用户在页面上点击"构建索引"之后,后端同步做了哪些事情?从 Controller 接到请求开始,一直到把消息丢进 Kafka 队列,整条同步链路我们一步步拆开来看。 ![[FlyQLV3MwDgfg2dnm2ovIia_F6Wf-2db85ac9.png]] 先上一张总览流程图,有个整体印象之后再逐段看源码。 ## 同步链路总览 ![[FrQBLr7GLauAsQ044tumWh3OlItw-9e0aac8a.png]] ## Controller 层:接收构建请求 入口非常简单,就是一个标准的 POST 接口,接收前端传过来的 `DocumentIndexBuildDto`,然后直接委托给 Service 层处理。 ```java //DocumentManageController.java @Operation(summary = "执行文档索引构建") @PostMapping("/index/build") public ApiResponse buildIndex(@Valid @RequestBody DocumentIndexBuildDto dto) { return ApiResponse.ok(documentManageService.buildIndex(dto)); } ``` 请求参数也很简洁,就三个字段: ```java //DocumentIndexBuildDto.java @Data public class DocumentIndexBuildDto { @NotNull(message = "文档id不能为空") private Long documentId; @NotNull(message = "方案id不能为空") private Long planId; private Long operatorId; } ``` - `documentId`:要构建索引的文档 ID - `planId`:当前生效的策略方案 ID - `operatorId`:操作人 ID,可以为空(为空表示系统自动触发) ## Service 层:校验 + 建任务 + 投递消息 这是整个同步链路的核心,代码虽然不短,但逻辑很清晰——先做一堆校验,然后创建任务记录,最后把消息丢给 Kafka。 ### 文档状态校验 **DocumentManageServiceImpl.java — buildIndex()** ```java // 先拿到文档主记录;如果文档不存在或已被逻辑删除,会在这里直接抛错拦截。 SuperAgentDocument document = getDocumentOrThrow(dto.getDocumentId()); // 索引构建必须建立在"文档解析完成 + 策略已经人工确认"这两个前提之上, // 否则异步链路无法拿到稳定的结构化结果和最终生效的处理策略。 if (!Objects.equals(document.getParseStatus(), DocumentParseStatusEnum.PARSE_SUCCESS.getCode()) || !Objects.equals(document.getStrategyStatus(), DocumentStrategyStatusEnum.CONFIRMED.getCode())) { throw new SuperAgentFrameException(DocumentManageCode.DOCUMENT_STATUS_INVALID.getCode(), "当前文档尚未完成“解析成功 + 策略确认”,不能构建索引。"); } ``` 这段代码做了两件事: - 先确认文档存在 - 再确认文档已经走完了"解析成功"和"策略确认"两个前置步骤 > 为什么要同时校验这两个状态? > > 索引构建依赖解析阶段产出的纯文本,也依赖策略方案里定义的切块规则。如果解析没完成,就没有文本可切;如果策略没确认,就不知道该怎么切。两个条件缺一不可。 ### 方案一致性校验 **DocumentManageServiceImpl.java — buildIndex()** ```java // 请求里传入的 planId 必须和文档当前生效方案一致, // 避免前端缓存旧方案、重复点击,导致索引基于过期策略生成。 if (!Objects.equals(document.getCurrentPlanId(), dto.getPlanId())) { throw new SuperAgentFrameException(DocumentManageCode.STRATEGY_PLAN_NOT_FOUND.getCode(), "当前文档的生效方案与请求方案不一致。"); } ``` 这个校验防的是一种很常见的场景:用户打开了页面,但方案后来被别人改过了,这时候如果还用旧的 planId 去构建索引,产出的结果就和预期不一致了。 ### 防重复提交 **DocumentManageServiceImpl.java — buildIndex()** ```java // 查询当前文档是否已经存在"新建中/执行中"的索引任务, // 用于防止重复点击构建按钮后生成多个并发任务,造成重复切块和重复写索引。 long runningTaskCount = taskMapper.selectCount(new LambdaQueryWrapper() .eq(SuperAgentDocumentTask::getDocumentId, dto.getDocumentId()) .eq(SuperAgentDocumentTask::getTaskType, DocumentTaskTypeEnum.BUILD_INDEX.getCode()) .in(SuperAgentDocumentTask::getTaskStatus, DocumentTaskStatusEnum.NEW.getCode(), DocumentTaskStatusEnum.RUNNING.getCode()) .eq(SuperAgentDocumentTask::getStatus, BusinessStatus.YES.getCode())); if (runningTaskCount > 0) { throw new SuperAgentFrameException(DocumentManageCode.INDEX_TASK_RUNNING.getCode(), DocumentManageCode.INDEX_TASK_RUNNING.getMsg()); } ``` 直接查数据库看有没有正在跑的索引任务,有的话就拒绝。简单粗暴但很有效。 ### 创建索引任务 **DocumentManageServiceImpl.java — buildIndex()** ```java // 重新从数据库读取方案快照,而不是完全信任请求参数, // 确保后面写入任务的 strategySnapshot 来自一条真实存在且有效的方案记录。 SuperAgentDocumentStrategyPlan plan = planMapper.selectById(dto.getPlanId()); if (plan == null || !Objects.equals(plan.getStatus(), BusinessStatus.YES.getCode())) { throw new SuperAgentFrameException(DocumentManageCode.STRATEGY_PLAN_NOT_FOUND.getCode(), DocumentManageCode.STRATEGY_PLAN_NOT_FOUND.getMsg()); } // 为本次索引构建生成独立任务记录;后续异步链路会围绕这个 taskId 更新阶段、日志和结果。 Long taskId = uidGenerator.getUid(); SuperAgentDocumentTask task = new SuperAgentDocumentTask(); task.setId(taskId); task.setDocumentId(document.getId()); task.setPlanId(dto.getPlanId()); task.setTaskType(DocumentTaskTypeEnum.BUILD_INDEX.getCode()); task.setTaskStatus(DocumentTaskStatusEnum.NEW.getCode()); // 索引链路的首个执行阶段从 chunk 开始,后续消费者会继续推进阶段流转。 task.setCurrentStage(DocumentTaskStageEnum.CHUNK_EXECUTE.getCode()); // 操作人 ID 允许为空:为空时表示系统触发,否则表示用户主动触发。 Long operatorId = parseOptionalLong(dto.getOperatorId()); task.setTriggerSource(resolveTriggerSource(operatorId)); // 将策略快照固化到任务上,保证异步执行时即使方案后续被修改,也仍然以当时确认的版本执行。 task.setStrategySnapshot(plan.getStrategySnapshot()); task.setRetryCount(0); task.setStatus(BusinessStatus.YES.getCode()); taskMapper.insert(task); ``` 这里有个很关键的设计:**策略快照会被固化到任务记录上**。为什么要这么做?因为从任务创建到 Kafka 消费端真正执行,中间可能隔了几秒甚至几分钟,这段时间里方案有可能被修改。把快照存下来,就能保证异步执行时用的一定是当时确认的那个版本。 ### 更新状态 + 投递 Kafka 消息 **DocumentManageServiceImpl.java — buildIndex()** ```java // 任务创建成功后,立即把文档索引状态切换成"构建中", // 这样列表页和详情页能第一时间反映当前文档正处于索引处理中。 document.setIndexStatus(DocumentIndexStatusEnum.BUILDING.getCode()); documentMapper.updateById(document); // 记录一条启动日志 taskLogService.saveLog(taskId, document.getId(), DocumentTaskStageEnum.CHUNK_EXECUTE.getCode(), DocumentTaskEventTypeEnum.START.getCode(), DocumentLogLevelEnum.INFO.getCode(), resolveOperatorType(operatorId), operatorId, "索引构建任务已创建,等待异步执行。", Map.of("planId", dto.getPlanId(), "strategySnapshot", plan.getStrategySnapshot())); // 这里只负责投递消息,不同步执行业务重逻辑,避免接口阻塞太久; // 真正的索引构建会由 MQ 消费端接手。 kafkaProducer.sendIndexBuild(new DocumentIndexBuildMessage(document.getId(), taskId, dto.getPlanId())); ``` 最后返回给前端一个即时结果,告诉它"任务已受理": **DocumentManageServiceImpl.java — buildIndex()** ```text return new DocumentIndexBuildVo( document.getId(), taskId, task.getTaskType(), enumMsg(DocumentTaskTypeEnum.getRc(task.getTaskType())), task.getTaskStatus(), enumMsg(DocumentTaskStatusEnum.getRc(task.getTaskStatus())), document.getIndexStatus(), enumMsg(DocumentIndexStatusEnum.getRc(document.getIndexStatus())) ); ``` ## Kafka 消息投递细节 消息体很简单,就三个 ID: **DocumentIndexBuildMessage.java** ```java @Data @NoArgsConstructor @AllArgsConstructor public class DocumentIndexBuildMessage { private Long documentId; private Long taskId; private Long planId; } ``` 发送方法也很直接,序列化成 JSON 然后同步发送: **DocumentKafkaProducer.java** ```java public void sendIndexBuild(DocumentIndexBuildMessage message) { send(SpringUtil.getPrefixDistinctionName() + "-" + properties.getKafka().getIndexTopic(), String.valueOf(message.getDocumentId()), message); } private void send(String topic, String key, Object message) { try { String payload = objectMapper.writeValueAsString(message); kafkaTemplate.send(topic, key, payload).get(); } catch (Exception exception) { throw new SuperAgentFrameException(DocumentManageCode.KAFKA_SEND_FAILED.getCode(), "Kafka 消息发送失败: " + exception.getMessage(), exception); } } ``` > 注意这里用了 .get() > > `kafkaTemplate.send()` 返回的是 `Future`,调用 `.get()` 会阻塞等待发送结果。这样做的好处是:如果 Kafka 发送失败,异常会直接抛出来,整个事务会回滚,不会出现"任务创建了但消息没发出去"的不一致状态。 ## 消费端入口 Kafka 消费端的代码非常简洁,就是反序列化消息然后调用异步处理服务: **DocumentKafkaConsumer.java** ```java @KafkaListener( topics = SPRING_INJECT_PREFIX_DISTINCTION_NAME + "-" + "${app.manage.kafka.index-topic}", groupId = "${app.manage.kafka.group-id}-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); } } ``` 真正的重头戏在 `handleIndexBuild()` 方法里,这个我们下一篇详细讲。 --- **企业级项目导航**:⬅️ [[12-解析结果统计与异步收尾|12-解析结果统计与异步收尾]] | 01-索引构建入口与Kafka消息投递 | ➡️ [[02-初步父子切块|02-初步父子切块]]