一口气读完 RocketMQ 架构

上一篇我讲解了 RabbitMQ 架构(可以通过文末的合集查看),这次我们继续讲解另外一款消息队列 RocketMQ 的架构。

RocketMQ 在功能、稳定性、性能层面都比 RabbitMQ 的表现更好

1. RocketMQ 系统架构简介

RocketMQ 由 Broker、NameServer、Producer、Consumer 四大组件组成。如下图所示:

图片

简单概括下四大组件的核心要点:

组件 技术要点
NameServer 轻量级注册中心
Broker 消息存储
Producer 同步、异步、单向多种发射方式
Consumer Push/Pull 双模式

RocketMQ 的 NameServer 负责元数据的存储

它是一个独立的进程,扮演着集群“中枢神经系统”的角色,其核心作用是为生产者(Producer)和消费者(Consumer)提供路由信息,帮助它们找到对应的 Broker 地址。

Broker 在启动的时候会主动连接 NameServer,将自己的元数据信息上报给 NameServer,每隔 30 秒还会上报一次元数据(心跳包)。

核心内容包含 Broker 的地址、名称、BrokerId、主节点地址、该 Broker 上的所有 Topic 的队列配置等。

NameServer 会将 Broker 的元数据信息缓存到本地路由表,供 Producer/Consumer 拉取,实现动态路由与故障感知。

生产者在发送数据的时候,会指定 Topic 或 MessageQueue,只指定 Topic 时,Producer 用负载均衡算法挑一个 MessageQueue;也可以直接指定 MessageQueue。

Broker 收到消息后,将消息顺序追加到 CommitLog 文件,如果文件大小超过固定大小(默认 1G),则会生成新的 CommitLog 文件,避免单个文件过大。

MessageQueue 只是逻辑分片,不会存储消息。ConsumeQueue 做逻辑分片,是消息的索引,指向 CommitLog 中消息的具体位置。

Producer 发送消息,启动时先跟 NameServer 集群中的其中一台建立长连接,并从 NameServer 中获取当前发送的 Topic 存在哪些 Broker 上,通过负载均衡从队列列表中选择一个 MessageQueue,然后与 MessageQueue 所在的 Broker 建立长连接从而向 Broker 发消息。

Consumer 跟 Producer 类似,跟其中一台 NameServer 建立长连接,获取当前订阅 Topic 存在哪些 Broker 上,然后直接跟 Broker 建立连接通道,开始消费消息。

Consumer 通过声明的 Group 进行分组拉取消息,消费者每拉一批消息,SDK 把最新 offset 写回 Broker(定时或手动 commit)。

2. RocketMQ 的网络协议

  • RocketMQ 5.0 之前:客户端与 Broker / NameServer 只支持 RocketMQ 私有 Remoting 协议(基于 Netty 的二进制协议,固定帧格式、序列化用 RocketMQ 自己的编码)。
  • RocketMQ 5.0 开始:官方 SDK 新增 gRPC 协议实现,同时保留旧的 Remoting 协议实现。

如下图所示:

图片

完整交互流程:

(1)Broker 与 NameServer 通信:通过 Remoting 协议向 NameServer(9876 端口)注册和发送心跳。

(2)客户端启动:通过 Remoting 协议连接 NameServer(9876 端口)获取 Broker 路由表。

(3)客户端与 Broker 通信:使用 Remoting 协议连接 Broker(10911 端口),或者直连 Proxy 的 8081 端口,将数据通过 gRPC 协议发送给 Proxy,最后还是通过 Remoting 协议转发给 Broker。

值得注意的是 Remoting 协议直接基于四层的 TCP 协议通信gRPC 基于七层的 HTTP/2 协议通信,不过 HTTP/2 底层也是基于 TCP。

Remoting 协议 和 gRPC 协议对比

维度 Remoting(私有协议) gRPC(开发库丰富,推荐)
性能 极致(私有协议优化) 稍低(HTTP/2 头部开销)
多语言支持 高成本(需重复实现) 低成本(官方/社区实现)
云原生集成 困难(需额外适配) 原生支持(Istio/K8s)
可观测性 需额外开发 原生支持(OpenTelemetry)
生态连接 封闭 开放(Service Mesh 等)

Remoting 适合 RocketMQ 内部高性能、低延迟的场景(如 Broker 间同步),而 gRPC 更适合面向用户和云原生的场景,两者不是替代关系,而是互补。

3. RocketMQ 的网络模块

RocketMQ 是基于 Netty 扩展出来的高性能网络通信框架,接下来我们来看看下面的原理图。

图片

RocketMQ 的 RPC 通信采用 Netty 作为底层通信库,并基于 Reactor 多线程模型进行了深度扩展和优化

Broker 中有一个Reactor 主线程(Netty BossGroup),Producer 和 Broker 建立 TCP 长连接时,Reactor 主线程会在端口上监听到客户端建立的请求。

然后处理 TCP 三次握手建立连接,创建并注册 SocketChannel 将 SocketChannel 注册到 selector 上。

Producer 和 Broker 里面都通过各自的 SocketChannel 维持长连接。Producer 通过 SocketChannel 发送消息给 Broker 中的 SocketChannel。

Broker 中还有一个Reactor 线程池(Netty WorkerGroup),里面的线程会监听到 SocketChannel 的网络数据,并将数据传递给 Worker 线程池中的一个线程进行预处理。

Reactor 线程池默认有三个线程。

Worker 线程池中的线程的核心职责:SSL 加密验证、编解码工作、检查空闲连接、管理网络连接等等。Worker 线程池默认有 8 个线程。

另外 Broker 有一个业务线程池:SendMessage 线程池,里面的线程专门用于处理 Worker 消息写入到磁盘。

该线程池的线程数是动态的,根据服务器的 CPU 核心数自动调整。这种设计使得 RocketMQ 能够高效处理大量并发请求,同时保持系统的稳定性和可扩展性。

4. RocketMQ 的存储模块

RocketMQ 的数据存储分为元数据存储和消息数据存储

4.1. 元数据存储

什么是 Broker 的元数据?

首先元数据是用来描述其他数据的属性。就像一本书的元数据是用来描述这本书的信息,如书名、作者、出版社、目录等。

那么 Broker 的元数据其实就是描述 Broker 中的核心数据的,如 Broker 的基础元数据: 名称、ID、集群名称、地址等,另外还有 Topic 的相关元数据。

如何存储 Broker 中这些重要的元数据呢?我们来看下原理图。

图片

首先每个 Broker 节点都会存储自己的元数据,然后他们会将自己的这些元数据上传到每个 NameServer 上(如果 NameServer 采用集群部署的方式,就会有多个 NameServer,各 NameServer 实例之间不进行相互通信)。

如果 Broker 有从节点,也会和主节点一样上传元数据。

所以每个 NameServer 上都会有所有 Broker 的元数据。即使某个 NameServer 宕机了,但其他 NameServer 有所有 Broker 元数据信息,整个集群还是能正常对外提供服务的。

另外需要注意的是这些数据都是放在 NameServer 的内存中,不会持久化存储。

NameServer 如何感知某个 Broker 宕机了?

Broker 会每隔 30 秒向所有 NameServer 发送心跳包,告诉每个 NameServer 自己还存活着。

每个 NameServer 在收到心跳包后都会更新这个 Broker 节点的最近一次的心跳时间。

NameServer 还会每隔 10 秒运行一个任务,用来检查每个 Broker 节点最近一次的心跳时间是否超过 120 秒没有更新过,如果超过了,则说明这个 Broker 节点宕机了。

Broker 主从节点的机制说明:

机制 说明
数据同步 主节点异步/同步复制消息到从节点
故障接管 主节点宕机时,从节点不会自动升主(需运维干预或依赖 RocketMQ DLedger 自动选主)
读扩散 消费者默认从主节点读取,高负载时可配置为从从节点读取数据,分担主节点压力。

4.2. 消息数据存储

跟 RocketMQ 存储相关的有三种文件:CommitLog、ConsumeQueue、IndexFile。如下图所示,CommitLog 和 ConsumeQueue 是 RocketMQ 存储体系的核心组成部分。

图片

(1)CommitLog:存储消息的主体内容,消息的内容不定长,单个文件默认 1G 大小。当文件写满后,写入下一个文件。每个 Broker 节点都有各自的 CommitLog 文件。

(2)ConsumeQueue:消息消费索引,用于提高消息消费的性能。

因 RocketMQ 基于主题的订阅模式,消息消费也是针对主题进行的,但是如果每次消费都要遍历 CommitLog 文件来检索对应主题的消息是非常低效的,所以才有了基于主题的 CommitLog 索引文件,也就是 ConsumeQueue 文件。

它的文件夹采用 topic/queue/file 三层组织结构。

每个文件采取定长设计,每一个条目共 20 个字节,分别为 8 字节的 CommitLog 物理偏移量、4 字节的消息长度、8 字节 tag hashcode,单个文件由 30 万个条目组成,可以像数组一样随机访问每一个条目,每个 ConsumeQueue 文件大小约 5.72MB。

(3)IndexFile:索引文件,提供了一种可以通过 key 或时间区间来查询消息的方法。

Index 文件的存储位置是:$HOME/store/index/{fileName},文件名 fileName 是以创建时的时间戳命名的,单个 IndexFile 文件大小固定约为 400MB,一个 IndexFile 可以保存 2000 万个索引。

4.3. 消息的刷盘机制

(1)同步刷盘:当 Broker 端收到消息后,只有将消息真正持久化至磁盘后,Broker 端才会真正返回给 Producer 端一个成功的 ACK 响应。同步刷盘对 MQ 消息可靠性来说是一种不错的保障,但是性能上会有较大影响,一般适用于金融业务,应用该模式较多。

(2)异步刷盘:当 Broker 端收到消息后,只要消息写入 PageCache 即可将成功的 ACK 返回给 Producer 端。消息刷盘采用后台异步线程提交的方式进行,降低了读写延迟,提高了 MQ 的性能和吞吐量。

5. RocketMQ 生产消费机制详解

5.1. Producer 生产消息

Producer 启动时先跟 NameServer 集群中的其中一台建立长连接,当需要发送消息到 Topic 或 MessageQueue 时,会根据从 NameServer 中拿到的路由信息找到要发送给哪个 Broker,然后通过负载均衡算法选择一个 MessageQueue,然后与队列所在的 Broker 建立长连接,并将消息发送给该 MessageQueue。

Producer 支持三种发送消息的形式

(1)单向发送(Oneway):发送消息后立即返回,Producer 不关心是否发送成功,也不会处理响应。

(2)同步发送(Sync):发送消息后,Producer 等待响应。

(3)异步发送(Async):发送消息后立即返回,Producer 会在自己提供的回调方法中处理响应。

5.2. Consumer 消费消息

Consumer 跟 Producer 类似,也是和 NameServer 建立长连接,获取路由信息,然后通过订阅的 Topic 找到对应的 Broker,然后直接跟 Broker 建立连接,开始消费消息,最后通过提交消费位点的形式来保存消费进度。

Consumer 支持三种消费消息的模式

(1)拉取模式(Pull):消费者主动向 Broker 发送拉取请求,指定要拉取的消息数量和偏移量(或时间戳),Broker 响应包含消息或空结果。

(2)推模式(Push):客户端与 Broker 建立长连接,并发送拉取消息的请求。如果当前没有新消息,Broker 不会立即响应,而是等待一段时间或直到有新消息到达再返回。

(3)无状态模式(Pop):在 RocketMQ 5.0 中,Pop 消费模式的设计核心在于将重平衡、位点管理及消息重试等任务转移至服务端处理,有效避免单点故障引起的消息积压,优化了整体消息处理效率和系统的水平扩展能力,且提升了系统的灵活性和扩展性。

消费者组概念

消费组在其中有着非常重要的作用,如果多个消费者设置了相同的 Consumer Group,我们认为这些消费者在同一个消费组内。

Apache RocketMQ 支持两种消费模式

(1)集群消费模式:当使用集群消费模式时,RocketMQ 认为任意一条消息只需要被消费组内的任意一个消费者处理即可。可以通过扩缩消费者数量,来提升或降低消费能力。

(2)广播消费模式:当使用广播消费模式时,RocketMQ 会将每条消息推送给消费组所有的消费者,保证消息至少被每个消费者消费一次。即使扩缩消费者数量也无法提升或降低消费能力。

6. RocketMQ 事务消息

RocketMQ 提供了事务消息的功能,采用 2PC(两段式协议)+ 补偿机制(事务回查)的分布式事务功能,通过这种方式能达到分布式事务的最终一致。原理如下图所示:

图片

事务消息发送步骤如下:

(1)发送方将半事务消息发送至消息队列 RocketMQ 版服务端。

(2)消息队列 RocketMQ 版服务端将消息持久化成功之后,向发送方返回 Ack 确认消息已经发送成功,此时消息为半事务消息。

(3)发送方开始执行本地事务逻辑。

(4)发送方根据本地事务执行结果向服务端提交二次确认(Commit 或是 Rollback),服务端收到 Commit 状态则将半事务消息标记为可投递,订阅方最终将收到该消息;服务端收到 Rollback 状态则删除半事务消息,订阅方将不会接收到该消息。

事务消息回查步骤如下:

(5)在断网或者是应用重启的特殊情况下,上述步骤 4 提交的二次确认最终未到达服务端,经过固定时间后服务端将对该消息发起消息回查。

(6)发送方收到消息回查后,需要检查对应消息的本地事务执行的最终结果。

(7)发送方根据检查得到的本地事务的最终状态再次提交二次确认,服务端仍按照步骤 4 对半事务消息进行操作。

7. RocketMQ 架构小结

RocketMQ 作为一款优秀的消息中间件,其架构设计充分考虑了高可用、高性能和高扩展性。通过 NameServer、Broker、Producer 和 Consumer 四大组件的协同工作,实现了可靠的消息传递和灵活的消费模式。

核心优势包括

  1. 高可用性:通过多副本、主从切换和故障转移机制,确保系统稳定运行。
  2. 高性能:基于 Netty 的网络通信和优化的存储机制,支持高并发处理。
  3. 高扩展性:支持水平扩展,可根据业务需求动态调整资源。
  4. 事务支持:提供分布式事务解决方案,保证数据一致性。
  5. 多种消费模式:支持集群消费、广播消费等多种模式,满足不同场景需求。

RocketMQ 的架构设计为大规模分布式系统提供了可靠的消息传递基础,是企业级应用中常用的消息中间件之一。