ARTICLE DETAIL

资讯详情

深耕郑州网站建设与运营推广的一线实战洞察。

MQ重复消费问题怎么解决

MQ重复消费问题怎么解决 MQ重复消费问题怎么解决重复消费的根本原因 核心原因消息确认机制与网络不可靠性的矛盾简单说就是消费者已经处理完消息但 ACK 没成功回到 MQ 服务器导致 MQ 认为消息处理失败触发重发机制。通用解决方案幂等性设计 ✅解决重复消费的唯一正确思路是让消费操作本身具备幂等性即同一条消息消费多次和消费一次的效果完全相同。1. 数据库唯一索引法最简单高效⭐技术亮点利用数据库原子性无需额外查询性能最优ServiceSlf4jpublicclassOrderPayConsumer{AutowiredprivateDuplicateRecordMapperduplicateRecordMapper;AutowiredprivateOrderServiceorderService;Transactional(rollbackForException.class)publicvoidonMessage(MessageExtmessage){// 提取业务唯一标识优先使用业务ID而非MQ的messageIdStringorderIdJSON.parseObject(message.getBody(),OrderMessage.class).getOrderId();try{// 先插入去重记录唯一索引冲突直接抛出异常// 技术亮点利用数据库原子性避免查询-插入的竞态条件duplicateRecordMapper.insertSelective(DuplicateRecord.builder().businessId(orderId).businessType(ORDER_PAY).messageId(message.getMsgId()).createTime(newDate()).build());// 执行业务逻辑orderService.updateOrderStatusToPaid(orderId);log.info(订单支付消息处理成功orderId: {},orderId);}catch(DuplicateKeyExceptione){// 重复消息直接返回成功告诉MQ不要再重发log.warn(检测到重复消息直接ACKorderId: {}, messageId: {},orderId,message.getMsgId());}catch(Exceptione){// 业务异常抛出异常让MQ重试log.error(订单支付消息处理失败orderId: {},orderId,e);thrownewRuntimeException(消息处理失败,e);}}}整理了面试真题、每日技术知识点、系统学习路线√X搜「Rain的Java大神之路」每天拆一个知识点陪你悄悄变强有空来坐坐。2. Redis 分布式锁法高并发场景技术亮点使用 Lua 脚本保证原子性防止锁误删ComponentSlf4jpublicclassRedisIdempotentUtil{AutowiredprivateStringRedisTemplateredisTemplate;// Lua脚本原子性判断锁是否存在并删除privatestaticfinalStringUNLOCK_SCRIPTif redis.call(get, KEYS[1]) ARGV[1] then return redis.call(del, KEYS[1]) else return 0 end;/** * 尝试获取幂等锁 * param businessId 业务唯一标识 * param expireTime 过期时间秒 * return 是否获取成功 */publicbooleantryLock(StringbusinessId,longexpireTime){Stringkeyidempotent:order_pay:businessId;// 使用UUID作为锁的值防止误删其他线程的锁StringvalueUUID.randomUUID().toString();BooleanresultredisTemplate.opsForValue().setIfAbsent(key,value,Duration.ofSeconds(expireTime));returnBoolean.TRUE.equals(result);}/** * 释放幂等锁 */publicvoidunlock(StringbusinessId,Stringvalue){Stringkeyidempotent:order_pay:businessId;redisTemplate.execute(newDefaultRedisScript(UNLOCK_SCRIPT,Long.class),Collections.singletonList(key),value);}}3. 业务状态机法最优雅技术亮点利用数据库行锁保证状态更新的原子性MapperpublicinterfaceOrderMapper{/** * 原子性更新订单状态 * 技术亮点通过WHERE条件保证只有待支付状态才能更新为已支付 * 同时利用数据库行锁防止并发更新 */Update(UPDATE orders SET status #{newStatus}, update_time NOW() WHERE id #{orderId} AND status #{oldStatus})intupdateStatus(Param(orderId)StringorderId,Param(oldStatus)StringoldStatus,Param(newStatus)StringnewStatus);}// 业务层调用TransactionalpublicvoidupdateOrderStatusToPaid(StringorderId){intaffectedRowsorderMapper.updateStatus(orderId,UNPAID,PAID);if(affectedRows0){// 订单状态不是待支付说明已经处理过log.warn(订单状态已更新无需重复处理orderId: {},orderId);return;}// 执行后续业务逻辑扣减库存、生成物流单等inventoryService.deductStock(orderId);logisticsService.createLogisticsOrder(orderId);}不同解决方案对比 解决方案实现难度性能适用场景缺点数据库唯一索引⭐⭐⭐⭐简单插入操作依赖数据库有性能瓶颈Redis 分布式锁⭐⭐⭐⭐⭐⭐高并发复杂业务实现复杂需处理锁过期消息去重表⭐⭐⭐⭐通用场景增加数据库压力业务状态机⭐⭐⭐⭐⭐⭐⭐⭐有状态流转的业务与业务强耦合核心技术难点与解决方案 ⚠️这是生产环境中最容易踩坑的地方也是面试官最关心的部分。技术难点问题描述解决方案“查询 - 插入” 竞态条件两个线程同时查询到消息未消费都执行业务逻辑✅ 使用数据库唯一索引✅ 使用 Redis 原子操作❌ 禁止先查询再插入分布式锁过期问题业务处理时间超过锁过期时间导致锁自动释放✅ 设置合理的过期时间业务耗时的 3-5 倍✅ 使用看门狗机制自动续期✅ 业务逻辑中增加二次校验死信队列处理消息多次重试失败后进入死信队列无人处理✅ 配置死信队列告警✅ 建立死信消息人工处理流程✅ 定期清理过期的死信消息去重表数据膨胀去重记录不断累积导致查询性能下降✅ 按时间分表✅ 定期归档历史数据✅ 设置数据过期时间如保留 30 天消息乱序问题重发的消息可能比后续消息先到达✅ 业务状态机天然支持乱序✅ 按业务 ID 分区消费✅ 增加版本号控制事务一致性问题业务逻辑执行成功但去重记录插入失败✅ 使用本地事务保证原子性✅ 先插入去重记录再执行业务逻辑✅ 异常时回滚整个事务大厂最佳实践 优先使用业务唯一标识而不是 MQ 生成的 messageId消费端先做幂等校验再执行业务逻辑设置合理的重试次数一般 3-5 次超过次数则转入死信队列定期清理过期的去重记录避免表数据量过大监控重复消费的次数及时发现异常情况所有 MQ 消费者必须实现幂等性这是分布式系统的基本要求总结 MQ 重复消费是分布式系统中不可避免的问题不能指望通过 “避免消息重发” 来解决必须从消费端入手通过幂等性设计来保证系统的正确性。在实际项目中我会优先选择数据库唯一索引或业务状态机这两种方案因为它们实现简单、性能好且不易出错。对于高并发场景会结合Redis 分布式锁来进一步提升性能。真实面试模拟真实面试模拟面试官我来看个场景设计题“MQ重复消费问题怎么解决”你之前项目里应该踩过这个坑吧先说说为什么会发生重复消费不用背八股就讲你的理解。候选人你好的我遇到的基本是三种情况一是生产者那端因为网络超时没收到ACK以为没发成功结果重发了同一条消息二是消费者挂了偏移量没提交重启后重复拉取三是Kafka的rebalance分区重新分配时会有一段重叠消费。简单说MQ保证的是至少一次投递不承诺恰好一次所以幂等必须业务侧自己扛。面试官对抓得很准。 那给你一个场景订单支付回调同一条消息可能被回调两次你打算怎么设计直接聊聊思路。候选人这种对一致性要求高的我第一选择是数据库去重表。核心思路消息里一定有个业务流水号比如支付回调的outTradeNo。我建一张payment_dedup表这个流水号字段加唯一索引。消费时先尝试 insert 这条流水号。如果插入成功说明第一次处理接着做业务比如更新订单状态然后和去重记录放在同一个本地事务里提交如果插入的时候报了唯一键冲突那就是重复消息直接返回ACK什么都不做。面试官嗯逻辑很清晰。那事务边界你怎么把握如果把去重插入和业务放到两个事务里会怎样候选人那就有漏洞了。比如先去重插入成功但业务更新失败回滚这时候标记留下来了真正的重试过来却被误判为重复业务就丢了。所以必须用同一个数据库事务insert 去重表 更新订单状态要么一起成功要么一起回滚。这样才能保证原子性。面试官 不错防住了常见坑。那如果这个回调QPS很高担心数据库扛不住你怎么优化候选人可以引入Redis 轻量级幂等。用SET key value NX EX 1800这种原子命令key 还是outTradeNo。如果 set 成功才执行业务如果已存在就是重复消息直接跳过。但有个细节业务处理成功后不能立刻 del 这个 key因为重复消息可能在这个窗口期之后才来。要靠 TTL 自己过期保证一个足够长的去重窗口。另外Redis 这种方案在极端情况下有丢失标记的风险比如主从切换所以金融类核心链路我还是会保留数据库去重表做最终兜底。面试官不错权衡做得很细。那更新场景比如订单状态从“待支付”改成“已支付”除了去重表有没有更轻量的办法候选人可以用版本号 affected rows 判断。SQL 这么写UPDATE orders SET statusPAID, versionversion1 WHERE id#{id} AND version#{oldVersion}第一条消息过来更新成功影响行数1第二条重复消息带着旧版本号再执行影响行数就是0。代码里判断行数0就认为是重复直接丢弃。这比单纯用WHERE statusWAIT_PAY要好能防止并发覆盖也天然幂等。面试官赞这个状态机模式用得很熟。如果我让你画个图把这些方案串起来给别人讲清楚你怎么画候选人我习惯用脑图或者流程图。比如处理流程可以这样画方案选型大概这样面试官很直观一看就懂。最后问一句你踩过哪些幂等设计的坑候选人三个最典型的先删缓存再执行业务第二次消费发现缓存没了导致重复处理去重标记放到业务事务外面造成误判直接用MQ自身的messageId做幂等键结果MQ重发时消息ID变了导致防不住。所以我一直要求团队只用业务唯一标识做幂等键标记和业务必须原子落地。面试官:“前面你说的思路很清晰那咱们再往深挖一点。你能不能把刚才提到的去重表、Redis幂等、版本号更新各写一段核心代码出来 不用长篇大论就把关键逻辑和你的技术亮点表达清楚就行。然后再帮我总结一下这些方案的技术难点以及你是怎么解决的。” ‍ 核心代码展示含技术亮点1. 数据库去重表 业务事务原子性ServicepublicclassPaymentCallbackService{AutowiredprivateOrderDaoorderDao;AutowiredprivateDedupDaodedupDao;/** * 处理支付回调幂等保证 * param callbackMsg 包含 outTradeNo, status 等 */Transactional(rollbackForException.class)// 亮点1整体事务publicvoidhandlePaymentCallback(CallbackMsgcallbackMsg){StringoutTradeNocallbackMsg.getOutTradeNo();// 亮点2去重表插入唯一索引自动检测重复intinserteddedupDao.insertIfAbsent(outTradeNo);if(inserted0){// 重复消息直接返回幂等跳过logger.info(重复消息outTradeNo{},outTradeNo);return;}// 第一次处理执行业务intaffectedorderDao.updateOrderStatus(outTradeNo,OrderStatus.PAID.getCode(),OrderStatus.WAIT_PAY.getCode());if(affected!1){thrownewBizException(订单状态异常无法更新);}}}// DedupDao.javaMapperpublicinterfaceDedupDao{// INSERT IGNORE 或 ON DUPLICATE KEY 保证幂等Insert(INSERT IGNORE INTO payment_dedup(out_trade_no, create_time) VALUES(#{outTradeNo}, NOW()))intinsertIfAbsent(Param(outTradeNo)StringoutTradeNo);} 技术亮点解析去重表插入和业务更新在同一个Transactional内要么都成功要么都回滚防止标记残留。使用INSERT IGNORE而非SELECT INSERT避免并发竞争数据库唯一索引本身保证了原子性。返回值判断0表示唯一键冲突即重复消息业务不做处理。2. Redis 分布式幂等高并发版ComponentpublicclassRedisDedupService{AutowiredprivateStringRedisTemplateredisTemplate;// 亮点1SET key value NX EX 原子操作publicbooleantryDedup(StringbizId,intexpireSeconds){Stringkeymsg:dedup:bizId;BooleanresultredisTemplate.opsForValue().setIfAbsent(key,1,Duration.ofSeconds(expireSeconds));returnBoolean.TRUE.equals(result);}// 亮点2业务完成后不主动删除靠TTL过期// 防止删除早于重复消息到达的窗口期}// 使用示例ServicepublicclassOrderMessageListener{AutowiredprivateRedisDedupServicededupService;publicvoidonMessage(OrderEventevent){StringbizIdevent.getOrderId()event.getEventType();// 去重窗口30分钟大于消息重试周期if(!dedupService.tryDedup(bizId,30*60)){return;// 重复消费丢弃}// 执行业务...}} 技术亮点解析setIfAbsent天然原子性实现“检查-设置”一步完成无并发漏洞。关键思想只设置 key不删除 key依赖 TTL 窗口过期自动清除完美覆盖重复消息的到达时间窗口。窗口时间设定需要大于 MQ 的重试周期比如 Kafka 的delivery.timeout.ms避免“过期了但重复消息还在路上”。3. 乐观锁版本号更新状态机式幂等MapperpublicinterfaceOrderDao{// 亮点带版本号的更新受影响行数自动判重Update(UPDATE orders SET status #{newStatus}, version version 1 WHERE order_id #{orderId} AND version #{oldVersion})intupdateStatusWithVersion(Param(orderId)StringorderId,Param(newStatus)StringnewStatus,Param(oldVersion)intoldVersion);}// 业务层调用TransactionalpublicvoiddispatchOrder(StringorderId,intcurrentVersion){introwsorderDao.updateStatusWithVersion(orderId,OrderStatus.SHIPPED.getCode(),currentVersion);if(rows0){// 重复消息或版本冲突幂等忽略thrownewDedupSkipException(订单状态已变更跳过重复操作);}// 后续发货逻辑...} 技术亮点解析利用version字段实现乐观锁同时天然具备幂等特性同一条消息再次更新时版本号已经不匹配affected rows 0。相比WHERE status ?可以防止并发更新时被覆盖两个“支付成功”消息同时到达版本号让第二个直接失效。适合状态流转场景比如“支付→发货→确认”每次状态变化版本号递增。 技术难点与解决方案总结在实际项目中这几个问题最容易踩坑我也梳理了对应的解法技术难点问题表现解决方案事务边界不清去重标记和业务操作不在同一事务造成标记残留或数据不一致利用数据库事务或 TCC 模式标记与业务必须原子化见代码1并发竞态两个相同消息同时消费都判断“未处理”导致重复执行① 数据库唯一索引约束② RedisSETNX原子性③ 版本号比较见代码2/3标记清理时机过早删除幂等标记导致窗口期外的重复消息再次生效不主动删除标记用过期时间自然淘汰Redis TTL数据库去重表可定期归档但不删除MQ 消息 ID 不可靠MQ 重发时可能生成新的messageId导致基于此的幂等失效强制使用业务唯一标识订单号、流水号绝不依赖 MQ 的消息 ID缓存穿透与幂等冲突先删缓存再执行业务重复消息发现缓存缺失重新走业务流程永远先设标记再删缓存或者反过来缓存仅在业务成功后更新标记优先级最高高并发下的性能瓶颈数据库唯一索引成为写入热点影响 TPS引入 Redis 做一级幂等热点数据数据库做二级兜底冷数据分流与降级策略重复消息的不可预知性可能数天后才来重复消息导致去重窗口难以设定根据业务容忍度设定足够长的窗口如支付回调设置7天数据库持久化标记能支持长尾去重回滚与补偿业务执行一半失败标记已写入后续重试无法恢复引入业务状态检查点标记仅表示“消息已接收”业务是否成功由独立的状态机或对账系统弥补面试官我往后靠了靠点了点头“这套代码和难点拆解把理论和落地结合得很好。尤其是标记不删除、业务ID去重、事务原子性这几个点都是高级工程师的分水岭。能把这些细节说清楚说明你不是光背八股是真在项目里用过。这个场景题我可以给你打个高分。” 整理了面试真题、每日技术知识点、系统学习路线√X搜「Rain的Java大神之路」每天拆一个知识点陪你悄悄变强有空来坐坐。
返回列表