kafka消费采集数据完成升级
在上一章中讲解了项目中责任链的四个节点的执行过程:
- 数据的采集
- 数据的计算
- 数据的保存
- 数据的升级
在数据的升级中,将消息发送到了 kafka 中,本章节会详细讲解从 kafka 中消费到消息后,后续的执行过程
Kafka消费消息
org.javaup.kafka.UpDimensionConsumer#consumerUpDimensionMessage
@KafkaListener(topics = {SPRING_INJECT_PREFIX_DISTINCTION_NAME+"-"+"${spring.kafka.topic:up_dimension}"},
containerFactory = "kafkaListenerContainerFactory")
public void consumerUpDimensionMessage(ConsumerRecord<String,String> consumerRecord, Acknowledgment acknowledgment){
String messageKey = consumerRecord.key();
String messageValue = consumerRecord.value();
log.info("开始处理Kafka消息,topic: {}, partition: {}, offset: {}, key: {}",
consumerRecord.topic(), consumerRecord.partition(), consumerRecord.offset(), messageKey);
// 检查消息是否为空
if (StringUtil.isEmpty(messageValue)) {
log.warn("收到空消息,跳过处理,key: {}", messageKey);
acknowledgment.acknowledge();
return;
}
// 解析消息
UpgradeDimensionMessage upgradeDimensionMessage = JSON.parseObject(messageValue, UpgradeDimensionMessage.class);
if (Objects.isNull(upgradeDimensionMessage)) {
log.error("消息解析失败,无法转换为UpgradeDimensionMessage对象,key: {}", messageKey);
// 解析失败的消息可以选择跳过,避免无限重试
acknowledgment.acknowledge();
return;
}
upDimensionConsumerExecutor.doConsumer(upgradeDimensionMessage,acknowledgment);
}
执行流程拆解
1. 提取消息基础信息并记录可观测日志
- 从
consumerRecord读取key与value,并打点日志包含topic、partition、offset与key,便于后续审计与问题排查。 - 这一步不改变消费位移,仅为后续处理做信息准备。
2. 空消息快速跳过与确认
// 检查消息是否为空
if (StringUtil.isEmpty(messageValue)) {
log.warn("收到空消息,跳过处理,key: {}", messageKey);
acknowledgment.acknowledge();
return;
}
- 判空: 若
value为空,直接打印警告并调用acknowledgment.acknowledge()手动提交 offset。 - 影响: 该消息被视为已消费,避免对“空消息”进行无意义重试。
3. 反序列化消息体为领域对象
// 解析消息
UpgradeDimensionMessage upgradeDimensionMessage = JSON.parseObject(messageValue, UpgradeDimensionMessage.class);
- 使用 FastJSON 将
messageValue反序列化为UpgradeDimensionMessage,承载后续业务处理所需的参数包(如messageId、createTime、totalParamTransfers等)。
4. 解析失败的快速跳过与确认
if (Objects.isNull(upgradeDimensionMessage)) {
log.error("消息解析失败,无法转换为UpgradeDimensionMessage对象,key: {}", messageKey);
// 解析失败的消息可以选择跳过,避免无限重试
acknowledgment.acknowledge();
return;
}
- 若反序列化结果为
null,判定为“消息格式不合规或内容异常”,立即日志报错并手动确认 offset。 - 避免“格式错误”的消息在消费端形成无限重试,尽快跳过。
5. 委派给升级维度消息执行器
upDimensionConsumerExecutor.doConsumer(upgradeDimensionMessage,acknowledgment);
- 将已通过前置校验与解析的
UpgradeDimensionMessage交给upDimensionConsumerExecutor.doConsumer(...)。 - 重要: 本方法不再直接处理“业务逻辑、幂等控制、消费记录入库与最终 ack 决策”;这些职责在执行器里完成。也就是说,正常处理路径下的“是否 ack、何时 ack”由执行器按结果来决定。
执行路径与分支总结
- 正常路径: 提取消息 → 打点日志 → 反序列化 → 委派给执行器 →(由执行器在成功时 ack 或在错误策略下决定是否 ack)
- 异常路径 A(空消息): 判空 → warn → 立即 ack → 返回
- 异常路径 B(解析失败): 解析为
null→ error → 立即 ack → 返回
设计要点与策略解读
- 最小职责原则:
UpDimensionConsumer只做“输入校验 + 反序列化 + 委派”,业务细节与消费幂等、入库追踪、最终 ack 策略都放在执行器,职责边界清晰。 - 手动 ack 策略:
- 在“空消息/解析失败”两类非业务型异常,快速 ack,避免重复消费。
- 成功与否的业务性判断留给执行器,确保“只有处理正确时才提交 offset”,从而获得“至少一次”到“有效一次”的控制能力。
- 可观测性: 入口日志包含
topic/partition/offset/key,配合执行器内部日志,可完整追踪一次消息从进入到处理完毕的全链路。 - 健壮性: 解析失败即跳过,避免无效重试占用资源;业务异常的处理策略位移提交由执行器决定,可根据异常类型灵活处理。
以上是Kafka消费到消息后,进行解析的过程,真正的执行业务过程是交给了升级维度消息执行器upDimensionConsumerExecutor
付费内容提示
该文档的全部内容仅对「码力全开」项目实战&技术讲解 知识星球用户开放
加入星球,一次获得完整项目资料、全栈技术知识库和长期答疑服务。
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 解答项目和技术问题都可以提问
- 针对性补充没有讲清楚的内容会继续补充
- 面试与简历指导梳理回答技巧和项目亮点
- 中间件云环境项目依赖可以直接接入使用
- 面试后复盘被问住的问题可以继续交流
- 远程问题解决项目启动问题可协助排查
