索引构建入口与Kafka消息投递
这篇文档讲的是:用户在页面上点击"构建索引"之后,后端同步做了哪些事情?从 Controller 接到请求开始,一直到把消息丢进 Kafka 队列,整条同步链路我们一步步拆开来看。
先上一张总览流程图,有个整体印象之后再逐段看源码。
同步链路总览
Controller 层:接收构建请求
入口非常简单,就是一个标准的 POST 接口,接收前端传过来的 DocumentIndexBuildDto,然后直接委托给 Service 层处理。
//DocumentManageController.java
@Operation(summary = "执行文档索引构建")
@PostMapping("/index/build")
public ApiResponse<DocumentIndexBuildVo> buildIndex(@Valid @RequestBody DocumentIndexBuildDto dto) {
return ApiResponse.ok(documentManageService.buildIndex(dto));
}
请求参数也很简洁,就三个字段:
//DocumentIndexBuildDto.java
@Data
public class DocumentIndexBuildDto {
@NotNull(message = "文档id不能为空")
private Long documentId;
@NotNull(message = "方案id不能为空")
private Long planId;
private Long operatorId;
}
documentId:要构建索引的文档 IDplanId:当前生效的策略方案 IDoperatorId:操作人 ID,可以为空(为空表示系统自动触发)
Service 层:校验 + 建任务 + 投递消息
这是整个同步链路的核心,代码虽然不短,但逻辑很清晰——先做一堆校验,然后创建任务记录,最后把消息丢给 Kafka。
文档状态校验
DocumentManageServiceImpl.java — buildIndex()
// 先拿到文档主记录;如果文档不存在或已被逻辑删除,会在这里直接抛错拦截。
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()
// 请求里传入的 planId 必须和文档当前生效方案一致,
// 避免前端缓存旧方案、重复点击,导致索引基于过期策略生成。
if (!Objects.equals(document.getCurrentPlanId(), dto.getPlanId())) {
throw new SuperAgentFrameException(DocumentManageCode.STRATEGY_PLAN_NOT_FOUND.getCode(),
"当前文档的生效方案与请求方案不一致。");
}
这个校验防的是一种很常见的场景:用户打开了页面,但方案后来被别人改过了,这时候如果还用旧的 planId 去构建索引,产出的结果就和预期不一致了。
防重复提交
DocumentManageServiceImpl.java — buildIndex()
// 查询当前文档是否已经存在"新建中/执行中"的索引任务,
// 用于防止重复点击构建按钮后生成多个并发任务,造成重复切块和重复写索引。
long runningTaskCount = taskMapper.selectCount(new LambdaQueryWrapper<SuperAgentDocumentTask>()
.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()
// 重新从数据库读取方案快照,而不是完全信任请求参数,
// 确保后面写入任务的 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()
// 任务创建成功后,立即把文档索引状态切换成"构建中",
// 这样列表页和详情页能第一时间反映当前文档正处于索引处理中。
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()
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
@Data
@NoArgsConstructor
@AllArgsConstructor
public class DocumentIndexBuildMessage {
private Long documentId;
private Long taskId;
private Long planId;
}
发送方法也很直接,序列化成 JSON 然后同步发送:
DocumentKafkaProducer.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
@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-解析结果统计与异步收尾 | 01-索引构建入口与Kafka消息投递 | ➡️ 02-初步父子切块
💬 评论