RocketMQ 消息驱动
定位
RocketMQ 解决什么问题
在微服务架构中,服务之间的同步调用(如 HTTP + Feign)形成强依赖链条:订单服务调用通知服务 → 通知服务宕机 → 订单创建失败。RocketMQ 通过消息队列(Message Queue)将同步调用转为异步通信,核心解决三个问题:
- 异步解耦:生产者(Producer)发消息后立即返回,消费者(Consumer)异步消费。生产者不关心消费者是谁、是否在线。
- 削峰填谷:瞬时流量高峰(如秒杀)涌入的消息先在 Broker 中暂存,消费者按自身处理能力匀速消费,避免下游被冲垮。
- 最终一致性:通过事务消息保证本地事务与消息发送的原子性,确保跨服务数据最终一致。
与 RabbitMQ / Kafka 的定位对比
| 维度 | RocketMQ | RabbitMQ | Kafka |
|---|---|---|---|
| 设计定位 | 金融级业务消息 | 通用消息代理 | 流数据平台 |
| 消息模型 | Topic + Tag + Queue | Exchange + Binding + Queue | Topic + Partition |
| 事务消息 | 原生支持 | 需插件(补偿) | 1.0+ 支持 |
| 延时消息 | 18 个固定级别 | 死信队列模拟 | 不原生支持 |
| 顺序消息 | 分区有序 + 全局有序 | 单 Queue 有序 | 分区内有序 |
| 吞吐量 | 十万级 TPS | 万级 TPS | 百万级 TPS |
| 典型场景 | 订单、支付、异步通知 | 低延迟通用消息 | 日志采集、流计算 |
RocketMQ 在电商交易、金融支付等对事务消息和顺序消息要求高的场景更有优势。
在 Spring Cloud 生态中的位置
Spring Cloud Stream 提供了消息驱动的抽象层,核心概念为 Binder(绑定器)、Binding(绑定)和 Channel(通道)。开发者面向统一 API 编程,切换底层 MQ 只需更换 Binder 依赖:
- RabbitMQ:
spring-cloud-starter-stream-rabbit - Kafka:
spring-cloud-starter-stream-kafka - RocketMQ:
spring-cloud-starter-stream-rocketmq
切换 Binder 后,Producer 和 Consumer 代码 零修改,只需调整 application.yml 配置。
本教程只讲与 Spring Cloud Stream 的集成用法,不展开消息队列的深度原理(存储引擎、刷盘机制、主从同步等)。
核心概念
RocketMQ 原生概念
| 概念 | 英文全称 | 说明 |
|---|---|---|
| Topic | Topic | 消息主题,消息分类的一级标签。例如 order-topic 存放订单消息 |
| Tag | Tag | 消息标签,Topic 下的二级分类,Consumer 可按 Tag 过滤。例如 order-topic:create、order-topic:cancel |
| Consumer Group | Consumer Group | 消费者组。集群模式(Clustering)下同组消费者共同消费,每条消息只被组内一个实例处理;广播模式(Broadcasting)下每条消息被组内所有实例处理 |
| Message Queue | Message Queue | 消息队列,Topic 下的物理分区(Partition),数量决定并行消费度 |
| NameServer | Name Server | 路由注册中心,无状态且相互独立。Broker 注册路由信息,Producer/Consumer 查询路由表 |
消息类型
| 类型 | 说明 | 典型场景 |
|---|---|---|
| 普通消息(Normal Message) | 同步/异步/单向发送,无特殊顺序或延迟要求 | 订单支付通知、短信发送 |
| 顺序消息(Ordered Message) | 同一业务标识的消息进入同一 Queue,单线程按 FIFO 消费 | 订单状态流转(创建→支付→发货→完成) |
| 延时消息(Delay Message) | 指定延迟级别,到达延迟时间后才对 Consumer 可见 | 订单超时 30 分钟未支付自动取消 |
| 事务消息(Transaction Message) | 半消息机制:先发送半消息(对 Consumer 不可见),本地事务执行成功则提交(可见),失败则回滚(丢弃) | 订单创建 + 库存扣减的最终一致性 |
消费模式
| 模式 | 英文 | 说明 |
|---|---|---|
| 集群消费 | Clustering | 同一 Consumer Group 内所有实例共享消费进度,每条消息只被一个实例消费。默认模式 |
| 广播消费 | Broadcasting | 每条消息推送给 Consumer Group 内所有实例,各自独立维护消费进度 |
Spring Cloud Stream 抽象
Spring Cloud Stream 将底层 MQ 概念抽象为统一的编程模型:
| 概念 | 说明 | RocketMQ 对应 |
|---|---|---|
| Binder(绑定器) | 与底层 MQ 的适配层,负责连接、发送、消费 | RocketMQ Binder(spring-cloud-starter-stream-rocketmq) |
| Binding(绑定) | 连接通道和 Binder 的桥梁,声明 Channel 与 Topic 的映射关系 | spring.cloud.stream.bindings 配置 |
| Destination | 消息目的地,即 RocketMQ 的 Topic | destination: order-topic |
函数式编程模型(Spring Cloud Stream 3.x+)
Spring Cloud Stream 3.x 全面转向函数式编程,废弃 @Input / @Output / @EnableBinding 等注解。核心 Bean 定义方式:
- Supplier<T>:消息生产者,应用启动后持续生成消息
- Consumer<T>:消息消费者,接收并处理消息
- Function<T, R>:消息处理器,接收消息、处理后返回新消息(类似 Bridge/Kafka Streams)
实际开发中更常用 StreamBridge(动态发送)替代 Supplier,通过 streamBridge.send("bindingName", message) 在业务方法中按需发送。
核心原理
RocketMQ 核心组件 + Spring Cloud Stream 适配层
事务消息半消息流程
Spring Cloud Stream Binder 映射关系
Spring Cloud Stream Binder 将 RocketMQ 概念映射为抽象模型的过程如下:
application.yml Spring Cloud Stream RocketMQ
───────────────────────────────────── ──────────────────────────── ───────────────────
spring.cloud.stream.bindings:
output: Binding: output ─────────► Topic: order-topic
destination: order-topic MessageChannel ──────────► Producer Group
content-type: application/json Content Type ───────────► 序列化格式
input: Binding: input ──────────► Topic: order-topic
destination: order-topic Consumer<T> Bean ────────► Consumer Group
group: order-consumer-group Group ───────────────────► order-consumer-group
Binder 自动完成:
- 根据
destination创建/绑定 Topic - 根据
group创建/加入 Consumer Group - 管理 Producer 到 Broker 的连接池
- 管理 Consumer 的 Rebalance(重平衡)和 Offset(消费位点)提交
环境准备
第一步:下载并启动 NameServer
# 下载 RocketMQ 5.x 二进制包
wget https://dist.apache.org/repos/dist/release/rocketmq/5.3.1/rocketmq-all-5.3.1-bin-release.zip
unzip rocketmq-all-5.3.1-bin-release.zip
cd rocketmq-all-5.3.1-bin-release
# 启动 NameServer(默认监听 9876 端口)
# Windows
.\bin\mqnamesrv.cmd
# Linux/macOS
nohup sh bin/mqnamesrv &
第二步:启动 Broker
# Windows
.\bin\mqbroker.cmd -n localhost:9876 autoCreateTopicEnable=true
# Linux/macOS
nohup sh bin/mqbroker -n localhost:9876 autoCreateTopicEnable=true &
autoCreateTopicEnable=true 允许 Producer 自动创建 Topic(生产环境建议关闭)。
第三步(可选):启动 RocketMQ Console 可视化控制台
# 下载 rocketmq-console 源码或 JAR
# 启动(默认端口 8080)
java -jar rocketmq-console-ng-2.0.1.jar --rocketmq.config.namesrvAddr=localhost:9876
访问 http://localhost:8080,可在控制台中查看 Topic、Consumer Group、消息轨迹等。
第四步:引入 Maven 依赖
<dependency>
<groupId>com.alibaba.cloud</groupId>
<artifactId>spring-cloud-starter-stream-rocketmq</artifactId>
</dependency>
版本由
spring-cloud-alibaba-dependenciesBOM 统一管理,无需手动指定。
完整示例1:订单支付异步通知
场景
飞翔科技(Feixiang Tech)订单服务由小崔负责,通知服务由黄俪负责。
操作前(同步调用):
用户下单 → 订单服务创建订单 → 同步调用通知服务发短信
↓
通知服务挂了
↓
订单服务报错 → 订单创建失败
操作后(消息驱动解耦):
用户下单 → 订单服务创建订单 → 发送消息到 RocketMQ → 订单创建成功
↓
RocketMQ 暂存消息
↓
通知服务恢复后消费 → 发短信 + 推送
项目结构
rocketmq-demo/
├── feixiang-order-service/ # 订单服务(小崔)
│ └── pom.xml
├── feixiang-notify-service/ # 通知服务(黄俪)
│ └── pom.xml
application.yml 配置
订单服务(Producer):
spring:
application:
name: feixiang-order-service
cloud:
stream:
rocketmq:
binder:
name-server: 127.0.0.1:9876
bindings:
order-output:
producer:
group: order-producer-group
sync: true # 同步发送,保证可靠性
bindings:
order-output: # 绑定名称,需与 StreamBridge.send() 一致
destination: order-pay-topic # RocketMQ Topic
content-type: application/json
通知服务(Consumer):
spring:
application:
name: feixiang-notify-service
cloud:
stream:
rocketmq:
binder:
name-server: 127.0.0.1:9876
bindings:
order-input:
consumer:
group: notify-consumer-group # Consumer Group
maxReconsumeTimes: 3 # 消费失败最多重试 3 次
bindings:
order-input:
destination: order-pay-topic
content-type: application/json
group: notify-consumer-group
Producer 代码(小崔 — 订单服务)
@RestController
@RequestMapping("/order")
public class OrderController {
@Autowired
private StreamBridge streamBridge;
@PostMapping("/create")
public Result<String> createOrder(@RequestBody Order order) {
// 1. 本地业务:写入订单数据库
orderService.save(order);
// 2. 构建消息
Message<Order> message = MessageBuilder
.withPayload(order)
.setHeader(MessageHeaders.CONTENT_TYPE, MimeTypeUtils.APPLICATION_JSON)
.setHeader("orderId", order.getId())
.build();
// 3. 发送到 RocketMQ(bindingName 对应配置中的 order-output)
boolean sent = streamBridge.send("order-output", message);
// 4. 无论发送成功与否,订单已创建(最终一致性)
return Result.success("订单创建成功,通知将异步发送");
}
}
Consumer 代码(黄俪 — 通知服务)
@Component
public class NotifyConsumer {
@Bean
public Consumer<Message<Order>> orderInput() {
return message -> {
Order order = message.getPayload();
System.out.println("收到订单支付消息,订单号: " + order.getId());
// 1. 发送短信
smsService.send(order.getPhone(),
"您在飞翔科技的订单 " + order.getId() + " 已支付成功");
// 2. App 推送
pushService.push(order.getUserId(),
"订单支付成功", "订单 " + order.getId() + " 已确认");
};
}
}
方法名
orderInput必须与配置中bindings下的 key(去掉-in-/-out-后缀)对应。order-input映射为函数式 Bean 名称orderInput。
完整示例2:事务消息保证数据一致性
场景
大翔(飞翔科技技术总监)要求:订单创建(订单库)和库存扣减(库存库)必须最终一致,不能用分布式锁,不能影响下单性能。
无事务消息时的问题:
订单服务:
1. INSERT INTO orders (id=1001) → 成功
2. sendMessage("stock-deduct", order) → 发送成功
3. 本地下一条 SQL 失败 → 抛出异常
4. 订单服务事务回滚 → orders 表无 id=1001
库存服务:
5. 收到步骤 2 的消息 → 扣减库存 → 库存被错误扣减!
消息发出了但本地事务回滚 → 消费了不该消费的数据。
有事务消息时:
1. 发送半消息(Consumer 不可见)
2. 执行本地事务:INSERT INTO orders
- 成功 → COMMIT 半消息 → Consumer 可见 → 扣减库存
- 失败 → ROLLBACK 半消息 → 消息被丢弃 → 库存不变
- 超时 → Broker 回查事务状态 → 最终 COMMIT 或 ROLLBACK
application.yml 配置
spring:
cloud:
stream:
rocketmq:
binder:
name-server: 127.0.0.1:9876
bindings:
stock-input:
consumer:
group: stock-consumer-group
bindings:
stock-input:
destination: stock-deduct-topic
content-type: application/json
group: stock-consumer-group
Producer 代码(小崔 — 订单服务)
@RestController
@RequestMapping("/order")
public class OrderController {
@Autowired
private RocketMQTemplate rocketMQTemplate; // 用于事务消息
@Autowired
private OrderMapper orderMapper;
@PostMapping("/create-with-tx")
public Result<String> createOrderWithTx(@RequestBody Order order) {
// 构建 RocketMQ Message
Message<Order> message = MessageBuilder
.withPayload(order)
.build();
// 发送事务消息(重点:走半消息流程)
// 参数:事务生产者组、Topic:Tag、消息体、业务参数(arg)
TransactionSendResult result = rocketMQTemplate.sendMessageInTransaction(
"tx-order-create-group", // 事务生产者组
"stock-deduct-topic:deduct", // Topic:Tag
message,
order.getId() // 透传给事务监听器的 arg
);
if (result.getSendStatus() == SendStatus.SEND_OK) {
return Result.success("订单创建成功,库存扣减消息已投递");
} else {
return Result.error("消息发送失败");
}
}
}
事务监听器(核心)
@RocketMQTransactionListener(txProducerGroup = "tx-order-create-group")
public class OrderTransactionListener implements RocketMQLocalTransactionListener {
@Autowired
private OrderMapper orderMapper;
/**
* 半消息发送成功后,执行本地事务
*/
@Override
public RocketMQLocalTransactionState executeLocalTransaction(
Message msg, Object arg) {
Long orderId = (Long) arg;
try {
// 执行本地事务:写入订单
orderMapper.insert(new OrderDO(orderId, "CREATED"));
System.out.println("订单创建成功,事务消息提交: " + orderId);
return RocketMQLocalTransactionState.COMMIT; // 提交 → Consumer 可见
} catch (Exception e) {
System.out.println("订单创建失败,事务消息回滚: " + orderId);
return RocketMQLocalTransactionState.ROLLBACK; // 回滚 → 消息被丢弃
}
}
/**
* 当 executeLocalTransaction 未返回(超时/宕机)时,
* Broker 主动回调此方法查询事务状态
*/
@Override
public RocketMQLocalTransactionState checkLocalTransaction(Message msg) {
// 从消息中提取业务标识
Long orderId = (Long) msg.getHeaders().get("orderId");
// 查询数据库确认订单是否已创建
OrderDO order = orderMapper.selectById(orderId);
if (order != null && "CREATED".equals(order.getStatus())) {
System.out.println("回查:订单已创建,提交事务消息: " + orderId);
return RocketMQLocalTransactionState.COMMIT;
} else {
System.out.println("回查:订单未创建,回滚事务消息: " + orderId);
return RocketMQLocalTransactionState.ROLLBACK;
}
}
}
Consumer 代码(库存服务)
@Component
public class StockConsumer {
@Autowired
private StockMapper stockMapper;
@Bean
public Consumer<Message<Order>> stockInput() {
return message -> {
Order order = message.getPayload();
System.out.println("扣减库存: productId=" + order.getProductId()
+ ", count=" + order.getCount());
stockMapper.deduct(order.getProductId(), order.getCount());
};
}
}
操作前后对比
| 场景 | 无事务消息 | 有事务消息 |
|---|---|---|
| 本地事务成功 | 消息发送,Consumer 正常消费 | 半消息 COMMIT,Consumer 正常消费 |
| 本地事务失败 | 消息已发出 → Consumer 错误消费 | 半消息 ROLLBACK → 消息被丢弃 |
| 发送半消息后宕机 | — | Broker 回查 → 根据数据库状态决定 COMMIT/ROLLBACK |
| 网络超时 | 消息可能丢失或重复 | 回查机制保证最终一致 |
完整示例3:延时消息
场景
大翔提需求:用户下单后 30 分钟内未支付,自动取消订单并释放库存。
操作前(定时任务轮询):
每 5 分钟跑一次定时任务 → 查询所有 status='UNPAID' 且 create_time < now()-30min
→ 逐条更新 status='CANCELLED' → 释放库存
问题:频繁轮询数据库压力大、精度差(最多 5 分钟延迟)。
操作后(延时消息):
订单创建 → 发送 30 分钟延时消息 → 30 分钟后 Consumer 收到消息
→ 查询订单状态 → 未支付则取消
精准触发,零数据库轮询。
RocketMQ 延时级别
RocketMQ 不支持任意时间的延时,提供 18 个固定延迟级别(messageDelayLevel):
| Level | 延时 | Level | 延时 |
|---|---|---|---|
| 1 | 1s | 10 | 6min |
| 2 | 5s | 11 | 7min |
| 3 | 10s | 12 | 8min |
| 4 | 30s | 13 | 10min |
| 5 | 1min | 14 | 20min |
| 6 | 2min | 15 | 30min |
| 7 | 3min | 16 | 30min |
| 8 | 4min | 17 | 40min |
| 9 | 5min | 18 | 1h |
上文 Level 编号为简化对照,实际使用时直接使用对应 Level 数字即可。30 分钟对应
delayLevel=16。
Producer 代码(小崔 — 订单服务,发送延时消息)
@RestController
@RequestMapping("/order")
public class OrderController {
@Autowired
private StreamBridge streamBridge;
@PostMapping("/create")
public Result<String> createOrder(@RequestBody Order order) {
order.setStatus("UNPAID");
orderService.save(order);
// 构建延时消息
Message<Order> delayMsg = MessageBuilder
.withPayload(order)
.setHeader(MessageHeaders.CONTENT_TYPE, MimeTypeUtils.APPLICATION_JSON)
.setHeader("orderId", order.getId())
// 关键:设置延时级别 = 16(约 30 分钟后投递)
.setHeader("DELAY", 16)
.build();
streamBridge.send("order-delay-output", delayMsg);
return Result.success("订单已创建,请在 30 分钟内完成支付");
}
}
Consumer 代码(订单超时检查)
@Component
public class OrderTimeoutConsumer {
@Autowired
private OrderMapper orderMapper;
@Bean
public Consumer<Message<Order>> orderDelayInput() {
return message -> {
Order order = message.getPayload();
OrderDO orderDO = orderMapper.selectById(order.getId());
if (orderDO != null && "UNPAID".equals(orderDO.getStatus())) {
// 超时未支付 → 取消订单
orderMapper.updateStatus(order.getId(), "CANCELLED");
System.out.println("订单 " + order.getId() + " 超时未支付,已自动取消");
// 释放库存...(可再发一条消息给库存服务)
} else {
System.out.println("订单 " + order.getId() + " 已支付,无需取消");
}
};
}
}
application.yml 配置
spring:
cloud:
stream:
rocketmq:
binder:
name-server: 127.0.0.1:9876
bindings:
order-delay-output:
producer:
group: order-delay-producer-group
order-delay-input:
consumer:
group: order-timeout-consumer-group
bindings:
order-delay-output:
destination: order-delay-topic
content-type: application/json
order-delay-input:
destination: order-delay-topic
content-type: application/json
group: order-timeout-consumer-group
操作前后对比
| 维度 | 定时任务轮询 | RocketMQ 延时消息 |
|---|---|---|
| 实时性 | 取决于轮询间隔(≥ 30s) | 精准触发(延时级别定) |
| 数据库压力 | 每次扫描全表未支付订单 | 无额外查询 |
| 扩展性 | 多实例需分布式锁 | 天然分布式 |
| 复杂度 | 需引入调度框架 | 一行 Header(DELAY) |
完整示例4:顺序消息
场景
黄俪接到大翔的需求:订单状态流转(创建 → 支付 → 发货 → 完成)必须严格按顺序消费。若先收到"已发货"再收到"已支付",业务逻辑会乱。
RocketMQ 保证顺序的原理:
- 发送方:指定
hashKey,同一订单的所有消息通过哈希路由到同一 MessageQueue。 - 消费方:设置
consumeMode=ORDERLY,同一 Queue 单线程顺序消费,前一条 ACK 后才拉取下一条。
Producer 代码(小崔 — 订单服务,发送顺序消息)
@RestController
@RequestMapping("/order")
public class OrderController {
@Autowired
private StreamBridge streamBridge;
@PostMapping("/{orderId}/pay")
public Result<String> payOrder(@PathVariable Long orderId) {
OrderEvent event = new OrderEvent(orderId, "PAID", "用户已支付");
// 关键:用 orderId 作为 hashKey,同一订单消息进入同一 Queue
Message<OrderEvent> message = MessageBuilder
.withPayload(event)
.setHeader(MessageHeaders.CONTENT_TYPE, MimeTypeUtils.APPLICATION_JSON)
.setHeader("orderId", orderId)
// RocketMQ 用 TAGS 做 Tag 过滤
.setHeader("TAGS", "status-change")
// 关键:顺序消息的哈希键
.setHeader("rocketmq_KEYS", orderId.toString())
.build();
streamBridge.send("order-status-output", message);
return Result.success("支付成功");
}
@PostMapping("/{orderId}/ship")
public Result<String> shipOrder(@PathVariable Long orderId) {
OrderEvent event = new OrderEvent(orderId, "SHIPPED", "已发货");
Message<OrderEvent> message = MessageBuilder
.withPayload(event)
.setHeader("orderId", orderId)
.setHeader("TAGS", "status-change")
.setHeader("rocketmq_KEYS", orderId.toString())
.build();
streamBridge.send("order-status-output", message);
return Result.success("发货成功");
}
}
Consumer 代码(顺序消费 + 异常处理)
@Component
public class OrderStatusConsumer {
@Bean
public Consumer<Message<OrderEvent>> orderStatusInput() {
return message -> {
OrderEvent event = message.getPayload();
String currentStatus = event.getStatus();
System.out.println("处理订单状态变更: orderId=" + event.getOrderId()
+ ", status=" + currentStatus);
// 校验状态流转是否合法
OrderDO order = orderMapper.selectById(event.getOrderId());
if (!isValidTransition(order.getStatus(), currentStatus)) {
// 状态不合法 → 跳过该消息并继续消费后续消息
System.out.println("非法状态流转: " + order.getStatus()
+ " → " + currentStatus + ",跳过");
return; // 返回即 ACK,队列继续
}
// 更新订单状态
orderMapper.updateStatus(event.getOrderId(), currentStatus);
};
}
private boolean isValidTransition(String from, String to) {
// 合法流转:CREATED → PAID → SHIPPED → COMPLETED
String validFlow = "CREATED,PAID,SHIPPED,COMPLETED";
if (from == null) return "CREATED".equals(to);
return validFlow.indexOf(from) < validFlow.indexOf(to);
}
}
application.yml 顺序消费配置
spring:
cloud:
stream:
rocketmq:
binder:
name-server: 127.0.0.1:9876
bindings:
order-status-input:
consumer:
group: order-status-consumer-group
# 关键:顺序消费模式
orderly: true
# 消费失败后的处理策略(仅顺序消费生效)
# SUSPEND_CURRENT_QUEUE_A_MOMENT:挂起当前队列,稍后重试
suspendCurrentQueueTimeMillis: 1000
bindings:
order-status-input:
destination: order-status-topic
content-type: application/json
group: order-status-consumer-group
顺序消息消费失败的处理
顺序消息消费失败时,RocketMQ 提供三种处理策略:
| 策略 | 枚举值 | 说明 |
|---|---|---|
| 挂起重试 | SUSPEND_CURRENT_QUEUE_A_MOMENT | 暂停当前队列,等待一段时间后重试。不阻塞其他队列 |
| 跳过 | 不抛异常 / return | 消费方法正常返回即 ACK,跳过该消息继续消费 |
顺序消息不能直接抛异常,否则队列会一直阻塞在该条消息上,后续消息永远无法消费。建议在 Consumer 内部 catch 异常并以日志告警 + 人工介入。
易错场景
NameServer 地址配置错误
错误写法(逗号分隔):
name-server: 127.0.0.1:9876,127.0.0.1:9877
正确写法(分号分隔):
name-server: 127.0.0.1:9876;127.0.0.1:9877
RocketMQ Binder 使用分号 ; 分隔多个 NameServer 地址,与 Nacos 的逗号分隔不同。
Consumer Group 重复启动导致消息重复消费
同一 Consumer Group 下所有实例共享消费进度。集群模式下一条消息只会被组内一个实例消费。但以下场景会触发重复消费:
- 两个不同应用使用了相同的
group但destination(Topic)不完全一致 - Consumer 消费后未返回 ACK(异常抛出),消息被重新投递
- 广播模式下所有实例都会收到全部消息
排查方法:在 RocketMQ Console 中查看 Consumer Group 的实例列表和消费进度。
事务消息回查未实现 checkLocalTransaction
只实现 executeLocalTransaction() 而未重写 checkLocalTransaction() 时:
- 如果
executeLocalTransaction()正常返回 → 没问题 - 如果
executeLocalTransaction()中进程崩溃或超时未返回 → Broker 回查时调用空实现 → 返回 UNKNOWN → 消息一直处于半消息状态,永不投递也不丢弃
checkLocalTransaction()必须实现,且必须幂等(可能被重复调用)。
顺序消息消费抛出异常未正确处理
顺序消息 Consumer 抛出异常后,RocketMQ 默认行为是无限重试,导致当前 Queue 完全阻塞,后续消息无法消费。
// ❌ 错误:异常直接抛出,队列阻塞
@Bean
public Consumer<Message<OrderEvent>> orderStatusInput() {
return message -> {
process(message.getPayload()); // 可能抛异常
};
}
// ✅ 正确:catch 异常,告警后继续
@Bean
public Consumer<Message<OrderEvent>> orderStatusInput() {
return message -> {
try {
process(message.getPayload());
} catch (Exception e) {
log.error("顺序消息处理失败,人工介入: {}", message.getPayload(), e);
// 返回即 ACK,不阻塞队列
}
};
}
延时消息延迟级别理解错误
RocketMQ 的延时消息使用固定 18 个级别,不是任意时间:
delayLevel=1→ 1s,delayLevel=2→ 5s,...,delayLevel=16→ 30min- 无法设置"15 分钟"或"2 小时"等不在表中的延时值
- 可通过修改 Broker 配置
messageDelayLevel调整各级别时间,但会影响所有 Topic
面试考点
事务消息的"半消息"机制
半消息(Half Message)对 Consumer 不可见,只有 COMMIT 后才可见。
面试回答要点:
- Producer 发送半消息到 Broker,Broker 存储但标记为"半消息状态"
- Consumer 拉取消息时,Broker 过滤掉半消息,Consumer 看不到
- Producer 执行本地事务,成功后调用
COMMIT,Broker 将半消息标记为正常消息 - 如果本地事务失败,调用
ROLLBACK,Broker 删除半消息 - 如果 Producer 未返回(宕机),Broker 定时回查
checkLocalTransaction(),最多回查 15 次
追问:什么场景适合事务消息?
跨数据库操作(订单库 + 库存库分离),不要求实时强一致,但要求最终一致,且不能用分布式事务(Seata 等)接管时。
消费重试机制
RocketMQ 消费失败默认重试 16 次,重试间隔递增:
| 重试次数 | 间隔 |
|---|---|
| 第 1 次 | 10s |
| 第 2 次 | 30s |
| 第 3 次 | 1min |
| 第 4 次 | 2min |
| ... | ... |
| 第 16 次 | 约 2h |
16 次后仍然失败 → 消息进入死信队列(DLQ,Dead Letter Queue)。死信队列的 Topic 名称为 %DLQ%原ConsumerGroup名。可在 RocketMQ Console 中查看死信消息并手动重投。
顺序消息的原理
RocketMQ 保证同一 MessageQueue 内消息 FIFO:
- Producer 发送消息时指定
hashKey(如订单 ID),通过哈希算法将同一订单的消息路由到同一 Queue - Consumer 使用 Orderly 消费模式,对每个 Queue 分配一个锁,同一时刻只允许一个线程消费
- 前一条消息 ACK 后,才从 Broker 拉取下一条
- Broker 在 Queue 级别的消息投递也是 FIFO
全局顺序:Topic 下只有一个 Queue,牺牲并行度换取全局有序。
如何保证消息不丢失?
三阶段都有丢消息的风险和对策:
| 阶段 | 风险 | RocketMQ 对策 |
|---|---|---|
| 发送 | Producer 异步发送失败不感知 | 同步发送(sync: true)+ 发送重试 |
| 存储 | Broker 宕机内存消息丢失 | 同步刷盘(flushDiskType=SYNC_FLUSH)+ 同步复制(brokerRole=SYNC_MASTER) |
| 消费 | Consumer 异常未正确处理 | 手动确认 + 消费重试 |
同步刷盘 + 同步复制会显著降低吞吐量,生产环境需权衡。
如何保证消息不重复消费?
RocketMQ 无法 100% 保证 Exactly-Once 投递,必须消费端自行实现幂等:
- 数据库唯一索引:消息中携带业务唯一 ID(如订单号),插入消费记录表,利用唯一约束防重
- Redis SETNX:
SET order:msg:${msgId} 1 NX EX 3600,消费前尝试设置,失败则跳过 - 消息 ID 去重表:维护一张
consumer_record(msg_id, consumer_group, create_time)表,消费前检查msg_id是否已存在
@Bean
public Consumer<Message<Order>> orderInput() {
return message -> {
String msgId = (String) message.getHeaders().get("rocketmq_MESSAGE_ID");
// Redis SETNX 防重
Boolean lock = redisTemplate.opsForValue()
.setIfAbsent("msg:consume:" + msgId, "1", Duration.ofHours(1));
if (Boolean.FALSE.equals(lock)) {
System.out.println("消息 " + msgId + " 已消费,跳过");
return;
}
// 执行业务逻辑
processOrder(message.getPayload());
};
}
参考资源:
- RocketMQ 官方文档:https://rocketmq.apache.org
- Spring Cloud Stream 官方文档:https://spring.io/projects/spring-cloud-stream
- Spring Cloud Alibaba 参考文档:https://sca.aliyun.com