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