RocketMQNotes

第 10 章:事务消息

zjc 于 2026-01-10 发布

这是《RocketMQ 零基础实战指南》的独立章节版。本章从概念、实操和生产排查三个视角展开,代码块保留了原书可直接运行的版本。 事务消息解决“本地数据库事务”和“消息发送”之间的原子性问题:先写本地事务,再提交消息;如果本地事务成功但消息提交失败,通过回查补齐状态。它不是分布式事务协议,也不是 XA 事务,而是基于最终一致的事件发送机制。

10.1 问题背景

常见错误写法:

begin transaction
  insert order
  send message
commit

消息发送成功后数据库回滚,会产生“订单不存在但事件已发出”的脏事件。数据库提交成功后进程崩溃,则可能漏发事件。

正确思路:

half message
  -> execute local transaction
  -> commit or rollback message

10.2 核心流程

Producer                Broker                Consumer
  |--- half message --->  |
  |<-- half ok ---------- |
  |                       |
  | execute local tx      |
  | commit / rollback     |
  |--- commit ----------> |--- real message ---> |
  |                       |
  |<------ checkback ---- |
  |--- commit/rollback -> |

状态:

状态 含义
COMMIT_MESSAGE 半消息转为真实消息
ROLLBACK_MESSAGE 删除半消息,不投递
UNKNOW 稍后回查

10.3 生产者实现

public class OrderTransactionListener implements TransactionListener {
    private final OrderService orderService;
    private final EventRecordRepository eventRecordRepository;

    @Override
    public LocalTransactionState executeLocalTransaction(Message msg, Object arg) {
        OrderCreateCommand command = (OrderCreateCommand) arg;
        try {
            orderService.createWithEvent(command, msg.getTransactionId());
            return LocalTransactionState.COMMIT_MESSAGE;
        } catch (Exception e) {
            log.error("local transaction failed, orderNo={}", command.orderNo(), e);
            return LocalTransactionState.ROLLBACK_MESSAGE;
        }
    }

    @Override
    public LocalTransactionState checkLocalTransaction(MessageExt msg) {
        String transactionId = msg.getTransactionId();
        return eventRecordRepository.findByTransactionId(transactionId)
                .map(record -> switch (record.getStatus()) {
                    case COMMITTED -> LocalTransactionState.COMMIT_MESSAGE;
                    case ROLLED_BACK -> LocalTransactionState.ROLLBACK_MESSAGE;
                    case PENDING -> LocalTransactionState.UNKNOW;
                })
                .orElse(LocalTransactionState.UNKNOW);
    }
}

发送:

TransactionMQProducer producer = new TransactionMQProducer("order-tx-producer");
producer.setNamesrvAddr("rocketmq-namesrv:9876");
producer.setTransactionListener(new OrderTransactionListener(orderService, repository));
producer.start();

Message message = new Message("OrderTransactionTopic", "OrderCreated",
        event(command).getBytes(StandardCharsets.UTF_8));

TransactionSendResult result = producer.sendMessageInTransaction(message, command);

10.4 本地事务表

事务消息仍推荐与本地事件表配合:

CREATE TABLE outbox_event (
  id BIGINT PRIMARY KEY AUTO_INCREMENT,
  event_id VARCHAR(64) NOT NULL,
  transaction_id VARCHAR(64) NOT NULL,
  aggregate_id VARCHAR(64) NOT NULL,
  event_type VARCHAR(64) NOT NULL,
  status VARCHAR(20) NOT NULL,
  payload JSON NOT NULL,
  created_at DATETIME(6) NOT NULL,
  updated_at DATETIME(6) NOT NULL,
  UNIQUE KEY uk_event_id (event_id),
  UNIQUE KEY uk_transaction_id (transaction_id),
  KEY idx_status_time (status, created_at)
);

本地事务内:

INSERT INTO orders(...);
INSERT INTO outbox_event(...);

回查时以 outbox_event 为准,而不是根据内存对象或不确定的日志判断。

10.5 回查设计

回查必须满足:

  1. 幂等;
  2. 可持久查询;
  3. 返回确定状态;
  4. 不依赖已经丢失的内存上下文;
  5. 有超时和最大回查限制;
  6. 对异常返回 UNKNOW
  7. 状态最终落库。

推荐回查逻辑:

query outbox by transactionId
  found COMMITTED    -> COMMIT_MESSAGE
  found ROLLED_BACK  -> ROLLBACK_MESSAGE
  found PENDING      -> UNKNOW
  not found          -> check timeout
      before timeout -> UNKNOW
      after timeout  -> ROLLBACK_MESSAGE 或人工治理

10.6 消费端语义

事务消息只约束生产端提交关系,不保证消费端与下游数据库强一致:

Broker -> Consumer -> downstream DB

消费端仍需要:

  1. 幂等;
  2. 重试;
  3. 死信治理;
  4. 对账;
  5. 业务状态机;
  6. 补偿任务。

如果消费失败,生产端事务不会自动回滚。不要把事务消息当成同步分布式事务。

10.7 与本地消息表对比

方案 优点 限制
RocketMQ 事务消息 有回查机制,减少轮询 依赖客户端和 Broker 能力
本地消息表 简单直观,易审计 需要定时扫描和重发
Outbox + CDC 解耦清晰 引入 CDC 组件

中大型系统可以组合使用:事务消息负责即时提交,本地表负责审计和补偿。

10.8 生产注意事项

  1. 生产者实例与监听器正确绑定;
  2. 回查接口响应要快;
  3. 回查逻辑必须有持久状态;
  4. UNKNOW 不能无限返回;
  5. 记录 half message ID、transactionId 和业务单号;
  6. 监控半消息数量和回查次数;
  7. 升级客户端前验证事务 API 兼容性;
  8. 生产环境演练进程崩溃和主备切换。

10.9 常见错误

错误 后果
本地事务内存判断 重启后无法回查
捕获异常后返回 COMMIT 本地失败但消息投递
消费失败期待生产事务回滚 语义理解错误
回查无超时 半消息长期悬挂
未记录 transactionId 无法定位事件
未做幂等 补发导致重复消费

本章小结

事务消息通过半消息、本地事务和回查机制,降低数据库提交与消息发送不一致的概率。它提供的是生产端最终一致,不等于全局分布式事务。可靠的落地关键是本地事务表、确定性回查、幂等消费和完善的监控补偿。

思考题

  1. 半消息为什么对消费者不可见?
  2. 回查为什么必须查询持久化状态?
  3. 消费失败时生产端事务是否会回滚?
  4. 事务消息和本地消息表如何组合?
  5. 半消息数量持续增长说明什么问题?