RocketMQ
开篇:RocketMQ 的设计目标——金融级可靠性
如果说 Kafka 是一辆追求极速的跑车,RabbitMQ 是一辆路线灵活的出租车,那 RocketMQ 就是一辆安全系数拉满的装甲运钞车。
RocketMQ 由阿里巴巴开源,后捐赠给 Apache 基金会。它诞生于阿里双 11 的实战中,设计目标就是金融级可靠性——事务消息、延迟消息、消息回溯等能力都是原生支持的。在电商、金融、支付等对消息可靠性要求极高的领域,RocketMQ 是首选。
这篇文章会从架构到原理,带你把 RocketMQ 的核心机制吃透。
一、核心架构:四大角色
RocketMQ 的架构围绕四个角色展开:
NameServer——路由注册中心
NameServer 是一个非常轻量的组件,功能类似于"简化版的 ZooKeeper"。每个 Broker 启动后会向所有 NameServer 注册自己的信息(IP、端口、持有的 Topic 等),并定时发送心跳。
Producer 和 Consumer 启动时随机连接一台 NameServer,获取 Broker 的路由信息(哪个 Topic 在哪个 Broker 上,有几个 Queue),然后直接与 Broker 通信。
NameServer 之间互不通信,各自独立维护信息。这意味着:
- 某个 NameServer 挂了不影响整体(其他 NameServer 还在)
- 但如果某个 Broker 只注册到了部分 NameServer,连接其他 NameServer 的客户端可能获取不到完整路由信息
Broker——消息存储与转发
Broker 是 RocketMQ 的核心组件,负责消息的接收、存储和投递。通常采用一主多从部署:
- Master:接收生产者消息,处理消费者请求
- Slave:从 Master 同步数据,当 Master 不可用时消费者可以从 Slave 读取
RocketMQ 是分布式的,可以部署在多台机器上。利用 Topic 对业务进行隔离,Topic 下面会有多个 Queue(类似 Kafka 的 Partition),Queue 可以分布在不同 Broker 上。
Producer——消息生产者
支持同步发送、异步发送和单向发送三种方式。单向发送只管发不管结果,性能最高但不可靠,一般不推荐。
Consumer——消息消费者
支持 Push 和 Pull 两种消费模式。Push 本质上也是基于 Pull 实现的(长轮询)。5.0 版本新增了 Pop 模式,后面会详细介绍。
二、消息存储:CommitLog + ConsumeQueue
RocketMQ 的存储设计是其高性能的核心之一。
CommitLog——所有消息的"总账本"
不管消息属于哪个 Topic、哪个 Queue,所有消息都写到同一个 CommitLog 文件中,顺序追加写入。这样做的好处是充分利用了磁盘顺序写的性能优势。
CommitLog(单个文件默认 1GB)
| msg-A (topic=order) | msg-B (topic=log) | msg-C (topic=order) | ...ConsumeQueue——消费索引
CommitLog 里的消息是混在一起的,消费者不可能挨个遍历。所以 RocketMQ 为每个 Topic 的每个 Queue 维护了一个 ConsumeQueue——它是 CommitLog 的索引,记录了每条消息在 CommitLog 中的物理偏移量、消息大小和 Tag 的 hashcode。
ConsumeQueue (topic=order, queue=0)
| offset=0, commitlog_pos=100, size=256, tag_hash=xxx |
| offset=1, commitlog_pos=800, size=128, tag_hash=yyy |消费者消费时:先从 ConsumeQueue 找到消息在 CommitLog 中的位置,再去 CommitLog 读取完整消息。
为什么这样设计?
- 写入极快——所有消息顺序写同一个文件,磁盘顺序写接近内存速度
- 读取高效——ConsumeQueue 文件很小(每条索引只有 20 字节),可以常驻内存
- 天然支持消息回溯——CommitLog 按时间线存储,可以按时间点回溯消费
三、事务消息:分布式事务的利器
事务消息是 RocketMQ 最有特色的功能之一,也是它和 Kafka 最大的差异点。
解决什么问题?
想象一个场景:用户下单后需要扣减库存。订单服务先创建订单(本地数据库操作),然后发消息通知库存服务扣库存。
问题来了:如果订单创建成功了,但消息发送失败了怎么办?或者消息发出去了,但订单因为事务回滚没有创建成功呢?
RocketMQ 的事务消息保证的是:本地事务和消息发送要么都成功,要么都失败。
注意和 Kafka 的区别:Kafka 的事务消息保证的是一组消息的原子性(都发成功或都失败),不涉及本地事务。
实现原理
整个过程分为四步:
第一步:发送半消息
Producer 先向 Broker 发送一条"半消息"(Half Message)。半消息会被存储,但消费者看不到它——Broker 会把它暂存在一个特殊的内部 Topic 中。
第二步:执行本地事务
半消息发送成功后,Producer 开始执行本地业务逻辑(比如创建订单、写数据库)。
第三步:提交或回滚
- 本地事务成功:通知 Broker COMMIT,半消息变成正式消息,消费者可以消费了
- 本地事务失败:通知 Broker ROLLBACK,半消息被删除
第四步:事务回查
如果 Broker 迟迟收不到 COMMIT 或 ROLLBACK(比如 Producer 崩溃了),Broker 会主动回查 Producer 的本地事务状态。
这就是为什么通常建议在执行本地事务时,同时往一张事务状态表中插入一条记录。回查的时候直接查这张表就行了,不需要再去执行复杂的业务查询。
代码示例
// 创建事务消息生产者
TransactionMQProducer producer =
new TransactionMQProducer("tx-group");
producer.setNamesrvAddr("localhost:9876");
// 设置事务监听器
producer.setTransactionListener(new TransactionListener() {
@Override
public LocalTransactionState executeLocalTransaction(
Message msg, Object arg) {
try {
// 执行本地事务(创建订单、写事务表等)
orderService.createOrder(msg);
return LocalTransactionState.COMMIT_MESSAGE;
} catch (Exception e) {
return LocalTransactionState.ROLLBACK_MESSAGE;
}
}
@Override
public LocalTransactionState checkLocalTransaction(
MessageExt msg) {
// 回查:查询事务表,判断本地事务是否成功
boolean exists = txLogService.exists(msg.getTransactionId());
return exists ?
LocalTransactionState.COMMIT_MESSAGE :
LocalTransactionState.ROLLBACK_MESSAGE;
}
});
producer.start();
// 发送事务消息
Message msg = new Message("order-topic", "创建订单".getBytes());
producer.sendMessageInTransaction(msg, null);为什么不是先执行本地事务再发消息?
因为如果先执行本地事务(订单创建成功),再发消息,发消息这步失败了怎么办?本地事务已经提交了,但消息没发出去,数据就不一致了。
事务消息的方案就是为了解决这个问题:通过半消息 + 回查机制,确保本地事务和消息发送要么都成功,要么都失败。
四、延迟消息与时间轮
4.1 延迟消息
RocketMQ 原生支持延迟消息。消息写入 Broker 后不会立即被消费,需要等待指定时间后才可被消费。
5.0 之前支持的延迟级别是固定的:
1s 5s 10s 30s 1m 2m 3m 4m 5m 6m 7m 8m 9m 10m 20m 30m 1h 2h共 18 个级别,通过 setDelayTimeLevel() 设置:
Message msg = new Message("topic", "Hello".getBytes());
msg.setDelayTimeLevel(3); // 级别3 = 延迟10秒
producer.send(msg);4.2 时间轮实现(5.0+)
5.0 之前用 Timer 定时器扫描延迟消息,大量任务时性能会下降。RocketMQ 5.0 引入了基于时间轮的定时消息,带来了两个关键改进:
- 支持任意延迟时间——不再限于固定的 18 个级别,可以指定秒级、毫秒级精度
- O(1) 的调度效率——时间轮能在常数时间内找到下一个要执行的任务
时间轮的工作原理:
时间轮(假设 60 个槽位,每个槽位 1 秒)
┌──┬──┬──┬──┬──┬──┬──┬──┬──┐
│s0│s1│s2│s3│ │ │ │s58│s59│
└──┴──┴──┴──┴──┴──┴──┴──┴──┘
↑ 当前指针
消息延迟 5 秒 → 放入 s5 槽位
指针走到 s5 → 投递该槽位中的所有消息五、顺序消息与消息重试
5.1 顺序消息
和 Kafka 一样,RocketMQ 只保证同一个 Queue 内的消息有序。要实现业务层面的顺序消费,需要两步:
发送端:把相关消息发到同一个 Queue
SendResult result = producer.send(msg,
new MessageQueueSelector() {
@Override
public MessageQueue select(List<MessageQueue> mqs,
Message msg, Object arg) {
Integer orderId = (Integer) arg;
int index = orderId % mqs.size();
return mqs.get(index); // 相同订单号路由到同一队列
}
}, orderId);注意:顺序消息必须用同步发送,异步发送无法保证顺序。
消费端:使用有序消费模式
consumer.registerMessageListener(
new MessageListenerOrderly() {
@Override
public ConsumeOrderlyStatus consumeMessage(
List<MessageExt> msgs,
ConsumeOrderlyContext context) {
// 按顺序处理消息
return ConsumeOrderlyStatus.SUCCESS;
}
});顺序消费的代价
RocketMQ 为了保证顺序消费,背后做了三重加锁:
- Broker 端分布式锁——锁住 MessageQueue,确保同一时间只有一个消费者消费
- 本地 MessageQueue 锁——确保同一时间只有一个线程处理该队列
- ProcessQueue 锁——防止重平衡时消息被重复消费
三重锁意味着:
- 吞吐量大幅下降(串行处理)
- 如果前面的消息阻塞,后面的消息全部卡住
- 出错时排查困难
所以,顺序消息要慎用。实际上很多场景可以让消费者自己做排序,而不是依赖 MQ 的顺序消费。
5.2 消息重试
消费者消费失败时,RocketMQ 会自动重试。在集群模式下:
- 消费失败的消息会进入 %RETRY%+ConsumerGroup 的重试队列
- 默认最多重试 16 次,重试间隔逐渐增大(10s、30s、1min、2min...2h)
- 超过最大重试次数后,消息进入死信队列(%DLQ%+ConsumerGroup)
死信队列中的消息不会再被自动投递,需要人工干预处理。
六、RocketMQ vs Kafka vs RabbitMQ
这是面试中最常被问到的对比题。我们从多个维度来看:
| 维度 | RocketMQ | Kafka | RabbitMQ |
|---|---|---|---|
| 吞吐量 | 十万级/秒 | 百万级/秒 | 万级/秒 |
| 延迟 | 毫秒级 | 毫秒级 | 微秒级 |
| 事务消息 | 本地事务+消息原子性 | 仅消息原子性 | 基本事务 |
| 延迟队列 | 原生支持(18级/时间轮) | 不原生支持 | 死信/插件 |
| 死信队列 | 原生支持 | 不支持 | 原生支持 |
| 消息重试 | 自动重试 | 需手动实现 | 手动配置 |
| 消息回溯 | 支持按时间回溯 | 支持按 Offset 回溯 | 不支持 |
| 消费模式 | Push/Pull/Pop | Pull | Push 为主 |
| 协议 | 自定义协议 | 自定义协议 | AMQP |
| 开发语言 | Java | Scala/Java | Erlang |
| 最佳场景 | 电商、金融、分布式事务 | 日志、大数据、流处理 | 中小规模、复杂路由 |
| 社区 | 国内活跃 | 全球活跃 | 全球活跃 |
简单总结:
- 追求极致吞吐量——选 Kafka
- 需要事务消息、延迟消息等丰富功能——选 RocketMQ
- 系统规模不大但路由需求复杂——选 RabbitMQ
七、常见面试题精选
1. RocketMQ 如何保证消息不丢失?
需要生产者、Broker、消费者三方配合:
生产者端——使用同步发送或带回调的异步发送,确保发送结果可感知。
// 同步发送
SendResult result = producer.send(msg);
// 异步发送
producer.sendAsync(msg, new SendCallback() {
@Override
public void onSuccess(SendResult sendResult) { }
@Override
public void onException(OnExceptionContext context) {
// 重试逻辑
}
});避免使用单向发送(sendOneway),它不等待任何确认。
Broker 端——两个关键配置:
# 同步刷盘(默认异步,异步有丢消息风险)
flushDiskType=SYNC_FLUSH
# 同步复制(默认异步,异步在 Master 挂了时可能丢数据)
brokerRole=SYNC_MASTER同步刷盘 + 同步复制最安全,但吞吐量最低。根据业务需求在可靠性和性能之间做权衡。
消费者端——确保业务处理成功后才返回 CONSUME_SUCCESS:
consumer.registerMessageListener(
new MessageListenerConcurrently() {
@Override
public ConsumeConcurrentlyStatus consumeMessage(
List<MessageExt> msgs,
ConsumeConcurrentlyContext context) {
try {
// 业务处理
processMessages(msgs);
return ConsumeConcurrentlyStatus.CONSUME_SUCCESS;
} catch (Exception e) {
// 返回 RECONSUME_LATER,消息会重试
return ConsumeConcurrentlyStatus.RECONSUME_LATER;
}
}
});2. 消息丢了,可能是什么原因?
| 环节 | 常见原因 |
|---|---|
| 生产者 | 使用了单向发送(sendOneway) |
| 生产者 | 同步/异步发送但未处理失败(吞掉了异常) |
| 生产者 | 发送过程中 Producer 崩溃 |
| Broker | 异步刷盘时宕机(flushDiskType=ASYNC_FLUSH) |
| Broker | 异步复制时 Master 挂了(brokerRole=ASYNC_MASTER) |
| Broker | 消息超过保留时间未消费(默认 72 小时) |
| 消费者 | 消费失败但错误地返回了 CONSUME_SUCCESS |
3. 消息重复消费了怎么办?
常见导致重复消费的原因:
- 消费者返回
RECONSUME_LATER,消息自动重投 - 消费超时,Broker 判定失联,消息重投
- 生产者因网络异常重复发送了同一条消息
- 重平衡导致 Offset 未及时提交
解决方案:消费端做好幂等。在消息中携带唯一业务标识,消费者端用"一锁、二判、三更新"的方式防止重复处理。
4. 消息堆积了怎么解决?
消息堆积是正常的——MQ 的削峰填谷功能本身就意味着消息可能暂时堆积。处理步骤:
- 定位问题:哪个 Topic 堆积了?当前堆积量是增加还是在减少?
- 评估影响:延迟是否在业务可接受范围内?
- 临时方案:扩容消费者数量(注意不要把下游打挂)
- 长期方案:
- 批量消费 + 线程池并发处理
- 优化慢 SQL,减少单条消息处理时间
- "先收单再处理"模式——消费者先把消息存到数据库,快速 ACK,然后由定时任务慢慢处理
5. RocketMQ 的事务消息和 Kafka 的有什么区别?
这是个高频对比题:
| 维度 | RocketMQ 事务消息 | Kafka 事务消息 |
|---|---|---|
| 保证对象 | 本地事务 + 消息发送的原子性 | 一组消息的原子性 |
| 核心机制 | 半消息 + 回查 | 事务日志 + 两阶段提交 |
| 典型用途 | 分布式事务(如订单+库存) | Exactly-Once 语义 |
| 复杂度 | 需实现回查逻辑 | 需配置事务 ID |
一句话:RocketMQ 管的是"本地事务和消息要一起成或一起败",Kafka 管的是"这批消息要一起成或一起败"。
6. Push 和 Pull 模式怎么选?Pop 模式是什么?
Push 模式:看似 Broker 主动推消息,实际底层是长轮询(Pull 的封装)。使用简单,适合大多数场景。
Pull 模式:消费者主动拉取,灵活度高,适合需要精确控制消费节奏的场景。
Pop 模式(5.0 新增):解决 Push 模式的两个痛点:
- Push 模式下消费者需要做负载均衡,且 Queue 数量限制了消费者数量
- 某个消费者 hang 住会导致分配给它的 Queue 消息堆积
Pop 模式中,消费者不再绑定特定 Queue,直接通过 Pop 接口获取消息。负载均衡和 Offset 管理都由 Broker 端完成。某个消费者 hang 住了,其他消费者照样能消费那些消息。
八、集群部署方式
RocketMQ 支持三种集群模式:
| 模式 | 特点 | 适用场景 |
|---|---|---|
| 单 Master | 最简单,Master 挂了集群不可用 | 开发测试 |
| 多 Master | 多个 Master,无 Slave。配合 RAID10 可靠性较高 | 对延迟敏感但可接受短暂不可用 |
| 多 Master 多 Slave | 最高可用性,Master 挂了 Slave 可继续提供读服务 | 生产环境推荐 |
多 Master 多 Slave 模式下还要区分同步复制和异步复制:
# 同步复制——Master 等 Slave 同步完才返回成功,最安全
brokerRole=SYNC_MASTER
# 异步复制——Master 写入即返回,性能更好但有丢数据风险
brokerRole=ASYNC_MASTER九、消息分发模式
RocketMQ 支持两种消费模式:
集群消费(默认)——同一个 Consumer Group 内,一条消息只会被一个消费者处理。适合大部分业务场景。
广播消费——同一个 Consumer Group 内,每个消费者都会收到全量消息。适合配置刷新、缓存同步等场景。
// 集群消费(默认)
props.put(PropertyKeyConst.MessageModel,
PropertyValueConst.CLUSTERING);
// 广播消费
props.put(PropertyKeyConst.MessageModel,
PropertyValueConst.BROADCASTING);注意:广播模式下消费失败不会重投,客户端重启会从最新消息开始消费(停止期间的消息会被跳过)。
十、重平衡:和 Kafka 有什么不同?
RocketMQ 在集群消费模式下也有重平衡机制,但和 Kafka 有明显区别:
定时触发:RocketMQ 每 20 秒做一次重平衡检查(Kafka 依赖心跳实时触发)。如果消费者宕机,最多 20 秒后才会重新分配。
STW 影响更小:RocketMQ 的消费者通过异步拉取消息到本地队列,即使重平衡期间,本地队列中已拉取的消息仍在正常处理。只要重平衡时间够短(通常毫秒级),消费者几乎感受不到中断。
局部调整:RocketMQ 默认只调整受影响的消费者,其他消费者不受影响(类似 Kafka 的渐进式重平衡,但 RocketMQ 天生就是这样的)。
小结
RocketMQ 的核心竞争力可以用四个词概括:事务消息、延迟消息、消息回溯、金融级可靠。
它的架构设计围绕"可靠性优先"展开:CommitLog 顺序写保证写入性能,ConsumeQueue 索引保证消费效率,半消息 + 回查机制保证分布式事务的一致性,多级延迟 + 时间轮保证延迟消息的精度。
如果你的系统处于电商、金融、支付等对消息可靠性要求极高的领域,需要事务消息、延迟消息、消息重试等开箱即用的能力,RocketMQ 是当之无愧的首选。如果单纯追求吞吐量做日志和大数据管道,那 Kafka 更合适。
最后附一张关键配置速查表:
| 参数 | 推荐值 | 作用 |
|---|---|---|
flushDiskType | SYNC_FLUSH | 同步刷盘,防止宕机丢消息 |
brokerRole | SYNC_MASTER | 同步复制,Master 挂了数据不丢 |
sendMsgTimeout | 3000 | 发送超时时间(毫秒) |
retryTimesWhenSendFailed | 3 | 同步发送失败重试次数 |
messageDelayLevel | 1s 5s 10s... | 延迟级别(5.0 前) |
consumeMessageBatchMaxSize | 32 | 批量消费每批最大消息数 |
使用消息队列的三大核心价值同样适用于 RocketMQ 的选型判断:解耦(服务间通过消息通信,互不依赖)、异步(耗时操作放队列后台处理)、削峰填谷(高峰请求先进队列缓冲)。
但需要注意一点:用了 MQ 不等于自动实现了削峰。如果使用 Push 模式,消息还是会被快速推送给消费者,消费者可能被打崩。真正想做好削峰填谷,应该使用 Pull 模式让消费者按自己的能力拉取消息,让消息在 Broker 中堆积缓冲,而不是在客户端堆积。
另外,RocketMQ 5.0 中的 Pop 模式提供了一个更优的选择——消费者不绑定 Queue,负载均衡由 Broker 完成,既简化了客户端逻辑,又突破了 Queue 数量对消费者扩容的限制。如果是新项目,推荐直接使用 5.0 的 Pop 模式。
RocketMQ 的工作流程概览:NameServer 启动等待连接 -> Broker 启动,向 NameServer 注册并定时心跳 -> Producer 连接 NameServer 获取路由,向 Broker 发送消息 -> Broker 接收消息,根据刷盘和复制策略持久化 -> Consumer 连接 NameServer 获取路由,从 Broker 消费消息。整个过程中 NameServer 是无状态的路由表,Broker 是有状态的存储引擎,Producer 和 Consumer 都是 Broker 的客户端。
关于"是否用了 RocketMQ 就一定能实现削峰"这个问题值得额外注意。答案是不一定。Push 模式下消息会被快速推送给消费者,如果消费能力不足,消息会堆积在消费者端甚至导致消费者崩溃,然后消息重投带来更大的压力。所以想做好削峰填谷,务必使用 Pull/Pop 模式,由消费者控制拉取速度,让 Broker 队列成为天然的缓冲池。同时注意配合
consumeMessageBatchMaxSize和pullBatchSize等参数调优消费能力。
最后,关于顺序消息的实际使用:大多数场景其实不需要 MQ 级别的顺序保证。顺序消息的代价(三重加锁、吞吐量下降、阻塞传播)往往大于收益。更常见的做法是让消费者自己排序处理——比如在消费端按照消息中的时间戳或序号做排序,或者通过状态机模式确保业务逻辑的正确性。只有在状态流转间隔极短(毫秒级)、且无法在消费端排序的场景下,才需要使用 RocketMQ 的顺序消息功能。
以上就是 RocketMQ 的核心知识体系。掌握了架构、存储、事务消息、延迟消息和消费模式这五大模块,面对大部分面试和实战场景都能从容应对。