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
*/