首页 / 资讯中心 / 文章详情

面试官:“如何保证RabbitMQ消息不丢失?”全链路排查+实战代码,一篇讲透

面试官:“如何保证RabbitMQ消息不丢失?”全链路排查+实战代码,一篇讲透 ★ FEATURED ARTICLE
很多面试Java后端岗位的人都会遇到这道经典面试题今天我也来讲讲我对这道题的理解。其实面试官并不单单想听你背八股文说出“持久化”、“手动ACK”这些关键词我在面试中直接被问你实际项目中是怎么做的配置怎么配代码怎么写的等这些实际问题。面试官是想知道你是否真正理解消息从生产到消费的完整链路你是否在实际项目中踩过坑、解决过问题你是否能给出系统性的解决方案而不是零散的知识点所以回答这个问题的关键不是简单的堆砌概念而是要展示你的全链路思维。一条消息从生产者发出到被消费者消费需要经历三段路生产者 → 交换机 → 队列 → 消费者每一段都有可能丢失消息对应三个核心问题生产者 → 交换机消息发出去了但是交换机不存在或者路由规则配置错误消息直接丢了而生产者这边却不知道。交换机 → 队列 →Broker存储消息到了Broker但是Broker宕机重启内存里的消息全没了。队列 → 消费者消费者拉取到了消息还没处理完就崩了但默认自动ACK消息已经被删除了。第一关生产者端——让每条消息都有回执如果面试官追问“生产者怎么知道消息有没有发送成功?这时候你要抛出两个核心机制Publisher Confirm和Return Callback。Publisher Confirm消息到达交换机后Broker会给生产者一个确认acktrue或拒绝ackfalse。Return Callback消息到了交换机但是没有匹配到任何队列时Broker会把消息退回给生产者。两者配合就能确保生产者端不会丢失消息。在代码层面先进行配置spring: rabbitmq: publisher-confirm-type: correlated # 开启Confirm4 publisher-returns: true # 开启Return然后配置回调ConfigurationSlf4jpublicclassRabbitConfig{AutowiredprivateRabbitTemplaterabbitTemplate;PostConstructpublicvoidinit(){rabbitTemplate.setConfirmCallback((correlationData,ack,cause)-{if(ack){log.info(消息已经到达交换机,id{},correlationData.getId());//更新本地消息表状态已送达}else{log.info(消息未到达交换机,cause{},cause);//触发重试或告警}});rabbitTemplate.setReturnsCallback(returnedMessage-{log.error(消息未路由到队列, exchange{}, routingKey{},returnedMessage.getExchange(),returnedMessage.getRoutingKey());// 记录到数据库后续人工处理});}}发送消息时别忘了设置 mandatorytrue否则Return回调不会触发rabbitTemplate.setMandatory(true);rabbitTemplate.convertAndSend(order-exchange,order.create,messageBody,newCorrelationData(UUID.randomUUID().toString()));第二关Broker端——三重持久化一个都不能少面试官追问“如果Broker挂了怎么办”这时候可以告诉面试官Broker要做到消息不丢失需要三样东西同时持久化要素配置交换机durable true队列durable true消息deliveryMode 2BeanpublicQueueorderQueue(){returnQueueBuilder.durable(order-queue).withArgument(x-dead-letter-exchange,dlx-exchange).withArgument(x-dead-letter-routing-key,dlx-routing-key).build();}第三关消费者端——手动ACK 死信队列面试官追问“如果消费者处理到一半突然挂了怎么办”这里是最容易丢失消息的环节也是开发最容易踩坑的地方。默认情况下消费者是自动ACK的消息一拉取到还没处理完Broker就认为已消费并删除消息如果此时消费者突然挂了那么消息就会永久丢失。解决方案关闭自动ACK手动确认pring:rabbitmq:listener:simple:acknowledge-mode:manual# 关闭自动ACK6prefetch:1# 每次只拉一条防止消息堆积消费者代码RabbitListener(queuesorder-queue)publicvoidhandleOrder(Messagemessage,Channelchannel)throwsIOException{longdeliveryTagmessage.getMessageProperties().getDeliveryTag();try{StringbodynewString(message.getBody(),StandardCharsets.UTF_8);OrderorderJSON.parseObject(body,Order.class);// 幂等校验防止重复消费if(isAlreadyProcessed(order.getOrderId())){channel.basicAck(deliveryTag,false);return;}// 执行业务逻辑stockService.deductStock(order.getGoodsId(),order.getQuantity());// 业务成功手动ACKchannel.basicAck(deliveryTag,false);}catch(Exceptione){log.error(消费失败,e);// 重试次数控制IntegerretryCountmessage.getMessageProperties().getHeader(x-retry-count);intcurrent(retryCountnull)?0:retryCount;if(current3){// 未超限重新入队message.getMessageProperties().setHeader(x-retry-count,current1);channel.basicNack(deliveryTag,false,true);}else{// 超限拒绝消息进入死信队列log.error(重试耗尽转入死信队列);channel.basicNack(deliveryTag,false,false);}}}终极兜底本地消息表面试官如果继续追问“如果Confirm回调本身也丢了呢”这时候我们要说出终极方案——本地消息表核心思路很简单业务数据和消息记录在同一个本地事务中写入数据库。定时任务轮询“待发送”状态的消息调用MQ发送。发送成功后更新为“已发送。如果失败则进入重试超过重试次数标记为“失败”人工介入。Transactional(rollbackForException.class)publicvoidcreateOrderReliably(Orderorder){// 1. 业务入库orderMapper.insert(order);// 2. 消息记录入库同事务MessageLoglognewMessageLog();log.setMsgId(UUID.randomUUID().toString());log.setMsgBody(JSON.toJSONString(order));log.setStatus(PENDING);messageLogMapper.insert(log);}Scheduled(fixedDelay30000)publicvoidretryPendingMessages(){ListMessageLogpendingmessageLogMapper.selectByStatus(PENDING);for(MessageLogmsg:pending){if(msg.getRetryCount()5){messageLogMapper.updateStatus(msg.getMsgId(),FAILED);continue;} rabbitTemplate.convertAndSend(order-exchange,order.create,msg.getMsgBody(),newCorrelationData(msg.getMsgId()));messageLogMapper.incrementRetryCount(msg.getMsgId());}}本地消息表的作用在于把“发消息”这个动作从同步调用变成了“异步补偿”即使MQ短暂不可用消息也不会丢失。综上这个面试题可以作如下回答” 消息从生产到消费有三个环节可能丢失我的方案是全链路闭环生产者端 我开启了Publisher Confirm和Return Callback确保消息到达交换机未路由的消息也能被回收处理。对于核心业务还会配合本地消息表做最终兜底。Broker端 我确保交换机、队列、消息三者都做了持久化生产环境使用集群模式避免单点故障。消费者端 我关闭了自动ACK改为手动确认业务处理成功后才ACK。处理失败的消息会有限次重试重试耗尽后转入死信队列避免消息丢失和毒消息循环。同时业务层做好幂等处理防止重复消费。这套方案在我们项目中已经跑了很久核心业务消息零丢失。“【Java笔记 小李版】我们一起进步
阅读完成 · 觉得有帮助?
咨询建站