Skip to content

基于Java的延迟队列

任务类实现Delayed接口

定义任务类实现Delayed接口并实现getDelay()和compareTo()方法来实现延迟队列

java
/**
 * @ClassName DelayQueue
 * @Author liupengxiang
 * @Date 2024/5/15 16:44
 * @Description 门票信息延迟队列任务
 */
public class TicketTask implements Delayed {
    // 订单号
    private String orderNo;
    // 延迟时间
    private long delayTime;
    // 到期时间
    private long expire;
    
    public TicketTask(String orderNo, long delay, TimeUnit unit) {
        this.orderNo = orderNo;
        this.delayTime = unit.toNanos(delay);
        this.expire = System.nanoTime() + this.delayTime;
    }
    public TicketTask(String orderNo, long delayTime, long expire) {
        this.orderNo = orderNo;
        this.delayTime = delayTime;
        this.expire = expire;
    }
    @Override
    public long getDelay(@NotNull TimeUnit unit) {
        return unit.convert(this.expire - System.nanoTime(), TimeUnit.NANOSECONDS);
    }

    @Override
    public int compareTo(@NotNull Delayed o) {
        TicketTask other = (TicketTask) o;
        long diff = this.expire - other.expire;
        if (diff < 0) {
            return -1;
        } else if (diff > 0) {
            return 1;
        } else {
            return 0;
        }
    }
    public String getOrderNo() {
        return orderNo;
    }
    public long getDelayTime() {
        return delayTime;
    }
    public long getExpire() {
        return expire;
    }

    @Override
    public String toString() {
        return "TicketTask{" +
                "orderNo='" + orderNo + '\'' +
                ", delayTime=" + delayTime +
                ", expire=" + expire +
                '}';
    }
}

声明全局队列对象

java
/**
 * @ClassName OverQueue
 * @Author liupengxiang
 * @Date 2024/5/15 16:58
 * @Description 延迟队列类变量
 */
public class OverQueue {
    private static volatile OverQueue queue;
    // 门票状态延迟队列
    private DelayQueue<TicketTask> ticketQueue = new DelayQueue<>();
    // 出票延迟队列
    private DelayQueue<SubmitTicketTask> submitTicketTasks = new DelayQueue<>();

    private OverQueue() {
    }
    public static synchronized OverQueue getInstance() {
        if (queue == null) {
            synchronized (OverQueue.class) {
                if (queue == null) {
                    queue = new OverQueue();
                }
            }
        }
        return queue;
    }
    public DelayQueue<TicketTask> getTicketQueue() {
        return ticketQueue;
    }

    public DelayQueue<SubmitTicketTask> getSubmitTicketTasks() {
        return submitTicketTasks;
    }
}

队列消费

配合redis实现队列持久化

java

/**
 * @ClassName DelayQueueConsume
 * @Author liupengxiang
 * @Date 2024/5/15 17:26
 * @Description 消费延迟队列并恢复 Redis 中持久化的任务
 */
@Component
@Slf4j
public class TicketTaskQueueConsumer {

    @Resource
    private ITicketService ticketService;
    @Resource
    StringRedisTemplate stringRedisTemplate;
    @PostConstruct
    public void ticketCheck() {
        DelayQueue<TicketTask> delayQueue = OverQueue.getInstance().getTicketQueue();
        new Thread(() -> {
            log.info("启动门票状态任务队列");
            // 恢复持久化数据
            String orderNoStr = stringRedisTemplate.opsForValue().get(RedisPrefix.KEY_TICKET_DELAY_QUEUE);
            if (StringUtils.isNotBlank(orderNoStr)) {
                log.info("缓存门票状态任务入队列");
                List<String> orderNos = JSONUtil.parseArray(orderNoStr).stream().map(String::valueOf).collect(Collectors.toList());
                for (String orderNo : orderNos) {
                    TicketTask ticketTask = new TicketTask(orderNo, 2, TimeUnit.SECONDS);
                    delayQueue.add(ticketTask);
                }
            }
            while (true) {
                try {
                    Thread.sleep(1000);
                    if (delayQueue.isEmpty()) {
                        continue;
                    }
                    TicketTask ticketTask = delayQueue.poll();
                    if (ticketTask == null) {
                        continue;
                    }
                    String orderNo = ticketTask.getOrderNo();
                    log.info("门票状态任务队列---->处理门票队列数据");
                    ticketService.ticketStatus(orderNo, true);
                } catch (InterruptedException e) {
                    log.info("门票状态任务队列---->延迟队列中断异常");
                    e.printStackTrace();
                }
            }
        }).start();
    }
    @PreDestroy
    public void destroy() {
        // 持久化延迟队列数据
        log.info("缓存门票数据");
        DelayQueue<TicketTask> delayQueue = OverQueue.getInstance().getTicketQueue();
        Object[] array = delayQueue.toArray();
        log.info("将要入缓存的门票状态延迟队列:{}",array);
        List<String> list = new ArrayList<>();
        for (Object o : array) {
            JSONObject jsonObject = JSONUtil.parseObj(o);
            log.info("延迟队列数据:{}", jsonObject);
            String orderNo = jsonObject.getStr("orderNo");
            list.add(orderNo);
        }
        stringRedisTemplate.opsForValue().set(RedisPrefix.KEY_TICKET_DELAY_QUEUE,JSONUtil.toJsonStr(list));
    }
}