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 讲解文档核心业务讲解视频从基础项目到复杂业务场景,项目资料会持续更新。
- 01Nexus Agent AI 智能体
- 02Nexus Agent Pro 完全版
- 03黑马点评Plus
- 04大麦
- 05大麦Pro
- 06大麦AI
- 07流量切换
- 08数据中台
加入后还能获得
进入星球后,即可享受上述所有服务,保证不会再有其他隐藏费用。从学习、面试到项目启动,都可以继续获得支持。
- 1 对 1 解答项目和技术问题都可以提问
- 针对性补充没有讲清楚的内容会继续补充
- 面试与简历指导梳理回答技巧和项目亮点
- 中间件云环境项目依赖可以直接接入使用
- 面试后复盘被问住的问题可以继续交流
- 远程问题解决项目启动问题可协助排查
