RocketMQNotes

第 09 章:延迟消息

zjc 于 2026-01-09 发布

这是《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);

注意事项:

  1. 使用 Broker 和客户端都支持的版本;
  2. 确认 Topic 的定时消息能力已开启;
  3. 了解最大可延迟时间;
  4. 了解精度和调度周期;
  5. 避免瞬时大量相同到期时间造成消费洪峰;
  6. 集群升级时验证存量定时消息兼容性。

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 幂等与补偿

延迟消息可能重复投递,也可能因迁移和重放再次出现:

  1. 以订单号或任务 ID 做幂等键;
  2. 状态机只允许合法转换;
  3. 记录任务处理结果;
  4. 定时对账未关闭订单;
  5. 人工处理死信;
  6. 对重复到期保留审计。

补偿任务示例:

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

优化方式:

  1. 创建时间随机化;
  2. 提醒时间加随机抖动;
  3. 拆分 Topic 和消费组;
  4. 控制消费线程;
  5. 下游限流排队;
  6. 预估队列数和实例数;
  7. 提前压测。

9.8 版本与迁移

迁移建议:

  1. 固定延迟与定时消息分 Topic 管理;
  2. 新旧生产者灰度切换;
  3. 验证 Broker 支持能力和容量;
  4. 对比消息轨迹和到期时间;
  5. 保留补偿扫描;
  6. 回滚方案保持旧路径可用。

不要只改客户端 API 就完成迁移,Broker 端能力和运维策略同样关键。

9.9 常见故障

现象 排查
消息不到期 版本不支持、延迟时间超出限制
全部延迟很久 定时调度积压、消费积压
提前消费 版本兼容、时间解析错误
大量重复 重试、重放或补偿并发
到期洪峰 相同到期时间集中
状态错误 未做状态机条件更新

本章小结

延迟消息把未来的触发点交给消息系统,减少扫表和轮询。4.x 以固定级别为主,5.x 提供更强的定时能力,但都必须确认版本限制、到期精度、容量和兼容性。业务侧仍要依靠状态机、幂等和补偿对账保证最终正确。

思考题

  1. 延迟消息为什么不能保证精确到毫秒触发?
  2. 固定级别延迟如何近似 25 分钟?
  3. 到期后为什么还要检查订单状态?
  4. 相同到期时间为什么会造成洪峰?
  5. 延迟消息迁移时为什么要验证 Broker 能力?