
系列文章目錄
準備篇 RabbitMQ安裝檔案
第一章 RabbitMQ快速入門篇
第二章 RabbitMQ的Web管理界面詳解
第三章 RabbitMQ進階篇之死信佇列
第四章 RabbitMQ進階篇之通過插件實作延遲佇列
文章目錄
- 系列文章目錄
- 前言
- 一、什么是延時佇列
- 二、延時佇列使用場景
- 三、RabbitMQ中的TTL
- 四、安裝延時佇列插件(rabbitmq_delayed_message_exchange)
- 五、實作插件版的延時佇列的實體
- 5.1 新增場景
- 5.2 調整需求
- 5.3 根據新需求修改代碼
前言
恭喜所有看到本篇文章的小伙伴,成功解鎖了RabbitMQ系列之高級特性插件版 延遲佇列的內容🎁通過本文,你將清楚的了解到:什么是延時佇列?延時佇列使用場景?如何安裝安裝延時佇列插件(rabbitmq_delayed_message_exchange)?😄本文最后,小名將上一篇文章的實體做一些修改來實作新的效果😁

一、什么是延時佇列
什么是延時佇列?顧名思義:首先它是一種佇列,再給它附加一個延遲消費佇列訊息的功能,也就是說延時佇列中的元素是都是帶時間屬性的,可以指定佇列中的訊息在哪個時間點被消費,
簡單來說,延時佇列就是用來存放需要在指定時間被處理的元素的佇列,
二、延時佇列使用場景
我們常見的延時佇列應用場景:
1、訂單成功后,在30分鐘內沒有支付,自動取消訂單并通知用戶
2、用戶注冊成功后,如果三天內沒有登陸則進行短信提醒,
3、電商平臺新建商戶一個月內還沒上傳商品資訊,將凍結商鋪等
4、預定會議后,需要在預定的時間點前十分鐘通知各個與會人員參加會議,
2.1 分析上述場景
這些場景都是我們常見的,所以我們思考一下:
拿第一個場景來說,系統創建訂單之后,需要取消所有超過30分鐘沒有支付的訂單,拿有“重度選擇恐懼癥”的小名來說吧,也許在加購物車,去支付這些操作都還好好的,突然在付款界面看到價格后停住了,發現這東西我不是那么的需要,就放棄支付了,相信一天之內這樣的小伙伴不在少數,比如:小名03:15放棄支付了,小紅03:16放棄支付了,小剛03:15放棄支付了,小王04:45放棄支付了……我們如何讓系統知道在03:45通知小名和小剛,3:46通知小紅,05:15通知小王呢?
再如后面的幾個場景:發生新用戶注冊事件,三天后檢查新注冊用戶的活動資料,然后通知沒有任何活動記錄的用戶;發生商戶一個月內還沒上傳商品資訊事件,則凍結該商戶的商鋪;發生預定會議事件,判斷離會議開始是否只有十分鐘了,如果是,則通知各個與會人員,
2.2 省時省力的解決方法
這幾種場景,你是不是感覺使用定時任務,輪詢所有資料,每秒查一次,取出需要被處理的資料,然后運行相應的業務代碼處理就可以了?的確如果資料量比較少,這樣做即省時又省力,
比如:“用戶注冊成功后,如果三天內沒有登陸則進行短信提醒”這樣對于時間不是嚴格限制的需求,
我們完全可以每天晚上跑個定時任務檢查一下所有注冊后三天沒登陸的用戶,這是的確一個可行的方案,
2.3 上述做法的缺點
如果資料量比較大,并且時效性較強的場景“30分鐘內沒有支付的訂單,自動取消訂單并通知用戶”,在很短時間內,沒有支付的訂單資料可能多到驚人,如果是活動期間甚至會達到百萬甚至千萬級別,對于這么龐大的資料量依舊使用上述的輪詢方式,顯然是不可取的,很可能在一秒內無法完成所有訂單的檢查,同時會給資料庫帶來很大壓力,無法滿足業務要求而且性能低下,
2.4 使用延時佇列的必要性
如果使用延時訊息佇列,我們在創建訂單的同時將時間推遲30分鐘放入訊息中間件中,等時間一到再取出消費即可,
三、RabbitMQ中的TTL
上文中小名已經解釋過了,這里呢,幫大家簡單回憶一下
需求:
模擬用戶商城購買商品時的兩種情況:1. 成功下單,2. 超時提醒
- 用戶下單
- 用戶下單后展示等待付款頁面
- 在頁面上點擊付款的按鈕,如果不超時,則跳轉到付款成功頁面
- 如果超時,則給用戶發送訊息通知,詢問用戶尚未付款,是否還需要?
配置代碼:
//訂單最多存在10s
args.put("x-message-ttl", 10 * 1000);
//宣告當前佇列系結的死信交換機
args.put("x-dead-letter-exchange", "ex.go.dlx");
//宣告當前佇列系結的死信路由key
args.put("x-dead-letter-routing-key", "go.dlx");
Time To Live(TTL)
RabbitMQ可以針對佇列設定x-message-ttl(對訊息進行單獨設定,每條訊息TTL可以不同),來控制訊息的生存時間,如果超時,則訊息變為dead letter(死信),簡單來說:就是如果設定了佇列的TTL屬性,那么一旦訊息過期,就會被佇列丟棄,
四、安裝延時佇列插件(rabbitmq_delayed_message_exchange)
由于是外網,可能下載速度有些慢,小名在這里幫大家準備好了安裝包,大家可以直接下載使用
地址:
https://wwp.lanzouq.com/ifWwA007nwmf
密碼:
eamon
第一步: 下載好先將檔案解壓到本地電腦上(網盤要求,無法上傳無法識別的*.ez檔案,小名壓縮了一下上傳的)
第二步: 上傳到服務器的RabbitMQ的插件目錄(/rabbitmq/plugins)中
第三步: 進入RabbitMQ的安裝目錄下的sbin目錄,執行下面命令讓該插件生效
rabbitmq-plugins enable rabbitmq_delayed_message_exchange
第四步: 重啟RabbitMQ
關閉服務:
rabbitmqctl stop
啟動服務:
rabbitmq-server -detached
五、實作插件版的延時佇列的實體
5.1 新增場景
假設上篇文章的方式只是在某app內資訊推送,后續添加新需求,比如1分鐘發郵件,1小時短信提醒等等,我們就需要創建很多的佇列用來接收不同的訊息,而且我們并不能保證這些訂單的是按順序提醒的( 即:有可能存在佇列中”A單“提醒時間長于佇列中”B單“的時間戳),這時我們就需要一個更通用的方式來發送此類訊息,這里我們用到了上述的延時佇列插件rabbitmq_delayed_message_exchange
5.2 調整需求
上一篇文章的需求:
模擬用戶商城購買商品時的兩種情況:1. 成功下單,2. 超時提醒
- 用戶下單
- 用戶下單后展示等待付款頁面
- 在頁面上點擊付款的按鈕,如果不超時,則跳轉到付款成功頁面
- 如果超過10s,則給用戶發送系統訊息通知,詢問用戶尚未付款,是否還需要?
上文中小名已經實作了一個超時10s給用戶發送訊息的功能,接下來,我們對上篇文章的代碼做如下
5.3 根據新需求修改代碼
- 新增佇列系結
@Configuration
public class DelayedConfig {
public static final String DELAYED_QUEUE_NAME = "q.delay.plugin";
public static final String DELAYED_EXCHANGE_NAME = "ex.delay.plugin";
public static final String DELAYED_ROUTING_KEY = "delay.plugin";
@Bean
public Queue queueDelay() {
return new Queue(DELAYED_QUEUE_NAME);
}
@Bean
public CustomExchange exchangeDelay() {
Map<String, Object> args = new HashMap<>();
args.put("x-delayed-type", "direct");
return new CustomExchange(DELAYED_EXCHANGE_NAME, "x-delayed-message", true, false, args);
}
@Bean
public Binding bindingDelayPlugin(@Qualifier("queueDelay") Queue queue,
@Qualifier("exchangeDelay") CustomExchange customExchange) {
return BindingBuilder.bind(queue).to(customExchange).with(DELAYED_ROUTING_KEY).noargs();
}
}
- 監聽器做一些修改
小名先說下需要修改的部分,翻遍大家對比,文末貼出完整版,
1)新增一個消費者
//插件延遲佇列,監聽
@RabbitListener(queues = DELAYED_QUEUE_NAME)
public void receiveD(Message message, Channel channel) throws IOException {
String msg = new String(message.getBody());
System.out.println("【※】當前時間:"+new Date().toString()+",延時佇列收到訊息:"+msg);
}
2)新增生產者
@RabbitListener(queues = "q.go.dlx")
public void dlxListener(Message message, Channel channel) throws IOException {
//省略……
//未支付,1min后給用戶發郵箱資訊
long t = System.currentTimeMillis();
String delayOneMin = String.valueOf(this.dateRoll(new Date(), Calendar.MINUTE, 1).getTime() - t);
sendDelayMsgByPlugin(message.getBody()+"【郵箱訊息】", delayOneMin);
//未支付,1小時后給用戶發短信
String delayOneHour = String.valueOf(this.dateRoll(new Date(), Calendar.HOUR, 1).getTime() - t);
sendDelayMsgByPlugin(message.getBody()+"【短信訊息】", delayOneHour);
}
}
public void sendDelayMsgByPlugin(String msg, String delayTime) {
System.out.println("延遲時間"+delayTime);
rabbitTemplate.convertAndSend(DELAYED_EXCHANGE_NAME, DELAYED_ROUTING_KEY, msg, a ->{
a.getMessageProperties().setDelay(Integer.valueOf(delayTime));//60*1000和Integer.valueOf(delayTime)的區別
return a;
});
}
【完整版代碼】
@Component
@Slf4j
public class MqListener {
@Autowired
IPracticeDlxOrderService iPracticeDlxOrderService;
@Autowired
private RabbitTemplate rabbitTemplate;
@RabbitListener(queues = "q.go.dlx")
public void dlxListener(Message message, Channel channel) throws IOException {
System.out.println("支付超時");
Long id = Long.valueOf(new String(message.getBody(), "utf-8"));
PracticeDlxOrder order = iPracticeDlxOrderService.lambdaQuery().eq(PracticeDlxOrder::getId, id).one();
Boolean payStatue = order.getPay();
//判斷是否支付
if (!payStatue) {//未支付,修改未超時
UpdateWrapper<PracticeDlxOrder> dlxOrder = new UpdateWrapper<>();
dlxOrder.eq("id", id);
dlxOrder.set("timeout", 1);
iPracticeDlxOrderService.update(dlxOrder);
log.info("當前時間:{},收到請求,msg:{},delayTime:{}", new Date(), message, new Date().toString());
//未支付,10后給用戶發app資訊
sendDelayMsg(id);
//未支付,1min后給用戶發郵箱資訊
long t = System.currentTimeMillis();
String delayOneMin = String.valueOf(this.dateRoll(new Date(), Calendar.MINUTE, 1).getTime() - t);
sendDelayMsgByPlugin(message.getBody()+"【郵箱訊息】", delayOneMin);
//未支付,1小時后給用戶發短信
String delayOneHour = String.valueOf(this.dateRoll(new Date(), Calendar.HOUR, 1).getTime() - t);
sendDelayMsgByPlugin(message.getBody()+"【短信訊息】", delayOneHour);
}
}
public Date dateRoll(Date date, int i, int d) {
// 獲取Calendar物件并以傳進來的時間為準
Calendar calendar = Calendar.getInstance();
calendar.setTime(date);
// 將現在的時間滾動固定時長,轉換為Date型別賦值
calendar.add(i, d);
// 轉換為Date型別再賦值
date = calendar.getTime();
return date;
}
//死信佇列監聽
@RabbitListener(queues = "q.delay")
public void delayListener(Message message, Channel channel) throws IOException {
System.out.println(new String(message.getBody()));
}
//插件延遲佇列,監聽
@RabbitListener(queues = DELAYED_QUEUE_NAME)
public void receiveD(Message message, Channel channel) throws IOException {
String msg = new String(message.getBody());
System.out.println("【※】當前時間:"+new Date().toString()+",延時佇列收到訊息:"+msg);
}
/**
* 未支付,10s后給用戶發資訊
*/
public void sendDelayMsg(Long id){
rabbitTemplate.setMandatory(true);
//id + 時間戳 全域唯一
Date date = DateUtil.getDate(new Date(),1,10);
CorrelationData correlationData = new CorrelationData(date.toString());
//發送訊息時指定 header 延遲時間
rabbitTemplate.convertAndSend("ex.delay", "q.delay", "您的訂單號:" + id + "尚未付款,是否還需要?",
new MessagePostProcessor() {
@Override
public Message postProcessMessage(Message message) throws AmqpException {
//設定訊息持久化
message.getMessageProperties().setDeliveryMode(MessageDeliveryMode.PERSISTENT);
message.getMessageProperties().setDelay(10*1000);
return message;
}
}, correlationData);
}
/**
*
* @param msg
* @param delayTime
*/
public void sendDelayMsgByPlugin(String msg, String delayTime) {
System.out.println("延遲時間"+delayTime);
rabbitTemplate.convertAndSend(DELAYED_EXCHANGE_NAME, DELAYED_ROUTING_KEY, msg, a ->{
a.getMessageProperties().setDelay(Integer.valueOf(delayTime));//60*1000和Integer.valueOf(delayTime)的區別
return a;
});
}
}
- 運行結果:
支付超時
2022-02-17 17:28:10.650 INFO 18324 --- [ntContainer#2-1] e.mq.dlx.modules.listener.MqListener : 當前時間:Thu Feb 17 17:28:10 CST 2022,收到請求,msg:(Body:'1494242214543482881' MessageProperties [headers={x-first-death-exchange=ex.go, x-death=[{reason=expired, count=1, exchange=ex.go, time=Thu Feb 17 17:28:08 CST 2022, routing-keys=[go], queue=q.go}], x-first-death-reason=expired, x-first-death-queue=q.go}, contentType=text/plain, contentEncoding=UTF-8, contentLength=0, receivedDeliveryMode=PERSISTENT, priority=0, redelivered=false, receivedExchange=ex.go.dlx, receivedRoutingKey=go.dlx, deliveryTag=1, consumerTag=amq.ctag-SasPqfbiS6-pt-e54uV5Hw, consumerQueue=q.go.dlx]),delayTime:Thu Feb 17 17:28:10 CST 2022
延遲時間60000
您的訂單號:1494242214543482881尚未付款,是否還需要?
延遲時間3616616
2022-02-17 17:28:27.268 WARN 18324 --- [nectionFactory1] o.s.amqp.rabbit.core.RabbitTemplate : Returned message but no callback available
【※】當前時間:Thu Feb 17 17:29:27 CST 2022,延時佇列收到訊息:[B@36c9cd1【郵箱訊息】
【※】當前時間:Thu Feb 17 18:28:44 CST 2022,延時佇列收到訊息:[B@36c9cd1【短信訊息】
轉載請註明出處,本文鏈接:https://www.uj5u.com/qita/436411.html
標籤:其他
