Java延遲消息隊列DelayQueue介紹和使用

1、DelayQueue

DelayQueue繼承AbstractQueue父類,實現(xiàn)了BlockingQueue接口(BlockingQueue基于ReentrantLock實現(xiàn)),是一個無界的有序阻塞隊列,其隊列中必須放置實現(xiàn)了Delayed接口的對象。隊列中元素的順序,是由Delayed實現(xiàn)類中compareTo方法決定的,compareTo方法要保證,到期時間越小的,越在前面。

BlockingQueue入隊出隊方法介紹

  • 入隊
方法 隊列已滿 隊列未滿 阻塞后情況
offer(E e) 直接返回false 隊尾插入元素,返回true 不阻塞
offer(E e, long timeout, TimeUnit unit) 進入等待,阻塞 隊尾插入元素,返回true 被喚醒、等待時間超時或者當前線程被中斷
put(E e) 進入等待,阻塞 隊尾插入元素,返回true 隊列有空余了,插入成功,返回true?;蛘呔€程被中斷
  • 出隊
方法 隊列不為空 隊列為空 阻塞后情況
poll() 返回隊首元素 返回null 不阻塞
poll(long timeout, TimeUnit unit) 返回隊首元素 進入等待,阻塞 隊列不為空了,且超時時間未到,返回隊首元素。超時時間到了,返回null
put(E e) 返回隊首元素 進入等待,阻塞 直到隊列不為空了,返回隊首元素。或者線程被中斷

2、Delayed接口

DelayQueue隊列中放置的對象,必須實現(xiàn)該接口。該接口主要有兩個方法:getDelay(TimeUnit unit)、compareTo(Delayed other)。下面是兩個方法的說明:

  • long getDelay(TimeUnit unit);返回到期時間,從DelayQueue中獲取元素時,會根據(jù)該方法判斷隊首的元素,是否到了到期時間,未到到期時間,則相當于隊列為空,獲取不到元素。
  • int compareTo(Delayed other);排序比較方法,DelayQueue中入隊元素時,會根據(jù)該方法進行排序,所以該方法的實現(xiàn),要保證到期時間越近的,越靠近前面。

3、Delayed自定義實現(xiàn)類DelayMessage

自己實現(xiàn)的一個DelayMessage,實現(xiàn)了Delayed接口,使用了泛型,來包裝具體的消息對象body。實現(xiàn)了getDelay方法和compareTo方法。

3.1、重要屬性介紹:

  • uuid:延遲消息對象的唯一ID,主要用于在后面工具類中,取消延遲消息時使用。
  • atomic、n:n這個屬性,使用atomic,在構造方法中,保證每一個對象的n都是在遞增的,后續(xù)接入的n更大。在compareTo方法中,如果延遲時間一樣,則用n來比較,越早加入的,越在前面。
  • body:延遲消息的內容
  • executeTime:到期時間

3.2、方法實現(xiàn):

  • getDelay() 方法比較好實現(xiàn),返回預計執(zhí)行的時間 - 當前時間,即為到期剩余時間。
  • compareTo() 方法,通過過期時間和序號n的對比,保證越早執(zhí)行的,越在前面。

3.3、實現(xiàn)代碼

package com.emdata.videomonitor.common.utils.delayed;

import lombok.Data;
import lombok.Getter;

import java.util.concurrent.Delayed;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicLong;

/**
 * java 延遲隊列消息
 *
 * @version 1.0
 * @date 2021/3/29 19:13
 */
@Gette
public class DelayMessage<T> implements Delayed {

    private static final AtomicLong atomic = new AtomicLong(0);

    private final long n;

    private String uuid;

    /**
     * 消息內容
     */
    private T body;

    /**
     * 到期時間,這個是必須的屬性因為要按照這個判斷延時時長。
     */
    private long executeTime;

    /**
     * 延遲毫秒數(shù)
     */
    private long delayTime;

    public DelayMessage(String uuid, T body, long delayTime) {
        this.uuid = uuid;
        this.n = atomic.getAndIncrement();
        this.body = body;
        this.delayTime = delayTime;
        this.executeTime = TimeUnit.NANOSECONDS.convert(delayTime, TimeUnit.MILLISECONDS) + System.nanoTime();
    }

    @Override
    public long getDelay(TimeUnit unit) {
        return unit.convert(this.executeTime - System.nanoTime(), TimeUnit.NANOSECONDS);
    }

    @Override
    public int compareTo(Delayed other) {
        if (other == this) {
            return 0;
        }
        if (other instanceof DelayMessage) {
            DelayMessage x = (DelayMessage) other;
            long diff = executeTime - x.executeTime;
            if (diff < 0) {
                return -1;
            } else if (diff > 0) {
                return 1;
            } else if (n < x.n) {
                return -1;
            } else {
                return 1;
            }
        }
        long d = (getDelay(TimeUnit.NANOSECONDS) - other.getDelay(TimeUnit.NANOSECONDS));
        return (d == 0) ? 0 : (d < 0 ? -1 : 1);
    }
}

4、 延遲消息管理工具類

自己實現(xiàn)的DelayQueueUtil,利用DelayQueue來管理DelayMessage延遲消息。內部使用線程池,來實現(xiàn)消費的線程,不會影響調用代碼的執(zhí)行。

4.1、方法結束

代碼中使用Consumer<T>接口,來實現(xiàn)消息到期后的回調方法。Consumer<T>具體用法不再做詳細介紹了。

  • submit(String uuid, T msg, Consumer<T> consumer, long delayTime) ,提交一個延遲消息。
    uuid,該條消息的uuid;
    T,為實際的消息;
    Consumer<T>對象,為延遲時間到期后的回調方法,延遲時間到期后,會調用該方法,將提交的T msg,作為參數(shù),傳遞過去;
    delayTime,為延遲的毫秒數(shù)。
  • cancel(String uuid)
    取消該uuid對應的消息,如果取消失敗,或者該條消息已到期執(zhí)行過了,則返回false

4.2、實現(xiàn)代碼

package com.emdata.videomonitor.biz.thread.camera;

import com.google.common.util.concurrent.ThreadFactoryBuilder;
import com.emdata.videomonitor.common.utils.delayed.DelayMessage;
import lombok.extern.slf4j.Slf4j;

import java.util.Map;
import java.util.concurrent.*;
import java.util.concurrent.atomic.AtomicBoolean;
import java.util.function.Consumer;

/**
 * java延遲隊列
 *
 * @version 1.0
 * @date 2021/3/29 20:11
 */
@Slf4j
public class DelayQueueUtil {

    private static final Map<String, Consumer<?>> CONSUMER_MAP = new ConcurrentHashMap<>();

    private static final AtomicBoolean STARTING = new AtomicBoolean();

    /**
     * 延遲隊列
     */
    private static final DelayQueue<DelayMessage<?>> DELAY_QUEUE = new DelayQueue<>();


    private static final int CORE_POOL_SIZE = 4;
    private static final int MAXIMUM_POOL_SIZE = 10;
    private static final long KEEP_ALIVE_TIME = 20;
    private static final TimeUnit UNIT = TimeUnit.SECONDS;
    private static final int MAXIMUM_ARRAY_SIZE = 10;
    private static final ThreadFactory NAMED_FACTORY = new ThreadFactoryBuilder().setNameFormat("java_delay_thread_%d").build();

    /**
     * 執(zhí)行讀取任務的線程池
     */
    private static final ExecutorService THREAD_POOL = new ThreadPoolExecutor(
            CORE_POOL_SIZE,
            MAXIMUM_POOL_SIZE,
            KEEP_ALIVE_TIME,
            UNIT,
            new ArrayBlockingQueue<>(MAXIMUM_ARRAY_SIZE),
            NAMED_FACTORY);


    /**
     * 提交一個延遲消息
     * @param uuid 消息的uuid
     * @param msg 消息對象
     * @param consumer 延遲到期后回到方法
     * @param delayTime 延遲時間,毫秒
     * @param <T> 消息對象類型
     * @return true: 提交成功
     */
    public static <T> boolean submit(String uuid, T msg, Consumer<T> consumer, long delayTime) {
        DelayMessage<T> delayMessage = new DelayMessage<>(uuid, msg, delayTime);
        addTask(uuid, consumer);

        return DELAY_QUEUE.offer(delayMessage);
    }

    /**
     * 取消一個延遲消息
     * @param uuid 消息的uuid
     * @return true: 取消成功
     */
    public static boolean cancel(String uuid) {
        return CONSUMER_MAP.remove(uuid) != null;
    }

    /**
     * 添加任務,懶加載開啟消費線程
     * @param uuid 消息的uuid
     * @param consumer 回調方法
     * @param <T> 消息對象類型
     */
    private static <T> void addTask(String uuid, Consumer<T> consumer) {
        CONSUMER_MAP.put(uuid, consumer);

        // STARTING 是false,則開啟監(jiān)聽隊列的線程
        if (!STARTING.compareAndSet(false, true)) {
            return;
        }
        THREAD_POOL.execute(() -> {
            while (STARTING.get()) {
                try {
                    DelayMessage<T> delayMessage = (DelayMessage<T>) DELAY_QUEUE.take();
                    // 只有當map里面有該uuid對應的消息,才執(zhí)行回調方法
                    if (CONSUMER_MAP.containsKey(delayMessage.getUuid())) {
                        // 執(zhí)行回調方法
                        execCall(consumer, delayMessage);
                    }
                } catch (InterruptedException e) {
                    STARTING.set(false);
                }
            }
        });
    }

    private static <T> void execCall(Consumer<T> consumer, DelayMessage<T> delayMessage) {
        CONSUMER_MAP.remove(delayMessage.getUuid());
        THREAD_POOL.execute(() -> consumer.accept(delayMessage.getBody()));
    }
}

5、測試一下延遲消息工具類

演示代碼,其中Guid工具類,封裝了一下java中uuid的實現(xiàn)方法,自己實現(xiàn)一個即可。

import com.emdata.videomonitor.common.utils.Guid;
public class Test {

    public static void main(String[] args) {
        String uuid = Guid.newGUID();
        long delayTime = 2000;
    
        log.info(System.currentTimeMillis() + "");
        
        // lambda表達式,將getMsg(),當做回調方法,來接收延遲消息
        DelayQueueUtil.submit(uuid, "消息1", Test::getMsg, delayTime);
    
        String uuid2 = Guid.newGUID();
        delayTime = 5000;
        DelayQueueUtil.submit(uuid2, "消息2", Test::getMsg, delayTime );
    
        log.info("main");
        // 測試cancel方法,或執(zhí)行cancel 方法,則 消息1 ,不會被打印
    //        boolean cancel = DelayQueueUtil.cancel(uuid);
    //        log.info(cancel + "");
    }
    
    private static void getMsg(String msg) {
        log.info("{}", msg)
    }
}


/*
輸出:
09:31:43.478 [main] INFO com.emdata.videomonitor.biz.thread.camera.DelayQueueUtil - 1617154303476
09:31:43.513 [main] INFO com.emdata.videomonitor.biz.thread.camera.DelayQueueUtil - main
09:31:45.514 [java_delay_thread_2] INFO com.emdata.videomonitor.biz.thread.camera.DelayQueueUtil - 消息1
09:31:48.514 [java_delay_thread_3] INFO com.emdata.videomonitor.biz.thread.camera.DelayQueueUtil - 消息2
*/

?著作權歸作者所有,轉載或內容合作請聯(lián)系作者
【社區(qū)內容提示】社區(qū)部分內容疑似由AI輔助生成,瀏覽時請結合常識與多方信息審慎甄別。
平臺聲明:文章內容(如有圖片或視頻亦包括在內)由作者上傳并發(fā)布,文章內容僅代表作者本人觀點,簡書系信息發(fā)布平臺,僅提供信息存儲服務。

相關閱讀更多精彩內容

友情鏈接更多精彩內容