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() 算出物理偏移,再通过 msgIdSupplier 把 storeHost 地址与 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 不能无限增长。
方案:
- 定时清理(每天凌晨清理 3 天前的记录)。
- 过期时间(Redis 直接设 TTL,如 24 小时。前提:消息不会超过 24 小时后重投——这在 RocketMQ 中不成立!消息可能被放在重试队列中延迟很久。)
- 消息时间戳(在去重表中记录消息的
bornTimestamp,只保留 N 天内的去重记录,N 远超最大重试间隔)。
推荐:数据库去重表用定时任务清理(如保留 7 天),Redis 方案设置合理的 TTL(如 72 小时,确保覆盖所有重试场景)。
🤔 思考穿插:Redis 去重 key 设 TTL 为什么有坑?—— 因为 RocketMQ 的重复消息可能被放进重试队列延迟很久才再次到达,若 TTL 比「最大重试间隔」短,第二次到达时 key 已过期,就会漏判重复;所以 TTL 必须覆盖所有重试场景(如 72h),而不是随意设 1 小时。
总结
RocketMQ 消息幂等性是一个”必须由消费端解决”的问题:
- RocketMQ 只保证 at-least-once,不保证不重复。
- MessageId 不能作为去重依据——生产者重试会产生不同的 MessageId。
- 业务唯一标识是关键——在消息体或消息属性中携带订单号、流水号等。
- 三种主流去重方案:数据库唯一索引(最可靠)、Redis SETNX(最高性能)、业务状态机(最优雅)。
- 去重记录需要定期清理,避免无限增长。
记住:消息队列是可靠的管道,但不是绝对精确的管道。 假设消息一定会重复,这是分布式系统的基本素养。
章末提问
追问 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 与一致性的权衡),业务状态机则是零额外成本的最优雅方案。