首页
学习
活动
专区
圈层
工具
发布

RabbitMQ 怎么实现延迟队列?

订单创建成功,半小时没支付,系统自动关闭。

这种需求看着简单,代码里真要出现下面这行,我一般会直接打回去:

Thread.sleep(30 * 60 * 1000L);

线程不是闹钟。让业务线程睡半小时,不仅浪费线程资源,服务一重启,睡着的任务也没了。

定时任务轮询数据库能做,但数据量上来以后,类似这样的 SQL 会一直扫:

select order_no

from trade_order

where pay_status = 'WAITING'

and expire_at <= now()

limit 500;

索引没设计好,扫描会越来越重;定时任务跑慢了,订单关闭又会延后。RabbitMQ 的延迟队列更适合处理这种“消息现在发,过一段时间再消费”的场景。

不过 RabbitMQ 默认没有一个名字就叫“延迟队列”的队列。常见实现有两套:TTL 加死信队列,或者使用延迟消息插件。

TTL 加死信队列

这个方案利用了 RabbitMQ 的两个能力。

消息进入等待队列后先不消费,等 TTL 到期,RabbitMQ 把它当成死信,转发到真正的业务队列。消费者监听业务队列,执行订单关闭。

等待队列配置我一般这样写:

@Configuration

public class OrderDelayTopology {

  static final String WAIT_EXCHANGE = "trade.order.wait.ex";

  static final String WAIT_QUEUE = "trade.order.wait.q";

  static final String CLOSE_EXCHANGE = "trade.order.close.ex";

  static final String CLOSE_QUEUE = "trade.order.close.q";

  @Bean

  DirectExchange waitExchange() {

      return new DirectExchange(WAIT_EXCHANGE, true, false);

  }

  @Bean

  DirectExchange closeExchange() {

      return new DirectExchange(CLOSE_EXCHANGE, true, false);

  }

  @Bean

  Queue waitQueue() {

      return QueueBuilder.durable(WAIT_QUEUE)

              .deadLetterExchange(CLOSE_EXCHANGE)

              .deadLetterRoutingKey("order.close")

              .build();

  }

  @Bean

  Queue closeQueue() {

      return QueueBuilder.durable(CLOSE_QUEUE).build();

  }

  @Bean

  Binding waitBinding() {

      return BindingBuilder.bind(waitQueue())

              .to(waitExchange())

              .with("order.wait");

  }

  @Bean

  Binding closeBinding() {

      return BindingBuilder.bind(closeQueue())

              .to(closeExchange())

              .with("order.close");

  }

}

创建订单后,把关闭消息发进等待队列,并给消息设置过期时间:

@Service

public class OrderDelaySender {

  private final RabbitTemplate rabbitTemplate;

  public OrderDelaySender(RabbitTemplate rabbitTemplate) {

      this.rabbitTemplate = rabbitTemplate;

  }

  public void scheduleClose(String orderNo, Duration delay) {

      OrderCloseCommand command =

              new OrderCloseCommand(orderNo, System.currentTimeMillis());

      rabbitTemplate.convertAndSend(

              OrderDelayTopology.WAIT_EXCHANGE,

              "order.wait",

              command,

              message -> {

                  long delayMillis = Math.max(delay.toMillis(), 1000L);

                  message.getMessageProperties()

                          .setExpiration(Long.toString(delayMillis));

                  return message;

              });

  }

}

消息过期后会进入trade.order.close.q,消费者开始处理:

@Component

public class OrderCloseConsumer {

  private final TradeOrderService tradeOrderService;

  public OrderCloseConsumer(TradeOrderService tradeOrderService) {

      this.tradeOrderService = tradeOrderService;

  }

  @RabbitListener(queues = OrderDelayTopology.CLOSE_QUEUE)

  public void close(OrderCloseCommand command) {

      boolean changed = tradeOrderService.closeIfUnpaid(command.orderNo());

      if (!changed) {

          // 已支付、已关闭或者消息重复投递,不再继续处理

          return;

      }

      System.out.printf(

              "delayed order closed, orderNo=%s, createdAt=%d%n",

              command.orderNo(),

              command.createdAt()

      );

  }

}

record OrderCloseCommand(String orderNo, long createdAt) {

}

这里有个坑,挺多人配完能跑就不管了。

如果同一个队列里的消息 TTL 不一样,RabbitMQ 通常要等排在队头的消息过期后,才会继续处理后面的消息。比如第一条延迟一小时,第二条只延迟一分钟,第二条也可能被第一条挡住。

所以 TTL 加死信更适合延迟时间固定的场景。不同延迟时间,可以拆成多个等待队列,例如一分钟、十分钟、半小时各一个。别为了省三个队列,把消费时间搞得不可控。

延迟消息插件

延迟时间比较散,我更愿意使用x-delayed-message插件。消息不用先进入一个等待队列,交换机会暂存消息,到时间后再路由。

交换机配置如下:

@Bean

CustomExchange delayedExchange() {

  Map<String, Object> arguments = new HashMap<>();

  arguments.put("x-delayed-type", "direct");

  return new CustomExchange(

          "trade.order.delay.ex",

          "x-delayed-message",

          true,

          false,

          arguments

  );

}

@Bean

Binding delayedBinding() {

  return BindingBuilder.bind(closeQueue())

          .to(delayedExchange())

          .with("order.close")

          .noargs();

}

发送时把延迟毫秒数放进x-delay请求头:

public void scheduleByPlugin(String orderNo, Duration delay) {

  OrderCloseCommand command =

          new OrderCloseCommand(orderNo, System.currentTimeMillis());

  rabbitTemplate.convertAndSend(

          "trade.order.delay.ex",

          "order.close",

          command,

          message -> {

              message.getMessageProperties()

                      .setHeader("x-delay", Math.toIntExact(delay.toMillis()));

              return message;

          });

}

这个方案没有前面那个队头阻塞问题,代码也干净一些。但它依赖额外插件,RabbitMQ 集群中的节点都要正确安装和启用,升级前也得做兼容验证。这个事情不能只让开发本地跑通,然后把锅扔给运维。

还有两件事比延迟方案本身更重要。

第一,消费者必须幂等。RabbitMQ 可能重复投递消息,关闭订单时不能直接无条件更新,应该带上原状态:

update trade_order

set pay_status = 'CLOSED',

  closed_at = now()

where order_no = ?

and pay_status = 'WAITING';

受影响行数为零,就说明订单已经支付、关闭,或者消息被重复处理了。

第二,延迟队列不是精确到毫秒的定时器。Broker 繁忙、消息堆积、消费者处理变慢,都可能让实际消费时间晚一点。业务如果要求某个时间点一到就绝对执行,单靠 MQ 不够,还要有数据库补偿任务兜底。

固定延迟,TTL 加死信队列够用;延迟时间动态变化,优先考虑延迟消息插件。RabbitMQ 只负责把消息晚点交出来,订单到底该不该关闭,最后还是要重新查状态。

这一步不能省。

  • 发表于:
  • 原文链接https://page.om.qq.com/page/Oz7hCZ9laKAjRS-s3OaRzfSQ0
  • 腾讯「腾讯云开发者社区」是腾讯内容开放平台帐号(企鹅号)传播渠道之一,根据《腾讯内容开放平台服务协议》转载发布内容。
  • 如有侵权,请联系 cloudcommunity@tencent.com 删除。
领券