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:49阅读更多 →
AI 数字员工 OpenClaw 实操分享 本地运行保障文件数据隐私安全(含安装包)

AI 数字员工 OpenClaw 实操分享 本地运行保障文件数据隐私安全(含安装包)

🦞OpenClaw(小龙虾)Windows 一键部署保姆级教程|10 分钟搭建专属数字员工 适配平台:Windows 10/11(64 位)|零基础友好|全图形界面操作|无开发门槛 &#x1…

2026/7/22 16:04:49阅读更多 →
RabbitMQ Exporter高级功能:BERT编码与no_sort特性提升监控效率

RabbitMQ Exporter高级功能:BERT编码与no_sort特性提升监控效率

RabbitMQ Exporter高级功能:BERT编码与no_sort特性提升监控效率 【免费下载链接】rabbitmq_exporter Prometheus exporter for RabbitMQ 项目地址: https://gitcode.com/gh_mirrors/ra/rabbitmq_exporter RabbitMQ Exporter是一款针对RabbitMQ的Prometheus监…

2026/7/22 16:02:48阅读更多 →
嵌入式USB接收端点寄存器深度解析:从RXMAXP到RXCSR的配置与调试

嵌入式USB接收端点寄存器深度解析:从RXMAXP到RXCSR的配置与调试

1. 项目概述与核心价值搞嵌入式USB开发,尤其是基于TI这类厂商的专用控制器,最让人头疼的往往不是协议栈本身,而是那一堆密密麻麻的寄存器手册。手册里每个位域都写得清清楚楚,但真到了写驱动、调通信的时候,怎么把这些…

2026/7/22 16:58:58阅读更多 →
嵌入式开发实战:USB端点与看门狗寄存器配置避坑指南

嵌入式开发实战:USB端点与看门狗寄存器配置避坑指南

1. 项目概述与核心价值在嵌入式系统开发里,和硬件打交道是绕不开的基本功。CPU要指挥USB收发数据,或者让看门狗定时器在关键时刻“踢”系统一脚,靠的都是直接读写那一组组看似冰冷的寄存器。很多新手觉得寄存器配置就是对着手册填地址和数值&…

2026/7/22 16:58:58阅读更多 →
多厂区标签模板统一与可追溯管理改进路径

多厂区标签模板统一与可追溯管理改进路径

作为一名企业QA,我深知标签虽小,却是产品身份的唯一标识,也是质量追溯的第一道防线。然而,在过去很长一段时间里,标签管理恰恰是我们质量管理体系中最头疼的环节。最核心的问题就是标准不统一、模板版本混乱。集团下多…

2026/7/22 16:58:58阅读更多 →
50mm小板一焊,100dB回音消失、90dB噪音退散、8米拾音稳了

50mm小板一焊,100dB回音消失、90dB噪音退散、8米拾音稳了

硬件工程师福音: 一颗33mm 22mm的语音处理模组,把降噪、AEC、波束成形、功放、USB声卡全部装了进去。焊上去,上电,音频部分就搞定了——不用调参数,不用写算法,不用外挂DSP。前言做通话类产品的硬件工程师…

2026/7/22 16:58:58阅读更多 →
深度解析 AP-0316:集 AI降噪 + AEC回消 + USB + 波束拾音于一体的“全能语音处理模组“

深度解析 AP-0316:集 AI降噪 + AEC回消 + USB + 波束拾音于一体的“全能语音处理模组“

导读:在智能门禁、车载通话、会议设备、安防监控等场景中,"噪音大、回音重、拾音差"一直是困扰开发者的三大音频痛点。本文深度解析 AP-0316 多功能语音处理模组——一款将 AI ENC 降噪、100dB AEC 回音消除、双麦波束拾音、USB/I2S/模拟多接口…

2026/7/22 16:58:58阅读更多 →
TI USBSS中断寄存器深度解析与嵌入式USB驱动开发实践

TI USBSS中断寄存器深度解析与嵌入式USB驱动开发实践

1. USB中断机制与TI USBSS架构深度解析 在嵌入式系统开发中,中断机制是连接硬件事件与软件响应的核心桥梁。它就像一位高效的秘书,当有紧急事件(如数据到达、设备连接)发生时,会立刻打断CPU当前的工作,递上…

2026/7/22 16:56:57阅读更多 →
Go语言静态资源打包方案对比与实践指南

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

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

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

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

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

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

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

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

2026/7/22 0:53:59阅读更多 →
中小企业小程序开发公司怎么选:预算、上手和售后避坑指南

中小企业小程序开发公司怎么选:预算、上手和售后避坑指南

中小企业做小程序,最常见的矛盾是预算有限,但又不希望功能太单薄;没有技术团队,但又希望后续能自己运营;想快速上线,又担心隐性收费和售后失联。选型时如果只看“低价套餐”或“案例数量”,很容…

2026/7/22 0:01:17阅读更多 →
GEO优化如何沉淀长期内容资产?广拓时代谈AI搜索时代的内容ROI

GEO优化如何沉淀长期内容资产?广拓时代谈AI搜索时代的内容ROI

企业做营销,最怕钱花完了,资产没有留下。 效果广告能带来一段时间的曝光,但预算停止后,流量往往也随之停止。短视频内容可能在几天内冲高,也可能很快沉下去。AI搜索时代,企业需要重新思考一个问题&#xff…

2026/7/22 0:01:17阅读更多 →
Agent 终态判定:何时该停止思考、给出最终回复

Agent 终态判定:何时该停止思考、给出最终回复

Agent 终态判定:何时该停止思考、给出最终回复 一、你的 Agent 在"再想想"的循环里绕了 12 轮,用户已经关窗口了 Agent 与人最大的区别是:人知道什么时候该停下来给答案,Agent 会一直"想"下去。你给 Agent 接…

2026/7/22 0:01:17阅读更多 →
YOLOv8推理性能优化:从1.2FPS到35FPS的全链路加速实践

YOLOv8推理性能优化:从1.2FPS到35FPS的全链路加速实践

如果你在部署 YOLOv8 时,发现推理速度只有可怜的 1-2 FPS,而别人的演示视频却能跑到 30 FPS 以上,那么问题很可能不在模型本身,而在于你的整个处理链路。很多开发者拿到一个训练好的 YOLOv8 模型后,会直接使用官方示例…

2026/7/21 22:53:50阅读更多 →
Coze与Dify对比指南:低代码AI应用开发从入门到实战

Coze与Dify对比指南:低代码AI应用开发从入门到实战

1. 从零到一:为什么你需要了解 Coze 和 Dify?如果你对 AI 应用开发感兴趣,但一看到“大模型”、“智能体”、“工作流”这些词就头疼,觉得门槛太高,那这篇文章就是为你准备的。很多开发者,包括我自己&#…

2026/7/21 18:53:30阅读更多 →
AI生图工具怎么选?2026年6月版实测对比

AI生图工具怎么选?2026年6月版实测对比

做自媒体的朋友应该都有体会:配图一直是个让人头疼的问题。2026年,AI生图工具已经非常成熟了,但工具太多反而不知道怎么选。以下是截至2026年6月我对主流AI生图工具的实测对比。Midjourney V8.1:速度之王2026年6月11日&#xff0c…

2026/7/21 18:53:30阅读更多 →