關于SpringBoot整合RabbitMQ實現(xiàn)死信隊列
概念介紹
什么是死信
死信可以理解成沒有被正常消費的消息,在RabbitMQ中以下幾種情況會被認定為死信:
- 消費者使用basic.reject或basic.nack(重新排隊參數(shù)設置為false)對消息進行否定確認。
- 消息到達生存時間還未被消費。
- 隊列超過長度限制,消息被丟棄。
這些消息會被發(fā)送到死信交換機并路由到死信隊列中(在RabbitMQ中死信交換機和死信隊列就是普通的交換機和隊列)。其流轉(zhuǎn)過程如下圖

死信隊列應用
- 作為消息可靠性的一個擴展。比如,在隊列已滿的情況下也不會丟失消息。
- 可以實現(xiàn)延遲消費功能。比如,訂單15分鐘內(nèi)未支付。
注意事項:基于死信隊列實現(xiàn)的延遲消費不適合時間過于復雜的場景。比如,一個隊列中第一條消息TTL為10s,第二條消息TTL為5s,由于RabbitMQ只會監(jiān)聽第一條消息,所以本應第二條消息先達到TTL會在第一條消息的TTL之后。對于該現(xiàn)象有兩種解決方案:
- 維護多個隊列,每個隊列維護一個TTL時間。
- 使用延遲交換機。這種方式需要下載插件支持
工程搭建
環(huán)境說明
- RabbitMQ環(huán)境
- Java版本:JDK1.8
- Maven版本:apache-maven-3.6.3
- 開發(fā)工具:IntelliJ IDEA
搭建步驟
1.創(chuàng)建SpringBoot項目。
2.pom.xml文件導入RabbitMQ依賴。
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-amqp</artifactId>
</dependency>3.application.yml文件添加RabbitMQ配置。
spring:
# rabbitmq配置信息 RabbitProperties類
rabbitmq:
host: 127.0.0.1
port: 5672
username: guest
password: guest
virtual-host: /
# 開啟confirm機制
publisher-confirm-type: correlated
# 開啟return機制
publisher-returns: true
#全局配置,局部配置存在就以局部為準
listener:
simple:
acknowledge-mode: manual # 手動ACK實現(xiàn)死信
準備Exchange&Queue
@Configuration
public class RabbitMQConfig {
/**
* 正常隊列
*/
public static final String EXCHANGE = "boot-exchange";
public static final String QUEUE = "boot-queue";
public static final String ROUTING_KEY = "boot-rout";
/**
* 死信隊列
*/
public static final String DEAD_EXCHANGE = "dead-exchange";
public static final String DEAD_QUEUE = "dead-queue";
public static final String DEAD_ROUTING_KEY = "dead-rout";
/**
* 聲明死信交換機
*
* @return
*/
@Bean
public Exchange deadExchange() {
return ExchangeBuilder.directExchange(DEAD_EXCHANGE).build();
}
/**
* 聲明死信隊列
*
* @return
*/
@Bean
public Queue deadQueue() {
return QueueBuilder.durable(DEAD_QUEUE).build();
}
/**
* 綁定死信的隊列和交換機
*
* @param deadExchange
* @param deadQueue
* @return
*/
@Bean
public Binding deadBind(Exchange deadExchange, Queue deadQueue) {
return BindingBuilder.bind(deadQueue).to(deadExchange).with(DEAD_ROUTING_KEY).noargs();
}
/**
* 聲明交換機,同channel.exchangeDeclare(EXCHANGE, BuiltinExchangeType.DIRECT);
*
* @return
*/
@Bean
public Exchange bootExchange() {
return ExchangeBuilder.directExchange(EXCHANGE).build();
}
/**
* 聲明隊列,同channel.queueDeclare(QUEUE, true, false, false, null);
* 綁定死信交換機及路由key
*
* @return
*/
@Bean
public Queue bootQueue() {
return QueueBuilder.durable(QUEUE)
.deadLetterExchange(DEAD_EXCHANGE)
.deadLetterRoutingKey(DEAD_ROUTING_KEY)
//聲明隊列屬性有更改時需要刪除隊列
//給隊列設置消息時長
//.ttl(10000)
//隊列最大長度
.maxLength(1)
.build();
}
/**
* 綁定隊列和交換機,同 channel.queueBind(QUEUE, EXCHANGE, ROUTING_KEY);
*
* @param bootExchange
* @param bootQueue
* @return
*/
@Bean
public Binding bootBind(Exchange bootExchange, Queue bootQueue) {
return BindingBuilder.bind(bootQueue).to(bootExchange).with(ROUTING_KEY).noargs();
}
}監(jiān)聽死信隊列
@RabbitListener(queues = RabbitMQConfig.DEAD_QUEUE)
public void listener_dead(String msg, Channel channel, Message message) throws IOException {
System.out.println("死信接收到消息" + msg);
System.out.println("唯一標識:" + message.getMessageProperties().getCorrelationId());
System.out.println("messageID:" + message.getMessageProperties().getMessageId());
channel.basicAck(message.getMessageProperties().getDeliveryTag(), false);
}方式一、消費者拒絕&否認
- 拒絕消息
@RabbitListener(queues = RabbitMQConfig.QUEUE)
public void listener(String msg, Channel channel, Message message) throws IOException {
System.out.println("接收到消息" + msg);
channel.basicReject(message.getMessageProperties().getDeliveryTag(), false)
}- 否認消息
@RabbitListener(queues = RabbitMQConfig.QUEUE)
public void listener(String msg, Channel channel, Message message) throws IOException {
System.out.println("接收到消息" + msg);
channel.basicNack(message.getMessageProperties().getDeliveryTag(), false, false);
}方式二、超過消息TTL 發(fā)送消息時設置TTL
@SpringBootTest
public class Publisher {
@Autowired
private RabbitTemplate template;
/**
* 5秒未被消費會路由到死信隊列
*/
@Test
public void publish_expir() {
template.convertAndSend(RabbitMQConfig.EXCHANGE, RabbitMQConfig.ROUTING_KEY, "hello expir dead", message -> {
message.getMessageProperties().setExpiration("5000");
return message;
});
}
}- 設置隊列所有消息的TTL
更新RabbitMQConfig類中bootQueue() ,更新后需要刪除隊列,因為隊列屬性有更改。
@Bean
public Queue bootQueue() {
return QueueBuilder.durable(QUEUE)
.deadLetterExchange(DEAD_EXCHANGE)
.deadLetterRoutingKey(DEAD_ROUTING_KEY)
//聲明隊列屬性有更改時需要刪除隊列
//給隊列設置消息時長
.ttl(10000)
.build();
}方式三、超過隊列長度限制
設置隊列長度限制,當隊列長度超過設置的閾值,消息便會路由到死信隊列。
@Bean
public Queue bootQueue() {
return QueueBuilder.durable(QUEUE)
.deadLetterExchange(DEAD_EXCHANGE)
.deadLetterRoutingKey(DEAD_ROUTING_KEY)
//聲明隊列屬性有更改時需要刪除隊列
.maxLength(1)
.build();
}代碼倉庫 點我
到此這篇關于關于SpringBoot整合RabbitMQ實現(xiàn)死信隊列的文章就介紹到這了,更多相關RabbitMQ實現(xiàn)死信隊列內(nèi)容請搜索腳本之家以前的文章或繼續(xù)瀏覽下面的相關文章希望大家以后多多支持腳本之家!
- Springboot使用RabbitMq延遲隊列和死信隊列詳解
- Springboot使用Rabbitmq的延時隊列+死信隊列實現(xiàn)消息延期消費
- springboot整合RabbitMQ中死信隊列的實現(xiàn)
- SpringBoot整合RabbitMQ實現(xiàn)延遲隊列和死信隊列
- springboot中RabbitMQ死信隊列的實現(xiàn)示例
- Springboot結(jié)合rabbitmq實現(xiàn)的死信隊列
- SpringBoot+RabbitMQ?實現(xiàn)死信隊列的示例
- SpringBoot整合RabbitMQ處理死信隊列和延遲隊列
- Springboot集成RabbitMQ死信隊列的實現(xiàn)
- SpringBoot集成RabbitMQ的方法(死信隊列)
- SpringBoot4.0整合RabbitMQ死信隊列詳解
相關文章
Spring線程池ThreadPoolExecutor配置并且得到任務執(zhí)行的結(jié)果
今天小編就為大家分享一篇關于Spring線程池ThreadPoolExecutor配置并且得到任務執(zhí)行的結(jié)果,小編覺得內(nèi)容挺不錯的,現(xiàn)在分享給大家,具有很好的參考價值,需要的朋友一起跟隨小編來看看吧2019-03-03
ArrayList和LinkedList區(qū)別及使用場景代碼解析
這篇文章主要介紹了ArrayList和LinkedList區(qū)別及使用場景代碼解析,小編覺得還是挺不錯的,具有一定借鑒價值,需要的朋友可以參考下2018-01-01
IntelliJ IDEA 2020.3.3現(xiàn)已發(fā)布!新增“受信任項目”功能
這篇文章主要介紹了IntelliJ IDEA 2020.3.3現(xiàn)已發(fā)布!新增“受信任項目”功能,本文給大家分享了idea2020.3.3激活碼的詳細破解教程,每種方法都很好用,使用idea2020.3以下所有版本,需要的朋友可以參考下2021-03-03
idea2020安裝MybatisCodeHelper插件的圖文教程
這篇文章主要介紹了idea2020安裝MybatisCodeHelper插件的方法,本文通過圖文并茂的形式給大家介紹的非常詳細,對大家的學習或工作具有一定的參考借鑒價值,需要的朋友可以參考下2020-09-09
java使用google身份驗證器實現(xiàn)動態(tài)口令驗證的示例
本篇文章主要介紹了java使用google身份驗證器實現(xiàn)動態(tài)口令驗證的示例,具有一定的參考價值,有興趣的可以了解一下2017-08-08

