Kafka 消费与文本内容解析

上一篇讲完了文档上传的同步链路——文件校验、MinIO 上传、数据库入库,最后把消息丢进了 Kafka。

那 Kafka 消费方拿到消息之后到底做了什么?

这篇就来拆这条异步链路的前半段:从消费消息开始,到文本提取与清洗完成为止。

先上总览流程图。

异步解析链路总览

Fhrq35ZYem8aVEVCLc9ctrbtIP-Z-b0ba0c72

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吗