RocketMQ 面试题
RocketMQ 面试题
导语:RocketMQ 的面试重点在存储设计(CommitLog / ConsumeQueue / IndexFile)、事务消息、延迟消息、顺序消费,以及主从切换与 Dledger。存储设计是它相对 Kafka 最独特的部分,几乎必问。共 12 题。
一、架构与存储
1. RocketMQ 的整体架构由哪些组件组成?
答: 四大组件:
| 组件 | 职责 |
|---|---|
| NameServer | 无状态的路由注册中心。Broker 定时(默认 30s)上报路由信息,Producer/Consumer 从这里拉取 Topic 的路由;多个 NameServer 节点彼此独立、互不通信 |
| Broker | 消息存储与转发的核心。分 Master / Slave,负责接收生产者消息、持久化、投递给消费者 |
| Producer | 生产者。先从 NameServer 获取 Topic 的路由,再按负载均衡策略选择队列发送 |
| Consumer | 消费者。同样先获取路由,再按分配策略消费;支持 PushConsumer(底层长轮询)与 PullConsumer |
核心概念:
| 概念 | 说明 |
|---|---|
| Topic | 消息的逻辑分类 |
| MessageQueue | Topic 的分区,是并行与扩容的基本单位(类比 Kafka 的 Partition);一个 Topic 的多个队列分布在多个 Broker 上 |
| ConsumerGroup | 消费组,组内竞争消费、组间广播 |
| Offset | 消费位点 |
| Tag | 消息的二级分类,用于 Broker 端高效过滤 |
| Key | 业务键(如 orderId),用于查询消息(Kafka 原生不具备的能力) |
消息流转:Producer 向 NameServer 拉取路由 → 按策略选择某个 MessageQueue → 发送到对应 Broker(写入 CommitLog)→ Consumer 拉取路由 → 从 Broker 拉取消息 → 处理并上报消费位点。
2. RocketMQ 的存储设计是怎样的?(CommitLog / ConsumeQueue / IndexFile)
答: 这是 RocketMQ 最具特色的设计,也是与 Kafka 的最大架构差异。
三个核心文件:
| 文件 | 作用 | 结构 |
|---|---|---|
| CommitLog | 消息真正存储的地方——所有 Topic 的所有消息都顺序追加写入同一个(组)CommitLog | 单文件默认 1GB,写满后新建下一个;采用 mmap 内存映射 |
| ConsumeQueue | 消费队列,是 CommitLog 的"轻量级索引",按 Topic + QueueId 组织 | 每个条目固定 20 字节:8 字节 CommitLog 物理偏移 + 4 字节消息长度 + 8 字节 Tag 的 hashCode |
| IndexFile | 按 Key / 时间查询消息的哈希索引文件 | IndexHeader(40B) + Slot Table(哈希槽) + Index Linked List(每条 20B) |
消息写入与读取的完整链路:
写入:Producer → Broker → 顺序追加到 CommitLog(mmap)
↘ 同时异步构建 ConsumeQueue(记录"这条消息在 CommitLog 的位置")
↘ 同时异步构建 IndexFile(记录 key → 物理位置)
读取:Consumer → 先读 ConsumeQueue(顺序读,极快)→ 拿到 CommitLog 物理偏移
→ 再按偏移去 CommitLog 读取真正的消息内容为什么这样设计(面试重点):
| 设计点 | 收益 |
|---|---|
| CommitLog 只有一份、全局顺序写 | 写入性能极稳定:无论 Topic/队列有多少,写入都是纯顺序追加,不存在随机写 |
| ConsumeQueue 顺序读 | 消费时按队列顺序读小文件(20 字节/条),读取也高效;且 ConsumeQueue 体积极小(是消息体的 1/100 量级),可大量缓存在 Page Cache |
| 读写分离 | 写入只认 CommitLog,读取先查 ConsumeQueue,两者解耦、互不干扰 |
| IndexFile | 支持按业务 key 查询消息(Kafka 原生不支持) |
与 Kafka 的关键对比(高频追问):
| 维度 | RocketMQ | Kafka |
|---|---|---|
| 存储组织 | 一个全局 CommitLog + 每个队列的 ConsumeQueue | 每个 Partition 一个独立的日志目录 |
| 分区数很多时 | 写入性能稳定(始终顺序追加一个文件) | 分区数过多时写放大与随机写加剧(每个分区独立文件、独立刷盘) |
| 按 key 查询 | ✅ 原生支持(IndexFile) | ❌ 不支持(需自己建外部索引) |
| 消息队列数量上限 | 单机可支撑数万队列 | 分区数过多会导致性能下降 |
一句话总结:RocketMQ 用"一份 CommitLog 顺序写 + 多份 ConsumeQueue 逻辑索引"实现了"高吞吐写入"与"多队列并发消费"的兼得;Kafka 用"每分区一个文件"换来了更简单的模型与更强的流处理生态。
二、可靠投递
3. RocketMQ 如何保证消息不丢失?
答: 与所有 MQ 一样分三层,但 RocketMQ 有两个独有的概念需要说清:
① 消息生产端
- 使用同步发送
producer.send(msg),并配置retryTimesWhenSendFailed(默认 2 次,即最多发送 3 次); - ⚠️ 必须检查返回的
SendStatus(这是最常见的生产 bug):
| SendStatus | 含义 | 是否丢消息 |
|---|---|---|
SEND_OK | 发送成功 | 不一定!还需看刷盘与复制配置(见下) |
FLUSH_DISK_TIMEOUT | 同步刷盘超时 | 可能丢,应重试/告警 |
FLUSH_SLAVE_TIMEOUT | 同步复制到 Slave 超时 | 可能丢,应重试/告警 |
SLAVE_NOT_AVAILABLE | 无可用 Slave | 可能丢 |
很多代码只做
try-catch而不检查SendStatus,导致"刷盘超时"等半失败被当作成功——这是实际项目里最常见的丢消息原因。
② Broker 端:刷盘与复制两个维度
| 维度 | 配置 | 含义 | 可靠性 |
|---|---|---|---|
| 刷盘 | flushDiskType=SYNC_FLUSH | 消息写入磁盘后才返回成功 | 最高(牺牲性能) |
flushDiskType=ASYNC_FLUSH(默认) | 先写 Page Cache 即返回,由后台异步刷盘 | 有窗口(Broker 宕机/断电会丢) | |
| 复制 | brokerRole=SYNC_MASTER | 主从都写成功才返回 | 最高(RT 升高;Slave 挂会影响写) |
brokerRole=ASYNC_MASTER | 主写成功即返回,Slave 异步同步 | 有窗口(Master 宕机会丢未同步数据) |
金融级可靠的标准组合:SYNC_FLUSH + SYNC_MASTER(即"同步刷盘 + 同步双写")。
③ 消费端
- 业务处理成功返回
CONSUME_SUCCESS,失败返回RECONSUME_LATER触发重试; - 注意:顺序消费(Orderly)失败是"原地重试",不会进入重试队列(见第 7 题)。
4. RocketMQ 的顺序消息有哪几种?
答: 两种,命名很直观:
| 类型 | 语义 | 实现 |
|---|---|---|
| 普通顺序消息(Normal Ordered) | 同一 MessageQueue 内有序,不同队列之间无序 | 生产者用 MessageQueueSelector 按业务 key(如 orderId)选择队列,保证同一 key 进同一队列 |
| 严格顺序消息(Strictly Ordered) | 全局严格有序 | Topic 只设置一个队列(或所有消息进同一队列);可用性与并发极差 |
消费端的配合:
- 必须使用
MessageListenerOrderly(顺序消费监听器)——Broker 会对该 MessageQueue 加分布式锁,保证同一时刻只有一个消费者在消费这个队列(消费完一批才释放); - 若用
MessageListenerConcurrently(并发消费),即使生产端保证了队列内有序,消费端并发处理也会打乱顺序。
生产建议:用"普通顺序消息 + key 选队列"即可满足绝大多数业务。严格顺序消息的可用性代价太大,几乎不用。
⚠️ 顺序消费的坑:失败后原地重试(默认间隔 1s),不进入重试队列——一条"毒消息"会阻塞整个队列。所以顺序消费必须做好"重试次数上限 + 落库/告警 + 人工跳过"的兜底。
5. RocketMQ 事务消息的原理?
答: 基于 两阶段提交 + 事务回查,目标是解决"本地事务执行"与"消息发送"的原子性问题。
完整流程:
① Producer 发送【半消息 Half Message】→ Broker
(半消息被存到内部 Topic RMQ_SYS_TRANS_HALF_TOPIC,【对消费者完全不可见】)
② Broker 确认半消息写入成功,返回 Producer
③ Producer 执行【本地事务】(如:在本地库插入订单)
④ Producer 根据本地事务结果,向 Broker 发送:
- CommitMessage → 事务成功
- RollbackMessage → 事务失败
- UNKNOWN → 状态未知(Broker 会发起回查)
⑤ 若 Broker 【长时间未收到④的确认】(Producer 宕机、网络异常等),
会定时【回查】Producer 的本地事务状态(默认最多回查 15 次,间隔按等级递增)
→ Producer 需要实现 checkLocalTransaction 返回真实状态
⑥ Broker 根据最终结果决定【Commit(消息对消费者可见)】或【Rollback(丢弃/标记删除)】三个关键点:
- 半消息阶段对消费者不可见——这是"消息与本地事务一致"的关键:只有本地事务成功、Broker 收到 Commit,消息才会被真正投递;
- 回查机制保证了"最终一致"——即使 Producer 在第 ④ 步宕机,Broker 也会反复询问,直到拿到明确状态;
- 回查接口必须幂等且能查到真实事务状态——通常做法是"把本地事务的结果写进业务表,回查时查这张表"。
面试追问:"为什么不能用普通的"先发消息再执行本地事务"?"——因为如果本地事务失败,消息已经发出去了,下游会处理一条"根本不存在"的业务数据。半消息的本质是"延迟消息的可见性",让消息的可见性由本地事务的结果决定。
三、延迟、重试与死信
6. RocketMQ 的延迟消息是如何实现的?
答: 4.x 与 5.0 的实现完全不同,这是版本差异的考点:
4.x:18 个固定延迟级别(不支持任意时间)
- 默认级别:
1s 5s 10s 30s 1m 2m 3m 4m 5m 6m 7m 8m 9m 10m 20m 30m 1h 2h(delayTimeLevel从 1 到 18); - 实现原理(面试重点):
- 生产者发送带
delayTimeLevel的消息时,Broker 不会直接投递到目标 Topic; - Broker 会把消息原本的 Topic 与 QueueId 备份到消息属性里,然后把消息的 Topic 改写为内部 Topic
SCHEDULE_TOPIC_XXXX,并按delayTimeLevel投递到对应的一个延迟队列(18 个级别 → 18 个队列); - 有一个后台定时任务
ScheduledMessageService按级别轮询这 18 个队列,每个级别维护一个"下一次到期时间"(用一个时间轮式的逻辑); - 到时间后,把消息从属性中恢复原来的 Topic 与 QueueId,重新写入 CommitLog(即真正投递);
- 生产者发送带
- 注意:这不是精确延迟——它是"级别 + 定时扫描"的近似实现,实际延迟会比配置值略长。
5.0:支持任意精度的定时消息
- 生产者可直接指定投递时间戳(秒级精度),不再受 18 个级别限制;
- 底层基于时间轮 + 定时文件(TimerWheel / TimerLog)实现,延迟精度到秒;
- 对长延迟消息做了专门的存储与索引设计(避免内存膨胀)。
与其他 MQ 的对比:
| MQ | 延迟消息 |
|---|---|
| RocketMQ | 原生支持(4.x 固定 18 级;5.0 任意秒级) |
| RabbitMQ | 需 TTL+DLX(有队头阻塞坑)或安装延迟插件 |
| Kafka | 原生不支持,需自研(如时间轮 + 内部 Topic) |
一句话:"延迟消息"是 RocketMQ 相对 Kafka/RabbitMQ 的一个明确优势功能,也是很多业务选择它的直接原因(如"订单 30 分钟未支付自动取消")。
7. RocketMQ 的消息重试与死信机制是怎样的?
答: 生产端与消费端的重试是两套完全不同的机制:
① 生产者重试(同步发送)
retryTimesWhenSendFailed(默认 2,即最多发 3 次);异步发送用retryTimesWhenSendAsyncFailed;- 重试时会尽量换一个 Broker/队列发送(
sendLatencyFaultEnable开启后会做故障规避,避免打到已知故障节点)。
② 消费者重试——两种消费模式的行为完全不同(高频考点):
| 消费模式 | 失败后的行为 |
|---|---|
并发消费(MessageListenerConcurrently) | 返回 RECONSUME_LATER → 消息进入内部 Topic %RETRY%<consumerGroup> → 按延迟级别递增重试,默认共 16 次:10s、30s、1m、2m、3m、4m、5m、6m、7m、8m、9m、10m、20m、30m、1h、2h |
顺序消费(MessageListenerOrderly) | 原地重试(默认间隔 1s,可用 suspendCurrentQueueTimeMillis 调整),不进入重试队列 → ⚠️ 毒消息会阻塞整个队列 |
③ 死信(DLQ)
- 并发消费重试达到最大次数(默认 16 次)仍失败 → 消息进入死信 Topic
%DLQ%<consumerGroup>; - 进入死信后不再自动重投,需要人工处理或写工具重新投递;
%RETRY%与%DLQ%都是系统 Topic(不会被自动删除),且每个消费组各有一份。
实践要点:
- 消费端必须监控重试 Topic 的堆积与死信数量,否则消息会"静默"沉到 DLQ 里;
- 顺序消费场景要额外实现"重试上限 + 跳过机制",否则一条脏数据会永久阻塞队列。
四、查询、过滤与高可用
8. RocketMQ 如何实现按 msgId / key 查询消息?
答: 两种查询走完全不同的路径:
① 按 MsgId 查询——不需要索引
MsgId由 Broker 生成,其结构里已经包含了"存储位置信息":Broker 地址 + CommitLog 物理偏移;- 因此拿到 MsgId 后,客户端可以直接定位到具体 Broker 的 CommitLog 物理位置读取,无需任何索引,效率极高。
② 按 Key 查询——依赖 IndexFile 哈希索引
Key是业务自定义的(如orderId),Broker 在写消息时异步构建 IndexFile:
IndexFile 结构:
IndexHeader(40 字节:起止时间戳、slot 数量、索引条目数量等)
Slot Table(哈希槽,默认 500 万 slot,每 slot 4 字节 → 存"索引条目链表头"的位置)
Index Linked List(每个条目 20 字节:
4 字节 key 的 hashCode
8 字节 CommitLog 物理偏移
4 字节 与文件起始时间的时间差
4 字节 指向"同 slot 上一条索引"的偏移 → 构成链表)- 查询流程:对 key 求 hashCode → 取模定位 slot → 沿链表遍历比对 hashCode(并校验 key 是否相同)→ 命中后按 CommitLog 偏移读取消息;
- 两个限制:
- 查询必须指定时间范围(beginTime / endTime)——因为索引文件是按时间滚动创建的(默认每个文件 400MB、可存 2000 万条索引、保留 7 天);
- key 不要设太多(每个 key 都会产生索引条目),一般每个消息设 1 个业务 key 即可。
对比记忆:"MsgId 自带坐标,key 需要查索引"。另外这也解释了为什么 RocketMQ 能提供"消息轨迹"与"按业务 key 追踪消息"的能力,而 Kafka 原生做不到。
9. RocketMQ 的 Tag 与 SQL92 消息过滤有什么区别?
答:
| 维度 | Tag 过滤 | SQL92 属性过滤 |
|---|---|---|
| 过滤依据 | 消息的 Tag 字符串 | 消息的用户属性(UserProperty) |
| 语法 | TagA、TagA || TagB、*(全部) | SQL92 表达式,如 price > 100 AND region = '上海' |
| Broker 端实现 | 哈希比对:ConsumeQueue 的 20 字节条目里内嵌了 Tag 的 hashCode,Broker 只需比对哈希值即可过滤 | 需要解析并执行 SQL 表达式,还要读取消息属性 |
| 性能 | 极高(O(1) 比对,在索引层就能过滤) | 较低(CPU 开销大,需开启 enablePropertyFilter=true) |
| 表达能力 | 一个消息只能打一个 Tag,表达能力有限 | 多属性任意组合,表达能力最强 |
| 使用建议 | ✅ 优先使用 | ⚠️ 仅当 Tag 无法满足时使用 |
为什么 Tag 过滤能做到如此高效(加分点):
RocketMQ 在设计 ConsumeQueue 条目时,就预留了 8 字节专门存 Tag 的 hashCode。这样一来:
- 消费者订阅时,把订阅的多个 Tag 预先算成 hashCode 集合;
- Broker 在遍历 ConsumeQueue 的过程中就能直接比对 hashCode 完成过滤,完全不需要读取 CommitLog 中的真实消息内容;
- 这正是"Broker 端过滤 + 索引层过滤"的极致优化——比"拉到消费者再过滤"节省了大量网络与 CPU。
注意:哈希比对存在极小的碰撞概率(哈希相同但 Tag 不同),所以 RocketMQ 在客户端还会再做一次精确的 Tag 字符串比对作为兜底。
10. RocketMQ 的高可用是怎么实现的?
答: 分四个层面:
① NameServer 层
- 多节点部署,彼此无状态、互不通信,各自维护全量路由;
- 任一节点宕机,客户端换一个节点即可,天然高可用;
- 代价:路由信息变更最长有 30s 延迟(Broker 上报周期),可能出现短暂的"读到旧路由"。
② Broker 层:主从架构
- 支持多 Master 多 Slave(一主多从 / 多主多从);Slave 不提供写,但可提供读(分担读压力);
- 复制方式:
- 异步复制(
ASYNC_MASTER):性能高,但 Master 宕机可能丢少量未同步数据; - 同步双写(
SYNC_MASTER):Master + Slave 都写成功才返回,不丢数据,但 RT 升高且Slave 故障会影响写。
- 异步复制(
③ 关键短板:传统主从不支持自动切换
- Master 宕机后,Slave 不会自动成为 Master,需要人工介入或借助外部组件(如运维脚本、监控系统)触发切换;
- 这是 RocketMQ 相对 Kafka 的一个明确劣势(Kafka 由 Controller 自动选主);
- 切换期间该 Topic 队列不可写,直到切换完成。
④ Dledger 模式(4.5+):Raft 自动选主
- 引入 Dledger(基于 Raft 协议)后,一组 Broker(至少 3 个节点)可自动选举 Leader、自动故障切换,无需人工干预;
- 代价:与 Kafka 类似,只有 Leader 提供读写(Slave 不再分担读),吞吐会有一定下降;
- 适用于对可用性要求高、且不能接受人工切换的场景。
⑤ RocketMQ 5.0 的演进
- 引入 Proxy(无状态网关,支持 gRPC 与多语言) 与 Controller 组件,把"自动选主与容灾"从 Dledger 中抽象出来,向"存算分离、自动容灾"演进。
11. RocketMQ 的并发消费与顺序消费有什么区别?
答: 这是消费者端的核心选择,也是"顺序消息"的落地关键:
| 维度 | 并发消费(Concurrently) | 顺序消费(Orderly) |
|---|---|---|
| 监听器 | MessageListenerConcurrently | MessageListenerOrderly |
| 消费方式 | 同一队列的消息可被多线程并发消费(默认 20 个线程) | 同一队列串行消费 |
| 加锁机制 | 无(依赖消费位点管理) | Consumer 消费前需向 Broker 申请该 MessageQueue 的分布式锁(lockMQ),保证同一队列同一时刻只被一个消费者消费 |
| Rebalance 行为 | 直接重新分配队列 | 只分配加锁成功的队列,加锁失败的队列等待下次 |
| 失败处理 | 返回 RECONSUME_LATER → 进入 %RETRY% 按延迟级别重试(共 16 次) | 原地重试(默认 1s 间隔,可配置),不进入重试队列 |
| 风险 | 消费位点管理复杂(失败重试可能与后续消息乱序,需幂等) | ⚠️ 毒消息阻塞整个队列(必须自己实现重试上限 + 跳过) |
| 吞吐 | 高 | 低(串行 + 加锁开销) |
| 适用 | 绝大多数场景 | 订单状态流转、增量同步等必须保序的场景 |
实践建议:
- 默认用并发消费,只对确实需要保序的业务用顺序消费;
- 顺序消费一定要配
suspendCurrentQueueTimeMillis+ 重试次数上限 + 死信/人工兜底,否则一条坏数据会卡死队列;- 更好的方案是"业务键分片 + 并发消费"——把同一 key 的消息用
MessageQueueSelector送到同一队列,消费端对每个队列用单线程处理,其他队列仍然并发。
12. RocketMQ 5.0 有哪些重要变化?
答: 5.0 是一次较大的架构升级,四个方向:
| 方向 | 变化 |
|---|---|
| 存算分离 | 引入无状态 Proxy(支持 gRPC 协议与多语言 SDK),Broker 更专注于存储;客户端不再直连 Broker,而是连 Proxy,架构更清晰、扩缩容更灵活 |
| 任意精度定时消息 | 用时间轮 + 定时文件替代 4.x 的 18 个固定延迟级别,支持秒级任意延迟 |
| 自动容灾 | 引入 Controller 组件,把自动选主与容灾能力从 Dledger 中抽象出来,部署与运维更简单 |
| 多语言与生态 | 提供基于 gRPC 的轻量级 SDK(Java / Go / C++ / Rust / Python 等),社区生态与可观测性(消息轨迹、OpenTelemetry 集成)显著增强 |
升级注意事项:
- 5.0 的 Broker 兼容 4.x 客户端的 Remoting 协议,可以平滑升级;
- 但新特性(如任意精度定时消息、Proxy)需要 5.x 的新 SDK;
- 若使用 Proxy 模式,需评估多一跳网络带来的延迟与 Proxy 自身的高可用。
参考对比:RocketMQ 5.0 与 Kafka 的 KRaft 演进方向相似——都在"去中心化协调、简化部署、增强扩展性"上发力,只是 RocketMQ 更强调"计算与存储分离",Kafka 更强调"元数据自管理"。
