java利用delayedQueue实现本地的延迟队列

一、了解DelayQueue

DelayQueue是什么?

DelayQueue是一个无界的BlockingQueue,用于放置实现了Delayed接口的对象,其中的对象只能在其到期时才能从队列中取走。这种队列是有序的,即队头对象的延迟到期时间最长。

注意:不能将null元素放置到这种队列中。

DelayQueue能做什么?

在我们的业务中通常会有一些需求是这样的:

  • 淘宝订单业务:下单之后如果三十分钟之内没有付款就自动取消订单。
  • 饿了吗订餐通知:下单成功后60s之后给用户发送短信通知。

那么这类业务我们可以总结出一个特点:需要延迟工作。
由此的情况,就是我们的DelayQueue应用需求的产生。

二、怎么用DelayQueue来解决这类的问题

先声明一个Delayed的对象

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

/**
 * <p>
 * [任务调度系统]
 * <br>
 * [队列中要执行的任务]
 * </p>
 *
 * @author wangguangdong
 * @version 1.0
 * @Date 2015年11月22日19:46:39
 */
public class Task<T extends Runnable> implements Delayed {
 /**
  * 到期时间
  */
 private final long time;

 /**
  * 问题对象
  */
 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();
 }

 /**
  * 返回与此对象相关的剩余延迟时间,以给定的时间单位表示
  */
 @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;
 }

}

再实现一个管理延迟任务的类

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>
 * [任务调度系统]
 * <br>
 * [后台守护线程不断的执行检测工作]
 * </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);
 /**
  * 守护线程
  */
 private Thread daemonThread;

 /**
  * 初始化守护线程
  */
 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 {
    //从延迟队列中取值,如果没有对象过期则队列一直等待,
    Task t1 = t.take();
    if (t1 != null) {
     //修改问题的状态
     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;
   }
  }
 }

 /**
  * 创建一个最初为空的新 DelayQueue
  */
 private DelayQueue<Task> t = new DelayQueue<>();

 /**
  * 添加任务,
  * time 延迟时间
  * task 任务
  * 用户为问题设置延迟时间
  */
 public void put(long time, Runnable task) {
  //转换成ns
  long nanoTime = TimeUnit.NANOSECONDS.convert(time, TimeUnit.MILLISECONDS);
  //创建一个任务
  Task k = new Task(nanoTime, task);
  //将任务放在延迟的队列中
  t.put(k);
 }

 /**
  * 结束订单
  * @param task
  */
 public boolean endTask(Task<Runnable> task){
  return t.remove(task);
 }
}

使用方法

  • 在容器初始化的时候调用init方法.
  • 实现一个runnable接口的类,调用TaskQueueDaemonThread的put方法传入进去.
  • 如果需要实现动态的取消任务的话,需要task任务的类重新hashcode方法,最好用业务限制hashcode的冲突发生.

总结

以上就是这篇文章的全部内容了,希望本文的内容对大家的学习或者工作能带来一定的帮助,如果有疑问大家可以留言交流,谢谢大家对我们的支持。

(0)

相关推荐

  • Java多线程并发开发之DelayQueue使用示例

    在学习Java 多线程并发开发过程中,了解到DelayQueue类的主要作用:是一个无界的BlockingQueue,用于放置实现了Delayed接口的对象,其中的对象只能在其到期时才能从队列中取走.这种队列是有序的,即队头对象的延迟到期时间最长.注意:不能将null元素放置到这种队列中. Delayed,一种混合风格的接口,用来标记那些应该在给定延迟时间之后执行的对象.此接口的实现必须定义一个 compareTo 方法,该方法提供与此接口的 getDelay 方法一致的排序. 在网上看到了一些

  • java利用delayedQueue实现本地的延迟队列

    一.了解DelayQueue DelayQueue是什么? DelayQueue是一个无界的BlockingQueue,用于放置实现了Delayed接口的对象,其中的对象只能在其到期时才能从队列中取走.这种队列是有序的,即队头对象的延迟到期时间最长. 注意:不能将null元素放置到这种队列中. DelayQueue能做什么? 在我们的业务中通常会有一些需求是这样的: 淘宝订单业务:下单之后如果三十分钟之内没有付款就自动取消订单. 饿了吗订餐通知:下单成功后60s之后给用户发送短信通知. 那么这类

  • JAVA 实现延迟队列的方法

    延迟队列的需求各位应该在日常开发的场景中经常碰到.比如: 用户登录之后5分钟给用户做分类推送: 用户多少天未登录给用户做召回推送: 定期检查用户当前退款账单是否被商家处理等等场景. 一般这种场景和定时任务还是有很大的区别,定时任务是你知道任务多久该跑一次或者什么时候只跑一次,这个时间是确定的.延迟队列是当某个事件发生的时候需要延迟多久触发配套事件,引子事件发生的时间不是固定的. 业界目前也有很多实现方案,单机版的方案就不说了,现在也没有哪个公司还是单机版的服务,今天我们一一探讨各种方案的大致实现

  • Java利用Redis实现消息队列的示例代码

    本文介绍了Java利用Redis实现消息队列的示例代码,分享给大家,具体如下: 应用场景 为什么要用redis? 二进制存储.java序列化传输.IO连接数高.连接频繁 一.序列化 这里编写了一个java序列化的工具,主要是将对象转化为byte数组,和根据byte数组反序列化成java对象; 主要是用到了ByteArrayOutputStream和ByteArrayInputStream; 注意:每个需要序列化的对象都要实现Serializable接口; 其代码如下: package Utils

  • Java延迟队列原理与用法实例详解

    本文实例讲述了Java延迟队列原理与用法.分享给大家供大家参考,具体如下: 延时队列,第一他是个队列,所以具有对列功能第二就是延时,这就是延时对列,功能也就是将任务放在该延时对列中,只有到了延时时刻才能从该延时对列中获取任务否则获取不到-- 应用场景比较多,比如延时1分钟发短信,延时1分钟再次执行等,下面先看看延时队列demo之后再看延时队列在项目中的使用: 简单的延时队列要有三部分:第一实现了Delayed接口的消息体.第二消费消息的消费者.第三存放消息的延时队列,那下面就来看看延时队列dem

  • java并发中DelayQueue延迟队列原理剖析

    介绍 DelayQueue队列是一个延迟队列,DelayQueue中存放的元素必须实现Delayed接口的元素,实现接口后相当于是每个元素都有个过期时间,当队列进行take获取元素时,先要判断元素有没有过期,只有过期的元素才能出队操作,没有过期的队列需要等待剩余过期时间才能进行出队操作. 源码分析 DelayQueue队列内部使用了PriorityQueue优先队列来进行存放数据,它采用的是二叉堆进行的优先队列,使用ReentrantLock锁来控制线程同步,由于内部元素是采用的Priority

  • Java 延迟队列的常用的实现方式

    延迟队列的使用场景还比较多,例如: 1.超时未收到支付回调,主动查询支付状态: 2.规定时间内,订单未支付,自动取消: ... 总之,但凡需要在未来的某个确定的时间点执行检查的场景中都可以用延迟队列. 常见的手段主要有:定时任务扫描.RocketMQ延迟队列.Java自动的延迟队列.监听Redis Key过期等等 1.  DelayQueue 首先,定义一个延迟任务 package com.cjs.example; import lombok.Data; import java.util.con

  • Java Kafka实现延迟队列的示例代码

    目录 基于kafka如何实现延迟队列 完善细节 Java代码实现 还需要做什么 kafka作为一个使用广泛的消息队列,很多人都不会陌生,但当你在网上搜索“kafka 延迟队列”,出现的都是一些讲解时间轮或者只是提供了一些思路,并没有一份真实可用的代码实现,今天我们就来打破这个现象,提供一份可运行的代码,抛砖引玉,吸引更多的大神来分享. 基于kafka如何实现延迟队列 想要解决一个问题,我们需要先分解问题.kafka作为一个高性能的消息队列,只要消费能力足够,发出的消息都是会立刻收到的,因此我们需

  • 详解Java线程池队列中的延迟队列DelayQueue

    目录 DelayQueue延迟队列 DelayQueue使用场景 DelayQueue属性 DelayQueue构造方法 实现Delayed接口使用示例 DelayQueue总结 在阻塞队里中,除了对元素进行增加和删除外,我们可以把元素的删除做一个延迟的处理,即使用DelayQueue的方法.本文就来和大家聊聊Java线程池队列中的DelayQueue—延迟队列 public enum QueueTypeEnum { ARRAY_BLOCKING_QUEUE(1, "ArrayBlockingQ

  • Redis延迟队列和分布式延迟队列的简答实现

    最近,又重新学习了下Redis,Redis不仅能快还能慢,简直利器,今天就为大家介绍一下Redis延迟队列和分布式延迟队列的简单实现. 在我们的工作中,很多地方使用延迟队列,比如订单到期没有付款取消订单,制订一个提醒的任务等都需要延迟队列,那么我们需要实现延迟队列.我们本文的梗概如下,同学们可以选择性阅读. 1. 实现一个简单的延迟队列. 我们知道目前JAVA可以有DelayedQueue,我们首先开一个DelayQueue的结构类图.DelayQueue实现了Delay.BlockingQue

  • java开发微服务架构设计消息队列的水有多深

    目录 消息队列的作用 消息队列的设计难题 处理并发和顺序消息 处理重复消息 编写幂等消息处理器 跟踪消息并丢弃重复消息 处理事务性消息 使用数据库表作为消息队列 使用事务日志发布事件 RocketMQ事务消息解决方案 很多人在做架构设计时往往会"过度设计",简单问题复杂化,上来就引一堆中间件,我想大概原因主要有下面两点: 为了秀(学)技术而架构 我们常说技术是为业务服务的,不能为了技术而技术,为了秀技术引入一堆复杂架构这是要不得的. 考虑问题不全面,或者说广度不够,不知道如何简单化 举

随机推荐