RabbitMQ重复消费问题
RabbitMQ 出现重复消费的核心原因在于网络抖动导致 ACK 确认丢失。当消费者处理完业务但尚未发送 ACK,或者 ACK 在传输过程中丢失时,RabbitMQ 会认为消息未被成功处理,从而将消息重新投递给消费者。此外,消费者处理失败后手动将消息重新入队(Requeue)也会导致重复消费。
一,RabbitMQ重复消费问题
解决重复消费问题的核心思路是幂等性设计,即确保同一消息无论被消费多少次,对业务数据产生的最终影响与消费一次完全一致。以下是几种主流的落地方案:
一、核心解决方案
1. 唯一 ID + Redis 去重(推荐高性能场景)
这是最常用且性能较好的方案。
实现逻辑:
生产者在发送消息时生成一个全局唯一的业务 ID(如 UUID 或雪花算法 ID),并放入消息头或消息体中。
消费者接收到消息后,先提取该唯一 ID。
使用 Redis 的 SETNX(Set if Not Exists)命令尝试写入该 ID。
如果返回 1,说明是第一次消费,执行业务逻辑,并在业务完成后保留该 ID(可设置合理过期时间以防内存溢出)。
如果返回 0,说明该 ID 已存在,直接丢弃消息或返回成功 ACK,不再执行业务逻辑。
优势:Redis 读写速度极快,适合高并发场景。
注意:需为 Redis Key 设置过期时间,避免内存无限增长。
2. 数据库唯一索引/去重表(推荐强一致性场景)
利用数据库的唯一约束机制保证幂等性。
实现逻辑:
在业务表中增加一个唯一字段(如 message_id 或 biz_no),专门存储消息的唯一标识。
或者建立一张独立的“消息去重表”,包含 message_id 主键。
消费者在处理业务前,先尝试插入该唯一 ID。
如果插入成功,继续执行业务逻辑。
如果抛出“唯一键冲突”异常,说明消息已处理,直接捕获异常并 ACK 确认。
优势:依靠数据库事务保证强一致性,可靠性最高。
缺点:频繁查询或插入数据库可能成为性能瓶颈。
3. 业务状态机判断(推荐状态流转场景)
适用于具有明确状态变更的业务,如订单状态更新。
实现逻辑:
在执行更新操作时,带上前置状态条件。例如:UPDATE orders SET status = ‘PAID’ WHERE id = 1001 AND status = ‘UNPAID’。
如果重复消费,由于状态已经变为 ‘PAID’,SQL 执行影响的行数为 0,业务逻辑自然跳过,不会产生副作用。
**优势无需额外存储组件,代码侵入小。
二、辅助优化措施
开启手动 ACK 模式
务必关闭自动 ACK(Auto Ack),改为在业务逻辑完全执行成功后再手动发送 basicAck。
若业务执行失败,可根据策略选择 basicNack 重新入队或转入死信队列,避免消息静默丢失或无限重试导致的数据混乱。
合理设置重试机制
如果因临时故障(如数据库连接超时)导致消费失败,不要立即无限重试。建议结合指数退避算法或设置最大重试次数,超过次数后转入死信队列人工介入,防止重复消费风暴。
消息去重表配合过期清理
若使用 Redis 或数据库去重,需定期清理过期的去重记录,以节省存储空间。
三、方案对比总结
| 方案 | 适用场景 | 优点 | 缺点 |
|---|---|---|---|
| Redis SETNX | 高并发、对性能要求高 | 速度快,支持高吞吐 | 需维护 Redis,存在短暂不一致风险 |
| 数据库唯一索引 | 金融、订单等强一致性场景 | 可靠性最高,强一致 | 数据库压力大,性能相对较低 |
| 状态机判断 | 订单状态变更、审批流 | 无额外组件依赖,逻辑简单 | 仅适用于有状态流转的业务 |
在实际项目中,建议组合使用多种方案。例如,先用 Redis 进行快速去重拦截大部分重复请求,再在数据库层面通过唯一索引做最终兜底,从而兼顾性能与数据安全性。
二,RabbitMQ重复消费的具体案例
下面我们以一个“用户积分增加”的业务场景为例,展示如何使用唯一 ID + Redis 去重方案来防止重复消费。
1. 项目结构与依赖
首先,确保你的pom.xml中包含以下依赖:
<dependencies><!-- Spring Boot Starter for AMQP (RabbitMQ) --><dependency><groupId>org.springframework.boot</groupId><artifactId>spring-boot-starter-amqp</artifactId></dependency><!-- Spring Boot Starter for Data Redis --><dependency><groupId>org.springframework.boot</groupId><artifactId>spring-boot-starter-data-redis</artifactId></dependency><!-- 其他必要依赖,如 Lombok, Web 等 --><dependency><groupId>org.projectlombok</groupId><artifactId>lombok</artifactId><optional>true</optional></dependency></dependencies>2. 消息生产者(Producer)
生产者在发送消息时,需要生成一个全局唯一的业务 ID(bizId)并放入消息头。
importorg.springframework.amqp.core.Message;importorg.springframework.amqp.core.MessageBuilder;importorg.springframework.amqp.core.MessageProperties;importorg.springframework.amqp.rabbit.core.RabbitTemplate;importorg.springframework.beans.factory.annotation.Autowired;importorg.springframework.stereotype.Service;importjava.util.UUID;@ServicepublicclassPointsProducerService{@AutowiredprivateRabbitTemplaterabbitTemplate;/** * 发送增加积分的消息 * @param userId 用户ID * @param points 增加的积分数 */publicvoidsendPointsMessage(LonguserId,Integerpoints){// 1. 构造业务数据PointsMessagepointsMessage=newPointsMessage(userId,points);// 2. 生成全局唯一的业务ID (这里使用UUID,生产环境建议用雪花算法)StringbizId=UUID.randomUUID().toString();// 3. 构建消息,将 bizId 放入消息头Messagemessage=MessageBuilder.withBody(pointsMessage.toString().getBytes()).setContentType(MessageProperties.CONTENT_TYPE_JSON).setHeader("bizId",bizId)// 关键:设置唯一标识.build();// 4. 发送消息到指定交换机和路由键rabbitTemplate.send("points.exchange","points.add",message);System.out.println("消息发送成功,bizId: "+bizId+", 内容: "+pointsMessage);}@Data@AllArgsConstructorstaticclassPointsMessage{privateLonguserId;privateIntegerpoints;// 省略 toString 方法}}3. 消息消费者(Consumer)与幂等性处理
消费者在消费前,先通过 Redis 检查bizId是否已处理。
importorg.springframework.amqp.core.Message;importorg.springframework.amqp.rabbit.annotation.RabbitListener;importorg.springframework.beans.factory.annotation.Autowired;importorg.springframework.data.redis.core.StringRedisTemplate;importorg.springframework.stereotype.Service;importorg.springframework.transaction.annotation.Transactional;importjava.nio.charset.StandardCharsets;importjava.util.concurrent.TimeUnit;@ServicepublicclassPointsConsumerService{@AutowiredprivateStringRedisTemplateredisTemplate;@AutowiredprivateUserPointsServiceuserPointsService;// 假设的业务服务// Redis Key 的前缀privatestaticfinalStringPOINTS_MSG_PREFIX="points:msg:id:";// 去重记录过期时间,例如 24 小时privatestaticfinallongEXPIRE_HOURS=24;/** * 监听积分增加队列 */@RabbitListener(queues="points.add.queue")@Transactional(rollbackFor=Exception.class)publicvoidhandlePointsMessage(Messagemessage){// 1. 从消息头中提取唯一业务IDStringbizId=message.getMessageProperties().getHeader("bizId");if(bizId==null||bizId.isEmpty()){// 没有 bizId,消息格式错误,可以记录日志并拒绝消息(不入队)System.err.println("消息缺少 bizId,拒绝处理。消息体: "+newString(message.getBody()));// 这里可以根据策略选择 basicNack 并 requeue=falsereturn;}StringredisKey=POINTS_MSG_PREFIX+bizId;// 2. 使用 SETNX 尝试在 Redis 中设置 KeyBooleanisFirstConsume=redisTemplate.opsForValue().setIfAbsent(redisKey,"PROCESSED",EXPIRE_HOURS,TimeUnit.HOURS);if(Boolean.TRUE.equals(isFirstConsume)){// 2.1 第一次消费,执行业务逻辑try{// 解析消息体StringmessageBody=newString(message.getBody(),StandardCharsets.UTF_8);PointsProducerService.PointsMessagepointsMessage=parseMessage(messageBody);// 核心业务:为用户增加积分userPointsService.addPoints(pointsMessage.getUserId(),pointsMessage.getPoints());System.out.println("业务执行成功,bizId: "+bizId+", userId: "+pointsMessage.getUserId());// 3. 业务成功,可以手动发送 ACK (如果配置了手动ACK)// channel.basicAck(deliveryTag, false);}catch(Exceptione){// 业务执行失败System.err.println("业务执行失败,bizId: "+bizId+", 错误: "+e.getMessage());// 删除 Redis 中的记录,允许消息重试(根据业务决定)redisTemplate.delete(redisKey);// 抛出异常,让消息重回队列或进入死信队列(根据配置)thrownewRuntimeException("处理消息失败",e);}}else{// 2.2 重复消费,直接确认消息,不执行业务System.out.println("检测到重复消息,bizId: "+bizId+",已跳过处理。");// 直接发送 ACK,避免消息堆积// channel.basicAck(deliveryTag, false);}}privatePointsProducerService.PointsMessageparseMessage(Stringbody){// 简化的 JSON 解析,实际使用 Jackson/Gson// 示例:{"userId":123,"points":10}// 这里返回一个模拟对象returnnewPointsProducerService.PointsMessage(123L,10);}}4. 业务服务层(Service)
importorg.springframework.stereotype.Service;@ServicepublicclassUserPointsService{/** * 为用户增加积分(幂等操作) * @param userId 用户ID * @param points 增加的积分数 */publicvoidaddPoints(LonguserId,Integerpoints){// 这里模拟数据库操作// 实际应包含事务、校验等逻辑System.out.println("为用户 "+userId+" 增加积分 "+points+" 点。");// 执行 UPDATE user_points SET points = points + ? WHERE user_id = ?}}5. 配置示例(application.yml)
spring:rabbitmq:host:localhostport:5672username:guestpassword:guest# 开启手动确认模式(ACK)listener:simple:acknowledge-mode:manual# 关键配置prefetch:1# 每次只预取一条消息,避免堆积redis:host:localhostport:6379# password: 你的密码6. 流程总结与测试要点
流程:
- 生产者发送消息,携带唯一
bizId。 - 消费者收到消息,用
bizId作为 Key 尝试写入 Redis。 - 写入成功(SETNX 返回 true)→ 执行业务 → 业务成功则完成。
- 写入失败(SETNX 返回 false)→ 消息重复 → 直接 ACK 丢弃。
- 生产者发送消息,携带唯一
测试重复消费:
- 在消费者业务逻辑中(
addPoints方法)模拟一个较长的处理时间或手动抛出异常。 - 由于配置了手动 ACK 且未发送,RabbitMQ 会在连接断开或 Channel 关闭后将消息重新投递。
- 观察日志:第一次会打印“业务执行成功”,第二次及以后会打印“检测到重复消息,已跳过处理”。
- 在消费者业务逻辑中(
关键点:
- Redis 键过期:必须设置过期时间,防止内存无限增长。
- 异常处理:业务失败时应删除 Redis 键,允许消息重试(根据业务决定是否重试)。
- 手动 ACK:确保业务成功后才确认消息,这是防止消息丢失的第一道防线。
- bizId 生成:生产环境建议使用分布式 ID 生成器(如雪花算法),确保全局唯一和高性能。
这个案例展示了从消息生产、幂等性判断到业务处理的完整闭环,你可以根据实际业务需求调整 Redis 操作、异常处理策略和重试机制。