1. MQ如何保证消息不丢失

1.1 哪些环节可能会丢消息

⾸先分析下MQ的整个消息链路中,有哪些步骤是可能会丢消息的。

在这里插入图片描述
其中,1、2、4 三个场景都是跨⽹络的,⽽跨⽹络就肯定会有丢消息的可能。

然后关于3这个环节,通常MQ存盘时都会先写⼊操作系统的缓存page cache中,然后再由操作系统异步的将消息写⼊硬盘。这个中间有个时间差,就可能会造成消息丢失。如果服务挂了,缓存中还没有来得及写⼊硬盘的消息就会丢失。

1.2 ⽣产者发送消息如何保证不丢失

⽣产者发送消息之所以可能会丢消息,都是因为⽹络。因为⽹络的不稳定性,容易造成请求丢失。怎么解决这样的问题呢?其实⼀个统⼀的思路就是⽣产者确认。简单来说,就是⽣产者发出消息后,给⽣产者⼀个确定的通知,这个消息在Broker端是否写⼊完成了。就好⽐打电话,不确定电话通没通,那就互相说个“喂”,具体确认⼀下。只不过基于这个同样的思路,各个MQ产品有不同的实现⽅式。

1.2.1 ⽣产者发送消息确认机制

在RocketMQ中,提供了三种不同的发送消息的⽅式:
在这里插入图片描述

/**
 * RocketMQ 原生Producer三种发送方式
 * 1. oneway单向发送:只管发送,不等待broker响应,性能最高,存在消息丢失风险
 * 2. sync同步发送:阻塞等待broker返回确认结果,消息可靠性最高,吞吐量较低
 * 3. async异步发送:不阻塞业务线程,通过回调接收结果;兼顾性能与可靠性,消耗客户端线程资源
 */

// 1、单向发送 OneWay
// 特点:发送之后直接返回,不需要 Broker 返回确认;性能极高,但有丢消息风险
// 使用场景:日志收集、统计埋点,允许少量消息丢失的场景
producer.sendOneway(msg);


// 2、同步发送 Sync
// 特点:业务线程阻塞,等待 Broker 返回发送结果;消息可靠性最高,吞吐效率偏低
// send第二个参数为超时时间,单位毫秒,这里设置20秒超时
// 使用场景:重要业务消息,订单、支付,要求确认消息已经投递成功
SendResult sendResult = producer.send(msg, 20 * 1000);


// 3、异步发送 Async
// 特点:业务线程不会阻塞,内部新开线程等待Broker响应,通过SendCallback回调接收成功/失败结果
// 权衡:兼顾可靠性与发送效率;但会占用客户端额外线程,大量消息时注意线程池压力
// 使用场景:接口不能长时间阻塞,同时又要保证消息不丢失的业务场景
producer.send(msg, new SendCallback() {
    /**
     * 消息发送成功回调
     * @param sendResult 发送结果,包含msgId、队列信息等
     */
    @Override
    public void onSuccess(SendResult sendResult) {
        // 发送成功后的业务处理,例如记录msgId日志
    }

    /**
     * 消息发送异常回调
     * @param e 发送失败异常:网络、broker异常、超时等
     */
    @Override
    public void onException(Throwable e) {
        // 发送失败处理,可做重试、告警、落盘记录
    }
});

  1. 禁用 sendOneway 单向发送—— 它不等待 Broker 应答,只适合日志埋点。业务消息用同步发送 send()异步发送 + SendCallback 回调
  2. 开启发送失败重试:设置 retryTimesWhenSendFailed(同步)和 retryTimesWhenSendAsyncFailed(异步)。
  3. 异步必须处理 onException:失败时做重试或本地落盘告警,不能静默吞掉。
  4. 校验 SendResult 状态:只有 SendStatus.SEND_OK 才算发送成功。
  5. 强一致场景用事务消息或本地消息表:保证 “业务数据库操作” 和 “发消息” 的原子性,避免库成功但消息没发出去。

1.2.2. Broker 端:保证消息持久化不丢

  1. 同步刷盘 SYNC_FLUSH:消息真正写入磁盘后才返回成功;默认的异步刷盘 ASYNC_FLUSH 只写 PageCache,机器断电会丢消息。
  2. 同步主从复制 SYNC_MASTER:Master 写完后等待 Slave 同步完成再返回;默认异步复制下 Master 宕机会丢失尚未同步的消息。
  3. 多副本集群部署,避免单点故障。

高可靠组合 = 同步刷盘 + 同步主从复制,代价是 TPS 下降。

1.2.3 消费者端:保证消息被真正处理

  1. 业务执行成功后才提交 offset:RocketMQ 默认在 onMessage 正常返回后才提交消费位点;抛出异常则不提交,消息会重试投递。
  2. 禁止吞异常:catch 住异常只打日志不抛出,框架会误以为消费成功,offset 直接提交,消息就丢了。正确做法是打印日志后继续向上抛出异常
  3. 监控死信队列(DLQ):超过最大重试次数 maxReconsumeTimes 的消息会进入死信队列,不监控等同于消息丢失。
  4. 业务幂等:RocketMQ 是 At-Least-Once(至少一次) 投递,重试会导致重复消费,业务侧必须做幂等。

2. MQ如何保证消息的顺序性

同一组需要有序的消息,发送到同一个 MessageQueue;消费者对这个 MessageQueue 串行消费。

https://blog.csdn.net/Blue_Pepsi_Cola/article/details/163996069?spm=1001.2014.3001.5501
在这里插入图片描述

3. RocketMQ如何保证消息幂等

https://rocketmq.apache.org/zh/docs/bestPractice/01bestpractice
RocketMQ 的投递语义是 At-Least-Once(至少一次),消息可能被重复投递,但RocketMQ 本身不提供 Exactly-Once(恰好一次)语义。幂等性必须由业务消费端自己保证。

重复消费的来源:消费者处理完消息但提交 offset 前宕机、网络超时导致 Broker 认为消费失败、消费异常触发重试等。

幂等的核心思路:利用消息的唯一标识,在执行业务前判断 “这条消息是否已经处理过”,处理过则直接跳过。
在这里插入图片描述

3.1 四种常用幂等方案

3.1.1 方案一:唯一消息 ID + 数据库唯一索引(最常用)

发送消息时携带全局唯一 ID(如订单号、业务流水号),消费端用这个 ID 做去重。

@Override
public void onMessage(MessageExt msg) {
    // 1. 从消息中取出唯一业务ID(推荐用业务自身的唯一键,如订单号)
    String orderNo = msg.getKeys(); // 生产者发送时 setKeys("订单号")

    // 2. 插入去重表,利用唯一索引防重
    //    插入成功 → 第一次消费,执行业务
    //    插入失败(主键冲突) → 重复消费,直接跳过
    try {
        dedupMapper.insert(orderNo); // INSERT INTO dedup_log(biz_id) VALUES(?)
    } catch (DuplicateKeyException e) {
        log.warn("重复消息,跳过消费: {}", orderNo);
        return; // 幂等返回,不执行业务
    }

    // 3. 执行业务逻辑
    doBusiness(orderNo);
}

CREATE TABLE msg_dedup_log (
    id BIGINT PRIMARY KEY AUTO_INCREMENT,
    biz_id VARCHAR(64) NOT NULL COMMENT '业务唯一ID',
    create_time DATETIME DEFAULT CURRENT_TIMESTAMP,
    UNIQUE KEY uk_biz_id (biz_id)
) COMMENT '消息消费去重表';

3.1.2 方案二:Redis 去重(高性能场景)

用 Redis 的 SETNX(不存在才设置)做幂等判断,适合高并发、对性能要求高的场景。

@Override
public void onMessage(MessageExt msg) {
    String bizId = msg.getKeys();
    String dedupKey = "rocketmq:dedup:" + bizId;

    // SETNX:key不存在才设置成功,返回1;已存在返回0
    Boolean isFirst = redisTemplate.opsForValue()
            .setIfAbsent(dedupKey, "1", 24, TimeUnit.HOURS); // 设置过期时间,避免key无限增长

    if (Boolean.FALSE.equals(isFirst)) {
        log.warn("重复消息,跳过: {}", bizId);
        return;
    }

    // 执行业务
    doBusiness(bizId);
}

优点:性能极高;缺点:Redis 宕机或 key 过期可能导致幂等失效,不能用于强一致业务。建议 Redis + DB 双保险。

3.1.3 方案三:状态机幂等(订单 / 流程类业务)

对于有状态流转的业务(如订单:待支付 → 已支付 → 已发货),通过状态判断天然实现幂等。

@Override
public void onMessage(MessageExt msg) {
    String orderNo = msg.getKeys();
    Order order = orderMapper.selectByOrderNo(orderNo);

    // 状态判断:只有"待支付"状态才能更新为"已支付"
    // 已经是"已支付"或更后状态,说明重复消息,直接返回
    if (order.getStatus() != OrderStatus.PENDING_PAY) {
        log.warn("订单状态不匹配,跳过: current={}", order.getStatus());
        return;
    }

    // 带条件更新,影响行数=0说明被并发消费了,也是幂等
    int rows = orderMapper.updateStatusToPaid(orderNo);
    if (rows == 0) {
        log.warn("并发更新,跳过");
        return;
    }

    // 后续业务
    doBusiness(order);
}

UPDATE t_order 
SET status = 'PAID' 
WHERE order_no = #{orderNo} AND status = 'PENDING_PAY';

这是最优雅的幂等方案,不需要额外去重表,利用业务自身状态约束。

3.1.4 方案四:乐观锁 / 版本号

在业务表加 version 字段,更新时携带版本号,版本不匹配则更新失败。

UPDATE t_account 
SET balance = balance - #{amount}, version = version + 1
WHERE account_id = #{accountId} AND version = #{version};

更新影响行数为 0 → 重复消息或并发冲突 → 跳过。

3.1.5 生产者端需要做什么?

幂等虽然主要在消费端,但生产者必须提供唯一标识

// 发送消息时设置 keys,作为幂等判断依据
Message<String> message = MessageBuilder
        .withPayload(msgBody)
        .setHeader(RocketMQHeaders.KEYS, orderNo)  // 业务唯一ID
        .build();
        
rocketMQTemplate.syncSend(topic, message);

推荐用业务自身的唯一键(订单号、流水号),而不是 RocketMQ 的 msgId—— 因为重试时 msgId 可能变化,业务键才是稳定的。

RocketMQ 是至少一次投递,会重复消费,幂等由业务端保证。核心思路是用消息的唯一业务 ID 做去重:常用方案有四种 —— 数据库唯一索引去重表(最通用)、Redis SETNX(高性能)、状态机判断(有状态业务首选)、乐观锁版本号(并发更新)。生产者发送时要把业务唯一 ID 放进 message keys,消费者执行业务前先判断是否已处理,处理过则直接返回。

Logo

openEuler 是由开放原子开源基金会孵化的全场景开源操作系统项目,面向数字基础设施四大核心场景(服务器、云计算、边缘计算、嵌入式),全面支持 ARM、x86、RISC-V、loongArch、PowerPC、SW-64 等多样性计算架构

更多推荐