发布时间:2026-07-21阅读(1)
消息系统我们在项目中经常使用,但是如何自己设计一个消息系统呢?一个消息系统需要有消息的生产者、消费者,还需要具备消息的存储功能,设计时要考虑如下问题:
我们从学习RocketMQ入手,来学习消息系统的设计。
二、RocketMQ架构RocketMQ的逻辑部署图:

RocketMQ逻辑部署图
RocketMQ的角色包含Broker、NameServer、Producer和Consumer。

RocketMQ角色
三、RocketMQ功能3.1 路由注册路由注册是指将Broker信息注册到NameServer中,以便生产者、消费者可以从NameServer中获取到Broker的信息,进行消息的发送或接收。这样做的好处是Broker在扩容、缩容时,开发者对Broker的变化无感知。

Broker注册流程
Broker端特殊说明
如果Broker宕机,NameServer无法收到心跳包,此时NameServer如何来剔除这些失效的Broker呢?
RocktMQ有两个入口来触发路由删除。
1)定时任务:NameServer每隔10s扫描路由表,检测上次心跳包时间戳与当前系统时间的时间差,如果时间差大于120s,会从路由表中移除该Broker相关的信息并关闭Socket连接。
2)Broker在正常被关闭的情况下,会执行unregisterBroker指令,NameServer收到后会移除该Broker相关的信息并关闭Socket连接。
为什么路由变化不会马上通知消息生产者,而要等120s呢?这是为了降低NameServer实现的复杂性,路由变化由发送端提供容错机制来保证消息发送的高可用性。
3.3 路由发现路由发现是指让生产者、消费者找到消息服务Broker的过程。RocketMQ路由发现是非实时的,当Topic路由出现变化后,NameServer不主动推送给客户端,而是由客户端定时拉取Topic最新的路由。

路由发现的流程
1)定时任务
生产者启动时会启动定时任务,每30s执行一次路由表的动态更新,流程如下:
2)消息发送时
生产者和消费者获取了broker地址和队列信息后,如何发送消息呢?
3.4 消息发送消息队列如何进行负载均衡Producer的代码在RocketMQ Client jar中,Producer启动后,会创建MQClientInstance实例,同时启动生产者和消费者。如果配置的NameServer地址相同,同一个JVM中的不同消费者和不同生产者在启动时获取到的MQClientInstane实例都是同一个。
一个Topic下可以配置多个消息队列,以提高服务吞吐量,生产者拿到路由信息后,需要确定发送的队列,通过过滤数据得到Master角色且具有写权限的队列,使用轮询的方式,用自增1的值对队列大小取模,确认要发送的队列,让消息平均落在不同的消息队列上。返回的消息队列按照broker、序号排序,格式如下:
[{"broker-Name":"broker-a","queueId":0},{"brokerName":"broker-a","queueId":1},{"brokerName":"broker-a","queueId":2},{"brokerName":"broker-a","queueId":3},{"brokerName":"broker-b","queueId":0},{"brokerName":"broker-b","queueId":1},{"brokerName":"broker-b","queueId":2},{"brokerName":"broker-b","queueId":3}]
选择消息队列算法:Math.abs(index ) % 消息队列大小。
消息发送如何实现高可用消息发送高可用主要通过两个手段:消息重试与Broker规避。
规避有故障的Broker
如果上一次根据路由算法选择的是宕机的Broker的第一个队列,那么随后的下次选择的是宕机Broker的第二个队列,消息发送很有可能会失败,再次引发重试,带来不必要的性能损耗。此时可以将该Broker进行规避,不再选择该Broker,提高发送消息的成功率。
批量消息发送如何实现一致性批量消息发送是将同一主题的多条消息封装成MessageBatch对象,一起打包发送到消息服务端,减少网络调用次数,提高网络传输效率。服务端按照同样的结构进行解析即可。
3.5 消息存储RocketMQ主要存储的文件包括CommitLog文件、ConsumeQueue文件、IndexFile文件。
CommitLog文件
CommitLog文件

ConsumeQueue文件
消息索引文件,主要存储消息Key与Offset的对应关系。生产者发送的消息包含key值,会用IndexFile存储消息索引,主要用于使用key或时间戳来查询消息。

存储文件
3.6 消息消费消息消费有两种模式:广播模式与集群模式。
消费者消费消息,是主动从服务端获取消息,通过“长轮询”方式达到Push效果的方法。消费者从Broker查询消息,当Broker服务端接到请求后,如果队列里没有新消息,并不立刻返回,而是通过循环HOLD住客户端一小段时间,在这个时间内有新消息到达,就利用现有的连接立刻返回消息给Consumer,如果没有消息则返回空结果。
消息确认如果消息监听器返回的消费结果为RECONSUME_LATER,则需要将这些消息发送给Broker延迟消息。如果发送ACK消息失败,将延迟5s后提交线程池进行消费。
消息过滤消费者使用tag对消息进行过滤。如果不需要消费某个Topic下的所有消息,可以通过指定消息的Tag进行消息过滤,比如:Consumer.subscribe("TopicTest", "tag1 || tag2 || tag3"),表示这个Consumer要消费“TopicTest”下带有tag1或tag2或tag3的消息(Tag是在发送消息时设置的标签)。在填写Tag参数的位置,用null或者“*”表示要消费这个Topic的所有消息。
参考书籍《RocketMQ技术内幕:RocketMQ架构设计与实现原理》
《RocketMQ实战与原理解析》
Copyright © 2024 有趣生活 All Rights Reserve吉ICP备19000289号-5 TXT地图HTML地图XML地图