中文字幕av专区_日韩电影在线播放_精品国产精品久久一区免费式_av在线免费观看网站

溫馨提示×

溫馨提示×

您好,登錄后才能下訂單哦!

密碼登錄×
登錄注冊×
其他方式登錄
點擊 登錄注冊 即表示同意《億速云用戶服務條款》

SpringBoot中如何基于RabbitMQ實現消息延遲隊列

發布時間:2021-09-19 17:27:47 來源:億速云 閱讀:197 作者:小新 欄目:大數據

小編給大家分享一下SpringBoot中如何基于RabbitMQ實現消息延遲隊列,相信大部分人都還不怎么了解,因此分享這篇文章給大家參考一下,希望大家閱讀完這篇文章后大有收獲,下面讓我們一起去了解一下吧!

延時隊列使用場景

>在很多的業務場景中,延時隊列可以實現很多功能,此類業務中,一般上是非實時的,需要延遲處理的,需要進行重試補償的。

  1. 訂單超時關閉:在支付場景中,一般上訂單在創建后30分鐘或1小時內未支付的,會自動取消訂單。

  2. 短信或者郵件通知:在一些注冊或者下單業務時,需要在1分鐘或者特定時間后進行短信或者郵件發送相關資料的。本身此類業務于主業務是無關聯性的,一般上的做法是進行異步發送。

  3. 重試場景:比如消息通知,在第一次通知出現異常時,會在隔幾分鐘之后進行再次重試發送。

RabbitMQ實現延時隊列

>本身在RabbitMQ中是未直接提供延時隊列功能的,但可以使用TTL(Time-To-Live,存活時間)DLX(Dead-Letter-Exchange,死信隊列交換機)的特性實現延時隊列的功能。

存活時間(Time-To-Live 簡稱 TTL)

>RabbitMQ中可以對隊列和消息分別設置TTL,TTL表明了一條消息可在隊列中存活的最大時間。當某條消息被設置了TTL或者當某條消息進入了設置了TTL的隊列時,這條消息會在TTL時間后**死亡成為Dead Letter**。如果既配置了消息的TTL,又配置了隊列的TTL,那么較小的那個值會被取用。

死信交換(Dead Letter Exchanges 簡稱 DLX)

>上個知識點也提到了,設置了TTL的消息或隊列最終會成為Dead Letter,當消息在一個隊列中變成死信之后,它能被重新發送到另一個交換機中,這個交換機就是DLX,綁定此DLX的隊列就是死信隊列。

一個消息變成死信一般上是由于以下幾種情況;

  1. 消息被拒絕

  2. 消息過期

  3. 隊列達到了最大長度。

所以,通過TTLDLX的特性可以模擬實現延時隊列的功能。當隊列中的消息超時成為死信后,會把消息死信重新發送到配置好的交換機中,然后分發到真實的消費隊列。故簡單來說,我們可以創建2個隊列,一個隊列用于發送消息,一個隊列用于消息過期后的轉發的目標隊列。

SpringBoot集成RabbitMQ實現延時隊列實戰

>以下使用SpringBoot集成RabbitMQ進行實戰說明,在進行http消息通知時,若通知失敗(地址不可用或者連接超時)時,將此消息轉入延時隊列中,待特定時間后進行重新發送。

0.引入pom依賴

    <!-- rabbit -->
    <dependency>
        <groupid>org.springframework.boot</groupid>
        <artifactid>spring-boot-starter-amqp</artifactid>
    </dependency>
    <!-- 簡化http操作 -->
    <dependency>
        <groupid>cn.hutool</groupid>
        <artifactid>hutool-http</artifactid>
        <version>4.5.16</version>
    </dependency>
    <dependency>
        <groupid>cn.hutool</groupid>
        <artifactid>hutool-json</artifactid>
        <version>4.5.16</version>
    </dependency>
    <dependency>
        <groupid>org.springframework.boot</groupid>
        <artifactid>spring-boot-starter-web</artifactid>
    </dependency>

1.編寫rabbitmq配置文件(關鍵配置) RabbitConfig.java

/** 
*
* @ClassName   類名:RabbitConfig 
* @Description 功能說明:
* <p>
* TODO
*</p>
************************************************************************
* @date        創建日期:2019年7月17日
* @version     版本號:V1.0
*<p>
***************************修訂記錄*************************************
* 
*   2019年7月17日   oKong   創建該類功能。
*
***********************************************************************
*</p>
*/
@Configuration
public class RabbitConfig {
    
    @Autowired
    ConnectionFactory connectionFactory;
    
    /**
     * 消費者線程數 設置大點 大概率是能通知到的
     */
    @Value("${http.notify.concurrency:50}")
    int concurrency;
    
    /**
     * 延遲隊列的消費者線程數 可設置小點
     */
    @Value("${http.notify.delay.concurrency:20}")
    int delayConcurrency;
    
    @Bean
    public RabbitAdmin rabbitAdmin() {
        return new RabbitAdmin(connectionFactory);
    }
    
    @Bean
    public DirectExchange httpMessageNotifyDirectExchange(RabbitAdmin rabbitAdmin) {
        //durable 是否持久化
        //autoDelete 是否自動刪除,即服務端或者客服端下線后 交換機自動刪除
        DirectExchange directExchange = new DirectExchange(ApplicationConstant.HTTP_MESSAGE_EXCHANGE,true,false);
        directExchange.setAdminsThatShouldDeclare(rabbitAdmin);
        return directExchange;
    }
    
    //設置消息隊列
    @Bean
    public Queue httpMessageStartQueue(RabbitAdmin rabbitAdmin) {
         
       /*
                       創建接收隊列,4個參數
         name - 隊列名稱
         durable - false,不進行持有化
         exclusive - true,獨占性
         autoDelete - true,自動刪除*/
        Queue queue = new Queue(ApplicationConstant.HTTP_MESSAGE_START_QUEUE_NAME, true, false, false);
        queue.setAdminsThatShouldDeclare(rabbitAdmin);
        return queue;
    }
    
    //隊列綁定交換機
    @Bean
    public Binding bindingStartQuene(RabbitAdmin rabbitAdmin,DirectExchange httpMessageNotifyDirectExchange, Queue httpMessageStartQueue) {
        Binding binding = BindingBuilder.bind(httpMessageStartQueue).to(httpMessageNotifyDirectExchange).with(ApplicationConstant.HTTP_MESSAGE_START_RK);
        binding.setAdminsThatShouldDeclare(rabbitAdmin);
        return binding;
    }
    
    @Bean
    public Queue httpMessageOneQueue(RabbitAdmin rabbitAdmin) {
        Queue queue = new Queue(ApplicationConstant.HTTP_MESSAGE_ONE_QUEUE_NAME, true, false, false);
        queue.setAdminsThatShouldDeclare(rabbitAdmin);
        return queue;
    }
    
    @Bean
    public Binding bindingOneQuene(RabbitAdmin rabbitAdmin,DirectExchange httpMessageNotifyDirectExchange, Queue httpMessageOneQueue) {
        Binding binding = BindingBuilder.bind(httpMessageOneQueue).to(httpMessageNotifyDirectExchange).with(ApplicationConstant.HTTP_MESSAGE_ONE_RK);
        binding.setAdminsThatShouldDeclare(rabbitAdmin);
        return binding;
    }
    
    //-------------設置延遲隊列--開始--------------------
    @Bean
    public Queue httpDelayOneQueue() {
        //name - 隊列名稱
        //durable - true
        //exclusive - false
        //autoDelete - false
        return QueueBuilder.durable("http.message.dlx.one")
                //以下是重點:當變成死信隊列時,會轉發至 路由為x-dead-letter-exchange及x-dead-letter-routing-key的隊列中
                .withArgument("x-dead-letter-exchange", ApplicationConstant.HTTP_MESSAGE_EXCHANGE)
                .withArgument("x-dead-letter-routing-key", ApplicationConstant.HTTP_MESSAGE_ONE_RK)
                .withArgument("x-message-ttl", 1*60*1000)//1分鐘 過期時間(單位:毫秒),當過期后 會變成死信隊列,之后進行轉發
                .build();
    }
    //綁定到交換機上
    @Bean
    public Binding bindingDelayOneQuene(RabbitAdmin rabbitAdmin, DirectExchange httpMessageNotifyDirectExchange, Queue httpDelayOneQueue) {
        Binding binding = BindingBuilder.bind(httpDelayOneQueue).to(httpMessageNotifyDirectExchange).with("delay.one");
        binding.setAdminsThatShouldDeclare(rabbitAdmin);
        return binding;
    }
    //-------------設置延遲隊列--結束--------------------

    //建議將正常的隊列和延遲處理的隊列分開
    //設置監聽容器
    @Bean("notifyListenerContainer")
    public SimpleRabbitListenerContainerFactory httpNotifyListenerContainer() {
        SimpleRabbitListenerContainerFactory factory = new SimpleRabbitListenerContainerFactory();
        factory.setAcknowledgeMode(AcknowledgeMode.MANUAL);// 手動ack
        factory.setConnectionFactory(connectionFactory);
        factory.setPrefetchCount(1);
        factory.setConcurrentConsumers(concurrency);
        return factory;
    }

    // 設置監聽容器
    @Bean("delayNotifyListenerContainer")
    public SimpleRabbitListenerContainerFactory httpDelayNotifyListenerContainer() {
        SimpleRabbitListenerContainerFactory factory = new SimpleRabbitListenerContainerFactory();
        factory.setAcknowledgeMode(AcknowledgeMode.MANUAL);// 手動ack
        factory.setConnectionFactory(connectionFactory);
        factory.setPrefetchCount(1);
        factory.setConcurrentConsumers(delayConcurrency);
        return factory;
    }
}

ApplicationConstant.java

public class ApplicationConstant {
    
    /**
     * 發送http通知的 exchange 隊列
     */
    public static final String HTTP_MESSAGE_EXCHANGE = "http.message.exchange";
    
    /**
     * 配置消息隊列和路由key值
     */
    public static final String HTTP_MESSAGE_START_QUEUE_NAME = "http.message.start";
    public static final String HTTP_MESSAGE_START_RK = "rk.start";
    
    public static final String HTTP_MESSAGE_ONE_QUEUE_NAME = "http.message.one";
    public static final String HTTP_MESSAGE_ONE_RK = "rk.one";
    
    /**
     * 通知隊列對應的延遲隊列關系,即過期隊列之后發送到下一個的隊列信息,可以根據實際情況添加,當然也可以根據一定規則自動生成
     */
    public static final Map<string,string> delayRefMap = new HashMap<string, string>() {
        /**
         * 
         */
        private static final long serialVersionUID = -779823216035682493L;

        {
            put(HTTP_MESSAGE_START_QUEUE_NAME, "delay.one");
        }
    };    
}

簡單來說,就是創建一個正常消息發送隊列,用于接收http消息請求的參數,同時進行http請求。同時,創建一個延時隊列,設置其x-dead-letter-exchangex-dead-letter-routing-keyx-message-ttl值,將其轉發到正常的隊列中。使用一個map對象維護一個關系,當正常消息異常時,需要發送的延時隊列的隊列名稱,當然時間場景匯總,根據需要可以進行動態配置或者根據一定規則進行動態映射。

2.創建監聽類,用于消息的消費操作,此處使用@RabbitListener來消費消息(當然也可以使用SimpleMessageListenerContainer進行消息配置的),創建了一個正常消息監聽和延時隊列監聽,由于一般上異常通知是低概率事件,可根據不同的監聽容器進行差異化配置。

/** 
*
* @ClassName   類名:HttpMessagerLister 
* @Description 功能說明:http通知消費監聽接口
* <p>
* TODO
*</p>
************************************************************************
* @date        創建日期:2019年7月17日
* @version     版本號:V1.0
*<p>
***************************修訂記錄*************************************
* 
*   2019年7月17日   oKong   創建該類功能。
*
***********************************************************************
*</p>
*/
@Component
@Slf4j
public class HttpMessagerLister {
    
    @Autowired
    HttpMessagerService messagerService;
    
    @RabbitListener(id = "httpMessageNotifyConsumer", queues = {ApplicationConstant.HTTP_MESSAGE_START_QUEUE_NAME}, containerFactory = "notifyListenerContainer")
    public void httpMessageNotifyConsumer(Message message, Channel channel) throws Exception {
        doHandler(message, channel);
    }
    
    @RabbitListener(id= "httpDelayMessageNotifyConsumer", queues = {
            ApplicationConstant.HTTP_MESSAGE_ONE_QUEUE_NAME,}, containerFactory = "delayNotifyListenerContainer")
    public void httpDelayMessageNotifyConsumer(Message message, Channel channel) throws Exception {
        doHandler(message, channel);
    }
    
    private void doHandler(Message message, Channel channel) throws Exception {
        String body = new String(message.getBody(),"utf-8");
        String queue = message.getMessageProperties().getConsumerQueue();
        log.info("接收到通知請求:{},隊列名:{}",body, queue);
        //消息對象轉換
        try {
            HttpEntity httpNotifyDto = JSONUtil.toBean(body, HttpEntity.class);
            channel.basicAck(message.getMessageProperties().getDeliveryTag(), false);
            //發送通知
            messagerService.notify(queue, httpNotifyDto);
        } catch(Exception e) {
            log.error(e.getMessage());
            //ack
            channel.basicAck(message.getMessageProperties().getDeliveryTag(), false);
        }
    }
}

HttpMessagerService.java:消息真正處理的類,此類是關鍵,這里未進行日志記錄,真實場景中,強烈建議進行消息通知的日志存儲,防止日后信息的查看,同時也能通過發送狀態,在重試次數都失敗后,進行定時再次發送功能,同時也有據可查。

@Component
@Slf4j
public class HttpMessagerService {
    
    @Autowired
    AmqpTemplate mqTemplate;    
    
    public void notify(String queue,HttpEntity httpEntity) {
        //發起請求
        log.info("開始發起http請求:{}", httpEntity);
        try {
            switch(httpEntity.getMethod().toLowerCase()) {
            case "POST":
                  HttpUtil.post(httpEntity.getUrl(), httpEntity.getParams());
                  break;
            case "GET":
            default:
                HttpUtil.get(httpEntity.getUrl(), httpEntity.getParams());
            }
        } catch (Exception e) {
            //發生異常,放入延遲隊列中
            String nextRk = ApplicationConstant.delayRefMap.get(queue);
            if(ApplicationConstant.HTTP_MESSAGE_ONE_QUEUE_NAME.equals(queue)) {
                //若已經是最后一個延遲隊列的消息隊列了,則后續可直接放入數據庫中 待后續定時策略進行再次發送
                log.warn("http通知已經通知N次失敗,進入定時進行發起通知,url={}", httpEntity.getUrl());
            } else {
               log.warn("http重新發送通知:{}, 通知隊列rk為:{}, 原隊列:{}", httpEntity.getUrl(), nextRk, queue);
               mqTemplate.convertAndSend(ApplicationConstant.HTTP_MESSAGE_EXCHANGE, nextRk, cn.hutool.json.JSONUtil.toJsonStr(httpEntity));
            }
        }
    }
}

3.創建控制層服務(真實場景中,如SpringCloud微服務中,一般上是創建個api接口,供其他服務進行調用)

@Slf4j
@RestController
@Api(tags = "http測試接口")
public class HttpDemoController {

    @Autowired
    AmqpTemplate mqTemplate;
    
    @PostMapping("/send")
    @ApiOperation(value="send",notes = "發送http測試")
    public String sendHttp(@RequestBody HttpEntity httpEntity) {
        //發送http請求
        log.info("開始發起http請求,發布異步消息:{}", httpEntity);
        mqTemplate.convertAndSend(ApplicationConstant.HTTP_MESSAGE_EXCHANGE, ApplicationConstant.HTTP_MESSAGE_START_RK, cn.hutool.json.JSONUtil.toJsonStr(httpEntity));
        return "發送成功:url=" + httpEntity.getUrl();        
    }
}

4.配置文件添加RabbitMQ相關配置信息

spring.rabbitmq.host=127.0.0.1
spring.rabbitmq.port=5672
spring.rabbitmq.username=guest
spring.rabbitmq.password=guest
spring.rabbitmq.virtual-host=/

# 通知-消費者線程數 設置大點 大概率是能通知到的
http.notify.concurrency=150
# 延遲隊列的消費者線程數 可設置小點 
http.notify.delay.concurrency=10

5.編寫啟動類。

@SpringBootApplication
@Slf4j
public class DelayQueueApplication {

    public static void main(String[] args) throws Exception {
        SpringApplication.run(DelayQueueApplication.class, args);
        log.info("spring-boot-rabbitmq-delay-queue-chapter38服務啟動!");
    }
}

6.啟動服務。使用swagger進行簡單調用測試。

  • 正常通知:

SpringBoot中如何基于RabbitMQ實現消息延遲隊列

2019-07-20 23:52:23.792  INFO 65216 --- [nio-8080-exec-1] c.l.l.s.c.controller.HttpDemoController  : 開始發起http請求,發布異步消息:HttpEntity(url=www.baidu.com, params={a=1}, method=get)
2019-07-20 23:52:23.794  INFO 65216 --- [TaskExecutor-97] c.l.l.s.chapter38.mq.HttpMessagerLister  : 接收到通知請求:{"method":"get","params":{"a":1},"url":"www.baidu.com"},隊列名:http.message.start
2019-07-20 23:52:23.794  INFO 65216 --- [TaskExecutor-97] c.l.l.s.c.service.HttpMessagerService    : 開始發起http請求:HttpEntity(url=www.baidu.com, params={a=1}, method=get)
  • 異常通知:訪問一個不存在的地址

SpringBoot中如何基于RabbitMQ實現消息延遲隊列

2019-07-20 23:53:14.699  INFO 65216 --- [nio-8080-exec-4] c.l.l.s.c.controller.HttpDemoController  : 開始發起http請求,發布異步消息:HttpEntity(url=www.baidu.com1, params={a=1}, method=get)
2019-07-20 23:53:14.705  INFO 65216 --- [TaskExecutor-84] c.l.l.s.chapter38.mq.HttpMessagerLister  : 接收到通知請求:{"method":"get","params":{"a":1},"url":"www.baidu.com1"},隊列名:http.message.start
2019-07-20 23:53:14.705  INFO 65216 --- [TaskExecutor-84] c.l.l.s.c.service.HttpMessagerService    : 開始發起http請求:HttpEntity(url=www.baidu.com1, params={a=1}, method=get)
2019-07-20 23:53:14.706  WARN 65216 --- [TaskExecutor-84] c.l.l.s.c.service.HttpMessagerService    : http重新發送通知:www.baidu.com1, 通知隊列rk為:delay.one, 原隊列:http.message.start

RabbitMQ后臺中,可以看見http.message.dlx.one隊列中存在這需要延時處理的消息,在一分鐘后會轉發至http.message.one隊列中。

SpringBoot中如何基于RabbitMQ實現消息延遲隊列

在一分鐘后,可以看見消息本再次消費了。

2019-07-20 23:54:14.722  INFO 65216 --- [TaskExecutor-16] c.l.l.s.chapter38.mq.HttpMessagerLister  : 接收到通知請求:{"method":"get","params":{"a":1},"url":"www.baidu.com1"},隊列名:http.message.one
2019-07-20 23:54:14.723  INFO 65216 --- [TaskExecutor-16] c.l.l.s.c.service.HttpMessagerService    : 開始發起http請求:HttpEntity(url=www.baidu.com1, params={a=1}, method=get)
2019-07-20 23:54:14.723  WARN 65216 --- [TaskExecutor-16] c.l.l.s.c.service.HttpMessagerService    : http通知已經通知N次失敗,進入定時進行發起通知,url=www.baidu.com1

一些最佳實踐

>在正式場景中,一般上補償或者重試機制大概率是不會發送的,倘若發生時,一般上是第三方業務系統出現了問題,故一般上在進行補充時,應該在非高峰期進行操作,故應該對延時監聽器,應該在高峰期時停止消費,在非高峰期時進行消費。同時,還可以根據不同的通知類型,放入不一樣的延時隊列中,保障業務的正常。這里簡單說明下,動態停止或者啟動演示監聽器的方式。一般上是使用RabbitListenerEndpointRegistry對象獲取延時監聽器,之后進行動態停止或者啟用。可設置@RabbitListener的id屬性,直接進行獲取,當然也可以直接獲取所有的監聽器,進行自定義判斷了。

    @Autowired
    RabbitListenerEndpointRegistry registry;
    
   @GetMapping("/set")
    @ApiOperation(value = "set", notes = "設置消息監聽器的狀態")
    public String setSimpleMessageListenerContainer(String status) {
        if("1".equals(status)) {
            registry.getListenerContainer("httpDelayMessageNotifyConsumer").start();
        } else {
            registry.getListenerContainer("httpDelayMessageNotifyConsumer").stop();
        }
        return status;
    }

這里,只是簡單進行演示說明,在真實場景下,可以使用定時器,判斷當前是否為高峰期,進而進行動態設置監聽器的狀態。

以上是“SpringBoot中如何基于RabbitMQ實現消息延遲隊列”這篇文章的所有內容,感謝各位的閱讀!相信大家都有了一定的了解,希望分享的內容對大家有所幫助,如果還想學習更多知識,歡迎關注億速云行業資訊頻道!

向AI問一下細節

免責聲明:本站發布的內容(圖片、視頻和文字)以原創、轉載和分享為主,文章觀點不代表本網站立場,如果涉及侵權請聯系站長郵箱:is@yisu.com進行舉報,并提供相關證據,一經查實,將立刻刪除涉嫌侵權內容。

AI

镇坪县| 崇左市| 宁安市| 惠安县| 海丰县| 龙南县| 饶河县| 息烽县| 宝应县| 师宗县| 大连市| 嘉禾县| 曲阳县| 句容市| 南靖县| 花垣县| 秦安县| 杭锦后旗| 璧山县| 固原市| 隆化县| 合肥市| 巫溪县| 乌什县| 乐都县| 稻城县| 湾仔区| 金沙县| 海林市| 读书| 元朗区| 达州市| 茂名市| 蓬莱市| 本溪市| 明水县| 阿勒泰市| 四会市| 伊金霍洛旗| 洛阳市| 察哈|