索引构建入口与Kafka消息投递

这篇文档讲的是:用户在页面上点击"构建索引"之后,后端同步做了哪些事情?从 Controller 接到请求开始,一直到把消息丢进 Kafka 队列,整条同步链路我们一步步拆开来看。

FlyQLV3MwDgfg2dnm2ovIia_F6Wf-2db85ac9

先上一张总览流程图,有个整体印象之后再逐段看源码。

同步链路总览

FrQBLr7GLauAsQ044tumWh3OlItw-9e0aac8a

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:要构建索引的文档 ID
  • planId:当前生效的策略方案 ID
  • operatorId:操作人 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-初步父子切块