Administrator
发布于 2020-07-14 / 3384 阅读
46

消息顺序性保证:全局有序与分区有序的取舍

同事问:订单状态怎么倒着走了

做微服务改造时,订单状态变更要通过 MQ 广播给下游(积分、优惠券、物流、通知)。上线一周后,物流组的同事找过来:"你们发的消息顺序不对,我这边先收到'已发货',后收到'已付款',状态机直接报错。"

我看了一眼他的日志,确实是反的:

14:03:21.114 收到订单消息 orderNo=SO20200714140319001 status=SHIPPED
14:03:21.119 状态机异常: WAIT_SHIP -> PAID 非法流转, 忽略
14:03:21.203 收到订单消息 orderNo=SO20200714140319001 status=PAID

两条消息间隔 89 毫秒,顺序颠倒了。原因不复杂:生产者是并发发的,Broker 上消息落在不同队列,消费者又是多线程拉的。

先说结论:全局有序基本做不到,也没必要

消息顺序性分两层:

  • 全局有序:整个 topic 的消息严格按发送顺序消费。实现方式只有一个——topic 只留 1 个队列,消费者只开 1 个线程。这等于主动放弃 MQ 的水平扩展能力,吞吐量退化到单机串行。
  • 分区有序(局部有序):同一组内有序,组间可以并行。比如同一个订单的消息有序,不同订单之间无所谓。这个才是生产上真正用的。

我们压测过:RocketMQ 4.7 单机,8 个队列、20 消费线程,能达到 4.2 万 TPS;改成 1 队列 1 线程后,TPS 掉到 2300,不到原来的 6%。有序性是用吞吐量换的,这个账得算清楚。

RocketMQ 的顺序消息怎么用

RocketMQ 的顺序消息靠两件事配合:发送端用 MessageQueueSelector 把同一组消息路由到同一个队列消费端用 MessageListenerOrderly 单线程消费这个队列

发送端,按订单号取模选队列:

rocketMQTemplate.asyncSendOrderly(
    "ORDER_STATUS_TOPIC",
    MessageBuilder.withPayload(event).build(),
    orderNo,                                  // hashKey,RocketMQ 用它算队列
    new SendCallback() {
        @Override public void onSuccess(SendResult r) { ... }
        @Override public void onException(Throwable e) { log.error("发送失败", e); }
    });

底层是 SelectMessageQueueByHash,对 hashKey 取哈希再对队列数取模:

// org.apache.rocketmq.client.producer.selector.SelectMessageQueueByHash
public MessageQueue select(List<MessageQueue> mqs, Message msg, Object arg) {
    int value = arg.hashCode();
    if (value < 0) value = Math.abs(value);
    value = value % mqs.size();
    return mqs.get(value);
}

消费端要用 MessageListenerOrderly,不能用并发的那个:

@RocketMQMessageListener(
        topic = "ORDER_STATUS_TOPIC",
        consumerGroup = "logistics-consumer-group",
        consumeMode = ConsumeMode.ORDERLY,          // 关键:顺序消费
        consumeThreadMax = 4)
public class OrderStatusConsumer implements RocketMQListener<MessageExt> {
    @Override
    public void onMessage(MessageExt msg) { ... }
}

ORDERLY 模式下 RocketMQ 会做三件事:向 Broker 申请锁定该队列(防止被其他消费者同时消费)、消费成功才提交 offset、失败时挂起当前队列无限重试而不是跳过(suspendCurrentQueueTimeMillis 默认 1 秒)。

顺序消息的两个大坑

坑一:发送端必须单线程同步发

这是文档里着墨不多但最容易踩的。如果发送端是多线程的,即使 hashKey 相同,线程 A 先调用 send、线程 B 后调用,也完全可能 B 的消息先落盘。

// 错误写法:多线程并发发,顺序不保证
orderIds.parallelStream().forEach(id -> {
    send(id, "PAID");
    send(id, "SHIPPED");
});

// 正确:同一订单的消息串行发
for (String id : orderIds) {
    send(id, "PAID");
    send(id, "SHIPPED");
}

而且要用 send 的同步版本。异步发送(asyncSendOrderly)虽然也能保证路由到同一队列,但回调线程执行顺序不定,落盘顺序仍然不保证。我们后来的做法是:在订单服务的业务线程池里,同一个订单的状态变更事件按 orderNo 做 hash 分派到同一个单线程 executor,异步发送只用于不要求顺序的场景(比如发通知)。

坑二:消费失败会阻塞整个队列

顺序消费不允许跳过消息。一条消息消费失败,RocketMQ 会暂停这个队列,等 suspendCurrentQueueTimeMillis 后重试同一条。如果这个队列里卡了一条毒丸消息,后面所有消息全部堆积。

我们上线第二周就遇到了:有个订单的收货地址字段超长(前端没限制,数据库字段是 VARCHAR(32),实际传了 60 个字符),物流服务 insert 报错。这条消息所在的队列整个卡死,20 分钟堆了 3 万多条。

现在的处理是加最大重试次数和死信投递:

@Override
public void onMessage(MessageExt msg) {
    int times = msg.getReconsumeTimes();
    if (times > 3) {
        // 别再卡队列了,扔到死信 topic 人工处理
        log.error("消息重试 {} 次仍失败,转死信, msgId={}", times, msg.getMsgId());
        mqTemplate.send("DLQ_ORDER_STATUS_TOPIC", msg);
        return;
    }
    try {
        handle(msg);
    } catch (Exception e) {
        throw new RuntimeException(e);     // 触发 orderly 的队列挂起重试
    }
}

更划算的做法:业务上规避顺序问题

折腾了一圈顺序消息之后,我心里的结论是:能不用顺序消息就别用。上面的方案让物流组改了状态机,改成"基于版本号"的处理,反而简单多了。

具体做法是:订单每次状态变更都带一个自增的 version,下游存下"当前处理到的最大版本号",乱序到达的旧版本直接丢弃。

CREATE TABLE t_order_snapshot (
    order_no    VARCHAR(32) NOT NULL PRIMARY KEY,
    status      VARCHAR(20) NOT NULL,
    version     BIGINT      NOT NULL DEFAULT 0,   /* 乐观锁版本号 */
    payload     JSON,
    update_time DATETIME    NOT NULL
) ENGINE = InnoDB;
<update id="upsertIfNewer">
    INSERT INTO t_order_snapshot (order_no, status, version, payload, update_time)
    VALUES (#{orderNo}, #{status}, #{version}, #{payload}, NOW())
    ON DUPLICATE KEY UPDATE
        status      = IF(VALUES(version) > version, VALUES(status), status),
        payload     = IF(VALUES(version) > version, VALUES(payload), payload),
        version     = GREATEST(version, VALUES(version)),
        update_time = IF(VALUES(version) > version, NOW(), update_time)
</update>

或者更常见的一种写法,消费时先查一次当前版本号:

public void handle(OrderEvent event) {
    Long current = snapshotMapper.getVersion(event.getOrderNo());
    if (current != null && event.getVersion() <= current) {
        log.info("旧版本消息丢弃, orderNo={}, msgVer={}, curVer={}",
                 event.getOrderNo(), event.getVersion(), current);
        return;                                  // 丢弃,返回消费成功
    }
    snapshotMapper.upsert(event);                // upsert 且更新 version
    doBusiness(event);
}

这个方案的好处是消费端可以随便并发,随便扩容,吞吐量不受影响。代价是业务上要接受"最终一致"——中间可能短暂出现状态回退,但只要版本对,最终一定收敛到最新状态。

我们改造前后的对比:

指标顺序消息方案版本号方案
消费 TPS(单消费者)230013500
消费线程数4(受队列数限制)32
扩容队列数固定 8,扩消费者无效加机器即可
异常消息影响阻塞整个队列只丢自己,进重试队列

什么时候才真的需要顺序消息

只有一种情况我认为绕不过去:业务上真的不允许中间状态被看到。比如 binlog 同步(MySQL 的 binlog 订阅必须严格按序重放,乱序执行 INSERT/UPDATE/DELETE 会导致数据错乱),比如数据库的增量同步到 ES 这类场景。这类场景没法靠版本号补救,因为操作本身不可交换。

订单状态机、积分发放、库存扣减这些,本质上都是"用最新值覆盖"的语义,用版本号完全够。

小结

  1. 全局有序 = 单队列单线程,吞吐量掉一个数量级,不要轻易用。
  2. RocketMQ 顺序消息要生效,发送端必须同步 + 同 key 串行发,消费端必须 ConsumeMode.ORDERLY
  3. 顺序消费的失败重试会挂起整个队列,一定要配最大重试次数和死信队列。
  4. 优先考虑版本号 / 状态机幂等,把顺序问题转成幂等问题,性能和运维都好得多。

参考