Administrator
发布于 2021-08-07 / 1011 阅读
17

延迟消息的实现方案对比

需求:30 分钟未支付自动关单

产品提了个很常见的需求:用户下单后 30 分钟没支付,自动关闭订单并释放库存。我们日均订单 12 万,粗略统计落在 30 分钟窗口内未支付的约 3.4 万单。

这个需求本质是"延迟任务",实现方式有好几种,我基本都试过,把各自的坑记一下。

方案一:定时扫表

最朴素的写法,起个定时任务每分钟跑一次:

@Scheduled(fixedDelay = 60_000)
public void closeTimeoutOrders() {
    LocalDateTime deadline = LocalDateTime.now().minusMinutes(30);
    List<Long> ids = orderMapper.selectTimeoutUnpaid(deadline, 1000);
    for (Long id : ids) {
        orderService.closeOrder(id);
    }
}
<select id="selectTimeoutUnpaid" resultType="java.lang.Long">
    SELECT id FROM t_order
     WHERE order_status = 0
       AND created_at < #{deadline}
     LIMIT #{limit}
</select>

优点是非常好理解,出问题直接查表。缺点也很明确:

  • 有延迟。平均延迟是扫描间隔的一半,间隔 1 分钟的话平均 30 秒,最坏 60 秒。设为 10 秒扫一次的话,数据库压力又上来了。
  • 扫描本身有成本。我们 t_order 有 4200 万行,虽然 (order_status, created_at) 有联合索引,但每分钟扫一次,高峰时单次 12 ms,一天 1440 次,累计占用不算什么,但总感觉是无效功。
  • 分库分表之后会很难受。我们订单分了 8 个库,这个定时任务要扫 8 个分片再合并,代码复杂度直接翻倍。

小项目(日订单几千)用这个完全够,别过度设计。我们没选,是因为量摆在这儿。

方案二:RocketMQ 延迟消息

我们本来就在用 RocketMQ 4.9.1,直接用它的延迟消息:

// 生产者:下单成功后发一条延迟消息
Message msg = new Message("ORDER_TIMEOUT_TOPIC",
        "CLOSE",
        orderId.toString().getBytes(StandardCharsets.UTF_8));
// 延迟级别 16 对应 30 分钟
msg.setDelayTimeLevel(16);
SendResult result = producer.send(msg);

消费者这边正常消费:

@RocketMQMessageListener(topic = "ORDER_TIMEOUT_TOPIC",
                         consumerGroup = "order-timeout-consumer")
public class OrderTimeoutListener implements RocketMQListener<String> {

    @Override
    public void onMessage(String orderId) {
        Order order = orderMapper.selectById(Long.valueOf(orderId));
        if (order == null || order.getOrderStatus() != 0) {
            return;                       // 已支付或已关闭,忽略
        }
        orderService.closeOrder(order.getId());
    }
}

这里有个大坑必须先说:RocketMQ 4.x 的延迟消息不支持任意时间,只支持 18 个固定级别。定义在 broker 配置里:

# broker.conf,默认值
messageDelayLevel=1s 5s 10s 30s 1m 2m 3m 4m 5m 6m 7m 8m 9m 10m 20m 30m 1h 2h

级别从 1 开始数,所以 setDelayTimeLevel(16) 是 30 分钟。产品要是说"改成 25 分钟",你得改 broker 配置加一个级别,而且改完要重启 broker 才生效,已经发出去的消息不受影响。这个限制在我们这没造成问题,因为 30 分钟刚好有现成的级别。

原理也可以简单说一下,理解了才能知道它的边界:broker 收到延迟消息后,不会立刻写进真正的 topic,而是改写 topic 为 SCHEDULE_TOPIC_XXXX、把原来的 topic 和 queueId 塞进消息属性,然后按延迟级别投递到对应队列。每个级别有一个定时任务(DeliverDelayedMessageTimerTask),每秒扫一次,把到期的消息重新写回原 topic。所以延迟是有误差的,我们实测平均偏差 1.2 秒,最大 3.8 秒,对关单这种场景完全够用。

方案三:Redis ZSet

有些场景需要任意延迟时间(比如"用户设置 2 小时 15 分钟后提醒我"),RocketMQ 的固定级别就顶不住了。这时候用 Redis ZSet 最灵活:score 存执行时间戳,member 存业务 ID。

// 下单时写入,score = 到期时间戳
public void scheduleClose(Long orderId, long delaySeconds) {
    long score = System.currentTimeMillis() + delaySeconds * 1000;
    stringRedisTemplate.opsForZSet().add("delay:order:close",
            orderId.toString(), score);
}

// 定时扫描,每 5 秒拉一次到期任务
@Scheduled(fixedDelay = 5000)
public void consumeDelayTasks() {
    long now = System.currentTimeMillis();
    Set<String> ids = stringRedisTemplate.opsForZSet()
            .rangeByScore("delay:order:close", 0, now, 0, 500);
    if (CollectionUtils.isEmpty(ids)) {
        return;
    }
    for (String id : ids) {
        // 用 Lua 保证原子性:只有 zrem 成功的那个实例才执行任务
        Long removed = stringRedisTemplate.execute(REMOVE_SCRIPT,
                Collections.singletonList("delay:order:close"),
                id);
        if (removed != null && removed == 1) {
            orderService.closeOrder(Long.valueOf(id));
        }
    }
}

关键点在那个 Lua 脚本。多实例部署时,两个实例可能同时 zrangeByScore 拉到同一批 ID,不加保护就会重复执行。用 ZREM 的返回值做抢占(返回 1 表示是自己删掉的):

private static final RedisScript<Long> REMOVE_SCRIPT = RedisScript.of(
    "if redis.call('zrem', KEYS[1], ARGV[1]) == 1 then return 1 else return 0 end",
    Long.class);

这个方案的精度取决于扫描间隔,我们设 5 秒,实测平均延迟 2.6 秒。缺点是Redis 挂了任务就没了,所以我们用它做兜底而不是主力。

方案四:时间轮

上面两种都是"借助外部组件"的方案。如果要在进程内自己实现高精度定时,标准答案就是时间轮。

Netty 的 HashedWheelTimer 是最常见的实现,我写过一个小工具用它管理心跳超时:

HashedWheelTimer timer = new HashedWheelTimer(
        new DefaultThreadFactory("delay-task"),
        100, TimeUnit.MILLISECONDS,     // 每格 100 ms
        512);                            // 512 格,一圈 51.2 秒

timer.newTimeout(timeout -> {
    // 到期执行
    closeOrder(orderId);
}, 30, TimeUnit.MINUTES);

时间轮的思路是:把时间轴分成固定数量的格子(tick),每个格子上挂一个待执行任务链表。指针每 tick 走一格,处理当前格子里的任务。任务延迟超过一圈的,记一个 round 数,指针转够圈数再执行。

它的优势是插入和到期检测都是 O(1),不像 ScheduledThreadPoolExecutor 那样靠优先级队列(插入 O(log n))。管理几万个连接的心跳超时,时间轮是唯一合理的选择。

RocketMQ 内部的延迟消息实现(ScheduleMessageService)用的就是类似结构:每个延迟级别一个 DeliverDelayedMessageTimerTask,本质是个按秒推进的定时器 + 每个级别维护一个 offset 队列。

缺点是数据全在内存,进程重启任务全丢。所以时间轮适合"丢了也无所谓"的场景(心跳检测、连接超时、请求超时),不适合"必须执行"的业务任务。

我们最后的选择

方案精度可靠性任意延迟适用场景
定时扫表分钟级支持日订单万级以下,最简单
RocketMQ 延迟消息秒级(实测偏差 1-4 秒)高(消息持久化)不支持,18 个固定级别延迟时间固定的业务
Redis ZSet取决于扫描间隔中(依赖 Redis 持久化)支持任意延迟 + 允许少量丢失
时间轮毫秒级低(内存)支持心跳、超时控制等可容忍丢失的场景

主力用 RocketMQ 延迟消息(30 分钟正好有现成级别),同时保留一个每 10 分钟跑一次的对账任务扫表兜底,处理消息丢失或者消费失败的情况。上线三个月,主链路触发关单 287 万次,兜底任务实际补偿了 412 单,占比 0.014%。

小结

  • 延迟时间固定、量又大,优先 RocketMQ 延迟消息,记得 setDelayTimeLevel 是从 1 开始数,级别对应的是 broker 配置里的 messageDelayLevel
  • RocketMQ 4.x 改延迟级别要重启 broker,上线前和产品确认延迟时间不会频繁变。
  • 需要任意延迟就用 Redis ZSet,多实例必须用 ZREM 返回值抢占,否则重复执行。
  • 消费者一定要做幂等:关单前先查订单状态,延迟消息可能重复投递(RocketMQ 只保证至少一次)。
  • 时间轮适合内存级、可丢的定时任务,别拿来存业务数据。
  • 不管用哪个方案,都留一个扫表兜底。消息系统不是 100% 可靠的。

参考