RabbitMQ重复消费问题 RabbitMQ 出现重复消费的核心原因在于‌网络抖动导致 ACK 确认丢失‌。当消费者处理完业务但尚未发送 ACK或者 ACK 在传输过程中丢失时RabbitMQ 会认为消息未被成功处理从而将消息重新投递给消费者。此外消费者处理失败后手动将消息重新入队Requeue也会导致重复消费。一RabbitMQ重复消费问题解决重复消费问题的核心思路是‌幂等性设计‌即确保同一消息无论被消费多少次对业务数据产生的最终影响与消费一次完全一致。以下是几种主流的落地方案一、核心解决方案1. 唯一 ID Redis 去重推荐高性能场景这是最常用且性能较好的方案。实现逻辑‌生产者在发送消息时生成一个全局唯一的业务 ID如 UUID 或雪花算法 ID并放入消息头或消息体中。消费者接收到消息后先提取该唯一 ID。使用 Redis 的 SETNXSet 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 模式‌务必关闭自动 ACKAuto Ack改为在业务逻辑完全执行成功后再手动发送 basicAck。若业务执行失败可根据策略选择 basicNack 重新入队或转入死信队列避免消息静默丢失或无限重试导致的数据混乱。合理设置重试机制‌如果因临时故障如数据库连接超时导致消费失败不要立即无限重试。建议结合指数退避算法或设置最大重试次数超过次数后转入死信队列人工介入防止重复消费风暴。消息去重表配合过期清理‌若使用 Redis 或数据库去重需定期清理过期的去重记录以节省存储空间。三、方案对比总结方案适用场景优点缺点‌Redis SETNX‌高并发、对性能要求高速度快支持高吞吐需维护 Redis存在短暂不一致风险‌数据库唯一索引‌金融、订单等强一致性场景可靠性最高强一致数据库压力大性能相对较低‌状态机判断‌订单状态变更、审批流无额外组件依赖逻辑简单仅适用于有状态流转的业务在实际项目中建议‌组合使用‌多种方案。例如先用 Redis 进行快速去重拦截大部分重复请求再在数据库层面通过唯一索引做最终兜底从而兼顾性能与数据安全性。二RabbitMQ重复消费的具体案例下面我们以一个“用户积分增加”的业务场景为例展示如何使用唯一 ID Redis 去重方案来防止重复消费。1. 项目结构与依赖首先确保你的pom.xml中包含以下依赖dependencies!-- Spring Boot Starter for AMQP (RabbitMQ) --dependencygroupIdorg.springframework.boot/groupIdartifactIdspring-boot-starter-amqp/artifactId/dependency!-- Spring Boot Starter for Data Redis --dependencygroupIdorg.springframework.boot/groupIdartifactIdspring-boot-starter-data-redis/artifactId/dependency!-- 其他必要依赖如 Lombok, Web 等 --dependencygroupIdorg.projectlombok/groupIdartifactIdlombok/artifactIdoptionaltrue/optional/dependency/dependencies2. 消息生产者Producer生产者在发送消息时需要生成一个全局唯一的业务 IDbizId并放入消息头。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. 构造业务数据PointsMessagepointsMessagenewPointsMessage(userId,points);// 2. 生成全局唯一的业务ID (这里使用UUID生产环境建议用雪花算法)StringbizIdUUID.randomUUID().toString();// 3. 构建消息将 bizId 放入消息头MessagemessageMessageBuilder.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);}DataAllArgsConstructorstaticclassPointsMessage{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_PREFIXpoints:msg:id:;// 去重记录过期时间例如 24 小时privatestaticfinallongEXPIRE_HOURS24;/** * 监听积分增加队列 */RabbitListener(queuespoints.add.queue)Transactional(rollbackForException.class)publicvoidhandlePointsMessage(Messagemessage){// 1. 从消息头中提取唯一业务IDStringbizIdmessage.getMessageProperties().getHeader(bizId);if(bizIdnull||bizId.isEmpty()){// 没有 bizId消息格式错误可以记录日志并拒绝消息不入队System.err.println(消息缺少 bizId拒绝处理。消息体: newString(message.getBody()));// 这里可以根据策略选择 basicNack 并 requeuefalsereturn;}StringredisKeyPOINTS_MSG_PREFIXbizId;// 2. 使用 SETNX 尝试在 Redis 中设置 KeyBooleanisFirstConsumeredisTemplate.opsForValue().setIfAbsent(redisKey,PROCESSED,EXPIRE_HOURS,TimeUnit.HOURS);if(Boolean.TRUE.equals(isFirstConsume)){// 2.1 第一次消费执行业务逻辑try{// 解析消息体StringmessageBodynewString(message.getBody(),StandardCharsets.UTF_8);PointsProducerService.PointsMessagepointsMessageparseMessage(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. 业务服务层Serviceimportorg.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.ymlspring:rabbitmq:host:localhostport:5672username:guestpassword:guest# 开启手动确认模式ACKlistener: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 操作、异常处理策略和重试机制。