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

RabbitMQ延遲隊列及消息延遲推送實(shí)現(xiàn)詳解

 更新時間:2019年12月10日 10:32:28   作者:海向  
這篇文章主要介紹了RabbitMQ延遲隊列及消息延遲推送實(shí)現(xiàn)詳解,文中通過示例代碼介紹的非常詳細(xì),對大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價值,需要的朋友可以參考下

這篇文章主要介紹了RabbitMQ延遲隊列及消息延遲推送實(shí)現(xiàn)詳解,文中通過示例代碼介紹的非常詳細(xì),對大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價值,需要的朋友可以參考下

應(yīng)用場景

目前常見的應(yīng)用軟件都有消息的延遲推送的影子,應(yīng)用也極為廣泛,例如:

  • 淘寶七天自動確認(rèn)收貨。在我們簽收商品后,物流系統(tǒng)會在七天后延時發(fā)送一個消息給支付系統(tǒng),通知支付系統(tǒng)將款打給商家,這個過程持續(xù)七天,就是使用了消息中間件的延遲推送功能。
  • 12306 購票支付確認(rèn)頁面。我們在選好票點(diǎn)擊確定跳轉(zhuǎn)的頁面中往往都會有倒計時,代表著 30 分鐘內(nèi)訂單不確認(rèn)的話將會自動取消訂單。其實(shí)在下訂單那一刻開始購票業(yè)務(wù)系統(tǒng)就會發(fā)送一個延時消息給訂單系統(tǒng),延時30分鐘,告訴訂單系統(tǒng)訂單未完成,如果我們在30分鐘內(nèi)完成了訂單,則可以通過邏輯代碼判斷來忽略掉收到的消息。

在上面兩種場景中,如果我們使用下面兩種傳統(tǒng)解決方案無疑大大降低了系統(tǒng)的整體性能和吞吐量:

  • 使用 redis 給訂單設(shè)置過期時間,最后通過判斷 redis 中是否還有該訂單來決定訂單是否已經(jīng)完成。這種解決方案相較于消息的延遲推送性能較低,因為我們知道 redis 都是存儲于內(nèi)存中,我們遇到惡意下單或者刷單的將會給內(nèi)存帶來巨大壓力。
  • 使用傳統(tǒng)的數(shù)據(jù)庫輪詢來判斷數(shù)據(jù)庫表中訂單的狀態(tài),這無疑增加了IO次數(shù),性能極低。
  • 使用 jvm 原生的 DelayQueue ,也是大量占用內(nèi)存,而且沒有持久化策略,系統(tǒng)宕機(jī)或者重啟都會丟失訂單信息。

消息延遲推送的實(shí)現(xiàn)

在 RabbitMQ 3.6.x 之前我們一般采用死信隊列+TTL過期時間來實(shí)現(xiàn)延遲隊列,我們這里不做過多介紹,可以參考之前文章來了解:TTL、死信隊列

在 RabbitMQ 3.6.x 開始,RabbitMQ 官方提供了延遲隊列的插件,可以下載放置到 RabbitMQ 根目錄下的 plugins 下。延遲隊列插件下載

首先我們創(chuàng)建交換機(jī)和消息隊列

import org.springframework.amqp.core.*;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;

import java.util.HashMap;
import java.util.Map;

@Configuration
public class MQConfig {

  public static final String LAZY_EXCHANGE = "Ex.LazyExchange";
  public static final String LAZY_QUEUE = "MQ.LazyQueue";
  public static final String LAZY_KEY = "lazy.#";

  @Bean
  public TopicExchange lazyExchange(){
    //Map<String, Object> pros = new HashMap<>();
    //設(shè)置交換機(jī)支持延遲消息推送
    //pros.put("x-delayed-message", "topic");
    TopicExchange exchange = new TopicExchange(LAZY_EXCHANGE, true, false, pros);
    exchange.setDelayed(true);
    return exchange;
  }

  @Bean
  public Queue lazyQueue(){
    return new Queue(LAZY_QUEUE, true);
  }

  @Bean
  public Binding lazyBinding(){
    return BindingBuilder.bind(lazyQueue()).to(lazyExchange()).with(LAZY_KEY);
  }
}

我們在 Exchange 的聲明中可以設(shè)置exchange.setDelayed(true)來開啟延遲隊列,也可以設(shè)置為以下內(nèi)容傳入交換機(jī)聲明的方法中,因為第一種方式的底層就是通過這種方式來實(shí)現(xiàn)的。

    //Map<String, Object> pros = new HashMap<>();
    //設(shè)置交換機(jī)支持延遲消息推送
    //pros.put("x-delayed-message", "topic");
    TopicExchange exchange = new TopicExchange(LAZY_EXCHANGE, true, false, pros);

發(fā)送消息時我們需要指定延遲推送的時間,我們這里在發(fā)送消息的方法中傳入?yún)?shù) new MessagePostProcessor() 是為了獲得 Message對象,因為需要借助 Message對象的api 來設(shè)置延遲時間。

import com.anqi.mq.config.MQConfig;
import org.springframework.amqp.AmqpException;
import org.springframework.amqp.core.Message;
import org.springframework.amqp.core.MessageDeliveryMode;
import org.springframework.amqp.core.MessagePostProcessor;
import org.springframework.amqp.rabbit.connection.CorrelationData;
import org.springframework.amqp.rabbit.core.RabbitTemplate;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.stereotype.Component;

import java.util.Date;

@Component
public class MQSender {

  @Autowired
  private RabbitTemplate rabbitTemplate;

  //confirmCallback returnCallback 代碼省略,請參照上一篇
 
  public void sendLazy(Object message){
    rabbitTemplate.setMandatory(true);
    rabbitTemplate.setConfirmCallback(confirmCallback);
    rabbitTemplate.setReturnCallback(returnCallback);
    //id + 時間戳 全局唯一
    CorrelationData correlationData = new CorrelationData("12345678909"+new Date());

    //發(fā)送消息時指定 header 延遲時間
    rabbitTemplate.convertAndSend(MQConfig.LAZY_EXCHANGE, "lazy.boot", message,
        new MessagePostProcessor() {
      @Override
      public Message postProcessMessage(Message message) throws AmqpException {
        //設(shè)置消息持久化
        message.getMessageProperties().setDeliveryMode(MessageDeliveryMode.PERSISTENT);
        //message.getMessageProperties().setHeader("x-delay", "6000");
        message.getMessageProperties().setDelay(6000);
        return message;
      }
    }, correlationData);
  }
}

我們可以觀察 setDelay(Integer i)底層代碼,也是在 header 中設(shè)置 x-delay。等同于我們手動設(shè)置 header

message.getMessageProperties().setHeader("x-delay", "6000");

/**
 * Set the x-delay header.
 * @param delay the delay.
 * @since 1.6
 */
public void setDelay(Integer delay) {
  if (delay == null || delay < 0) {
    this.headers.remove(X_DELAY);
  }
  else {
    this.headers.put(X_DELAY, delay);
  }
}

消費(fèi)端進(jìn)行消費(fèi)

import com.rabbitmq.client.Channel;
import org.springframework.amqp.rabbit.annotation.*;
import org.springframework.amqp.support.AmqpHeaders;
import org.springframework.stereotype.Component;

import java.io.IOException;
import java.util.Map;

@Component
public class MQReceiver {

  @RabbitListener(queues = "MQ.LazyQueue")
  @RabbitHandler
  public void onLazyMessage(Message msg, Channel channel) throws IOException{
    long deliveryTag = msg.getMessageProperties().getDeliveryTag();
    channel.basicAck(deliveryTag, true);
    System.out.println("lazy receive " + new String(msg.getBody()));

  }

測試結(jié)果

import org.junit.Test;
import org.junit.runner.RunWith;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.boot.test.context.SpringBootTest;
import org.springframework.test.context.junit4.SpringRunner;

@SpringBootTest
@RunWith(SpringRunner.class)
public class MQSenderTest {

  @Autowired
  private MQSender mqSender;

  @Test
  public void sendLazy() throws Exception {
    String msg = "hello spring boot";

    mqSender.sendLazy(msg + ":");
  }
}

果然在 6 秒后收到了消息 lazy receive hello spring boot:

以上就是本文的全部內(nèi)容,希望對大家的學(xué)習(xí)有所幫助,也希望大家多多支持腳本之家。

相關(guān)文章

  • Java對象級別與類級別的同步鎖synchronized語法示例

    Java對象級別與類級別的同步鎖synchronized語法示例

    這篇文章主要為大家介紹了Java對象級別與類級別的同步鎖synchronized語法示例,有需要的朋友可以借鑒參考下,希望能夠有所幫助,祝大家多多進(jìn)步
    2022-03-03
  • Spring Boot應(yīng)用啟動時自動執(zhí)行代碼的五種方式(常見方法)

    Spring Boot應(yīng)用啟動時自動執(zhí)行代碼的五種方式(常見方法)

    Spring Boot為開發(fā)者提供了多種方式在應(yīng)用啟動時執(zhí)行自定義代碼,這些方式包括注解、接口實(shí)現(xiàn)和事件監(jiān)聽器,本文我們將探討一些常見的方法,以及如何利用它們在應(yīng)用啟動時執(zhí)行初始化邏輯,感興趣的朋友一起看看吧
    2024-04-04
  • 使用SpringBoot內(nèi)置web服務(wù)器

    使用SpringBoot內(nèi)置web服務(wù)器

    這篇文章主要介紹了使用SpringBoot內(nèi)置web服務(wù)器操作,具有很好的參考價值,希望對大家有所幫助。如有錯誤或未考慮完全的地方,望不吝賜教
    2021-09-09
  • Spring解決循環(huán)依賴問題及三級緩存的作用

    Spring解決循環(huán)依賴問題及三級緩存的作用

    這篇文章主要介紹了Spring解決循環(huán)依賴問題及三級緩存的作用,所謂的三級緩存只是三個可以當(dāng)作是全局變量的Map,Spring的源碼中大量使用了這種先將數(shù)據(jù)放入容器中等使用結(jié)束再銷毀的代碼風(fēng)格
    2022-07-07
  • 解決配置Feign時報錯PathVariable annotation was empty on param 0.

    解決配置Feign時報錯PathVariable annotation was empty

    在配置Feign客戶端時,如果遇到`@PathVariable`注解為空的問題,是因為在聲明接口方法時沒有為`@PathVariable`注解提供`value`屬性,解決方法是為`@PathVariable`注解添加`value`屬性,這樣就可以避免報錯,并成功啟動Feign客戶端
    2024-11-11
  • java原碼補(bǔ)碼反碼關(guān)系解析

    java原碼補(bǔ)碼反碼關(guān)系解析

    這篇文章主要為大家詳細(xì)介紹了java原碼補(bǔ)碼反碼的關(guān)系,文中示例代碼介紹的非常詳細(xì),具有一定的參考價值,感興趣的小伙伴們可以參考一下
    2020-02-02
  • 使用Backoff策略提高HttpClient連接管理的效率

    使用Backoff策略提高HttpClient連接管理的效率

    這篇文章主要為大家介紹了Backoff策略提高HttpClient連接管理的效率使用解析,有需要的朋友可以借鑒參考下,希望能夠有所幫助,祝大家多多進(jìn)步,早日升職加薪
    2023-10-10
  • Java編程之繼承問題代碼示例

    Java編程之繼承問題代碼示例

    這篇文章主要介紹了Java編程之繼承問題代碼示例,具有一定參考價值,需要的朋友可以了解下。
    2017-11-11
  • springboot+vue實(shí)現(xiàn)websocket配置過程解析

    springboot+vue實(shí)現(xiàn)websocket配置過程解析

    這篇文章主要介紹了springboot+vue實(shí)現(xiàn)websocket配置過程解析,文中通過示例代碼介紹的非常詳細(xì),對大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價值,需要的朋友可以參考下
    2020-04-04
  • 在Spring Boot中加載XML配置的完整步驟

    在Spring Boot中加載XML配置的完整步驟

    這篇文章主要給大家介紹了關(guān)于在Spring Boot中加載XML配置的完整步驟,文中通過示例代碼介紹的非常詳細(xì),對大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價值,需要的朋友們下面隨著小編來一起學(xué)習(xí)學(xué)習(xí)吧
    2020-09-09

最新評論

宝坻区| 慈利县| 西平县| 平塘县| 柘城县| 徐汇区| 白山市| 姚安县| 穆棱市| 会宁县| 阿合奇县| 云南省| 孟州市| 茌平县| 南乐县| 巩留县| 海伦市| 淄博市| 阿尔山市| 南乐县| 万山特区| 斗六市| 镇平县| 隆昌县| 灵川县| 永吉县| 肃宁县| 容城县| 伊金霍洛旗| 洛川县| 吉隆县| 布拖县| 资兴市| 淳安县| 武川县| 霍林郭勒市| 夏津县| 安义县| 高碑店市| 汽车| 屏边|