Administrator
发布于 2021-08-19 / 8091 阅读
180

RocketMQ 事务消息实现分布式最终一致

起因:财务对账时发现的几十笔差异

8 月初,财务同事甩过来一张表:7 月份有 63 笔订单,用户积分扣了但优惠券没发。我们日均 12 万单,63 笔占比万分之五,但财务按笔核对,一笔都不能有。

查下来问题出在下单流程。下单成功后要发一条消息给营销服务,让它发优惠券:

@Transactional
public void createOrder(CreateOrderReq req) {
    // 1. 写订单表
    orderMapper.insert(order);
    // 2. 扣库存
    inventoryClient.deduct(order.getItems());
    // 3. 扣积分
    pointClient.deduct(order.getUserId(), order.getPointAmount());
    // 4. 发消息通知营销服务发券
    rocketMQTemplate.convertAndSend("ORDER_CREATED_TOPIC", new OrderCreatedEvent(order.getId()));
}

这段代码有个经典问题。@Transactional 的提交发生在方法返回之后,而 convertAndSend 在方法内部就已经把消息发出去了。于是:

  • 情况 A:消息发出去了,但本地事务因为库存不足或者 pointClient 抛异常回滚了。订单不存在,营销服务却发了券。这就是财务看到的"扣了积分没发券"的反向情况,实际两种都有。
  • 情况 B:本地事务提交了,消息发送失败(网络抖动、broker 短暂不可用)。订单存在,券没发。

我们统计了 7 月的日志,63 笔里 41 笔是 A(事务回滚但消息已发),22 笔是 B(消息发送超时)。

有人说把 convertAndSend 挪到事务外面就行。确实能解决 A,但解决不了 B:事务提交成功后、发消息之前进程挂了,消息就永久丢了。这个窗口很小,但在 12 万单/天的量级上,一个月出现二十几次完全合理。

两个可选方案

业界成熟的就两条路:

本地消息表:在业务库建一张 t_message,和业务数据写在同一个本地事务里,保证"业务数据变了消息一定在表里"。然后起一个独立的投递任务扫表发消息,发成功改状态,失败一直重试。

RocketMQ 事务消息:利用 RocketMQ 的两阶段提交 + 状态回查。

我们选了后者,原因是前者要新建一个扫描任务、要考虑扫描频率、消息表归档,运维成本更高。但后面会写到,事务消息方案其实离不开本地事务表,这点很多人没意识到。

RocketMQ 事务消息的流程

先搞清楚它到底做了什么。整个流程分四步:

  1. 生产者发一条半消息(Half Message)给 broker。broker 收到后,把消息的 topic 换成 RMQ_SYS_TRANS_HALF_TOPIC、queueId 改成 0,原 topic 和 queueId 存进消息属性里,然后持久化。此时消费者订阅不到这条消息,因为它根本不在目标 topic 上。
  2. broker 返回 SEND_OK,生产者开始执行本地事务。
  3. 生产者根据本地事务结果,向 broker 发二次确认:COMMITROLLBACKUNKNOWN
  4. COMMIT 时 broker 从半消息恢复出原消息,写回原 topic,消费者可见;ROLLBACK 时写一条 OP 消息标记删除,不投递。

如果在第 3 步生产者挂了、或者返回了 UNKNOWN,broker 会主动回查生产者:调用 TransactionListener.checkLocalTransaction,问"这条半消息对应的本地事务到底成功了没有"。

回查的关键参数:

# broker.conf
transactionCheckInterval=60000        # 回查间隔,默认 60 秒
transactionCheckMax=15                # 单条消息最多回查 15 次

超过 15 次还拿不到结果,broker 默认把这条半消息回滚(写 OP 标记删除)并打印日志,消息就丢了。所以要保证回查逻辑一定能给出确定答案。

代码怎么写的

第一步,建本地事务表。这是必须的,因为回查时你要能查到"那次本地事务到底成没成":

CREATE TABLE t_transaction_log (
    id           BIGINT UNSIGNED NOT NULL AUTO_INCREMENT,
    tx_id        VARCHAR(64)  NOT NULL COMMENT '事务ID,与消息绑定',
    biz_type     VARCHAR(32)  NOT NULL,
    biz_id       VARCHAR(64)  NOT NULL,
    status       TINYINT      NOT NULL DEFAULT 0 COMMENT '0-未知 1-已提交 2-已回滚',
    created_at   DATETIME     NOT NULL DEFAULT CURRENT_TIMESTAMP,
    updated_at   DATETIME     NOT NULL DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP,
    PRIMARY KEY (id),
    UNIQUE KEY uk_tx_id (tx_id),
    KEY idx_status_created (status, created_at)
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4;

第二步,写 TransactionListener

@Slf4j
@Component
public class OrderTransactionListener implements TransactionListener {

    @Autowired
    private OrderService orderService;
    @Autowired
    private TransactionLogMapper transactionLogMapper;

    /**
     * 半消息发送成功后执行本地事务
     */
    @Override
    public LocalTransactionState executeLocalTransaction(Message msg, Object arg) {
        String txId = msg.getProperty(MessageConst.PROPERTY_TRANSACTION_ID);
        String orderJson = new String(msg.getBody(), StandardCharsets.UTF_8);

        try {
            // 本地事务里:写订单 + 写事务日志,同一个 @Transactional
            orderService.createOrderWithTxLog(orderJson, txId);
            return LocalTransactionState.COMMIT_MESSAGE;
        } catch (Exception e) {
            log.error("local transaction failed, txId={}", txId, e);
            return LocalTransactionState.ROLLBACK_MESSAGE;
        }
    }

    /**
     * broker 回查:查事务日志表给出确定答案
     */
    @Override
    public LocalTransactionState checkLocalTransaction(MessageExt msg) {
        String txId = msg.getProperty(MessageConst.PROPERTY_TRANSACTION_ID);
        TransactionLog log = transactionLogMapper.selectByTxId(txId);

        if (log == null) {
            // 查不到:本地事务可能还没开始,或者刚开始就崩了
            // 这里不能返回 ROLLBACK,否则会误杀。返回 UNKNOWN 让 broker 下次再问
            return LocalTransactionState.UNKNOW;
        }
        switch (log.getStatus()) {
            case 1:  return LocalTransactionState.COMMIT_MESSAGE;
            case 2:  return LocalTransactionState.ROLLBACK_MESSAGE;
            default: return LocalTransactionState.UNKNOW;
        }
    }
}

executeLocalTransaction 里有个细节:不能依赖它的返回值作为唯一信号,因为它可能返回之前进程就挂了。真正的关键是事务日志表里的记录,它和订单在同一个本地事务,要么都在要么都不在,这才是回查能给出正确答案的基础。

@Transactional(rollbackFor = Exception.class)
public void createOrderWithTxLog(String orderJson, String txId) {
    Order order = JSON.parseObject(orderJson, Order.class);
    orderMapper.insert(order);
    inventoryClient.deduct(order.getItems());
    pointClient.deduct(order.getUserId(), order.getPointAmount());

    // 和业务数据同一个事务
    transactionLogMapper.insert(new TransactionLog(txId, "ORDER",
            order.getOrderNo(), 1));
}

第三步,生产者:

@PostConstruct
public void init() throws MQClientException {
    producer = new TransactionMQProducer("order-tx-producer");
    producer.setNamesrvAddr(namesrvAddr);
    producer.setTransactionListener(orderTransactionListener);
    // 回查线程池,broker 回查请求走这里,别用默认的小线程池
    producer.setExecutorService(new ThreadPoolExecutor(
            4, 8, 60, TimeUnit.SECONDS,
            new LinkedBlockingQueue<>(2000),
            new ThreadFactoryBuilder().setNameFormat("tx-check-%d").build()));
    producer.start();
}

public void sendOrderCreated(Order order) {
    Message msg = new Message("ORDER_CREATED_TOPIC", "CREATE",
            JSON.toJSONString(order).getBytes(StandardCharsets.UTF_8));
    // 注意:事务消息会忽略延时属性,也不要用批量发送
    TransactionSendResult result = producer.sendMessageInTransaction(msg, null);
    log.info("send tx msg, txId={}, status={}",
             result.getLocalTransactionState(), result.getSendStatus());
}

两个必须注意的限制

第一,事务消息不支持延迟级别setDelayTimeLevel 在事务消息上会被忽略。我们要"下单 30 分钟未支付关单",只能另发一条普通延迟消息,不由事务消息承担。

第二,事务消息不支持批量发送。一条一条发,我们实测单次 RT 从原来的 3.2 ms 涨到 15.4 ms(多了一次半消息持久化和一次二次确认的网络往返)。下单接口整体 P99 从 168 ms 涨到 180 ms,可以接受,因为发消息可以挪到异步线程里。

和本地消息表方案的对比

维度RocketMQ 事务消息本地消息表
业务侵入要写 TransactionListener + 事务日志表要建消息表 + 写扫描任务
数据库压力事务日志表,写一次消息表,写一次 + 反复扫描
时效性事务提交后毫秒级投递取决于扫描间隔,通常秒级到分钟级
回查/重试broker 主动回查,最多 15 次自己控制重试次数和退避
组件依赖强依赖 RocketMQ 的半消息机制,换 MQ 就废与 MQ 无关,换 Kafka 也能用
排查难度半消息状态在 broker 里,要用 mqadmin 查,不直观直接查表,一目了然
表膨胀事务日志表只留近期,可定期清理消息表要归档,量大时麻烦

我个人的判断:如果团队已经深度使用 RocketMQ,用事务消息;如果 MQ 选型还没定死、或者业务对"可查可控"要求很高,本地消息表更稳

顺带说一个真相:这两个方案都要建本地表。区别在于本地消息表把"待发送的消息"也存下来(存全量消息体),事务消息只存"事务状态"(存一个 txId + status)。前者信息更全、能自己重投,后者更省空间。不存在"用了事务消息就不用建表"这回事。

线上效果

8 月 19 日上线之后,统计到 9 月底的一个完整月:

指标改造前(7 月)改造后(9 月)
账实不符笔数630
下单接口 P99168 ms180 ms
事务回查触发次数不适用1,842 次
回查后 COMMIT不适用1,797 次
回查后 ROLLBACK不适用45 次

一个月触发了 1842 次回查,占订单量的万分之五左右。触发原因主要是生产者在二阶段确认时网络抖动,或者正好赶上发版重启。这说明回查不是摆设,是真的在兜底。

另外补一个监控:我们给 RMQ_SYS_TRANS_HALF_TOPIC 的堆积加了告警,超过 200 条就通知。正常情况下半消息几秒内就会被确认掉,堆积起来说明回查逻辑出问题了。

$ sh mqadmin topicStatus -n 10.0.1.5:9876 -t RMQ_SYS_TRANS_HALF_TOPIC
#Broker Name  #QID  #Min Offset  #Max Offset  #Last Updated
broker-a        0      10238452     10238460   2021-09-28 14:22:31

消费端别忘了幂等

解决了"消息一定发得出去",还有一半问题是"消息可能重复"。RocketMQ 保证至少一次投递,加上回查机制,同一条消息被消费两次是完全可能的。营销服务的消费者必须做幂等:

// 用订单号做幂等键,MySQL 唯一索引兜底
try {
    couponGrantMapper.insertSelective(grant);   // uk_order_no
} catch (DuplicateKeyException e) {
    log.info("duplicate consume, orderNo={}", orderNo);
    return;                                     // 直接返回成功
}

我们在发券表上加了 uk_order_no 唯一索引,靠数据库兜底。比 Redis 幂等表更可靠,因为 Redis 的 SETNX 和数据库写入本身也有原子性问题。

小结

  • 先理解流程:半消息(改写 topic 到 RMQ_SYS_TRANS_HALF_TOPIC)→ 本地事务 → 二次确认 → 必要时 broker 回查。半消息对消费者不可见是理解这套机制的钥匙。
  • 回查逻辑必须查持久化状态,也就是必须建事务日志表,且这张表要和业务数据在同一个本地事务里。返回 UNKNOW 让 broker 重试,比猜一个结果安全。
  • 两个硬限制:事务消息不支持延迟级别,不支持批量发送。要延迟就另发普通消息。
  • 回查默认最多 15 次、间隔 60 秒,超过就丢弃。回查线程池要自己配,默认的可能扛不住。
  • 本地消息表和事务消息不是二选一的对立关系,两者都要建表。选哪个主要看 MQ 是否锁死、以及团队更想要"broker 托管"还是"自己可控"。
  • 消费端幂等是配套的另一半,用数据库唯一索引比 Redis 更可靠。

改造从设计到上线用了九天,其中三天在和财务一起核对历史数据、写补数脚本。事后看,最值的不是那套代码,而是把"跨服务一致性"这件事从口头约定变成了可验证的机制。

参考