上一期《7天学会Redis》已经完结,本周开始整理RocketMQ,欢迎指正,点赞关注哦~

《RocketMQ研读》Day1:RocketMQ 基础架构与核心概念

《RocketMQ研读》Day2:生产者的设计与实践

《RocketMQ研读》Day3:消费者的设计与实践

《RocketMQ研读》Day4:消息存储与复制机制

《RocketMQ研读》Day5:事务机制与顺序消息

《RocketMQ研读》Day6:高级特性与实战案例

《RocketMQ研读》Day7:高可用与集群部署


Day1:RocketMQ 基础架构与核心概念

一、消息中间件概述与选型对比

1.1 为什么需要消息队列?

消息队列(Message Queue)是一种跨进程的通信机制,用于上下游传递消息。在分布式系统中,消息队列主要解决以下问题:

  • 应用解耦:将不同业务逻辑解耦,例如订单系统与库存系统、物流系统等,通过消息队列传递订单消息,各系统独立处理,互不影响。

  • 异步处理:将非核心流程异步化,提高系统响应速度。例如用户注册后,发送邮件、短信等操作可以异步处理。

  • 流量削峰:在高并发场景下,将请求放入消息队列,系统按照处理能力消费消息,避免系统被压垮。

  • 数据分发:一对多消息发布,多个系统订阅同一Topic获取数据。

1.2 主流消息中间件对比

维度 RocketMQ 4.9+ Kafka 2.8+ RabbitMQ 3.9+ Pulsar 2.9+
存储模型 CommitLog + 索引文件 Partition 分区日志 Queue + Exchange 分层存储(BookKeeper)
消息协议 自定义二进制协议 自定义协议 AMQP 0.9.1/1.0 自定义协议
消息类型 事务/顺序/延时/批量 普通/批量 普通/事务 多种语义支持
消费模式 推/拉 推/拉 推/拉
延时消息 固定级别(18级) 不支持原生 插件支持 原生支持
消息回溯 按时间/偏移量 按偏移量 不支持 按时间

二、核心组件深度解析

2.1 整体架构设计哲学

        RocketMQ采用"轻中心化"架构,NameServer无状态,Broker主从分离,实现高可用和高性能的平衡。

2.2 NameServer:轻量级路由中心

NameServer是RocketMQ的"轻量级注册中心",设计目标是简单、高效、无状态。

职责边界:

  • 只做路由注册与发现,不存消息、不处理心跳逻辑、不进行选主

  • AP 架构:节点间不通信,通过客户端轮询实现最终一致

心跳注册流程

  • 无状态 NameServer:仅维护 Topic → Broker IP 路由表(内存 Map),无选举、无状态、无持久化
  • Broker 主动注册:Broker 启动时向所有 NameServer 发送 RegisterBrokerRequest,之后每 30s 发送心跳(HeartbeatRequest)。

2.3 Broker:消息存储与转发的引擎

Broker 内部架构:

┌─────────────────────────────────────────────────────────┐
│                    Broker Server                        │
├─────────────────────────────────────────────────────────┤
│  ┌─────────────┐  ┌─────────────┐  ┌─────────────┐     │
│  │ 通信层       │  │ 业务处理层  │  │ 存储层       │     │
│  │ Netty Server│  │ 消息处理器   │  │ CommitLog   │     │
│  │ Remoting    │  │ Producer    │  │ ConsumeQueue│     │
│  │ 协议编解码   │  │ Consumer    │  │ IndexFile   │     │
│  └─────────────┘  └─────────────┘  └─────────────┘     │
├────────────────────────────────────────────────────────┤
│  ┌──────────────────────────────────────────────────┐  │
│  │                HA 复制模块                       │   │
│  │  Master/Slave 同步/异步复制                      │   │
│  │  DLedger 多副本强一致                            │   │
│  └──────────────────────────────────────────────────┘  │
├────────────────────────────────────────────────────────┤
│  ┌─────────────┐  ┌─────────────┐  ┌─────────────┐     │
│  │ 过滤服务     │  │ 事务回查    │  │ 统计监控    │      │
│  │ Tag/SQL过滤  │  │ 定时检查    │  │ 指标上报    │      │
│  └─────────────┘  └─────────────┘  └─────────────┘     │
└────────────────────────────────────────────────────────┘

职责边界:

  • 消息持久化:顺序写 CommitLog,异步构建 ConsumeQueue 与 IndexFile

  • 复制与 HA:Master-Slave 同步/异步复制,支持 Dledger 强一致

  • 消费位点管理:集群模式 Offset 存 Broker,广播模式 Offset 存本地

  • 限流与降级:根据内存、磁盘、队列深度进行生产/消费流控

存储架构全景图:

Broker 配置

# Broker 核心配置
brokerClusterName=DefaultCluster
brokerName=broker-a
brokerId=0  # 0=Master, >0=Slave
brokerRole=SYNC_MASTER  # SYNC_MASTER/ASYNC_MASTER/SLAVE
flushDiskType=ASYNC_FLUSH  # SYNC_FLUSH/ASYNC_FLUSH

# 存储路径
storePathRootDir=/data/rocketmq/store-a
storePathCommitLog=/data/rocketmq/store-a/commitlog
storePathConsumeQueue=/data/rocketmq/store-a/consumequeue

# 文件大小
mapedFileSizeCommitLog=1073741824  # 1GB,SSD 可设为 2GB
mapedFileSizeConsumeQueue=6000000  # 6MB,约 30 万条索引

# 刷盘参数
syncFlushTimeout=5000  # 同步刷盘超时
putMsgIndexHightWater=600000  # ConsumeQueue 构建水位

# Broker 对外服务端口
listenPort=10911  # Broker 端口
haListenPort=10912  # HA 端口

2.4 生产者核心要点详解

  1. 动态寻址与路由

    • 生产者不直接连接Broker,而是定时(默认 30 秒)从NameServer集群拉取路由信息,确保能感知到Broker的上下线。

    • 发送时,根据路由表将消息投递到Topic对应的具体MessageQueue上,而MessageQueue分布在不同的Broker上。这是实现分布式和水平扩展的基础。

  2. 队列选择与负载均衡

    • 默认策略是轮询(Round Robin) 选择队列,这保证了消息在Broker集群间尽可能均匀分布,避免单个Broker过热。

    • 自定义策略:可依据消息Key(如订单ID)哈希到特定队列,这是保证局部顺序的前提。

  3. 高可用发送机制

    • 失败重试:当向某个Broker发送失败时,会自动重试到其他Broker的队列。在同步/异步模式下默认重试 2 次(共 3 次),可通过 retryTimesWhenSendFailed 配置。

    • 发送方式三选一

      • 同步发送:等待Broker响应确认,可靠性最高。

      • 异步发送:立即返回,通过回调处理响应,兼顾吞吐和可靠性。

      • 单向发送:只管发,不等待响应,性能最高,可能丢失。

2.5 消费者核心要点详解

  1. 订阅与消费模式

    • 订阅:消费者必须订阅一个Topic,并可使用TagSQL表达式进行消息过滤。

    • 消费模式

      • 集群模式(Clustering)默认且最常用。同一个Consumer Group下的多个消费者共同分担消费所有消息,每条消息只被组内一个消费者消费。系统内部会自动进行消费负载均衡,在消费者增减时重新分配队列。

      • 广播模式(Broadcasting):组内每个消费者都消费全部消息。

  2. Pull VS Push

    • 主流的DefaultMQPushConsumer实际上是基于长轮询的Pull模拟的“推”效果。它平衡了实时性和服务端压力。

    • DefaultMQPullConsumer 已废弃(自 RocketMQ 4.6+),推荐使用 LitePullConsumer(更轻量、可控)。
  3. 并发与顺序

    • MessageListenerConcurrently:并发消费,性能高,但顺序无法保证。

    • MessageListenerOrderly:顺序消费,它会锁定当前队列,确保一个队列在同一时刻只被一个消费线程处理。需要注意

      • 如果消费失败,会阻塞该队列后续消息(直到成功或跳过)。
      • 需要消费者主动返回 SUSPEND_CURRENT_QUEUE_A_MOMENT 来控制重试
  4. 可靠性基石:ACK与重试

    • 消费成功后,必须向Broker返回CONSUME_SUCCESS进行确认。

    • 若消费失败(异常、超时或返回RECONSUME_LATER),消息会进入重试队列,延迟一段时间后再次投递。重试16次失败后,消息转入死信队列,等待人工干预。

三、核心概念体系深度解析

3.1 Topic(主题)—— 消息的逻辑分类

Topic 的关键特性

特性 说明
逻辑隔离 Topic 之间完全隔离。发送到 Topic_A 的消息,订阅了 Topic_B 的消费者绝对收不到
一对多通信 一个 Topic 可以被多个生产者发送,也可以被多个消费者组订阅。这是发布/订阅模式的精髓。
命名规范 Topic 名称需全局唯一,通常由系统名+业务领域组成,如 TRADE_ORDEROMS_PAYMENT
物理实现 Topic 的物理承载物是队列。创建 Topic 时,需指定其下的 MessageQueue 数量。消息实际存储和分布在各个队列中。

Topic 创建策略

3.2 Message Queue(消息队列)—— Topic的物理分区

        Message Queue是Topic物理实现单元。一个Topic中的所有消息,并非存储在一个大文件中,而是分散存储在该Topic下的多个Queue中

比喻示例

  • Topic(主题):一个名为“订单处理”的任务清单(逻辑分类)。

  • Message Queue(消息队列):这份清单被拆分到4个并行的“工作篮”里。每个工作篮里都有一批按顺序排列的待办任务。

  • 生产者:负责将新的“订单任务”(消息)均匀地放入这4个工作篮。

  • 消费者:多个“处理员”(消费者实例)可以同时从不同的工作篮中领取并处理任务,实现并行作业。

Topic 与 MessageQueue 的关系:逻辑与物理的映射

所以,Topic到Queue是“一对多”,Queue到Broker是“多对多”分布

3.3 Tag(标签)—— 二级消息过滤

Tag 的关键特性与设计原则

特性 说明 类比与意义
轻量级过滤 Tag 设计为一个简短的字符串(通常1个单词),存储在消息属性中,过滤在Broker端完成,网络传输和客户端负载极小。 快递分拣员看一眼包裹上的城市标签(Tag)就能快速分拣,无需拆箱检查内容(消息体)。
订阅时过滤 过滤逻辑发生在消息从Broker投递给消费者之前,无关消息不会到达消费者,节约了带宽和客户端资源。 上海配送站从一开始就不会接收到送往杭州的包裹,无需自己再丢弃。
灵活表达式 支持 *(全部)、TAGA(单个)、TAGA || TAGB(多个,关系)的订阅表达式。 上海站可以只要“上海”的包裹,也可以订阅“上海 || 苏州”,接收两个城市的包裹。
非必须字段 生产者和消费者都可以不使用 Tag。消息可以没有 Tag(null),消费者也可以订阅全部(*)。 有些包裹不需要按城市分拣(无Tag),或者有个总仓需要处理所有包裹(订阅*)。

Tag 的数据流

流程解读

  1. 消息存储Tag 与消息体一同被持久化存储,但Broker会为其建立额外的过滤索引,使得后续匹配非常快速。

  2. 过滤点:关键的过滤动作发生在 Broker 端,在准备向消费者推送消息时即时进行判断。这是一个服务端过滤模型,是最高效的方式。

  3. 过滤结果:只有完全匹配的消息才会被放入网络包中发送给消费者。不匹配的消息对消费者完全透明,仿佛不存在一样。

3.4 Offset(偏移量)—— 消费位置管理

Offset 的本质:消息队列的消费坐标

你可以将 Offset 理解为消息在队列中的“唯一地址”或“页码”

  • 在队列中:每个 MessageQueue 中的消息都按到达顺序存储。Offset 就是一个从 0 开始单调递增的 长整型数字,唯一标识了每条消息在队列中的逻辑位置。Offset=5 表示这条消息是这个队列里的第6条消息(从0开始计数)。

  • 对消费者而言:消费者需要记录自己当前已经成功消费到了哪个位置。这个记录下来的 Offset,就是 “消费位点” 。它回答了“我下次应该从哪儿开始消费?”这个关键问题。

核心比喻

  • MessageQueue:一本只有页码、没有目录的小说(消息流)。

  • Offset:这本小说每一页的页码

  • 消费者:一个读者

  • 消费位点:读者夹在书里的书签,记录了他上次读到了哪一页。下次他只需从书签的下一页继续读即可。

        消费位点的管理是消费者客户端和Broker协同工作的结果,其核心目标是保证 “至少一次” 的可靠消费语义。下图展示了在集群模式下,位点管理的两种主要方式及其数据流:

核心流程解读

  1. 位点提交 (Commit):消费者成功处理一批消息后,必须主动向Broker报告自己最新的消费位置。这是一个 “确认消费” 的关键动作。提交可以是同步的,但通常是异步定时批量提交以提升性能。

  2. 位点恢复 (Fetch):当消费者启动,或发生重平衡(例如新增了一个消费者实例)后,新的消费者需要知道自己应该从哪个位置开始消费。它会向Broker发起查询,请求该消费者组在该 MessageQueue 上最后提交的位点,然后从这个位置开始(或往后一点)拉取。

  3. 两种存储模式

    • Broker端存储(集群模式默认):如上图场景一,位点集中存储在Broker上。优点是全局一致,同一个消费者组内的任何实例都能看到统一的进度。这是最常用、最安全的方式。

    • 客户端本地存储(广播模式默认):如上图场景二,每个消费者实例独立地将位点保存在自己的磁盘上。因为广播模式下每个实例消费全量消息,进度彼此独立。

3.5 Consumer Group(消费者组)—— 负载均衡单元

Consumer Group 的本质:逻辑订阅单位与负载均衡边界

你可以将 Consumer Group 理解为一组完成相同任务的“工人团队”

  • 逻辑定义:一个 Consumer Group 是一个由多个消费者实例(进程或线程)组成的逻辑集合。这些实例共享同一个 Group ID

  • 核心规则:在 RocketMQ 的集群模式下,订阅的同一个Topic的每条消息,只会被投递给同一个Consumer Group内的任意一个实例消费。这是实现“负载均衡”而非“广播”的基础。

  • 团队比喻

    • Consumer Group = “华东区订单处理组

    • 组内的每个 Consumer Instance = 组里的一个个“订单处理员

    • Topic = “订单消息流

    • MessageQueue = 订单流被分成的多个“订单分包

    • 工作方式:订单流(Topic)被分成多个包(Queue),组内的处理员们(Instances)共同瓜分这些包,每人负责处理其中几个,每个包只由一个人处理。这样,整个团队就能并行处理所有订单。

Consumer Group 的关键特性与核心作用

特性 说明
负载均衡单元 负载均衡发生在Consumer Group内部。系统会自动将Topic下的所有MessageQueue,尽可能平均地分配给组内的各个消费者实例。
消息投递语义 定义了消息是“集群消费”(单播)还是“广播消费”。
进度管理单元 消费进度(Offset)以Consumer Group为单位进行管理和持久化。组内的所有实例共享和提交同一个消费位点。
扩缩容的基本单位 通过增加或减少组内的消费者实例数,即可动态调整该组的整体消费能力。

Consumer Group 与重平衡(Rebalance)机制

重平衡是 Consumer Group 实现动态负载均衡的关键过程。

  • 触发时机:当 Consumer Group 内的消费者实例数量发生变化时(如实例启动、宕机、网络隔离、手动缩容)。

  • 核心动作:RocketMQ 会立即触发重平衡,根据当前在线的实例数,按照既定策略(如平均分配、一致性哈希等)重新分配每个实例负责的 MessageQueue 集合。

  • 重要性:这个过程是自动且透明的,它保证了:

    • 高可用:一个实例挂了,它负责的队列会立刻被分配给其他活着的实例,消费不中断。

    • 弹性伸缩:新增实例时,它能自动分担一部分队列,提升整体消费能力。

  • 注意点:重平衡期间,队列的分配会短暂失效,可能导致短暂的消费暂停。过于频繁的重平衡(如实例网络不稳定)会影响服务稳定性。

Logo

汇聚全球AI编程工具,助力开发者即刻编程。

更多推荐