rabbitmq消息ACK確認(rèn)機(jī)制及發(fā)送失敗處理方式
rabbitmq為確保消息發(fā)送和接收成功,采用ack機(jī)制。
(1)生產(chǎn)者producter發(fā)送消息到mq時(shí),mq會發(fā)送ack給producter告知消息是否投遞成功;
(2)消費(fèi)者consumer接收處理消息后,consumer會發(fā)送ack給mq告知消息是否處理成功;
通過ack機(jī)制,確保消息能夠被producter成功發(fā)送和consumer成功接收處理,保證消息不丟失。
1、消息發(fā)送
rabbitmq消息發(fā)送分為兩個(gè)階段:
(1)producter將消息發(fā)送到broker,即發(fā)送到exchage交換機(jī);
(2)消息通過交換機(jī)exchange被路由到隊(duì)列queue;
消息只有被正確投遞到隊(duì)列queue中,才算發(fā)送成功。
消息發(fā)送代碼:
public boolean send(String queueName, String json, String msgId){
Message message = MessageBuilder.withBody(json.getBytes()).setCorrelationId(msgId).build();
message.getMessageProperties().setDeliveryMode(MessageDeliveryMode.PERSISTENT);//設(shè)置消息持久化
CorrelationDataExt correlationData = new CorrelationDataExt();
correlationData.setId(msgId);
correlationData.setData(json);
rabbitTemplate.setEncoding("UTF-8");
rabbitTemplate.setMandatory(true);//設(shè)置手工ack確認(rèn)
rabbitTemplate.setConfirmCallback(this);//ack回調(diào)
rabbitTemplate.setReturnCallback(this);//回退回調(diào)
rabbitTemplate.convertAndSend(queueName, message, correlationData);
return true;
}
在消息發(fā)送之前,我們要設(shè)置ack機(jī)制相關(guān)參數(shù):
- setMandatory:設(shè)置手工確認(rèn)ack;
- setConfirmCallback:設(shè)置消息發(fā)送到exchange結(jié)果回調(diào);
- setReturnCallback:設(shè)置消息投遞到queue失敗回退時(shí)回調(diào);
通過上述兩個(gè)回調(diào)方法,我們能夠?qū)Πl(fā)送失敗的消息進(jìn)行重發(fā)處理,確保消息不丟失。
2、消息發(fā)送失敗
根據(jù)rabbitmq發(fā)送過程,消息發(fā)送失敗的有三種情況會出現(xiàn):
(1)producter連接mq失敗,消息沒有發(fā)送到mq
(2)producter連接mq成功,但是發(fā)送到exchange失敗
(3)消息發(fā)送到exchange成功,但是路由到queue失??;
3、發(fā)送失敗處理
(1)producter連接mq失敗,消息沒有發(fā)送到mq
這種情況下,在發(fā)送消息時(shí)可以通過捕捉AmqpException異常,將消息保存db中后續(xù)進(jìn)行重發(fā)處理。
try{
rabbitTemplate.convertAndSend(queueName, message, correlationData);
}catch (Exception e){
logger.error("連接MQ失敗", e);
//todo 存儲到db中進(jìn)行重發(fā)
}
(2)producter連接mq成功,但是發(fā)送到exchange失敗
通過實(shí)現(xiàn)ConfirmCallback接口,對發(fā)送結(jié)果進(jìn)行處理。
@Override
public void confirm(CorrelationData correlationData, boolean ack, String cause) {
String msgId = correlationData.getId();
if(ack){
//發(fā)送成功
logger.debug("ack,消息投遞到exchange成功,msgId:{}",msgId);
}else{
//發(fā)送失敗,重試
logger.error("ack,消息投遞exchange失敗,msgId:{},原因{}" ,msgId, cause);
}
}
confirm方法有3個(gè)參數(shù),correlationData是消息發(fā)送時(shí)攜帶的數(shù)據(jù)對象,ack消息是否成功發(fā)送到exchange,cause是發(fā)送失敗時(shí)的原因。
通過ack我們可以判斷發(fā)送到exchange是否成功,如果ack=false,則我們進(jìn)行失敗處理。
但是這里存在一個(gè)問題,correlationData里面只有一個(gè)id屬性,沒有關(guān)于消息內(nèi)容的屬性,對于數(shù)據(jù)失敗處理非常不方便。
為解決此問題,我們可以自定義一個(gè)CorrelationData擴(kuò)展對象,繼承CorrelationData,并添加自己想要保存數(shù)據(jù)的屬性,在消息發(fā)送時(shí),攜帶相關(guān)數(shù)據(jù)在該對象上即可。
自定義CorrelationData對象:
/**
* CorrelationData的自定義實(shí)現(xiàn),用于拿到消息內(nèi)容
*/
public class CorrelationDataExt extends CorrelationData {
//數(shù)據(jù)
private volatile Object data;
//隊(duì)列
private String queueName;
public Object getData() {
return data;
}
public void setData(Object data) {
this.data = data;
}
public String getQueueName() {
return queueName;
}
public void setQueueName(String queueName) {
this.queueName = queueName;
}
}
重寫發(fā)送方法,使用CorrelationDataExt對象攜帶數(shù)據(jù):
public boolean send(String queueName, String json, String msgId){
Message message = MessageBuilder.withBody(json.getBytes()).setCorrelationId(msgId).build();
message.getMessageProperties().setDeliveryMode(MessageDeliveryMode.PERSISTENT);//設(shè)置消息持久化
//使用自定義的數(shù)據(jù)對象
CorrelationDataExt correlationData = new CorrelationDataExt();
correlationData.setId(msgId);
correlationData.setData(json);
correlationData.setQueueName(queueName);
rabbitTemplate.setEncoding("UTF-8");
rabbitTemplate.setMandatory(true);//設(shè)置手工ack確認(rèn)
rabbitTemplate.setConfirmCallback(this);//設(shè)置發(fā)送成功回調(diào)
rabbitTemplate.setReturnCallback(this);//設(shè)置消息回退回調(diào)
try{
rabbitTemplate.convertAndSend(queueName, message, correlationData);//使用amqp default exchange direct
}catch (Exception e){
logger.error("MQ連接失敗,請聯(lián)系管理員處理!!!!");
//保存到db重發(fā)
saveToDB(msgId, json, queueName, "90");
}
return true;
}
重寫confirm方法,對CorrelationData進(jìn)行處理:
@Override
public void confirm(CorrelationData correlationData, boolean ack, String cause) {
String msgId = correlationData.getId();
if(ack){
//發(fā)送成功
logger.debug("ack,消息投遞到exchange成功,msgId:{}",msgId);
}else{
//發(fā)送失敗,重試
logger.error("ack,消息投遞exchange失敗,msgId:{},原因{}" ,msgId, cause);
if(correlationData instanceof CorrelationDataExt){
CorrelationDataExt correlationDataExt = (CorrelationDataExt) correlationData;
String message = (String) correlationDataExt.getData();
String queueName = ((CorrelationDataExt) correlationData).getQueueName();
saveToDB(msgId, message, queueName, "91");
}else{
logger.info("correlationData對象不包含數(shù)據(jù)");
}
}
}
(3)消息發(fā)送到exchange成功,但是路由到queue失敗
通過實(shí)現(xiàn)ReturnCallback接口,對回退消息進(jìn)行重發(fā)處理。
@Override
public void returnedMessage(Message message, int replyCode, String replyText, String exchange, String routingKey) {
logger.error("消息發(fā)送失敗-消息回退,應(yīng)答碼:{},原因:{},交換機(jī):{},路由鍵:{}", replyCode, replyText, exchange, routingKey);
String msgId = message.getMessageProperties().getCorrelationId();
String data = new String(message.getBody());
saveToDB(msgId, data, routingKey, "92");
}
關(guān)于對失敗消息的處理,我這里是統(tǒng)一保存到DB中,后續(xù)通過定時(shí)任務(wù)進(jìn)行重發(fā)處理的。
通過以上3個(gè)方面對失敗消息的處理,可以確保消息能夠成功發(fā)送到mq,確保不丟失。
總結(jié)
以上為個(gè)人經(jīng)驗(yàn),希望能給大家一個(gè)參考,也希望大家多多支持腳本之家。
相關(guān)文章
Javaweb mybatis接口開發(fā)實(shí)現(xiàn)過程詳解
這篇文章主要介紹了Javaweb mybatis接口開發(fā)實(shí)現(xiàn)過程詳解,文中通過示例代碼介紹的非常詳細(xì),對大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價(jià)值,需要的朋友可以參考下2020-07-07
解析Spring框架中的XmlBeanDefinitionStoreException異常情況
這篇文章主要介紹了解析Spring框架中的XmlBeanDefinitionStoreException異常情況,具有很好的參考價(jià)值,希望對大家有所幫助,如有錯(cuò)誤或未考慮完全的地方,望不吝賜教2024-04-04
全面解析Java中常見Exception異常的錯(cuò)誤排查與代碼修正
這篇文章主要為大家詳細(xì)介紹了Java開發(fā)中最常見的異常類型及其處理方法,包括運(yùn)行時(shí)異常(如NullPointerException,ArrayIndexOutOfBoundsException)和受檢異常,下面小編就和大家詳細(xì)介紹一下吧2026-03-03
Java 必知必會的 URL 和 URLConnection使用
這篇文章主要介紹了Java 必知必會的 URL 和 URLConnection使用,文中通過示例代碼介紹的非常詳細(xì),對大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價(jià)值,需要的朋友們下面隨著小編來一起學(xué)習(xí)學(xué)習(xí)吧2019-10-10
SpringBoot和Vue.js項(xiàng)目中實(shí)現(xiàn)文件壓縮下載功能的具體方案
在Spring Boot和Vue.js項(xiàng)目中實(shí)現(xiàn)文件壓縮下載功能,主要思路是后端負(fù)責(zé)將多個(gè)文件壓縮成ZIP包,前端負(fù)責(zé)觸發(fā)下載并處理文件流,以下是具體的實(shí)現(xiàn)方案,需要的朋友可以參考下2025-09-09
idea中創(chuàng)建新類時(shí)自動添加注釋的實(shí)現(xiàn)
在每次使用idea創(chuàng)建一個(gè)新類時(shí),過了一段時(shí)間發(fā)現(xiàn)看不懂這個(gè)類是用來干嘛的,為了解決這個(gè)問題,我們可以設(shè)置在創(chuàng)建一個(gè)新類時(shí)自動添加注釋,幫助我們理解這個(gè)類的用處,本文主要介紹了在idea中創(chuàng)建新類時(shí)自動添加注釋的實(shí)現(xiàn),感興趣的可以了解一下2025-03-03
詳解Servlet3.0新特性(從注解配置到websocket編程)
Servlet3.0的出現(xiàn)是servlet史上最大的變革,其中的許多新特性大大的簡化了web應(yīng)用的開發(fā),為廣大勞苦的程序員減輕了壓力,提高了web開發(fā)的效率。2017-04-04
Java設(shè)計(jì)模式之狀態(tài)模式State Pattern詳解
這篇文章主要介紹了Java設(shè)計(jì)模式之狀態(tài)模式State Pattern,狀態(tài)模式允許一個(gè)對象在其內(nèi)部狀態(tài)改變的時(shí)候改變其行為。這個(gè)對象看上去就像是改變了它的類一樣2022-11-11

