Kafka 消费与文本内容解析
上一篇讲完了文档上传的同步链路——文件校验、MinIO 上传、数据库入库,最后把消息丢进了 Kafka。
那 Kafka 消费方拿到消息之后到底做了什么?
这篇就来拆这条异步链路的前半段:从消费消息开始,到文本提取与清洗完成为止。
先上总览流程图。
异步解析链路总览
Kafka 消费者:消息入口
消费端的入口在 DocumentKafkaConsumer,它的职责很单一——把 JSON 消息反序列化,然后转交给异步处理服务。
/**
* 消费"解析路由"消息。
* <p>
* 这一步是上传完成后的异步链入口:收到消息后,会把 documentId 和 taskId 交给异步处理服务,
* 继续执行文档下载、正文解析、结构节点生成、策略推荐等步骤。
* </p>
*/
@KafkaListener(
topics = SPRING_INJECT_PREFIX_DISTINCTION_NAME + "-" + "${app.manage.kafka.parse-topic}",
groupId = "${app.manage.kafka.group-id}-parse")
public void consumeParseRoute(String payload) {
try {
// 先把 JSON 还原成强类型消息对象,避免后续处理层直接面对原始字符串。
DocumentParseRouteMessage message = objectMapper.readValue(payload,
DocumentParseRouteMessage.class);
// 真正的业务推进放到异步处理服务中,这里只承担"消费并转发"的职责。
asyncProcessService.handleParseRoute(message.getDocumentId(), message.getTaskId());
}
catch (Exception exception) {
// 消费失败只记录日志,不让异常继续向外冒泡破坏监听线程。
log.error("消费解析路由消息失败,payload={}", payload, exception);
}
}
几个要点:
- topic 名称是
环境前缀-配置的parseTopic,和上传端发送时用的是同一个 topic groupId带了-parse后缀,和索引构建的消费组区分开- 整个方法用 try-catch 包住,消费失败只打日志不抛异常——这是为了防止一条坏消息把整个消费线程搞挂
- Consumer 本身不做任何业务逻辑,纯粹是"反序列化 + 转发"
handleParseRoute:异步解析主方法
真正干活的是 DocumentAsyncProcessServiceImpl.handleParseRoute()。这个方法比较长,我们按阶段拆开来看。
加载上下文 + 推进任务状态
/**
* 处理"解析路由"任务。
* <p>
* 这是文档上传成功后的第一条异步业务主链,整体顺序如下:
* 1. 读取文档记录和任务记录,确认异步处理上下文存在;
* 2. 把任务状态切到 RUNNING,把当前阶段推进到 CONTENT_PARSE;
* 3. 从对象存储下载原始文件,并调用解析器提取纯文本和结构信息;
* 4. 把解析后的纯文本重新上传为 txt,便于后续索引构建直接复用;
* 5. 用结构节点服务替换文档结构节点,并同步导航索引、图投影和画像;
* 6. 基于解析结果调用策略服务生成推荐切块方案;
* 7. 把推荐方案和步骤写入数据库,同时更新文档的解析状态、策略状态和统计信息;
* 8. 以成功或失败状态收尾任务,并记录任务日志。
* </p>
* <p>
* 这个方法本身不直接执行向量化和 chunk 落库,那是后续"索引构建任务"的职责;
* 当前阶段的目标是把文档从"原始文件"推进到"解析完成并拿到推荐策略"。
* </p>
*/
public void handleParseRoute(Long documentId, Long taskId) {
SuperAgentDocument document = documentMapper.selectById(documentId);
SuperAgentDocumentTask task = taskMapper.selectById(taskId);
if (document == null || task == null) {
log.warn("解析任务对应的文档或任务不存在,documentId={}, taskId={}", documentId, taskId);
return;
}
Date startTime = new Date();
try {
// 进入异步解析后,先把任务状态切成 RUNNING,当前阶段设为"内容解析"。
task.setTaskStatus(DocumentTaskStatusEnum.RUNNING.getCode());
task.setCurrentStage(DocumentTaskStageEnum.CONTENT_PARSE.getCode());
task.setStartTime(startTime);
taskMapper.updateById(task);
// 文档主表的解析状态也同步切到 PARSING,方便前端查询时看到实时状态。
document.setParseStatus(DocumentParseStatusEnum.PARSING.getCode());
documentMapper.updateById(document);
第一件事就是从数据库把文档记录和任务记录捞出来。如果查不到(比如上传事务回滚了、或者数据被删了),直接 return,不做任何处理。
然后把任务状态从 NEW 切到 RUNNING,当前阶段设为 CONTENT_PARSE。文档主表的 parseStatus 也同步更新为 PARSING。这样前端轮询文档状态时就能看到"正在解析中"。
记录开始日志 + 下载原始文件
// 先记一条开始日志,把原始对象名落下来,便于后续排查下载与解析问题。
taskLogService.saveLog(taskId, documentId,
DocumentTaskStageEnum.CONTENT_PARSE.getCode(),
DocumentTaskEventTypeEnum.START.getCode(),
DocumentLogLevelEnum.INFO.getCode(),
DocumentOperatorTypeEnum.SYSTEM.getCode(),
null,
"开始解析文档内容。",
Map.of("objectName", document.getObjectName()));
// 重新从对象存储下载原始文件,确保异步链处理的是"上传后实际持久化成功"的那份文件。
byte[] fileBytes = storageService.downloadObject(document.getObjectName());
为什么要重新下载文件?
上传阶段的
byte[]是在 HTTP 请求线程里的内存数据,请求结束后就释放了。异步消费方运行在另一个线程(甚至另一个进程),所以必须从 MinIO 重新下载。这也顺便验证了文件确实已经成功持久化到对象存储。
调用 Tika 解析文档
// parserService 会负责提取纯文本、结构节点候选、字符数、token 数、质量等级等分析结果。
DocumentAnalysisResult analysisResult = parserService.parse(fileBytes,
document.getOriginalFileName(),
document.getMimeType(),
DocumentFileTypeEnum.getRc(document.getFileType()));
这一步是整条链路的核心——把原始二进制文件变成结构化的分析结果。我们深入看看 TikaDocumentParserService.parse() 的实现。
TikaDocumentParserService:文档解析器
parse 方法总览
/**
* 解析文档原始字节并生成标准化分析结果。
* <p>
* 这是上传后异步解析链里的核心一步,整体顺序如下:
* 1. 先根据文件类型提取原始文本;
* 2. 对文本做换行、空白字符和乱码清洗;
* 3. 提取结构节点候选,并估算标题数量;
* 4. 统计段落、字符数、token 数;
* 5. 评估文档结构等级和内容质量等级;
* 6. 把这些结果打包成 DocumentAnalysisResult 返回给上游。
* </p>
*/
public DocumentAnalysisResult parse(byte[] bytes, String originalFileName,
String mimeType, DocumentFileTypeEnum fileType) {
// 第一步先把不同格式文件转换成统一的原始文本。
String rawText = extractRawText(bytes, originalFileName, mimeType, fileType);
// 再做文本清洗,为后续标题识别、段落切分和结构抽取提供更稳定的输入。
String cleanedText = cleanupText(rawText);
// 结构节点由专门的提取器负责抽取,这些节点后面会参与导航和切块策略判断。
List<DocumentStructureNodeCandidate> structureNodes =
structureNodeExtractor.extract(originalFileName, cleanedText);
int headingCount = countHeadings(cleanedText, structureNodes);
// 段落统计既用于结构判断,也用于后续策略推荐时判断是否适合语义切块。
List<String> paragraphList = extractParagraphs(cleanedText);
int maxParagraphLength = paragraphList.stream().mapToInt(String::length).max().orElse(0);
int charCount = cleanedText.length();
// token 数这里是估算值,不是精确 tokenizer 结果,但足够用于策略判断和粗粒度统计。
int tokenCount = estimateTokenCount(cleanedText);
int structureLevel = evaluateStructureLevel(headingCount, paragraphList.size());
int contentQualityLevel = evaluateContentQuality(cleanedText, charCount);
return new DocumentAnalysisResult(
cleanedText, charCount, tokenCount, structureLevel, contentQualityLevel,
headingCount, paragraphList.size(), maxParagraphLength, structureNodes
);
}
整个 parse 方法可以看成一条"分析流水线",每一步都在为后续的策略推荐积累判断依据。接下来我们按执行顺序逐个展开。
extractRawText:按文件类型提取原始文本
/**
* 按文件类型提取原始文本。
* <p>
* Office/PDF/HTML 等复杂格式优先走 Tika;
* 纯文本格式则直接按 UTF-8 解码。
* 如果 Tika 失败但 MIME/文件类型仍明显是文本,则退回到直接字符串解码作为兜底。
* </p>
*/
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) {
// 某些文本文件虽然 Tika 解析失败,但原始字节本身仍然可以直接按文本解码读取。
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);
}
}
提取策略很清晰:
- PDF / DOC / DOCX / HTML:走 Apache Tika,它能处理各种复杂格式
- TXT / MD:直接 UTF-8 解码,不需要 Tika
- 兜底:如果 Tika 解析失败,但文件类型明显是文本类的,就退化成直接字符串解码
cleanupText:文本清洗
private String cleanupText(String rawText) {
if (StrUtil.isBlank(rawText)) {
return "";
}
String cleaned = rawText
.replace("\r\n", "\n") // Windows 换行统一成 Unix 换行
.replace('\r', '\n') // 老 Mac 换行也统一
.replace('\u0000', ' ') // 空字符替换成空格
.replaceAll("[\\t\\x0B\\f]+", " ") // 制表符等控制字符替换成空格
.replaceAll("\\n{3,}", "\n\n") // 连续 3 个以上换行压缩成 2 个
.replaceAll("[ ]{2,}", " ") // 连续多个空格压缩成 1 个
.trim();
return cleaned;
}
这步看着简单,但很重要。Tika 提取出来的文本经常会有各种奇怪的控制字符、多余的空行,如果不清洗,后面的标题识别和段落切分都会受影响。
小结
到这里,异步解析链路的入口和文本提取阶段就走完了。整个过程可以概括为:
消费 Kafka 消息 → 加载上下文 → 推进任务状态 → 下载原始文件 → Tika 提取文本 → 文本清洗
拿到清洗后的纯文本之后,parse() 方法接下来要做的第一件事就是结构节点提取——这是整条解析流水线中最复杂的一环,下一篇单独展开。
企业级项目导航:⬅️ 07-知识路由的功能 | 01-Kafka 消费与文本内容解析 | ➡️ 02-一定得用Kafka吗

💬 评论