列:異步解耦與業(yè)務(wù)削峰)
RabbitMQ消息隊(duì)列異步解耦與業(yè)務(wù)削峰同步調(diào)用就像你打電話(huà)等對(duì)方接——對(duì)方不接你就一直卡著異步消息就像發(fā)微信——發(fā)完該干嘛干嘛對(duì)方有空了自然回你。一、消息隊(duì)列解決了什么問(wèn)題在單體架構(gòu)時(shí)代所有功能揉在一個(gè)項(xiàng)目里方法之間直接調(diào)用簡(jiǎn)單粗暴。但一旦系統(tǒng)變大問(wèn)題就來(lái)了異步處理用戶(hù)注冊(cè)后要發(fā)郵件、發(fā)短信、發(fā)優(yōu)惠券……同步調(diào)用的話(huà)用戶(hù)得等半天體驗(yàn)極差。丟到消息隊(duì)列里注冊(cè)接口秒回后續(xù)操作慢慢消費(fèi)。應(yīng)用解耦訂單系統(tǒng)直接調(diào)用庫(kù)存系統(tǒng)庫(kù)存掛了訂單也跟著掛。中間加個(gè)隊(duì)列訂單只管發(fā)消息庫(kù)存恢復(fù)了繼續(xù)消費(fèi)即可。流量削峰秒殺場(chǎng)景瞬間涌入10萬(wàn)請(qǐng)求數(shù)據(jù)庫(kù)直接被干趴。隊(duì)列做個(gè)緩沖消費(fèi)者按自己的節(jié)奏處理系統(tǒng)穩(wěn)如老狗。日志收集分布式系統(tǒng)中各服務(wù)把日志推到隊(duì)列由統(tǒng)一的日志服務(wù)消費(fèi)存儲(chǔ)EFK/ELK的經(jīng)典套路。二、RabbitMQ核心概念RabbitMQ的消息流轉(zhuǎn)模型如下Producer → Exchange → (Binding) → Queue → Consumer 生產(chǎn)者 交換機(jī) 綁定 隊(duì)列 消費(fèi)者Producer生產(chǎn)者產(chǎn)生消息的應(yīng)用程序Exchange交換機(jī)接收生產(chǎn)者發(fā)送的消息根據(jù)路由規(guī)則分發(fā)到隊(duì)列Queue隊(duì)列存放消息的緩沖區(qū)消息在這里排隊(duì)等消費(fèi)Binding綁定交換機(jī)和隊(duì)列之間的關(guān)聯(lián)關(guān)系附帶路由鍵Consumer消費(fèi)者從隊(duì)列中獲取消息并處理的應(yīng)用程序三、交換機(jī)四種類(lèi)型RabbitMQ提供了四種Exchange類(lèi)型理解清楚就知道消息怎么路由了。3.1 Direct直連最簡(jiǎn)單的模式消息的路由鍵routing key和綁定的鍵完全匹配消息才會(huì)被投遞到對(duì)應(yīng)隊(duì)列。routing key order.create → 只匹配綁定 order.create 的隊(duì)列3.2 Fanout扇出廣播模式忽略路由鍵消息被投遞到與該交換機(jī)綁定的所有隊(duì)列。適合廣播通知場(chǎng)景。3.3 Topic主題支持通配符匹配靈活性最高*匹配一個(gè)單詞#匹配零個(gè)或多個(gè)單詞綁定鍵 order.* → 匹配 order.create、order.cancel不匹配 order.create.detail 綁定鍵 order.# → 匹配 order.create、order.create.detail 全都匹配3.4 Headers頭部不靠路由鍵而是根據(jù)消息頭headers中的鍵值對(duì)匹配。用的少了解即可。四、SpringBoot整合RabbitMQ4.1 引入依賴(lài)dependencygroupIdorg.springframework.boot/groupIdartifactIdspring-boot-starter-amqp/artifactId/dependency4.2 yml配置spring:rabbitmq:host:127.0.0.1port:5672username:guestpassword:guest# 消息確認(rèn)機(jī)制publisher-confirm-type:correlated# 發(fā)布確認(rèn)publisher-returns:true# 消息返回listener:simple:acknowledge-mode:manual# 手動(dòng)ACKprefetch:1# 每次拉取消息數(shù)4.3 隊(duì)列與交換機(jī)配置ConfigurationpublicclassRabbitMQConfig{// 隊(duì)列名稱(chēng)publicstaticfinalStringEMAIL_QUEUEemail.queue;publicstaticfinalStringSMS_QUEUEsms.queue;publicstaticfinalStringORDER_EXCHANGEorder.exchange;publicstaticfinalStringORDER_ROUTING_KEYorder.notify;BeanpublicDirectExchangeorderExchange(){returnnewDirectExchange(ORDER_EXCHANGE,true,false);}BeanpublicQueueemailQueue(){returnnewQueue(EMAIL_QUEUE,true);}BeanpublicQueuesmsQueue(){returnnewQueue(SMS_QUEUE,true);}BeanpublicBindingemailBinding(QueueemailQueue,DirectExchangeorderExchange){returnBindingBuilder.bind(emailQueue).to(orderExchange).with(ORDER_ROUTING_KEY);}BeanpublicBindingsmsBinding(QueuesmsQueue,DirectExchangeorderExchange){returnBindingBuilder.bind(smsQueue).to(orderExchange).with(ORDER_ROUTING_KEY);}}五、發(fā)送消息RabbitTemplateServicepublicclassOrderService{AutowiredprivateRabbitTemplaterabbitTemplate;publicvoidcreateOrder(OrderDTOorderDTO){// 1. 保存訂單數(shù)據(jù)庫(kù)操作省略// ...// 2. 異步發(fā)送通知消息StringmsgJSON.toJSONString(orderDTO);rabbitTemplate.convertAndSend(RabbitMQConfig.ORDER_EXCHANGE,RabbitMQConfig.ORDER_ROUTING_KEY,msg);// 3. 直接返回不等郵件/短信發(fā)送完成return;}}六、接收消息RabbitListenerComponentpublicclassEmailConsumer{RabbitListener(queuesRabbitMQConfig.EMAIL_QUEUE)RabbitHandlerpublicvoidreceive(Stringmessage,Channelchannel,MessagemessageObj)throwsIOException{longdeliveryTagmessageObj.getMessageProperties().getDeliveryTag();try{OrderDTOorderJSON.parseObject(message,OrderDTO.class);// 發(fā)送郵件邏輯System.out.println(發(fā)送郵件到order.getEmail());// 手動(dòng)確認(rèn)channel.basicAck(deliveryTag,false);}catch(Exceptione){// 消費(fèi)失敗拒絕并重新入隊(duì)channel.basicNack(deliveryTag,false,true);}}}短信消費(fèi)者結(jié)構(gòu)同理監(jiān)聽(tīng)SMS_QUEUE即可。一個(gè)交換機(jī)綁定了兩個(gè)隊(duì)列同一條消息會(huì)同時(shí)投遞到郵件隊(duì)列和短信隊(duì)列實(shí)現(xiàn)并行處理。七、消息可靠性保障消息從生產(chǎn)到消費(fèi)要經(jīng)過(guò)多個(gè)環(huán)節(jié)任何一個(gè)環(huán)節(jié)都可能丟消息。7.1 生產(chǎn)者確認(rèn)機(jī)制publisher-confirm-type:correlated# 異步確認(rèn)性能好rabbitTemplate.setConfirmCallback((correlationData,ack,cause)-{if(!ack){System.err.println(消息未到達(dá)Exchange原因cause);// 記錄日志重發(fā)等處理}});7.2 消費(fèi)者手動(dòng)ACK默認(rèn)是自動(dòng)確認(rèn)auto消息一拿到就標(biāo)記消費(fèi)成功但如果業(yè)務(wù)代碼報(bào)異常消息就丟了。改為手動(dòng)確認(rèn)manual業(yè)務(wù)成功后調(diào)basicAck失敗調(diào)basicNack。八、死信隊(duì)列消息變成死信的三種情況消息被消費(fèi)者rejectbasicReject/basicNack且不重新入隊(duì)消息TTL過(guò)期隊(duì)列或消息設(shè)置了過(guò)期時(shí)間隊(duì)列達(dá)到最大長(zhǎng)度新消息被擠出去死信隊(duì)列的配置思路給正常隊(duì)列綁定一個(gè)死信交換機(jī)DLX消息變成死信后自動(dòng)轉(zhuǎn)發(fā)到DLX再由DLX路由到死信隊(duì)列。BeanpublicQueuenormalQueue(){MapString,ObjectargsnewHashMap();args.put(x-message-ttl,60000);// 消息60秒過(guò)期args.put(x-dead-letter-exchange,dlx.exchange);args.put(x-dead-letter-routing-key,dlx.routing.key);returnnewQueue(normal.queue,true,false,false,args);}死信隊(duì)列常用于延遲任務(wù)消息過(guò)期→死信→消費(fèi)、失敗消息重試、訂單超時(shí)取消等場(chǎng)景。九、常見(jiàn)問(wèn)題與解決方案9.1 消息重復(fù)消費(fèi)冪等性網(wǎng)絡(luò)抖動(dòng)導(dǎo)致ACK沒(méi)及時(shí)到達(dá)RabbitMQ會(huì)重投消息消費(fèi)者就重復(fù)處理了。解決方案業(yè)務(wù)唯一鍵校驗(yàn)消費(fèi)前查數(shù)據(jù)庫(kù)/Redis已處理則直接ACK跳過(guò)樂(lè)觀鎖update語(yǔ)句加where status 0條件Redis分布式鎖setnx保證同一消息只處理一次publicvoidreceive(Stringmessage){StringmsgIdextractMsgId(message);// Redis標(biāo)記已處理則跳過(guò)BooleanisNewredisTemplate.opsForValue().setIfAbsent(msg:processed:msgId,1,24,TimeUnit.HOURS);if(Boolean.FALSE.equals(isNew)){return;// 已處理過(guò)}// 正常消費(fèi)邏輯}9.2 消息積壓處理消費(fèi)速度跟不上生產(chǎn)速度隊(duì)列堆積越來(lái)越多的消息。應(yīng)對(duì)策略臨時(shí)擴(kuò)容消費(fèi)者增加消費(fèi)者實(shí)例數(shù)量批量消費(fèi)一個(gè)消費(fèi)者一次拉取多條消息處理消息轉(zhuǎn)存緊急將積壓消息轉(zhuǎn)存到另一個(gè)隊(duì)列后續(xù)慢慢消費(fèi)根因排查消費(fèi)者是不是有慢查詢(xún)是不是依賴(lài)的外部服務(wù)超時(shí)了RabbitMQ用好了就是系統(tǒng)穩(wěn)定性的護(hù)城河用不好就是給自己挖坑。把可靠性保障和冪等性設(shè)計(jì)到位消息隊(duì)列才能真正發(fā)揮價(jià)值。