Java实现延迟队列案例

wen java案例 3

本文目录导读:

Java实现延迟队列案例

  1. 基于 JDK DelayQueue 的实现
  2. 基于 Redis 的延迟队列实现
  3. 使用 ScheduledExecutorService 实现延迟任务
  4. 应用场景示例:订单超时处理
  5. 注意事项

我来为你提供几个Java实现延迟队列的完整案例,从简单到复杂。

基于 JDK DelayQueue 的实现

延迟消息实体类

import java.util.concurrent.Delayed;
import java.util.concurrent.TimeUnit;
/**
 * 延迟消息实体
 */
public class DelayMessage implements Delayed {
    private String messageId;
    private String content;
    private long delayTime; // 延迟时间(毫秒)
    private long expireTime; // 过期时间(毫秒)
    public DelayMessage(String messageId, String content, long delayTime) {
        this.messageId = messageId;
        this.content = content;
        this.delayTime = delayTime;
        this.expireTime = System.currentTimeMillis() + delayTime;
    }
    @Override
    public long getDelay(TimeUnit unit) {
        long diff = expireTime - System.currentTimeMillis();
        return unit.convert(diff, TimeUnit.MILLISECONDS);
    }
    @Override
    public int compareTo(Delayed o) {
        if (this.expireTime < ((DelayMessage) o).expireTime) {
            return -1;
        } else if (this.expireTime > ((DelayMessage) o).expireTime) {
            return 1;
        }
        return 0;
    }
    // getters and setters
    public String getMessageId() {
        return messageId;
    }
    public void setMessageId(String messageId) {
        this.messageId = messageId;
    }
    public String getContent() {
        return content;
    }
    public void setContent(String content) {
        this.content = content;
    }
    public long getDelayTime() {
        return delayTime;
    }
    public void setDelayTime(long delayTime) {
        this.delayTime = delayTime;
    }
    public long getExpireTime() {
        return expireTime;
    }
    public void setExpireTime(long expireTime) {
        this.expireTime = expireTime;
    }
    @Override
    public String toString() {
        return "DelayMessage{" +
                "messageId='" + messageId + '\'' +
                ", content='" + content + '\'' +
                ", delayTime=" + delayTime +
                ", expireTime=" + expireTime +
                '}';
    }
}

延迟队列服务类

import java.util.concurrent.DelayQueue;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;
import java.util.concurrent.TimeUnit;
/**
 * 延迟队列管理器
 */
public class DelayQueueManager {
    private static final DelayQueue<DelayMessage> delayQueue = new DelayQueue<>();
    private static final ExecutorService executorService = Executors.newFixedThreadPool(5);
    // 启动消费者
    public void start() {
        System.out.println("延迟队列服务启动...");
        for (int i = 0; i < 5; i++) {
            executorService.execute(() -> {
                while (true) {
                    try {
                        // 获取并移除延迟队列的头元素,如果没有到期则返回 null
                        DelayMessage message = delayQueue.poll();
                        if (message != null) {
                            System.out.println("消费者收到消息: " + message);
                            // 处理消息
                            processMessage(message);
                        } else {
                            // 没有消息,等待一段时间再检查
                            Thread.sleep(1000);
                        }
                    } catch (InterruptedException e) {
                        Thread.currentThread().interrupt();
                        break;
                    }
                }
            });
        }
    }
    // 处理消息
    private void processMessage(DelayMessage message) {
        // 这里可以放具体的业务处理逻辑
        System.out.println("处理延迟消息: " + message.getContent() + 
                          ", 当前时间: " + System.currentTimeMillis());
    }
    // 添加延迟消息
    public void addMessage(DelayMessage message) {
        System.out.println("添加延迟消息: " + message + 
                          ", 当前时间: " + System.currentTimeMillis());
        delayQueue.put(message);
    }
    // 停止服务
    public void stop() {
        executorService.shutdown();
        System.out.println("延迟队列服务已停止");
    }
    public static void main(String[] args) throws InterruptedException {
        DelayQueueManager manager = new DelayQueueManager();
        // 启动服务
        manager.start();
        // 添加测试消息,延迟时间分别为 3秒、5秒、10秒
        DelayMessage msg1 = new DelayMessage("001", "订单1超时关闭", 3000);
        DelayMessage msg2 = new DelayMessage("002", "订单2超时关闭", 5000);
        DelayMessage msg3 = new DelayMessage("003", "订单3超时关闭", 10000);
        manager.addMessage(msg1);
        manager.addMessage(msg2);
        manager.addMessage(msg3);
        // 运行一段时间后停止
        Thread.sleep(15000);
        manager.stop();
    }
}

基于 Redis 的延迟队列实现

依赖准备

<dependency>
    <groupId>redis.clients</groupId>
    <artifactId>jedis</artifactId>
    <version>4.4.0</version>
</dependency>

Redis延迟队列实现

import redis.clients.jedis.Jedis;
import redis.clients.jedis.JedisPool;
import redis.clients.jedis.Tuple;
import com.alibaba.fastjson.JSON;
import java.util.Set;
import java.util.UUID;
import java.util.concurrent.Executors;
import java.util.concurrent.ScheduledExecutorService;
import java.util.concurrent.TimeUnit;
/**
 * Redis延迟队列实现(使用 ZSet)
 */
public class RedisDelayQueue {
    private static final String DELAY_QUEUE_KEY = "delay:queue";
    private static final String READY_QUEUE_KEY = "ready:queue";
    private JedisPool jedisPool;
    private ScheduledExecutorService scheduledExecutor;
    public RedisDelayQueue() {
        // 初始化Redis连接池
        jedisPool = new JedisPool("localhost", 6379);
        scheduledExecutor = Executors.newScheduledThreadPool(4);
    }
    /**
     * 添加延迟任务
     * @param task 任务内容
     * @param delay 延迟时间
     * @param unit 时间单位
     */
    public String addTask(String task, long delay, TimeUnit unit) {
        String taskId = UUID.randomUUID().toString();
        long expireTime = System.currentTimeMillis() + unit.toMillis(delay);
        try (Jedis jedis = jedisPool.getResource()) {
            // 将任务添加到延迟队列,score为过期时间
            jedis.zadd(DELAY_QUEUE_KEY, expireTime, taskId);
            // 存储任务内容
            TaskData taskData = new TaskData(taskId, task, expireTime);
            jedis.hset("task:data", taskId, JSON.toJSONString(taskData));
            System.out.println("添加延迟任务: " + taskId + ", 过期时间: " + expireTime);
            return taskId;
        }
    }
    /**
     * 检查并转移到期任务
     */
    private void transferExpiredTasks() {
        try (Jedis jedis = jedisPool.getResource()) {
            // 获取到期的任务
            Set<String> expiredTasks = jedis.zrangeByScore(DELAY_QUEUE_KEY, 
                                                          0, System.currentTimeMillis());
            if (!expiredTasks.isEmpty()) {
                for (String taskId : expiredTasks) {
                    // 将任务从延迟队列移动到就绪队列
                    String taskData = jedis.hget("task:data", taskId);
                    if (taskData != null) {
                        // 添加到就绪队列
                        jedis.lpush(READY_QUEUE_KEY, taskData);
                        // 移除延迟队列中的任务
                        jedis.zrem(DELAY_QUEUE_KEY, taskId);
                        // 清理任务数据
                        jedis.hdel("task:data", taskId);
                        System.out.println("任务已到期,转移到就绪队列: " + taskId);
                    }
                }
            }
        }
    }
    /**
     * 启动延迟队列服务
     */
    public void start() {
        // 定时检查到期的任务(每秒执行一次)
        scheduledExecutor.scheduleAtFixedRate(this::transferExpiredTasks, 
                                              0, 1, TimeUnit.SECONDS);
        // 启动消费者线程处理就绪队列
        for (int i = 0; i < 3; i++) {
            scheduledExecutor.execute(() -> {
                while (!Thread.currentThread().isInterrupted()) {
                    try (Jedis jedis = jedisPool.getResource()) {
                        // 从就绪队列获取任务(阻塞式)
                        String taskData = jedis.brpop(0, READY_QUEUE_KEY);
                        if (taskData != null) {
                            // 处理任务
                            TaskData data = JSON.parseObject(taskData, TaskData.class);
                            processTask(data);
                        }
                    } catch (Exception e) {
                        e.printStackTrace();
                    }
                }
            });
        }
        System.out.println("Redis延迟队列服务已启动");
    }
    /**
     * 处理任务
     */
    private void processTask(TaskData taskData) {
        System.out.println("处理任务: " + taskData.getTaskId() + 
                          ", 内容: " + taskData.getTask() + 
                          ", 当前时间: " + System.currentTimeMillis());
        // 在这里执行具体的业务逻辑
    }
    /**
     * 停止服务
     */
    public void stop() {
        scheduledExecutor.shutdown();
        jedisPool.close();
        System.out.println("Redis延迟队列服务已停止");
    }
    // 任务数据类
    static class TaskData {
        private String taskId;
        private String task;
        private long expireTime;
        public TaskData() {}
        public TaskData(String taskId, String task, long expireTime) {
            this.taskId = taskId;
            this.task = task;
            this.expireTime = expireTime;
        }
        public String getTaskId() { return taskId; }
        public void setTaskId(String taskId) { this.taskId = taskId; }
        public String getTask() { return task; }
        public void setTask(String task) { this.task = task; }
        public long getExpireTime() { return expireTime; }
        public void setExpireTime(long expireTime) { this.expireTime = expireTime; }
    }
    public static void main(String[] args) throws InterruptedException {
        RedisDelayQueue delayQueue = new RedisDelayQueue();
        // 启动服务
        delayQueue.start();
        // 添加测试任务
        delayQueue.addTask("任务1:超时关闭订单", 5, TimeUnit.SECONDS);
        delayQueue.addTask("任务2:超时确认收货", 10, TimeUnit.SECONDS);
        delayQueue.addTask("任务3:定时发送提醒", 15, TimeUnit.SECONDS);
        // 运行一段时间
        Thread.sleep(20000);
        delayQueue.stop();
    }
}

使用 ScheduledExecutorService 实现延迟任务

import java.util.concurrent.Executors;
import java.util.concurrent.ScheduledExecutorService;
import java.util.concurrent.TimeUnit;
/**
 * 基于 ScheduledExecutorService 的延迟任务实现
 */
public class ScheduledDelayQueue {
    private ScheduledExecutorService executorService;
    private int threadPoolSize;
    public ScheduledDelayQueue(int threadPoolSize) {
        this.threadPoolSize = threadPoolSize;
        this.executorService = Executors.newScheduledThreadPool(threadPoolSize);
    }
    /**
     * 添加一次性延迟任务
     */
    public void addOneTimeTask(Runnable task, long delay, TimeUnit unit) {
        System.out.println("添加一次性延迟任务: " + delay + " " + unit);
        executorService.schedule(task, delay, unit);
    }
    /**
     * 添加周期性任务(固定频率)
     */
    public void addFixedRateTask(Runnable task, long initialDelay, 
                                 long period, TimeUnit unit) {
        System.out.println("添加固定频率任务");
        executorService.scheduleAtFixedRate(task, initialDelay, period, unit);
    }
    /**
     * 添加周期性任务(固定延迟)
     */
    public void addFixedDelayTask(Runnable task, long initialDelay, 
                                  long delay, TimeUnit unit) {
        System.out.println("添加固定延迟任务");
        executorService.scheduleWithFixedDelay(task, initialDelay, delay, unit);
    }
    /**
     * 停止服务
     */
    public void shutdown() {
        executorService.shutdown();
        try {
            if (!executorService.awaitTermination(60, TimeUnit.SECONDS)) {
                executorService.shutdownNow();
            }
        } catch (InterruptedException e) {
            executorService.shutdownNow();
            Thread.currentThread().interrupt();
        }
        System.out.println("调度服务已停止");
    }
    public static void main(String[] args) throws InterruptedException {
        ScheduledDelayQueue queue = new ScheduledDelayQueue(4);
        // 一次性延迟任务
        queue.addOneTimeTask(() -> {
            System.out.println("3秒后执行一次性延迟任务: " + 
                              System.currentTimeMillis());
        }, 3, TimeUnit.SECONDS);
        // 固定频率任务(每2秒执行一次)
        queue.addFixedRateTask(() -> {
            System.out.println("固定频率任务,每2秒执行: " + 
                              System.currentTimeMillis());
        }, 0, 2, TimeUnit.SECONDS);
        // 固定延迟任务(上一次执行完后延迟1秒)
        queue.addFixedDelayTask(() -> {
            System.out.println("固定延迟任务: " + System.currentTimeMillis());
            try {
                Thread.sleep(1000); // 模拟任务执行时间
            } catch (InterruptedException e) {
                Thread.currentThread().interrupt();
            }
        }, 0, 1, TimeUnit.SECONDS);
        // 运行10秒后停止
        Thread.sleep(10000);
        queue.shutdown();
    }
}

应用场景示例:订单超时处理

import java.util.concurrent.DelayQueue;
import java.util.concurrent.TimeUnit;
/**
 * 订单超时处理示例
 */
public class OrderTimeoutHandler {
    private static final DelayQueue<DelayMessage> orderDelayQueue = new DelayQueue<>();
    // 订单信息类
    static class OrderInfo {
        String orderId;
        String userId;
        double amount;
        long createTime;
        public OrderInfo(String orderId, String userId, double amount, long createTime) {
            this.orderId = orderId;
            this.userId = userId;
            this.amount = amount;
            this.createTime = createTime;
        }
        @Override
        public String toString() {
            return "订单{" +
                    "订单号='" + orderId + '\'' +
                    ", 用户ID='" + userId + '\'' +
                    ", 金额=" + amount +
                    ", 创建时间=" + createTime +
                    '}';
        }
    }
    // 启动延迟队列消费者
    public static void startConsumer() {
        Thread consumer = new Thread(() -> {
            while (true) {
                try {
                    DelayMessage message = orderDelayQueue.take();
                    String orderId = message.getMessageId();
                    System.out.println("订单超时处理: " + orderId);
                    // 检查订单状态
                    checkOrderStatus(orderId);
                } catch (InterruptedException e) {
                    Thread.currentThread().interrupt();
                    break;
                }
            }
        });
        consumer.start();
    }
    // 创建订单并加入延迟队列
    public static OrderInfo createOrder(String orderId, String userId, double amount) {
        OrderInfo order = new OrderInfo(orderId, userId, amount, System.currentTimeMillis());
        System.out.println("创建订单: " + order);
        // 添加延迟任务,30分钟超时
        DelayMessage message = new DelayMessage(orderId, 
                                               "订单超时关闭", 
                                               30 * 60 * 1000);
        orderDelayQueue.put(message);
        return order;
    }
    // 检查订单状态
    private static void checkOrderStatus(String orderId) {
        // 这里需要查询数据库获取订单状态
        System.out.println("检查订单 " + orderId + " 的支付状态");
        // 如果未支付,则关闭订单
        System.out.println("订单 " + orderId + " 超时未支付,已自动关闭");
    }
    public static void main(String[] args) throws InterruptedException {
        // 启动消费者
        startConsumer();
        // 模拟创建订单
        createOrder("O2024001", "U1001", 200.50);
        createOrder("O2024002", "U1002", 150.00);
        createOrder("O2024003", "U1003", 899.00);
        // 演示用,手动添加一个3秒过期的任务
        DelayMessage testMessage = new DelayMessage("TEST001", "测试过期", 3000);
        orderDelayQueue.put(testMessage);
        System.out.println("添加测试消息,3秒后过期");
        // 运行一段时间
        Thread.sleep(10000);
    }
}

注意事项

  1. 延迟队列选择

    • 小型应用、单机部署:使用 JDK DelayQueue
    • 分布式系统、需要持久化:使用 Redis ZSet
    • 简单定时任务:使用 ScheduledExecutorService
  2. 可靠性

    • DelayedQueue 不持久化,程序重启会丢失
    • Redis 延迟队列需要额外处理持久化
  3. 性能考虑

    • 大批量任务时推荐使用 Redis
    • 考虑使用时间轮算法(如 Netty 的 HashedWheelTimer)提高性能
  4. 生产环境建议

    • 结合消息队列(如 RabbitMQ、Kafka)的延迟消息功能
    • 使用成熟的框架(如 Quartz、XXL-JOB)

这些示例代码可以直接运行,根据实际业务需求进行修改即可。

抱歉,评论功能暂时关闭!