RabbitMQ延時(shí)隊(duì)列實(shí)現(xiàn)方法
Linux環(huán)境下(centOS,Rocky)docker部署,rabbitMQ3.8
RabbitMQ延時(shí)隊(duì)列實(shí)現(xiàn)的目的:
延時(shí)隊(duì)列主要應(yīng)用于延時(shí)任務(wù)(訂單的超時(shí)取消和獲取支付服務(wù)的支付狀態(tài))
需要知道什么是死信
死信是什么
死信(Dead Letter Message) 就是 RabbitMQ 中無法正常投遞或消費(fèi)的消息。
它不會(huì)被直接丟棄,而是被 RabbitMQ 自動(dòng)“打上死信標(biāo)簽”,然后重新發(fā)布到你提前配置好的 死信交換機(jī)(Dead Letter Exchange,簡(jiǎn)稱 DLX),再由 DLX 路由到一個(gè)專門的死信隊(duì)列(Dead Letter Queue,簡(jiǎn)稱 DLQ)。
死信觸發(fā)條件
| 序號(hào) | 觸發(fā)條件 | 具體說明 | 常見場(chǎng)景 |
|---|---|---|---|
| 1 | 消費(fèi)者拒絕消息 | 消費(fèi)者調(diào)用 basic.reject 或 basic.nack,并且設(shè)置 requeue=false | 業(yè)務(wù)處理失敗,不想重試 |
| 2 | 消息過期(TTL) | 消息設(shè)置了存活時(shí)間(x-message-ttl 或消息屬性 expiration),到期后 | 延遲消息最常用場(chǎng)景 |
| 3 | 隊(duì)列達(dá)到最大長(zhǎng)度限制 | 隊(duì)列設(shè)置了 x-max-length(最大消息數(shù)),新消息進(jìn)來時(shí)把最老的消息擠出去 | 隊(duì)列爆滿 |
| 4 | 隊(duì)列達(dá)到最大字節(jié)數(shù)限制 | 隊(duì)列設(shè)置了 x-max-length-bytes(最大占用字節(jié)),擠出最老的消息 | 大消息導(dǎo)致隊(duì)列容量超限 |
死信交換機(jī)與RepulishMessageRecoverer區(qū)別
| 維度 | 死信交換機(jī)(DLX) | RepublishMessageRecoverer |
|---|---|---|
| 所屬層級(jí) | RabbitMQ Broker(服務(wù)器端) 原生機(jī)制 | Spring AMQP(應(yīng)用層) 提供的工具類 |
| 觸發(fā)時(shí)機(jī) | 消息成為死信時(shí)(4種情況:拒絕+不重入隊(duì)、TTL過期、隊(duì)列長(zhǎng)度超限、字節(jié)數(shù)超限) | 消費(fèi)者本地重試次數(shù)耗盡后拋出異常時(shí) |
| 觸發(fā)者 | RabbitMQ 服務(wù)器自動(dòng)觸發(fā) | Spring 的 ErrorHandler + MessageRecoverer 觸發(fā) |
| 消息處理方式 | Broker 直接把原消息重新發(fā)布到 DLX | Spring 先 ACK 原消息(告訴 Broker 已消費(fèi)),然后用 RabbitTemplate 重新發(fā)布一份新消息 |
| 是否走 DLX 機(jī)制 | 直接走 DLX | 不走 DLX(因?yàn)橐呀?jīng) ACK 了) |
| 能否攜帶額外信息 | 只能帶 x-death header(記錄幾次死信) | 可以自動(dòng)添加 x-exception-stacktrace、x-exception-message 等豐富異常信息 |
| 靈活性 | 中等(只能配置在隊(duì)列上) | 很高(可以指定任意交換機(jī) + RoutingKey) |
| 適用場(chǎng)景 | 1. TTL 過期(延遲消息) 2. 隊(duì)列超長(zhǎng) 3. 手動(dòng) reject + 不重入隊(duì) | 消費(fèi)業(yè)務(wù)異常后,需要把失敗消息轉(zhuǎn)到專門的錯(cuò)誤隊(duì)列 |
| 配置 | 需要在yaml文件中配置不可重新入隊(duì) default-requeue-rejected: false | 需要在yaml文件中配置不可重新入隊(duì)和最大重試 default-requeue-rejected: false retry: enabled: true max-attempts: 5 |
1.死信隊(duì)列+TTL的實(shí)現(xiàn)
實(shí)現(xiàn)結(jié)構(gòu):

在consumer服務(wù)基于@RabbitListener注解來聲明隊(duì)列、交換機(jī)和綁定隊(duì)列和交換機(jī),并且設(shè)置交換機(jī)為死信交換機(jī):
@RabbitListener(bindings = @QueueBinding(
value=@Queue(name="dlx.queue",durable =" true"),
exchange = @Exchange(name="dlx.direct"),
key = {"hi"}
))
public void listenDlxQueue(String message)throws Exception{
log.info("消費(fèi)者監(jiān)聽到 dlx.queue的消息,{}",message);在normalConfiguration 下通過@Bean聲明任意隊(duì)列和任意交換機(jī)
import org.springframework.amqp.core.*;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
@Configuration
public class NormalConfiguration {
@Bean
public DirectExchange normalExchange(){
return new DirectExchange("normal.direct");
}
@Bean
public Queue normalQueue(){
return QueueBuilder
.durable("normal.queue")
.deadLetterExchange("dlx.direct")
.build();
}
//推薦依賴注入的方式
//需要注意代碼規(guī)范,方法名應(yīng)為小駝峰,反之,會(huì)找不到bean
@Bean
public Binding normalExchangeBinding(Queue normalQueue, DirectExchange normalExchange){
return BindingBuilder
.bind(normalQueue)
.to(normalExchange)
.with("hi");
}
//只能使用當(dāng)前類下的bean
/* @Bean
public Binding normalExchangeBinding(){
return BindingBuilder
.bind(normalQueue())
.to(normalExchange())
.with("hi");
}*/
}如何發(fā)送消息:通過setExpiration方法指定延時(shí)時(shí)長(zhǎng)
@Test
public void testSendDelayMessage() throws Exception {
rabbitTemplate.convertAndSend("normal.direct", "hi", "hello everyone__", message -> {
message.getMessageProperties().setExpiration( "10000");
return message;
});
}2.延時(shí)消息插件
使用死信隊(duì)列可以實(shí)現(xiàn)延遲消息,但這種方法過于繁瑣。為了簡(jiǎn)化這一過程,RabbitMQ的官方推出了一款插件,該插件原生支持延遲消息功能。該插件的運(yùn)作原理是設(shè)計(jì)了一種特殊的交換機(jī),當(dāng)消息投遞到這種交換機(jī)時(shí),它能夠暫存一段時(shí)間,直到達(dá)到設(shè)定的延遲時(shí)間后再將消息投遞到相應(yīng)的隊(duì)列。這種設(shè)計(jì)大大簡(jiǎn)化了延遲消息的處理過程,提高了系統(tǒng)的效率和可靠性。
官方文檔:
https://www.rabbitmq.com/blog/2015/04/16/scheduling-messages-with-rabbitmq
github下載地址:版本為3.8.17
rabbitmq/rabbitmq-delayed-message-exchange: Delayed Messaging for RabbitMQ
安裝插件
由于之前是基于Docker安裝的RabbitMQ,所以需要查看RabbitMQ插件目錄對(duì)應(yīng)的數(shù)據(jù)卷:
docker volume inspect mq-plugins
運(yùn)行結(jié)果:
[
{
"CreatedAt": "2023-12-15T09:57:39+08:00",
"Driver": "local",
"Labels": null,
"Mountpoint": "/var/lib/docker/volumes/mq-plugins/_data",
"Name": "mq-plugins",
"Options": null,
"Scope": "local"
}
]
切換到該數(shù)據(jù)卷的路徑下:
需要把下載的插件文件放在這個(gè)數(shù)據(jù)卷路徑下
cd /var/lib/docker/volumes/mq-plugins/_data
安裝插件docker命令:
docker exec -it mq rabbitmq-plugins enable rabbitmq_delayed_message_exchange
基于注解方式
在consumer服務(wù)基于@RabbitListener注解來聲明隊(duì)列、交換機(jī)和綁定隊(duì)列和交換機(jī),并且設(shè)置交換機(jī)為延遲交換機(jī):
@RabbitListener(bindings = @QueueBinding(
value = @Queue(value = "delay.queue", durable = "true"),
exchange = @Exchange(value = "delay.direct", delayed = "true"),
key = "hidelay"
))
public void listenDelayQueue(String msg) {
log.info("delay.queue:" + msg);
}
//通過 delayed = "true" 聲明為延遲交換機(jī)基于@Bean方式
在consumer服務(wù)基于@Bean注解來聲明交換機(jī)、隊(duì)列和綁定隊(duì)列和交換機(jī),并且設(shè)置交換機(jī)為延遲交換機(jī):
@Configuration
public class DirectConfiguration {
@Bean
public DirectExchange delayExchange() {
return ExchangeBuilder
.directExchange("delay.direct")
.delayed()//這里聲明為延遲交換機(jī)
.durable(true)
.build();
}
@Bean
public Queue delayedQueue() {
return new Queue("delay.queue");
}
@Bean
public Binding delayQueueBinding() {
return BindingBuilder.bind(delayedQueue()).to(delayExchange()).with("delay");
}
}如何發(fā)送消息:通過setDelay方法來設(shè)置延時(shí)時(shí)長(zhǎng)
@Test
public void testSendDelayMessageByplugin() {
rabbitTemplate.convertAndSend("delay.direct", "hidelay", "hello", new MessagePostProcessor() {
@Override
public Message postProcessMessage(Message message) throws AmqpException {
message.getMessageProperties().setDelay(10000);
return message;
}
});
log.info("消息發(fā)送成功");
}總結(jié):
| 維度 | 死信隊(duì)列 + TTL(DLX + TTL) | 延時(shí)消息插件(x-delayed-message) |
|---|---|---|
| 是否需要插件 | ? 不需要(原生功能) | ? 需要安裝 rabbitmq_delayed_message_exchange |
| 實(shí)現(xiàn)原理 | 消息設(shè)置 TTL 過期 → 成為死信 → 路由到 DLX → 進(jìn)入目標(biāo)隊(duì)列 | 發(fā)布消息時(shí)帶 x-delay 頭 → 插件內(nèi)部暫存 → 到期自動(dòng)投遞 |
| 延遲時(shí)間靈活性 | ? 固定延遲(通常每種延遲時(shí)間建一個(gè)隊(duì)列) | ? 支持任意/動(dòng)態(tài)延遲(毫秒級(jí) per-message) |
| 消息存儲(chǔ)位置 | 存放在普通隊(duì)列中(會(huì)占用隊(duì)列資源) | 存放在插件內(nèi)部(Mnesia 表) |
| 延遲精確度 | 一般(尤其是高并發(fā)時(shí)可能有誤差) | 較高(插件定時(shí)器更精確) |
| 消息量支持 | ? 極高(千萬級(jí)輕松支持) | ? 有上限(Mnesia 內(nèi)存/磁盤限制) |
| 性能影響 | 隊(duì)列會(huì)膨脹,內(nèi)存/磁盤壓力大 | 插件專用存儲(chǔ),普通隊(duì)列不膨脹 |
| 實(shí)現(xiàn)復(fù)雜度 | 中等(需配置 DLX、TTL、多個(gè)隊(duì)列) | 簡(jiǎn)單(一個(gè)特殊交換機(jī) + x-delay 參數(shù)) |
| 管理服務(wù)支持 | ? 幾乎所有云廠商都支持 | ? 部分云廠商(阿里云、騰訊云等)不支持插件 |
| 能否取消延遲 | ? 較難 | ? 相對(duì)容易(刪除未到期消息) |
| 適用場(chǎng)景 | 固定延遲、超高并發(fā)、大消息量(如訂單 30 分鐘超時(shí)) | 動(dòng)態(tài)延遲、少量精確延時(shí)(如 7 秒后、3 小時(shí) 15 分后) |
到此這篇關(guān)于RabbitMQ延時(shí)隊(duì)列實(shí)現(xiàn)方法的文章就介紹到這了,更多相關(guān)RabbitMQ延時(shí)隊(duì)列內(nèi)容請(qǐng)搜索腳本之家以前的文章或繼續(xù)瀏覽下面的相關(guān)文章希望大家以后多多支持腳本之家!
- RabbitMQ?延時(shí)隊(duì)列插件安裝與使用示例詳解(基于?Delayed?Message?Plugin)
- Springboot使用Rabbitmq的延時(shí)隊(duì)列+死信隊(duì)列實(shí)現(xiàn)消息延期消費(fèi)
- RabbitMQ延時(shí)消息隊(duì)列在golang中的使用詳解
- 關(guān)于Rabbitmq死信隊(duì)列及延時(shí)隊(duì)列的實(shí)現(xiàn)
- 使用Java實(shí)現(xiàn)RabbitMQ延時(shí)隊(duì)列
- 如何利用rabbitMq的死信隊(duì)列實(shí)現(xiàn)延時(shí)消息
- .NETCore基于RabbitMQ實(shí)現(xiàn)延時(shí)隊(duì)列的兩方法
- Docker安裝RabbitMQ并安裝延時(shí)隊(duì)列插件
- 運(yùn)用.NetCore實(shí)例講解RabbitMQ死信隊(duì)列,延時(shí)隊(duì)列
相關(guān)文章
Java中使用Files類的copy()方法實(shí)現(xiàn)復(fù)制文件
這篇文章主要介紹了Java中使用Files類的copy()方法實(shí)現(xiàn)復(fù)制文件,具有很好的參考價(jià)值,希望對(duì)大家有所幫助,如有錯(cuò)誤或未考慮完全的地方,望不吝賜教2025-05-05
mybatis實(shí)現(xiàn)表與對(duì)象的關(guān)聯(lián)關(guān)系_動(dòng)力節(jié)點(diǎn)Java學(xué)院整理
這篇文章主要介紹了mybatis實(shí)現(xiàn)表與對(duì)象的關(guān)聯(lián)關(guān)系_動(dòng)力節(jié)點(diǎn)Java學(xué)院整理,需要的朋友可以參考下2017-09-09
Windows下Java+MyBatis框架+MySQL的開發(fā)環(huán)境搭建教程
這篇文章主要介紹了Windows下Java+MyBatis框架+MySQL的開發(fā)環(huán)境搭建教程,Mybatis對(duì)普通SQL語句的支持非常好,需要的朋友可以參考下2016-04-04
springboot2.0如何通過fastdfs實(shí)現(xiàn)文件分布式上傳
這篇文章主要介紹了springboot2.0如何通過fastdfs實(shí)現(xiàn)文件分布式上傳,文中通過示例代碼介紹的非常詳細(xì),對(duì)大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價(jià)值,需要的朋友可以參考下2019-12-12
java實(shí)現(xiàn)對(duì)Hadoop的操作
這篇文章主要介紹了java實(shí)現(xiàn)對(duì)Hadoop的操作,通過非常完整詳細(xì)的代碼展示了如何去進(jìn)行一系列操作,包括基本操作,文件讀寫,需要的朋友可以參考下2021-07-07
Java多線程教程之如何利用Future實(shí)現(xiàn)攜帶結(jié)果的任務(wù)
Callable與Future兩功能是Java?5版本中加入的,這篇文章主要給大家介紹了關(guān)于Java多線程教程之如何利用Future實(shí)現(xiàn)攜帶結(jié)果任務(wù)的相關(guān)資料,需要的朋友可以參考下2021-12-12

