《RocketMQ研读》Day1:RocketMQ 基础架构与核心概念
上一期《7天学会Redis》已经完结,本周开始整理RocketMQ,欢迎指正,点赞关注哦~
《RocketMQ研读》Day1:RocketMQ 基础架构与核心概念
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 生产者核心要点详解
动态寻址与路由:
生产者不直接连接Broker,而是定时(默认 30 秒)从NameServer集群拉取路由信息,确保能感知到Broker的上下线。
发送时,根据路由表将消息投递到
Topic对应的具体MessageQueue上,而MessageQueue分布在不同的Broker上。这是实现分布式和水平扩展的基础。队列选择与负载均衡:
默认策略是轮询(Round Robin) 选择队列,这保证了消息在Broker集群间尽可能均匀分布,避免单个Broker过热。
自定义策略:可依据消息Key(如订单ID)哈希到特定队列,这是保证局部顺序的前提。
高可用发送机制:
失败重试:当向某个Broker发送失败时,会自动重试到其他Broker的队列。在同步/异步模式下默认重试 2 次(共 3 次),可通过
retryTimesWhenSendFailed配置。发送方式三选一:
同步发送:等待Broker响应确认,可靠性最高。
异步发送:立即返回,通过回调处理响应,兼顾吞吐和可靠性。
单向发送:只管发,不等待响应,性能最高,可能丢失。
2.5 消费者核心要点详解
订阅与消费模式:
订阅:消费者必须订阅一个
Topic,并可使用Tag或SQL表达式进行消息过滤。消费模式:
集群模式(Clustering):默认且最常用。同一个Consumer Group下的多个消费者共同分担消费所有消息,每条消息只被组内一个消费者消费。系统内部会自动进行消费负载均衡,在消费者增减时重新分配队列。
广播模式(Broadcasting):组内每个消费者都消费全部消息。
Pull VS Push:
主流的
DefaultMQPushConsumer实际上是基于长轮询的Pull模拟的“推”效果。它平衡了实时性和服务端压力。DefaultMQPullConsumer已废弃(自 RocketMQ 4.6+),推荐使用LitePullConsumer(更轻量、可控)。并发与顺序:
MessageListenerConcurrently:并发消费,性能高,但顺序无法保证。
MessageListenerOrderly:顺序消费,它会锁定当前队列,确保一个队列在同一时刻只被一个消费线程处理。需要注意
- 如果消费失败,会阻塞该队列后续消息(直到成功或跳过)。
- 需要消费者主动返回
SUSPEND_CURRENT_QUEUE_A_MOMENT来控制重试可靠性基石:ACK与重试:
消费成功后,必须向Broker返回
CONSUME_SUCCESS进行确认。若消费失败(异常、超时或返回
RECONSUME_LATER),消息会进入重试队列,延迟一段时间后再次投递。重试16次失败后,消息转入死信队列,等待人工干预。
三、核心概念体系深度解析
3.1 Topic(主题)—— 消息的逻辑分类
Topic 的关键特性
| 特性 | 说明 |
|---|---|
| 逻辑隔离 | Topic 之间完全隔离。发送到 Topic_A 的消息,订阅了 Topic_B 的消费者绝对收不到。 |
| 一对多通信 | 一个 Topic 可以被多个生产者发送,也可以被多个消费者组订阅。这是发布/订阅模式的精髓。 |
| 命名规范 | Topic 名称需全局唯一,通常由系统名+业务领域组成,如 TRADE_ORDER,OMS_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 的数据流

流程解读:
消息存储:
Tag与消息体一同被持久化存储,但Broker会为其建立额外的过滤索引,使得后续匹配非常快速。过滤点:关键的过滤动作发生在 Broker 端,在准备向消费者推送消息时即时进行判断。这是一个服务端过滤模型,是最高效的方式。
过滤结果:只有完全匹配的消息才会被放入网络包中发送给消费者。不匹配的消息对消费者完全透明,仿佛不存在一样。
3.4 Offset(偏移量)—— 消费位置管理
Offset 的本质:消息队列的消费坐标
你可以将 Offset 理解为消息在队列中的“唯一地址”或“页码”。
在队列中:每个
MessageQueue中的消息都按到达顺序存储。Offset就是一个从 0 开始单调递增的 长整型数字,唯一标识了每条消息在队列中的逻辑位置。Offset=5表示这条消息是这个队列里的第6条消息(从0开始计数)。对消费者而言:消费者需要记录自己当前已经成功消费到了哪个位置。这个记录下来的
Offset,就是 “消费位点” 。它回答了“我下次应该从哪儿开始消费?”这个关键问题。核心比喻:
MessageQueue:一本只有页码、没有目录的小说(消息流)。
Offset:这本小说每一页的页码。消费者:一个读者。
消费位点:读者夹在书里的书签,记录了他上次读到了哪一页。下次他只需从书签的下一页继续读即可。
消费位点的管理是消费者客户端和Broker协同工作的结果,其核心目标是保证 “至少一次” 的可靠消费语义。下图展示了在集群模式下,位点管理的两种主要方式及其数据流:
核心流程解读:
位点提交 (Commit):消费者成功处理一批消息后,必须主动向Broker报告自己最新的消费位置。这是一个 “确认消费” 的关键动作。提交可以是同步的,但通常是异步定时批量提交以提升性能。
位点恢复 (Fetch):当消费者启动,或发生重平衡(例如新增了一个消费者实例)后,新的消费者需要知道自己应该从哪个位置开始消费。它会向Broker发起查询,请求该消费者组在该
MessageQueue上最后提交的位点,然后从这个位置开始(或往后一点)拉取。两种存储模式:
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集合。重要性:这个过程是自动且透明的,它保证了:
高可用:一个实例挂了,它负责的队列会立刻被分配给其他活着的实例,消费不中断。
弹性伸缩:新增实例时,它能自动分担一部分队列,提升整体消费能力。
注意点:重平衡期间,队列的分配会短暂失效,可能导致短暂的消费暂停。过于频繁的重平衡(如实例网络不稳定)会影响服务稳定性。
更多推荐



所有评论(0)