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 操作、异常处理策略和重试机制。

相关新闻

RabbitMQ如何保证消息不丢失

RabbitMQ如何保证消息不丢失

RabbitMQ 保证消息不丢失需要从‌生产者、Broker、消费者‌三个核心环节同时配置,缺一不可,核心是开启持久化、生产者确认和消费者手动ACK机制。一,RabbitMQ如何保证消息不丢失一、生产者端:确保消息成功送达Broker‌开启生产者确…

2026/7/22 16:04:24 阅读更多 →
嵌入式DMA开发实战:EDMA3中断、队列与优先级机制深度解析

嵌入式DMA开发实战:EDMA3中断、队列与优先级机制深度解析

1. 项目概述与核心价值 在嵌入式系统,尤其是高性能处理器(如TI的C6000系列DSP)的开发中,数据搬移的效率直接决定了整个系统的性能天花板。当你在处理音频流、视频帧或者雷达回波数据时,如果让CPU亲自去搬运每一个字节&…

2026/7/22 16:03:24 阅读更多 →
EDMA3高级DMA技术:Ping-Pong缓冲与传输链原理及McBSP音频实战

EDMA3高级DMA技术:Ping-Pong缓冲与传输链原理及McBSP音频实战

1. 项目概述与核心价值 在嵌入式系统开发,尤其是涉及高速数据流处理的应用中,CPU的资源是极其宝贵的。无论是音频编解码、图像传感器数据采集,还是高速通信接口的数据搬移,如果让CPU亲自去处理每一个字节的传输,其计算…

2026/7/22 16:03:24 阅读更多 →

最新新闻

算清AI这笔账:3个指标衡量企业每一美元智能产出

算清AI这笔账:3个指标衡量企业每一美元智能产出

过去两年,企业在 AI 上的投入像潮水一样涌来。从对话生成到图像识别,从自动客服到代码助手,几乎每一家像样的公司都在采购算力、试点模型、培训团队。但到 2025 年底,越来越多的 CFO 开始抬头问一个尴尬的问题:钱花出去…

2026/7/22 16:46:49 阅读更多 →
彻底终结RAG全量重建!生产级增量更新:Hash精准过滤+向量相似度兜底(小白也能懂)

彻底终结RAG全量重建!生产级增量更新:Hash精准过滤+向量相似度兜底(小白也能懂)

🔥 原创|小白易懂|生产级落地|无歧义干货 🏷️ 标签:#RAG #大模型知识库 #向量数据库 #增量更新 #AI工程化 🙈 很多RAG新手踩坑核心:分不清Hash精准匹配 & 向量相似度! 🙋 看完本文你将彻底弄懂:为什么生产环境必须「Hash前置+向量兜底」,根治文档改顺序、…

2026/7/22 16:46:49 阅读更多 →
计算机毕业设计之基于springboot的商场智能停车管理系统

计算机毕业设计之基于springboot的商场智能停车管理系统

本文设计并实现了一款基于Spring Boot的智能停车场管理系统,旨在解决现代城市停车难、管理效率低下的问题。系统分为用户端和管理员端,用户端提供个人中心、优惠政策查看、公告信息浏览等功能,使用户能够方便地管理自己的停车事务并享受实惠的…

2026/7/22 16:46:49 阅读更多 →
Wand-Enhancer完整指南:免费解锁专业版功能,彻底告别游戏修改限制

Wand-Enhancer完整指南:免费解锁专业版功能,彻底告别游戏修改限制

Wand-Enhancer完整指南:免费解锁专业版功能,彻底告别游戏修改限制 【免费下载链接】Wand-Enhancer Advanced UX and interoperability extension for Wand (WeMod) app 项目地址: https://gitcode.com/GitHub_Trending/we/Wand-Enhancer 你是否厌…

2026/7/22 16:46:49 阅读更多 →
CHI协议验证中的异常及边界验证

CHI协议验证中的异常及边界验证

CHI协议验证中的异常及边界验证 针对 CHI 协议的错误注入工具、覆盖率衡量方法及实际项目中的投入平衡 CHI 协议作为多核系统中复杂的缓存一致性协议,验证其行为需要强大的工具和方法来执行错误注入和边界条件测试,并衡量测试覆盖率。以下详细讨论常用工具、覆盖率评估方法及…

2026/7/22 16:46:49 阅读更多 →
TI McASP音频接口实战:从Burst到TDM模式配置与调试指南

TI McASP音频接口实战:从Burst到TDM模式配置与调试指南

1. 项目概述:从芯片手册到工程实践,拆解McASP的硬核玩法如果你正在用TI的DSP或者某些高性能处理器做音频相关的嵌入式开发,那你大概率绕不开一个名字:McASP。这玩意儿全称叫Multichannel Audio Serial Port,翻译过来就…

2026/7/22 16:45:48 阅读更多 →

日新闻

TI DSP系统配置模块SYSCFG详解:中断机制与主设备优先级配置实战

TI DSP系统配置模块SYSCFG详解:中断机制与主设备优先级配置实战

1. 项目概述与SYSCFG模块的核心价值在嵌入式系统,尤其是像TI C6000系列这样的高性能DSP开发中,我们常常会与芯片手册里那些密密麻麻的寄存器打交道。很多开发者可能更关注算法实现、内存优化或者外设驱动,但对于一个稳定、高效的系统而言&…

2026/7/22 0:00:26 阅读更多 →
微信Server酱:高到达率的应急通知方案实践

微信Server酱:高到达率的应急通知方案实践

1. 为什么我们需要"最次"的通知方案? 在数字化协作环境中,消息通知系统的重要性不言而喻明。但现实情况是,企业级通知方案往往需要复杂的API对接(如企业微信、钉钉、飞书),个人开发者的小项目又经…

2026/7/22 0:00:26 阅读更多 →
甲方要的“简洁“PPT,到底是简洁还是省事?

甲方要的“简洁“PPT,到底是简洁还是省事?

甲方说"简洁一点",乙方听到的是"少做几页"。甲方说"不要太复杂",乙方理解成"别放图表了"。结果交过去,甲方说"我说的简洁不是这个意思"。"简洁"这个词在PPT语境里,是…

2026/7/22 0:00:26 阅读更多 →

周新闻

Go语言静态资源打包方案对比与实践指南

Go语言静态资源打包方案对比与实践指南

1. 项目背景与核心需求在Go语言开发中,我们经常需要处理静态资源文件的打包问题。无论是Web应用的模板文件、前端资源,还是配置文件、证书等,都需要随程序一起分发。传统做法是将这些文件与编译后的二进制文件放在同一目录下,但这…

2026/7/22 8:58:19 阅读更多 →
Go语言实现高性能LDAP认证服务的架构与实践

Go语言实现高性能LDAP认证服务的架构与实践

1. 项目背景与核心价值LDAP(轻量级目录访问协议)作为企业级身份认证的黄金标准,已经服务了超过80%的财富500强公司。我在金融科技领域实施统一认证体系时,发现传统Java方案存在启动慢、内存占用高等痛点。而Go语言凭借其协程并发模…

2026/7/21 5:34:47 阅读更多 →
【AI面试官实战指南】:用ChatGPT模拟10类高频技术岗面试,3天提升应答精准度92%

【AI面试官实战指南】:用ChatGPT模拟10类高频技术岗面试,3天提升应答精准度92%

更多请点击: https://intelliparadigm.com 第一章:AI面试官实战指南的核心价值与适用场景 AI面试官并非替代人类HR的“黑箱工具”,而是以可解释、可审计、可迭代的方式,赋能招聘全链路的关键基础设施。其核心价值在于将主观经验沉…

2026/7/22 12:54:44 阅读更多 →

月新闻