最新国产好看的视频,伊人天堂AV在线,国产Aaaaaa视频,蜜臀视频在线观看一区,人妻av色图,密臀久久久精品影片,青青视频免费观看毛片,久草在线观看视,国产三级精品色情在线

RabbitMQ延時(shí)隊(duì)列實(shí)現(xiàn)方法

 更新時(shí)間:2026年05月18日 10:33:23   作者:星辰大海2024  
文章主要介紹了在Linux環(huán)境下使用CentOS/Rocky使用Docker部署RabbitMQ 3.8版本,并實(shí)現(xiàn)現(xiàn)現(xiàn)了延時(shí)隊(duì)列的實(shí)現(xiàn)方式,總結(jié)了兩種方法在靈活性、性能和管理上的優(yōu)缺點(diǎ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.rejectbasic.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ā)布到 DLXSpring 先 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)文章希望大家以后多多支持腳本之家!

相關(guān)文章

最新評(píng)論

凤庆县| 新乐市| 杭州市| 都兰县| 内乡县| 诸暨市| 定州市| 翼城县| 旬阳县| 新竹市| 大英县| 扎兰屯市| 瑞安市| 定襄县| 岗巴县| 长兴县| 错那县| 望都县| 和田县| 右玉县| 长沙市| 靖西县| 高陵县| 昌邑市| 九龙坡区| 洮南市| 汾西县| 梅河口市| 原阳县| 加查县| 荣成市| 罗源县| 磐石市| 桦川县| 黑山县| 洛扎县| 新蔡县| 穆棱市| 水富县| 申扎县| 德化县|