目录
关于其他消息队列的文章: 带你轻松学习Kafka-CSDN博客https://blog.csdn.net/2401_88959292/article/details/149309549?sharetype=blogdetail&sharerId=149309549&sharerefer=PC&sharesource=2401_88959292&spm=1011.2480.3001.8118 带你轻松学习RabbitMQ-CSDN博客https://blog.csdn.net/2401_88959292/article/details/149294023?sharetype=blogdetail&sharerId=149294023&sharerefer=PC&sharesource=2401_88959292&spm=1011.2480.3001.8118
一、核心组件/概念
1.命名服务中心(NameServer)
主要是负责路由管理,负责建立和消费者所绑定的Topic到Broker之间的映射;然后就是维护Broker的状态,在Broker宕机时及时清理相关路由。
2.消息代理(Broker)
主要负责消息的存储以及转发。
3.主题(Topic)
消息的逻辑分类,每个Topic都有多个物理分片(Partition),用于。
4.标签(Tag)
过滤消息的标识。
二、消息流程


流程说明:
1.生产者发送消息
生产者连接NameServer并注册自己的信息,定期拉取Topic路由数据;
然后生产者根据负载均衡为消息的Topic选择一个Partition;
最后发送消息。
消息在发送前需要指定以下三个属性:
- Topic
- Tag
- 消息实体
2.Broker接收消息
Broker接收消息后,会先将消息元数据写入CommitLog Buffer(内存缓冲区),如果此时是同步模式则会在写入完成后返回ACK确认,异步模式则会直接返回;
因为是写在内存里的,所以也需要刷盘策略:
- 异步刷盘(默认):定期将CommitLog Buffer中的消息批量写入磁盘的CommitLog文件。
- 同步刷盘:每发送一条信息到CommitLog Buffer就立即写入磁盘中。
此环节的刷盘采用顺序写,每次都是追加模式写入CommitLog,这样会使得磁盘IO均为顺序IO,效率远高于RabbitMQ的随机写(因为RocketMQ的CommitLog是所有逻辑队列共用的,而RabbitMQ的存储文件是每个队列独立的)。 这是RocketMQ高吞吐量(百万级/秒)的核心原因之一。
3.Broker构建逻辑队列
为了让消费者快速找到指定Topic的Partition的消息,Broker会为每个Topic的Partition构建ConsumeQueue(也位于磁盘空间)。
Broker会解析CommitLog文件,为每条信息生成一个条目(包含其在CommitLog的起始偏移量offset、哈希值以及消息实体的大小)并写入ConsumeQueue文件中。
同时建立内存映射,将ConsumeQueue和CommitLog文件映射到PageCache(操作系统的内存缓存),后续消费者直接读取内存,加速索引查询效率。
4.消费者拉取消息
与生产者一致,消费者也要连接NameServer并注册自己的信息,定期拉取Topic路由数据;
然后根据负载均衡策略获取自己所要负责的Partition;
当消费者订阅Topic时,Broker会遍历ConsumeQueue,根据每个消息的Tag哈希值为其建立Bitset并存入filterMap;
每当消费者拉取信息时,会根据消费者指定的Tag的哈希值先查询filterMap,快速过滤掉不匹配的消息,然后查询到对应的offset,并到CommitLog文件中去读取对应的消息实体。
此处读取CommitLog文件采用了零拷贝的技术,这是RocketMQ高吞吐量的核心原因之一。
什么是零拷贝? 零拷贝是一种优化数据传输性能的技术,核心目标是减少不必要的 CPU 拷贝和上下文切换。它的名字听起来像“完全不拷贝”,实际是指避免数据在“用户态”和“内核态”之间来回复制(这些复制由 CPU 负责,非常耗时)。数据依然会拷贝,但主要由更高效的 DMA(直接内存访问) 技术完成。
传统 I/O 的问题:效率低下的 4 次拷贝 & 4 次切换 假设程序需要读取文件并发送到网络,传统流程如下:
4次上下文切换(用户态 ⇄ 内核态): read() 调用:用户态 → 内核态 → 用户态write() 调用:用户态 → 内核态 → 用户态 4 次数据拷贝: 步骤拷贝方向执行者① 磁盘 → 内核缓冲区磁盘 → 内核空间DMA② 内核缓冲区 → 用户缓冲区内核空间 → 用户空间CPU③ 用户缓冲区 → Socket 缓冲区用户空间 → 内核空间CPU④ Socket 缓冲区 → 网卡内核空间 → 网卡DMA
问题关键: 步骤②和③是完全多余的!数据只是从内核出去又回来,白白消耗 CPU 资源。 频繁切换状态(用户态/内核态)也会降低性能。而零拷贝的核心思路:跳过用户缓冲区,让数据直接在内核中流动。
2 次上下文切换: 只需一次 sendfile() 系统调用(替代 read() + write())。 2 次数据拷贝(均由硬件完成): 步骤拷贝方向执行者① 磁盘 → 内核缓冲区磁盘 → 内核空间DMA② 内核缓冲区 → 网卡内核空间 → 网卡SG-DMA
关键改进: 取消用户缓冲区参与:数据全程在内核中流动。 SG-DMA 的作用:网卡通过“描述符”(记录数据位置和长度)直接从内核缓冲区抓取数据,无需经 Socket 缓冲区。
以上是Kafka使用到的零拷贝技术; 而RocketMQ实际使用的是mmap() + write()的零拷贝技术仍有4次上下文切换以及3次数据拷贝操作,这也是RocketMQ吞吐量不如Kafka的原因之一。那么为什么RocketMQ不使用sendfile零拷贝技术呢? 主要是因为RocketMQ侧重于处理业务消息,这类消息通常大小更小,而mmap通常在小文件读写上性能更优;而且RocketMQ需要消息过滤、属性匹配等操作,因此可以使用 mmap 将文件映射到内存后,可直接在内存中解析消息属性,无需额外拷贝,而sendfile则无法做到这点。 可以说这是业务功能需求和吞吐量的权衡之下而做出的选择。
ps:如果想要进一步了解零拷贝技术,可以阅读以下文章:原来 8 张图,就可以搞懂「零拷贝」了 - 小林coding - 博客园https://www.cnblogs.com/xiaolincoding/p/13719610.html
当拉取消息成功后,消费者还会进行二次过滤(防止Tag哈希碰撞),并且此处被过滤掉的消息仍旧可以被其他消费者消费。
有同学可能会有疑问:
既然消息都已经出队到消费者这端了,为什么被过滤掉的消息还可以被其他消费者消费?
不要忘记了,ConsumeQueue只是逻辑队列,消息实体是存储在磁盘中的CommitLog文件里的,其他消费者仍可以通过更新offset来读取到这条消息。
那CommitLog文件中已经消费掉的消息实体什么时候才会删除?
与的消费一条删一条的策略不同,采用的是定期删除的策略,每条消息都有对应的有效时间,有效时间一过,不管是否被消费掉,都会直接从磁盘里删除。
5.消费者返回ACK
消费者收到消息后,需处理,并向Broker确认消息已处理(确保消息不丢失)。
Broker只有当收到了ACK之后才会更新offset来让消费者拉取下一条消息,因此如果消费失败消费者不需要做任何操作,下一条拉取的消息仍是相同消息(默认可重新投递16次)。
6.Broker维护offset
Broker在收到消费者发来的ACK确认之后,会持久化当前的offset,保存到磁盘中的OffsetStore文件中,使得消费者即使宕机或者重启后可以继续拉取后续的消息而以免重复消费。
三、优劣分析
(一)优点
1.高以及低延迟
- 顺序写:所有逻辑队列共享一份CommitLog文件,因此磁盘IO就是顺序读写的形式,大大加快了磁盘IO效率。
- 零拷贝:如上文所说。
- PageCache内存映射:CommitLog和ConsumeQueue均映射到操作系统PageCache,读取消息时直接从内存获取,避免磁盘IO。
2.事务的支持
RocketMQ支持发送半消息:
- 先发送半消息到Topic中,但此时的消息是不可被消费的
- 本地事务开始执行
- 生产者根据本地事务的执行情况向Broker提交对应的状态:
- COMMIT:半消息转为可消费状态,投递到真实Topic。
- ROLLBACK:删除半消息,流程终止。
- 若发送半消息后Broker未收到二次确认消息,则会触发事务回查
- Broker定时扫描滞留的半消息(默认间隔1分钟)。
- 向生产者发起回查请求,生产者需实现检查器返回本地事务最终状态。
- Broker根据回查结果执行COMMIT/ROLLBACK。
3.分布式架构扩展性良好
NameServer、Broker均采用无状态或可水平扩展的设计,避免了“单点瓶颈”,支持集群规模的无限扩展。
- 每个NameServer都存储了完整的路由数据,没有单点依赖,Broker通过轮询提交路由数据,这意味着可以直接添加NameServer节点而无需修改任何配置。
- 每个Broker都拥有自己的CommitLog和ConsumeQueue,Broker之间无共享存储,即使出现宕机也不会影响其他Broker节点正常运行。
- 每个Topic的Partition会均匀分布到多个Broker节点,因此即使出现宕机该Topic的消息传递仍可以通过其他Broker进行。
4.高可用机制
- NameServer、Broker均采用无状态设计,使得一台宕机也不会给其他节点造成任何影响。
- NameServer定期检查Broker心跳,及时剔除宕机Broker路由,避免消息的无效传递。
- Broker分为主从模式,可自由搭建多主多从的同步或异步模式,且主从都拥有读写权,使得即使主节点宕机了从节点也可以立即接替其职责。
- 定期刷盘使消息数据持久化。
5.消息高可靠性的保障
通过四重机制确保消息从生产者到消费者的全链路可靠,满足金融、电商等核心业务的“零丢失”需求:
- 持久化:
- CommitLog:所有消息都写入磁盘(顺序写),确保Broker宕机后消息不丢失;
- ConsumeQueue:Topic-Queue的索引文件,存储在磁盘,确保索引不丢失;
- offset:写入磁盘,确保消费者宕机或重启可以正常拉取后续消息;
- 生产者确认: 生产者支持同步发送、异步发送,确保消息到达Broker。
- 消费者确认: 消费者采用Offset确认机制,处理完消息后发送
ACK,Broker更新Offset;若消费者宕机,下次从上次的Offset继续消费,确保至少一次交付; - 高可用集群: 采用Master-Slave主从副本机制,Master处理读写请求,Slave异步同步Master的数据;Master宕机后,Slave自动晋升为新Master。
6.不同系统之间实现资源隔离
通过Namespace(类似RabbitMQ的Virtual Host)和水平扩展,,隔离不同业务系统的资源(Topic、Queue、用户权限),避免资源冲突。
7. 消息可回溯消费
消费者可通过以下方式重新消费过去的消息:
- 按时间戳回溯:调用
resetOffsetByTime方法,重置消费者的offset到该时间戳之后的第一条消息。 - 按Offset回溯:通过
seek方法,指定某个队列的offset,从该位置开始消费。 - 重置消费起点:在消费者启动时,通过
setConsumeFromWhere配置消费起点。
(二)缺点
1.场景下的顺序消息瓶颈
RocketMQ支持顺序消息,但需要启用单线程,并且生产者指定消息的唯一业务标识Key都要一致,使得消息都哈希到同一个队列中。
但在高并发场景下,单线程处理会成为性能瓶颈导致消息堆积。
虽然可以朝扩展队列分区数量、拆分业务逻辑等方向优化,但优化后的性能仍然是不容乐观的。
四、最佳实践
使用Tag进行过滤,且一个消费者组仅订阅一个Topic
这样可以实现业务解耦,而且水平扩展时也不会影响到其他业务的MQ,再加上Tag快速过滤无关消息,整个性能大大提升。
注意!必须避免以下情况:
- 一个消费者组订阅多个不相关Topic,会导致业务耦合、负载不均衡。
- 一个Topic被过多消费者组订阅会导致Broker压力大,增加Broker的网络与磁盘IO。
- 不使用Tag过滤,订阅所有消息,会导致无关消息传输,浪费网络带宽与消费资源。
- 一个消费者组里的消费者使用不同的Tag,Broker默认该组的全部消费者使用最后一个消费者的Tag,这会使需要的消息被过滤导致业务逻辑错误。
- 一个消费者组里的消费者订阅不同Topic,这是违反消费者组设计初衷的(统一处理同一业务场景的消息),并且会使业务耦合,而且使负载均衡失效,维护成本也大大增加。
五、总结
RocketMQ的特性决定了其适合大规模高并发、强可靠、需要复杂业务处理的场景,不适合小规模低延迟、简单通信或资源受限的场景。
1. 适合的场景
- 大规模高并发场景: 如电商秒杀、日志收集、用户行为数据采集等。吞吐量可达百万级/秒,且支持线性扩展,能支撑秒杀等高并发场景的消息传输需求。
- 强可靠要求的金融场景: 如交易清算、转账支付、账户余额变更等。RocketMQ通过四大机制和事务消息,确保消息零丢失和跨系统事务原子性,满足金融级可靠性要求。
- 需要顺序处理的场景: 如订单流程、物流状态更新等。RocketMQ通过Key哈希到同一Queue,保证消息按业务顺序处理,避免状态混乱。
- 分布式系统服务解耦场景: 如微服务间通信、跨系统消息传递等。RocketMQ实现了生产者与消费者的彻底解耦,且支持消费者组负载均衡,适合分布式系统的服务解耦需求。
2. 不适合的场景
- 小规模低并发场景: 如创业公司内部系统、小流量应用等。RocketMQ架构复杂,且资源占用高,小规模场景下会造成资源浪费。
- 极度低延迟场景: 如高频交易、实时视频等。RocketMQ的消息延迟为毫秒级,无法满足微秒级延迟需求,此类场景更适合ZeroMQ、NATS等轻量级消息中间件。
- 简单点对点通信场景: 如内部通知系统、简单消息传递等。RocketMQ的
Topic-Queue模型比RabbitMQ的Direct交换机更复杂,且不需要顺序消息、分布式事务等复杂特性,简单场景下用RabbitMQ更便捷。
码文不易,留个赞再走吧
原文链接: 带你轻松学习RocketMQ 作者: Yilena
如果这篇文章对你有帮助,欢迎分享给更多人!
部分信息可能已经过时










