第 14 章:RocketMQ
学习目标
- 理解 RocketMQ 的架构和特色
- 掌握 Producer / Consumer Java 代码
- 用好顺序消息、事务消息
- 了解 NameServer / Broker / Topic 关系
一、RocketMQ 是什么
RocketMQ 是阿里开源的分布式消息中间件,设计目标:高吞吐、低延迟、海量消息堆积。
适用场景:
- 电商交易(订单、支付、物流)
- 金融级事务消息
- 大规模日志收集
二、核心概念
| 概念 | 含义 |
|---|---|
| NameServer | 轻量级注册中心,Broker 注册、Producer/Consumer 发现 |
| Broker | 消息存储节点,主从架构 |
| Topic | 消息主题 |
| Queue | 消息队列,Topic 下分多个队列(类似 Kafka 分区) |
| Producer Group | 生产者组 |
| Consumer Group | 消费者组 |
| Tag | 消息子分类,用于过滤 |
三、RocketMQ vs Kafka vs RabbitMQ
| 维度 | RocketMQ | Kafka | RabbitMQ |
|---|---|---|---|
| 吞吐量 | 百万级 | 百万级 | 万级 |
| 延迟 | 毫秒级 | 毫秒级 | 微秒级 |
| 事务消息 | 原生支持 | 0.11+ 支持 | 不支持 |
| 顺序消息 | 支持 | 分区内支持 | 弱 |
| 消息堆积 | 强(CommitLog) | 强 | 弱 |
| 协议 | 自定义 | 自定义 | AMQP |
四、安装
bash
docker run -d --name rmqnamesrv -p 9876:9876 \
apache/rocketmq:5.1.0 sh mqnamesrv
docker run -d --name rmqbroker \
-p 10911:10911 -p 10909:10909 \
-e NAMESRV_ADDR=rmqnamesrv:9876 \
apache/rocketmq:5.1.0 sh mqbroker管理控制台:
bash
docker run -d --name rmqadmin \
-p 8080:8080 \
-e NAMESRV_ADDR=rmqnamesrv:9876 \
apacherocketmq/rocketmq-dashboard:latest访问 http://localhost:8080。
⚠️ 坑 1:Broker 默认配置自动创建 topic,生产环境必须关闭,避免误创建。
五、Spring Boot 整合
依赖
xml
<dependency>
<groupId>org.apache.rocketmq</groupId>
<artifactId>rocketmq-spring-boot-starter</artifactId>
<version>2.3.0</version>
</dependency>配置
yaml
rocketmq:
name-server: localhost:9876
producer:
group: order-producer
send-message-timeout: 10000
retry-times-when-send-failed: 3六、普通消息
java
@Service
public class OrderProducer {
@Resource
private RocketMQTemplate rocketMQTemplate;
public void send(OrderEvent event) {
// topic + tag
Message<OrderEvent> msg = MessageBuilder.withPayload(event).build();
rocketMQTemplate.send("orders:orderCreated", msg);
}
}java
@Component
public class OrderConsumer implements RocketMQListener<OrderEvent> {
@Override
public void onMessage(OrderEvent event) {
log.info("收到:{}", event);
// 业务处理
}
}java
@RocketMQMessageListener(
topic = "orders",
selectorExpression = "orderCreated", // 订阅 tag
consumerGroup = "order-consumer"
)
public class OrderConsumer implements RocketMQListener<OrderEvent> { ... }⚠️ 坑 2:
topic和selectorExpression必须对应生产端:send("topic:tag", msg)↔topic + tag订阅。
七、顺序消息
业务:同一订单的"创建 → 支付 → 完成"必须按顺序处理。
java
// 生产端:同一个 orderId 选同一个队列
rocketMQTemplate.syncSendOrderly(
"orders",
msg,
String.valueOf(event.getOrderId()) // hashKey
);java
// 消费端:MessageListenerOrderly(不是并发监听)
rocketMQTemplate.setMessageQueueSelector();java
@RocketMQMessageListener(
topic = "orders",
consumerGroup = "order-consumer",
consumeMode = ConsumeMode.ORDERLY // 单线程顺序
)⚠️ 坑 3:顺序消息不能并行消费,TPS 受限 —— 只对必须有序的部分用顺序消息,其他用普通消息。
八、事务消息
场景:下单 → 创建订单 → 扣库存 → 发消息,要求本地事务和消息发送要么都成功要么都失败。
RocketMQ 通过两阶段提交 + 反查机制实现:
java
@Resource
private RocketMQTemplate rocketMQTemplate;
public void createOrder(Order order) {
// 1. 发送 half 消息(对消费者不可见)
Message msg = MessageBuilder.withPayload(order).build();
rocketMQTemplate.sendMessageInTransaction(
"order-tx",
"orders:orderCreated",
msg,
order
);
}
// 2. 本地事务监听器
@RocketMQTransactionListener
class OrderTxListener implements RocketMQLocalTransactionListener {
@Override
public RocketMQLocalTransactionState executeLocalTransaction(Message msg, Object arg) {
try {
// 真正的事务:DB 写订单 + 扣库存
orderService.createOrder((Order) arg);
return RocketMQLocalTransactionState.COMMIT; // 提交
} catch (Exception e) {
return RocketMQLocalTransactionState.ROLLBACK; // 回滚
}
}
@Override
public RocketMQLocalTransactionState checkLocalTransaction(Message msg) {
// 反查:Broker 异步询问本地事务状态
return orderService.isOrderExists(extractId(msg))
? RocketMQLocalTransactionState.COMMIT
: RocketMQLocalTransactionState.ROLLBACK;
}
}⚠️ 坑 4:
checkLocalTransaction必须能查到事务最终状态(查 DB 标志位),不能假设"还在内存中"。
九、消息可靠性
| 阶段 | 策略 |
|---|---|
| 生产 | syncSend + 失败重试 + 落库补偿 |
| 存储 | Broker 主从同步 + 刷盘策略 |
| 消费 | 至少一次 + 业务幂等 |
十、本章小结
| 要点 | 关键 |
|---|---|
| 架构 | NameServer + Broker + Producer + Consumer |
| 普通消息 | send("topic:tag", msg) |
| 顺序消息 | syncSendOrderly + consumeMode = ORDERLY |
| 事务消息 | half 消息 + 本地事务 + 反查 |
| 可靠性 | syncSend + 落库补偿 + 业务幂等 |
| Tag | 替代多 topic,降低订阅复杂度 |
动手练习
- 用 Docker 启动 NameServer + Broker + Dashboard
- 发普通消息:生产者
orders:orderCreated、消费者订阅 tag - 实现一个事务消息:订单创建成功才发"已下单"消息,失败回滚
下一章:第 15 章:Nacos 配置中心 →