乐途乐途
主页
  • 计算机基础

    • TCP/IP
    • Linux
    • HTTP
  • 数据库

    • SQL
    • MySQL 5.7
  • 编程语言

    • C
    • C++
    • Java SE
    • Python2
    • Python3
  • 数据格式

    • JSON
    • XML
  • 认证与安全

    • JWT
  • 工具

    • Markdown
  • Git

    • GitFlow
  • Quartz

    • Quartz
  • Java

    • Maven 入门
    • Maven 进阶
    • MyBatis
    • Spring
    • Spring MVC
  • Java

    • Spring Boot
    • Spring Cloud
    • Spring Cloud Alibaba
    • Spring Security
    • Spring AI
    • Spring Batch
    • Kafka
    • Java 设计模式
  • 缓存

    • Redis
  • 搜索引擎

    • Elasticsearch
  • 分布式协调

    • ZooKeeper
联系
阿里云
主页
  • 计算机基础

    • TCP/IP
    • Linux
    • HTTP
  • 数据库

    • SQL
    • MySQL 5.7
  • 编程语言

    • C
    • C++
    • Java SE
    • Python2
    • Python3
  • 数据格式

    • JSON
    • XML
  • 认证与安全

    • JWT
  • 工具

    • Markdown
  • Git

    • GitFlow
  • Quartz

    • Quartz
  • Java

    • Maven 入门
    • Maven 进阶
    • MyBatis
    • Spring
    • Spring MVC
  • Java

    • Spring Boot
    • Spring Cloud
    • Spring Cloud Alibaba
    • Spring Security
    • Spring AI
    • Spring Batch
    • Kafka
    • Java 设计模式
  • 缓存

    • Redis
  • 搜索引擎

    • Elasticsearch
  • 分布式协调

    • ZooKeeper
联系
阿里云
  • 学习路径
  • Spring Cloud Alibaba概述与技术选型
  • Nacos注册中心
  • Nacos配置中心
  • Sentinel流量控制
  • Sentinel降级与熔断
  • Seata分布式事务
  • RocketMQ消息驱动
  • Spring AI Alibaba与AI集成
  • Dubbo RPC服务调用
  • Gateway服务网关
  • GraalVM静态编译
  • 其他组件速览
  • Alibaba最佳实践与面试考点

RocketMQ 消息驱动

定位

RocketMQ 解决什么问题

在微服务架构中,服务之间的同步调用(如 HTTP + Feign)形成强依赖链条:订单服务调用通知服务 → 通知服务宕机 → 订单创建失败。RocketMQ 通过消息队列(Message Queue)将同步调用转为异步通信,核心解决三个问题:

  • 异步解耦:生产者(Producer)发消息后立即返回,消费者(Consumer)异步消费。生产者不关心消费者是谁、是否在线。
  • 削峰填谷:瞬时流量高峰(如秒杀)涌入的消息先在 Broker 中暂存,消费者按自身处理能力匀速消费,避免下游被冲垮。
  • 最终一致性:通过事务消息保证本地事务与消息发送的原子性,确保跨服务数据最终一致。

与 RabbitMQ / Kafka 的定位对比

维度RocketMQRabbitMQKafka
设计定位金融级业务消息通用消息代理流数据平台
消息模型Topic + Tag + QueueExchange + Binding + QueueTopic + 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 原生概念

概念英文全称说明
TopicTopic消息主题,消息分类的一级标签。例如 order-topic 存放订单消息
TagTag消息标签,Topic 下的二级分类,Consumer 可按 Tag 过滤。例如 order-topic:create、order-topic:cancel
Consumer GroupConsumer Group消费者组。集群模式(Clustering)下同组消费者共同消费,每条消息只被组内一个实例处理;广播模式(Broadcasting)下每条消息被组内所有实例处理
Message QueueMessage Queue消息队列,Topic 下的物理分区(Partition),数量决定并行消费度
NameServerName 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 的 Topicdestination: 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-dependencies BOM 统一管理,无需手动指定。


完整示例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延时
11s106min
25s117min
310s128min
430s1310min
51min1420min
62min1530min
73min1630min
84min1740min
95min181h

上文 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 后才可见。

面试回答要点:

  1. Producer 发送半消息到 Broker,Broker 存储但标记为"半消息状态"
  2. Consumer 拉取消息时,Broker 过滤掉半消息,Consumer 看不到
  3. Producer 执行本地事务,成功后调用 COMMIT,Broker 将半消息标记为正常消息
  4. 如果本地事务失败,调用 ROLLBACK,Broker 删除半消息
  5. 如果 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:

  1. Producer 发送消息时指定 hashKey(如订单 ID),通过哈希算法将同一订单的消息路由到同一 Queue
  2. Consumer 使用 Orderly 消费模式,对每个 Queue 分配一个锁,同一时刻只允许一个线程消费
  3. 前一条消息 ACK 后,才从 Broker 拉取下一条
  4. Broker 在 Queue 级别的消息投递也是 FIFO

全局顺序:Topic 下只有一个 Queue,牺牲并行度换取全局有序。

如何保证消息不丢失?

三阶段都有丢消息的风险和对策:

阶段风险RocketMQ 对策
发送Producer 异步发送失败不感知同步发送(sync: true)+ 发送重试
存储Broker 宕机内存消息丢失同步刷盘(flushDiskType=SYNC_FLUSH)+ 同步复制(brokerRole=SYNC_MASTER)
消费Consumer 异常未正确处理手动确认 + 消费重试

同步刷盘 + 同步复制会显著降低吞吐量,生产环境需权衡。

如何保证消息不重复消费?

RocketMQ 无法 100% 保证 Exactly-Once 投递,必须消费端自行实现幂等:

  1. 数据库唯一索引:消息中携带业务唯一 ID(如订单号),插入消费记录表,利用唯一约束防重
  2. Redis SETNX:SET order:msg:${msgId} 1 NX EX 3600,消费前尝试设置,失败则跳过
  3. 消息 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
上一页
Seata分布式事务
下一页
Spring AI Alibaba与AI集成