这是《RocketMQ 零基础实战指南》的独立章节版。本章从概念、实操和生产排查三个视角展开,代码块保留了原书可直接运行的版本。 延迟消息用于“现在提交,未来某个时间触发”的场景,例如订单超时关闭、预约提醒、延迟重试和冷却期控制。RocketMQ 历史上提供固定级别延迟,5.x 引入任意时间延迟能力,具体实现和限制应以所使用版本文档为准。
9.1 典型场景
| 场景 | 延迟 |
|---|---|
| 未支付订单关闭 | 30 分钟 |
| 支付成功后发放券 | 5 分钟 |
| 失败任务重试 | 指数退避 |
| 会议提醒 | 固定时间 |
| 风控冷却 | 秒到小时 |
| 数据归档触发 | 天级 |
延迟消息比定时扫表更适合离散到期时间的业务,但到期精度通常受调度周期、集群负载和消费者处理能力影响。
9.2 固定级别延迟
RocketMQ 4.x 常见延迟级别:
1s 5s 10s 30s 1m 2m 3m 4m 5m 6m 7m 8m 9m 10m 20m 30m 1h 2h
Java 示例:
Message message = new Message("OrderTimeoutTopic", "OrderTimeout",
timeoutEvent(orderNo).getBytes(StandardCharsets.UTF_8));
message.setDelayTimeLevel(16); // 30m
producer.send(message);
如果业务需要 25 分钟,只能选择接近的级别,例如 30 分钟,或在消费端再判断实际到期时间。
9.3 5.x 定时消息
RocketMQ 5.x 支持按时间戳设置定时消息:
Message message = new Message("AppointmentTopic", "Remind",
remindEvent(appointmentId).getBytes(StandardCharsets.UTF_8));
message.setDeliverTimeMs(System.currentTimeMillis() + TimeUnit.MINUTES.toMillis(30));
producer.send(message);
注意事项:
- 使用 Broker 和客户端都支持的版本;
- 确认 Topic 的定时消息能力已开启;
- 了解最大可延迟时间;
- 了解精度和调度周期;
- 避免瞬时大量相同到期时间造成消费洪峰;
- 集群升级时验证存量定时消息兼容性。
9.4 到期处理
消费端不能假设消息一定准时:
long now = System.currentTimeMillis();
long deliverAt = msg.getDeliverTimeMs();
if (deliverAt > now) {
log.warn("timer message arrived early, waitMs={}", deliverAt - now);
return ConsumeConcurrentlyStatus.RECONSUME_LATER;
}
if (now - deliverAt > TimeUnit.MINUTES.toMillis(5)) {
metric.lateTimeout();
}
processTimeout(msg);
业务执行前应再次判断状态:
UPDATE orders
SET status = 'CLOSED', close_reason = 'TIMEOUT'
WHERE order_no = ? AND status = 'WAIT_PAY';
如果更新行数为 0,说明订单已支付或已关闭,应记录日志并确认消费。
9.5 幂等与补偿
延迟消息可能重复投递,也可能因迁移和重放再次出现:
- 以订单号或任务 ID 做幂等键;
- 状态机只允许合法转换;
- 记录任务处理结果;
- 定时对账未关闭订单;
- 人工处理死信;
- 对重复到期保留审计。
补偿任务示例:
SELECT order_no
FROM orders
WHERE status = 'WAIT_PAY'
AND created_at < NOW() - INTERVAL 35 MINUTE
LIMIT 1000;
9.6 设计模式
固定超时
create order -> send 30m delay -> close if WAIT_PAY
取消超时
下单 -> 延迟关单
支付 -> 状态变更为 PAID
到期 -> 关单条件失败,不执行
如果业务支持显式取消,需要保证取消事件和到期事件不会互相覆盖,通常以数据库状态机为最终裁决。
指数退避
retry delay = base * 2^attempt
delay = min(delay, maxDelay)
attempt = attempt + 1
每次消费时计算下一次延迟,并记录尝试次数。
9.7 容量与洪峰
大量任务在同一秒到期会产生集中消费:
10:00:00 100,000 reminders
优化方式:
- 创建时间随机化;
- 提醒时间加随机抖动;
- 拆分 Topic 和消费组;
- 控制消费线程;
- 下游限流排队;
- 预估队列数和实例数;
- 提前压测。
9.8 版本与迁移
迁移建议:
- 固定延迟与定时消息分 Topic 管理;
- 新旧生产者灰度切换;
- 验证 Broker 支持能力和容量;
- 对比消息轨迹和到期时间;
- 保留补偿扫描;
- 回滚方案保持旧路径可用。
不要只改客户端 API 就完成迁移,Broker 端能力和运维策略同样关键。
9.9 常见故障
| 现象 | 排查 |
|---|---|
| 消息不到期 | 版本不支持、延迟时间超出限制 |
| 全部延迟很久 | 定时调度积压、消费积压 |
| 提前消费 | 版本兼容、时间解析错误 |
| 大量重复 | 重试、重放或补偿并发 |
| 到期洪峰 | 相同到期时间集中 |
| 状态错误 | 未做状态机条件更新 |
本章小结
延迟消息把未来的触发点交给消息系统,减少扫表和轮询。4.x 以固定级别为主,5.x 提供更强的定时能力,但都必须确认版本限制、到期精度、容量和兼容性。业务侧仍要依靠状态机、幂等和补偿对账保证最终正确。
思考题
- 延迟消息为什么不能保证精确到毫秒触发?
- 固定级别延迟如何近似 25 分钟?
- 到期后为什么还要检查订单状态?
- 相同到期时间为什么会造成洪峰?
- 延迟消息迁移时为什么要验证 Broker 能力?