跳到主要内容

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 本身不做任何业务逻辑,纯粹是"反序列化 + 转发"

付费内容提示

该文档的全部内容仅对「码力全开」项目实战&技术讲解 知识星球用户开放

加入星球,一次获得完整项目资料、全栈技术知识库和长期答疑服务。

100万+字全栈技术知识库深入讲解技术核心、数据库、中间件和分布式等内容
8套热门的实战项目持续更新的企业级项目覆盖高并发、微服务、数据中台 和 AI Agent 等方向
AI 技术知识大模型面试详解覆盖 AI 模型原理、Agent、RAG、MCP、Skills、Harness 等核心知识点
文档 + 视频两种讲解形式既能系统阅读,也能跟随视频理解核心业务

完整项目实战资料

每套项目均包含从 0 到 1 讲解文档核心业务讲解视频

从基础项目到复杂业务场景,项目资料会持续更新。

8 套项目
  • 01Nexus Agent AI 智能体
  • 02Nexus Agent Pro 完全版
  • 03黑马点评Plus
  • 04大麦
  • 05大麦Pro
  • 06大麦AI
  • 07流量切换
  • 08数据中台

加入后还能获得

进入星球后,即可享受上述所有服务,保证不会再有其他隐藏费用。

从学习、面试到项目启动,都可以继续获得支持。

  • 1 对 1 解答项目和技术问题都可以提问
  • 针对性补充没有讲清楚的内容会继续补充
  • 面试与简历指导梳理回答技巧和项目亮点
  • 中间件云环境项目依赖可以直接接入使用
  • 面试后复盘被问住的问题可以继续交流
  • 远程问题解决项目启动问题可协助排查
知识星球二维码

扫码进入知识星球

  1. 打开微信,扫描左侧二维码,加入「码力全开」项目实战&技术讲解 知识星球
  2. 查看星球使用指导,获取完整项目讲解资料索引
解锁全部付费内容
🎁优惠