异步解析入口
我们接着上一篇“上传接口与文档主记录创建”往后学。上一篇的终点是:
上传接口完成文件校验、MinIO 存储、文档主记录入库、任务记录入库、任务日志记录
然后事务提交后发送 Kafka 消息
这篇的起点就是:
Kafka 消费者收到解析路由消息
然后进入异步解析链路
也就是说,上一节是同步上传入口,这一节是异步解析入口。
一、先建立整体链路认知
上一篇上传接口做完后,数据库里已经有了两类关键记录:
document 表:记录这份文档的基本信息、存储位置、解析状态、策略状态、索引状态
task 表:记录这次上传后产生的解析任务
同时,Kafka 里多了一条消息:
{ "documentId": 123, "taskId": 456}
这条消息的意义是:
这份文档已经上传完成并落库,可以开始异步解析了。
所以这篇要学的是:消费方拿到 documentId 和 taskId 后,如何把原始文件从 MinIO 下载下来,并用 Apache Tika 提取文本、清洗文本,为后续结构识别和切块策略推荐做准备。
完整流程可以先记成一句话:
Kafka 消费者收到解析路由消息后,反序列化出 documentId 和 taskId,调用异步处理服务加载文档和任务上下文,将任务状态推进到内容解析阶段,然后从 MinIO 下载原始文件,交给 Tika 解析器提取纯文本,并对文本做清洗、段落统计、标题识别、结构质量评估和内容质量评估,最终生成标准化的 DocumentAnalysisResult,供后续结构节点生成和切块策略推荐使用。
这一节先重点掌握前半段:
消费 Kafka 消息
→ 加载 document/task 上下文
→ 更新任务状态
→ 记录开始日志
→ 从 MinIO 下载原始文件
→ Tika 提取文本
→ 文本清洗
→ 初步统计分析
二、这篇在整个项目中的位置
你的知识库文档处理链路可以拆成几个阶段:
上一篇学到 E,这篇从 F 开始,重点讲到 I,并简单接触 J 前的准备。
这部分很重要,因为后续所有高级能力都依赖这里提取出来的文本质量。
比如:
结构化切块依赖标题和章节识别
语义切块依赖干净的段落文本
向量化依赖稳定的纯文本
Neo4j 图谱构建依赖结构节点
RAG 引用来源依赖原文片段
文档摘要和画像依赖解析文本
所以,文档解析不是简单把 PDF 转成字符串,而是整个知识库构建质量的地基。
三、Kafka Consumer:异步链路的消息入口
文档里的消费者入口是:
@KafkaListener(
topics = SPRING_INJECT_PREFIX_DISTINCTION_NAME + "-" + "",
groupId = "-parse")
public void consumeParseRoute(String payload) {
try {
DocumentParseRouteMessage message = objectMapper.readValue(
payload,
DocumentParseRouteMessage.class
);
asyncProcessService.handleParseRoute(
message.getDocumentId(),
message.getTaskId()
);
}
catch (Exception exception) {
log.error("消费解析路由消息失败,payload={}", payload, exception);
}
}
这个方法很短,但里面有几个很关键的设计点。
四、Consumer 只做“反序列化 + 转发”,不做业务
这里 Consumer 的职责非常克制:
收到 Kafka 原始字符串 payload
反序列化成 DocumentParseRouteMessage
取出 documentId 和 taskId
调用 asyncProcessService.handleParseRoute()
它没有在监听方法里直接写复杂的解析逻辑。
这是一个很好的分层设计。
DocumentKafkaConsumer
只负责消息接入
DocumentAsyncProcessService
负责异步业务推进
TikaDocumentParserService
负责文档解析
StorageService
负责对象存储下载
TaskLogService
负责任务日志
这样做的好处是:
第一,Kafka Consumer 更稳定。监听方法越简单,越不容易出问题。
第二,业务逻辑更容易测试。handleParseRoute(documentId, taskId) 可以脱离 Kafka 单独做单元测试或集成测试。
第三,后续如果消息入口换成别的方式,比如定时补偿任务、手动重试按钮、HTTP 管理接口,也可以复用同一个 handleParseRoute()。
比如未来可以有:
@PostMapping("/task/retry/{taskId}")
public void retry(@PathVariable Long taskId) {
asyncProcessService.handleParseRoute(documentId, taskId);
}
也就是说,Kafka 只是触发方式之一,真正的业务入口被抽象到了 Service。
面试里可以这样讲:
Kafka Consumer 层保持轻量,只做消息反序列化和业务转发,不在监听器里堆复杂逻辑。真正的解析流程下沉到异步处理服务,这样既降低了监听线程出错风险,也方便后续定时补偿、人工重试和单元测试复用。
五、topic 和 groupId 的意义
文档里提到:
topic 名称是:环境前缀 + parseTopic
groupId 带 -parse 后缀
这个设计和上一篇发送消息时是对应的。
上一篇生产者发送 Kafka 时,topic 是类似:
dev-document-parse-topic
test-document-parse-topic
prod-document-parse-topic
这样不同环境之间不会串消息。
如果不加环境前缀,可能会出现非常危险的问题:
测试环境上传的文档,被生产环境消费者消费
生产环境的任务,被开发环境误消费
所以环境隔离非常重要。
groupId = "-parse" 的作用是区分消费组。
Kafka 的消费组有一个关键规则:
同一个 topic 下,同一个 consumer group 内,一条消息通常只会被一个消费者实例处理;
不同 consumer group 之间,彼此独立消费。
也就是说,如果你有 3 个解析服务实例,它们都属于同一个 parse 消费组,那么同一条文档解析消息只会被其中一个实例处理。
这样可以实现横向扩展:
解析服务实例 1:处理文档 A
解析服务实例 2:处理文档 B
解析服务实例 3:处理文档 C
但如果另一个服务也想关心这条消息,比如审计服务,它可以使用另一个消费组:
parse-group:负责解析文档
audit-group:负责记录审计
stat-group:负责统计上传事件
同一条消息就可以被多个业务独立消费。
六、Consumer 中 try-catch 的意义和隐患
文档里 Consumer 用了:
try {
...
}
catch (Exception exception) {
log.error("消费解析路由消息失败,payload={}", payload, exception);
}
文档的解释是:
消费失败只记录日志,不让异常继续向外冒泡破坏监听线程。
这个思路可以理解:避免一条坏消息导致监听线程异常。
但是这里要进一步学习一个更深的点:捕获异常但不抛出,可能会让 Kafka 框架认为这条消息消费成功。
这取决于 Kafka 的 ack 配置和 Spring Kafka 的容器配置。
如果是自动提交 offset,或者监听方法正常返回后提交 offset,那么这里 catch 住异常不再抛出,就可能导致:
消息实际处理失败
但是 offset 已提交
后续不会自动重试
所以这种写法适合什么场景?
适合消费失败后由业务表补偿兜底的设计。
比如这个项目里有:
document 表
task 表
task_log 表
retryCount
parseStatus
taskStatus
如果消费失败只打日志,那么最好还要有补偿机制:
定时扫描 taskStatus = NEW 或 RUNNING 但长时间未推进的任务
重新投递 Kafka
或者直接调用 handleParseRoute 重试
否则一旦消费失败,就可能出现任务卡死。
更成熟的做法有几种:
方式一:catch 后更新任务状态为 FAILED,再记录失败日志
方式二:catch 后抛出异常,让 Kafka 错误处理器重试
方式三:进入重试 topic 或死信 topic
方式四:依靠任务表定时补偿
所以这段代码你不能只背“try-catch 防止监听线程挂掉”,还要知道它背后的代价。
面试里可以这样讲:
Consumer 里捕获异常可以避免异常直接打断监听线程,但如果捕获后不抛出,需要配套任务表补偿或失败状态更新,否则可能出现 offset 已提交但业务未成功的情况。更稳的设计是结合 Spring Kafka 的错误处理器、重试 topic、死信队列,或者通过 task 表定时扫描长时间未推进的任务进行补偿。
七、进入 handleParseRoute:真正的异步解析主流程
Consumer 只负责转发,真正的业务在:
asyncProcessService.handleParseRoute(documentId, taskId);
文档中的方法说明很完整,它大致做这些事:
1. 读取文档记录和任务记录,确认上下文存在
2. 把任务状态切到 RUNNING,把当前阶段推进到 CONTENT_PARSE
3. 从对象存储下载原始文件
4. 调用解析器提取纯文本和结构信息
5. 把解析后的纯文本重新上传为 txt
6. 替换文档结构节点,同步导航索引、图投影和画像
7. 基于解析结果生成推荐切块方案
8. 写入推荐方案和步骤,更新文档状态、策略状态、统计信息
9. 成功或失败收尾任务,并记录任务日志
这篇文档主要讲到第 4 步,也就是文本提取和清洗。
我们先学习前几步。
八、第一步:加载 document 和 task 上下文
代码是:
SuperAgentDocument document = documentMapper.selectById(documentId);
SuperAgentDocumentTask task = taskMapper.selectById(taskId);
if (document == null || task == null) {
log.warn("解析任务对应的文档或任务不存在,documentId={}, taskId={}", documentId, taskId);
return;
}
Kafka 消息里只有两个 ID:
documentId
taskId
消费方拿到 ID 后,必须回数据库查完整上下文。
这也对应上一篇提到的设计原则:
Kafka 消息只传轻量 ID
完整业务上下文以数据库为准
为什么不把完整文档信息都放到 Kafka 消息里?
因为文档信息可能很多,而且会变化。比如:
documentName
objectName
mimeType
fileType
parseStatus
strategyStatus
operatorId
knowledgeScope
如果都放消息里,消息体变大,还会带来数据一致性问题。
更好的方式是:
消息负责触发
数据库负责事实状态
所以这里先查 document 和 task。
如果查不到,直接 return。这个情况可能来自:
上传事务实际失败,但异常情况下消息被发出
文档或任务被人工删除
消息重复或历史脏消息
数据库数据异常
正常情况下,上一篇已经通过 TransactionTemplate 保证了“事务提交后再发 Kafka”,所以这里查不到的概率应该很低。
九、第二步:推进任务状态
查到上下文后,进入解析流程之前要先更新任务状态:
Date startTime = new Date();
task.setTaskStatus(DocumentTaskStatusEnum.RUNNING.getCode());
task.setCurrentStage(DocumentTaskStageEnum.CONTENT_PARSE.getCode());
task.setStartTime(startTime);
taskMapper.updateById(task);
document.setParseStatus(DocumentParseStatusEnum.PARSING.getCode());
documentMapper.updateById(document);
这里有两个对象的状态要更新。
第一个是任务状态:
taskStatus: NEW → RUNNING
currentStage: FILE_UPLOAD → CONTENT_PARSE
startTime: 当前时间
这表示这条任务已经被消费方领取,并进入内容解析阶段。
第二个是文档状态:
parseStatus: PARSING
这表示文档整体处于解析中。
你要理解 document 和 task 的职责差异:
document 表关注文档生命周期状态
task 表关注某一次异步任务的执行状态
文档是资产,任务是动作。
比如同一份文档未来可能有多个任务:
第一次上传后的解析任务
重新解析任务
重新索引任务
摘要生成任务
图谱重建任务
所以任务状态和文档状态不要混在一起。
十、这里可以引出状态机思想
异步任务最怕状态混乱,所以通常要设计状态机。
这篇里的任务状态大概可以理解成:
stateDiagram-v2
[*] --> NEW
NEW --> RUNNING
RUNNING --> SUCCESS
RUNNING --> FAILED
FAILED --> NEW: retry
SUCCESS --> [*]
阶段状态则是:
stateDiagram-v2
[*] --> FILE_UPLOAD
FILE_UPLOAD --> CONTENT_PARSE
CONTENT_PARSE --> STRUCTURE_BUILD
STRUCTURE_BUILD --> STRATEGY_RECOMMEND
STRATEGY_RECOMMEND --> WAIT_CONFIRM
文档主状态则可能是:
parseStatus:
WAIT_PARSE / PARSING / PARSE_SUCCESS / PARSE_FAILED
strategyStatus:
WAIT_RECOMMEND / RECOMMENDING / WAIT_CONFIRM / CONFIRMED
indexStatus:
WAIT_BUILD / BUILDING / BUILT / BUILD_FAILED
这就是为什么异步系统里一定要有任务表和状态字段。
如果没有状态字段,用户问:
我上传的文档现在处理到哪了?
系统没法回答。
如果有状态字段和任务日志,系统就可以展示:
文件已上传
内容解析中
结构识别完成
切块策略已推荐
等待用户确认
索引构建中
处理完成
十一、第三步:记录开始日志
代码是:
taskLogService.saveLog(
taskId,
documentId,
DocumentTaskStageEnum.CONTENT_PARSE.getCode(),
DocumentTaskEventTypeEnum.START.getCode(),
DocumentLogLevelEnum.INFO.getCode(),
DocumentOperatorTypeEnum.SYSTEM.getCode(),
null,
"开始解析文档内容。",
Map.of("objectName", document.getObjectName())
);
这条日志的作用是告诉系统:
这个任务已经进入内容解析阶段
并且准备从这个 objectName 下载原始文件
日志里记录了 objectName,非常重要。
因为后面如果下载失败,可以排查:
objectName 是否为空
MinIO 里是否存在这个对象
bucket 是否正确
权限是否正常
文件是否被误删
任务日志是异步链路的“黑匣子”。
你可以把每个阶段都记录成:
开始日志 START
成功日志 COMPLETE
失败日志 FAILED
关键上下文 metadata
这样后台排查问题时,不需要猜任务执行到哪一步,而是直接看时间线。
十二、第四步:从 MinIO 重新下载原始文件
代码是:
byte[] fileBytes = storageService.downloadObject(document.getObjectName());
这里有一个很多初学者容易疑惑的问题:
上传接口里不是已经读过 byte[] 了吗?为什么消费方还要重新下载?
原因是:上传接口和 Kafka 消费方不是同一个执行上下文。
上一篇上传阶段里:
MultipartFile → byte[] → 上传 MinIO
这个 byte[] 是 HTTP 请求线程里的内存数据。
请求结束后,内存就释放了。
Kafka 消费方可能在:
另一个线程
另一个服务实例
另一台机器
甚至几分钟后
执行。
所以它不能依赖上传接口里的内存对象,只能从持久化存储 MinIO 重新拿文件。
这也体现了 MinIO 的作用:
MinIO 是原始文件的持久化存储
Kafka 只是通知
MySQL 保存文件位置
消费方根据 objectName 从 MinIO 下载
完整关系是:
flowchart LR
A[Kafka 消息 documentId taskId] --> B[查询 MySQL document]
B --> C[拿到 objectName]
C --> D[从 MinIO 下载原始文件]
D --> E[Tika 解析]
这里也有一个设计好处:
消费方处理的是上传后实际持久化成功的那份文件
而不是某个临时内存对象。
十三、第五步:调用 parserService.parse()
下载到文件字节后,调用:
DocumentAnalysisResult analysisResult = parserService.parse(
fileBytes,
document.getOriginalFileName(),
document.getMimeType(),
DocumentFileTypeEnum.getRc(document.getFileType())
);
这里传了四个参数:
| 参数 | 作用 |
|---|---|
fileBytes |
原始文件二进制内容 |
originalFileName |
用于辅助识别文件类型和结构 |
mimeType |
用于判断内容类型 |
fileType |
系统内部枚举化后的文件类型 |
最终返回:
DocumentAnalysisResult
这个结果不是单纯一个字符串,而是一个标准化分析结果对象。
它包含:
cleanedText:清洗后的纯文本
charCount:字符数
tokenCount:估算 token 数
structureLevel:结构化程度
contentQualityLevel:内容质量等级
headingCount:标题数量
paragraphCount:段落数量
maxParagraphLength:最大段落长度
structureNodes:结构节点候选
这就说明解析服务不是只负责“提取文本”,还要做“初步分析”。
为什么要做这些统计?
因为后续切块策略推荐需要这些信息。
例如:
标题多、层级清晰:适合结构化切块
段落很长:需要递归切块兜底
文本很乱、结构差:可能需要 LLM 智能切块
内容太短:不适合复杂切块
token 数很大:需要更严格预算控制
十四、TikaDocumentParserService 的整体职责
解析器入口是:
public DocumentAnalysisResult parse(
byte[] bytes,
String originalFileName,
String mimeType,
DocumentFileTypeEnum fileType
)
文档里把它拆成六步:
1. 根据文件类型提取原始文本
2. 对文本做换行、空白字符和乱码清洗
3. 提取结构节点候选,并估算标题数量
4. 统计段落、字符数、token 数
5. 评估文档结构等级和内容质量等级
6. 打包成 DocumentAnalysisResult 返回
可以画成:
flowchart TD
A[原始文件 bytes] --> B[extractRawText 提取原始文本]
B --> C[cleanupText 清洗文本]
C --> D[structureNodeExtractor 提取结构节点候选]
C --> E[extractParagraphs 提取段落]
D --> F[countHeadings 统计标题数量]
E --> G[统计段落数和最大段落长度]
C --> H[统计 charCount 和 tokenCount]
F --> I[evaluateStructureLevel 结构等级]
H --> J[evaluateContentQuality 内容质量等级]
I --> K[DocumentAnalysisResult]
J --> K
D --> K
G --> K
这一点很关键:Tika 只是底层文本提取工具,真正的业务解析器还要在 Tika 结果上做清洗、统计和结构分析。
十五、extractRawText:按文件类型提取原始文本
代码核心是:
private String extractRawText(
byte[] bytes,
String originalFileName,
String mimeType,
DocumentFileTypeEnum fileType
) {
try {
if (fileType == DocumentFileTypeEnum.PDF
|| fileType == DocumentFileTypeEnum.DOC
|| fileType == DocumentFileTypeEnum.DOCX) {
return tika.parseToString(new ByteArrayInputStream(bytes));
}
if (fileType == DocumentFileTypeEnum.TXT
|| fileType == DocumentFileTypeEnum.MD) {
return new String(bytes, StandardCharsets.UTF_8);
}
if (fileType == DocumentFileTypeEnum.HTML) {
return tika.parseToString(new ByteArrayInputStream(bytes));
}
return tika.parseToString(new ByteArrayInputStream(bytes));
}
catch (Exception exception) {
if (fileType == DocumentFileTypeEnum.TXT
|| fileType == DocumentFileTypeEnum.MD
|| mimeType != null && mimeType.startsWith("text/")) {
return new String(bytes, StandardCharsets.UTF_8);
}
throw new IllegalStateException("Tika 解析失败: " + exception.getMessage(), exception);
}
}
这里可以分成三类处理。
1. PDF / Word / HTML 等复杂格式走 Tika
tika.parseToString(new ByteArrayInputStream(bytes))
Apache Tika 的作用是屏蔽底层文件格式差异。
对于 PDF、DOC、DOCX、HTML 等格式,文件内部不是简单文本,需要解析结构、编码、压缩内容、元数据等,所以交给 Tika。
Tika 底层能适配很多格式,比如:
PDF
Word
PPT
Excel
HTML
RTF
OpenDocument
纯文本
部分压缩包内文本
项目描述里说支持 100+ 格式,就是借助 Tika 的统一解析能力。
2. TXT / MD 直接 UTF-8 解码
return new String(bytes, StandardCharsets.UTF_8);
纯文本和 Markdown 本质就是文本文件,不需要复杂解析。
直接解码更快,也避免 Tika 对 Markdown 格式做一些不必要处理。
3. Tika 失败时对文本类文件降级
catch 里有一段兜底:
if (fileType == TXT || fileType == MD || mimeType.startsWith("text/")) {
return new String(bytes, StandardCharsets.UTF_8);
}
这个设计很好。
因为有些文本类文件,Tika 可能因为编码、MIME、格式判断问题解析失败,但文件内容本身还是可以直接读的。
此时直接 UTF-8 解码,可以提高系统容错性。
但如果是 PDF、DOCX 这类复杂格式,Tika 失败后直接用字符串解码通常没有意义,得到的可能是一堆乱码,所以直接抛异常。
十六、为什么不能所有文件都直接 UTF-8 解码?
这是一个常见问题。
TXT 和 MD 可以直接解码,是因为它们本身就是文本。
但 PDF、DOCX、PPTX、XLSX 不是纯文本文件。
比如 DOCX 本质上是一个压缩包,里面包含 XML 文件、样式、关系文件等。PDF 也有自己的对象结构、字体编码、页面布局。
如果你直接:
new String(bytes, StandardCharsets.UTF_8)
解析 PDF 或 DOCX,得到的通常是乱码、二进制符号或者无意义内容。
所以复杂格式必须用 Tika、PDFBox、POI 这类解析工具。
十七、cleanupText:文本清洗为什么重要?
Tika 提取出来的文本还不能直接用于切块和检索。
原因是原始文本经常会有这些问题:
Windows 换行和 Unix 换行混杂
连续多个空行
制表符、控制字符
空字符
多余空格
页面页眉页脚残留
PDF 换行异常
表格内容错位
文档里的清洗代码是:
private String cleanupText(String rawText) {
if (StrUtil.isBlank(rawText)) {
return "";
}
String cleaned = rawText
.replace("\r\n", "\n")
.replace('\r', '\n')
.replace('\u0000', ' ')
.replaceAll("[\t\\x0B\\f]+", " ")
.replaceAll("\n{3,}", "\n\n")
.replaceAll("[ ]{2,}", " ")
.trim();
return cleaned;
}
这段逻辑可以逐行理解。
十八、清洗规则逐条讲解
1. 统一换行符
.replace("\r\n", "\n")
.replace('\r', '\n')
不同系统的换行符不一样:
Windows: \r\n
Unix/Linux/macOS: \n
老 Mac: \r
如果不统一,后续按 \n 切段落时会不稳定。
统一成 \n 后,段落切分和标题识别就简单很多。
2. 空字符替换为空格
.replace('\u0000', ' ')
有些解析结果里会出现空字符 \u0000。
这类字符对展示、存储、ES 分词、向量化都没有意义,甚至可能导致一些库处理异常。
替换成空格比较安全。
3. 制表符和控制字符替换为空格
.replaceAll("[\t\\x0B\\f]+", " ")
这里处理的是:
\t:制表符
\x0B:垂直制表符
\f:换页符
这些符号在 PDF 或 Office 解析结果中很常见。
它们对语义帮助不大,统一替换成空格,能让文本更平滑。
4. 连续多个换行压缩成两个换行
.replaceAll("\n{3,}", "\n\n")
连续 3 个以上换行压缩成 2 个。
为什么不是压缩成 1 个?
因为两个换行通常表示段落边界。
例如:
第一段内容。
第二段内容。
如果全部压成一个换行,段落边界会弱化。
所以保留两个换行是为了兼顾清洁度和结构信息。
5. 连续多个空格压缩成一个空格
.replaceAll("[ ]{2,}", " ")
多余空格对检索、向量化和切块都没什么好处。
压缩后文本更稳定。
6. trim 去掉首尾空白
.trim()
避免文本开头和结尾有无效空白。
十九、文本清洗对后续模块的影响
这一段不要小看。
文本清洗质量会直接影响后续所有模块。
1. 影响标题识别
如果换行混乱,标题识别可能失败。
例如:
第一章 总则
第一条 xxx
如果中间有奇怪控制字符,标题提取器可能识别不出章节。
2. 影响段落切分
段落是语义切块和递归切块的重要基础。
连续空行不处理,会导致很多空段落;换行压缩过度,又会导致段落边界丢失。
3. 影响向量化效果
向量模型输入如果包含大量控制字符、乱码、无意义空白,会降低 embedding 表示质量。
4. 影响 ES 关键词检索
异常字符会影响分词和关键词匹配。
5. 影响 RAG 回答质量
如果召回的证据文本很脏,模型生成回答时也容易受到干扰。
所以这一层虽然看起来简单,但它是后续 RAG 质量的基础保障。
二十、结构节点候选提取:下一篇的重点,但这里先建立概念
parse 方法里接下来做了:
List<DocumentStructureNodeCandidate> structureNodes =
structureNodeExtractor.extract(originalFileName, cleanedText);
int headingCount = countHeadings(cleanedText, structureNodes);
这部分文档说下一篇会单独展开,但我们现在先理解它的意义。
所谓结构节点候选,就是从清洗后的文本中识别出类似这些节点:
第一章 总则
1. 项目背景
1.1 系统目标
1.1.1 技术架构
二、功能设计
(一)上传模块
这些节点后面会用于:
构建文档目录
生成 Neo4j 文档结构图谱
支持章节导航
辅助结构化切块
辅助策略推荐
这和你的项目亮点“基于 Neo4j 构建文档层级结构图谱”直接相关。
如果一篇文档结构清晰,系统就可以做:
文档节点 → 章节节点 → 条目节点
并支持:
章节编号定位
标题路径查找
前后章节浏览
子节点展开
语义最佳匹配
所以文本解析阶段不是孤立的,它在给图谱构建和导航能力准备基础数据。
二十一、段落统计:为策略推荐准备特征
parse 方法里还有:
List<String> paragraphList = extractParagraphs(cleanedText);
int maxParagraphLength = paragraphList.stream()
.mapToInt(String::length)
.max()
.orElse(0);
int charCount = cleanedText.length();
int tokenCount = estimateTokenCount(cleanedText);
这里统计了:
段落列表
最大段落长度
字符数
token 数
这些指标会影响后续策略推荐。
比如:
1. 字符数很少
如果文档很短,比如只有 300 字,就不需要复杂切块。
可能直接一个块就够了。
2. 段落很多且长度适中
适合按段落或结构化方式切块。
3. 最大段落特别长
说明文档中可能有大段连续文本,需要递归切块兜底。
4. token 数很大
说明后续要严格控制证据预算、切块大小和摘要压缩。
5. 标题数量多
说明文档结构清晰,适合结构化切块,并适合构建目录图谱。
二十二、tokenCount 为什么是估算值?
文档里写:
int tokenCount = estimateTokenCount(cleanedText);
这里是估算 token 数,不是精确 tokenizer。
为什么可以估算?
因为这个阶段主要是做策略判断和粗粒度统计,不需要像模型调用前那样精确。
比如系统只需要知道:
这个文档大概是 1k token、10k token,还是 100k token
用于判断:
是否需要切块
切块策略选哪种
是否需要摘要
是否可能超出模型上下文
如果每次都引入精确 tokenizer,成本会更高,性能也会下降。
常见估算方式可能是:
中文:按字符数粗略折算
英文:按单词数或字符数折算
混合文本:综合估算
在工程上,只要用于粗粒度判断,估算是合理的。
二十三、结构等级和内容质量等级
parse 方法最后还有:
int structureLevel = evaluateStructureLevel(headingCount, paragraphList.size());
int contentQualityLevel = evaluateContentQuality(cleanedText, charCount);
这两个字段非常有业务价值。
1. structureLevel:结构化程度
它可能根据:
标题数量
段落数量
标题密度
章节层级规律
判断文档结构是否清晰。
例如:
结构等级高:有明确章节、标题、编号
结构等级中:有段落,但标题不明显
结构等级低:大段纯文本或解析混乱
结构等级会影响切块策略:
结构等级高 → 优先结构化切块
结构等级中 → 递归切块 + 语义优化
结构等级低 → 可能使用 LLM 智能切块
2. contentQualityLevel:内容质量等级
它可能根据:
文本是否为空
字符数是否太少
乱码比例
可读字符比例
重复字符比例
控制字符比例
评估解析文本质量。
比如扫描版 PDF 如果没有 OCR,Tika 可能提取不到有效文本,内容质量就会很低。
内容质量低时,系统后续可能:
标记解析失败
提示用户文件质量较差
走 OCR 解析
走 LLM 修复
不推荐自动切块
所以这两个等级不是装饰字段,而是后续自动化决策的依据。
二十四、DocumentAnalysisResult 是解析阶段的标准输出
最后返回:
return new DocumentAnalysisResult(
cleanedText,
charCount,
tokenCount,
structureLevel,
contentQualityLevel,
headingCount,
paragraphList.size(),
maxParagraphLength,
structureNodes
);
这个对象可以理解为“文档解析报告”。
它把原始文件解析后的核心信息统一封装起来。
后续服务不需要关心:
这个文件是 PDF 还是 DOCX
Tika 怎么解析
换行怎么清洗
标题怎么识别
段落怎么统计
后续只需要面对统一的:
DocumentAnalysisResult
这是典型的抽象设计。
也就是说:
输入层:多格式文件
解析层:统一成 DocumentAnalysisResult
后续层:基于统一结果做结构节点、切块推荐、索引构建
这就是“多格式文档统一解析”的核心。
二十五、这一节和上一篇上传链路的衔接
现在我们把上一篇和这一篇连起来。
上一篇:
用户上传文件
Controller 接收 multipart
Service 校验文件
读取 byte[]
生成 documentId
上传 MinIO
构建 document 主记录
构建 task 任务记录
TransactionTemplate 提交入库
事务提交后发 Kafka
返回 documentId/taskId
这一篇:
Kafka Consumer 收到消息
反序列化 DocumentParseRouteMessage
调用 handleParseRoute
查询 document/task
任务状态 NEW → RUNNING
阶段 FILE_UPLOAD → CONTENT_PARSE
记录开始日志
从 MinIO 下载 objectName 对应的原始文件
调用 TikaDocumentParserService.parse()
按文件类型提取原始文本
清洗文本
提取结构节点候选
统计段落、字符、token
评估结构等级和内容质量等级
返回 DocumentAnalysisResult
可以整合成一个完整流程图:
flowchart TD
A[上传接口收到 multipart 文件] --> B[校验文件类型和内容]
B --> C[上传原始文件到 MinIO]
C --> D[写 document 主表]
D --> E[写 task 任务表]
E --> F[写 task_log 上传完成日志]
F --> G[事务提交]
G --> H[发送 Kafka 解析路由消息]
H --> I[Kafka Consumer 消费消息]
I --> J[反序列化 documentId taskId]
J --> K[查询 document 和 task]
K --> L[更新任务为 RUNNING CONTENT_PARSE]
L --> M[记录内容解析开始日志]
M --> N[从 MinIO 下载原始文件]
N --> O[Apache Tika 提取原始文本]
O --> P[文本清洗]
P --> Q[结构节点候选提取]
P --> R[段落/字符/token 统计]
Q --> S[生成 DocumentAnalysisResult]
R --> S
二十六、这一块的核心技术点
这一篇你要重点掌握这些技术点。
1. Kafka 消费者职责要轻
消费者监听方法只做:
接消息
反序列化
转发到业务服务
不要把复杂业务都写在 listener 方法里。
2. Kafka 消息只传 ID
消息里只需要:
documentId
taskId
完整上下文从数据库查。
这能降低消息体复杂度,也避免消息数据和数据库状态不一致。
3. 异步任务必须推进状态
消费开始后要更新:
taskStatus = RUNNING
currentStage = CONTENT_PARSE
parseStatus = PARSING
startTime = 当前时间
这对前端进度展示和后台排查都非常重要。
4. 原始文件必须从 MinIO 重新下载
异步消费方不能依赖上传接口内存里的 byte[]。
上传阶段保存文件,消费阶段通过 objectName 下载文件。
5. Tika 解决多格式解析问题
不同文档格式统一转成文本,是知识库构建的第一步。
Tika 是底层解析工具,但业务上还要做清洗和分析。
6. 文本清洗是 RAG 质量基础
换行、空白、控制字符、空段落这些问题如果不处理,会影响:
结构识别
切块
向量化
关键词检索
模型回答
7. DocumentAnalysisResult 是解析层统一输出
它统一承载:
清洗文本
统计信息
结构候选
质量评估
为后续结构图谱和切块策略推荐提供输入。
二十七、这块可以怎么写进简历
可以提炼成这样:
负责实现文档上传后的 Kafka 异步解析链路,Consumer 端接收解析路由消息后反序列化出 documentId 和 taskId,并转交异步处理服务加载文档与任务上下文。处理过程中将任务状态从 NEW 推进到 RUNNING,将当前阶段更新为 CONTENT_PARSE,并记录任务日志用于全链路追踪。随后基于文档主记录中的 objectName 从 MinIO 下载原始文件,调用 Apache Tika 按文件类型提取原始文本,并对换行、控制字符、多余空白和异常字符进行标准化清洗,生成包含纯文本、字符数、token 估算、段落统计、标题数量、结构等级和内容质量等级的 DocumentAnalysisResult,为后续结构节点生成、切块策略推荐和索引构建提供统一输入。
如果要突出架构性,可以加一句:
通过“Kafka 消息轻量触发 + MySQL 状态表追踪 + MinIO 原始文件持久化 + Tika 统一解析”的设计,实现了上传链路与解析链路解耦,并为多格式文档处理提供了稳定的异步处理基础。
二十八、面试官可能会问的问题
问题 1:Kafka 消费者为什么不直接处理所有业务?
可以回答:
Consumer 层保持轻量,只负责消息反序列化和转发,真正的业务流程放在 Service 层。这样可以降低监听线程复杂度,同时让解析流程可以被 Kafka 消费、定时补偿、人工重试等多种入口复用。
问题 2:为什么 Kafka 消息只传 documentId 和 taskId?
可以回答:
因为完整上下文已经在数据库中,Kafka 消息只需要做轻量触发。这样消息体更小,也避免消息中的冗余字段和数据库状态不一致。消费方拿到 ID 后再查 document 和 task,数据库作为事实来源。
问题 3:为什么消费方还要从 MinIO 下载文件?上传时不是已经有 byte[] 吗?
可以回答:
上传时的 byte[] 只存在于 HTTP 请求线程内存中,请求结束后就释放了。Kafka 消费方可能运行在另一个线程、进程甚至机器上,所以必须从 MinIO 重新下载已经持久化的原始文件。这样也能保证消费方处理的是上传成功后的正式文件。
问题 4:为什么用 Apache Tika?
可以回答:
因为项目要支持 PDF、Word、HTML、Markdown、TXT 等多种格式,Tika 可以屏蔽底层文件格式差异,把复杂文档统一提取成文本。对于 TXT 和 MD 这类纯文本文件,则可以直接 UTF-8 解码,提高效率;如果 Tika 解析文本类文件失败,也可以降级为直接解码。
问题 5:文本清洗为什么重要?
可以回答:
Tika 提取出的文本通常存在换行混乱、控制字符、多余空格、空字符等问题。如果不清洗,会影响标题识别、段落切分、切块策略、向量化质量和关键词检索效果。因此解析后要统一换行符、压缩空行和空格、去除控制字符,为后续结构识别和 RAG 检索提供稳定输入。
问题 6:Consumer catch 异常不抛出有没有问题?
可以回答:
有潜在问题。如果 catch 后不抛出,Kafka 框架可能认为消息已经消费成功并提交 offset,导致业务失败但消息不再重试。因此需要配套任务表补偿、失败状态记录、重试机制或死信队列。更稳妥的方案是结合 Spring Kafka 的错误处理器和业务幂等机制处理失败消息。
二十九、这一节你需要背下来的主线
可以用这段伪代码记忆:
@KafkaListener(...)
public void consumeParseRoute(String payload) {
DocumentParseRouteMessage message = objectMapper.readValue(payload, DocumentParseRouteMessage.class);
asyncProcessService.handleParseRoute(message.getDocumentId(), message.getTaskId());
}
业务主流程:
public void handleParseRoute(Long documentId, Long taskId) {
// 1. 加载上下文
SuperAgentDocument document = documentMapper.selectById(documentId);
SuperAgentDocumentTask task = taskMapper.selectById(taskId);
if (document == null || task == null) {
return;
}
// 2. 推进任务状态
task.setTaskStatus(RUNNING);
task.setCurrentStage(CONTENT_PARSE);
task.setStartTime(new Date());
taskMapper.updateById(task);
document.setParseStatus(PARSING);
documentMapper.updateById(document);
// 3. 记录开始日志
taskLogService.saveLog(..., "开始解析文档内容", objectName);
// 4. 下载原始文件
byte[] fileBytes = storageService.downloadObject(document.getObjectName());
// 5. Tika 解析 + 文本清洗 + 统计分析
DocumentAnalysisResult result = parserService.parse(
fileBytes,
document.getOriginalFileName(),
document.getMimeType(),
fileType
);
// 6. 后续:保存 txt、结构节点、策略推荐、更新状态
}
解析器主流程:
public DocumentAnalysisResult parse(byte[] bytes, String fileName, String mimeType, FileType fileType) {
String rawText = extractRawText(bytes, fileName, mimeType, fileType);
String cleanedText = cleanupText(rawText);
List<DocumentStructureNodeCandidate> structureNodes =
structureNodeExtractor.extract(fileName, cleanedText);
List<String> paragraphs = extractParagraphs(cleanedText);
int charCount = cleanedText.length();
int tokenCount = estimateTokenCount(cleanedText);
int headingCount = countHeadings(cleanedText, structureNodes);
int structureLevel = evaluateStructureLevel(headingCount, paragraphs.size());
int contentQualityLevel = evaluateContentQuality(cleanedText, charCount);
return new DocumentAnalysisResult(...);
}
你不需要逐行背代码,但要能讲清楚:
Kafka 消费消息
查上下文
更新状态
下载文件
Tika 提取文本
清洗文本
统计分析
输出 DocumentAnalysisResult
企业级项目导航:⬅️ 03-上传接口与文档主记录创建 | 04-异步解析入口 | ➡️ 05-文本转文档结构树
💬 评论