起因:财务对账时发现的几十笔差异
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 事务消息的流程
先搞清楚它到底做了什么。整个流程分四步:
- 生产者发一条半消息(Half Message)给 broker。broker 收到后,把消息的 topic 换成
RMQ_SYS_TRANS_HALF_TOPIC、queueId 改成 0,原 topic 和 queueId 存进消息属性里,然后持久化。此时消费者订阅不到这条消息,因为它根本不在目标 topic 上。 - broker 返回
SEND_OK,生产者开始执行本地事务。 - 生产者根据本地事务结果,向 broker 发二次确认:
COMMIT、ROLLBACK或UNKNOWN。 - 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 月) |
|---|---|---|
| 账实不符笔数 | 63 | 0 |
| 下单接口 P99 | 168 ms | 180 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 更可靠。
改造从设计到上线用了九天,其中三天在和财务一起核对历史数据、写补数脚本。事后看,最值的不是那套代码,而是把"跨服务一致性"这件事从口头约定变成了可验证的机制。