Springboot使用RabbitMq延遲隊(duì)列和死信隊(duì)列詳解
前言
在最近的項(xiàng)目中,結(jié)合minio文件服務(wù)器的一些特性。
需要做一個(gè)分片上傳的功能:用戶(hù)上傳文件到md5的桶下,合并文件后刪除這個(gè)臨時(shí)桶。
會(huì)出現(xiàn)這樣一種情況,用戶(hù)上傳文件傳到一半就不再上傳了,那么如何去刪除,什么時(shí)候去刪除時(shí)需要解決問(wèn)題。
一、業(yè)務(wù)解決方案
1.quartz定時(shí)器
如果是單體項(xiàng)目,可以考慮使用quartz定時(shí)器。在創(chuàng)建桶的時(shí)候加入到定時(shí)任務(wù)里。
2.redis定時(shí)器
redis定時(shí)器需要修改配置文件,并且對(duì)redis進(jìn)行監(jiān)聽(tīng),在創(chuàng)建桶時(shí),設(shè)置過(guò)時(shí)時(shí)間,一旦時(shí)間超時(shí),可以對(duì)key進(jìn)行捕捉,最好對(duì)名字進(jìn)行規(guī)范設(shè)計(jì)以便于業(yè)務(wù)
3.mq消息隊(duì)列
使用延遲隊(duì)列和死信隊(duì)列進(jìn)行定時(shí)任務(wù)
這篇主要講解mq的方式解決問(wèn)題
二、RabbitMq延遲隊(duì)列
1.延遲隊(duì)列
延遲隊(duì)列也是一個(gè)普通的隊(duì)列,和普通的隊(duì)列相比,他多了幾個(gè)屬性,比如:
1)延遲的時(shí)間:表示隊(duì)列中消息的生命周期,在指定時(shí)間后,要么拋棄這個(gè)消息,要么投遞到死信隊(duì)列中
2)指定死信交換機(jī):如果不希望丟棄這個(gè)消息,那么可以將這個(gè)過(guò)期的消息丟到死信隊(duì)列中
定義一個(gè)延遲隊(duì)列:
//桶延遲隊(duì)列
@Bean(BUCKET_TTL_QUEUE)
public Queue bucketTtlQueue(){
Map<String,Object> deadParamsMap = new HashMap<>();
// 設(shè)置死信隊(duì)列的Exchange
deadParamsMap.put("x-dead-letter-exchange",BUCKET_DEAD_EXCHANGE);
//設(shè)置死信隊(duì)列的RouteKey
deadParamsMap.put("x-dead-letter-routing-key",BUCKET_DEAD_QUEUE);
// 設(shè)置對(duì)接過(guò)期時(shí)間"x-message-ttl"
deadParamsMap.put("x-message-ttl",60000*5);//5分鐘
// 設(shè)置對(duì)接可以存儲(chǔ)的最大消息數(shù)量
//deadParamsMap.put("x-max-length",10);
return new Queue(BUCKET_TTL_QUEUE,true,false,false,deadParamsMap);
}
延遲隊(duì)列交換機(jī)
如上所說(shuō),延遲隊(duì)列本就是一個(gè)普通的隊(duì)列,如果你想更細(xì)粒的對(duì)他進(jìn)行控制,那么需要綁定交換機(jī),如果不綁定交換機(jī),會(huì)綁定到默認(rèn)交換機(jī),在發(fā)送消息時(shí),交換機(jī)寫(xiě)""就行,默認(rèn)交換機(jī)為直連交換機(jī);
我這里指定了延遲隊(duì)列的交換機(jī),因?yàn)闆](méi)有做消息冪等性,所以采用直連交換機(jī)應(yīng)對(duì)在集群下消息只被消費(fèi)一次
//桶延遲交換機(jī)
@Bean(BUCKET_TTL_EXCHANGE)
public DirectExchange bucketTtlExchange() {
return new DirectExchange(BUCKET_TTL_EXCHANGE,true,false);
}
// 綁定
@Bean
public Binding bucketTtlBinding() {
return BindingBuilder.bind(bucketTtlQueue())
.to(bucketTtlExchange())
.with(BUCKET_TTL_QUEUE);
}
2.死信交換機(jī)
DLX也是一個(gè)正常的Exchange,和一般的Exchange沒(méi)有區(qū)別,它能在任何的隊(duì)列上被指定,實(shí)際上就是設(shè)置某個(gè)隊(duì)列的屬性。
當(dāng)這個(gè)隊(duì)列中有死信時(shí),RabbitMQ就會(huì)自動(dòng)的將這個(gè)消息重新發(fā)布到設(shè)置的Exchange上去,進(jìn)而被路由到另一個(gè)隊(duì)列
死信交換機(jī)也是普通交換機(jī),他只是你指定接收過(guò)期消息的交換機(jī)而已
/**
* 死信隊(duì)列
*
* @return
*/
@Bean(BUCKET_DEAD_QUEUE)
public Queue bucketDeadQueue() {
//屬性參數(shù) 隊(duì)列名稱(chēng) 是否持久化
return new Queue(BUCKET_DEAD_QUEUE, true);
}
/**
* 死信隊(duì)列交換機(jī)
*
* @return
*/
@Bean(BUCKET_DEAD_EXCHANGE)
public DirectExchange bucketDeadExchange() {
return new DirectExchange(BUCKET_DEAD_EXCHANGE,true,false);
}
/**
* 給死信隊(duì)列綁定交換機(jī)
*
* @return
*/
@Bean
public Binding bucketDeadBinding() {
return BindingBuilder.bind(bucketDeadQueue()).to(bucketDeadExchange()).with(BUCKET_DEAD_QUEUE);
}
3.監(jiān)聽(tīng)器
消息處理的邏輯,在消息過(guò)期后,送到死信交換機(jī)里,監(jiān)聽(tīng)器監(jiān)聽(tīng)到死信交換機(jī)的消息進(jìn)行刪除桶以及文件的業(yè)務(wù)邏輯處理
/**
* @description:死信隊(duì)列監(jiān)聽(tīng)器,用來(lái)刪除過(guò)期的桶
* @author manchao
* @date 2022/2/17 9:52
*/
@Configuration
public class BucketDeadConsumer {
@Autowired
private CachingConnectionFactory cachingConnectionFactory;
@Autowired
private MinioTemplate minioTemplate;
@Bean
public SimpleMessageListenerContainer BucketDeadListenerContainer() {
SimpleMessageListenerContainer container = new SimpleMessageListenerContainer(cachingConnectionFactory);
// 監(jiān)聽(tīng)隊(duì)列名
container.setQueueNames(MyMqConfig.BUCKET_DEAD_QUEUE);
// 當(dāng)前消費(fèi)者數(shù)量 開(kāi)啟幾個(gè)線(xiàn)程去處理數(shù)據(jù) 支持運(yùn)行時(shí)動(dòng)態(tài)修改
container.setConcurrentConsumers(5);
// 最大消費(fèi)者數(shù)量 , 消息堵塞太多的時(shí)候,會(huì)幫我自動(dòng)擴(kuò)展到我的最大消費(fèi)者數(shù)量
container.setMaxConcurrentConsumers(10);
// 是否重回隊(duì)列
container.setDefaultRequeueRejected(true);
// 手動(dòng)確認(rèn)
container.setAcknowledgeMode(AcknowledgeMode.MANUAL);
// 設(shè)置監(jiān)聽(tīng)器
container.setMessageListener(new ChannelAwareMessageListener(){
@Override
public void onMessage(Message message, Channel channel) throws IOException {
// 消息的唯一性ID deliveryTag:該消息的index 自增長(zhǎng)
long deliveryTag = message.getMessageProperties().getDeliveryTag();
byte[] messageBody = message.getBody();
String s = new String(messageBody);
System.out.println("消息: " + s);
System.out.println("消息來(lái)自: "+message.getMessageProperties().getConsumerQueue());
System.out.println("交換機(jī): "+message.getMessageProperties().getReceivedExchange());
//刪除桶
try {
minioTemplate.removeBuckets(s, "");
channel.basicAck(deliveryTag, false);
} catch (IOException e) {
e.printStackTrace();
channel.basicReject(deliveryTag, false);
}
}
});
return container;
}
}
注:我設(shè)置了消息發(fā)送和確認(rèn)的回調(diào)函數(shù),為什么沒(méi)有觸發(fā)這個(gè)函數(shù)?因?yàn)槲沂菑暮笈_(tái)管理頁(yè)面發(fā)的消息,沒(méi)有通過(guò)rabbitteplate進(jìn)行發(fā)送,不會(huì)是這個(gè)原因吧!
總結(jié)
業(yè)務(wù)的解決方法有太多種了,找到一個(gè)高可用以及簡(jiǎn)便的方法才是解決問(wèn)題的關(guān)鍵
以上為個(gè)人經(jīng)驗(yàn),希望能給大家一個(gè)參考,也希望大家多多支持腳本之家。
- Springboot使用Rabbitmq的延時(shí)隊(duì)列+死信隊(duì)列實(shí)現(xiàn)消息延期消費(fèi)
- springboot整合RabbitMQ中死信隊(duì)列的實(shí)現(xiàn)
- SpringBoot整合RabbitMQ實(shí)現(xiàn)延遲隊(duì)列和死信隊(duì)列
- springboot中RabbitMQ死信隊(duì)列的實(shí)現(xiàn)示例
- Springboot結(jié)合rabbitmq實(shí)現(xiàn)的死信隊(duì)列
- 關(guān)于SpringBoot整合RabbitMQ實(shí)現(xiàn)死信隊(duì)列
- SpringBoot+RabbitMQ?實(shí)現(xiàn)死信隊(duì)列的示例
- SpringBoot整合RabbitMQ處理死信隊(duì)列和延遲隊(duì)列
- Springboot集成RabbitMQ死信隊(duì)列的實(shí)現(xiàn)
- SpringBoot集成RabbitMQ的方法(死信隊(duì)列)
- SpringBoot4.0整合RabbitMQ死信隊(duì)列詳解
相關(guān)文章
Java修改eclipse中web項(xiàng)目的server部署路徑問(wèn)題
這篇文章主要介紹了Java修改eclipse中web項(xiàng)目的server部署路徑,本文給大家介紹的非常詳細(xì),對(duì)大家的學(xué)習(xí)或工作具有一定的參考借鑒價(jià)值,需要的朋友可以參考下2020-11-11
如何使用spring-ws發(fā)布webservice服務(wù)
文章介紹了如何使用Spring-WS發(fā)布Web服務(wù),包括添加依賴(lài)、創(chuàng)建XSD文件、生成JAXB實(shí)體、配置Endpoint、啟動(dòng)服務(wù)等步驟,結(jié)合實(shí)例代碼給大家介紹的非常詳細(xì),感興趣的朋友跟隨小編一起看看吧2024-11-11
詳解java中面向?qū)ο笤O(shè)計(jì)模式類(lèi)與類(lèi)的關(guān)系
這篇文章主要介紹了java面向?qū)ο笤O(shè)計(jì)模式中類(lèi)與類(lèi)之間的關(guān)系,下面小編和大家一起來(lái)學(xué)習(xí)一下吧2019-05-05
springmvc json類(lèi)型轉(zhuǎn)換錯(cuò)誤解決方案
這篇文章主要介紹了springmvc json類(lèi)型轉(zhuǎn)換錯(cuò)誤解決方案,文中通過(guò)示例代碼介紹的非常詳細(xì),對(duì)大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價(jià)值,需要的朋友可以參考下2019-12-12
Java幸運(yùn)28系統(tǒng)搭建數(shù)組的使用實(shí)例詳解
在本篇文章里小編給大家整理了關(guān)于Java幸運(yùn)28系統(tǒng)搭建數(shù)組的使用實(shí)例內(nèi)容,有需要的朋友們可以參考學(xué)習(xí)下。2019-09-09
使用XSD校驗(yàn)Mybatis的SqlMapper配置文件的方法(1)
這篇文章以前面對(duì)SqlSessionFactoryBean的重構(gòu)為基礎(chǔ),簡(jiǎn)單的介紹了相關(guān)操作知識(shí),然后在給大家分享使用XSD校驗(yàn)Mybatis的SqlMapper配置文件的方法,感興趣的朋友參考下吧2016-11-11
spring的構(gòu)造函數(shù)注入屬性@ConstructorBinding用法
這篇文章主要介紹了關(guān)于spring的構(gòu)造函數(shù)注入屬性@ConstructorBinding用法,具有很好的參考價(jià)值,希望對(duì)大家有所幫助,如有錯(cuò)誤或未考慮完全的地方,望不吝賜教2023-12-12

