一、 基础理论:
- 消息队列应用场景:削峰填谷、异步解耦、可靠消息传递(延迟消息、事务消息)
- 通信协议:JMS(Java专属,不支持跨语言。代表:ActiveMQ)、AMQP(定义线路层协议,支持多种路由模式。代表:RabbitMQ、RocketMQ)、MQTT(轻量级发布/订阅模型,适合IoT设备)、Kafka Protocol(自定义二进制协议,高吞吐,基于TCP。代表:Kafka)
- 消息模型:点对点(核心载体:Queue,即队列模型)、发布订阅(核心载体:Topic,即主题模型)
二、问题与解决方案
- 消息可靠性(消息丢失可能发生在生产、存储、消费三个阶段)
- 生产端:网络波动,导致消息未送达(解决方案:开启生产者确认机制)
- 存储端:MQ宕机导致内存消息丢失(解决方案:开启持久化)
- 消费端:消费者未处理完消息就确认(解决方案:手动ACK,消费完成后调用BasicAck)
- 消息幂等性(因网络抖动导致消息重发,从而消息被重复消费。消费者端实现)
- 唯一标识:消息携带业务唯一ID,消费前检查redis/数据库确认消息是否已处理
- 幂等方案:数据库(唯一键约束或者乐观锁)、缓存(用redis的setnx或hash记录消费状态)
- 消息顺序性(即需要按顺序处理消息)
- 队列区分:同一个业务ID,消息发送到同一个队列/分区,一个队列对应一个消费者
- 顺序发送/消费:RocketMQ通过MessageQueueSelector将消息路由到指定分区,消费者按顺序拉取
- 消息堆积(生产者发送消息大于消费者处理速度)
- 提高消费:增加消费者实例、优化消费逻辑(批量处理)、调整prefetch(RabbitMQ)或consumeThreadMin(RocketMQ)
- 扩大存储:使用惰性队列(RabbitMQ)或调整消息保留策略(Kafka)
- 死信队列与延迟队列
- 死信队列:处理无法消费的消息(如拒绝消费、消息过期、队列满),配置x-dead-letter-exchange,死信自动发送到指定队列
- 延迟队列:实现定时任务(如订单超时关闭)。通过“死信+TTL“(消息过期后转死信队列)或者插件(如RabbitMQ使用rabbitmq_delayed_message_exchange)实现
- 选型建议:RabbitMQ(服务解耦,任务队列)、RocketMQ(高并发,电商等场景)、Kafka(高吞吐,日志采集、实时流处理)
三、RabbitMQ(插件机制)
- 核心组件
- Exchange(交换器):接收生产者消息(生产者不会把消息直接投递到队列中),根据路由规则转发至队列,核心类型包括Direct、Topic、Fanout、Headers
- Queue:存储消息的容器,支持持久化、自动删除等属性,消息按FIFO原则等待消费
- Broker:消息中间件的服务节点/服务实例
- Producer:消息的发送方,通过信道(Channel)将消息发送至交换机,每条消息包含路由键(routing key)
- Consumer:消费者,从队列中拉取消息兵处理,需手动确认(Ack)消息完成消费
- 消息模型(流程总结:生产者通过信道将消息发送到交换机,交换机根据路由键和绑定关系(Binding)将消息路由到队列,消费者监听队列并处理消息):Producer -> Exchange -> routing key(Binding规则,告诉交换器应该投递到哪个队列) -> Queue(绑定多个Consumer) -> Consumer。
- Exchange的4种路由规则(交换器类型,优先推荐Direct和Topic)
- Direct:完全匹配routing key。适用于点到点精确路由
- Fanout:广播到所有绑定的队列。适用于消息广播
- Topic:按routing key模式匹配(*匹配一个、#匹配零个或多个)。适用于模糊匹配
- Headers:按消息头属性匹配。性能较差,不推荐
- AMQP的三大组件
- 交换器:把消息**路由(即投递)**到队列的组件
- 队列:用来存储消息的数据结构,位于硬盘或内存中。
- 绑定:一套规则,交换器根据路由规则将消息投递到对应的队列
- 工作模式:简单模式、work模式、pub/sub 发布订阅模式、Routing 路由模式、Topic主题模式
- 消息的可靠性(确保消息不丢失):
- 生产者端:confirm模式(开启后每条消息分配唯一ID,RabbitMQ接收并持久化后返回ACK,失败则返回Nack,生产者可重试)、Mandatory参数(设置为true是,若消息无法路由到队列,RabbitMQ会通过ReturnListener返回消息,避免丢失)
- 服务端:队列与消息持久化(队列声明时,durable=true,消息发送时delivery_mode=2,确保MQ重启后数据不丢失)、死信交换机(处理无法消费的消息,如TTL过期,队列达到最大长度等,绑定死信队列用于后续排查)
- 消费者端:手动ACK(关闭自动确认,处理完消息后调用basicAck,避免消费者崩溃导致消息丢失)、消费限流(控制消费者诶次接收消息的数量,防止过载)
- 消息的顺序性(多线程消费、重试导致顺序错乱)
- RabbitMQ 仅保证单个 Queue 内的 FIFO 顺序,但多消费者场景下可能出现乱序
- 解决方案:单个Consumer模式、分区有序(推荐,但注意失效模式)、内部内存队列(慎重)
- 死信队列
- 死信队列:处理无法消费的消息(如拒绝消费、消息过期、队列到达最大长度),配置x-dead-letter-exchange,死信自动发送到指定队列
- 消息进入死信队列的条件:消息者拒绝消费且不重新入队、消息TTL过期、队列到达最大长度
- 高可用与集群模式
- 普通集群:仅同步元数据(队列配置),消息存储在单节点,节点宕机则消息丢失。适合提升吞吐量但不保证高可用
- 镜像集群:队列同步到多个节点,主节点宕机后,从节点自动切换。保证高可用,但性能开销大(消息需同步到所有节点)
- Quorum队列:基于raft协议实现分布式一直,代替惊喜那个队列,支持更高的可靠性和扩展,是RabbitMQ3.8+的推荐方案。
- 常见问题与解决方案
- 重复消费:消费者实现幂等性,如使用数据库唯一索引、redis setnx或布隆过滤器判重
- 消息堆积:增加消费者节点、开启线程池异步处理、使用惰性队列将消息存储到磁盘
- 延迟队列:TTL+死信交换机(消息过期后进入死信队列被消费)、安装rabbitmq_delayed_message_exchange插件,直接诶发送延迟消息
- 消息顺序性:同一个业务的消息需发送到同一个队列,且消费者单线程处理(或者通过全局ID排序)
四、RocketMQ(基于主题模型,即发布/订阅模型)
- 核心组件
- NameServer:轻量级注册中心,存储Broker集群的路由细腻系,支持Broker动态注册与心跳检测
- Broker:消息存储与转发的核心,负责消息持久化、投递与查询。分Master/Slave,支持同步/异步复制(同步复制保证数据不丢失)
- Proxy:可选组件,计算和存储进行分离(客户端和Broker的代理,让Broker更专注存储消息)
- Topic:消息的一级分类,每条消息必须属于一个Topic。
- Tag:消息的二级分类,用于同一个Topic下区分不同类型的消息(优化查询与消费逻辑)。
- Queue:存储消息的物理实体,一个Topic包含多个Queue,每个Queue中的消息按顺序存储。
- Producer:生产者,同步/异步/单向多种发送方式
- Consumer:Push/Pull/Simple 三种消费模式
- 网络模块:RocketMQ 的 RPC 通信采用 Netty 作为底层通信库
- 消息类型:普通消息、定时/延时消息(两种消息本质相同,都是服务端根据消息设置的定时时间在某一固定时刻将消息投递给消费者消费)、顺序消息、事务消息
- 消息的存储机制(采用混合存储架构,通过CommitLog和ConsumeQueue实现读写分离)
- CommitLog(写):所有消息的物理存储文件,顺序写入,单个文件默认1GB
- ConsumeQueue(读):消息消费索引,存储Topic下每个Queue中的消息在CommitLog中的物理偏离量、长度和Tag Hash。
- IndexFile:提供基于消息key或时间分区的查询能力,采用hashMap的索引结构。
- 消费者分类:PushConsumer(推模式消费者,SDK 自动管理拉取和提交)、SimpleConsumer、PullConsumer(拉模式消费者,应用程序主动控制拉取过程,仅推荐在流处理框架场景下集成使用)
- RocketMQ特性
- 顺序消息:通过MessageQueueSelector确保同一个业务键的消息发往同一个Queue,消费者按Queue顺序消费。实现需要发送端指定messageGroup,消费端对Queue加锁(避免并发消费)。
- 事务消息:半消息机制,内部隐藏队列对消费者不可见,本地事务执行成功后半消息转为全消息,消费者可消费。
- 延时消息:预设18个延时级别,消息先存入系统Topic,时间轮定期检查到期后,投递至目标Topic
- 消息重试与死信队列:消费失败的消息会进入重试队列,转入重试16次后转至死信队列
- 负载均衡设计
- 生产者:通过轮询、哈希等策略选择Queue,优先避免失败Broker,提升发送效率
- 消费者:同一组消费者均分Queue(支持平均分配,一致性哈希等策略),若消费者数量超过Queue数,多余消费者将空转。
- 高可用设计
- 主从架构:Master-Slave模式下,Slave异步复制Master数据;RocketmQ4.5引入Dledger(基于Raft协议)实现主从自动切换,故障恢复时间缩短至秒级。
- 异步/同步复制:同步确保消息不丢失(需等待Slave确认),异步复制追求更高性能
- 零拷贝读写
- mmap(内存映射):RocketMQ 采用 mmap + write 的零拷贝组合,而非 sendfile。主要原因是需要在 Broker 层对消息进行处理(如延时消息、过滤、死信队列等),而 sendfile 仅支持文件描述符(fd)级别的传输,无法直接操作数据内容。
- 网络传输的零拷贝优化:RocketMQ 基于 Netty 实现网络通信,利用 Netty 的零拷贝特性进一步减少数据拷贝。
- 刷盘机制(把消息写入磁盘确保数据持久化,消息可靠性的实现策略):
- 同步刷盘(SYNC_FLUSH):强可靠优先,同步刷盘要求消息写入磁盘后才向Producer返回成功响应,适用于金融交易等零丢失场景。核心实现为GroupCommitService线程,通过组提交(Group Commit)优化性能,避免每条消息单独刷盘的低效问题。
- 异步刷盘(ASYNC_FLUSH):异步刷盘是RocketMQ默认策略,消息写入PageCache后立即返回成功,由后台线程FlushRealTimeService定时或定量刷盘,兼顾性能与可靠性
五、Kafka
- 核心组件:
- Producer:消息生产者,负责将数据发送到指定Topic。
- Consumer:消息消费者,从Topic中拉取消息兵处理。通过消费者组实现负载均衡,同一组内消费者不会重复消费同一分区数据(组内消费者竞争:同一个消费者组内的多个消费者会通过分区分配策略竞争Topic的分区资源,每个分区只能被组内一个消费者独占消费)。
- Broker:一个Kafka实例,存储Topic的分区副本。每个Broker可管理多个分区,处理读写请求并同步数据(多个 Kafka Broker 组成一个 Kafka Cluster)
- Topic:Producer 将消息发送到特定的主题(消息的逻辑分类,类似文件夹)。每个Topic可分为多个Partition,消息按顺序追加到Partition中。
- Partition(分区) : Partition 属于 Topic 的一部分(Topic的物理存储单元),每个Partition是一个有序、不可变的日志文件。一个 Topic 可以有多个 Partition ,并且同一 Topic 下的 Partition 可以分布在不同的 Broker 上,这也就表明一个 Topic 可以横跨多个 Broker (也就是同一个Topic会存在多个Broker中)。
- Replica:Partition的副本,分为Leader(读写数据)和Follower(同步Leader数据)。ISR机制确保只有同步的副本才能参与Leader的选举。
- Zookeepr:传统Kafka依赖zk管理元数据(Broker的存货、Topic-Partition的映射),Kafka4.0引入KRaft代替zk以提升性能。
- 架构:Producer -> Topic(包含多个分区) -> Broker、Cusomer(消费者组分配消费) -> Partition(分区,消费者竞争消费分区)
- 核心工作原理
- 消息生产流程:选择分区Partition(通过轮询、按key哈希或自定义策略,确定消息要写入哪个分区)、批量发送(累计多条消息后发送,减少网络请求)、ACK机制确认(根据acks参数控制可靠性,0:不等待Broker确认,但可能会丢失数据;1:仅Leader写入成功即返回;-1:Leader和ISR中的所有副本写入成功才返回,最可靠但延迟较高)
- 消息存储机制(消息持久化到磁盘):顺序写入(消息追加到日志文件末尾,避免随机IO,速度接近内存。补充:随机IO会导致更多的磁头寻道开销)、分段存储(将每个Partition分为多个Segment,后续通过二分法查找快速定位消息)、零拷贝技术(使用sendfile系统调用,数据直接从磁盘发送到网络,减少内核态/用户态的数据拷贝过程,提升吞吐量)
- 消费消费机制:消费者组(同一个组内消费者分工消费Topic的Partition,实现消息并行处理、消费者组实现负载均衡、不同组可重复消费同一个Topic)、offset管理(消费者记录到已消费的偏移量,消息处理完成后自动/手动提交,避免丢失)、重平衡(消费者数量发生变化/Topic分区数量调整时,重新分配Partition给消费者)
- 可靠性和高可用
- 多副本:每个Partition配置多个副本(如3个),分布在不同Broker上,避免单点故障。
- ISR机制:仅与Leader保持同步的副本(ISR)有资格参与Leader选举。若ISR中副本数低于min.insync.replicas,生产者写入分区Leader副本会失败
- Leader选举:Leader宕机时,从ISR中选举新Leader,禁用unclean.leader.election.enable,可避免*非同步的副本成为Leader导致数据丢失。
- 写入流程:Producer发送消息 -> 根据key计算决定写入哪个分区 -> 获取Parition的Leader节点 -> 消息写入Leader中的Partition分区 -> Follower副本同步消息 -> Leader返回ACK给Producer(leader确认)
- 消费者组与分区分配:一个Partition只能被一个Consumer Group内的一个Consumer消费,不同的Consumer Group互相独立,都能消费到全量消息。
- 零拷贝原理:使用系统调用
sendfile,数据直接从磁盘发送到网络,减少内核态/用户态之间的拷贝过程(减少拷贝次数) - 重试机制
- 生产者发送重试:解决发送未确认问题。生产者重试由 retries 和 retry.backoff.ms 控制,默认 retries=0(不重试),需手动开启。
- 消费者消费重试:消费者默认无自动重试(仅靠未提交偏移量重复拉取),但直接重试可能导致阻塞后续消息或顺序错误(重试需平衡消息处理失败和顺序性。解决方案:重试主题+死信队列)。
- Kafka特性:顺序I/O(磁盘顺序写入,远高于随机I/O)、PageCache(利用操作系统的缓存PageCache存储数据,减少JVM内存占用,提升读写速度)、批量与压缩(生产者批量发送并压缩消息,减少网络传输量)、分区与并行(通过增加Partition数量线性提升吞吐量,建议分区数等于消费者或生产者数量)