解決RocketMQ的冪等性問(wèn)題
造成重復(fù)消費(fèi)的原因
- 當(dāng)系統(tǒng)的調(diào)用鏈路比較長(zhǎng)的時(shí)候,比如系統(tǒng)A調(diào)用系統(tǒng)B,系統(tǒng)B再把消息發(fā)送到RocketMQ中,在系統(tǒng)A調(diào)用系統(tǒng)B的時(shí)候,如果系統(tǒng)B處理成功,但是遲遲沒(méi)有將調(diào)用成功的結(jié)果返回給系統(tǒng)A的時(shí)候,系統(tǒng)A就會(huì)嘗試重新發(fā)起請(qǐng)求給系統(tǒng)B,造成系統(tǒng)B重復(fù)處理,發(fā)起多條消息給RocketMQ造成重復(fù)消費(fèi);
- 在系統(tǒng)B發(fā)送給RocketMQ的時(shí)候,也有可能會(huì)發(fā)生和上面一樣的問(wèn)題,消息發(fā)送超時(shí),結(jié)果系統(tǒng)B重試,導(dǎo)致RocketMQ接收到了重讀消息;
- 當(dāng)RocketMQ成功接收到消息,并將消息交給消費(fèi)者處理,如果消費(fèi)者消費(fèi)完成后還沒(méi)來(lái)得及提交CONSUME_SUCCESS給RocketMQ,自己宕機(jī)或者重啟了,那么RocketMQ沒(méi)有接收到CONSUME_SUCCESS,就會(huì)認(rèn)為消費(fèi)失敗了,會(huì)重發(fā)消息給消費(fèi)者再次消費(fèi);

通過(guò)冪等性來(lái)保證,只要保證重復(fù)消息不對(duì)結(jié)果產(chǎn)生影響,就完美地解決這個(gè)問(wèn)題。
解決方法
生產(chǎn)者端
- RocketMQ支持消息查詢(xún)的功能,只要去RocketMQ查詢(xún)一下是否已經(jīng)發(fā)送過(guò)該條消息就可以了,不存在則發(fā)送,存在則不發(fā)送,也就是message.setKeys();
- 引入Redis,在發(fā)送消息到RocketMQ成功之后,向Redis中插入一條數(shù)據(jù),如果發(fā)送重試,則先去Redis查詢(xún)一個(gè)該條消息是否已經(jīng)發(fā)送過(guò)了,存在的話(huà)就不重復(fù)發(fā)送消息了;

缺點(diǎn)
方法一:RocketMQ消息查詢(xún)的性能不是特別好,如果在高并發(fā)的場(chǎng)景下,每條消息在發(fā)送到RocketMQ時(shí)都去查詢(xún)一下,可能會(huì)影響接口的性能;
方法二:在一些極端的場(chǎng)景下,Redis也無(wú)法保證消息發(fā)送成功之后,就一定能寫(xiě)入Redis成功,比如寫(xiě)入消息成功而Redis此時(shí)宕機(jī),那么再次查詢(xún)Redis判斷消息是否已經(jīng)發(fā)送過(guò),是無(wú)法得到正確結(jié)果的;
消費(fèi)者端
- 建立一個(gè)消息表,拿到這個(gè)消息做數(shù)據(jù)庫(kù)的insert操作。給這個(gè)消息做一個(gè)唯一主鍵(primary key)或者唯一約束,那么就算出現(xiàn)重復(fù)消費(fèi)的情況,就會(huì)導(dǎo)致主鍵沖突。
- 拿到這個(gè)消息做redis的set的操作.redis就是天然冪等性
代碼實(shí)現(xiàn)
方式一:
生產(chǎn)者
public class MQProducer {
public static void main(String[] args) throws MQClientException {
//創(chuàng)建生產(chǎn)者
DefaultMQProducer producer=new DefaultMQProducer("rmq-group");
//設(shè)置NameServer地址
producer.setNamesrvAddr("192.168.138.187:9876;192.168.138.188:9876");
//設(shè)置生產(chǎn)者實(shí)例名稱(chēng)
producer.setInstanceName("producer");
//啟動(dòng)生產(chǎn)者
producer.start();
try {
//發(fā)送消息
for (int i=0;i<1;i++){
Thread.sleep(1000); //每秒發(fā)送一次
//創(chuàng)建消息
Message msg = new Message("wn04", // topic 主題名稱(chēng)
"TagA", // tag 臨時(shí)值
("w-"+i).getBytes()// body 內(nèi)容
);
//消息的唯一標(biāo)識(shí)
msg.setKeys(System.currentTimeMillis() + "");
//發(fā)送消息
SendResult sendResult=producer.send(msg);
System.out.println(sendResult.toString());
}
} catch (Exception e) {
e.printStackTrace();
}
producer.shutdown();
}
}
消費(fèi)者端:
public class MQConsumer {
//保存標(biāo)識(shí)的集合
static private Map<String, String> logMap = new HashMap<>();
public static void main(String[] args) throws MQClientException {
//創(chuàng)建消費(fèi)者
DefaultMQPushConsumer consumer=new DefaultMQPushConsumer("rmq-group");
//設(shè)置NameServer地址
consumer.setNamesrvAddr("192.168.138.187:9876;192.168.138.188:9876");
//設(shè)置消費(fèi)者實(shí)例名稱(chēng)
consumer.setInstanceName("consumer");
//訂閱topic
consumer.subscribe("wn04","TagA");
//監(jiān)聽(tīng)消息
consumer.registerMessageListener(new MessageListenerConcurrently() {
@Override
public ConsumeConcurrentlyStatus consumeMessage(List<MessageExt> list, ConsumeConcurrentlyContext consumeConcurrentlyContext) {
String key = null;
String msgId = null;
try {
for (MessageExt msg : list) {
key = msg.getKeys();
//判斷集合當(dāng)中有沒(méi)有存在key,存在就不需要重試,不存在先存key再回來(lái)重試后消費(fèi)消息
if (logMap.containsKey(key)) {
// 無(wú)需繼續(xù)重試。
System.out.println("key:"+key+",無(wú)需重試...");
return ConsumeConcurrentlyStatus.CONSUME_SUCCESS;
}
msgId = msg.getMsgId();
System.out.println("key:" + key + ",msgid:" + msgId + "---" + new String(msg.getBody()));
//模擬異常
int i = 1 / 0;
}
} catch (Exception e) {
//e.printStackTrace();
//重試
return ConsumeConcurrentlyStatus.RECONSUME_LATER;
} finally {
//保存key
logMap.put(key, msgId);
}
return ConsumeConcurrentlyStatus.CONSUME_SUCCESS;
}
});
consumer.start();
System.out.println("Consumer Started...");
}
}
到此這篇關(guān)于解決RocketMQ的冪等性問(wèn)題的文章就介紹到這了,更多相關(guān)RocketMQ 冪等性?xún)?nèi)容請(qǐng)搜索腳本之家以前的文章或繼續(xù)瀏覽下面的相關(guān)文章希望大家以后多多支持腳本之家!
相關(guān)文章
Java構(gòu)建JDBC應(yīng)用程序的實(shí)例操作
在本篇文章里小編給大家整理了一篇關(guān)于Java構(gòu)建JDBC應(yīng)用程序的實(shí)例操作,有興趣的朋友們可以學(xué)習(xí)參考下。2021-03-03
關(guān)于Idea的Invalidate Caches/Restart使用
項(xiàng)目類(lèi)導(dǎo)入爆紅可能因Idea緩存異常導(dǎo)致Maven依賴(lài)識(shí)別失敗,解決方法為通過(guò)Invalidate Caches/Restart清除緩存,等待重新構(gòu)建索引后重新進(jìn)入項(xiàng)目2025-07-07
SpringBoot中@ComponentScan注解過(guò)濾排除不加載某個(gè)類(lèi)的3種方法
這篇文章主要給大家介紹了關(guān)于SpringBoot中@ComponentScan注解過(guò)濾排除不加載某個(gè)類(lèi)的3種方法,文中通過(guò)實(shí)例代碼介紹的非常詳細(xì),對(duì)大家學(xué)習(xí)或者使用SpringBoot具有一定的參考學(xué)習(xí)價(jià)值,需要的朋友可以參考下2023-07-07
詳解Java中ExceptionInInitializer錯(cuò)誤的解決方法
ExceptionInInitializerError 是 Java 中的未經(jīng)檢查的異常,它是 Error 類(lèi)的子類(lèi), 它屬于運(yùn)行時(shí)異常的類(lèi)別,下面我們就來(lái)看看它的具體解決方法吧2024-02-02
JSON反序列化Long變Integer或Double的問(wèn)題及解決
這篇文章主要介紹了JSON反序列化Long變Integer或Double的問(wèn)題及解決方案,具有很好的參考價(jià)值,希望對(duì)大家有所幫助。如有錯(cuò)誤或未考慮完全的地方,望不吝賜教2022-01-01
Spring?Boot網(wǎng)絡(luò)配置之server.address超詳細(xì)講解
在Spring Boot中,我們可以通過(guò)配置server.address屬性來(lái)指定應(yīng)用程序綁定的IP地址,下面這篇文章主要介紹了Spring?Boot網(wǎng)絡(luò)配置之server.address的相關(guān)資料,需要的朋友可以參考下2025-09-09
Java并發(fā)編程之阻塞隊(duì)列(BlockingQueue)詳解
這篇文章主要介紹了詳解Java阻塞隊(duì)列(BlockingQueue)的實(shí)現(xiàn)原理,阻塞隊(duì)列是Java util.concurrent包下重要的數(shù)據(jù)結(jié)構(gòu),有興趣的可以了解一下2021-09-09
eclipse配置tomcat10的詳細(xì)步驟總結(jié)
今天給大家?guī)?lái)的是關(guān)于Java的相關(guān)知識(shí),文章圍繞著eclipse配置tomcat10的詳細(xì)步驟展開(kāi),文中有非常詳細(xì)的介紹及圖文示例,需要的朋友可以參考下2021-06-06
詳解mybatis多對(duì)一關(guān)聯(lián)查詢(xún)的方式
這篇文章主要給大家介紹了關(guān)于mybatis多對(duì)一關(guān)聯(lián)查詢(xún)的相關(guān)資料,文中將關(guān)聯(lián)方式以及配置方式介紹的很詳細(xì),需要的朋友可以參考下2021-06-06

