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

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

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

MQTT 协议

📡 MQTT 协议:角色体系、Broker 原理与 QoS 分级机制全解析 问题切入:一个智能家居的消息困境 假设你要开发一个智能家居系统,包含以下设备: 10 个温湿度传感器,每 5 秒上报一次数据 5 个智能插座,需要接收开关指令并上报当前功率 1 个手机 App,需要实时看到所有设备的状态,并能下发控制指令 你的第一反应可能是用 HTTP:传感器 POST 数据到服务端,App 轮询拉取最新状态。但很快问题就来了: 传感器数量 × 上报频率 = 10 × (1 / 5s) = 2 QPS 的上报请求 App 轮询最新状态 = 1 × (1 / 2s) = 0.5 QPS 的查询请求 设备控制指令 = App POST 到服务端,服务端再推给设备... HTTP 是请求-响应模式,服务端无法主动向设备推送指令。如果让设备轮询指令,延迟高且浪费带宽。而且温湿度传感器是低功耗设备(电池供电的 ESP8266),HTTP 的 TCP 三次握手 + Header 开销太大。 这就是 MQTT(Message Queuing Telemetry Transport,消息队列遥测传输协议)解决的问题:它是一个 发布-订阅模式 的轻量级消息协议,专为低带宽、高延迟、不可靠网络下的物联网设备通信而设计。 MQTT 的角色体系 MQTT 协议定义了三种角色。大部分文章对它们的介绍含糊其词,这里逐个讲清楚。 ...

十月 1, 2022 · 14 分钟 · 2882 字 · yaomingye
Cat Radio