RocketMQ 消息可靠性

📖 前置阅读:本文假设读者已掌握 SpringBoot RocketMQ 的发送和消费操作。如果还不熟悉,建议先阅读 SpringBoot RocketMQ 全操作指南

一、⚡ 消息可能丢在哪?

和 RabbitMQ 一样,RocketMQ 的消息丢失也分三个环节:

Producer  →  [网络]  →  Broker  →  [网络]  →  Consumer
   ① 发送丢失      ② Broker 宕机丢失      ③ 消费失败丢失

在 RocketMQ 中,生产者端没有 RabbitMQ 的 Publisher Confirm——替代方案是同步发送 + 返回值判断。Broker 端的持久化取决于刷盘策略主从同步。消费端靠消费状态返回 + 重试 + 死信

二、① 生产者端:同步发送判断返回值

RocketMQ 没有 RabbitMQ 的 ConfirmCallback 异步通知机制,但 syncSend 本身就是同步等待 Broker 确认——返回 SEND_OK 说明消息已写入 CommitLog:

@Service
public class ReliableProducer {

    @Autowired
    private RocketMQTemplate rocketMQTemplate;

    public void sendReliably(OrderMessage msg) {
        // syncSend 是同步阻塞的——返回 SEND_OK 才说明 Broker 已接收
        SendResult result = rocketMQTemplate.syncSend(
            "order-topic:created", msg);

        if (SendStatus.SEND_OK.equals(result.getSendStatus())) {
            log.info("消息已确认到达 Broker: msgId={}", result.getMsgId());
        } else {
            // 只有 SEND_OK 才认为成功——FLUSH_DISK_TIMEOUT 等状态说明刷盘超时
            log.error("消息发送未确认: status={}", result.getSendStatus());
            // 补偿:写入 DB 重试表
            saveToRetryTable(msg);
        }
    }
}

SendStatus 的四种返回值:

状态含义可靠性
SEND_OKBroker 已收到消息并写入 CommitLog
FLUSH_DISK_TIMEOUTBroker 收到但刷盘超时(仅 flushDiskType=SYNC_FLUSH 时可能)中(在内存中,宕机丢失)
FLUSH_SLAVE_TIMEOUTMaster 收到但同步到 Slave 超时(仅 brokerRole=SYNC_MASTER 时可能)中(Slave 无副本)
SLAVE_NOT_AVAILABLEMaster 收到但 Slave 不可用低(无备份)

生产建议:判断 SEND_OK 才认为消息可靠到达。其他状态一律走补偿逻辑——写入重试表或直接告警。

三、② Broker 端:刷盘策略与主从同步

3.1 同步刷盘 vs 异步刷盘

RocketMQ 的消息先写入内存的 CommitLog,然后异步或同步刷到磁盘:

刷盘模式Broker 配置行为性能可靠性
ASYNC_FLUSH(默认)flushDiskType=ASYNC_FLUSH消息写入 OS PageCache 后立即返回 ACK,后台线程定期刷盘(默认 500ms)最高宕机可能丢失 500ms 内的消息
SYNC_FLUSHflushDiskType=SYNC_FLUSH消息写入磁盘后才返回 ACK低(约 1/10)最高——写入磁盘后才确认
# broker.conf——同步刷盘
flushDiskType = SYNC_FLUSH

绝大多数业务用 ASYNC_FLUSH 足够——配合主从同步,Master 宕机后 Slave 上也有消息副本。只有金融交易、支付确认这类"一条都不能丢"的场景才开 SYNC_FLUSH。

3.2 主从同步——刷盘只保证单机不丢,主从保证节点宕机不丢

CommitLog 刷盘只保证写入该 Broker 的磁盘。如果整台机器宕机(磁盘坏了、电源炸了),需要另一个节点有副本

主从模式Broker 配置行为可靠性
ASYNC_MASTERbrokerRole=ASYNC_MASTERMaster 写入后立即返回 ACK,后台异步复制到 Slave宕机可能丢失少量未同步的消息
SYNC_MASTERbrokerRole=SYNC_MASTERMaster 等待 Slave 确认收到后才返回 ACK最高——Master 宕机 Slave 有完整副本
# broker.conf——同步主从
brokerRole = SYNC_MASTER

生产建议:关键业务的 Master 配 SYNC_MASTER + 至少一台 Slave。非关键业务(日志)用 ASYNC_MASTERSYNC_FLUSHSYNC_MASTER 通常不同时开——前者拖慢单机吞吐,后者保证跨节点冗余。

3.3 可靠性配置组合

flowchart LR
classDef startEnd fill:#701a4c,stroke:#e11d48,stroke-width:2px,color:#fce7f3,font-weight:bold;
classDef condition fill:#2a1147,stroke:#a855f7,stroke-width:1.5px,color:#ede9fe,font-weight:bold;
classDef process fill:#1e1e24,stroke:#6b7280,stroke-width:1.5px,color:#e5e7eb;
classDef data fill:#052e16,stroke:#16a34a,stroke-width:1.5px,color:#bbf7d0,font-weight:bold;
classDef highlight fill:#450a0a,stroke:#dc2626,stroke-width:1.5px,color:#fecaca,font-weight:bold;

    START([业务消息类型]) --> Q1{消息丢一条\n会怎样?}
    Q1 -- "严重(金融/支付)" --> SYNC[SYNC_FLUSH + SYNC_MASTER\n每条消息同步刷盘 + 同步复制\n吞吐量约 1000 msg/s]
    Q1 -- "一般(订单/通知)" --> ASYNC[ASYNC_FLUSH + SYNC_MASTER\n异步刷盘 + 同步复制\n吞吐量约 10000 msg/s]
    Q1 -- "不重要(日志/埋点)" --> SIMPLE[ASYNC_FLUSH + ASYNC_MASTER\n单向发送 + 异步刷盘 + 异步复制\n吞吐量最高]

    class START startEnd;
    class Q1 condition;
    class SYNC,ASYNC,SIMPLE highlight;

四、③ 消费者端:重试与死信队列

4.1 消费重试机制

消费者返回 RECONSUME_LATER 或抛异常时,RocketMQ 自动将消息送入重试队列——不需要手动调用任何 NACK 方法,这和 RabbitMQ 的 basicNack(requeue=true) 完全不同。

@Component
@RocketMQMessageListener(
    topic = "order-topic",
    consumerGroup = "order-consumer-group"
)
public class OrderRetryListener
        implements RocketMQListener<OrderMessage> {

    @Override
    public void onMessage(OrderMessage msg) {
        try {
            processOrder(msg);
            // 正常返回 → 消费成功 → offset 推进
        } catch (RetryableException e) {
            // 可重试异常 → 抛出去 → RocketMQ 自动重试
            throw new RuntimeException("临时失败,重试", e);
        } catch (NonRetryableException e) {
            // 不可重试异常(数据格式错误等)→ 记录日志并吞掉
            log.error("消息格式错误,跳过: orderId={}", msg.getOrderId(), e);
            // 正常返回 → 消费成功 → 这条坏消息被消费掉(不再重试)
        }
    }
}

重试间隔是递增的

重试次数     等待时间
第 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   ← 默认最大重试 16 次

⚠️ 新手提示:RocketMQ 的重试消息实际发送到了一个特殊的 Topic——%RETRY%{consumerGroup}。Broker 在这个 Topic 上设置了延迟级别(对应上表的等待时间)。消费者无感知——重试消息和正常消息一样进入 onMessage

4.2 死信队列 —— 重试 16 次后的归宿

重试 16 次仍然失败时,消息不再投递——进入死信 Topic%DLQ%{consumerGroup}

不需要手动配置死信队列——RocketMQ 自动创建。但需要编写消费者监听死信 Topic 来发现问题:

@Component
@RocketMQMessageListener(
    topic = "%DLQ%order-consumer-group",   // ← 死信 Topic
    consumerGroup = "order-dlq-consumer-group"
)
public class OrderDLQListener
        implements RocketMQListener<MessageExt> {

    @Override
    public void onMessage(MessageExt msg) {
        // MessageExt 包含原始消息的全部信息
        log.error("死信消息: msgId={}, topic={}, tag={}, body={}, reconsumeTimes={}",
                msg.getMsgId(),
                msg.getTopic(),      // 原始 Topic
                msg.getTags(),       // 原始 Tag
                new String(msg.getBody()),
                msg.getReconsumeTimes()  // 重试次数
        );
        // 发告警、记录 DB、通知人工处理...
    }
}

4.3 与 RabbitMQ 的可靠性机制对比

机制RabbitMQRocketMQ
生产者确认Publisher Confirm(NACK → 重新发送)同步发送后判断 SendStatus
Broker 持久化消息 + Exchange + Queue 三者持久化CommitLog 顺序写 + 刷盘策略
消费确认手动 basicAck / basicNack返回 CONSUME_SUCCESS 或抛异常
重试机制手动 basicNack(requeue=true)自动进 %RETRY% Topic,16 次递增间隔
死信手动配置 DLX + DLQ自动进入 %DLQ% Topic
跨节点冗余镜像队列 / 仲裁队列主从同步(SYNC_MASTER / ASYNC_MASTER)

五、消息幂等 —— 重复消费的防线

5.1 RocketMQ 什么情况下会重复

场景原因
Producer 超时重发syncSend 超时但 Broker 实际已写入——Producer 重发同一条
Consumer Rebalance消费者实例增减时,Queue 重新分配——正在处理的 offset 可能回退
主从切换Master 宕机 → Slave 提升为新 Master,offset 可能回退

RocketMQ 不保证 exactly-once——消费者端必须自己做幂等。

5.2 幂等实现

核心思路和 RabbitMQ 一样——用消息的唯一 Key 去重:

@Component
@RocketMQMessageListener(
    topic = "order-topic",
    consumerGroup = "order-consumer-group"
)
public class IdempotentOrderListener
        implements RocketMQListener<MessageExt> {  // 注意:泛型是 MessageExt 以获取 msgId

    @Autowired
    private RedisTemplate<String, String> redisTemplate;

    @Override
    public void onMessage(MessageExt msg) {
        // RocketMQ 每条消息有全局唯一的 msgId
        String msgId = msg.getMsgId();
        String orderId = msg.getKeys();    // 发送时 setKeys 设置的业务 Key
        String idempotentKey = orderId != null ? orderId : msgId;

        // SETNX 原子判重
        Boolean firstTime = redisTemplate.opsForValue()
                .setIfAbsent("rocketmq:consumed:" + idempotentKey,
                        "1", Duration.ofHours(24));

        if (Boolean.FALSE.equals(firstTime)) {
            log.warn("重复消息,跳过: msgId={}, keys={}", msgId, orderId);
            return;  // 正常返回 → offset 推进
        }

        try {
            OrderMessage orderMsg = JSON.parseObject(
                    new String(msg.getBody()), OrderMessage.class);
            processOrder(orderMsg);
        } catch (Exception e) {
            // 处理失败 → 删除幂等标记,让重试时可以重新处理
            redisTemplate.delete("rocketmq:consumed:" + idempotentKey);
            throw new RuntimeException("消费失败,回滚幂等标记", e);
        }
    }
}

建议用业务 Key 而非 msgId

// 发送时设置业务 Key
Message<String> message = MessageBuilder
        .withPayload(JSON.toJSONString(orderMsg))
        .setHeader(MessageConst.PROPERTY_KEYS, "order:10001:created")
        .build();
rocketMQTemplate.syncSend("order-topic:created", message);

因为 Producer 超时重发时,两条消息的 msgId 不同但业务相同——用 msgId 去重就无效了。用业务 Key(order:10001:created)去重更可靠。

六、🎯 三个防线总结

生产者端                          Broker端                          消费者端
① 同步发送                ② 刷盘 + 主从同步                   ③ 重试 + 死信 + 幂等
syncSend                    ASYNC/SYNC_FLUSH                   抛异常自动重试
判断 SEND_OK                + ASYNC/SYNC_MASTER                16次后进%DLQ%
失败写重试表                按业务重要性选配置                    SETNX去重
防线机制配置
① 生产者同步发送 + 判断 SendStatussyncSend,判断 SEND_OK
② BrokerCommitLog 刷盘 + 主从同步flushDiskType + brokerRole
③ 消费者抛异常自动重试 + 死信 Topic + 幂等去重16 次递增重试,%DLQ% 自动创建,业务 Key 去重

📖 下一步阅读:消费端的可靠性搞定了,但"集群消费还是广播消费"、“消息什么时候推什么时候拉”、“100 种 Tag 怎么过滤"这些问题还没讲。继续阅读 消费者模式与过滤器,一篇覆盖集群/广播、Push/PULL、Tag/SQL 过滤和 Rebalance 机制。