SpringBoot整合RabbitMQ實(shí)現(xiàn)延遲隊(duì)列和死信隊(duì)列
一、死信隊(duì)列
RabbitMQ的死信隊(duì)列(Dead Letter Queue,DLQ)是一種特殊的隊(duì)列,用于接收其他隊(duì)列中的“死信”消息。所謂“死信”,是指滿足一定條件而無(wú)法被消費(fèi)者正確處理的消息,這些條件包括消息被拒絕、消息過(guò)期、消息達(dá)到最大重試次數(shù)等。
當(dāng)消息成為死信時(shí),RabbitMQ會(huì)將其重新發(fā)送到指定的死信隊(duì)列,而不是丟棄它們。這樣做的好處是可以對(duì)死信進(jìn)行分析和處理,例如記錄日志、重新入隊(duì)或者進(jìn)一步處理。
死信隊(duì)列通常與RabbitMQ的延遲隊(duì)列(Delayed Message Queue)一起使用,通過(guò)延遲隊(duì)列延遲消息的處理時(shí)間,可以更容易地觸發(fā)消息成為死信的條件,從而進(jìn)行測(cè)試和調(diào)試。
死信隊(duì)列在消息中間件中有許多實(shí)際應(yīng)用場(chǎng)景,主要用于處理無(wú)法被正常消費(fèi)的消息,增強(qiáng)了消息的可靠性和處理能力。以下是一些常見(jiàn)的應(yīng)用場(chǎng)景:
延遲消息處理:通過(guò)將消息發(fā)送到延遲隊(duì)列,在指定的時(shí)間后再將消息發(fā)送到目標(biāo)隊(duì)列,實(shí)現(xiàn)延遲處理消息的功能。
消息重試:當(dāng)消費(fèi)者無(wú)法處理消息時(shí),消息可以被重新發(fā)送到隊(duì)列并設(shè)置重試次數(shù),達(dá)到最大重試次數(shù)后轉(zhuǎn)發(fā)到死信隊(duì)列,以便進(jìn)行進(jìn)一步處理。
異常處理:當(dāng)消息無(wú)法被消費(fèi)者正常處理時(shí)(如格式錯(cuò)誤、業(yè)務(wù)異常等),將消息轉(zhuǎn)發(fā)到死信隊(duì)列,用于記錄日志、報(bào)警或人工處理。
消息超時(shí)處理:當(dāng)消息在隊(duì)列中等待時(shí)間過(guò)長(zhǎng)時(shí),可以設(shè)置消息的過(guò)期時(shí)間(TTL),超過(guò)時(shí)間后將消息轉(zhuǎn)發(fā)到死信隊(duì)列。
消息路由失敗:當(dāng)消息無(wú)法被正確路由到目標(biāo)隊(duì)列時(shí),可以將消息發(fā)送到死信隊(duì)列,避免消息丟失。
消息版本兼容性處理:當(dāng)消息的格式或內(nèi)容發(fā)生變化時(shí),通過(guò)死信隊(duì)列可以處理老版本消息,確保新版本系統(tǒng)的兼容性。
RabbitMQ的工作模式

死信隊(duì)列的工作模式
今天我要實(shí)現(xiàn)的就是這個(gè)延遲隊(duì)列和死信隊(duì)列。生產(chǎn)者首先向延遲隊(duì)列發(fā)送消息,待達(dá)到TTL后消息會(huì)被轉(zhuǎn)送到死信隊(duì)列當(dāng)中,消費(fèi)者會(huì)從死信隊(duì)列中獲取消息進(jìn)行消費(fèi)。

二、RabbitMQ相關(guān)的安裝
win10安裝rabbitMQ的詳細(xì)步驟_java_腳本之家 (jb51.net)
我這里直接引用別人的文章了,下載需要大家去看一看。
RabbitMQ延遲插件的安裝。
三、SpringBoot引入RabbitMQ
1.引入依賴
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter</artifactId>
</dependency>
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-web</artifactId>
</dependency>
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-amqp</artifactId>
</dependency>
<dependency>
<groupId>org.slf4j</groupId>
<artifactId>slf4j-log4j12</artifactId>
</dependency>2.創(chuàng)建隊(duì)列和交換器
這一步是很重要的,如果你配置錯(cuò)誤了,消息很可能無(wú)法正確的傳送。要實(shí)現(xiàn)延遲隊(duì)列和死信隊(duì)列,我們一共要?jiǎng)?chuàng)建以下幾個(gè)組件:
- 延遲隊(duì)列
- 延遲隊(duì)列的交換器
- 死信隊(duì)列
- 死信隊(duì)列的交換器
在我們創(chuàng)建了這幾個(gè)組件之后,我們還要干一些事情,我們需要把這些組件進(jìn)行組裝,如果你不了解RabbitMQ的基礎(chǔ),你可以先看看基礎(chǔ)教學(xué),我這里簡(jiǎn)單的說(shuō)一下。RabbitMQ中有一種綁定方式,這種綁定方式會(huì)把BindingKey和RoutingKey完全匹配的進(jìn)行綁定,如下圖所示,生產(chǎn)者發(fā)送了一個(gè)BindingKey為“warning”的消息,那么這個(gè)消息就會(huì)被發(fā)送到Queue1和Queue2,這并不難理解。
我們要做的就是把隊(duì)列和交換器通過(guò)一個(gè)RoutingKey綁定在一起。

2.1 變量聲明
接下來(lái)的代碼要好好看了,首先我們把我們后邊要用到的名稱變量全部定義出來(lái)。因?yàn)檫@個(gè)名稱起的很長(zhǎng),我們不方便直接使用。創(chuàng)建DeadRabbitConfig。在類中定義如下變量,延遲隊(duì)列交換器名稱、延遲隊(duì)列名稱、延遲隊(duì)列Routing名稱。除此之外還有死信隊(duì)列交換器名稱、死信隊(duì)列名稱和死信Routing名稱。
// 延遲隊(duì)列交換器名稱
public static final String DELAY_EXCHANGE_NAME = "delay.queue.demo.business.exchange";
// 延遲隊(duì)列A名稱
public static final String DELAY_QUEUE_A_NAME = "delay.queue.demo.business.queue_a";
// 延遲隊(duì)列B名稱
public static final String DELAY_QUEUE_B_NAME = "delay.queue.demo.business.queue_b";
// 延遲隊(duì)列routingA名稱
public static final String DELAY_QUEUE_ROUTING_A_NAME = "delay.queue.demo.business.queue_a.routing_key";
// 延遲隊(duì)列routingB名稱
public static final String DELAY_QUEUE_ROUTING_B_NAME = "delay.queue.demo.business.queue_b.routing_key";
// 死信隊(duì)列
public static final String DEAD_LETTER_EXCHANGE = "delay.queue.demo.deadletter.exchange";
public static final String DEAD_LETTER_QUEUE_A_ROUTING_KEY = "delay.queue.demo.deadletter.delay_10s.routing_key";
public static final String DEAD_LETTER_QUEUE_B_ROUTING_KEY = "delay.queue.demo.deadletter.delay_60s.routing_key";
public static final String DEAD_LETTER_QUEUE_A_NAME = "delay.queue.demo.deadletter.queue_a";
public static final String DEAD_LETTER_QUEUE_B_NAME = "delay.queue.demo.deadletter.queue_b";2.2 創(chuàng)建延遲交換器
// 注冊(cè)延遲交換器delayExchange
@Bean("delayExchange")
public DirectExchange delayExchange(){
return new DirectExchange(DELAY_EXCHANGE_NAME);
}2.3 創(chuàng)建延遲隊(duì)列
這里的延遲隊(duì)列需要我們額外的配置一些參數(shù),用于和死信隊(duì)列進(jìn)行信息發(fā)送。這里我是用了兩種不同的方式構(gòu)建延遲隊(duì)列A和延遲隊(duì)列B,在延遲隊(duì)列A種我沒(méi)有設(shè)置TTL參數(shù),而是通過(guò)RabbitMQ的延遲插件實(shí)現(xiàn)的,而延遲隊(duì)列B我設(shè)置了TTL為10000ms,也就是十秒,十秒內(nèi)消息如果沒(méi)有被消費(fèi)掉就會(huì)發(fā)送到死信隊(duì)列。
// 注冊(cè)延遲隊(duì)列A 還要綁定死信交換器和死信routingA
@Bean("delayQueueA")
public Queue delayQueueA(){
Map<String,Object> args = new HashMap<>();
args.put("x-dead-letter-exchange",DEAD_LETTER_EXCHANGE);
args.put("x-dead-letter-routing-key",DEAD_LETTER_QUEUE_A_ROUTING_KEY);
//args.put("x-message-ttl",6000);
return QueueBuilder.durable(DELAY_QUEUE_A_NAME).withArguments(args).build();
}
// 注冊(cè)延遲隊(duì)列B 還要綁定死信交換器和死信routingB
@Bean("delayQueueB")
public Queue delayQueueB(){
Map<String,Object> args = new HashMap<>();
args.put("x-dead-letter-exchange",DEAD_LETTER_EXCHANGE);
args.put("x-dead-letter-routing-key",DEAD_LETTER_QUEUE_B_ROUTING_KEY);
args.put("x-message-ttl",10000);
return QueueBuilder.durable(DELAY_QUEUE_B_NAME).withArguments(args).build();
}2.4 延遲隊(duì)列綁定延遲交換器
// 延遲隊(duì)列A綁定交換器
@Bean
public Binding delayQueueABinding(@Qualifier("delayQueueA") Queue queue, @Qualifier("delayExchange") DirectExchange delayExchange){
return BindingBuilder.bind(queue).to(delayExchange).with(DELAY_QUEUE_ROUTING_A_NAME);
}
// 延遲隊(duì)列B綁定交換器
@Bean
public Binding delayQueueBBinding(@Qualifier("delayQueueB") Queue queue,@Qualifier("delayExchange") DirectExchange delayExchange){
return BindingBuilder.bind(queue).to(delayExchange).with(DELAY_QUEUE_ROUTING_B_NAME);
}2.5 死信隊(duì)列配置
與延遲隊(duì)列不同的是,死信隊(duì)列并沒(méi)有配置延遲參數(shù)。
// 注冊(cè)死信隊(duì)列A
@Bean("deadLetterQueueA")
public Queue deadLetterQueueA(){
return new Queue(DEAD_LETTER_QUEUE_A_NAME);
}
// 注冊(cè)死信隊(duì)列B
@Bean("deadLetterQueueB")
public Queue deadLetterQueueB(){
return new Queue(DEAD_LETTER_QUEUE_B_NAME);
}
// 注冊(cè)死信交換器
@Bean
public DirectExchange deadLetterExchange(){
return new DirectExchange(DEAD_LETTER_EXCHANGE);
}
// 死信隊(duì)列A綁定死信交換器
@Bean
public Binding deadLetterQueueABinding(@Qualifier("deadLetterQueueA") Queue queue, @Qualifier("deadLetterExchange") DirectExchange deadLetterExchange){
return BindingBuilder.bind(queue).to(deadLetterExchange).with(DEAD_LETTER_QUEUE_A_ROUTING_KEY);
}
// 死信隊(duì)列B綁定死信交換器
@Bean
public Binding deadLetterQueueBBinding(@Qualifier("deadLetterQueueB") Queue queue, @Qualifier("deadLetterExchange")DirectExchange deadLetterExchange){
return BindingBuilder.bind(queue).to(deadLetterExchange).with(DEAD_LETTER_QUEUE_B_ROUTING_KEY);
}到此為止,RabbitMQ的組件配置完成。
3. 添加application.yml
server:
port: 8081
spring:
application:
name: test-rabbitmq-producer
rabbitmq:
host: 127.0.0.1
port: 5672
username: guest
password: guest4. 添加RabbitMQListener (消費(fèi)者)
下方的代碼一共有兩個(gè)消費(fèi)者,一個(gè)消費(fèi)者獲取死信隊(duì)列A中的消息,另一個(gè)消費(fèi)者獲取死信隊(duì)列B中的消息。
@Component
public class DeadLetterQueueConsumer {
public static final Logger LOGGER = LoggerFactory.getLogger(DeadLetterQueueConsumer.class);
@RabbitListener(queues = DeadRabbitConfig.DEAD_LETTER_QUEUE_A_NAME,ackMode = "MANUAL")
public void receiveA(Message message, Channel channel) throws IOException {
String msg = new String(message.getBody());
LOGGER.info("當(dāng)前時(shí)間:{},死信隊(duì)列A收到消息:{}", new Date().toString(), msg);
System.out.println(message.getMessageProperties().getDeliveryTag());
channel.basicAck(message.getMessageProperties().getDeliveryTag(), false);
}
@RabbitListener(queues = DeadRabbitConfig.DEAD_LETTER_QUEUE_B_NAME,ackMode = "MANUAL")
public void receiveB(Message message, Channel channel) throws IOException {
String msg = new String(message.getBody());
LOGGER.info("當(dāng)前時(shí)間:{},死信隊(duì)列B收到消息:{}", new Date().toString(), msg);
channel.basicAck(message.getMessageProperties().getDeliveryTag(), false);
}
}5. 創(chuàng)建DelayMessageSender
這里采用的就是兩種不同的方式,一種方式是使用插件來(lái)延遲消息的發(fā)送,另一種是通過(guò)TTL參數(shù)。
@Component
public class DelayMessageSender {
@Resource
RabbitTemplate rabbitTemplate;
public void sendMessage(String msg,Integer delayTimes){
switch (delayTimes){
case 6:
rabbitTemplate.convertAndSend(DeadRabbitConfig.DELAY_EXCHANGE_NAME, DeadRabbitConfig.DELAY_QUEUE_ROUTING_A_NAME,msg,new MessagePostProcessor() {
@Override
public Message postProcessMessage(Message message) throws AmqpException {
message.getMessageProperties().setExpiration(String.valueOf(6000));
return message;
}
});
break;
case 10:
rabbitTemplate.convertAndSend(DeadRabbitConfig.DELAY_QUEUE_B_NAME,msg);
break;
}
}
}6. 創(chuàng)建Controller
@RestController
@RequestMapping("/student")
public class StudentController {
@Autowired
DelayMessageSender messageSender;
@RequestMapping("/send-message")
public String sendMessage(String msg,Integer delayTimes){
System.out.println(new Date());
messageSender.sendMessage(msg,delayTimes);
return "發(fā)送成功";
}
}7.測(cè)試
在瀏覽器中輸入以下地址進(jìn)入RabbitMQ界面。賬號(hào)密碼都是guest。
http://localhost:15672/
先來(lái)看看我們的初始隊(duì)列。這里是什么都沒(méi)有的。

然后我們啟動(dòng)項(xiàng)目后在看。我們剛才創(chuàng)建出來(lái)的四個(gè)隊(duì)列全部都被加載了出來(lái)。

使用PostMan發(fā)送一次請(qǐng)求。

我們的請(qǐng)求在17s的時(shí)候發(fā)送到后端,消息打印在23s,說(shuō)明我們的延遲隊(duì)列有效果。

接下來(lái)我們測(cè)試10s的延遲隊(duì)列。

10s后死信隊(duì)列B成功的接收到了消息。

四、死信隊(duì)列的應(yīng)用場(chǎng)景
延遲隊(duì)列通常用于需要延遲執(zhí)行某些任務(wù)或觸發(fā)某些事件的場(chǎng)景。例如,在電子商務(wù)中,可以使用延遲隊(duì)列實(shí)現(xiàn)訂單超時(shí)未支付自動(dòng)取消功能。
1.訂單創(chuàng)建:
用戶下單后,系統(tǒng)生成訂單,并將訂單信息發(fā)送到一個(gè)普通隊(duì)列,同時(shí)設(shè)置一個(gè)TTL(Time-To-Live)為30分鐘。這個(gè)隊(duì)列配置了死信交換機(jī)(Dead Letter Exchange, DLX),當(dāng)消息過(guò)期后會(huì)被轉(zhuǎn)發(fā)到死信隊(duì)列。2.等待支付:
在30分鐘內(nèi),用戶可以完成支付。如果用戶在30分鐘內(nèi)支付完成,系統(tǒng)會(huì)從普通隊(duì)列中移除對(duì)應(yīng)的消息并正常處理訂單。3.訂單超時(shí)處理:
如果用戶未在30分鐘內(nèi)完成支付,消息會(huì)自動(dòng)過(guò)期并轉(zhuǎn)發(fā)到死信交換機(jī),進(jìn)而轉(zhuǎn)發(fā)到死信隊(duì)列。4.取消訂單:
系統(tǒng)有一個(gè)專門的消費(fèi)者監(jiān)聽(tīng)死信隊(duì)列。當(dāng)有消息進(jìn)入死信隊(duì)列時(shí),消費(fèi)者會(huì)自動(dòng)處理這些消息,即取消訂單、釋放庫(kù)存,并通知用戶訂單已取消。5.定時(shí)任務(wù)(可選):
雖然死信隊(duì)列已經(jīng)提供了超時(shí)訂單的處理,但為了防止消息丟失或處理延遲,可以設(shè)置一個(gè)定時(shí)任務(wù)定期檢查訂單狀態(tài),確保所有超時(shí)未支付的訂單都得到了處理。
以上就是SpringBoot整合RabbitMQ實(shí)現(xiàn)延遲隊(duì)列和死信隊(duì)列的詳細(xì)內(nèi)容,更多關(guān)于SpringBoot RabbitMQ隊(duì)列的資料請(qǐng)關(guān)注腳本之家其它相關(guān)文章!
- Springboot使用RabbitMq延遲隊(duì)列和死信隊(duì)列詳解
- Springboot使用Rabbitmq的延時(shí)隊(duì)列+死信隊(duì)列實(shí)現(xiàn)消息延期消費(fèi)
- springboot整合RabbitMQ中死信隊(duì)列的實(shí)現(xiàn)
- 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)文章
教你怎么通過(guò)IDEA設(shè)置堆內(nèi)存空間
這篇文章主要介紹了教你怎么通過(guò)IDEA設(shè)置堆內(nèi)存空間,文中有非常詳細(xì)的代碼示例,對(duì)正在使用IDEA的小伙伴們很有幫助喲,需要的朋友可以參考下2021-05-05
解讀為什么不推薦使用keySet()進(jìn)行遍歷HashMap
這篇文章主要介紹了我為什么不推薦使用keySet()進(jìn)行遍歷HashMap的問(wèn)題,具有很好的參考價(jià)值,希望對(duì)大家有所幫助,如有錯(cuò)誤或未考慮完全的地方,望不吝賜教2025-05-05
Java工廠模式之簡(jiǎn)單工廠,工廠方法,抽象工廠模式詳解
這篇文章主要為大家詳細(xì)介紹了Java工廠模式之簡(jiǎn)單工廠、工廠方法、抽象工廠模式,文中示例代碼介紹的非常詳細(xì),具有一定的參考價(jià)值,感興趣的小伙伴們可以參考一下,希望能夠給你帶來(lái)幫助2022-02-02
java 中Object與Objects的區(qū)別在哪里
這篇文章主要介紹了java 中Object與Objects的區(qū)別說(shuō)明,具有很好的參考價(jià)值,希望對(duì)大家有所幫助。如有錯(cuò)誤或未考慮完全的地方,望不吝賜教2021-05-05
詳解SpringBoot注冊(cè)Windows服務(wù)和啟動(dòng)報(bào)錯(cuò)的原因
這篇文章主要介紹了詳解SpringBoot注冊(cè)Windows服務(wù)和啟動(dòng)報(bào)錯(cuò)的原因,小編覺(jué)得挺不錯(cuò)的,現(xiàn)在分享給大家,也給大家做個(gè)參考。一起跟隨小編過(guò)來(lái)看看吧2019-03-03
maven多個(gè)plugin相同phase的執(zhí)行順序
這篇文章主要介紹了maven多個(gè)plugin相同phase的執(zhí)行順序,本文給大家介紹的非常詳細(xì),對(duì)大家的學(xué)習(xí)或工作具有一定的參考借鑒價(jià)值,需要的朋友可以參考下2020-12-12
SpringBoot引入Redis報(bào)Redis?command?timed?out兩種異常情況
這篇文章主要給大家介紹了關(guān)于SpringBoot引入Redis報(bào)Redis?command?timed?out兩種異常情況的相關(guān)資料,文中通過(guò)示例代碼介紹的非常詳細(xì),需要的朋友可以參考下2023-08-08
mybatis報(bào)錯(cuò)元素內(nèi)容必須由格式正確的字符數(shù)據(jù)或標(biāo)記組成異常的解決辦法
今天小編就為大家分享一篇關(guān)于mybatis查詢出錯(cuò)解決辦法,小編覺(jué)得內(nèi)容挺不錯(cuò)的,現(xiàn)在分享給大家,具有很好的參考價(jià)值,需要的朋友一起跟隨小編來(lái)看看吧2018-12-12
Java反應(yīng)式框架Reactor中的Mono和Flux
這篇文章主要介紹了Java反應(yīng)式框架Reactor中的Mono和Flux,本文給大家介紹的非常詳細(xì),對(duì)大家的學(xué)習(xí)或工作具有一定的參考借鑒價(jià)值,需要的朋友可以參考下2021-07-07

