消息队列选型与 RocketMQ 原理
消息队列概述
MQ 核心作用
消息队列(Message Queue,MQ)是一种采用生产者-消费者模式的中间件,发送方(Producer)将消息投递到队列中,接收方(Consumer)从队列中拉取消息进行处理。它在分布式系统中扮演着三大核心角色:
解耦
系统之间通过消息队列建立间接通信,生产方无需关心消费方的具体实现。当新增或下线消费方时,生产方零改动。
示例场景:订单服务在创建订单后,需要通知库存服务扣减库存、通知积分服务增加积分、通知物流服务生成运单。如果使用 RPC 直连,每新增一个下游都要修改订单服务代码;引入 MQ 后,订单服务只需发送一条"订单已创建"的消息,所有订阅方各自消费即可。
// 订单服务 — 解耦后只需发送一条消息
OrderCreatedEvent event = new OrderCreatedEvent(orderId, userId, amount);
producer.send("order-topic", event);异步
将同步耗时的调用链路改造为异步处理,显著降低接口响应时间。
示例场景:用户注册需要发送欢迎邮件 + 短信验证 + 初始化空间,这些操作合计耗时 3s。通过 MQ 异步化后,主线程仅需 50ms 写入消息即可返回。
// 同步耗时 ~3s
userService.register(user); // 50ms
emailService.sendWelcome(user); // 500ms
smsService.sendNotify(user); // 300ms
initService.createWorkspace(user); // 2000ms
// 异步化后 ~50ms
userService.register(user);
producer.send("notify-topic", user);削峰填谷
应对瞬时流量洪峰,将突增请求暂存在队列中,下游按自身处理能力匀速消费,防止系统被流量冲垮。
示例场景:秒杀活动中,瞬间涌入 10 万请求,后端数据库只能承受 1000 QPS。MQ 将请求全部接入队列,后端以 1000 TPS 匀速处理。
流量曲线:
请求量 ^
10万 | ████████▁▁▁▁▁▁▁▁▁▁ ← 削峰:消息暂存在 MQ
| ████████▁▁▁▁▁▁▁▁▁▁
1000 | ▁▁▁▁▁▁▁▁▁████████▁ ← 填谷:后端匀速消费
+─────────────────────────> 时间使用场景总结
| 场景 | 说明 | 案例 |
|---|---|---|
| 异步处理 | 非关键路径异步化,提升响应速度 | 注册通知、日志上报 |
| 流量削峰 | 缓冲瞬时高流量 | 秒杀、抢红包 |
| 系统解耦 | 减少服务间直接依赖 | 订单 → 履约流水线 |
| 日志处理 | 海量日志收集与分发 | 业务审计、监控采集 |
| 事件驱动 | 基于事件的 CQRS/EDA 架构 | 领域事件、变更数据捕获 |
| 最终一致性 | 分布式事务的可靠异步通信 | 跨服务状态同步 |
消息队列对比
RocketMQ vs Kafka vs RabbitMQ vs Pulsar
| 对比维度 | RocketMQ | Kafka | RabbitMQ | Pulsar |
|---|---|---|---|---|
| 开发语言 | Java | Scala/Java | Erlang | Java |
| 一致性模型 | 最终一致性 / 强一致(DLedger) | 最终一致性 | 可选(强一致/最终) | 强一致性(BookKeeper) |
| 顺序消息 | ✅ 分区顺序 + 全局顺序 | ✅ 分区顺序 | ❌ 不保证全局 | ✅ 分区顺序 |
| 事务消息 | ✅ 完整支持(半消息 + 回查) | ⚠️ 事务性 API(Kafka 0.11+) | ⚠️ 基于 AMQP 事务 | ✅ 支持 |
| 延迟消息 | ✅ 内置 18 个级别 | ❌ 需应用层实现 | ✅ 插件支持 | ❌ 需应用层实现 |
| 死信队列 | ✅ 内置 DLQ 机制 | ✅ DLQ 需 Streams 配置 | ✅ 死信交换机 | ❌ 需插件或自定义 |
| 消息回溯 | ✅ 按时间/偏移量回溯 | ✅ 按偏移量回溯 | ❌ 消费后删除 | ⚠️ 有限支持(如开启) |
| 批量能力 | 支持 | 极强 | 弱 | 支持 |
| 吞吐量 | 10 万级/秒 | 百万级/秒 | 万级/秒 | 百万级/秒 |
| 延迟 | 毫秒级 | 毫秒级 | 微秒级 | 毫秒级 |
| 运维复杂度 | 中等(依赖 NameServer) | 中等(依赖 ZooKeeper/KRaft) | 简单 | 较高(依赖 BookKeeper + ZooKeeper) |
| 客户端语言 | Java 生态为主(C++/Go 次之) | 多语言丰富 | 多语言丰富 | 多语言丰富 |
| 公司/社区 | 阿里巴巴 → Apache | Apache/LinkedIn | VMware/Pivotal | Apache/Splunk |
| 协议 | 自定义 TCP | 自定义 TCP | AMQP 0-9-1 | 自定义 TCP + 兼容 Kafka |
选型建议
- 追求高吞吐、海量日志 → 选 Kafka
- 需要事务消息、延迟消息、顺序消息,且技术栈以 Java 为主 → 选 RocketMQ
- 功能丰富、运维简单、中小规模 → 选 RabbitMQ
- 云原生、多租户、计算存储分离 → 选 Pulsar
RocketMQ 架构
核心组件
┌─────────────────────────────────────────────────────┐
│ NameServer │
│ 注册中心、路由管理、轻量级协调 │
└──────────┬──────────────────┬───────────────────────┘
│ 注册 / 心跳 │ 路由发现
┌──────▼──────┐ ┌─────▼──────┐
│ Broker │ │ Producer │
│ │ │ │
│ Topic A Q0 │ │ 消息发送者 │
│ Topic A Q1 │ └────────────┘
│ Topic B Q0 │
└──────┬──────┘
│ 拉取 / 推送
┌──────▼──────┐
│ Consumer │
│ 消息消费者 │
└─────────────┘NameServer
- 无状态的注册中心,各个 NameServer 节点不互相通信
- Broker 启动时向所有 NameServer 注册自身信息(IP、端口、持有的 Topic)
- 提供轻量级的路由发现 API,Producer 和 Consumer 通过 NameServer 获取 Broker 地址
- 支持多节点部署防止单点故障
NameServer 配置示例(namesrv.properties):
# 监听端口
listenPort=9876
# 存储路径
storePathRootDir=/data/rocketmq/namesrvBroker
- 消息存储和中转的核心节点
- 负责消息的接收、存储、投递
- 每个 Broker 可以持有多个 Topic 的多个 Queue(读写队列)
- Broker 启动后会向所有 NameServer 注册
- 支持主从架构(Master/Slave)和 DLedger 多副本
Broker 配置示例(broker.properties):
# Broker 集群名称
brokerClusterName=DefaultCluster
# Broker 名称
brokerName=broker-a
# Broker 角色: ASYNC_MASTER / SYNC_MASTER / SLAVE
brokerRole=ASYNC_MASTER
# 刷盘方式: ASYNC_FLUSH / SYNC_FLUSH
flushDiskType=ASYNC_FLUSH
# Namesrv 地址
namesrvAddr=192.168.1.1:9876;192.168.1.2:9876
# 存储路径
storePathRootDir=/data/rocketmq/store
# 自动创建 Topic 开关
autoCreateTopicEnable=trueProducer
- 消息的生产者,从 NameServer 获取 Topic 路由信息
- 将消息发送到 Broker 上对应的 Queue
- 支持同步发送、异步发送、单向发送三种方式
Consumer
- 消息的消费者,从 Broker 拉取消息进行消费
- 支持集群消费和广播消费两种模式
- 通过消费组(Consumer Group)组织,同组内的消费者分摊消息
Topic 与 Queue
Topic 是消息的逻辑分类,生产者将消息发送到指定 Topic,消费者订阅 Topic 进行消费。
Queue(也称为 MessageQueue)是 Topic 下物理存储的最小单元。一个 Topic 可以包含多个 Queue,Queue 分布在不同的 Broker 上,是水平扩展和负载均衡的基础。
Topic "OrderTopic"
├── Queue 0 (Broker-A)
├── Queue 1 (Broker-A)
├── Queue 2 (Broker-B)
└── Queue 3 (Broker-B)- 每个 Queue 在磁盘上对应一组 CommitLog + ConsumeQueue 文件
- Queue 数量决定了消息的并发吞吐上限
- 生产者轮询或指定 Queue 发送消息
- 消费者通过负载均衡分配到不同的 Queue 上
Offset
Offset(偏移量) 是 Consumer 在 Queue 中的消费进度标识。
- 本地 Offset 模式:由 Consumer 客户端在本地内存维护
- 远程 Offset 模式(推荐):Offset 存储在 Broker 端,便于消费者重启后恢复进度
- 每个 Consumer Group 在每个 Queue 上维护一个独立的 Offset
// 查看消费进度(Admin CLI)
$ mqadmin consumerProgress -g consumer-group-a -n 127.0.0.1:9876
# Broker Name Queue Consumer Offset Max Offset Diff
# broker-a 0 150 200 50
# broker-a 1 180 200 20消息发送
发送方式
RocketMQ 提供三种发送方式,分别适用于不同的场景。
同步发送(Sync)
发送后等待 Broker 返回写入确认结果,保证消息绝对可靠,适用于重要业务消息。
// 同步发送 — 订单支付成功通知
SendResult result = producer.send(
new Message("order-topic", "pay", orderJson.getBytes())
);
if (result.getSendStatus() == SendStatus.SEND_OK) {
log.info("消息发送成功, msgId={}", result.getMsgId());
} else {
log.warn("消息发送失败, status={}", result.getSendStatus());
}异步发送(Async)
发送后不阻塞,通过回调获取结果,适用于对延迟敏感的场景。
producer.send(
new Message("log-topic", "", logJson.getBytes()),
new SendCallback() {
@Override
public void onSuccess(SendResult sendResult) {
log.info("异步发送成功, msgId={}", sendResult.getMsgId());
}
@Override
public void onException(Throwable e) {
log.error("异步发送失败", e);
// 执行补偿逻辑
}
}
);
// 此处继续执行其他操作,不阻塞单向发送(One-way)
只发送不等待任何确认,吞吐量最高,适用于可靠性要求低的场景。
// 单向发送 — 日志、监控数据,丢弃不影响业务
producer.sendOneway(new Message("metrics-topic", "", metricData));三种方式对比:
| 方式 | 延迟 | 吞吐量 | 可靠性 | 适用场景 |
|---|---|---|---|---|
| 同步 | 高 | 低 | 最高 | 订单、支付、通知 |
| 异步 | 低 | 高 | 高 | 日志、监控 |
| 单向 | 最低 | 最高 | 低 | 非关键日志、遥测 |
消息类型
普通消息
不保证顺序,RocketMQ 默认的消息类型,吞吐量最高。
Message msg = new Message("normal-topic", "tagA", body);
producer.send(msg);顺序消息
保证同一队列内消息按照发送顺序被消费。详见 消息顺序 章节。
事务消息
用于实现分布式事务。详见 事务消息 章节。
延迟消息
消息不会立即投递到消费者,而是在指定的延迟时间后才可见。
RocketMQ 内置 18 个延迟级别:
| 级别 | 1 | 2 | 3 | 4 | 5 | 6 | 7 | 8 | 9 |
|---|---|---|---|---|---|---|---|---|---|
| 延迟 | 1s | 5s | 10s | 30s | 1m | 2m | 3m | 4m | 5m |
| 级别 | 10 | 11 | 12 | 13 | 14 | 15 | 16 | 17 | 18 |
|---|---|---|---|---|---|---|---|---|---|
| 延迟 | 6m | 7m | 8m | 9m | 10m | 20m | 30m | 1h | 2h |
Message msg = new Message("order-topic", "", body);
// 设置延迟级别 3 — 延迟 10s 后投递
msg.setDelayTimeLevel(3);
producer.send(msg);应用场景:订单超时未支付自动取消、定时提醒、重试间隔。
消息消费
Push vs Pull
RocketMQ 的消费模式本质上是 Long Polling(长轮询),客户端轮询 Broker 拉取消息。
| 模式 | 实现 | 特点 |
|---|---|---|
| Push | DefaultMQPushConsumer | 消费端注册监听器,框架自动拉取并回调;使用方感知为"推送" |
| Pull | DefaultMQPullConsumer | 使用方自行控制拉取时机、Offset 管理,灵活度高 |
Push 模式(推荐):
DefaultMQPushConsumer consumer = new DefaultMQPushConsumer("consumer-group");
consumer.setNamesrvAddr("127.0.0.1:9876");
consumer.subscribe("order-topic", "*");
consumer.registerMessageListener((MessageListenerConcurrently) (msgs, context) -> {
for (MessageExt msg : msgs) {
System.out.printf("消费消息: %s%n", new String(msg.getBody()));
}
return ConsumeConcurrentlyStatus.CONSUME_SUCCESS;
});
consumer.start();Pull 模式:
DefaultLitePullConsumer consumer = new DefaultLitePullConsumer("consumer-group");
consumer.setNamesrvAddr("127.0.0.1:9876");
consumer.subscribe("order-topic", "*");
consumer.start();
while (true) {
List<MessageExt> msgs = consumer.poll(1000);
for (MessageExt msg : msgs) {
// 自行处理消息和 Offset
process(msg);
}
}集群消费与广播消费
集群消费(Clustering)
同一条消息在同一消费组内只会被一个消费者实例消费,消息按负载均衡规则分配到组内各实例。这是默认消费模式。
Producer ──→ Topic ──→ Queue 0 ──→ Consumer A (group-a)
Queue 1 ──→ Consumer B (group-a)
Queue 2 ──→ Consumer A (group-a)
Queue 3 ──→ Consumer B (group-a)每条消息只会被 group-a 中的一个消费者处理。
consumer.setMessageModel(MessageModel.CLUSTERING); // 默认广播消费(Broadcasting)
同一条消息在消费组内被所有消费者实例消费一次。适用于配置同步、缓存刷新等场景。
consumer.setMessageModel(MessageModel.BROADCASTING);重试机制
消息消费失败后,RocketMQ 会在重试队列(%RETRY%{consumerGroup})中进行延时重试。
重试级别递增:
第1次重试:10s 第2次:30s 第3次:1m
第4次:2m 第5次:3m 第6次:4m
第7次:5m 第8次:6m 第9次:7m
第10次:8m 第11次:9m 第12次:10m
第13次:20m 第14次:30m 第15次:1h
第16次:2h// 消费端处理重试
public ConsumeConcurrentlyStatus consumeMessage(
List<MessageExt> msgs, ConsumeConcurrentlyContext context) {
try {
process(msgs);
return ConsumeConcurrentlyStatus.CONSUME_SUCCESS;
} catch (Exception e) {
// 返回 RECONSUME_LATER 触发重试
return ConsumeConcurrentlyStatus.RECONSUME_LATER;
}
}死信队列
当消息重试次数超过最大阈值(默认 16 次)后,消息会被投递到死信队列(Dead Letter Queue, DLQ)。死信队列的 Topic 名称为 %DLQ%{consumerGroup}。
# Broker 配置:设置最大重试次数
maxReconsumeTimes=16死信队列处理策略:
- 人工介入:通过
mqadmin查询死信消息,确定原因后重新投递 - 自动化处理:编写单独的消费者订阅
%DLQ%队列,记录日志并告警
# 查询死信消息
$ mqadmin queryMsgByKey -t %DLQ%consumer-group-a -k msgKey -n 127.0.0.1:9876
# 重新投递死信消息到原 Topic
$ mqadmin sendMsgBack -t order-topic -g consumer-group-a -m <msgId>事务消息
概述
RocketMQ 的事务消息用于解决分布式事务问题,典型场景是跨服务的数据一致性(如订单创建 + 积分扣减 + 库存预占)。
RocketMQ 通过半消息(Half Message)+ 事务回查机制,实现了最终一致性的分布式事务方案。
半消息机制
事务消息的执行流程分为三个阶段:
Producer Broker
│ │
├── 1. 发送半消息 ──────────→ │ 半消息对 Consumer 不可见
│ │
├── 2. 执行本地事务 │
│ (如扣库存、写订单) │
│ │
├── 3. 提交/回滚事务消息 ───→ │
│ COMMIT → Consumer 可见 │
│ ROLLBACK → 消息丢弃 │
│ │
├── 4. [回查] ←─────────── │ Broker 未收到 3,主动回查
│ 检查本地事务结果 │
└── 5. 提交/回滚 ──────────→ │关键步骤说明:
- 发送半消息:Producer 向 Broker 发送半消息,该消息对 Consumer 不可见(标记为 Prepared)
- 执行本地事务:Producer 执行本地业务逻辑(如数据库操作)
- 提交或回滚:根据本地事务结果向 Broker 发送 Commit 或 Rollback
- 事务回查:如果 Broker 长时间未收到 Commit/Rollback,会主动回调 Producer 检查事务状态
代码示例
public class OrderTransactionProducer {
public static void main(String[] args) throws MQClientException {
// 1. 创建事务消息生产者
TransactionMQProducer producer = new TransactionMQProducer("tx-producer-group");
producer.setNamesrvAddr("127.0.0.1:9876");
// 2. 设置事务监听器(包含本地事务执行 + 回查逻辑)
producer.setTransactionListener(new TransactionListener() {
@Override
public LocalTransactionState executeLocalTransaction(Message msg, Object arg) {
// 步骤 2:执行本地事务(扣库存、创建订单)
String orderId = arg.toString();
try {
orderService.createOrder(orderId);
// 本地事务成功 → 提交消息
return LocalTransactionState.COMMIT_MESSAGE;
} catch (Exception e) {
// 本地事务失败 → 回滚消息
return LocalTransactionState.ROLLBACK_MESSAGE;
}
}
@Override
public LocalTransactionState checkLocalTransaction(MessageExt msg) {
// 步骤 4:回查 — Broker 未收到 Commit/Rollback 时触发
String orderId = msg.getKeys();
if (orderService.isOrderExists(orderId)) {
return LocalTransactionState.COMMIT_MESSAGE;
}
return LocalTransactionState.UNKNOW;
}
});
producer.start();
// 3. 发送半消息
String orderId = "ORD-" + System.currentTimeMillis();
Message msg = new Message("order-tx-topic", "tagA", orderId,
orderJson.getBytes());
msg.setKeys(orderId);
// 发送事务消息
SendResult result = producer.sendMessageInTransaction(msg, orderId);
log.info("事务消息发送结果: {}", result);
}
}消费者处理事务消息
// 消费者无需特殊处理,直接消费事务提交后的消息
DefaultMQPushConsumer consumer = new DefaultMQPushConsumer("tx-consumer-group");
consumer.subscribe("order-tx-topic", "*");
consumer.registerMessageListener((MessageListenerConcurrently) (msgs, context) -> {
for (MessageExt msg : msgs) {
String orderId = msg.getKeys();
// 处理订单后续逻辑(积分、物流等)
fulfillmentService.processOrder(orderId);
}
return ConsumeConcurrentlyStatus.CONSUME_SUCCESS;
});
consumer.start();分布式事务示意图
┌───────────┐ ┌───────────┐ ┌───────────┐
│ 订单服务 │ │ RocketMQ │ │ 积分服务 │
│ (Producer) │ │ (Broker) │ │ (Consumer) │
└─────┬─────┘ └─────┬─────┘ └─────┬─────┘
│ │ │
│ 1. 发送半消息 │ │
│────────────────────→│ │
│ │ │
│ 2. 执行本地事务 │ │
│ (写订单表) │ │
│───────┐ │ │
│ │ 成功 │ │
│←──────┘ │ │
│ │ │
│ 3. COMMIT │ │
│────────────────────→│ │
│ │ 4. 投递消息 │
│ │────────────────────→│
│ │ │ 5. 增加积分
│ │ │───────┐
│ │ │ │ 成功
│ │ │←──────┘适用场景
| 场景 | 说明 |
|---|---|
| 订单 — 库存 | 创建订单同时预占库存,任一失败则回滚 |
| 账户 — 积分 | 资金变动同时更新积分 |
| 交易 — 结算 | 交易成功触发结算流水写入 |
注意事项
- 幂等性:事务消息可能被重复投递,消费方必须实现幂等处理
- 回查超时:合理设置回查间隔和最大回查次数
- 本地事务与消息发送:确保本地事务与消息 Commit 的原子性,半消息机制正是为此设计
消息顺序
全局顺序
全局顺序要求一个 Topic 下所有消息都严格按发送顺序消费。
实现方式:Topic 只设置一个 Queue,所有消息写入同一条队列,Consumer 单线程消费。
// Producer 端
producer.setDefaultTopicQueueNums(1);
producer.send(new Message("global-order-topic", body));
// Consumer 端 — 顺序消费
consumer.registerMessageListener((MessageListenerOrderly) (msgs, context) -> {
// MessageListenerOrderly 保证单线程串行消费
for (MessageExt msg : msgs) {
process(msg);
}
return ConsumeOrderlyStatus.SUCCESS;
});权衡:吞吐量受限于单 Queue 性能,适用于要求不高的小流量场景。
分区顺序
分区顺序保证同一业务标识(如 OrderId)的消息按顺序消费,不同标识的消息可并行处理。这是生产中最常用的方案。
实现要点:
- Producer 端:相同业务标识的消息发送到同一个 Queue
- Consumer 端:使用
MessageListenerOrderly消费
// Producer 端 — 保证同一订单的消息进入同一 Queue
String orderId = "ORD-12345";
Message msg = new Message("order-seq-topic", body);
// 使用订单 ID 作为选择 HASH 的 Key
producer.send(msg, new MessageQueueSelector() {
@Override
public MessageQueue select(List<MessageQueue> mqs, Message msg, Object arg) {
String id = (String) arg;
// 对订单 ID 取 HASH,确保同一订单进入同一 Queue
int index = Math.abs(id.hashCode()) % mqs.size();
return mqs.get(index);
}
}, orderId);// Consumer 端 — 顺序消费
DefaultMQPushConsumer consumer = new DefaultMQPushConsumer("order-seq-group");
consumer.setNamesrvAddr("127.0.0.1:9876");
consumer.subscribe("order-seq-topic", "*");
consumer.registerMessageListener((MessageListenerOrderly) (msgs, context) -> {
// MessageListenerOrderly 保证队列级别的串行消费
// 内部使用锁机制,确保同一 Queue 内的消息不会被并行消费
for (MessageExt msg : msgs) {
Order order = JSON.parseObject(msg.getBody(), Order.class);
orderService.processOrder(order);
}
return ConsumeOrderlyStatus.SUCCESS;
});
consumer.start();顺序保证的原理
| 层面 | 机制 |
|---|---|
| Producer 端 | 通过 MessageQueueSelector 将相同标识的消息路由到同一 Queue |
| Broker 端 | 同一 Queue 内消息写入 CommitLog,天然有序 |
| Consumer 端 | MessageListenerOrderly 采用锁机制(Broker 锁 + 本地锁),保证同一 Queue 串行消费 |
MessageListenerOrderly 与 MessageListenerConcurrently 的区别:
| 特征 | MessageListenerOrderly | MessageListenerConcurrently |
|---|---|---|
| 并发消费 | 单线程串行(队列级别) | 多线程并行 |
| 重试行为 | 无限重试(阻塞) | 有限重试 + 死信队列 |
| 锁机制 | Broker 分布式锁 + 本地锁 | 无锁 |
| 适用场景 | 顺序要求严格 | 吞吐优先 |
高可用
DLedger 主从切换
RocketMQ 4.5+ 引入 DLedger 实现基于 Raft 协议的自动主从切换,替代传统的人工主从切换模式。
传统主从架构
Broker-A (Master) ──── 异步同步 ────→ Broker-A (Slave)
↑ ↑
手动切换 仅备份缺点:Master 宕机后需要手动切换 Slave 为 Master,存在分钟级不可用窗口。
DLedger 架构
┌─────────────┐ ┌─────────────┐ ┌─────────────┐
│ Broker-A │ │ Broker-B │ │ Broker-C │
│ (Leader) │━━━━━│ (Follower) │━━━━━│ (Follower) │
└──────┬──────┘ └─────────────┘ └─────────────┘
│ 读写 ↑ ↑
│ └── Raft 日志复制 ────┘
Producer/Consumer核心特性:
- 基于 Raft 协议实现强一致性多副本
- 使用 DLedger CommitLog 替代原生的 CommitLog
- 自动 Leader 选举,故障时秒级切换
- 多数派(Quorum)写入确认,保证数据不丢失
# broker.conf — DLedger 模式配置
brokerClusterName=DLedgerCluster
brokerName=RaftNode
# 启用 DLedger
enableDLegerCommitLog=true
# DLedger 集群节点组
dLegerGroup=DLedgerGroup
# 当前节点 ID(要求集群内唯一,递增)
dLegerPeers=n0-192.168.1.1:40911;n1-192.168.1.2:40911;n2-192.168.1.3:40911
# 本节点 ID
dLegerSelfId=n0
# 发送消息超时
sendMessageThreadPoolNums=16多副本机制对比
| 特性 | 传统主从 | DLedger |
|---|---|---|
| 副本同步 | 异步 / 同步 | Raft 多数派写入 |
| 故障切换 | 手动切换或 VIP 漂移 | 自动选举 |
| 数据一致性 | 最终一致性 | 强一致性 |
| 写入性能 | 高(异步刷盘) | 较低(多数派确认) |
| 适用场景 | 允许少量丢数据 | 数据零丢失 |
| 部署节点 | 2 个(1M + 1S) | 3 个(多数派) |
刷盘策略
| 策略 | 配置 | 说明 |
|---|---|---|
| 异步刷盘 | flushDiskType=ASYNC_FLUSH | 消息写入 PageCache 即返回,后台线程刷盘。吞吐量高,故障可能丢失少量数据 |
| 同步刷盘 | flushDiskType=SYNC_FLUSH | 消息写入物理磁盘后才返回。吞吐量较低,数据零丢失 |
# 同步刷盘配置
flushDiskType=SYNC_FLUSH
# 同步刷盘间隔(毫秒)
flushIntervalConsumeQueue=1000主从复制
| 策略 | 配置 | 说明 |
|---|---|---|
| 异步复制 | brokerRole=ASYNC_MASTER | Master 写入后立即返回,Slave 异步拉取。性能最好 |
| 同步复制 | brokerRole=SYNC_MASTER | Master 等待 Slave 写入确认后才返回。数据更安全 |
高可用部署建议
- NameServer:至少部署 2 个节点,无状态,前方挂 SLB
- Broker:至少 2 主 2 从,或 3 节点 DLedger 集群
- 刷盘策略:核心业务使用同步刷盘,日志类使用异步刷盘
- 主从复制:核心业务使用同步复制,确保故障时数据不丢
Spring Boot 集成
依赖引入
使用 rocketmq-spring-boot-starter 简化集成。
<dependency>
<groupId>org.apache.rocketmq</groupId>
<artifactId>rocketmq-spring-boot-starter</artifactId>
<version>2.3.0</version>
</dependency>配置文件
# application.properties
rocketmq.name-server=127.0.0.1:9876
# Producer 配置
rocketmq.producer.group=spring-producer-group
rocketmq.producer.send-message-timeout=3000
rocketmq.producer.retry-times-when-send-failed=2
# Consumer 配置(可在代码中覆盖)
rocketmq.consumer.group=spring-consumer-group
rocketmq.consumer.consume-thread-min=10
rocketmq.consumer.consume-thread-max=20消息生产者
@Component
public class OrderProducer {
@Autowired
private RocketMQTemplate rocketMQTemplate;
/**
* 发送普通消息
*/
public void sendOrderMessage(Order order) {
rocketMQTemplate.convertAndSend("order-topic:tagA", order);
log.info("发送订单消息成功: orderId={}", order.getOrderId());
}
/**
* 发送同步消息(带结果)
*/
public SendResult sendSync(Order order) {
Message<Order> msg = MessageBuilder
.withPayload(order)
.setHeader("orderId", order.getOrderId())
.build();
SendResult result = rocketMQTemplate.syncSend("order-topic:tagB", msg);
return result;
}
/**
* 发送异步消息
*/
public void sendAsync(Order order) {
rocketMQTemplate.asyncSend("order-topic:tagC", order,
new SendCallback() {
@Override
public void onSuccess(SendResult result) {
log.info("异步发送成功: {}", result.getMsgId());
}
@Override
public void onException(Throwable e) {
log.error("异步发送失败", e);
}
});
}
/**
* 发送顺序消息 — 同一订单 ID 进入同一 Queue
*/
public void sendOrderly(Order order) {
rocketMQTemplate.syncSendOrderly(
"order-seq-topic", order, order.getOrderId());
}
/**
* 发送延迟消息 — 30 分钟后超时取消
*/
public void sendDelay(Order order) {
// delayLevel=16 对应 30 分钟
rocketMQTemplate.syncSend(
"order-delay-topic",
MessageBuilder.withPayload(order).build(),
3000,
16);
}
/**
* 发送事务消息
*/
public void sendTransaction(Order order) {
rocketMQTemplate.sendMessageInTransaction(
"order-tx-topic",
MessageBuilder.withPayload(order).build(),
order.getOrderId());
}
}消息消费者
@Component
@RocketMQMessageListener(
topic = "order-topic",
consumerGroup = "order-consumer-group",
selectorExpression = "tagA || tagB",
consumeMode = ConsumeMode.CONCURRENTLY, // 并发消费
messageModel = MessageModel.CLUSTERING // 集群模式
)
public class OrderConsumer implements RocketMQListener<Order> {
@Override
public void onMessage(Order order) {
log.info("收到订单消息: {}", order.getOrderId());
// 处理订单
orderService.handleOrder(order);
}
}顺序消息消费者:
@Component
@RocketMQMessageListener(
topic = "order-seq-topic",
consumerGroup = "order-seq-consumer-group",
consumeMode = ConsumeMode.ORDERLY // 顺序消费模式
)
public class OrderSeqConsumer implements RocketMQListener<Order> {
@Override
public void onMessage(Order order) {
// 串行处理同 Queue 内的消息
orderService.processSeqOrder(order);
}
}事务消息监听器:
@Component
@RocketMQTransactionListener(rocketMQTemplateBeanName = "rocketMQTemplate")
public class OrderTransactionListener implements RocketMQLocalTransactionListener {
@Autowired
private OrderService orderService;
@Override
public RocketMQLocalTransactionState executeLocalTransaction(Message msg, Object arg) {
// 执行本地事务
String orderId = (String) arg;
try {
orderService.createOrder(orderId);
return RocketMQLocalTransactionState.COMMIT;
} catch (Exception e) {
log.error("本地事务失败", e);
return RocketMQLocalTransactionState.ROLLBACK;
}
}
@Override
public RocketMQLocalTransactionState checkLocalTransaction(Message msg) {
// 事务回查
String orderId = new String((byte[]) msg.getPayload());
if (orderService.isOrderExists(orderId)) {
return RocketMQLocalTransactionState.COMMIT;
}
return RocketMQLocalTransactionState.UNKNOWN;
}
}完整配置示例(YAML)
rocketmq:
name-server: 192.168.1.1:9876;192.168.1.2:9876
producer:
group: shop-producer
send-message-timeout: 5000
retry-times-when-send-failed: 3
retry-times-when-send-async-failed: 2
consumer:
group: shop-consumer
consume-thread-min: 5
consume-thread-max: 20
consume-concurrently-max-span: 2000
pull-batch-size: 32注意事项
- 客户端版本兼容:确保
rocketmq-spring-boot-starter版本与 RocketMQ 服务端版本匹配 - 消费幂等性:RocketMQ 可能重复投递消息,消费端需通过业务键(
msg.getKeys())去重 - 线程池隔离:不同业务类型的消费者使用不同的线程池配置,防止相互影响
- 监控:集成 RocketMQ Console 或 Prometheus + Grafana 监控消费积压和吞吐