构建索引

我们接着上一篇“策略推荐与方案持久化”往下走。

上一篇的终点是:

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;
}

为什么消息体这么瘦?为什么不把 strategySnapshotparseTextPath 这些都塞进去?

原因 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-索引构建链路白话讲解