rocketmq的面试题及原理

基础定位

Apache RocketMQ 分布式消息中间件,阿里开源,支持异步通信、削峰填谷、解耦、最终一致性事务消息、顺序消息、延迟消息、死信队列
三大应用场景:流量削峰、系统解耦、异步通知、分布式事务。

核心角色

  1. Producer 生产者:发送消息
  2. Consumer 消费者:消费消息
  3. NameServer:注册中心(轻量,无状态)
  4. Broker 消息服务器:消息存储、转发,核心节点
    • Broker 分为 Master(读写)、Slave(只读、同步备份)

完整消息发送流程

  1. Producer 启动,定时从 NameServer 拉取 Topic 路由(Broker 列表 + 队列信息)
  2. 根据负载均衡策略选出队列(MessageQueue)
  3. 建立 TCP 长连接,发送消息到对应 Broker
  4. Broker 接收消息:
    • 写入 CommitLog(所有消息统一顺序文件)
    • 异步分发到 ConsumeQueue(消息索引文件)
    • 执行刷盘机制
  5. 返回发送结果给生产者

三大存储文件

CommitLog

  • 所有 Topic 消息全部顺序写入同一个文件
  • 优点:顺序写磁盘,性能极高;避免多个文件随机写
  • 文件固定大小 1G,滚动创建

ConsumeQueue(消费队列索引)

每个 Topic 下每个队列对应一个 ConsumeQueue;
不存储完整消息,只保存索引:CommitLog 偏移量、消息长度、tag hashcode。
消费者消费时:查询 ConsumeQueue → 根据偏移量去 CommitLog 读取真实消息。

IndexFile

消息索引文件,根据msgKey快速检索消息,运维排查使用。

刷盘机制

  1. 同步刷盘 SYNC_FLUSH 消息写入内存,立刻持久化磁盘,返回成功;宕机不丢消息,性能偏低。金融场景选用。
  2. 异步刷盘 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 严格有序(生产极少用,吞吐量极低)
分区顺序(队列顺序):同一队列内有序;业务主流使用
延迟消息:消息投递后不立即消费,等待指定时间;只支持预设级别,不支持任意时间
事务消息:实现分布式事务(半消息、二阶段确认)
批量消息:多条消息打包发送,提升吞吐量

消费模式

  1. 集群消费(默认):同一个 ConsumerGroup 内,一条消息只会被一个消费者消费
  2. 广播消费:同一个 ConsumerGroup 所有实例都收到这条消息

Tag & Key

  • Tag:消息标签,消息过滤(简单过滤)
  • Key:业务唯一标识,方便运维根据 key 查询消息轨迹

分布式事务消息原理

事务消息三阶段(半消息机制)

  1. 第一阶段:发送半消息(Prepared 消息) Broker 接收半消息,对消费者不可见;回调生产者执行本地事务。

  2. 执行本地事务 生产者执行业务本地事务,返回三种状态:

  • Commit:Broker 将半消息转为有效消息,消费者可见

  • Rollback:Broker 删除半消息,消息丢弃

  • Unknown:未知状态,等待 Broker 回查

  1. 事务回查(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内置的缓存队列中,再分别提交到消费线程中,触发监听器执行本地消费逻辑。


rocketmq的面试题及原理
http://hanqichuan.com/2023/11/21/架构与分布式/mq/rocketmq/rocketmq的面试题及原理/
作者
韩启川
发布于
2023年11月21日
许可协议