使用delayedQueue實(shí)現(xiàn)你本地的延遲隊(duì)列

了解DelayQueue

DelayQueue是什么?

DelayQueue是一個(gè)無(wú)界的BlockingQueue,用于放置實(shí)現(xiàn)了Delayed接口的對(duì)象,其中的對(duì)象只能在其到期時(shí)才能從隊(duì)列中取走。這種隊(duì)列是有序的,即隊(duì)頭對(duì)象的延遲到期時(shí)間最長(zhǎng)。注意:不能將null元素放置到這種隊(duì)列中。

DelayQueue能做什么?

在我們的業(yè)務(wù)中通常會(huì)有一些需求是這樣的:

  1. 淘寶訂單業(yè)務(wù):下單之后如果三十分鐘之內(nèi)沒(méi)有付款就自動(dòng)取消訂單。
  2. 餓了嗎訂餐通知:下單成功后60s之后給用戶發(fā)送短信通知。

那么這類業(yè)務(wù)我們可以總結(jié)出一個(gè)特點(diǎn):需要延遲工作。
由此的情況,就是我們的DelayQueue應(yīng)用需求的產(chǎn)生。

怎么用DelayQueue來(lái)解決這類的問(wèn)題

先聲明一個(gè)Delayed的對(duì)象


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

/**
 * <p>
 * [任務(wù)調(diào)度系統(tǒng)]
 * <br>
 * [隊(duì)列中要執(zhí)行的任務(wù)]
 * </p>
 *
 * @author wangguangdong
 * @version 1.0
 * @Date 2015年11月22日19:46:39
 */
public class Task<T extends Runnable> implements Delayed {
    /**
     * 到期時(shí)間
     */
    private final long time;

    /**
     * 問(wèn)題對(duì)象
     */
    private final T task;
    private static final AtomicLong atomic = new AtomicLong(0);

    private final long n;

    public Task(long timeout, T t) {
        this.time = System.nanoTime() + timeout;
        this.task = t;
        this.n = atomic.getAndIncrement();
    }

    /**
     * 返回與此對(duì)象相關(guān)的剩余延遲時(shí)間,以給定的時(shí)間單位表示
     */
    @Override
    public long getDelay(TimeUnit unit) {
        return unit.convert(this.time - System.nanoTime(), TimeUnit.NANOSECONDS);
    }

    @Override
    public int compareTo(Delayed other) {
        // TODO Auto-generated method stub
        if (other == this) // compare zero ONLY if same object
            return 0;
        if (other instanceof Task) {
            Task x = (Task) other;
            long diff = time - x.time;
            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);
    }

    public T getTask() {
        return this.task;
    }

    @Override
    public int hashCode() {
        return task.hashCode();
    }

    @Override
    public boolean equals(Object object) {
        if (object instanceof Task) {
            return object.hashCode() == hashCode() ? true : false;
        }
        return false;
    }


}

再實(shí)現(xiàn)一個(gè)管理延遲任務(wù)的類

import org.apache.log4j.Logger;

import java.util.concurrent.DelayQueue;
import java.util.concurrent.Executor;
import java.util.concurrent.Executors;
import java.util.concurrent.TimeUnit;

/**
 * <p>
 * [任務(wù)調(diào)度系統(tǒng)]
 * <br>
 * [后臺(tái)守護(hù)線程不斷的執(zhí)行檢測(cè)工作]
 * </p>
 *
 * @author wangguangdong
 * @version 1.0
 * @Date 2015年11月23日14:19:40
 */
public class TaskQueueDaemonThread {

    private static final Logger LOG = Logger.getLogger(TaskQueueDaemonThread.class);

    private TaskQueueDaemonThread() {
    }

    private static class LazyHolder {
        private static TaskQueueDaemonThread taskQueueDaemonThread = new TaskQueueDaemonThread();
    }

    public static TaskQueueDaemonThread getInstance() {
        return LazyHolder.taskQueueDaemonThread;
    }

    Executor executor = Executors.newFixedThreadPool(20);
    /**
     * 守護(hù)線程
     */
    private Thread daemonThread;

    /**
     * 初始化守護(hù)線程
     */
    public void init() {
        daemonThread = new Thread(() -> execute());
        daemonThread.setDaemon(true);
        daemonThread.setName("Task Queue Daemon Thread");
        daemonThread.start();
    }

    private void execute() {
        System.out.println("start:" + System.currentTimeMillis());
        while (true) {
            try {
                //從延遲隊(duì)列中取值,如果沒(méi)有對(duì)象過(guò)期則隊(duì)列一直等待,
                Task t1 = t.take();
                if (t1 != null) {
                    //修改問(wèn)題的狀態(tài)
                    Runnable task = t1.getTask();
                    if (task == null) {
                        continue;
                    }
                    executor.execute(task);
                    LOG.info("[at task:" + task + "]   [Time:" + System.currentTimeMillis() + "]");
                }
            } catch (Exception e) {
                e.printStackTrace();
                break;
            }
        }
    }

    /**
     * 創(chuàng)建一個(gè)最初為空的新 DelayQueue
     */
    private DelayQueue<Task> t = new DelayQueue<>();

    /**
     * 添加任務(wù),
     * time 延遲時(shí)間
     * task 任務(wù)
     * 用戶為問(wèn)題設(shè)置延遲時(shí)間
     */
    public void put(long time, Runnable task) {
        //轉(zhuǎn)換成ns
        long nanoTime = TimeUnit.NANOSECONDS.convert(time, TimeUnit.MILLISECONDS);
        //創(chuàng)建一個(gè)任務(wù)
        Task k = new Task(nanoTime, task);
        //將任務(wù)放在延遲的隊(duì)列中
        t.put(k);
    }

    /**
     * 結(jié)束訂單
     * @param task
     */
    public boolean endTask(Task<Runnable> task){
        return t.remove(task);
    }
}

使用方法

  1. 在容器初始化的時(shí)候調(diào)用init方法.
  2. 實(shí)現(xiàn)一個(gè)runnable接口的類,調(diào)用TaskQueueDaemonThread的put方法傳入進(jìn)去.
  3. 如果需要實(shí)現(xiàn)動(dòng)態(tài)的取消任務(wù)的話,需要task任務(wù)的類重新hashcode方法,最好用業(yè)務(wù)限制hashcode的沖突發(fā)生.
最后編輯于
?著作權(quán)歸作者所有,轉(zhuǎn)載或內(nèi)容合作請(qǐng)聯(lián)系作者
【社區(qū)內(nèi)容提示】社區(qū)部分內(nèi)容疑似由AI輔助生成,瀏覽時(shí)請(qǐng)結(jié)合常識(shí)與多方信息審慎甄別。
平臺(tái)聲明:文章內(nèi)容(如有圖片或視頻亦包括在內(nèi))由作者上傳并發(fā)布,文章內(nèi)容僅代表作者本人觀點(diǎn),簡(jiǎn)書(shū)系信息發(fā)布平臺(tái),僅提供信息存儲(chǔ)服務(wù)。

相關(guān)閱讀更多精彩內(nèi)容

友情鏈接更多精彩內(nèi)容