RocketMQ 延迟消息:30 分钟未支付自动关单怎么做
作者:鱼宵 | 实战驱动系列 · 第 7 篇
完整课程与可运行源码已开源在 Gitee:https://gitee.com/j67mk2/rocketmq-journey (本文对应 lesson-07/)
一、先想一个谁都躲不开的问题
用户在你的商城拍下一单,30 分钟内不付款,库存一直被占着。这种情况怎么办?
总不能在应用里起个线程Thread.sleep(30, TimeUnit.MINUTES)干等吧——一秒钟进来几百个订单,就是几百个干等的线程,机器当场就扛不住了。
也不能每分钟扫一遍订单表:库被扫得嗷嗷叫,还容易漏、容易重复关单。
这种"下单 30 分钟后,提醒我看看他到底付没付款"的需求,在 RocketMQ 里有个现成答案:延迟消息。一行代码给消息加上"晚 10 秒再投递"的属性,Broker(存消息、管投递的那个服务器)替你记着时间,到点才把消息推给消费者。
今天就把它跑通:发一条延迟 10 秒的消息,亲眼算出"从发送到消费正好约 10 秒",再把它套到"30 分钟未支付自动关单"这个经典场景上。代码就在仓库lesson-07/里,clone 下来按第四节跑一遍,现象自己就跳出来了。
二、18 个固定级别:为什么不能"延迟任意 37 秒"
先说一个最容易踩的坑:RocketMQ 4.x不支持延迟任意秒数,只提供 18 个写死的级别。你不是说"延迟 37 秒",而是选一个级别号,对应一个固定延迟:
| 级别 | 延迟 | 级别 | 延迟 |
|---|---|---|---|
| 1 | 1 秒 | 10 | 6 分钟 |
| 2 | 5 秒 | 11 | 7 分钟 |
| 3 | 10 秒 | 12 | 8 分钟 |
| 4 | 30 秒 | 13 | 9 分钟 |
| 5 | 1 分钟 | 14 | 10 分钟 |
| 6 | 2 分钟 | 15 | 20 分钟 |
| 7 | 3 分钟 | 16 | 30 分钟 |
| 8 | 4 分钟 | 17 | 1 小时 |
| 9 | 5 分钟 | 18 | 2 小时 |
级别号和秒数不是一一对应的:级别 3 不是 3 秒,是 10 秒;级别 16 才是 30 分钟。要用哪一档,先查表,别凭感觉数。
为什么不能延迟任意秒?因为服务端是按级别预先建好 18 个定时队列的(内部名字叫SCHEDULE_TOPIC_XXXX),每个级别对应一个队列。到点时,Broker 后台每秒扫一遍这些队列,把"到点了"的消息批量搬出来投递。
如果允许"延迟 37 秒""延迟 2 小时 15 分"这种任意值,每条消息都得单独记一个精确时间点、单独排序、单独定时,成本高得多。用固定级别,换来的是"批量定时扫描"的高性能。想要任意时刻?RocketMQ 5.x 才支持;本系列用的是 4.9.7,开发时选最接近的一档就行。
三、它不是"晚发",是先睡在 Broker 的队列里
很多人第一反应:"延迟消息是不是客户端压着不发?"不是。整个过程是这样的:
Producer(生产者) Broker(服务器) │ 发消息 setDelayTimeLevel(3) │ │ ───────────────────────────▶ │ ① 不投递给真实 Topic, │ │ 先写进内部定时队列 │ │ 顺便算好"投递时刻 = 发送时刻 + 10 秒" │ │ │ │ ② 后台定时任务每秒扫描: │ │ 哪些消息到点了? │ │ │ │ ③ 到点的消息搬出来,投到真实 Topic │ │ ──────────────────────────▶ Consumer(消费者)才收到一句话记住:延迟消息是先睡在 Broker 的定时队列里,到点才被搬出来投递。延迟这件事由服务端定时扫描实现,不依赖任何客户端;它和普通消息一样落盘,Broker 重启也不丢。
四、动手:两个类,跑通延迟 10 秒
环境照第 1 课:Docker 起好 RocketMQ 三件套,控制台能打开即可。工程在lesson-07/,核心就两个主类,代码可以直接复制。
先看生产者。唯一的主角是setDelayTimeLevel(3)那一行:
packagecom.example.lesson07;importorg.apache.rocketmq.client.producer.DefaultMQProducer;importorg.apache.rocketmq.client.producer.SendResult;importorg.apache.rocketmq.common.message.Message;importjava.nio.charset.StandardCharsets;importjava.text.SimpleDateFormat;importjava.util.Date;/** * 第 7 课 Demo:延迟消息生产者 * 往 TopicLesson07 发 1 条"延迟 10 秒"的消息(delayTimeLevel=3)。 */publicclassDelayProducer{publicstaticvoidmain(String[]args)throwsException{SimpleDateFormatsdf=newSimpleDateFormat("HH:mm:ss.SSS");// 1. 用普通生产者即可,延迟消息不需要专用类DefaultMQProducerproducer=newDefaultMQProducer("lesson07_producer_group");producer.setNamesrvAddr("127.0.0.1:9876");// 连 Namesrv(服务发现的"登记簿")producer.start();System.out.println(sdf.format(newDate())+" 生产者已启动。");// 2. 构造消息:TopicLesson07(主题,消息的类别)+ TagDelay + 消息体Messagemsg=newMessage("TopicLesson07","TagDelay","订单 12345:拍下已 10 秒还没支付,系统来提醒你赶紧付款".getBytes(StandardCharsets.UTF_8));// 3. 设置延迟级别(本课主角):3 = 延迟 10 秒// 级别号和秒数不是一一对应的!3 不是 3 秒,是 10 秒msg.setDelayTimeLevel(3);// 4. 发送:send() 立刻返回成功,不代表消息立刻被消费!// 此刻消息正躺在 Broker 的定时队列里睡觉System.out.println(sdf.format(newDate())+" 发送延迟消息(级别=3,即延迟 10 秒)……");SendResultresult=producer.send(msg);System.out.println(sdf.format(newDate())+" 发送成功:"+result);producer.shutdown();System.out.println(sdf.format(newDate())+" 生产者已关闭。");}}再看消费者。它干的事很简单:收到消息后,拿消息自带的"出生时间"和现在对比,算一算这条消息到底睡了多久:
packagecom.example.lesson07;importorg.apache.rocketmq.client.consumer.DefaultMQPushConsumer;importorg.apache.rocketmq.client.consumer.listener.ConsumeConcurrentlyContext;importorg.apache.rocketmq.client.consumer.listener.ConsumeConcurrentlyStatus;importorg.apache.rocketmq.client.consumer.listener.MessageListenerConcurrently;importorg.apache.rocketmq.common.message.MessageExt;importjava.nio.charset.StandardCharsets;importjava.text.SimpleDateFormat;importjava.util.Date;importjava.util.List;importjava.util.concurrent.atomic.AtomicInteger;/** * 第 7 课 Demo:延迟消息消费者 * 收到消息后用 bornTimestamp(生产者发送时刻)算一算"这条消息睡了多久"。 */publicclassDelayConsumer{publicstaticvoidmain(String[]args)throwsException{SimpleDateFormatsdf=newSimpleDateFormat("HH:mm:ss.SSS");// 1. 创建 Push 消费者,消费组 lesson07_consumer_groupDefaultMQPushConsumerconsumer=newDefaultMQPushConsumer("lesson07_consumer_group");consumer.setNamesrvAddr("127.0.0.1:9876");// 2. 订阅 TopicLesson07 的全部消息(* 表示不按 Tag 过滤)consumer.subscribe("TopicLesson07","*");// 3. 注册监听器:收到消息后算时间差AtomicIntegerreceivedCount=newAtomicInteger(0);consumer.registerMessageListener(newMessageListenerConcurrently(){@OverridepublicConsumeConcurrentlyStatusconsumeMessage(List<MessageExt>msgs,ConsumeConcurrentlyContextcontext){for(MessageExtmsg:msgs){Stringbody=newString(msg.getBody(),StandardCharsets.UTF_8);longnow=System.currentTimeMillis();// bornTimestamp = 生产者发送这条消息那一刻打的时间戳doubleelapsedSeconds=(now-msg.getBornTimestamp())/1000.0;System.out.println(sdf.format(newDate(now))+" 收到消息: "+body);System.out.println(" 消息出生于 "+sdf.format(newDate(msg.getBornTimestamp()))+",从发送到现在经过了 "+String.format("%.1f",elapsedSeconds)+" 秒(≈延迟级别 3 的 10 秒,另有零点几秒的网络/拉取耗时)");receivedCount.incrementAndGet();}returnConsumeConcurrentlyStatus.CONSUME_SUCCESS;}});// 4. 启动消费者。新消费组默认只读"启动之后投递进 Topic 的新消息",// 所以一定要先启动消费者、再去另一个窗口发延迟消息consumer.start();System.out.println("消费者已启动,正在等待延迟消息……(预计 10 秒后收到;收到 1 条后自动退出,最多等 60 秒)");// 5. 演示用退出逻辑:收到 1 条,或 60 秒超时// 超时要留够:mvn 启动、等消费者就绪、再加 10 秒延迟,30 秒容易把消息饿死在窗外longdeadline=System.currentTimeMillis()+60_000;while(receivedCount.get()<1&&System.currentTimeMillis()<deadline){Thread.sleep(500);}consumer.shutdown();System.out.println(sdf.format(newDate())+" 观察结束,消费者已关闭(共收到 "+receivedCount.get()+" 条)。");}}运行顺序有讲究:先起消费者,再起生产者。因为新消费组默认只收"启动之后才投递进 Topic 的新消息",而延迟消息要 10 秒后才进真实 Topic——你要是先发消息再起消费者,它正好错过,干等 60 秒超时。
两个 PowerShell 窗口的启动命令(先设好编码和 JDK,路径换成你本机的工程位置):
# 窗口一:先起消费者[Console]::OutputEncoding=[System.Text.Encoding]::UTF8$env:JAVA_TOOL_OPTIONS="-Dfile.encoding=UTF-8"$env:JAVA_HOME="C:\Program Files\Java\jdk-17"cd rocketmq-journey\lesson-07 mvn exec:java"-Dexec.mainClass=com.example.lesson07.DelayConsumer"# 窗口二:看到消费者打印"正在等待延迟消息"后,再起生产者[Console]::OutputEncoding=[System.Text.Encoding]::UTF8$env:JAVA_TOOL_OPTIONS="-Dfile.encoding=UTF-8"$env:JAVA_HOME="C:\Program Files\Java\jdk-17"cd rocketmq-journey\lesson-07 mvn exec:java"-Dexec.mainClass=com.example.lesson07.DelayProducer"五、实测输出:发送到消费,正好 10.0 秒
下面是本教程同款环境(Windows 11 + Docker 29 + JDK 17)的真实运行输出,你跑出来应该长这样:
生产者输出:
18:33:07.642 生产者已启动。 18:33:07.643 发送延迟消息(级别=3,即延迟 10 秒)…… 18:33:07.852 发送成功:SendResult [sendStatus=SEND_OK, msgId=7F000001..., messageQueue=MessageQueue [topic=TopicLesson07, brokerName=broker-a, queueId=2], queueOffset=16] 18:33:07.864 生产者已关闭。消费者输出:
消费者已启动,正在等待延迟消息……(预计 10 秒后收到;收到 1 条后自动退出,最多等 60 秒) 18:33:17.860 收到消息: 订单 12345:拍下已 10 秒还没支付,系统来提醒你赶紧付款 消息出生于 18:33:07.837,从发送到现在经过了 10.0 秒(≈延迟级别 3 的 10 秒,另有零点几秒的网络/拉取耗时) 18:33:18.143 观察结束,消费者已关闭(共收到 1 条)。盯着两个时间点看:18:33:07.84 生产者发完立刻SEND_OK;18:33:17.86 消费者才收到——正好 10.0 秒。要是普通消息不设延迟,消费者一般一两百毫秒内就收到了。这中间差的 10 秒,就是延迟消息在 Broker 定时队列里"睡"出来的。
六、套回真实业务:30 分钟未支付自动关单
Demo 跑通了,回到开头那个问题。落地其实就三步:
- 用户下单成功 → 发一条延迟消息,级别选 16(30 分钟),消息体里带上订单号;
- 消费者 30 分钟后收到消息 → 拿订单号去查订单最新状态;
- 已支付 → 什么都不做;未支付 → 关单、释放库存。
这里有个面试常追问的细节:用户要是第 20 秒就付款了怎么办?
答案是:那条 30 分钟的关单消息照样会到消费者手里,挡不住。所以消费者收到消息时不能直接关单,必须先查一次最新订单状态——已支付就当没这条消息,未支付才执行关单。延迟消息不会因为用户支付了就凭空消失,"到点查状态、再决定动不动手"这个思路,才是这套方案的关键。
七、别用错地方:延迟消息 vs xxl-job vs 时间轮
| MQ 延迟消息 | xxl-job 定时任务 | 时间轮 | |
|---|---|---|---|
| 定位 | 消息的副产品:“这条消息晚 N 分钟再投” | 通用任务调度平台:“每天凌晨跑个报表” | JVM 内存级定时器:心跳、连接超时 |
| 持久化 | 落盘,Broker 重启不丢 | 靠数据库/调度中心 | 纯内存,进程重启就没 |
| 适合 | 订单超时关单、延迟通知 | 报表同步、定时数据清理 | 框架内部的短周期计时 |
| 类比 | 快递预约送达 | 每天固定响的闹钟 | 秒表倒计时 |
一句话:MQ 延迟消息是"消息的一种属性",不是通用调度平台。"每天凌晨 2 点跑报表"这种活,别拿它硬凑。
八、挑战题(答案在仓库源码里,跑起来才知道)
- ⭐ 把生产者里的
setDelayTimeLevel(3)改成setDelayTimeLevel(2)(5 秒),重新跑一遍,消费者算出的时间差约等于几秒? - ⭐⭐ 改成
setDelayTimeLevel(5)(1 分钟)再跑一遍,体会长延迟——注意:消费者那个 60 秒等待窗口要不要调大?不调会发生什么? - ⭐⭐⭐ 把 Demo 改造成"下单后 30 秒未支付自动关单"的简化版:消息体带订单号,消费者收到后模拟"查订单状态",未支付就打印"自动关闭订单 XXX"。想一想:用户提前支付后,这条消息的消费逻辑该怎么写,才不会误关单?
九、面试回答模板
面试官:RocketMQ 延迟消息有几个级别?为什么不能延迟任意秒?
4.x 是 18 个固定级别:1s/5s/10s/30s/1m/2m/3m/4m/5m/6m/7m/8m/9m/10m/20m/30m/1h/2h。因为服务端按级别预建了 18 个定时队列,到点批量扫描投递;任意秒要每条消息单独定时,成本高。5.x 才支持任意时刻的定时消息。(见本文第二节)
追问:实现原理是什么?
延迟消息不会立刻投到真实 Topic,而是先写进内部定时队列(SCHEDULE_TOPIC_XXXX),记录投递时刻;Broker 后台每秒扫描,到点才把消息搬出来投递。延迟由服务端定时扫描实现,消息落盘,重启不丢。(见本文第三节)
追问:典型场景?和 xxl-job 有什么区别?
订单超时未支付自动关单、超时未确认自动收货这类。MQ 延迟消息是"消息的属性"、落盘持久化;xxl-job 是通用任务调度平台,适合跑报表这种定时任务。(见本文第六、七节)
关于这个系列
本文是「Java 后端实战精通营」系列第 7 篇,原则:实战驱动、由浅到深、面试向,每篇文章的结论都可以亲手验证。
👉RocketMQ 实战精通营(10 课):https://gitee.com/j67mk2/rocketmq-journey
- 本文对应源码位置:
lesson-07/(两个主类:DelayProducer 发延迟消息、DelayConsumer 算时间差)
系列文章一览(按发布顺序):
| 篇 | 主题 |
|---|---|
| 1 | RocketMQ 入门:Docker 一行起三件套,跑通你的第一条消息 |
| 2 | RocketMQ 发送方式:同步异步批量单向,消息都怎么发出去 |
| 3 | RocketMQ 消费模式:集群、广播与重试,消息怎么被吃掉 |
| 4 | RocketMQ 可靠性:发送重试加幂等,消息一条都不丢 |
| 5 | RocketMQ 顺序消息:订单流程不乱套的秘密 |
| 6 | RocketMQ 事务消息:订单与积分的最终一致 |
| 7 | RocketMQ 延迟消息:30 分钟未支付自动关单怎么做 |
| 8 | RocketMQ 积压治理:百万消息堵在队列怎么办 |
| 9 | RocketMQ 集群高可用与过滤:主从架构 + Tag 精准投递 |
| 10 | RocketMQ 面试冲刺:高频考点一口气背完 |
下一篇预告:《RocketMQ 积压治理:百万消息堵在队列怎么办》——消费者挂了一阵、流量突然暴增,消息在 Broker 里堆成山怎么发现、怎么治。
跑完有任何报错,把终端输出发评论区,一起排查。