Skip to content
Go back

RocketMQ 消息幂等性:为什么 MessageId 不可靠,以及如何真正去重

RocketMQ 消息幂等性:为什么 MessageId 不可靠

一句话结论(30s)

RocketMQ 消息幂等本质上是「必须由消费端兜底」的问题,因为 RocketMQ 只保证 at-least-once、不保证不重复,而重复又来自生产重试、消费 ACK 丢失、Rebalance 三个来源。关键设计是业务唯一标识 + 去重存储(数据库唯一索引 / Redis SETNX / 业务状态机),而不是用 MessageId——因为生产者重试会产生不同的 MessageId,无法识别「同一条业务数据被发送了两次」。核心权衡:exactly-once 代价极高,RocketMQ 选择了 at-least-once + 消费端幂等这条更务实的路。

核心原理(2min)

主流程:在消息体或消息属性里携带业务唯一 ID(订单号/流水号/traceId),消费端处理前先检查该 ID 是否已处理过,再决定执行还是直接返回成功。关键机制有四种方案:数据库唯一索引最可靠(INSERT 去重记录,与业务共享事务,DuplicateKey 即判定重复);Redis SETNX 最高性能但和业务库不在同一事务(可能 Redis 记了而业务失败);业务状态机最优雅(用 WHERE status='PENDING' 的乐观更新,affected==0 即重复);RocketMQ 5.0 的 Broker 端 messageKey 去重开销大、未大规模生产验证。同时去重记录要定期清理(定时任务或设 TTL,注意消息可能被放在重试队列延迟很久)。

底层深入(5-10min)

消息重复,是设计如此

RocketMQ 的投递语义是 at-least-once(至少一次)。这意味着:每条消息至少被消费一次,但在故障场景下可能被消费多次。

这不是 Bug,是设计的取舍。要实现 exactly-once(精确一次),代价极高——需要两阶段提交、分布式事务、Broker 端去重存储。RocketMQ 选择了更务实的方案:保证 at-least-once,让消费者自己处理幂等。

🤔 思考穿插:为什么 RocketMQ 不做 exactly-once,而坚持 at-least-once?—— 因为精确一次要两阶段提交、分布式事务、Broker 端去重存储,代价极高且吞吐大降;而消息重复大多可通过消费端幂等廉价地解决。所以它把「不重复」这个难事交给业务侧兜底,是更务实的工程取舍。

消息重复的三个来源

来源一:生产者发送重复

生产者                              Broker
  │                                   │
  │────── 发送消息 (msgId=abc) ──────▶│ Broker 收到,写入 CommitLog
  │                                   │
  │   ⏰ 网络闪断,ACK 未返回           │
  │                                   │
  │────── 重试发送 (msgId=def) ──────▶│ Broker 再次收到(同一条业务数据,不同的 MessageId)
  │                                   │
  │◀──── ACK ─────────────────────────│

生产者发送消息后,Broker 返回 ACK 之前网络中断。生产者不知道 Broker 是否已成功写入,只能重试。如果 Broker 实际已写入,这条消息就在 CommitLog 中存在两份——MessageId 不同,但业务内容是重复的。

来源二:消费者应答失败导致重复投递

Consumer                           Broker
  │                                   │
  │◀──── 拉取消息 (offset=100) ───────│
  │                                   │
  │      处理消息......                │
  │                                   │
  │────── ACK (offset=101) ──────────▶│  ❌ 网络超时,ACK 丢失
  │                                   │
  │◀──── 拉取消息 (offset=100) ───────│  Broker 以为 offset 还是 100,重新投递

Consumer 处理完消息,提交 offset 时 ACK 丢失。Broker 认为该消息未消费成功,再次投递。

来源三:Rebalance 导致重复

Consumer-A 持有 Queue-0,已拉取 [msg1, msg2, msg3],正在消费 msg2
⟹ 触发 Rebalance(如 Consumer-B 上线)
Consumer-A 释放 Queue-0,msg2 处理结果丢失
Queue-0 分配给 Consumer-B → 重新拉取 [msg1, msg2, msg3] → msg2 再次被消费

为什么 MessageId 不能用来去重?

很多人第一反应:用 MessageId 去重。但这是有坑的。

RocketMQ 的 MessageId 在每条消息写入 CommitLog 时由 Broker 生成,包含 Broker 地址和 CommitLog offset。 所以:

同一条业务数据,第一次发送:
  MessageId = C0A80101-0000000000000001  (offset=1)

重试发送(网络闪断后):
  MessageId = C0A80101-0000000000000005  (offset=5)

不同的 MessageId,相同的业务内容 → MessageId 无法去重!

消息重试消费时 MessageId 不变(同一条消息重试,CommitLog offset 没变),这是可以用于去重的场景。但生产者重试发送产生的重复,MessageId 不同,无法用 MessageId 去重。

换句话说,MessageId 只能识别”同一条消息被消费了两次”(Consumer 端重复),不能识别”同一条业务数据被发送了两次”(Producer 端重复)。

🤔 思考穿插:为什么说 MessageId 只能识别「同一条消息消费两次」、识别不了「同一条业务数据发送两次」?—— 因为 MessageId = Broker 地址 + CommitLog offset,消费重试时 offset 不变所以 id 不变;但生产者重试发送会写进新的 offset,id 就变了。所以业务去重必须用业务唯一标识,而不是 MessageId。

源码印证:MessageId 到底是怎么生成的

MessageId 的内容结构,用三段真实源码就能看清:它本质上是「Broker 地址 + CommitLog 物理 offset」两个字段拼出来的十六进制串。

第一处:MessageId 组装器(MessageDecoder.java)

public static String createMessageId(final ByteBuffer input, final ByteBuffer addr, final long offset) {
    input.flip();
    int msgIDLength = addr.limit() == 8 ? 16 : 28;
    input.limit(msgIDLength);

    input.put(addr);
    input.putLong(offset);

    return UtilAll.bytes2string(input.array());
}

createMessageId 是 msgId 的组装器:input.put(addr) 先写入 Broker 地址(IPv4 时 4 字节 IP + 4 字节端口共 8 字节,msgId 总长 16 字节;IPv6 时 28 字节),再 putLong(offset) 写入 8 字节的 CommitLog 物理偏移。也就是说,msgId 里除了 Broker 地址,剩下的就是 offset,没有携带任何业务信息。

第二处:Broker 写入 CommitLog 时生成(CommitLog.java)

// PHY OFFSET
long wroteOffset = fileFromOffset + byteBuffer.position();

Supplier<String> msgIdSupplier = () -> {
    int sysflag = msgInner.getSysFlag();
    int msgIdLen = (sysflag & MessageSysFlag.STOREHOSTADDRESS_V6_FLAG) == 0 ? 4 + 4 + 8 : 16 + 4 + 8;
    ByteBuffer msgIdBuffer = ByteBuffer.allocate(msgIdLen);
    MessageExt.socketAddress2ByteBuffer(msgInner.getStoreHost(), msgIdBuffer);
    msgIdBuffer.clear();//because socketAddress2ByteBuffer flip the buffer
    msgIdBuffer.putLong(msgIdLen - 8, wroteOffset);
    return UtilAll.bytes2string(msgIdBuffer.array());
};

Broker 在把消息写进 CommitLog 时,用 wroteOffset = fileFromOffset + byteBuffer.position() 算出物理偏移,再通过 msgIdSupplierstoreHost 地址与 wroteOffset 拼成 msgId。这直接解释了为什么同一份业务数据重试发送会得到不同 msgId——两次写入落到了 CommitLog 的不同 offset 上,msgId 自然随之变化。

第三处:反向解析(MessageDecoder.java)

public static MessageId decodeMessageId(final String msgId) throws UnknownHostException {
    byte[] bytes = UtilAll.string2bytes(msgId);
    ByteBuffer byteBuffer = ByteBuffer.wrap(bytes);

    // address(ip+port)
    byte[] ip = new byte[msgId.length() == 32 ? 4 : 16];
    byteBuffer.get(ip);
    int port = byteBuffer.getInt();
    SocketAddress address = new InetSocketAddress(InetAddress.getByAddress(ip), port);

    // offset
    long offset = byteBuffer.getLong();

    return new MessageId(address, offset);
}

decodeMessageId 是反向解析:从 msgId 里依次读出 IP、端口、offset,还原成 MessageId(address, offset)。offset 每次写都可能不同,而业务内容却完全相同——这就是「MessageId 无法对业务去重」的底层原因,也是必须引入业务唯一标识的根本动机。

真正的幂等方案:业务唯一标识 + 数据库/Redis

核心思路

在消息体(Message Body)中携带一个业务唯一标识(如订单号、交易流水号、请求 traceId),消费端在处理前检查这个标识是否已经处理过。

方案一:数据库唯一索引(最可靠)

CREATE TABLE message_dedup (
    biz_id   VARCHAR(64)  NOT NULL,
    status   TINYINT      DEFAULT 0,  -- 0=处理中, 1=已处理
    created_at DATETIME   DEFAULT CURRENT_TIMESTAMP,
    PRIMARY KEY (biz_id)
);

-- 消费端
INSERT INTO message_dedup (biz_id, status) VALUES ('ORDER-10086', 0);
-- 如果 INSERT 成功 → 第一次处理 → 执行业务逻辑 → UPDATE status=1
-- 如果 INSERT 失败(DuplicateKeyException)→ 消息已处理过 → 直接返回 SUCCESS

幂等流程

@Transactional
public ConsumeConcurrentlyStatus consume(MessageExt msg) {
    String bizId = msg.getProperty("BIZ_ID");  // 从消息属性中取业务ID

    // 1. 尝试插入去重记录
    try {
        jdbcTemplate.update(
            "INSERT INTO message_dedup (biz_id, status) VALUES (?, 0)", bizId);
    } catch (DuplicateKeyException e) {
        // 已经处理过
        return ConsumeConcurrentlyStatus.CONSUME_SUCCESS;
    }

    // 2. 执行业务逻辑
    doBusinessLogic(msg);

    // 3. 标记已处理
    jdbcTemplate.update(
        "UPDATE message_dedup SET status=1 WHERE biz_id=?", bizId);

    return ConsumeConcurrentlyStatus.CONSUME_SUCCESS;
}

优点:与业务数据库共享事务,强一致。缺点:多了一张去重表。

方案二:Redis 原子操作(高性能)

public ConsumeConcurrentlyStatus consume(MessageExt msg) {
    String bizId = msg.getProperty("BIZ_ID");
    String dedupKey = "msg:dedup:" + bizId;

    // SETNX: 只有 key 不存在时才 set,返回 1
    // EX 3600: 设置过期时间 1 小时
    String result = jedis.set(dedupKey, "1", "NX", "EX", 3600);

    if (!"OK".equals(result)) {
        // key 已存在 → 重复消息
        return ConsumeConcurrentlyStatus.CONSUME_SUCCESS;
    }

    // 第一次处理
    doBusinessLogic(msg);
    return ConsumeConcurrentlyStatus.CONSUME_SUCCESS;
}

优点:高性能,不依赖业务数据库。缺点:Redis 和业务数据库不在同一事务中——可能业务处理成功但 Redis 记录设置了、也可能 Redis 设置了但业务处理失败。

方案三:业务状态机(最优雅)

不额外建表,利用业务自身的状态字段:

public void processOrder(OrderMessage msg) {
    // 状态机:PENDING → PAID → SHIPPED → COMPLETED
    // 只有当前状态允许时才更新
    int affected = db.update(
        "UPDATE orders SET status='PAID', paid_at=NOW() " +
        "WHERE order_id=? AND status='PENDING'", msg.getOrderId());

    if (affected == 0) {
        // 订单已经不是 PENDING 状态 → 已经被处理过了
        log.info("Duplicate message ignored: orderId={}", msg.getOrderId());
        return;
    }

    // 后续逻辑...
}

优点:零额外成本,精确实现业务幂等。缺点:依赖业务设计,不是万能方案。

方案四:RocketMQ 5.0 新特性——消息幂等

RocketMQ 5.0 引入了 Broker 端去重能力,使用 messageKey 作为去重键:

Message msg = new Message("OrderTopic", "TagA", body);
msg.setKeys("ORDER-10086");  // 设置去重键
producer.send(msg);

在 Consumer 端,同一个 messageKey 的消息不会重复投递。但这需要 Broker 端存储去重信息,对 Broker 有额外开销,目前还未大规模生产验证。

三阶幂等性保证

阶段机制解决的问题
发送时业务唯一 ID 透传(消息属性/keys)标记每一条业务消息
Broker 端(可选)5.0 幂等机制Broker 端去重
消费时DB 唯一索引 / Redis SETNX / 业务状态机最终防线

核心准则:消费端的幂等性是最后一道防线。不能依赖 Broker、不能依赖 MessageId、不能依赖”消息中间件保证不重复”。必须假设重复一定会发生,并在消费端兜底。

🤔 思考穿插:为什么「消费端幂等」是最后一道防线,而不是 Broker 端去重?—— 因为重复的三个来源(生产重试、ACK 丢失、Rebalance)都发生在 Broker 去重覆盖不到或不可靠的环节,且 5.0 的 Broker 去重尚未大规模生产验证;只有消费端自己校验业务唯一标识,才能覆盖所有重复路径。所以「假设重复一定会发生」是分布式系统的基本素养。

去重记录的清理

去重表/Redis key 不能无限增长。

方案

  1. 定时清理(每天凌晨清理 3 天前的记录)。
  2. 过期时间(Redis 直接设 TTL,如 24 小时。前提:消息不会超过 24 小时后重投——这在 RocketMQ 中不成立!消息可能被放在重试队列中延迟很久。)
  3. 消息时间戳(在去重表中记录消息的 bornTimestamp,只保留 N 天内的去重记录,N 远超最大重试间隔)。

推荐:数据库去重表用定时任务清理(如保留 7 天),Redis 方案设置合理的 TTL(如 72 小时,确保覆盖所有重试场景)。

🤔 思考穿插:Redis 去重 key 设 TTL 为什么有坑?—— 因为 RocketMQ 的重复消息可能被放进重试队列延迟很久才再次到达,若 TTL 比「最大重试间隔」短,第二次到达时 key 已过期,就会漏判重复;所以 TTL 必须覆盖所有重试场景(如 72h),而不是随意设 1 小时。

总结

RocketMQ 消息幂等性是一个”必须由消费端解决”的问题:

  1. RocketMQ 只保证 at-least-once,不保证不重复。
  2. MessageId 不能作为去重依据——生产者重试会产生不同的 MessageId。
  3. 业务唯一标识是关键——在消息体或消息属性中携带订单号、流水号等。
  4. 三种主流去重方案:数据库唯一索引(最可靠)、Redis SETNX(最高性能)、业务状态机(最优雅)。
  5. 去重记录需要定期清理,避免无限增长。

记住:消息队列是可靠的管道,但不是绝对精确的管道。 假设消息一定会重复,这是分布式系统的基本素养。

章末提问

追问 1:RocketMQ 为什么不能保证消息不重复?重复都有哪些来源?

结论先行:因为 RocketMQ 只保证 at-least-once(至少一次),重复来自生产重试、消费 ACK 丢失、Rebalance 三个来源。

因为:生产者在 ACK 返回前网络中断会重试发送,同一条业务数据被写两份;消费者提交 offset 时 ACK 丢失,Broker 以为未消费成功而重新投递;Rebalance 时未完成的消费结果丢失、新消费者重新拉取整批。三种情况都可能在故障窗口内产生重复,这是「可靠投递」和「精确一次」之间的固有矛盾。

追问 2:用 MessageId 做幂等去重为什么不可靠?

结论先行:因为 MessageId 由 Broker 地址 + CommitLog offset 拼成,生产者重试发送会落到不同 offset,得到不同的 MessageId。

因为:MessageId 只能识别「同一条消息被消费两次」(重试消费 offset 不变、id 不变),识别不了「同一条业务数据被发送两次」(重试发送写进新 offset、id 变了)。所以去重必须引入业务唯一标识(订单号/流水号/traceId),MessageId 只在消费重试场景勉强可用。

追问 3:数据库唯一索引和 Redis SETNX 做幂等,各自的优缺点是什么?怎么选?

结论先行:DB 唯一索引最可靠、与业务共享事务,但多一张去重表;Redis SETNX 高性能、不依赖业务库,但和业务不在同一事务、可能不一致。

因为:DB 方案靠 INSERT 触发 DuplicateKeyException 判定重复,能和业务逻辑随同一事务一起提交或回滚;Redis 方案 SETNX 原子去重,但 Redis 记录和业务提交是两次独立操作,可能「Redis 记了但业务失败」或「业务成功但 Redis 没记」。强一致场景选 DB,高吞吐场景选 Redis(并接受 TTL 与一致性的权衡),业务状态机则是零额外成本的最优雅方案。


Share this post on:

Previous Post
RocketMQ 消息积压处理:监控、流控、扩容与 DLQ 的全链路方案
Next Post
RocketMQ死信队列——16次重试后的最终兜底