Java中RabbitMQ延遲隊(duì)列實(shí)現(xiàn)詳解
一、RabbitMQ延遲隊(duì)列實(shí)現(xiàn)
1.1、RabbitMQ延遲隊(duì)列實(shí)現(xiàn)流程

- 生產(chǎn)者生產(chǎn)一條延遲消息,根據(jù)延遲時(shí)間的不同,利用不同的routing-key將消息路由到不同的延遲隊(duì)列,每個(gè)隊(duì)列都設(shè)置了不同的 TTL 屬性 ( TTL ( Time To Live ) 生存時(shí)間 ),并綁定到同一個(gè)死信交換機(jī)中。
- 消息過期后,根據(jù)routing-key的不同,又會(huì)被死信交換機(jī)路由到不同的死信隊(duì)列中,消費(fèi)者只需要監(jiān)聽對(duì)應(yīng)的死信隊(duì)列進(jìn)行消費(fèi)即可。
1.2、配置RabbitMQ連接
#[ RabbitMQ相關(guān)配置 ] #rabbitmq服務(wù)器IP spring.rabbitmq.host=安裝RabbitMQ的服務(wù)器IP #rabbitmq服務(wù)器端口(默認(rèn)為5672) spring.rabbitmq.port=5672 #用戶名 spring.rabbitmq.username=guest #用戶密碼 spring.rabbitmq.password=guest #虛擬主機(jī)(一個(gè)RabbitMQ服務(wù)可以配置多個(gè)虛擬主機(jī),每一個(gè)虛擬機(jī)主機(jī)之間是相互隔離,相互獨(dú)立的,授權(quán)用戶到指定的virtual-host就可以發(fā)送消息到指定隊(duì)列) #vhost虛擬主機(jī)地址( 默認(rèn)為/ ) spring.rabbitmq.virtual-host=/
1.3、創(chuàng)建配置類
配置兩個(gè)交換機(jī)、四個(gè)隊(duì)列、以及根據(jù)路由鍵配置交換機(jī)和隊(duì)列的綁定關(guān)系
import org.springframework.amqp.core.*;
import org.springframework.beans.factory.annotation.Qualifier;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import java.util.HashMap;
import java.util.Map;
@Configuration
public class RabbitMQConfiguration {
//延遲交換機(jī)
public static final String DELAY_EXCHANGE = "delay_exchange";
//延遲隊(duì)列A
public static final String DELAY_QUEUE_A = "delay_queue_a";
//延遲隊(duì)列B
public static final String DELAY_QUEUE_B = "delay_queue_b";
//延遲路由鍵10S
public static final String DELAY_QUEUE_10S_ROUTING_KEY = "delay_queue_10s_routing_key";
//延遲路由鍵60S
public static final String DELAY_QUEUE_60S_ROUTING_KEY = "delay_queue_60s_routing_key";
//死信交換機(jī)
public static final String DEAD_LETTER_EXCHANGE = "dead_letter_exchange";
//死信隊(duì)列A
public static final String DEAD_LETTER_QUEUE_A = "dead_letter_queue_a";
//死信隊(duì)列B
public static final String DEAD_LETTER_QUEUE_B = "dead_letter_queue_b";
//死信路由鍵10S
public static final String DEAD_LETTER_QUEUE_10S_ROUTING_KEY = "dead_letter_queue_10s_routing_key";
//死信路由鍵60S
public static final String DEAD_LETTER_QUEUE_60S_ROUTING_KEY = "dead_letter_queue_60s_routing_key";
//延遲交換機(jī)
@Bean("delayExchange")
public DirectExchange delayExchange(){
return new DirectExchange(DELAY_EXCHANGE, true, false);
}
//延遲隊(duì)列A
@Bean("delayQueueA")
public Queue delayQueueA(){
Map<String, Object> args = new HashMap<>();
//設(shè)置延遲隊(duì)列綁定的死信交換機(jī)
args.put("x-dead-letter-exchange", DEAD_LETTER_EXCHANGE);
//設(shè)置延遲隊(duì)列綁定的死信路由鍵
args.put("x-dead-letter-routing-key", DEAD_LETTER_QUEUE_10S_ROUTING_KEY);
//設(shè)置延遲隊(duì)列的 TTL 消息存活時(shí)間
args.put("x-message-ttl", 10*1000);
return new Queue(DELAY_QUEUE_A, true, false, false, args);
}
//延遲隊(duì)列B
@Bean("delayQueueB")
public Queue delayQueueB(){
Map<String, Object> args = new HashMap<>();
//設(shè)置延遲隊(duì)列綁定的死信交換機(jī)
args.put("x-dead-letter-exchange", DEAD_LETTER_EXCHANGE);
//設(shè)置延遲隊(duì)列綁定的死信路由鍵
args.put("x-dead-letter-routing-key", DEAD_LETTER_QUEUE_60S_ROUTING_KEY);
//設(shè)置延遲隊(duì)列的 TTL 消息存活時(shí)間
args.put("x-message-ttl", 60*1000);
return new Queue(DELAY_QUEUE_B, true, false, false, args);
}
//延遲隊(duì)列A的綁定關(guān)系
@Bean("delayBindingA")
public Binding delayBindingA(@Qualifier("delayQueueA")Queue queue,
@Qualifier("delayExchange")DirectExchange exchange){
return BindingBuilder.bind(queue).to(exchange).with(DELAY_QUEUE_10S_ROUTING_KEY);
}
//延遲隊(duì)列B的綁定關(guān)系
@Bean("delayBindingB")
public Binding delayBindingB(@Qualifier("delayQueueB")Queue queue,
@Qualifier("delayExchange")DirectExchange exchange){
return BindingBuilder.bind(queue).to(exchange).with(DELAY_QUEUE_60S_ROUTING_KEY);
}
//死信交換機(jī)
@Bean("deadLetterExchange")
public DirectExchange deadLetterExchange(){
return new DirectExchange(DEAD_LETTER_EXCHANGE, true, false);
}
//死信隊(duì)列A
@Bean("deadLetterQueueA")
public Queue deadLetterQueueA(){
return new Queue(DEAD_LETTER_QUEUE_A, true);
}
//死信隊(duì)列B
@Bean("deadLetterQueueB")
public Queue deadLetterQueueB(){
return new Queue(DEAD_LETTER_QUEUE_B, true);
}
//死信隊(duì)列A的綁定關(guān)系
@Bean("deadLetterBindingA")
public Binding deadLetterBindingA(@Qualifier("deadLetterQueueA")Queue queue,
@Qualifier("deadLetterExchange")DirectExchange exchange){
return BindingBuilder.bind(queue).to(exchange).with(DEAD_LETTER_QUEUE_10S_ROUTING_KEY);
}
//死信隊(duì)列B的綁定關(guān)系
@Bean("deadLetterBindingB")
public Binding deadLetterBindingB(@Qualifier("deadLetterQueueB")Queue queue,
@Qualifier("deadLetterExchange")DirectExchange exchange){
return BindingBuilder.bind(queue).to(exchange).with(DEAD_LETTER_QUEUE_60S_ROUTING_KEY);
}
}1.4、創(chuàng)建一個(gè)枚舉類來配置延遲類型
@Getter
@AllArgsConstructor
public enum DelayTypeEnum {
//10s
DELAY_10s(1),
//60s
DELAY_60s(2);
private Integer type;
/**
* 延遲類型
* @param type
* @return 延遲類型
*/
public static DelayTypeEnum getDelayTypeEnum(Integer type){
if(Objects.equals(type, DELAY_10s.type)){
return DELAY_10s;
}
if(Objects.equals(type, DELAY_60s.type)){
return DELAY_60s;
}
return null;
}
}1.5、創(chuàng)建生產(chǎn)者類發(fā)送消息
import com.cd.springbootrabbitmq.enums.DelayTypeEnum;
import org.springframework.amqp.rabbit.core.RabbitTemplate;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.stereotype.Component;
import static com.cd.springbootrabbitmq.config.RabbitMQConfiguration.DELAY_EXCHANGE;
import static com.cd.springbootrabbitmq.config.RabbitMQConfiguration.DELAY_QUEUE_10S_ROUTING_KEY;
import static com.cd.springbootrabbitmq.config.RabbitMQConfiguration.DELAY_QUEUE_60S_ROUTING_KEY;
/**
* 延遲消息生產(chǎn)者
*/
@Component
public class DelayMessageProducer {
@Autowired
private RabbitTemplate rabbitTemplate;
/**
* 發(fā)送延遲消息
* @param message 要發(fā)送的消息
* @param type 延遲類型(延時(shí)10s的延遲隊(duì)列 或 延時(shí)60s的延遲隊(duì)列)
*/
public void sendDelayMessage(String message, DelayTypeEnum type){
switch (type){
case DELAY_10s:
rabbitTemplate.convertAndSend(DELAY_EXCHANGE, DELAY_QUEUE_10S_ROUTING_KEY, message);
break;
case DELAY_60s:
rabbitTemplate.convertAndSend(DELAY_EXCHANGE, DELAY_QUEUE_60S_ROUTING_KEY, message);
break;
default:
break;
}
}
}1.6、創(chuàng)建消費(fèi)者類消費(fèi)消息
import lombok.extern.slf4j.Slf4j;
import org.springframework.amqp.core.Message;
import org.springframework.amqp.rabbit.annotation.RabbitListener;
import org.springframework.amqp.rabbit.core.RabbitTemplate;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.stereotype.Component;
import java.time.LocalDateTime;
import static com.cd.springbootrabbitmq.config.RabbitMQConfiguration.DEAD_LETTER_QUEUE_A;
import static com.cd.springbootrabbitmq.config.RabbitMQConfiguration.DEAD_LETTER_QUEUE_B;
@Slf4j
@Component
public class DeadLetterQueueConsumer {
@Autowired
private RabbitTemplate rabbitTemplate;
/**
* 監(jiān)聽死信隊(duì)列A
* @param message 接收的信息
*/
//@RabbitListener(queues = "dead_letter_queue_a")
@RabbitListener(queues = DEAD_LETTER_QUEUE_A)
public void receiveA(Message message) {
String msg = new String(message.getBody());
// 記錄日志
log.info("當(dāng)前時(shí)間:{},死信隊(duì)列A收到的消息:{}", LocalDateTime.now(), msg);
}
/**
* 監(jiān)聽死信隊(duì)列B
* @param message 接收的信息
*/
//@RabbitListener(queues = "dead_letter_queue_b")
@RabbitListener(queues = DEAD_LETTER_QUEUE_B)
public void receiveB(Message message){
String msg = new String(message.getBody());
// 記錄日志
log.info("當(dāng)前時(shí)間:{},死信隊(duì)列B收到的消息:{}", LocalDateTime.now(), msg);
}
}1.7、創(chuàng)建控制類
import com.cd.springbootrabbitmq.enums.DelayTypeEnum;
import com.cd.springbootrabbitmq.producer.DelayMessageProducer;
import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.web.bind.annotation.RequestMapping;
import org.springframework.web.bind.annotation.RestController;
import java.time.LocalDateTime;
import java.util.Objects;
@Slf4j
@RestController
@RequestMapping("/rabbitmq")
public class RabbitMQController {
@Autowired
private DelayMessageProducer producer;
@RequestMapping("/send")
public void send(String message, Integer delayType){
// 記錄日志
log.info("當(dāng)前時(shí)間:{},消息:{},延遲類型:{}", LocalDateTime.now(), message, delayType);
// 發(fā)送延遲消息
producer.sendDelayMessage(message, Objects.requireNonNull(DelayTypeEnum.getDelayTypeEnum(delayType)));
}
}1.8、測(cè)試
在瀏覽器中先后提交下面兩個(gè)請(qǐng)求:
1)localhost:8080/rabbitmq/send?message=測(cè)試自定義延遲處理60s&delayType=2
2)localhost:8080/rabbitmq/send?message=測(cè)試自定義延遲處理10s&delayType=1
查看idea控制臺(tái):

到此這篇關(guān)于Java中RabbitMQ延遲隊(duì)列實(shí)現(xiàn)詳解的文章就介紹到這了,更多相關(guān)RabbitMQ延遲隊(duì)列內(nèi)容請(qǐng)搜索腳本之家以前的文章或繼續(xù)瀏覽下面的相關(guān)文章希望大家以后多多支持腳本之家!
相關(guān)文章
Java線程池如何實(shí)現(xiàn)精準(zhǔn)控制每秒API請(qǐng)求
這篇文章主要介紹了Java線程池如何實(shí)現(xiàn)精準(zhǔn)控制每秒API請(qǐng)求問題,具有很好的參考價(jià)值,希望對(duì)大家有所幫助,如有錯(cuò)誤或未考慮完全的地方,望不吝賜教2024-08-08
Spring Cloud之服務(wù)監(jiān)控turbine的示例
這篇文章主要介紹了Spring Cloud之服務(wù)監(jiān)控turbine的示例,小編覺得挺不錯(cuò)的,現(xiàn)在分享給大家,也給大家做個(gè)參考。一起跟隨小編過來看看吧2018-05-05
SpringBoot集成easy-rules規(guī)則引擎流程詳解
這篇文章主要介紹了SpringBoot集成easy-rules規(guī)則引擎流程,合理的使用規(guī)則引擎可以極大的減少代碼復(fù)雜度,提升代碼可維護(hù)性。業(yè)界知名的開源規(guī)則引擎有Drools,功能豐富,但也比較龐大2023-03-03
Spring 靜態(tài)變量/構(gòu)造函數(shù)注入失敗的解決方案
我們經(jīng)常會(huì)遇到一下問題:Spring對(duì)靜態(tài)變量的注入為空、在構(gòu)造函數(shù)中使用Spring容器中的Bean對(duì)象,得到的結(jié)果為空。不要擔(dān)心,本文將為大家介紹如何解決這些問題,跟隨小編來看看吧2021-11-11
java 遍歷Map及Map轉(zhuǎn)化為二維數(shù)組的實(shí)例
這篇文章主要介紹了java 遍歷Map及Map轉(zhuǎn)化為二維數(shù)組的實(shí)例的相關(guān)資料,希望通過本文能幫助到大家,實(shí)現(xiàn)這樣的功能,需要的朋友可以參考下2017-08-08
Echarts+SpringMvc顯示后臺(tái)實(shí)時(shí)數(shù)據(jù)
這篇文章主要為大家詳細(xì)介紹了Echarts+SpringMvc顯示后臺(tái)實(shí)時(shí)數(shù)據(jù),文中示例代碼介紹的非常詳細(xì),具有一定的參考價(jià)值,感興趣的小伙伴們可以參考一下2019-12-12

