SpringBoot Kafka 全操作指南

SpringBoot Kafka 实战 📖 前置阅读:本文假设读者已理解 Kafka 的核心概念(Broker、Topic、Partition、ConsumerGroup、Offset)。如果还不熟悉,建议先阅读 Kafka 核心架构与日志存储模型。 🎯 第一步:目标说明 上一篇用原版 Kafka Java Client 写了 KafkaProducer + KafkaConsumer。和 RabbitMQ、RocketMQ 一样——真实的 SpringBoot 项目里不需要那些样板代码。spring-kafka 帮我们处理了连接管理、Producer 生命周期、Consumer 线程池、Offset 提交。 读完这篇会掌握: KafkaTemplate 三种发送方式(同步/异步/回调) @KafkaListener 注解消费——单条和批量 JSON 序列化全链路配置——Producer 端 JsonSerializer + Consumer 端 JsonDeserializer Producer 配置:acks、retries、batch.size、linger.ms、compression.type Consumer 配置:group.id、auto.offset.reset、enable.auto.commit、max.poll.records 📋 第二步:前置条件 前置项 具体要求 验证命令 JDK 17+(8+ 也兼容) java -version SpringBoot 3.x(文中用 3.2) mvn dependency:tree | grep spring-boot Kafka 3.7.0 KRaft 模式(单节点即可) docker ps | grep kafka 前置知识 Broker/Topic/Partition/ConsumerGroup/Offset 概念 — 确认 Kafka 在跑: ...

十一月 14, 2022 · 8 分钟 · 1619 字 · yaomingye

Kafka 核心架构与日志存储模型

Kafka:分布式提交日志,不是消息队列 📖 前置阅读:本文假设读者已理解消息队列的基本价值(异步、解耦、削峰填谷),最好读过 RabbitMQ 或 RocketMQ 的任意一篇基础文章。有了 MQ 概念再学 Kafka 事半功倍。 一、⚡ 问题切入:RabbitMQ 和 RocketMQ 有什么共同的"毛病"? 先回顾 RabbitMQ 和 RocketMQ 的消费模型: RabbitMQ: Consumer 收到消息 → 手动 basicAck → Broker 删除消息 RocketMQ: Consumer 收到消息 → 返回 CONSUME_SUCCESS → offset 推进 共同点:消息被消费者确认(ACK)后,Broker 就把它删了。消息在 Broker 上的生命周期是"暂存"——它存在只是为了等待消费者拿走。 这个模型有一个隐含的限制:一条消息只能被消费一次。想重放消息?RabbitMQ 做不到(消息已经删了),RocketMQ 可以重置 offset 但受 CommitLog 保留时间限制。 这时候再看 Kafka 的设计:消息消费后不删除。消息存在磁盘上,按时间或大小策略统一过期,消费者想从哪个位置读就从哪个位置读。 Kafka: Consumer 自己管 offset,随时可以回到过去的某个位置重读 消息不是被消费掉的——是按时间自然过期的 这就是 Kafka 和 RabbitMQ/RocketMQ 本质上的不同——Kafka 不是一个消息队列,它是一个分布式提交日志(Distributed Commit Log)。 二、🧬 Kafka 是什么:分布式提交日志 2.1 核心定义 Kafka 的官方定位:分布式、分区化、多副本的提交日志服务。 ...

十一月 13, 2022 · 5 分钟 · 1020 字 · yaomingye

RocketMQ 生产环境部署与调优

RocketMQ 生产部署 📖 前置阅读:本文是 RocketMQ 系列的终篇,假设读者已经掌握前五篇的全部内容(核心架构、SpringBoot 集成、高级消息类型、可靠性、消费者模式)。 一、⚡ 问题切入:单机 Docker 的瓶颈 第一篇搭的单机 RocketMQ(一个 NameServer + 一个 Broker)只能用来学习——生产环境中: 单点 后果 一台 NameServer 挂了 Producer/Consumer 无法获取路由——整个 MQ 瘫痪 一台 Broker 挂了 所有消息不可用——消息无法发送/消费 JVM 内存不足 Full GC 频繁 → 消息延迟抖动 → 超时重试雪崩 磁盘写满 CommitLog 无法写入 → 生产者阻塞 生产最低配:2 台 NameServer + 至少 2 台 Broker(主从)。 二、双主双从高可用集群搭建 2.1 架构设计 NameServer 集群:2 台(互不通信,各自独立) Broker 集群:2 组主从(Master-A + Slave-A、Master-B + Slave-B) NameServer-1 NameServer-2 (192.168.1.10:9876) (192.168.1.11:9876) ↑ ↑ ┌─────────┴───────┬───────────────┘ │ │ Broker-A (Master) Broker-A (Slave) 192.168.1.20:10911 192.168.1.21:10911 brokerId=0 brokerId=1 Broker-B (Master) Broker-B (Slave) 192.168.1.22:10911 192.168.1.23:10911 brokerId=0 brokerId=1 路由发现:Producer/Consumer 配置所有的 NameServer 地址——只要有一台 NameServer 活着,路由就能工作。 ...

十一月 12, 2022 · 5 分钟 · 909 字 · yaomingye

RocketMQ 消费者模式与过滤器

RocketMQ 消费者模式 📖 前置阅读:本文假设读者已掌握 SpringBoot RocketMQ 的基本消费操作。如果还不熟悉,建议先阅读 SpringBoot RocketMQ 全操作指南。 一、⚡ 问题切入:一条消息,谁来消费? 前面五篇的消费者代码都默认了一件事——一条消息只被一个消费者实例处理。但实际业务中: 订单消息——只能被一个实例消费(一个订单不能被两个服务处理两次) 配置刷新消息——所有实例都要收到(所有缓存节点刷新缓存) 部分 Tag 的消息——只关心"订单创建",对"订单支付"不感兴趣 消费不过来——10 个 Queue,2 个实例,怎么分? 这四大问题的答案都在这一篇里。 二、集群消费 vs 广播消费 2.1 集群消费(CLUSTERING)——默认模式 同一个 ConsumerGroup 内的所有实例共享消费一个 Topic 的消息——每条消息只被组内一个实例消费。 @Component @RocketMQMessageListener( topic = "order-topic", consumerGroup = "order-consumer-group", messageModel = MessageModel.CLUSTERING // 集群模式(默认值,可以不写) ) public class OrderClusteringListener implements RocketMQListener<OrderMessage> { @Override public void onMessage(OrderMessage msg) { System.out.printf("[实例A] 处理订单: orderId=%d%n", msg.getOrderId()); } } 部署 3 个实例,Topic 有 8 个 Queue: ...

十一月 11, 2022 · 4 分钟 · 730 字 · yaomingye

RocketMQ 消息可靠性与容错

RocketMQ 消息可靠性 📖 前置阅读:本文假设读者已掌握 SpringBoot RocketMQ 的发送和消费操作。如果还不熟悉,建议先阅读 SpringBoot RocketMQ 全操作指南。 一、⚡ 消息可能丢在哪? 和 RabbitMQ 一样,RocketMQ 的消息丢失也分三个环节: Producer → [网络] → Broker → [网络] → Consumer ① 发送丢失 ② Broker 宕机丢失 ③ 消费失败丢失 在 RocketMQ 中,生产者端没有 RabbitMQ 的 Publisher Confirm——替代方案是同步发送 + 返回值判断。Broker 端的持久化取决于刷盘策略和主从同步。消费端靠消费状态返回 + 重试 + 死信。 二、① 生产者端:同步发送判断返回值 RocketMQ 没有 RabbitMQ 的 ConfirmCallback 异步通知机制,但 syncSend 本身就是同步等待 Broker 确认——返回 SEND_OK 说明消息已写入 CommitLog: @Service public class ReliableProducer { @Autowired private RocketMQTemplate rocketMQTemplate; public void sendReliably(OrderMessage msg) { // syncSend 是同步阻塞的——返回 SEND_OK 才说明 Broker 已接收 SendResult result = rocketMQTemplate.syncSend( "order-topic:created", msg); if (SendStatus.SEND_OK.equals(result.getSendStatus())) { log.info("消息已确认到达 Broker: msgId={}", result.getMsgId()); } else { // 只有 SEND_OK 才认为成功——FLUSH_DISK_TIMEOUT 等状态说明刷盘超时 log.error("消息发送未确认: status={}", result.getSendStatus()); // 补偿:写入 DB 重试表 saveToRetryTable(msg); } } } SendStatus 的四种返回值: ...

十一月 10, 2022 · 5 分钟 · 853 字 · yaomingye

RocketMQ 顺序消息、延迟消息与事务消息

顺序消息、延迟消息与事务消息 📖 前置阅读:本文假设读者已掌握 SpringBoot RocketMQ 的基本操作(RocketMQTemplate、@RocketMQMessageListener)。如果还不熟悉,建议先阅读 SpringBoot RocketMQ 全操作指南。 一、⚡ 问题切入:三种 RabbitMQ 做不到或做不好的事 RabbitMQ 六篇系列学完时留了几个坑——有些场景 RabbitMQ 不是不能用,而是做起来别扭: 需求 RabbitMQ 方案 痛点 订单创建→支付→发货严格按序 单队列 + 单消费者,关并发 吞吐量压到一条线;一旦重试入队顺序全乱 30 分钟后自动取消 Delayed Message 插件或 TTL+DLX 插件生产不可靠;TTL+DLX 有消息时序问题 下单 + 扣库存 + 发消息三件事原子执行 自己实现本地消息表 + 定时补偿 代码量大,维护麻烦 RocketMQ 对这三种场景都有原生支持——不是插件,不是 workaround,是设计时就考虑进去了。 二、顺序消息 —— 深度篇 2.1 上一篇回顾 + 补充 上一篇讲了基本用法:syncSendOrderly 用 orderId 哈希选 Queue,同一个 orderId 进同一个 Queue → 该 Queue 内 FIFO。消费端 consumeMode = ConsumeMode.ORDERLY。 但这只讲了正常流程。重试会破坏顺序——这是最容易踩的坑。 2.2 顺序消费的重试机制:挂起而非重入队 并发消费中,失败的消息通过 RECONSUME_LATER 进入重试 Topic,然后延迟重新投递。但顺序消费不能这么干——如果第 2 条消息失败后进了重试队列,第 3 条消息先被消费,顺序就乱了。 ...

十一月 9, 2022 · 6 分钟 · 1151 字 · yaomingye

SpringBoot RocketMQ 全操作指南

SpringBoot 集成 RocketMQ:从发送到消费 📖 前置阅读:本文假设读者已理解 RocketMQ 的核心概念(NameServer、Broker、Topic、Queue、ConsumerGroup)。如果还不熟悉,建议先阅读 RocketMQ 核心架构与消息模型。 Part 1:概念与前置 1.1 本文目标 上一篇用原版 RocketMQ Java Client 写了 DefaultMQProducer + DefaultMQPushConsumer。真实的 SpringBoot 项目里不需要那么多样板代码——rocketmq-spring-boot-starter 帮你处理了 NameServer 连接、Producer 启动、Consumer 注册。 读完这篇会掌握: MqHelper 封装——为什么要在 RocketMQTemplate 上再包一层,asyncSend + SendCallback 的真实用法 30 分钟延迟取消订单——RocketMQ 内置 delayLevel 的完整实战流程,和 RabbitMQ 方案的对比 @RocketMQMessageListener 消费者——MessageExt 手动反序列化 vs 泛型自动解析,以及三个真实业务消费者 Domain Entity 直传——为什么不做 DTO 转换,以及什么情况下不能这样做 双 MQ 基础设施先行——RabbitMQ 拓扑已就绪但全用 RocketMQ 的设计决策 1.2 前置条件 前置项 具体要求 验证命令 JDK 17+(8+ 也兼容) java -version SpringBoot 3.x(文中用 3.2) mvn dependency:tree | grep spring-boot RocketMQ 5.1.4(NameServer + Broker 都在运行) docker ps | grep rocketmq 前置知识 NameServer/Broker/Topic/Queue 概念 — 确认 RocketMQ 在跑: ...

十一月 8, 2022 · 16 分钟 · 3233 字 · yaomingye

RocketMQ 核心架构与消息模型——领域模型、存储引擎与消息生命周期全景

领域模型与存储引擎 📖 前置阅读:本文假设读者已理解消息队列的基本价值(异步、解耦、削峰填谷)。如果还不熟悉消息队列,建议先阅读 RabbitMQ 核心概念与 AMQP 协议。 一、问题切入:为什么 RocketMQ 的概念比 RabbitMQ 多? RabbitMQ 学完六篇,Exchange / Binding / Queue 的路由模型印象深刻——概念不多,全靠灵活组合。翻开 RocketMQ 的文档,一眼扫过去:Producer、Consumer、Topic、Queue、ConsumerGroup、Subscription、Broker、NameServer……光领域概念就七个,外加两个部署组件。 这不是设计过度,而是 RocketMQ 把"谁负责发、谁负责收、怎么分组、怎么扩容、消息存哪里、谁管路由"全部显式拆开了。 RabbitMQ 用少数概念的组合来表达这些维度,RocketMQ 选择每件事都定义一个独立概念。 好处是每个概念职责单一,坏处是初学者一看就晕——概念之间谁包谁、谁管谁、谁和谁是平等的,不看图根本理不清。 所以这篇不讲"先记住七个概念"——先看一张全景图,把七者的层级关系钉在脑子里。 二、领域模型全景:一张图串起七个概念 Apache RocketMQ 官方把领域模型定义为七个核心概念。这是一张按层级包含关系组织的全景图: flowchart TD subgraph PRODUCTION["① 生产"] P(["Producer\n生产者"]) end subgraph STORAGE["② 存储"] MSG["Message\n消息体"] TOPIC["Topic\n主题(逻辑容器)"] Q0[("Queue-0\n物理分片")] Q1[("Queue-1")] Q2[("Queue-n")] MSG --> TOPIC TOPIC --> Q0 TOPIC --> Q1 TOPIC --> Q2 end subgraph CONSUMPTION["③ 消费"] CG["ConsumerGroup\n消费分组"] C1(["Consumer-1"]) C2(["Consumer-2"]) SUB["Subscription\n订阅关系"] CG --> C1 CG --> C2 CG --> SUB end P -->|"发送"| TOPIC Q0 & Q1 & Q2 -->|"拉取"| CG SUB -.->|"绑定"| TOPIC classDef startEnd fill:#fff5f5,stroke:#e53e3e,stroke-width:2.5px,font-weight:bold; classDef process fill:#f7fafc,stroke:#718096,stroke-width:2px; classDef data fill:#f0fff4,stroke:#38a169,stroke-width:2px,font-weight:bold; classDef root fill:#ebf8ff,stroke:#3182ce,stroke-width:2.5px,font-weight:bold; class P,C1,C2 startEnd; class MSG,TOPIC process; class Q0,Q1,Q2 data; class CG,SUB root; 这张图的阅读顺序:从左到右,消息从 Producer 出发 → 经过 Topic → 落入 Queue → 被 ConsumerGroup 内的 Consumer 拉取。关键层级关系: ...

十一月 7, 2022 · 9 分钟 · 1741 字 · yaomingye

RabbitMQ 生产环境部署与调优

生产环境部署与调优 📖 前置阅读:本文是 RabbitMQ 系列的终篇,假设读者已经掌握前五篇的全部内容(核心概念、交换机类型、SpringBoot 集成、消息可靠性、高级特性)。 一、⚡ 问题切入:单机 RabbitMQ 什么时候扛不住? 前三篇代码都在本地单机 RabbitMQ 上跑的——一个 Docker 容器,内存 1G,磁盘 10G。跑到生产环境会发生什么? 场景 单机 RabbitMQ 的后果 服务器宕机 整个 MQ 服务中断——所有生产者阻塞、消费者闲置 消息积压 50 万条 内存打满 → RabbitMQ 触发内存告警 → 阻塞所有生产者(flow control) 磁盘写满 RabbitMQ 拒绝所有写入→ 消息丢失 每秒 2 万条消息 单机 CPU 100%,消息延迟从 1ms 飙升到 500ms 单机不是不能用,但要知道它的边界。以下场景必须上集群: 消息不能丢(金融交易、订单处理) 服务不能停(7×24 在线业务) 吞吐量超过单机极限(> 5 万 msg/s) 二、RabbitMQ 集群架构 2.1 集群的基本原理 RabbitMQ 集群是多个 Erlang 节点组成的对等网络——每个节点运行一个 RabbitMQ 实例。 集群中共享的东西: Exchange、Queue、Binding 的元数据(定义信息)——在所有节点上自动同步 用户、vhost、权限——自动同步 集群中不共享的东西: ...

十一月 6, 2022 · 6 分钟 · 1085 字 · yaomingye

RabbitMQ 延迟队列与高级特性

延迟队列与高级特性 📖 前置阅读:本文假设读者已掌握 RabbitMQ 的基础操作(Exchange/Queue/Binding/手动ACK)和 SpringBoot 集成。如果还不熟悉,建议先阅读前三篇。 一、⚡ 问题切入:有些事不能现在做 看几个每天都在发生的业务需求: 用户下单后 30 分钟未支付,自动取消订单 用户注册后 7 天未登录,发一封"想你了"邮件 优惠券到期前 3 小时,发短信提醒用户使用 支付成功后立即通知商家,但每封邮件间隔 5 秒(防被邮箱限流) 这些需求的共同点:消息不能发出去就被立刻消费,需要在未来某个时间点才能被消费。 普通队列是"发了就收"——消息一进队就被消费者拿走。延迟队列是"发了先等着,时间到了再收"。 RabbitMQ 没有原生的延迟队列类型,但有两种方式可以实现。 二、延迟队列方案一:TTL + 死信队列 2.1 原理 这是利用已有的机制组合而成的"曲线救国"方案: Producer → 普通 Exchange → 死信队列(当"延迟缓冲区") → TTL 到期 → 死信 Exchange → 实际消费队列 → Consumer ↑ 消息过期变死信自动转入 核心思想:创建一个没有消费者的队列,设置 TTL。消息在这个队列中"等"TTL 时间后过期变成死信,被自动转发到真正的消费队列——消费者只监听消费队列。消费者感知不到延迟的存在。 flowchart TD classDef startEnd fill:#701a4c,stroke:#e11d48,stroke-width:2px,color:#fce7f3,font-weight:bold; classDef process fill:#1e1e24,stroke:#6b7280,stroke-width:1.5px,color:#e5e7eb; classDef data fill:#052e16,stroke:#16a34a,stroke-width:1.5px,color:#bbf7d0,font-weight:bold; classDef highlight fill:#450a0a,stroke:#dc2626,stroke-width:1.5px,color:#fecaca,font-weight:bold; subgraph DELAY_FLOW ["TTL + DLQ 延迟队列流程"] P([Producer]) -->|"发送\n带 TTL"| DE[delay.exchange] DE -->|"routing: order.delay.30m"| DQ[queue.order.delay.30m\nTTL=30分钟\ndeadLetterExchange=real.exchange\ndeadLetterRoutingKey=order.real] DQ -->|"30分钟后\n消息过期变死信"| DLX[real.exchange] DLX -->|"routingKey=order.real"| RQ[queue.order.real\n实际消费队列] RQ --> C([Consumer\n订单取消服务]) C -.->|"⚠️ 消费者不监听延迟队列\n只监听实际消费队列"| DQ end class P startEnd; class DE,DQ,DLX,RQ process; class C startEnd; 2.2 配置代码 每段延迟时间需要一个独立的队列——每个队列有固定的 TTL: ...

十一月 5, 2022 · 6 分钟 · 1205 字 · yaomingye
Cat Radio