日韩成人免费在线_国产成人一二_精品国产免费人成电影在线观..._日本一区二区三区久久久久久久久不

當前位置:首頁 > 科技  > 軟件

分布式延時消息的另外一種選擇 Redisson

來源: 責編: 時間:2024-05-16 09:09:51 160觀看
導讀前言因為工作中需要用到分布式的延時隊列,調研了一段時間,選擇使用 Redisson DelayedQueue,為了搞清楚內部運行流程,特記錄下來。總體流程大概是圖中的這個樣子,初看一眼有點不知從何下手,接下來我會通過以下幾點來分析流程

前言

Gcs28資訊網——每日最新資訊28at.com

因為工作中需要用到分布式的延時隊列,調研了一段時間,選擇使用 Redisson DelayedQueue,為了搞清楚內部運行流程,特記錄下來。Gcs28資訊網——每日最新資訊28at.com

總體流程大概是圖中的這個樣子,初看一眼有點不知從何下手,接下來我會通過以下幾點來分析流程,相信看完本文你能了解整個運行流程。Gcs28資訊網——每日最新資訊28at.com

  • 基本使用
  • 內部數據結構介紹
  • 基本流程
  • 發送延時消息
  • 獲取延時消息
  • 初始化延時隊列

圖片圖片Gcs28資訊網——每日最新資訊28at.com

基本使用

發送延遲消息代碼如下,發送了一條延遲時間為 5s 的消息。Gcs28資訊網——每日最新資訊28at.com

public void produce() {  String queuename = "delay-queue";  RBlockingQueue<String> blockingQueue = redissonClient.getBlockingQueue(queuename);  RDelayedQueue<String> delayedQueue = redissonClient.getDelayedQueue(blockingQueue);  delayedQueue.offer("測試延遲消息", 5, TimeUnit.SECONDS);}

接收消息代碼如下,可以看到 delayedQueue 是沒有用到的,那么為什么要加這一行呢,這個后面總結部分回答。Gcs28資訊網——每日最新資訊28at.com

public void consume() throws InterruptedException { String queuename = "delay-queue";  RBlockingQueue<String> blockingQueue = redissonClient.getBlockingQueue(queuename);  RDelayedQueue<String> delayedQueue = redissonClient.getDelayedQueue(blockingQueue);  String msg = blockingQueue.take();  //收到消息進行處理...}

這兩段代碼可以寫在兩個不同的 Java 工程里,只要連接的是同一個 Redis 就行。Gcs28資訊網——每日最新資訊28at.com

調用 comsume() 之后,如果隊列里沒有消息,會阻塞等待隊列里有消息并且取到了才會返回。之所以這么說是因為可能有別的 Java 進程也在跟你一樣取同一個隊列里的消息,如果消息被另一個搶完了,那這時就還得阻塞等待。Gcs28資訊網——每日最新資訊28at.com

這時看上去的原理是這樣的:Gcs28資訊網——每日最新資訊28at.com

生產者調用 offer() 后,自己內部開啟一個定時器,等到了時間再發送到 redis 的 list 里。Gcs28資訊網——每日最新資訊28at.com

圖片圖片Gcs28資訊網——每日最新資訊28at.com

如果是這樣設計的話,相信大家都能看出來一個很簡單的問題,要是延時時間還沒到,生產者自己掛了,那樣消息就丟了。所以還是讓我們接著往下看。Gcs28資訊網——每日最新資訊28at.com

內部數據結構介紹

redisson 源碼里一共創建了三個隊列:【消息延時隊列】、【消息順序隊列】、【消息目標隊列】。Gcs28資訊網——每日最新資訊28at.com

圖片圖片Gcs28資訊網——每日最新資訊28at.com

假設在同一時間按照 msg1、msg2、msg3 的順序發消息到延時隊列,這三條消息就會被保存在【消息延時隊列】和【消息順序隊列】。Gcs28資訊網——每日最新資訊28at.com

可以看到【消息延時隊列】的順序是按照到期時間升序排列的,而不是像【消息順序隊列】按照插入順序排。Gcs28資訊網——每日最新資訊28at.com

消息到期后會將消息從前兩個隊列移除(怎么移?誰來移?),插入【消息目標隊列】,也就是圖中第三個隊列。Gcs28資訊網——每日最新資訊28at.com

消費者也是阻塞在【消息目標隊列】上取消息。Gcs28資訊網——每日最新資訊28at.com

這時可以簡單說明下每個隊列的作用:Gcs28資訊網——每日最新資訊28at.com

  • 【消息延時隊列】利用按照到期時間排序的特性,可以很快找到下一個要到期的消息,客戶端內部自己定時到【消息目標隊列】取
  • 【消息順序隊列】這個隊列對分析的流程關聯不大,可以忽略
  • 【消息目標隊列】存放到期的消息,供消費端取

其實【消息延時隊列】隊列里存的時間(也就是 zet 的 score)是到期的時間戳,為了畫圖方便,圖里就畫的是延遲的時間,不過不影響理解。Gcs28資訊網——每日最新資訊28at.com

理解好這幾個隊列的名字和作用,后面還會一直用到,如果忘了可以翻回來回顧下。Gcs28資訊網——每日最新資訊28at.com

因為書寫理解方便和【消息順序隊列】在本文沒涉及到,后面部分好幾次提到的內容:把到期的消息從【消息延時隊列】移到【消息目標隊列】里,這句話實際的代碼邏輯是這樣:把【消息延時隊列】和【消息順序隊列】里的到期消息移除,把它們插入到【消息目標隊列】。Gcs28資訊網——每日最新資訊28at.com

基本流程

知道了內部所使用到的數據結構后,這里可以簡單說下整體的基本流程。Gcs28資訊網——每日最新資訊28at.com

先說發送延遲消息,發送的延遲消息會先存在【消息延時隊列】和【消息順序隊列】,如果【消息延時隊列】原本是空的,會發布訂閱信息提醒有新的消息。Gcs28資訊網——每日最新資訊28at.com

獲取延遲消息只需要從【消息目標隊列】阻塞的取就行了,因為里面都是到期數據。Gcs28資訊網——每日最新資訊28at.com

那么問題就只剩下怎么樣判斷時間到了,把【消息延時隊列】里的消息移動到【消息目標隊列】里呢?Gcs28資訊網——每日最新資訊28at.com

這部分工作交給了初始化延時隊列來處理。Gcs28資訊網——每日最新資訊28at.com

這里面會定時從【消息延時隊列】查詢最新到期時間,定時去把【消息延時隊列】里的消息移動到【消息目標隊列】里。Gcs28資訊網——每日最新資訊28at.com

如果【消息延時隊列】是空的,就不會再定時查,而是等待發布訂閱信息提醒,再定時把【消息延時隊列】里的消息移動到【消息目標隊列】里。Gcs28資訊網——每日最新資訊28at.com

剛開始看可能有點抽象,可以看完底下一節內容之后,再回頭來看這里對應的流程總結,可能會比較清晰。Gcs28資訊網——每日最新資訊28at.com

發送延時消息

發送延時消息的邏輯比較簡單,先看下發送的代碼。Gcs28資訊網——每日最新資訊28at.com

public void produce() {  String queuename = "delay-queue";  RBlockingQueue<String> blockingQueue = redissonClient.getBlockingQueue(queuename);  RDelayedQueue<String> delayedQueue = redissonClient.getDelayedQueue(blockingQueue);  delayedQueue.offer("測試延遲消息", 5, TimeUnit.SECONDS);}

從 delayedQueue.offer 方法開始,最終會執行到 RedissonDelayedQueue 的 offerAsync 方法里。Gcs28資訊網——每日最新資訊28at.com

offerAsync 方法的作用就是發送一段腳本給 redis 執行,腳本內容是:Gcs28資訊網——每日最新資訊28at.com

  1. 將消息和到期時間插入【消息延時隊列】和【消息順序隊列】
  2. 如果最近到期的消息是剛剛插入的消息,則對指定主題發布到期時間,目的是為了讓客戶端定時去把【消息延時隊列】里的到期數據移動到【消息目標隊列】
@Overridepublic RFuture<Void> offerAsync(V e, long delay, TimeUnit timeUnit) {  if (delay < 0) {   throw new IllegalArgumentException("Delay can't be negative");  }  long delayInMs = timeUnit.toMillis(delay);  long timeout = System.currentTimeMillis() + delayInMs;  long randomId = ThreadLocalRandom.current().nextLong();  return commandExecutor.evalWriteNoRetryAsync(getRawName(), codec, RedisCommands.EVAL_VOID,  "local value = struct.pack('dLc0', tonumber(ARGV[2]), string.len(ARGV[3]), ARGV[3]);"   + "redis.call('zadd', KEYS[2], ARGV[1], value);"  + "redis.call('rpush', KEYS[3], value);"  // if new object added to queue head when publish its startTime   // to all scheduler workers   + "local v = redis.call('zrange', KEYS[2], 0, 0); "  + "if v[1] == value then "  + "redis.call('publish', KEYS[4], ARGV[1]); "  + "end;",  Arrays.<Object>asList(getRawName(), timeoutSetName, queueName, channelName),  timeout, randomId, encode(e));}

獲取延時消息

獲取延時消息是本文最簡單的一部分。Gcs28資訊網——每日最新資訊28at.com

public void consume() throws InterruptedException {  String queuename = "delay-queue";  RBlockingQueue<String> blockingQueue = redissonClient.getBlockingQueue(queuename);  RDelayedQueue<String> delayedQueue = redissonClient.getDelayedQueue(blockingQueue);  String msg = blockingQueue.take();  //收到消息進行處理...}

blockingQueue.take() 方法其實只是對【消息目標隊列】執行 blpop 阻塞的獲取到期消息Gcs28資訊網——每日最新資訊28at.com

初始化延時隊列

看一下初始化的代碼。Gcs28資訊網——每日最新資訊28at.com

public void init() {    String queuename = "delay-queue";    RBlockingQueue<String> blockingQueue = redissonClient.getBlockingQueue(queuename);    RDelayedQueue<String> delayedQueue = redissonClient.getDelayedQueue(blockingQueue);}

入口就是在 redissonClient.getDelayedQueue(blockingQueue) 中,創建了 RedissonDelayedQueue 對象,并執行了構造方法里的邏輯。Gcs28資訊網——每日最新資訊28at.com

那么這里面主要做了什么事呢?Gcs28資訊網——每日最新資訊28at.com

主要是調用了 QueueTransferTask 的 start() 方法。Gcs28資訊網——每日最新資訊28at.com

public void start() {  RTopic schedulerTopic = getTopic();  statusListenerId = schedulerTopic.addListener(new BaseStatusListener() {      @Override    public void onSubscribe(String channel) {      pushTask();    } }); messageListenerId = schedulerTopic.addListener(Long.class, new MessageListener<Long>() {      @Override      public void onMessage(CharSequence channel, Long startTime) {     scheduleTask(startTime);   } });}

這段代碼主要是設置了指定主題(主題名:redisson_delay_queue_channel:{queuename})兩個發布訂閱的監聽器。Gcs28資訊網——每日最新資訊28at.com

  1. 當指定主題有新訂閱時調用 pushTask() 方法,里面又會調用 pushTaskAsync() 方法
  2. 當指定主題有新消息時調用 scheduleTask(startTime) 方法

需要注意的是,這里會先訂閱指定主題,然后觸發執行 onSubscribe() 方法。Gcs28資訊網——每日最新資訊28at.com

所以我們主要搞懂這三個方法都是做什么的,那么整個初始化流程就明白了。Gcs28資訊網——每日最新資訊28at.com

因為這三個方法是相互調用的,只看文字的話容易云里霧里,這里有個流程圖,看方法解釋文字的時候可以對照著流程圖看比較有印象。Gcs28資訊網——每日最新資訊28at.com

Gcs28資訊網——每日最新資訊28at.com

圖片圖片Gcs28資訊網——每日最新資訊28at.com

  • scheduleTask()這個方法看起來多,但核心內容就是根據方法參數指定的時間調用 pushTask()。
private void scheduleTask(final Long startTime) {  TimeoutTask oldTimeout = lastTimeout.get();  if (startTime == null) {    return;  }  if (oldTimeout != null) {    oldTimeout.getTask().cancel();  }  long delay = startTime - System.currentTimeMillis();  if (delay > 10) {    Timeout timeout = connectionManager.newTimeout(new TimerTask() {                          @Override      public void run(Timeout timeout) throws Exception {        pushTask();        TimeoutTask currentTimeout = lastTimeout.get();        if (currentTimeout.getTask() == timeout) {          lastTimeout.compareAndSet(currentTimeout, null);        }      }    }, delay, TimeUnit.MILLISECONDS);    if (!lastTimeout.compareAndSet(oldTimeout, new TimeoutTask(startTime, timeout))) {      timeout.cancel();    }  } else {    pushTask();  }}
  • pushTaskAsync()這個方法是抽象方法,在創建 RedissonDelayedQueue 對象的時候傳進來的,代碼如下:
@Overrideprotected RFuture<Long> pushTaskAsync() {  return commandExecutor.evalWriteAsync(getRawName(), LongCodec.INSTANCE, RedisCommands.EVAL_LONG,  "local expiredValues = redis.call('zrangebyscore', KEYS[2], 0, ARGV[1], 'limit', 0, ARGV[2]); "  + "if #expiredValues > 0 then "  + "for i, v in ipairs(expiredValues) do "  + "local randomId, value = struct.unpack('dLc0', v);"  + "redis.call('rpush', KEYS[1], value);"  + "redis.call('lrem', KEYS[3], 1, v);"  + "end; "  + "redis.call('zrem', KEYS[2], unpack(expiredValues));"  + "end; "  // get startTime from scheduler queue head task  + "local v = redis.call('zrange', KEYS[2], 0, 0, 'WITHSCORES'); "  + "if v[1] ~= nil then "  + "return v[2]; "  + "end "  + "return nil;",  Arrays.<Object>asList(getRawName(), timeoutSetName, queueName),  System.currentTimeMillis(), 100);}

看不懂也不要緊,聽我解釋下就明白了。Gcs28資訊網——每日最新資訊28at.com

這里發送了一段腳本給 redis 執行:Gcs28資訊網——每日最新資訊28at.com

我的理解就是初始化的時候Gcs28資訊網——每日最新資訊28at.com

1是為了處理舊的消息,比如生產者1發送了消息,然后時間沒到自己下線了,這時如果沒有其他客戶端在線,就沒有人能把數據從【消息目標隊列】移到【消息目標隊列】了。Gcs28資訊網——每日最新資訊28at.com

2是返回的這個時間戳,會拿這個定時,等時間到了去【消息目標隊列】拉取到期的消息。Gcs28資訊網——每日最新資訊28at.com

簡單總結就是這個方法是把到期消息從【消息延時隊列】放到【消息目標隊列】里,并且返回了最近要到期消息的時間戳。Gcs28資訊網——每日最新資訊28at.com

  1. 從【消息延時隊列】取出前一百條到期的消息,如果有的話,添加到【消息目標隊列】里,并將這些消息從【消息延時隊列】和【消息順序隊列】中移除
  2. 從【消息延時隊列】取出下一條要到期的消息,返回它的到期時間戳(如果隊列里沒消息返回空)。
  • pushTask()
private void pushTask() {  RFuture<Long> startTimeFuture = pushTaskAsync();  startTimeFuture.whenComplete((res, e) -> {    if (e != null) {      if (e instanceof RedissonShutdownException) {        return;      }      log.error(e.getMessage(), e);      scheduleTask(System.currentTimeMillis() + 5 * 1000L);      return;    }    if (res != null) {      scheduleTask(res);    }  });}

這個代碼看起來就比較簡單,調用了 pushTaskAsync() 獲取最近要到期消息的時間戳(異步封裝了一下)。Gcs28資訊網——每日最新資訊28at.com

有異常的話就調用 scheduleTask() 五秒后再執行一次 pushTask()。Gcs28資訊網——每日最新資訊28at.com

沒有異常的話如果有最近要到期消息的時間戳(說明【消息延時隊列】里還有未到期消息),用這個最新到期時間調用 scheduleTask(),在這個指定的時間調用 pushTask()。Gcs28資訊網——每日最新資訊28at.com

這個方法簡單總結就是決定了要不要調用、什么時候再調用 pushTask(),主要操作邏輯都在 pushTaskAsync() 里(把到期的消息從【消息延時隊列】移到【消息目標隊列】供消費端消費)。Gcs28資訊網——每日最新資訊28at.com

了解了上面幾個方法的流程和含義,還記得一開頭提到的添加了兩個發布訂閱的監聽器嗎?Gcs28資訊網——每日最新資訊28at.com

1.當指定主題有新訂閱時調用 pushTask() 方法,里面又會調用 pushTaskAsync() 方法Gcs28資訊網——每日最新資訊28at.com

2.當指定主題有新消息時調用 scheduleTask(startTime) 方法Gcs28資訊網——每日最新資訊28at.com

需要注意的是,這里會先訂閱指定主題,然后觸發執行 onSubscribe() 方法Gcs28資訊網——每日最新資訊28at.com

  1. 在初始化延時隊列剛啟動的時候,處理到期舊數據:把到期的消息從【消息延時隊列】移到【消息目標隊列】供消費端消費;處理新數據:獲取下次到期時間決定下次調用 pushTask() 的時間。上面講的這種情況是站在當前客戶端的視角,但畢竟這是監聽訂閱信息,如果啟動不止一個客戶端的話(就算是1個生產者1個消費者,也算兩個客戶端),總有一個客戶端的訂閱信息回調函數,會不會有問題?仔細想想是沒有的,處理到期舊數據:之前啟動的客戶端已經處理完了;處理新數據:獲取最近到期時間,在 scheduleTask() 里,如果之前有正在定時的任務,會把原來正在定時的任務取消掉。這個被取消的任務,時間要么就是當前這個時間,要么是之后的時間,取消掉不會影響邏輯。
  2. 為了應對原本【消息延時隊列】里沒消息了這種情況,流程結束了,重啟定時去調用 pushTask() ,把到期的消息從【消息延時隊列】移到【消息目標隊列】供消費端消費。

總結

再放一下開頭的圖總體流程圖:Gcs28資訊網——每日最新資訊28at.com

圖片圖片Gcs28資訊網——每日最新資訊28at.com

  1. 初始化延時隊列時會把【消息延時隊列】里的到期數據移動到【消息目標隊列】,沒有也有可能;然后是找最近要到期的消息時間,定時去拉,這個剛啟動也是可能沒有的,不過不要緊,這兩步是為了處理滯留在【消息延時隊列】的舊數據(在發送了延時消息后,還沒到期時所有客戶端都下線了,這樣就沒人能把【消息延時隊列】里的到期數據移動到【消息目標隊列】里,就會出現這種情況);最主要的還是設置了發布訂閱監聽器,當有人發送延時消息的時候能收到通知,定時去將【消息延時隊列】里的到期數據移動到【消息目標隊列】。
  2. 發送延時消息會先發送到【消息延時隊列】和【消息順序隊列】,如果【消息延時隊列】里沒有數據,則將剛發送的到期時間發布到指定主題,提醒其他客戶端有新消息。
  3. 初始化延時隊列時設置的發布訂閱監聽器把【消息延時隊列】里的到期數據移動到【消息目標隊列】里。
  4. 獲取延遲消息只需要執行 blpop 阻塞的獲取【消息目標隊列】的消息就可以了。

這里回答開頭部分說的問題,到這看完了本文,你可以試著自己想一想這個問題的答案。Gcs28資訊網——每日最新資訊28at.com

接收消息代碼如下,可以看到 delayedQueue 是沒有用到的,那么為什么要加這一行呢,這個后面總結部分回答。Gcs28資訊網——每日最新資訊28at.com

public void consume() throws InterruptedException {    String queuename = "delay-queue";    RBlockingQueue<String> blockingQueue = redissonClient.getBlockingQueue(queuename);    RDelayedQueue<String> delayedQueue = redissonClient.getDelayedQueue(blockingQueue);    String msg = blockingQueue.take();    //收到消息進行處理...}

其實這個問題也是我開發過程中遇到的一個奇怪的地方,接收方代碼沒有初始化延時隊列。Gcs28資訊網——每日最新資訊28at.com

首先再啰嗦一句,初始化延時隊列的作用是會定時去把【消息延時隊列】里的到期數據移動到【消息目標隊列】。Gcs28資訊網——每日最新資訊28at.com

如果只有發送方初始化延時隊列:Gcs28資訊網——每日最新資訊28at.com

  1. 發送方發送了延遲消息,在到期之前下線了(它就不能把【消息延時隊列】里的到期數據移動到【消息目標隊列】),而且沒有其他發送方。
  2. 接收方不管有多少個,都沒人能把【消息延時隊列】里的到期數據移動到【消息目標隊列】。

所以接收方代碼里也初始化延時隊列能夠避免一部分數據丟失問題。Gcs28資訊網——每日最新資訊28at.com

本文鏈接:http://www.www897cc.com/showinfo-26-88383-0.html分布式延時消息的另外一種選擇 Redisson

聲明:本網頁內容旨在傳播知識,若有侵權等問題請及時與本網聯系,我們將在第一時間刪除處理。郵件:2376512515@qq.com

上一篇: 聊聊Vue.js 基礎語法詳解

下一篇: 竟然還能這樣高效地操作 JSON 對象!

標簽:
  • 熱門焦點
  • 6月安卓手機好評榜:魅族20 Pro蟬聯冠軍

    性能榜和性價比榜之后,我們來看最后的安卓手機好評榜,數據來源安兔兔評測,收集時間2023年6月1日至6月30日,僅限國內市場。第一名:魅族20 Pro好評率:95%5月份的時候魅族20 Pro就是
  • 太卷!Redmi MAX 100英寸電視便宜了:12999元買Redmi史上最大屏

    8月5日消息,從小米商城了解到,Redmi MAX 100英寸巨屏電視日前迎來官方優惠,到手價12999元,比發布價便宜了7000元,在大屏電視市場開卷。據了解,Redmi MAX 100
  • 三言兩語說透設計模式的藝術-簡單工廠模式

    一、寫在前面工廠模式是最常見的一種創建型設計模式,通常說的工廠模式指的是工廠方法模式,是使用頻率最高的工廠模式。簡單工廠模式又稱為靜態工廠方法模式,不屬于GoF 23種設計
  • Java NIO內存映射文件:提高文件讀寫效率的優秀實踐!

    Java的NIO庫提供了內存映射文件的支持,它可以將文件映射到內存中,從而可以更快地讀取和寫入文件數據。本文將對Java內存映射文件進行詳細的介紹和演示。內存映射文件概述內存
  • 虛擬鍵盤 API 的妙用

    你是否在遇到過這樣的問題:移動設備上有一個固定元素,當激活虛擬鍵盤時,該元素被隱藏在了鍵盤下方?多年來,這一直是 Web 上的默認行為,在本文中,我們將探討這個問題、為什么會發生
  • 慕巖炮轟抖音,百合網今何在?

    來源:價值研究所 作者:Hernanderz&ldquo;難道就因為自己的一個產品牛逼了,從客服到總裁,都不愿意正視自己產品和運營上的問題,選擇逃避了嗎?&rdquo;這一番話,出自百合網聯合創
  • 質感不錯!OPPO K11渲染圖曝光:旗艦IMX890傳感器首次下放

    一直以來,OPPO K系列機型都保持著較為均衡的產品體驗,歷來都是2K價位的明星機型,去年推出的OPPO K10和OPPO K10 Pro兩款機型憑借各自的出色配置,堪稱有
  • 聯想的ThinkBook Plus下一版曝光,鍵盤旁邊塞個平板

    ThinkBook Plus 是聯想的一個特殊筆記本類別,它在封面放入了一塊墨水屏,也給人留下了較為深刻的印象。據有人爆料,聯想的下一款 ThinkBook Plus 可能更特殊,它
  • SN570 NVMe SSD固態硬盤 價格與性能兼具

    SN570 NVMe SSD固態硬盤是西部數據發布的最新一代WD Blue系列的固態硬盤,不僅閃存技術更為精進,性能也得到了進一步的躍升。WD Blue SN570 NVMe SSD的包裝外
Top 主站蜘蛛池模板: 镇江市| 建水县| 兴国县| 全椒县| 汝州市| 姚安县| 黄冈市| 县级市| 怀安县| 保靖县| 读书| 潮州市| 晋州市| 江永县| 阿克苏市| 武汉市| 赤水市| 平陆县| 民和| 望江县| 蓝田县| 株洲市| 六盘水市| 高雄县| 卢龙县| 嘉荫县| 宝清县| 南皮县| 林西县| 绵阳市| 凤山县| 滨海县| 沙田区| 葵青区| 庆云县| 龙江县| 禹州市| 铁岭县| 广汉市| 新竹县| 株洲市|