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

RabbitMQ 消息可靠性保障

消息可靠性保障:三道防线全覆盖 📖 前置阅读:本文假设读者已掌握 SpringBoot RabbitMQ 的基本操作(RabbitTemplate 发送、@RabbitListener 消费、手动 ACK)。如果还不熟悉,建议先阅读 SpringBoot RabbitMQ 全操作指南。 一、⚡ 问题切入:消息去哪儿了? 先看一段日常的订单处理代码: // 下单成功后发消息 @Service public class OrderService { @Transactional public void createOrder(OrderRequest req) { orderMapper.insert(req.toOrder()); // 1. 写 MySQL rabbitTemplate.convertAndSend( // 2. 发消息 "order.exchange", "order.created", req); } } 表面看起来没问题。但消息真的被消费了吗?在以下任何一个环节都可能丢: Producer → [网络] → RabbitMQ → [网络] → Consumer ① 发送丢失 ② Broker 宕机丢失 ③ 消费失败丢失 环节 丢失原因 后果 ① 生产者 → Broker 网络断连、Exchange 不存在、消息路由失败 消息根本没进队列 ② Broker 存储 RabbitMQ 进程崩溃、服务器断电 内存中的消息全部丢失 ③ Consumer 消费 消费者处理到一半挂了、代码异常没 ACK 消息被取走但实际没处理完 这三个环节必须逐一设防——RabbitMQ 提供了完整的机制,但需要生产者、Broker、消费者三端配合。 ...

十一月 4, 2022 · 8 分钟 · 1512 字 · yaomingye

SpringBoot RabbitMQ 全操作指南

SpringBoot 集成 RabbitMQ:从发送到消费 📖 前置阅读:本文假设读者已理解 RabbitMQ 的核心概念(Exchange、Queue、Binding、RoutingKey)和四种交换机类型。如果还不熟悉,建议先阅读前两篇: RabbitMQ 核心概念与 AMQP 协议 交换机类型完全指南 Part 1:概念与前置 1.1 本文目标 前两篇用 RabbitMQ 原生 Java Client 写了所有代码——channel.basicPublish、channel.basicConsume、手动 basicAck。理解底层是正确的,但真正进项目时,Spring AMQP 帮我们做了 90% 的重复工作。 读完这篇会掌握: 用 RabbitTemplate 一行代码发消息(替代 channel.basicPublish 那一大堆) 用 @RabbitListener 注解收消息(替代手动 basicConsume + DeliverCallback) 用 Jackson2JsonMessageConverter 自动序列化/反序列化 Java 对象 用 @Bean + 声明式配置 管理 Exchange/Queue/Binding(替代每次启动时 channel.exchangeDeclare) 三种交换机在 Spring 中的完整示例代码(Direct / Fanout / Topic) 手动 ACK 的配置和坑 1.2 前置条件 前置项 具体要求 验证命令 JDK 17+(8+ 也兼容) java -version Maven 3.6+ mvn -v SpringBoot 3.x(文中用 3.2) mvn dependency:tree | grep spring-boot RabbitMQ 3.12+(management 版) docker ps | grep rabbitmq 前置知识 前两篇的 Exchange/Queue/Binding/RoutingKey 概念 — 确认 RabbitMQ 在跑: ...

十一月 3, 2022 · 13 分钟 · 2677 字 · yaomingye

RabbitMQ 交换机类型完全指南

四种交换机:路由机制完全解析 📖 前置阅读:本文假设读者已理解上一篇中 Exchange、Queue、Binding、RoutingKey 的概念。如果还不清楚,建议先阅读 RabbitMQ 核心概念与 AMQP 协议。 一、⚡ 问题切入:同一条消息,为什么有人收到有人收不到? 上一篇结尾发了第一条 RabbitMQ 消息——消息发出去,消费者收到了。但实际业务远比这个复杂: 订单创建后,所有下游服务(短信、邮件、风控、日志)都要收到通知 商品价格变更后,只有关注了这个商品的搜索服务需要重建索引 用户行为日志中,一部分是购买行为(需要发优惠券),一部分是浏览行为(只需要统计) 这些需求的本质是路由——同一批消息,不同消费者按不同规则接收不同子集。RabbitMQ 用 Exchange(交换机)来承担这个角色。 Exchange 有四种类型。它们唯一的不同是如何匹配 RoutingKey 和 BindingKey: flowchart TD classDef startEnd fill:#701a4c,stroke:#e11d48,stroke-width:2px,color:#fce7f3,font-weight:bold; classDef condition fill:#2a1147,stroke:#a855f7,stroke-width:1.5px,color:#ede9fe,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; P([Producer\n消息 + RoutingKey]) --> EX{Exchange 类型?} EX -->|"Direct"| D[精确匹配\nRoutingKey == BindingKey] EX -->|"Fanout"| F[忽略 RoutingKey\n广播所有绑定队列] EX -->|"Topic"| T[通配符匹配\n* 单段 / # 多段] EX -->|"Headers"| H[消息头属性匹配\nx-match: all / any] class P startEnd; class EX condition; class D,F,T,H highlight; 去管理界面 Exchanges 页面点开一个 Exchange,看到 type 字段的值就是这四种之一。 ...

十一月 2, 2022 · 8 分钟 · 1549 字 · yaomingye

RabbitMQ 核心概念与 AMQP 协议

核心概念与 AMQP 协议 一、⚡ 问题切入:同步处理为什么不行? 先看一个电商系统里最常见的下单流程: @Service public class OrderService { @Transactional public Order createOrder(CreateOrderRequest request) { // 1. 扣减库存 inventoryService.deduct(request.getProductId(), request.getQuantity()); // 2. 创建订单 Order order = orderMapper.insert(request); // 3. 发送下单成功短信——这一步是同步的 smsService.sendOrderConfirm(request.getUserId(), order.getId()); // 4. 发送下单成功邮件——这一步也是同步的 emailService.sendOrderConfirm(request.getUserId(), order.getId()); // 5. 写入操作日志 operationLogService.record("CREATE_ORDER", order.getId()); return order; } } 一次下单请求,用户要等库存扣减、订单入库、短信发送、邮件发送、日志写入全部完成才能收到响应。短信调用运营商接口,邮件走 SMTP,日志写入数据库——这三步加起来可能要 500ms ~ 2s。用户在前端点完"提交订单"后盯着屏幕转圈,体验糟糕。 有人会说:“那简单,开个线程异步执行不就行了?” // 线程池异步——似乎解决了问题 executorService.submit(() -> smsService.sendOrderConfirm(userId, orderId)); executorService.submit(() -> emailService.sendOrderConfirm(userId, orderId)); 但这引入了一连串新问题: ...

十一月 1, 2022 · 7 分钟 · 1466 字 · yaomingye
Cat Radio