Java延迟消息队列DelayQueue介绍和使用

1、DelayQueue

DelayQueue继承AbstractQueue父类,实现了BlockingQueue接口(BlockingQueue基于ReentrantLock实现),是一个无界的有序阻塞队列,其队列中必须放置实现了Delayed接口的对象。队列中元素的顺序,是由Delayed实现类中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队列中放置的对象,必须实现该接口。该接口主要有两个方法:getDelay(TimeUnit unit)、compareTo(Delayed other)。下面是两个方法的说明:

  • long getDelay(TimeUnit unit);返回到期时间,从DelayQueue中获取元素时,会根据该方法判断队首的元素,是否到了到期时间,未到到期时间,则相当于队列为空,获取不到元素。
  • int compareTo(Delayed other);排序比较方法,DelayQueue中入队元素时,会根据该方法进行排序,所以该方法的实现,要保证到期时间越近的,越靠近前面。

3、Delayed自定义实现类DelayMessage

自己实现的一个DelayMessage,实现了Delayed接口,使用了泛型,来包装具体的消息对象body。实现了getDelay方法和compareTo方法。

3.1、重要属性介绍:

  • uuid:延迟消息对象的唯一ID,主要用于在后面工具类中,取消延迟消息时使用。
  • atomic、n:n这个属性,使用atomic,在构造方法中,保证每一个对象的n都是在递增的,后续接入的n更大。在compareTo方法中,如果延迟时间一样,则用n来比较,越早加入的,越在前面。
  • body:延迟消息的内容
  • executeTime:到期时间

3.2、方法实现:

  • getDelay() 方法比较好实现,返回预计执行的时间 - 当前时间,即为到期剩余时间。
  • compareTo() 方法,通过过期时间和序号n的对比,保证越早执行的,越在前面。

3.3、实现代码

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;

    /**
     * 延迟毫秒数
     */
    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、 延迟消息管理工具类

自己实现的DelayQueueUtil,利用DelayQueue来管理DelayMessage延迟消息。内部使用线程池,来实现消费的线程,不会影响调用代码的执行。

4.1、方法结束

代码中使用Consumer<T>接口,来实现消息到期后的回调方法。Consumer<T>具体用法不再做详细介绍了。

  • submit(String uuid, T msg, Consumer<T> consumer, long delayTime) ,提交一个延迟消息。
    uuid,该条消息的uuid;
    T,为实际的消息;
    Consumer<T>对象,为延迟时间到期后的回调方法,延迟时间到期后,会调用该方法,将提交的T msg,作为参数,传递过去;
    delayTime,为延迟的毫秒数。
  • cancel(String uuid)
    取消该uuid对应的消息,如果取消失败,或者该条消息已到期执行过了,则返回false

4.2、实现代码

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();

    /**
     * 执行读取任务的线程池
     */
    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,则开启监听队列的线程
        if (!STARTING.compareAndSet(false, true)) {
            return;
        }
        THREAD_POOL.execute(() -> {
            while (STARTING.get()) {
                try {
                    DelayMessage<T> delayMessage = (DelayMessage<T>) DELAY_QUEUE.take();
                    // 只有当map里面有该uuid对应的消息,才执行回调方法
                    if (CONSUMER_MAP.containsKey(delayMessage.getUuid())) {
                        // 执行回调方法
                        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的实现方法,自己实现一个即可。

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方法,或执行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
*/

©著作权归作者所有,转载或内容合作请联系作者
  • 序言:七十年代末,一起剥皮案震惊了整个滨河市,随后出现的几起案子,更是在滨河造成了极大的恐慌,老刑警刘岩,带你破解...
    沈念sama阅读 216,287评论 6 498
  • 序言:滨河连续发生了三起死亡事件,死亡现场离奇诡异,居然都是意外死亡,警方通过查阅死者的电脑和手机,发现死者居然都...
    沈念sama阅读 92,346评论 3 392
  • 文/潘晓璐 我一进店门,熙熙楼的掌柜王于贵愁眉苦脸地迎上来,“玉大人,你说我怎么就摊上这事。” “怎么了?”我有些...
    开封第一讲书人阅读 162,277评论 0 353
  • 文/不坏的土叔 我叫张陵,是天一观的道长。 经常有香客问我,道长,这世上最难降的妖魔是什么? 我笑而不...
    开封第一讲书人阅读 58,132评论 1 292
  • 正文 为了忘掉前任,我火速办了婚礼,结果婚礼上,老公的妹妹穿的比我还像新娘。我一直安慰自己,他们只是感情好,可当我...
    茶点故事阅读 67,147评论 6 388
  • 文/花漫 我一把揭开白布。 她就那样静静地躺着,像睡着了一般。 火红的嫁衣衬着肌肤如雪。 梳的纹丝不乱的头发上,一...
    开封第一讲书人阅读 51,106评论 1 295
  • 那天,我揣着相机与录音,去河边找鬼。 笑死,一个胖子当着我的面吹牛,可吹牛的内容都是我干的。 我是一名探鬼主播,决...
    沈念sama阅读 40,019评论 3 417
  • 文/苍兰香墨 我猛地睁开眼,长吁一口气:“原来是场噩梦啊……” “哼!你这毒妇竟也来了?” 一声冷哼从身侧响起,我...
    开封第一讲书人阅读 38,862评论 0 274
  • 序言:老挝万荣一对情侣失踪,失踪者是张志新(化名)和其女友刘颖,没想到半个月后,有当地人在树林里发现了一具尸体,经...
    沈念sama阅读 45,301评论 1 310
  • 正文 独居荒郊野岭守林人离奇死亡,尸身上长有42处带血的脓包…… 初始之章·张勋 以下内容为张勋视角 年9月15日...
    茶点故事阅读 37,521评论 2 332
  • 正文 我和宋清朗相恋三年,在试婚纱的时候发现自己被绿了。 大学时的朋友给我发了我未婚夫和他白月光在一起吃饭的照片。...
    茶点故事阅读 39,682评论 1 348
  • 序言:一个原本活蹦乱跳的男人离奇死亡,死状恐怖,灵堂内的尸体忽然破棺而出,到底是诈尸还是另有隐情,我是刑警宁泽,带...
    沈念sama阅读 35,405评论 5 343
  • 正文 年R本政府宣布,位于F岛的核电站,受9级特大地震影响,放射性物质发生泄漏。R本人自食恶果不足惜,却给世界环境...
    茶点故事阅读 40,996评论 3 325
  • 文/蒙蒙 一、第九天 我趴在偏房一处隐蔽的房顶上张望。 院中可真热闹,春花似锦、人声如沸。这庄子的主人今日做“春日...
    开封第一讲书人阅读 31,651评论 0 22
  • 文/苍兰香墨 我抬头看了看天上的太阳。三九已至,却和暖如春,着一层夹袄步出监牢的瞬间,已是汗流浃背。 一阵脚步声响...
    开封第一讲书人阅读 32,803评论 1 268
  • 我被黑心中介骗来泰国打工, 没想到刚下飞机就差点儿被人妖公主榨干…… 1. 我叫王不留,地道东北人。 一个月前我还...
    沈念sama阅读 47,674评论 2 368
  • 正文 我出身青楼,却偏偏与公主长得像,于是被迫代替她去往敌国和亲。 传闻我的和亲对象是个残疾皇子,可洞房花烛夜当晚...
    茶点故事阅读 44,563评论 2 352

推荐阅读更多精彩内容