Skip to content

RocketMQ 深度解析 ​

#组件 · #RocketMQ · #消息队列 · #事务消息 · #阿里 · #分布式

阿里开源的万亿级消息引擎——高可靠存储、分布式事务、延迟消息与金融级特性。

一、架构 ​

1.1 四大角色 ​

mermaid
flowchart TB
    subgraph NS["NameServer Cluster<br/>无状态, 路由注册与发现"]
        NS1["NameServer 1"]
        NS2["NameServer 2"]
    end

    subgraph Brokers["Broker 集群"]
        subgraph B1["Broker 1"]
            BM1["Master (写)"]
            BS1["Slave (备, 读)"]
        end
        subgraph B2["Broker 2"]
            BM2["Master (写)"]
            BS2["Slave (备, 读)"]
        end
    end

    Brokers -->|"定时注册<br/>Topic → Broker 路由"| NS

    Producer["Producer<br/>生产者"] -->|"从 NameServer<br/>获取路由"| NS
    Producer -->|"发送消息"| Brokers

    Consumer["Consumer<br/>消费者"] -->|"从 NameServer<br/>获取路由"| NS
    Brokers -->|"拉取消息 (Pull)"| Consumer

与 Kafka 关键区别:Broker 主从靠手动/自动切换,没有 ZK/KRaft 选举 Leader。高可用靠 DLedger(Raft 变体,4.5+)或 Controller 模式(5.0+)。

角色职责
NameServer路由注册中心 (Topic → Broker 映射),无状态,互相不通信
Broker消息存储与投递,主从模式,定时向所有 NameServer 注册
Producer从 NameServer 拉路由 → 发送消息到 Broker
Consumer从 NameServer 拉路由 → 拉取消息(Pull 模型)

与 Kafka 架构关键不同:Broker 主从依赖手动/自动切换,没有 ZK/KRaft 选举 Leader。高可用靠 DLedger(Raft 变体,4.5+)或 Controller 模式(5.0+)。

1.2 Broker 集群模式 ​

模式特点
单 Master测试用,宕机不可用
多 Master无 Slave,宕机会丢少量消息
多 Master 多 Slave (异步)主挂备升,可能丢少量
多 Master 多 Slave (同步双写)主备同时写,不丢消息但性能下降
DLedger (Raft)自动选主,强一致,金融级推荐

二、消息存储(CommitLog + ConsumeQueue) ​

2.1 存储架构 ​

RocketMQ 文件布局 (单 Broker):
store/
├── commitlog/                    ← 所有 Topic 消息顺序追加写入
│   ├── 00000000000000000000      (1GB/文件)
│   ├── 00000000001073741824
│   └── ...
├── consumequeue/                 ← 按 Topic-Queue 的索引
│   └── TopicA/
│       ├── 0/00000000000000000000
│       └── 1/00000000000000000000
├── index/                        ← 按 Key/时间查询的哈希索引
│   └── ...
└── config/

2.2 写入与消费流程 ​

Producer 发送消息:
  消息 → CommitLog (顺序写, 单文件, 1GB轮转)
       ↓ (后台 ReputService 异步构建)
  ConsumeQueue (Topic-Queue 粒度索引)
       ↓ (消费者拉取)
  Consumer

ConsumeQueue 条目 (定长 20B):
┌────────────┬──────────┬──────────┐
│ CommitLog  │ 消息大小  │ Tag Hash │
│ Offset (8B)│  (4B)    │  (8B)    │
└────────────┴──────────┴──────────┘

优势:所有 Topic 共享一个 CommitLog,顺序写只需一个文件,写入极快。

对比 Kafka:Kafka 按 Partition 单独写文件,Partition 多时随机写增多;RocketMQ 所有消息合并顺序写,仅需少量随机读。

2.3 刷盘与 HA ​

参数说明
flushDiskType=ASYNC_FLUSH异步刷盘(默认,吞吐高)
flushDiskType=SYNC_FLUSH同步刷盘(每条消息等 fsync)
brokerRole=SYNC_MASTER同步双写 Slave(主等备确认)

可靠级别:异步刷盘+异步复制 < 同步复制 < 同步刷盘+同步复制(最安全但最低)


三、消息模型 ​

3.1 消费模式 ​

集群消费 (默认, CLUSTERING):
Topic: Order
  Queue 0 → Consumer A1  ─┐
  Queue 1 → Consumer A2   ├ Consumer Group A
  Queue 2 → Consumer A3  ─┘
  (每条消息只被组内一个消费者消费)

广播消费 (BROADCASTING):
Topic: Notify
  Queue 0 → Consumer B1, Consumer B2, Consumer B3
  Queue 1 → Consumer B1, Consumer B2, Consumer B3
  (每条消息被组内所有消费者都收到)

3.2 消费进度(Offset) ​

Consumer Offset 存储在 Broker (非 ZK):
  集群模式: 存 Broker 端
  广播模式: 存 Consumer 本地

消费进度: minOffset ←── 已消费 ──→ consumerOffset ←── 未消费 ──→ maxOffset

重新消费: 重置 consumerOffset 到某个时间点

3.3 Rebalance ​

触发条件:
- 消费者数量变化 (上下线)
- Topic Queue 数量变化
- 消费者心跳超时

策略 (5.0 统一为兼容 Kafka 的协作式):
  AllocateMessageQueueAveragely       ← 平均分配 (默认)
  AllocateMessageQueueAveragelyByCircle ← 环形平均
  AllocateMessageQueueConsistentHash  ← 一致性 Hash
  AllocateMachineRoomNearby           ← 同机房优先

四、事务消息 ​

4.1 分布式事务流程 ​

Producer                Broker                Consumer
   │                      │                      │
   │ 1. send half msg     │                      │
   │─────────────────────→│                      │
   │        (半消息, 对消费者不可见)              │
   │                      │                      │
   │ 2. 执行本地事务       │                      │
   │ (扣款/减库存...)     │                      │
   │                      │                      │
   │ 3. commit / rollback │                      │
   │─────────────────────→│                      │
   │                      │                      │
   │              ┌──────┴──────┐               │
   │              │ commit?     │               │
   │              │ 将半消息标记 │               │
   │              │ 为可投递     │               │
   │              └──────┬──────┘               │
   │                      │                      │
   │                      │ 4. 投递              │
   │                      │─────────────────────→│
   │                      │                      │
   │    如果步骤 3 超时     │                      │
   │    Broker 回查事务状态  │                      │
   │ ←───────────────────→│                      │

4.2 使用示例 ​

java
// 发送事务消息
TransactionMQProducer producer = new TransactionMQProducer("tx_group");
producer.setTransactionListener(new TransactionListener() {
    @Override
    public LocalTransactionState executeLocalTransaction(Message msg, Object arg) {
        // 执行本地事务
        try {
            orderService.createOrder();    // 本地 DB 操作
            return LocalTransactionState.COMMIT_MESSAGE;
        } catch (Exception e) {
            return LocalTransactionState.ROLLBACK_MESSAGE;
        }
    }

    @Override
    public LocalTransactionState checkLocalTransaction(MessageExt msg) {
        // Broker 回查: 检查本地事务是否提交
        String orderId = msg.getKeys();
        return orderService.isOrderExists(orderId)
            ? LocalTransactionState.COMMIT_MESSAGE
            : LocalTransactionState.ROLLBACK_MESSAGE;
    }
});
producer.start();

// 发送
Message msg = new Message("OrderTopic", "order.create", orderId, body);
producer.sendMessageInTransaction(msg, null);

五、延迟消息 ​

RocketMQ 开箱即用,支持 18 个等级:

1s 5s 10s 30s 1m 2m 3m 4m 5m 6m 7m 8m 9m 10m 20m 30m 1h 2h
java
Message msg = new Message("Topic", "Tag", body);
msg.setDelayTimeLevel(3);  // 等级 3 = 10s
producer.send(msg);

原理:延迟消息先存在 SCHEDULE_TOPIC_XXXX 内部 Topic,定时任务到点捞取投递到目标 Topic。


六、消息过滤 ​

方式说明
Tag 过滤消息标记 Tag,消费端订阅 TopicA || Tag1 || Tag2
SQL92 过滤按消息属性 a > 5 AND b = 'abc'(需 Broker 开启)
类过滤 (4.x)Java 序列化过滤(不推荐)

七、性能与最佳实践 ​

优化说明
消息大小< 512KB(默认 4MB 限制)
Queue 数量与消费者数量匹配,不宜过多(影响 CommitLog 随机读)
异步发送吞吐提升 3-5 倍
批量消费consumeMessageBatchMaxSize 一次拉多条
顺序消息同一业务 ID 发到同一 Queue(MessageQueueSelector)

顺序消息保证 ​

java
// 全局有序: 一个 Topic 只有一个 Queue
// 分区有序 (常见):
producer.send(msg, new MessageQueueSelector() {
    @Override
    public MessageQueue select(List<MessageQueue> mqs, Message msg, Object arg) {
        long orderId = (Long) arg;
        int index = (int) (orderId % mqs.size());
        return mqs.get(index);
    }
}, orderId);

八、Kafka vs RocketMQ 技术选型 ​

维度KafkaRocketMQ
事务消息✅ Kafka Streams 事务✅ 原生事务消息(半消息+回查)
延迟消息❌ 需手动实现✅ 原生 18 级延迟
顺序消息分区内有序Queue 内有序
消息重试需自行死信处理✅ 原生重试 + 死信队列
堆积能力强(千亿级)强(百亿级)
吞吐量极高高
运维复杂度中等较高
生态大数据/流处理阿里云商业版、Spring Cloud Stream
适用日志、流、大数据电商、金融、分布式事务

参考 ​


九、RocketMQ 5.0 新特性 ​

架构演进 ​

mermaid
flowchart TB
    subgraph V4["RocketMQ 4.x"]
        NS4["NameServer<br/>(路由)"]
        B4["Broker<br/>(存储+计算)"]
        P4["Producer"]
        C4["Consumer"]
        P4 --> B4
        B4 --> C4
        B4 --> NS4
    end

    subgraph V5["RocketMQ 5.0"]
        NS5["NameServer<br/>(路由)"]
        Proxy["Proxy<br/>(接入层, gRPC)"]
        B5["Broker<br/>(存储)"]
        Controller["Controller<br/>(选主, Raft)"]
        P5["Producer<br/>(gRPC SDK)"]
        C5["Consumer<br/>(Pop 模式)"]
        P5 -->|"gRPC"| Proxy
        Proxy --> B5
        B5 --> Controller
        Proxy --> C5
    end

5.0 核心变化 ​

特性4.x5.0说明
接入协议自定义 RemotinggRPC多语言 SDK 统一,跨语言更方便
消费模式Pull + PushPop 消费无状态消费,无需 Rebalance
高可用DLedger (可选)Controller 模式Raft 选主,自动故障转移
延迟消息18 个固定级别任意时间延迟基于时间轮,精度秒级
架构Broker 耦合Proxy + Broker 分离计算存储分离,独立扩展
可观测性基础OpenTelemetry原生链路追踪

Pop 消费模式 ​

传统 Pull 模式问题:
  - Consumer 与 Queue 绑定 → Rebalance 时消费暂停
  - Consumer 挂了 → 该 Queue 消息积压直到 Rebalance 完成

Pop 模式(5.0 新增):
  - Consumer 无状态,每次 Pop 请求从 Broker 获取消息
  - 无需 Rebalance,Consumer 随时加入/退出
  - Broker 端管理消费进度(类似 RabbitMQ 的 Push)
  - 适合 Serverless / 弹性伸缩场景
java
// 5.0 gRPC SDK 示例
ClientServiceProvider provider = ClientServiceProvider.loadService();
ClientConfiguration configuration = ClientConfiguration.newBuilder()
    .setEndpoints("proxy-host:8081")
    .build();

// Simple Consumer(Pop 模式)
SimpleConsumer consumer = provider.newSimpleConsumerBuilder()
    .setClientConfiguration(configuration)
    .setConsumerGroup("my-group")
    .setSubscriptionExpressions(Collections.singletonMap("TopicA", FilterExpression.SUB_ALL))
    .build();

// 主动拉取(无 Rebalance)
List<MessageView> messages = consumer.receive(32, Duration.ofSeconds(10));
for (MessageView msg : messages) {
    // 处理消息
    consumer.ack(msg);  // 确认
}

十、运维排查 ​

常用运维命令 ​

bash
# 集群状态
mqadmin clusterList -n localhost:9876

# Topic 列表
mqadmin topicList -n localhost:9876

# Topic 详情(Queue 分布)
mqadmin topicStatus -n localhost:9876 -t OrderTopic

# 消费者组状态(消费进度、积压)
mqadmin consumerProgress -n localhost:9876 -g my-consumer-group

# 查看消息积压(最重要的运维指标)
mqadmin consumerProgress -n localhost:9876 -g my-group
# 输出:
# Topic          Broker     QueueId  Broker Offset  Consumer Offset  Diff
# OrderTopic     broker-a   0        100000         99500            500
# OrderTopic     broker-a   1        100000         98000            2000  ← 积压!

# 按 MessageId 查询消息
mqadmin queryMsgById -n localhost:9876 -i 0A0A0A0A00002A9F000000000000001E

# 按 Key 查询消息
mqadmin queryMsgByKey -n localhost:9876 -t OrderTopic -k order_12345

# 按时间范围查询
mqadmin queryMsgByOffset -n localhost:9876 -t OrderTopic -b broker-a -i 0 -o 99000

# 重置消费位点(回溯消费)
mqadmin resetOffsetByTime -n localhost:9876 -g my-group -t OrderTopic -s "2024-01-01#00:00:00:000"

# 查看 Broker 运行状态
mqadmin brokerStatus -n localhost:9876 -b broker-a:10911

常见问题排查 ​

问题排查命令解决方案
消息积压consumerProgress 看 Diff增加消费者 / 增加线程 / 跳过非关键消息
发送失败查看 Producer 日志 + brokerStatus检查 Broker 磁盘/内存/网络
消费重复查看 Consumer 日志 + Rebalance 频率实现幂等 / 增大 maxReconsumeTimes
Broker 主从切换clusterList 看 Master 变化检查 DLedger/Controller 日志
消息丢失对比 Producer 发送数和 Consumer 消费数开启同步刷盘 + 同步复制
延迟消息不准检查 SCHEDULE_TOPIC_XXXX 积压Broker 负载过高,扩容

监控指标 ​

promql
# RocketMQ Exporter 关键指标

# 消费积压(最重要)
rocketmq_consumer_diff{group="my-group", topic="OrderTopic"}

# Broker 写入 TPS
rocketmq_broker_tps{broker="broker-a"}

# 发送延迟
rocketmq_producer_send_latency_bucket

# 消费延迟
rocketmq_consumer_consume_latency_bucket

# 告警规则
- alert: RocketMQConsumerLag
  expr: rocketmq_consumer_diff > 10000
  for: 5m
  labels:
    severity: warning
  annotations:
    summary: "RocketMQ 消费积压 {{ $value }} 条"
批注模式

💬 文章评论

暂无评论,来说点什么吧 👇

编程学习笔记