Skip to content
第 14 章 ⏱ 13 分钟阅读

第 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 ​

维度RocketMQKafkaRabbitMQ
吞吐量百万级百万级万级
延迟毫秒级毫秒级微秒级
事务消息原生支持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,降低订阅复杂度

动手练习 ​

  1. 用 Docker 启动 NameServer + Broker + Dashboard
  2. 发普通消息:生产者 orders:orderCreated、消费者订阅 tag
  3. 实现一个事务消息:订单创建成功才发"已下单"消息,失败回滚

下一章:第 15 章:Nacos 配置中心 →

本站基于 VitePress 构建 · 由 StackHub 团队维护