延迟消息案例

wen java案例 1

本文目录导读:

延迟消息案例

  1. 基于 Redis 的延迟队列(ZSET实现)
  2. RabbitMQ 延迟消息(死信队列实现)
  3. 使用 RocketMQ 延迟消息(内置延迟等级)
  4. Java 本地延迟队列(ScheduledExecutorService)
  5. 场景化完整案例
  6. 技术选型建议
  7. 注意事项

我将为您提供几个常见的延迟消息实现案例,涵盖不同技术栈和场景。

基于 Redis 的延迟队列(ZSET实现)

场景:订单超时自动取消

@Service
public class OrderDelayService {
    @Autowired
    private StringRedisTemplate redisTemplate;
    private static final String DELAY_QUEUE_KEY = "order:delay:queue";
    /**
     * 添加延迟任务
     * @param orderId 订单ID
     * @param delayTime 延迟时间(毫秒)
     */
    public void addDelayTask(String orderId, long delayTime) {
        // 到期时间 = 当前时间 + 延迟时间
        long expireTime = System.currentTimeMillis() + delayTime;
        redisTemplate.opsForZSet().add(DELAY_QUEUE_KEY, orderId, expireTime);
        log.info("订单[{}]已加入延迟队列,将在{}秒后执行", orderId, delayTime/1000);
    }
    /**
     * 消费者:定时扫描处理过期任务
     */
    @Scheduled(fixedDelay = 1000) // 每秒执行一次
    public void consumeDelayTask() {
        // 获取当前时间
        long currentTime = System.currentTimeMillis();
        // 获取所有到期的任务(分数 <= 当前时间)
        Set<String> expiredOrders = redisTemplate.opsForZSet()
            .rangeByScore(DELAY_QUEUE_KEY, 0, currentTime);
        if (CollectionUtils.isEmpty(expiredOrders)) {
            return;
        }
        for (String orderId : expiredOrders) {
            // 原子性移除任务,避免并发重复消费
            Long removed = redisTemplate.opsForZSet()
                .remove(DELAY_QUEUE_KEY, orderId);
            if (removed != null && removed > 0) {
                // 执行具体业务逻辑
                handleExpiredOrder(orderId);
                log.info("订单[{}]已超时,执行自动取消", orderId);
            }
        }
    }
    private void handleExpiredOrder(String orderId) {
        // 业务处理:检查订单状态,更新为已取消等
        Order order = orderMapper.selectById(orderId);
        if (order != null && "待支付".equals(order.getStatus())) {
            order.setStatus("已取消");
            orderMapper.updateById(order);
        }
    }
}

使用示例

@Test
public void testDelayMessage() throws InterruptedException {
    // 1. 创建订单
    String orderId = "ORDER_2024001";
    orderService.createOrder(orderId);
    // 2. 添加延迟任务:30分钟后自动取消
    orderDelayService.addDelayTask(orderId, 30 * 60 * 1000);
    // 3. 等待测试
    Thread.sleep(5000);
}

RabbitMQ 延迟消息(死信队列实现)

场景:支付超时提醒

@Configuration
public class RabbitDelayConfig {
    // 交换机
    public static final String ORDER_EXCHANGE = "order.exchange";
    public static final String ORDER_DELAY_EXCHANGE = "order.delay.exchange";
    // 队列
    public static final String ORDER_QUEUE = "order.queue";
    public static final String ORDER_DELAY_QUEUE = "order.delay.queue";
    // 路由键
    public static final String ORDER_ROUTING_KEY = "order.routing";
    public static final String ORDER_DELAY_ROUTING_KEY = "order.delay.routing";
    // 延迟时间
    public static final long DELAY_TIME = 30 * 60 * 1000; // 30分钟
    @Bean
    public Queue delayQueue() {
        Map<String, Object> args = new HashMap<>();
        // 设置死信交换机
        args.put("x-dead-letter-exchange", ORDER_EXCHANGE);
        // 设置死信路由键
        args.put("x-dead-letter-routing-key", ORDER_ROUTING_KEY);
        // 设置消息过期时间
        args.put("x-message-ttl", DELAY_TIME);
        return QueueBuilder.durable(ORDER_DELAY_QUEUE).withArguments(args).build();
    }
    @Bean
    public Queue orderQueue() {
        return QueueBuilder.durable(ORDER_QUEUE).build();
    }
    @Bean
    public DirectExchange exchange() {
        return new DirectExchange(ORDER_EXCHANGE);
    }
    @Bean
    public DirectExchange delayExchange() {
        return new DirectExchange(ORDER_DELAY_EXCHANGE);
    }
    // 绑定
    @Bean
    public Binding delayBinding() {
        return BindingBuilder.bind(delayQueue())
            .to(delayExchange()).with(ORDER_DELAY_ROUTING_KEY);
    }
    @Bean
    public Binding orderBinding() {
        return BindingBuilder.bind(orderQueue())
            .to(exchange()).with(ORDER_ROUTING_KEY);
    }
}

生产者与消费者

@Service
public class PayTimeoutService {
    @Autowired
    private RabbitTemplate rabbitTemplate;
    /**
     * 发送延迟消息
     */
    public void sendDelayMessage(String orderId) {
        // 创建消息
        Message message = MessageBuilder
            .withBody(orderId.getBytes(StandardCharsets.UTF_8))
            .setContentType(MessageProperties.CONTENT_TYPE_JSON)
            .build();
        // 发送到延迟队列
        rabbitTemplate.send(ORDER_DELAY_EXCHANGE, ORDER_DELAY_ROUTING_KEY, message);
        log.info("支付提醒延迟消息已发送,订单ID:{}", orderId);
    }
    /**
     * 消费者:处理超时未支付的订单
     */
    @RabbitListener(queues = ORDER_QUEUE)
    public void handlePayTimeout(String orderId) {
        log.info("收到延迟消息,订单ID:{}", orderId);
        // 业务处理:检查支付状态
        Order order = orderService.getById(orderId);
        if (order != null && "未支付".equals(order.getStatus())) {
            // 发送支付提醒短信或推送
            smsService.sendPayRemind(order.getMemberId(), orderId);
            log.info("已向用户发送支付提醒");
        }
    }
}

使用 RocketMQ 延迟消息(内置延迟等级)

场景:秒杀活动结束通知

@Service
public class SeckillNotificationService {
    @Autowired
    private RocketMQTemplate rocketMQTemplate;
    private static final String TOPIC = "seckill-notification";
    /**
     * 发送延迟通知
     */
    public void sendSeckillEndNotification(String activityId) {
        // RocketMQ延迟等级:1s 5s 10s 30s 1m 2m 3m 4m 5m 6m 7m 8m 9m 10m 20m 30m 1h 2h
        // 这里使用 30分钟 = 延迟等级 16
        Message<String> message = MessageBuilder
            .withPayload(activityId)
            .build();
        // 设置延迟等级
        rocketMQTemplate.syncSend(
            TOPIC, 
            message, 
            3000, 
            // 延迟等级,16 表示30分钟
            10  // 延迟等级3 => 10秒
        );
        log.info("秒杀活动结束通知已发送,活动ID:{}", activityId);
    }
    /**
     * 消费者
     */
    @Service
    @RocketMQMessageListener(
        topic = TOPIC,
        consumerGroup = "seckill-notification-group"
    )
    public static class NotificationConsumer implements RocketMQListener<String> {
        @Override
        public void onMessage(String activityId) {
            log.info("收到秒杀活动结束通知:{}", activityId);
            // 批量处理参与用户
            List<Long> userIds = seckillService.getParticipants(activityId);
            notificationService.batchSend(userIds, "秒杀活动已结束");
        }
    }
}

Java 本地延迟队列(ScheduledExecutorService)

场景:定时任务调度

@Component
public class LocalDelayQueueExample {
    private final ScheduledExecutorService scheduler = 
        Executors.newScheduledThreadPool(10);
    private final DelayQueue<DelayTask> delayQueue = new DelayQueue<>();
    /**
     * 任务实体
     */
    public static class DelayTask implements Delayed {
        private final String taskId;
        private final long executeTime;
        private final Runnable task;
        public DelayTask(String taskId, long delayMillis, Runnable task) {
            this.taskId = taskId;
            this.executeTime = System.currentTimeMillis() + delayMillis;
            this.task = task;
        }
        @Override
        public long getDelay(TimeUnit unit) {
            return unit.convert(executeTime - System.currentTimeMillis(), 
                TimeUnit.MILLISECONDS);
        }
        @Override
        public int compareTo(Delayed other) {
            return Long.compare(this.executeTime, 
                ((DelayTask) other).executeTime);
        }
        public void execute() {
            task.run();
        }
    }
    /**
     * 添加延迟任务
     */
    public void addTask(String taskId, long delayMillis, Runnable task) {
        DelayTask delayTask = new DelayTask(taskId, delayMillis, task);
        delayQueue.offer(delayTask);
        log.info("延迟任务[{}]已添加,延迟{}毫秒", taskId, delayMillis);
    }
    /**
     * 启动消费者线程
     */
    @PostConstruct
    public void startConsumer() {
        scheduler.scheduleAtFixedRate(() -> {
            try {
                DelayTask task = delayQueue.poll(1, TimeUnit.SECONDS);
                if (task != null) {
                    log.info("开始执行延迟任务:{}", task);
                    task.execute();
                }
            } catch (InterruptedException e) {
                Thread.currentThread().interrupt();
                log.error("延迟队列消费失败", e);
            }
        }, 0, 100, TimeUnit.MILLISECONDS);
    }
    /**
     * 使用示例
     */
    public void demo() {
        // 添加不同类型的延迟任务
        addTask("task-1", 5000, () -> {
            System.out.println("5秒后执行的定时任务");
        });
        addTask("task-2", 10000, () -> {
            System.out.println("10秒后执行的定时任务");
        });
    }
}

场景化完整案例

电商交易场景:订单状态流转

@Service
public class OrderStatusFlowService {
    // 不同业务的延迟时间
    private static final long PAY_TIMEOUT = 30 * 60 * 1000;      // 30分钟支付超时
    private static final long CONFIRM_TIMEOUT = 7 * 24 * 3600 * 1000; // 7天自动确认收货
    private static final long REFUND_TIMEOUT = 24 * 3600 * 1000;  // 24小时退款超时
    @Autowired
    private OrderDelayService orderDelayService;
    /**
     * 订单创建后设置支付超时
     */
    public void createOrderAndSetTimeout(String orderId) {
        // 创建订单
        orderService.create(orderId);
        // 设置30分钟支付超时
        orderDelayService.addDelayTask(
            orderId,
            PAY_TIMEOUT,
            DelayTaskType.PAY_TIMEOUT
        );
    }
    /**
     * 支付成功后设置自动确认收货
     */
    public void paySuccessAndSetConfirm(String orderId) {
        orderService.payMark(orderId);
        // 7天后自动确认收货
        delayTaskManager.addTask(
            orderId,
            CONFIRM_TIMEOUT,
            DelayTaskType.AUTO_CONFIRM
        );
    }
    /**
     * 统一处理延迟任务
     */
    @Scheduled(cron = "*/5 * * * * ?")
    public void processDelayTasks() {
        List<DelayTask> tasks = delayTaskManager.getExpiredTasks();
        for (DelayTask task : tasks) {
            switch (task.getTaskType()) {
                case PAY_TIMEOUT:
                    handlePayTimeout(task.getBizId());
                    break;
                case AUTO_CONFIRM:
                    handleAutoConfirm(task.getBizId());
                    break;
                case REFUND_TIMEOUT:
                    handleRefundTimeout(task.getBizId());
                    break;
            }
        }
    }
    /**
     * 处理支付超时
     */
    private void handlePayTimeout(String orderId) {
        Order order = orderService.getById(orderId);
        if ("待支付".equals(order.getStatus())) {
            orderService.cancelOrder(orderId, "支付超时");
            inventoryService.restoreStock(order.getSkuId(), order.getQuantity());
            log.info("订单[{}]支付超时,已自动取消", orderId);
        }
    }
}

技术选型建议

方案 优点 缺点 适用场景
Redis ZSET 实现简单、高性能 需自建轮询、数据可能丢失 小型系统、对实时性要求不高
RabbitMQ 死信 可靠、消息不丢失 需额外配置、延迟精度有限 金融、交易等关键业务
RocketMQ 内置延迟等级、高吞吐 依赖特定MQ 大规模分布式系统
Kafka 高吞吐、可靠 延迟精度不足 日志分析、异步批处理
Java DelayQueue 轻量、简单 单机、无法持久化 单体应用、内存任务

注意事项

  1. 消息可靠性:生产环境建议使用 MQ 方案,确保消息不丢失
  2. 延迟精度:Redis 方案受轮询频率影响,MQ 有最小延迟限制
  3. 幂等性:消费端要做幂等处理,防止重复消费
  4. 监控告警:对延迟消息队列进行监控,设置积压告警
  5. 优雅关闭:提供关闭钩子,处理内存中的延迟任务

选择方案时需要结合业务场景、系统规模、技术栈等因素综合评估。

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