基于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));
}
}