rocketmq的面试题及原理
基础定位
Apache RocketMQ 分布式消息中间件,阿里开源,支持异步通信、削峰填谷、解耦、最终一致性事务消息、顺序消息、延迟消息、死信队列。
三大应用场景:流量削峰、系统解耦、异步通知、分布式事务。
核心角色
- Producer 生产者:发送消息
- Consumer 消费者:消费消息
- NameServer:注册中心(轻量,无状态)
- Broker 消息服务器:消息存储、转发,核心节点
- Broker 分为 Master(读写)、Slave(只读、同步备份)
完整消息发送流程
- Producer 启动,定时从 NameServer 拉取 Topic 路由(Broker 列表 + 队列信息)
- 根据负载均衡策略选出队列(MessageQueue)
- 建立 TCP 长连接,发送消息到对应 Broker
- Broker 接收消息:
- 写入 CommitLog(所有消息统一顺序文件)
- 异步分发到 ConsumeQueue(消息索引文件)
- 执行刷盘机制
- 返回发送结果给生产者
三大存储文件
CommitLog
- 所有 Topic 消息全部顺序写入同一个文件
- 优点:顺序写磁盘,性能极高;避免多个文件随机写
- 文件固定大小 1G,滚动创建
ConsumeQueue(消费队列索引)
每个 Topic 下每个队列对应一个 ConsumeQueue;
不存储完整消息,只保存索引:CommitLog 偏移量、消息长度、tag hashcode。
消费者消费时:查询 ConsumeQueue → 根据偏移量去 CommitLog 读取真实消息。
IndexFile
消息索引文件,根据msgKey快速检索消息,运维排查使用。
刷盘机制
- 同步刷盘 SYNC_FLUSH 消息写入内存,立刻持久化磁盘,返回成功;宕机不丢消息,性能偏低。金融场景选用。
- 异步刷盘 ASYNC_FLUSH(默认) 消息写入 PageCache 直接返回成功;后台线程异步刷入磁盘。性能高;机器断电存在消息丢失风险。
消息负载均衡
生产者负载均衡
发送消息轮询分发到 Topic 多个 MessageQueue;顺序消息使用自定义哈希选择队列。
消费者负载均衡
Push 模式底层是 Pull 实现;客户端执行负载均衡算法分配队列 算法:AllocateMessageQueueAveragely(平均分配,默认)
规则:同一个消费组,队列平均分配给各个消费者
经典约束:队列数量是消费并发上限! Topic 有 4 个队列,集群消费最多 4 个线程并行消费;扩容消费者超过 4 台,多出实例闲置。
消息重试机制
消费失败(抛出异常,不 return CONSUME_SUCCESS):
消息不会立刻丢弃;Broker 将消息投递到 重试队列 % RETRY% ConsumerGroup 每条消息带有重试次数,次数递增,投递间隔逐步拉长;
超过最大重试次数 → 转入死信队列 DLQ。
死信队列:人工排查处理,默认不自动重试。
NameServer 作用?和 Nacos/ZK 区别
NameServer 核心能力:
- Broker 启动主动注册 Topic 路由信息
- Producer/Consumer 从 NameServer 拉取 Topic 路由(Broker 地址列表)
- 定时检测 Broker 存活状态
特点:无状态、不持久化路由、不主动推送、客户端定时轮询拉取路由
对比 ZK:
ZK 有强一致性,复杂选举;NameServer 简化设计,无分布式一致性协议,CAP 选择 AP,集群节点互相不通信,水平扩展简单。
重要:NameServer不存储消息,只存路由元数据。
Broker Master/Slave 区别
Master:支持读写,接收生产者消息
Slave:不接收生产者写入;负责数据备份;Master 宕机后可升级提供读请求(不能写)
同步方式两种:
同步刷盘 + 同步复制(SYNC_MASTER):消息落盘且同步到 Slave,返回成功,可靠性最高
异步刷盘 + 异步复制(ASYNC_MASTER):性能高,存在丢失风险
消息基础类型
普通消息:默认,并发消费,性能最高
顺序消息:
全局顺序:整个 Topic 严格有序(生产极少用,吞吐量极低)
分区顺序(队列顺序):同一队列内有序;业务主流使用
延迟消息:消息投递后不立即消费,等待指定时间;只支持预设级别,不支持任意时间
事务消息:实现分布式事务(半消息、二阶段确认)
批量消息:多条消息打包发送,提升吞吐量
消费模式
- 集群消费(默认):同一个 ConsumerGroup 内,一条消息只会被一个消费者消费
- 广播消费:同一个 ConsumerGroup 所有实例都收到这条消息
Tag & Key
- Tag:消息标签,消息过滤(简单过滤)
- Key:业务唯一标识,方便运维根据 key 查询消息轨迹
分布式事务消息原理
事务消息三阶段(半消息机制)
第一阶段:发送半消息(Prepared 消息) Broker 接收半消息,对消费者不可见;回调生产者执行本地事务。
执行本地事务 生产者执行业务本地事务,返回三种状态:
Commit:Broker 将半消息转为有效消息,消费者可见
Rollback:Broker 删除半消息,消息丢弃
Unknown:未知状态,等待 Broker 回查
- 事务回查(Check) 长时间未收到 Commit/Rollback,Broker 主动发起回查;生产者查询本地事务最终状态,再次提交指令。
消息队列怎么避免重复消息
消息队列一般都有重试机制,或者生产者多次发送。
所以需要保证消费者接口幂等性保障:
比如可以使用redis的setnx。
接口幂等性
一锁、二判断、三更新
加分布式锁,判断状态机,更新至持久化的数据源中。
rocketmq如何保证消息的可靠性传输
需要生产者、broker、消费者共同配置。
生产者使用同步或异步发送,重试。
broker使用多主多从架构,nameServer使用集群,broker采用同步复制及同步刷盘。硬盘使用RAID10。
消费者消费成功之后返回正确的ACK,异常返回未消费ack。
消息发生大量堆积应该怎么处理
消费者可以消费,非bug
1.增加消费者数量
2.如果是消费者代码慢,优化代码。
3.降低生产者生产的速度。
4.清理过期数据,比如有一些消息已经处理过了或者过期了。
5.调整消费者的配置参数,提高消费者的效率。
6.增加topic队列数。
消费者bug, 堆积消息
没有做好监控,做好监控,提前报警。
1.增加消费者,看是否能抗住压力。这种会影响正在运的用户。
2.改变offset,旧的队列直接从最新开始,从旧topic快速读取新的临时topic,新的topic增加多队列数,启动新的消费者消费新topic。
rocketmq如何保障顺序消息的顺序性
rocketmq保障消息队列的顺序。
慎用顺序消息。
生产者
4.0 生产者使用MessageQueueSelector,同一业务ID(比如订单ID)的发送到同一消息队列中。
4.0感觉也不太对。
5.0的文档:
单一生产者 + 串行发送。
比如说有两个订单服务,服务A支付成功,服务B接收到取消订单,因为网络问题有可能取消订单在支付成功之前,这时没有退款,所以需要保证同一订单,要访问一个服务。或者直接支侍时使用支付中,支付中不能取消,不使用顺序消息。
消费者
消费者使用消费模式使用ConsumeMode.ORDERLY 顺序消费模式。
rocketmq的消费逻辑是,申请分布式锁,这时只有这个消费者内的线程可以消费,去申请broker的messageQueue锁,然后只有这个线程可以消费,然后消费者线程去申请broker的processQueue进行加锁,当新增加消费者时,不会把消息重复给新的消费者。
延时消息如何实现的
当消息发送到broker后,根据延时级别存储消息(不能消费),然后使用不同延时级别的定时器把消息改为可以消费。
延时定时器:它只支持:1s 5s 10s 30s 1m 2m 3m 4m 5m 6m 7m 8m 9m 10m 20m 30m 1h 2h这几个时长。
广播消息与集群消息怎么做到的
是消息者设置的消息模式。
当广播消息时,是每个消费者都有offset。
当集群消息时,一个消息队列一个offset。
rocketmq的消息是推还是拉?
Push是服务端主动推送消息给客户端,pull是客户端需要主动到服务端轮询获取数据。
推:及时性比较好。
拉:流处理框架场景下集成使用。
推模式:
推模式其底层的实现还是基于pull实现的,SDK内置了一个长轮询线程,先将消息异步拉取到SDK内置的缓存队列中,再分别提交到消费线程中,触发监听器执行本地消费逻辑。