一、一个真实的生产事故
某电商平台上线了“下单送积分”功能:用户支付成功后,发送 MQ 消息,积分服务收到消息后给用户加积分。
上线第三天,客服接到投诉:“我买了1单,怎么加了2次积分?”
排查发现,MQ 消息被重复消费了。原因有两个:
原因一:生产者重试导致重复投递
订单服务发送消息 → 网络抖动,Broker未确认
↓
订单服务重试 → 消息被投递两次
原因二:消费者 ACK 丢失导致重投
消费者处理成功 → 调用ACK → 网络抖动,ACK丢失
↓
Broker认为消息未处理 → 重新投递
↓
消费者再次处理 → 重复加积分
事故复盘:
| 问题 |
说明 |
| 消息重复不可避免 |
MQ 只保证 At Least Once,不保证 Exactly Once |
| 无幂等设计 |
消费逻辑没有做去重,重复消息被重复执行 |
| 无唯一约束 |
数据库没有唯一索引兜底 |
| 无监控 |
重复消费发生了3天才被发现 |
二、问题根源:为什么消息一定会重复?
2.1 MQ的三种投递语义
| 语义 |
含义 |
实现难度 |
适用场景 |
| At Most Once |
最多一次,可能丢 |
简单 |
日志采集等允许丢失的场景 |
| At Least Once |
至少一次,可能重复 |
中等 |
绝大多数业务场景 |
| Exactly Once |
恰好一次 |
极难 |
金融级场景(通常通过幂等实现) |
主流 MQ(RabbitMQ、RocketMQ、Kafka)都只能保证 At Least Once。
2.2 消息重复的三个来源
来源一:生产者重试
生产者发送消息 → 未收到Broker确认 → 重试发送
→ 消息可能已经在Broker中,重试导致重复
来源二:消费者 ACK 丢失
消费者处理成功 → 发送ACK → 网络抖动,ACK丢失
→ Broker认为消息未处理 → 重新投递
→ 消费者再次收到消息
来源三:消费者处理超时
消费者拉取消息 → 处理时间超过超时时间
→ Broker认为消费者已死 → 重新投递消息给其他消费者
→ 消息被两个消费者同时处理
2.3 幂等性 vs 去重的区别
| 概念 |
说明 |
| 去重 |
通过某种方式识别重复消息,丢弃重复的 |
| 幂等 |
保证同一操作执行多次的结果与执行一次相同 |
去重是手段,幂等是目标。但很多场景下,去重就等于幂等。
2.4 幂等消费的核心原则
不要试图消除重复,而要容忍重复。
既然重复不可避免,正确的做法是让消费逻辑本身具有幂等性。
三、核心技术原理
3.1 幂等消费的四种方案
| 方案 |
实现方式 |
优点 |
缺点 |
适用场景 |
| Redis SET NX |
消息ID写入 Redis,重复则跳过 |
性能高 |
依赖 Redis |
高并发场景 |
| 数据库唯一索引 |
业务唯一ID插入唯一表 |
强一致 |
性能略低 |
强一致场景 |
| 状态机 |
检查业务状态,已处理则跳过 |
无额外存储 |
需要状态设计 |
状态流转场景 |
| 乐观锁 |
版本号控制 |
无锁 |
需要版本字段 |
更新场景 |
3.2 Redis SET NX 的原理
// 核心操作:SET message:id:{msgId} 1 NX EX 86400
Boolean success = redisTemplate.opsForValue()
.setIfAbsent(key, "1", 24, TimeUnit.HOURS);
if (Boolean.FALSE.equals(success)) {
// key已存在,说明消息已被处理,跳过
return;
}
// key不存在,首次处理
三个关键点:
- 原子性:
SET NX 是原子操作,并发安全
- 过期时间:设置 TTL,防止 Redis 内存无限增长
- 唯一Key:通常用
消息ID 或 业务ID + 操作类型
3.3 数据库唯一索引的原理
CREATE TABLE message_consume_log (
id BIGINT AUTO_INCREMENT,
message_id VARCHAR(64) NOT NULL,
consume_time DATETIME,
PRIMARY KEY (id),
UNIQUE KEY uk_message_id (message_id) -- 唯一索引
);
try {
// 插入消费日志,如果messageId已存在则抛异常
jdbcTemplate.update(
"INSERT INTO message_consume_log (message_id, consume_time) VALUES (?, NOW())",
messageId);
// 插入成功,执行业务逻辑
processBusiness();
} catch (DuplicateKeyException e) {
// 唯一索引冲突,说明已消费过
log.warn("消息已消费: messageId={}", messageId);
}
3.4 组合方案:Redis + 数据库
单一方案的问题:
| 方案 |
问题 |
| 只用 Redis |
Redis 故障后无法去重 |
| 只用数据库 |
高并发下数据库压力大 |
组合方案的优势:
第一层:Redis SET NX → 快速拦截99%的重复消息
第二层:数据库唯一索引 → 兜底,即使Redis故障也能保证幂等
3.5 状态机的幂等设计
很多业务场景天然支持幂等——通过状态检查:
待支付 → 已支付 → 已发货 → 已完成
消费"支付成功"消息时:
当前状态 = 待支付 → 更新为已支付 ✅
当前状态 = 已支付 → 跳过(已处理)✅
当前状态 = 已完成 → 跳过 ✅
核心SQL:
UPDATE orders
SET status = 1, pay_time = NOW()
WHERE id = ? AND status = 0; -- 关键:WHERE条件限制状态
如果订单已经是 status=1,则 WHERE status=0 不匹配,更新影响行数为0,天然幂等。
四、实现方案及场景
4.1 本文采用的方案
三层幂等防护:
| 层次 |
实现 |
作用 |
| 第一层 |
Redis SET NX(消息ID去重) |
快速拦截重复消息 |
| 第二层 |
业务状态检查(条件更新) |
业务层面幂等 |
| 第三层 |
数据库唯一索引(消费日志) |
最终兜底 |
4.2 整体流程
消费消息
↓
① Redis SET NX message:id:{msgId}
├── 已存在 → 跳过
└── 不存在 → 继续
↓
② 检查业务状态
├── 已处理 → 跳过
└── 未处理 → 继续
↓
③ 执行业务(条件更新 WHERE status = 预期值)
↓
④ 插入消费日志(唯一索引兜底)
↓
⑤ ACK
4.3 适用场景
| 场景 |
说明 |
| 订单支付通知 |
支付成功后更新订单状态 |
| 库存扣减 |
扣减库存需要幂等 |
| 积分发放 |
发放积分需要防重复 |
| 优惠券发放 |
优惠券只能发放一次 |
五、完整实现
5.1 项目结构
heyou-mq-idempotent/
├── pom.xml
└── src/main/java/com/heyou/idempotent/
├── IdempotentApplication.java
├── config/
│ └── RabbitMQConfig.java
├── consumer/
│ ├── OrderPayConsumer.java # 支付消息消费者
│ └── PointsConsumer.java # 积分消息消费者
├── service/
│ └── IdempotentService.java # 幂等服务
├── mapper/
│ └── OrderMapper.java
└── entity/
├── Order.java
└── MessageConsumeLog.java
5.2 Maven依赖
<dependencies>
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-web</artifactId>
</dependency>
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-amqp</artifactId>
</dependency>
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-data-redis</artifactId>
</dependency>
<dependency>
<groupId>org.apache.commons</groupId>
<artifactId>commons-pool2</artifactId>
</dependency>
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-jdbc</artifactId>
</dependency>
<dependency>
<groupId>com.mysql</groupId>
<artifactId>mysql-connector-j</artifactId>
<scope>runtime</scope>
</dependency>
<dependency>
<groupId>org.projectlombok</groupId>
<artifactId>lombok</artifactId>
<optional>true</optional>
</dependency>
</dependencies>
5.3 数据库表
-- 订单表
CREATE TABLE `orders` (
`id` BIGINT NOT NULL AUTO_INCREMENT,
`order_no` VARCHAR(64) NOT NULL,
`user_id` BIGINT NOT NULL,
`amount` DECIMAL(10,2) NOT NULL,
`status` TINYINT DEFAULT 0 COMMENT '0-待支付 1-已支付 2-已取消',
`pay_time` DATETIME DEFAULT NULL,
`create_time` DATETIME DEFAULT CURRENT_TIMESTAMP,
PRIMARY KEY (`id`),
UNIQUE KEY `uk_order_no` (`order_no`)
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4;
-- 消息消费日志(唯一索引兜底)
CREATE TABLE `message_consume_log` (
`id` BIGINT NOT NULL AUTO_INCREMENT,
`message_id` VARCHAR(64) NOT NULL,
`biz_type` VARCHAR(50) NOT NULL,
`consume_time` DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP,
PRIMARY KEY (`id`),
UNIQUE KEY `uk_message_id` (`message_id`)
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4;
-- 积分表(用于验证幂等)
CREATE TABLE `user_points` (
`id` BIGINT NOT NULL AUTO_INCREMENT,
`user_id` BIGINT NOT NULL,
`points` INT NOT NULL DEFAULT 0,
`update_time` DATETIME DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP,
PRIMARY KEY (`id`),
UNIQUE KEY `uk_user_id` (`user_id`)
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4;
5.4 幂等服务(核心)
package com.example.idempotent.service;
import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j;
import org.springframework.dao.DuplicateKeyException;
import org.springframework.data.redis.core.StringRedisTemplate;
import org.springframework.jdbc.core.JdbcTemplate;
import org.springframework.stereotype.Service;
import java.util.concurrent.TimeUnit;
@Slf4j
@Service
@RequiredArgsConstructor
public class IdempotentService {
private final StringRedisTemplate redisTemplate;
private final JdbcTemplate jdbcTemplate;
private static final String IDEMPOTENT_PREFIX = "mq:idempotent:";
private static final long EXPIRE_HOURS = 24;
/**
* 幂等性检查(第一层:Redis SET NX)
*
* @param messageId 消息唯一ID
* @return true-首次处理,false-重复消息
*/
public boolean checkRedisIdempotent(String messageId) {
String key = IDEMPOTENT_PREFIX + messageId;
Boolean success = redisTemplate.opsForValue()
.setIfAbsent(key, "1", EXPIRE_HOURS, TimeUnit.HOURS);
if (Boolean.FALSE.equals(success)) {
log.warn("消息已被消费(Redis拦截): messageId={}", messageId);
return false;
}
return true;
}
/**
* 记录消费日志(第三层:数据库唯一索引)
*
* @param messageId 消息唯一ID
* @param bizType 业务类型
* @return true-首次消费,false-重复消费
*/
public boolean recordConsumeLog(String messageId, String bizType) {
try {
jdbcTemplate.update(
"INSERT INTO message_consume_log (message_id, biz_type) VALUES (?, ?)",
messageId, bizType);
return true;
} catch (DuplicateKeyException e) {
log.warn("消息已被消费(数据库兜底): messageId={}", messageId);
return false;
}
}
/**
* 综合幂等检查(三层防护)
*/
public boolean checkIdempotent(String messageId, String bizType) {
// 第一层:Redis快速拦截
if (!checkRedisIdempotent(messageId)) {
return false;
}
// 第三层:数据库兜底
if (!recordConsumeLog(messageId, bizType)) {
// Redis通过了但数据库没通过,说明Redis数据过期
// 此时应该阻止处理,防止重复
return false;
}
return true;
}
}
5.5 支付消息消费者
package com.example.idempotent.consumer;
import com.example.idempotent.service.IdempotentService;
import com.rabbitmq.client.Channel;
import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j;
import org.springframework.amqp.core.Message;
import org.springframework.amqp.rabbit.annotation.RabbitListener;
import org.springframework.jdbc.core.JdbcTemplate;
import org.springframework.stereotype.Component;
import java.io.IOException;
@Slf4j
@Component
@RequiredArgsConstructor
public class OrderPayConsumer {
private final IdempotentService idempotentService;
private final JdbcTemplate jdbcTemplate;
/**
* 消费支付成功消息
*
* 三层幂等防护:
* 1. Redis SET NX(消息ID去重)
* 2. 业务状态检查(条件更新)
* 3. 数据库唯一索引(消费日志兜底)
*/
@RabbitListener(queues = "order.pay.queue")
public void onPaySuccess(Message message, Channel channel) throws IOException {
String messageId = message.getMessageProperties().getMessageId();
String content = new String(message.getBody());
long deliveryTag = message.getMessageProperties().getDeliveryTag();
log.info("收到支付消息: messageId={}, content={}", messageId, content);
try {
// 第一层 + 第三层:综合幂等检查
if (!idempotentService.checkIdempotent(messageId, "ORDER_PAY")) {
channel.basicAck(deliveryTag, false);
return;
}
// 解析消息
Long orderId = Long.parseLong(content.split(":")[0]);
// 第二层:业务状态检查(条件更新)
// 关键:WHERE status = 0 保证只有待支付订单才能被更新
int rows = jdbcTemplate.update(
"UPDATE orders SET status = 1, pay_time = NOW() " +
"WHERE id = ? AND status = 0",
orderId);
if (rows > 0) {
log.info("订单支付成功: orderId={}", orderId);
// 执行后续逻辑(发短信、发积分等)
} else {
// 影响行数为0,说明订单状态不是待支付
log.warn("订单状态不允许支付(可能已处理): orderId={}", orderId);
}
// 手动ACK
channel.basicAck(deliveryTag, false);
} catch (Exception e) {
log.error("支付消息处理失败: messageId={}", messageId, e);
// 处理失败,拒绝并重新入队
channel.basicNack(deliveryTag, false, true);
}
}
}
5.6 积分消息消费者
package com.example.idempotent.consumer;
import com.example.idempotent.service.IdempotentService;
import com.rabbitmq.client.Channel;
import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j;
import org.springframework.amqp.core.Message;
import org.springframework.amqp.rabbit.annotation.RabbitListener;
import org.springframework.jdbc.core.JdbcTemplate;
import org.springframework.stereotype.Component;
import java.io.IOException;
@Slf4j
@Component
@RequiredArgsConstructor
public class PointsConsumer {
private final IdempotentService idempotentService;
private final JdbcTemplate jdbcTemplate;
/**
* 消费积分发放消息
* 加积分是"非幂等"操作,必须严格去重
*/
@RabbitListener(queues = "points.queue")
public void onGrantPoints(Message message, Channel channel) throws IOException {
String messageId = message.getMessageProperties().getMessageId();
String content = new String(message.getBody());
long deliveryTag = message.getMessageProperties().getDeliveryTag();
log.info("收到积分消息: messageId={}, content={}", messageId, content);
try {
// 幂等检查:积分发放绝对不能重复
if (!idempotentService.checkIdempotent(messageId, "GRANT_POINTS")) {
log.warn("积分消息重复,跳过: messageId={}", messageId);
channel.basicAck(deliveryTag, false);
return;
}
// 解析消息:userId:points
String[] parts = content.split(":");
Long userId = Long.parseLong(parts[0]);
int points = Integer.parseInt(parts[1]);
// 加积分(使用 ON DUPLICATE KEY 保证幂等)
jdbcTemplate.update(
"INSERT INTO user_points (user_id, points) VALUES (?, ?) " +
"ON DUPLICATE KEY UPDATE points = points + ?",
userId, points, points);
log.info("积分发放成功: userId={}, points={}", userId, points);
channel.basicAck(deliveryTag, false);
} catch (Exception e) {
log.error("积分消息处理失败: messageId={}", messageId, e);
channel.basicNack(deliveryTag, false, true);
}
}
}
5.7 配置
spring:
rabbitmq:
host: localhost
port: 5672
username: guest
password: guest
listener:
simple:
acknowledge-mode: manual
prefetch: 10
concurrency: 5
redis:
host: localhost
port: 6379
datasource:
url: jdbc:mysql://localhost:3306/test?useUnicode=true&characterEncoding=UTF-8&useSSL=false
username: root
password: 123456
5.8 测试Controller
@RestController
@RequestMapping("/api/test")
@RequiredArgsConstructor
public class TestController {
private final RabbitTemplate rabbitTemplate;
/**
* 模拟发送支付消息(可重复发送验证幂等)
*/
@PostMapping("/send-pay")
public String sendPayMessage(@RequestParam Long orderId,
@RequestParam String messageId) {
rabbitTemplate.convertAndSend("order.pay.queue", orderId.toString(),
message -> {
message.getMessageProperties().setMessageId(messageId);
return message;
});
return "消息已发送: messageId=" + messageId;
}
/**
* 模拟发送积分消息
*/
@PostMapping("/send-points")
public String sendPointsMessage(@RequestParam Long userId,
@RequestParam int points,
@RequestParam String messageId) {
String content = userId + ":" + points;
rabbitTemplate.convertAndSend("points.queue", content,
message -> {
message.getMessageProperties().setMessageId(messageId);
return message;
});
return "积分消息已发送: messageId=" + messageId;
}
}
六、进阶优化
6.1 Redis 数据过期后的兜底
Redis 的 key 设置了24小时过期,24小时后的重复消息 Redis 无法拦截,此时依赖数据库唯一索引兜底。
注意:数据库消费日志会无限增长,需要定期归档:
-- 每月清理3个月前的消费日志
DELETE FROM message_consume_log
WHERE consume_time < DATE_SUB(NOW(), INTERVAL 3 MONTH);
6.2 业务唯一键替代消息ID
有些场景消息没有唯一ID,可以用业务唯一键替代:
// 用"订单ID + 操作类型"作为幂等key
String idempotentKey = "order:" + orderId + ":pay";
推荐:优先使用业务唯一键,因为它比消息ID更稳定(消息ID可能因为重发而变化)。
6.3 布隆过滤器优化
如果消息量特别大(>1亿),Redis 存储所有消息ID会占用大量内存。可以使用布隆过滤器:
优点:内存占用小(1亿条数据约100MB)
缺点:有假阳性(小概率把未处理的消息判为已处理)
适用:对假阳性容忍度高的场景(如日志去重)。
6.4 分布式锁 + 状态检查
对于并发消费场景,Redis SET NX 可能因为并发时序问题失效。此时需要:
1. 获取分布式锁(按业务ID)
2. 检查业务状态
3. 执行业务
4. 释放锁
其中 分布式锁 的选型与实现,在 Java 高并发面试中也经常被问到,建议结合具体业务深入理解。
6.5 幂等失败的处理
如果第一层 Redis 通过了但业务执行失败,Redis 已经被"占用",消息重试时会被误判为重复。
解决方案:Redis 标记为"处理中",成功后再标记为"已完成":
// 状态1:处理中
redisTemplate.opsForValue().set(key, "PROCESSING", 1, TimeUnit.MINUTES);
try {
processBusiness();
// 状态2:已完成(长时间保留)
redisTemplate.opsForValue().set(key, "COMPLETED", 24, TimeUnit.HOURS);
} catch (Exception e) {
// 处理失败,删除"处理中"标记,允许重试
redisTemplate.delete(key);
throw e;
}
七、踩坑指南
| 坑 |
表现 |
解决方案 |
| 幂等Key设置错误 |
不同消息判为重复 |
用消息ID或"业务ID+操作类型" |
| Redis与DB不一致 |
Redis 过期后重复消息无法拦截 |
Redis TTL + DB唯一索引双保险 |
| 消息ID变化 |
同一条消息重发ID不同 |
用业务唯一键替代消息ID |
| 处理失败后重复拦截 |
处理失败的消息被误判重复 |
状态标记:处理中/已完成 |
| 数据库日志无限增长 |
磁盘空间被占满 |
定期归档历史日志 |
| 并发消费穿透 |
两个消费者同时处理 |
分布式锁 |
| 事务与幂等冲突 |
幂等成功后事务回滚 |
幂等记录与业务在同一事务 |
八、生产环境推荐
| 场景 |
推荐方案 |
| 普通业务 |
Redis SET NX + 业务状态检查 |
| 强一致场景 |
Redis + DB唯一索引 + 状态检查 |
| 高并发场景 |
Redis SET NX + 分布式锁 |
| 海量消息 |
布隆过滤器 + DB兜底 |
| 金融级场景 |
状态机 + 条件更新 + 唯一索引 |
参数推荐:
| 参数 |
推荐值 |
说明 |
| Redis TTL |
24小时 |
覆盖消息重试窗口 |
| 消费日志保留 |
3个月 |
平衡存储和审计 |
| 状态标记TTL |
1分钟(处理中)/ 24小时(已完成) |
区分处理中与完成 |
| 最大重试次数 |
3次 |
超过则进入死信队列 |
九、总结与思考
核心三句话:
- 消息重复不可避免——MQ 只保证 At Least Once,必须靠幂等消费。
- 幂等的核心是用唯一ID去重:Redis SET NX(快)+ 数据库唯一索引(稳)。
- 最强的幂等是业务天然幂等——通过状态检查和条件更新实现。
方案速查:
| 你的需求 |
推荐方案 |
| 快速幂等 |
Redis SET NX |
| 强一致幂等 |
Redis + DB唯一索引 |
| 天然幂等 |
状态机 + 条件更新 |
| 海量消息 |
布隆过滤器 |
附录:完整Maven依赖
<dependencies>
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-web</artifactId>
</dependency>
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-amqp</artifactId>
</dependency>
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-data-redis</artifactId>
</dependency>
<dependency>
<groupId>org.apache.commons</groupId>
<artifactId>commons-pool2</artifactId>
</dependency>
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-jdbc</artifactId>
</dependency>
<dependency>
<groupId>com.mysql</groupId>
<artifactId>mysql-connector-j</artifactId>
<scope>runtime</scope>
</dependency>
<dependency>
<groupId>org.projectlombok</groupId>
<artifactId>lombok</artifactId>
<optional>true</optional>
</dependency>
</dependencies>
⚠️ 运行前请确认:
- □ 确保 RabbitMQ、Redis、MySQL 均已启动
- □ 创建
orders、message_consume_log、user_points 三张表
- □ 在RabbitMQ中创建
order.pay.queue 和 points.queue
- □ 使用相同的
messageId 多次调用 POST /api/test/send-pay,验证重复消息被拦截
- □ 使用相同的
messageId 多次调用 POST /api/test/send-points,验证积分不会重复发放
- □ 观察日志,验证 Redis 拦截和数据库兜底是否生效