欢迎光临
我们一直在努力

RabbitMQ消息队列:延迟消息

一、方案一:死信交换机 + TTL(Time-To-Live)

RabbitMQ本身并没有直接提供延迟消息的功能,但我们可以巧妙地利用 死信交换机(Dead Letter Exchange) 和 消息TTL(过期时间) 来模拟实现。

1.1 什么是死信交换机?

当一个消息在一个队列中变为“死信”时,它会被重新投递到指定的死信交换机,再由它路由到最终的队列。

消息成为死信的三种情况:

  • 消费者拒绝消费:使用 basic.reject 或 basic.nack 声明消费失败,且 requeue 参数设为 false。

  • 消息过期:消息在队列中存活时间超过了设置的 TTL。

  • 队列达到最大长度:队列满了,无法再接纳新消息。

  • 当一个队列配置了 dead-letter-exchange 属性,那么发生上述情况的消息就会被转发到该交换机。

    1.2 利用死信交换机实现延迟消息的核心思想

  • 消息投递:生产者将消息发送到一个没有消费者的普通队列,并设置消息的 TTL(例如5秒)。

  • 消息过期:消息在队列中存活到 TTL 结束后,变为死信。

  • 死信转发:该队列配置了死信交换机,因此死信被转发到死信交换机。

  • 最终消费:死信交换机根据路由规则将消息投递到最终的业务队列,由消费者处理。

  • 此时,从消息发送到消费者收到,刚好经历了5秒的延迟。

    1.3 方案总结

    优点:

    • 实现简单:利用RabbitMQ原生机制,无需安装额外插件。

    • 稳定性高:基于核心功能,可靠性强。

    缺点:

    • 配置繁琐:需要为每个延迟任务配置死信交换机和队列。

    • 时间精度不高:RabbitMQ的TTL是追溯检查的,只有当过期消息位于队首时才会被处理。如果队列前有其他消息积压,即便消息已过期,也无法及时被处理,导致延迟时间不准确。

    注意:由于“队首阻塞”问题,该方案不适合对延迟时间精度要求极高的场景。


    二、方案二:DelayExchange 插件(官方推荐)

    鉴于方案一的局限性,RabbitMQ官方推出了 延迟消息插件(rabbitmq-delayed-message-exchange),提供了更优雅、更精准的延迟消息实现。

    2.1 声明延迟交换机

    我们可以声明一种新型交换机,其 delayed 属性为 true。

    • 基于注解方式:

    @RabbitListener(bindings = @QueueBinding(
    value = @Queue(name = "delay.queue", durable = "true"),
    exchange = @Exchange(name = "delay.direct", delayed = "true"),
    key = "delay"
    ))
    public void listenDelayMessage(String msg){
    log.info("接收到delay.queue的延迟消息:{}", msg);
    }

    • 基于 @Bean 方式:

    @Bean
    public DirectExchange delayExchange(){
    return ExchangeBuilder
    .directExchange("delay.direct")
    .delayed() // 关键:开启延迟特性
    .durable(true)
    .build();
    }

    2.2 发送延迟消息

    发送消息时,通过设置消息头 x-delay 来指定延迟的毫秒数。

    @Test
    void testPublisherDelayMessage() {
    String message = "hello, delayed message";
    rabbitTemplate.convertAndSend("delay.direct", "delay", message, new MessagePostProcessor() {
    @Override
    public Message postProcessMessage(Message message) throws AmqpException {
    // 设置5秒延迟
    message.getMessageProperties().setDelay(5000);
    return message;
    }
    });
    }

    2.3 方案总结

    优点:

    • 使用简单:只需声明交换机类型,并在发送时指定延迟时间。

    • 精度更高:插件内部通过Erlang定时器实现,比基于死信队列的方案更准时。

    缺点:

    • 依赖插件:需要额外安装。

    • 性能开销:大量长延迟消息会占用插件内部数据库表和定时器资源,增加CPU开销。因此,不建议设置过长时间的延迟。


    三、实战:订单支付状态同步

    接下来,我们将基于 DelayExchange 插件 的方案,在“交易服务”中实现一个高可用的订单支付状态同步功能。

    3.1 业务场景优化思路

    30分钟的延迟消息在MQ中等待,资源消耗较大。更优方案是采用 “梯度延迟检测” 策略:
    在下单后的 10秒、30秒、1分钟、2分钟、5分钟……30分钟 等多个时间点设置延迟消息。一旦在某个时间点检测到订单已支付,后续的检测任务自然取消,从而减少无效的MQ资源占用。

    3.2 核心步骤

    3.2.1 定义延迟消息体

    为了支持“多级延迟”,我们定义一个 MultiDelayMessage 类,其中包含业务数据和一个 List<Long> 类型的延迟时间集合(单位:毫秒)。

    @Data
    public class MultiDelayMessage<T> {
    private T data;
    private List<Long> delayMillis;

    // 获取并移除第一个延迟时间,实现“消费一个,取一个”的效果
    public Long removeNextDelay(){
    return delayMillis.remove(0);
    }

    public boolean hasNextDelay(){
    return !delayMillis.isEmpty();
    }
    }

    3.2.2 服务改造与配置
  • 定义常量:明确交换机、队列、路由Key。

    public interface MqConstants {
    String DELAY_EXCHANGE = "trade.delay.topic";
    String DELAY_ORDER_QUEUE = "trade.order.delay.queue";
    String DELAY_ORDER_ROUTING_KEY = "order.query";
    }

  • 引入依赖:在交易服务中引入 Spring AMQP 依赖。

  • 共享MQ配置:将 RabbitMQ 的连接信息抽取到 Nacos 配置中心,方便统一管理。

  • 3.2.3 改造下单业务

    在用户下单成功后,立即发送第一条延迟消息(例如10秒后)。

    // 创建订单后…
    // 发送延迟消息,检查支付状态
    // 延迟时间数组:10秒、30秒、1分钟…
    MultiDelayMessage<Long> msg = MultiDelayMessage.of(orderId, 10000L, 30000L, 60000L, …);
    rabbitTemplate.convertAndSend(MqConstants.DELAY_EXCHANGE, MqConstants.DELAY_ORDER_ROUTING_KEY, msg);

    3.2.4 编写支付状态查询接口

    在 pay-service 中提供根据业务订单号查询支付状态的接口,并在 hm-api 模块中声明对应的 FeignClient,供交易服务远程调用。

    3.2.5 核心:监听器处理逻辑

    消息监听器是整个流程的大脑,其处理逻辑如下:

  • 消费消息:从 delay.queue 获取包含订单ID的延迟消息。

  • 检查本地订单状态:若订单已支付或已关闭,直接结束。

  • 查询支付服务:若本地订单仍为“未支付”,则远程调用支付服务,查询最新状态。

  • 状态判断:

    • 已支付:更新本地订单状态为“已支付”,流程结束。

    • 未支付:判断 MultiDelayMessage 中是否还有剩余延迟时间。

      • 有:取出下一个延迟时间,重新发送延迟消息。

      • 无:说明已超过最大等待时间(如30分钟),执行业务取消订单、恢复库存。

  • java

    @RabbitListener(bindings = @QueueBinding(…))
    public void listenOrderCheckDelayMessage(MultiDelayMessage<Long> msg) {
    // 1. 获取订单ID
    // 2. 本地订单状态检查
    // 3. 远程查询支付状态
    // 4. 支付成功,更新订单
    // 5. 未支付,判断是否继续延迟检测
    if (msg.hasNextDelay()) {
    int delayVal = msg.removeNextDelay().intValue();
    // 重新发送延迟消息,x-delay = delayVal
    } else {
    // 6. 超时未支付,取消订单
    orderService.cancelOrder(orderId);
    }
    }


    四、总结

    本文详细介绍了RabbitMQ实现延迟消息的两种主流方案,并深入讲解了其在电商订单超时处理场景下的实战应用。

    方案实现方式优点缺点适用场景
    死信交换机 + TTL 利用消息过期和死信转发机制 无需额外插件,基于核心功能 配置复杂,延迟时间可能不精确 对时间精度要求不高,且不想引入插件的场景
    DelayExchange 插件 使用官方插件,设置 x-delay 属性 使用简单,延迟精度高 需要安装插件,大量长延迟消息有性能开销 对时间精度有要求,且延迟时间不宜过长的场景。生产环境更推荐

    关键点回顾:

  • 延迟消息是解决分布式系统中定时任务的一种优雅方案。

  • “梯度延迟检测”策略能有效降低MQ资源消耗,是优化延迟任务的重要手段。

  • 结合 Feign 远程调用 与 RabbitMQ,可以实现服务间的松耦合和高效协作。

  • 赞(0)
    未经允许不得转载:171主机测评 » RabbitMQ消息队列:延迟消息
    分享到: 更多 (0)

    评论 抢沙发

    • 昵称 (必填)
    • 邮箱 (必填)
    • 网址