Skip to content

RocketMQ 事务消息机制与实现

问题

RocketMQ 的事务消息是如何实现分布式事务的?和 Kafka 的事务有什么区别?在生产环境中使用时有哪些坑?

分析

分布式事务是微服务架构中最棘手的问题之一。当一次业务操作涉及多个服务(如订单服务、库存服务、支付服务),如何保证数据的一致性?传统的 XA 两阶段提交虽然能保证强一致性,但性能差、实现复杂、容易造成资源锁定。而 RocketMQ 的事务消息提供了一种最终一致性的轻量级方案,核心思想是"消息发送和本地事务要么同时成功,要么同时失败"。

为什么需要事务消息?

先看一个典型的电商下单场景:

  1. 用户在下单页面提交订单
  2. 订单服务创建订单(写入订单表)
  3. 通知库存服务扣减库存
  4. 通知积分服务增加积分

如果直接用普通 MQ 消息,可能出现两种情况:

场景本地事务消息发送结果
先发消息后执行事务消息已发送执行失败库存扣了但订单没创建,数据不一致
先执行事务后发消息事务已提交发送失败(网络超时/Broker 宕机)订单创建了但库存没扣,超卖
发消息 + 事务都成功但消费端重复消费事务已提交消息被 Broker 重投库存多扣,少卖了

事务消息就是为了解决"发送消息"和"本地事务"之间的原子性问题。注意消费端幂等是"另一层问题",事务消息只保证 Producer 端的原子性,消费端幂等要自己实现。

RocketMQ 事务消息的核心机制

RocketMQ 的事务消息使用 半消息(Half Message) + 事务反查 机制,时序流程如下:

Producer                     Broker                     Consumer
   │                           │                           │
   │  1. 发送半消息              │                           │
   ├──────────────────────────►│                           │
   │                           │ 写入 RMQ_SYS_TRANS_HALF_TOPIC
   │                           │ 消息状态 = PREPARED        │
   │  2. 返回半消息结果          │                           │
   │◄──────────────────────────┤                           │
   │                           │                           │
   │  3. 回调 executeLocalTransaction()                     │
   │   └─ 执行本地事务         │                           │
   │     ├─ 订单表 INSERT ✅   │                           │
   │     ├─ 库存表扣减 ✅      │                           │
   │     └─ 积分表插入 ✅      │                           │
   │                           │                           │
   │  4. 返回 COMMIT/ROLLBACK   │                           │
   ├──────────────────────────►│                           │
   │                           │  COMMIT: 投递到目标 Topic  │
   │                           │  ROLLBACK: 删除半消息      │
   │                           │                           │
   │                           │  5. Consumer 消费全消息     │
   │                           ├──────────────────────────►│
   │                           │                           │
   │  ── 如果 Producer 在步骤 3 崩溃 ──                     │
   │                           │                           │
   │                           │  6. 定时扫描(60s 一次)     │
   │                           │  回调 checkLocalTransaction()
   │                           │◄──────────────────────────┤
   │                           │  7. 返回 COMMIT/ROLLBACK   │
   │                           ├──────────────────────────►│
   │                           │                           │

这个机制的关键在于半消息:消息先发送到 MQ,但处于"半可见"状态,Consumer 是看不到的。只有本地事务确认成功后,消息才会变成"全可见"状态,投递到目标 Topic。

半消息的存储实现

半消息写入 RocketMQ 的内部系统 Topic:RMQ_SYS_TRANS_HALF_TOPIC。这个 Topic 有 1 个 Queue(默认配置),所有事务消息的半消息都写入这个 Queue。

一条半消息的 CommitLog 记录结构如下:

字段内容说明
msgId系统生成的全局唯一 ID用于反查时定位
origTopic用户指定的目标 Topic如 "order-tx-topic"
queueId目标 Topic 的 Queue ID提交时投递到正确分区
queueOffset目标 Queue 的 Offset提交时消息顺序
body业务消息体序列化后的业务数据
transactionStateCOMMIT/ROLLBACK/PREPARED默认 PREPARED
preparedTransactionOffset半消息在 CommitLog 中的偏移提交时回填

事务提交时,RocketMQ 从 RMQ_SYS_TRANS_HALF_TOPIC 删除半消息(标记为已提交),同时将消息投递到 origTopicqueueId 对应 Queue。这种设计的好处是:对目标 Topic 的 ConsumeQueue 索引没有侵入性,正常消费逻辑无需感知事务消息的存在。代价是写入时多了一次 IO(半消息 Topic 的 CommitLog 写入),提交时又多了两次 IO(删除半消息记录 + 写入目标 Topic 的 CommitLog)。

事务反查的详细机制

反查是 RocketMQ 事务消息最核心的"保底"机制。当 Producer 在提交半消息后崩溃(或者网络分区导致超时),RocketMQ 无法确定本地事务的状态,就会主动回调 Producer 的 checkLocalTransaction() 方法获取事务状态。

反查的全流程时序:

TransactionMessageCheckService 线程

         │  每 60s 扫描一次
         ├── 扫描 RMQ_SYS_TRANS_HALF_TOPIC 的 ConsumeQueue

         ├── 发现一条 PREPARED 状态的半消息
         │   └─ 检查是否已超过 transactionTimeOut(默认 6s)

         ├── 检查反查次数是否超过 transactionCheckMax(默认 15 次)
         │   ├─ 未超限 → 发送反查请求到 Producer
         │   └─ 已超限 → 将消息转移到 DLQ(死信队列)

         ├── Producer 收到反查请求
         │   └─ 调用 checkLocalTransaction()
         │       ├─ COMMIT_MESSAGE → Broker 提交消息
         │       ├─ ROLLBACK_MESSAGE → Broker 回滚消息
         │       └─ UNKNOWN → Broker 等待下次反查(60s 后)

         └── 如果反查请求发送失败(Producer 仍不可达)
             └─ 等待下次扫描周期(60s 后再试)

反查的触发条件

  1. 半消息在 RMQ_SYS_TRANS_HALF_TOPIC 中存活超过 transactionTimeOut(默认 6s)
  2. 该半消息的反查次数 ≤ transactionCheckMax(默认 15 次)
  3. Broker 端的 TransactionMessageCheckService 线程正常运行

反查超时的兜底处理:如果反查达到 15 次仍无法确定事务状态,消息会被转移到死信队列。此时需要人工介入,通过 rmqadmin 工具手动提交或回滚:

bash
# 手动提交事务消息
mqadmin commitTransaction -n 127.0.0.1:9876 -t RMQ_SYS_TRANS_HALF_TOPIC -i msgId

# 手动回滚事务消息
mqadmin rollbackTransaction -n 127.0.0.1:9876 -t RMQ_SYS_TRANS_HALF_TOPIC -i msgId

生产事故案例:反查阻塞导致事务消息堆积

背景:某电商平台促销活动期间,订单服务突然响应变慢,订单创建失败率升高。

排查过程

  1. 监控看到 RocketMQ Broker 的 RMQ_SYS_TRANS_HALF_TOPIC 消息数从正常几百条飙升到 10 万+
  2. 查看 Broker 日志,发现大量 checkTransaction failed 的 WARN 日志
  3. 进一步排查发现,事务反查调用 checkLocalTransaction() 时,订单服务查询数据库的 SQL 走了慢查询(全表扫描)
  4. 反查响应超时,每次都返回 UNKNOWN,导致 Broker 不断重试,进一步加重了 Broker 的扫描线程负载

根因checkLocalTransaction() 反查接口没有做索引优化,SELECT 用 order_id 查但表没有索引,每次反查耗时 2-3 秒。Broker 的 TransactionMessageCheckService 单线程扫描,反查超时导致后续半消息排队等待,恶性循环。

修复

  1. order_id 字段加唯一索引,反查耗时降到 1ms
  2. checkLocalTransaction() 中增加缓存,30s 内查过的订单直接返回
  3. 调整 transactionCheckInterval 从 60s 降到 30s,加速反查周期
  4. 手动处理了已经堆积的 10 万+半消息:通过 mqadmin 批量查询数据库状态后提交

代码示例

事务消息发送端(修正版)

java
@Component
public class OrderTransactionProducer {

    @Autowired
    private RocketMQTemplate rocketMQTemplate;

    @Autowired
    private OrderService orderService;

    // 事务消息发送
    public void createOrder(OrderDTO orderDTO) {
        // 构建消息,orderId 放入 header 供反查时使用
        Message<OrderDTO> message = MessageBuilder
            .withPayload(orderDTO)
            .setHeader("orderId", orderDTO.getOrderId())
            .build();

        TransactionSendResult result = rocketMQTemplate.sendMessageInTransaction(
            "order-tx-producer-group",  // producer group,需与 listener 保持一致
            "order-tx-topic",            // 目标 Topic
            message,                     // 消息内容
            orderDTO.getOrderId()        // 事务参数,传给 executeLocalTransaction 的 arg
        );

        if (result.getLocalTransactionState() == LocalTransactionState.COMMIT_MESSAGE) {
            log.info("订单事务消息提交成功: orderId={}", orderDTO.getOrderId());
        } else if (result.getLocalTransactionState() == LocalTransactionState.ROLLBACK_MESSAGE) {
            log.warn("订单事务消息回滚: orderId={}", orderDTO.getOrderId());
        } else {
            log.warn("订单事务状态未知,等待反查: orderId={}", orderDTO.getOrderId());
        }
    }

    // 本地事务执行器 + 反查
    @RocketMQTransactionListener(txProducerGroup = "order-tx-producer-group")
    public class OrderTransactionListener implements RocketMQLocalTransactionListener {

        @Override
        @Transactional
        public RocketMQLocalTransactionState executeLocalTransaction(Message msg, Object arg) {
            String orderId = (String) arg;
            try {
                // 解析消息体
                OrderDTO orderDTO = (OrderDTO) ((RocketMQLocalTransactionMessage) msg).getPayload();

                // 执行本地事务:创建订单 + 扣减本地库存 + 增加积分
                orderService.createOrder(orderDTO);

                log.info("本地事务执行成功: orderId={}", orderId);
                return RocketMQLocalTransactionState.COMMIT;

            } catch (DataIntegrityViolationException e) {
                // 订单号重复(幂等 INSERT),无需回滚
                log.warn("订单已存在,视为成功: orderId={}", orderId);
                return RocketMQLocalTransactionState.COMMIT;
            } catch (Exception e) {
                log.error("本地事务执行失败,回滚消息: orderId={}", orderId, e);
                return RocketMQLocalTransactionState.ROLLBACK;
            }
        }

        @Override
        public RocketMQLocalTransactionState checkLocalTransaction(Message msg) {
            // 从 header 中获取 orderId
            // 注意:RocketMQ Spring 的 checkLocalTransaction 中 msg 是 generic Message
            // 如果使用 spring-messaging 的 Message,需要通过 HeaderAccessor 获取
            String orderId = null;
            try {
                // 方式一:直接转换(RocketMQ 原生 Message)
                if (msg instanceof org.apache.rocketmq.common.message.Message) {
                    org.apache.rocketmq.common.message.Message extMsg =
                        (org.apache.rocketmq.common.message.Message) msg;
                    orderId = extMsg.getProperty("orderId");
                }
                // 方式二:spring-messaging Message
                else if (msg instanceof org.springframework.messaging.Message) {
                    orderId = (String) ((org.springframework.messaging.Message<?>) msg)
                        .getHeaders().get("orderId");
                }

                if (orderId == null) {
                    log.warn("反查消息中无 orderId");
                    return RocketMQLocalTransactionState.UNKNOWN;
                }

                // 查询订单状态(带缓存,30s 过期)
                Order order = orderService.getOrderByIdWithCache(orderId);
                if (order == null) {
                    // 订单不存在,可能事务还没提交,返回 UNKNOWN 等下次反查
                    return RocketMQLocalTransactionState.UNKNOWN;
                }

                switch (order.getStatus()) {
                    case CREATED:
                    case PAID:
                        return RocketMQLocalTransactionState.COMMIT;
                    case CANCELLED:
                    case REFUNDED:
                        return RocketMQLocalTransactionState.ROLLBACK;
                    default:
                        return RocketMQLocalTransactionState.UNKNOWN;
                }

            } catch (Exception e) {
                log.error("反查异常: orderId={}", orderId, e);
                // 返回 UNKNOWN 让 Broker 下次重试,不要返回 ROLLBACK 导致误删
                return RocketMQLocalTransactionState.UNKNOWN;
            }
        }
    }
}

消费端实现(幂等消费 + 业务处理)

java
@Component
public class OrderConsumer {

    @Autowired
    private RedisTemplate<String, String> redisTemplate;

    @Autowired
    private InventoryService inventoryService;

    @Autowired
    private PointsService pointsService;

    @RocketMQMessageListener(
        topic = "order-tx-topic",
        consumerGroup = "order-consumer-group",
        // 消费者线程数,根据业务吞吐量调整
        consumeThreadNumber = 20,
        // 最大重试次数,默认 16
        maxReconsumeTimes = 3
    )
    public class OrderMessageListener implements RocketMQListener<OrderDTO> {

        @Override
        public void onMessage(OrderDTO message) {
            String orderId = message.getOrderId();
            String dedupKey = "order:dedup:" + orderId;

            // 幂等性检查:使用 Redis SETNX + 过期时间
            // 注意:过期时间要大于业务处理时间 + 重试窗口,避免"真"重复
            Boolean isFirst = redisTemplate.opsForValue()
                .setIfAbsent(dedupKey, "1", Duration.ofHours(24));

            if (Boolean.FALSE.equals(isFirst)) {
                log.info("消息已处理过,跳过幂等消费: orderId={}", orderId);
                return;
            }

            try {
                // 执行下游业务:扣库存 + 加积分
                // 库存服务
                inventoryService.deductStock(message.getProductId(), message.getQuantity());

                // 积分服务
                pointsService.addPoints(message.getUserId(), message.getTotalAmount() / 10);

                // 处理成功,业务完成

            } catch (Exception e) {
                // 处理失败,删除幂等标记,让 RocketMQ 重试
                redisTemplate.delete(dedupKey);
                log.error("处理订单消息失败,将重试: orderId={}", orderId, e);
                throw new RuntimeException("处理订单消息失败", e);
            }
        }
    }
}

事务消息配置参数详解

yaml
# application.yml
rocketmq:
  name-server: 192.168.1.100:9876;192.168.1.101:9876
  producer:
    group: order-tx-producer-group
    # 发送消息超时时间,默认 3000ms,事务消息建议调大
    send-message-timeout: 5000
    # 失败重试次数,默认 2
    retry-times-when-send-failed: 2
    # 事务消息超时时间,默认 6000ms
    transaction-timeout: 10000
    # 最大反查次数,默认 15
    transaction-check-max: 10
    # 反查间隔,默认 60000ms
    transaction-check-interval: 30000
properties
# broker.conf 关键配置
# 事务消息反查线程池大小,默认 4,堆消息多时需调大
transactionCheckThreadPoolNums=8
# 半消息 Topic 队列数,默认 1
halfTopicQueueNums=1
# 事务消息反查间隔,默认 60000ms
transactionCheckInterval=30000
# 事务消息超时时间,默认 6000ms
transactionTimeOut=10000
# 最大反查次数,默认 15
transactionCheckMax=10

总结

RocketMQ 的事务消息用半消息 + 回调反查机制,实现了 Producer 端本地事务和消息发送的原子性,是一种最终一致性的分布式事务方案。它和 Kafka 的事务有本质区别:Kafka 的事务是跨分区原子写入的 EOS(Exactly-Once Semantics),解决的是流处理中的精确一次语义;RocketMQ 的事务消息解决的是跨系统的分布式事务一致性,通过反查机制保证即使 Producer 崩溃也能恢复事务状态。

使用注意事项

  1. 事务超时时间transactionTimeOut 默认为 6 秒,如果本地事务执行超过该时间,RocketMQ 会开始回调反查。对于执行时间较长的本地事务(如跨库写入、文件上传),需要适当调大这个值。调大后同时调整 transactionCheckInterval,避免反查线程空转。

  2. 反查接口必须幂等且快checkLocalTransaction() 可能被多次调用,必须保证幂等性。实现方式:以订单 ID 为唯一标识,查询数据库判断事务状态。反查接口必须快速返回(< 100ms),否则会阻塞 Broker 的反查线程,导致大量半消息堆积。

  3. Consumer 侧必须幂等:事务消息可能被多次投递(如反查超时后消息被重新投递、Consumer 端处理超时后 Broker 重投),消费端必须做好幂等处理。推荐使用 Redis SETNX 或数据库唯一键做去重表。

  4. 半消息 Topic 堆积监控:所有事务消息的半消息都存储在 RMQ_SYS_TRANS_HALF_TOPIC,如果大量事务消息长时间处于"半消息"状态(反查超时或 Producer 长时间不可达),会导致这个 Topic 堆积,影响 Broker 性能。需要配置 Prometheus 告警,阈值建议:半消息数 > 5000 触发告警。

  5. 反查失败的处理transactionCheckMax 默认 15 次,超过后消息被转移到 DLQ。需要配置告警通知人工介入,或者通过 rmqadmin 手动提交/回滚。

  6. 性能影响:事务消息相比普通消息,多了半消息写入(一次 CommitLog 追加)、反查定时扫描(每 60s 一次)、消息恢复(提交时再写一次 CommitLog)等环节,Pub 端吞吐量大约下降 20-30%。对于高吞吐场景(如日志收集、埋点上报),应该用普通消息。对于低吞吐但对一致性要求高的场景(如订单、支付、资金),值得用事务消息。

与 Kafka 事务的对比

特性RocketMQ 事务消息Kafka 事务
解决的问题分布式事务最终一致性流处理 Exactly-Once
核心机制半消息 + 回调反查PID + 序列号 + 事务协调器
跨系统一致性支持(MQ ↔ 数据库)不支持(仅 Kafka 内部原子写入)
Consumer 隔离半消息对 Consumer 不可见isolation.level=read_committed 隔离事务消息
性能影响Pub 端下降约 20-30%生产端增加 1 次 RTT(事务协调器交互)
运维复杂度需关注反查超时和半消息堆积需关注事务协调器健康、日志清理
适用场景跨系统分布式事务(订单、支付)流式 ETL、Kafka Streams 状态一致性
反查依赖依赖 Producer 服务在线无(事务状态持久化在 Kafka 内部 Topic)

替代方案:本地消息表

如果不想依赖 MQ 的事务特性,可以考虑更轻量的本地消息表方案:

java
@Transactional
public void placeOrder(OrderDTO order) {
    // 1. 插入订单表
    orderDao.insert(order);
    // 2. 插入消息表(同一事务,相同的数据库连接)
    messageDao.insert(new MessageRecord(
        order.getOrderId(),
        "order-topic",
        JSON.toJSONString(order),
        MessageStatus.PENDING
    ));
}

// 3. 定时任务轮询消息表,发送未发送的消息
@Scheduled(fixedDelay = 5000)
public void sendPendingMessages() {
    // 批量拉取 PENDING 状态的消息,每次 100 条
    List<MessageRecord> pending = messageDao.findByStatus(
        MessageStatus.PENDING,
        PageRequest.of(0, 100)
    );

    for (MessageRecord record : pending) {
        try {
            SendResult result = rocketMQTemplate.syncSend(
                record.getTopic(),
                record.getContent()
            );
            messageDao.updateStatus(record.getId(), MessageStatus.SENT);
            log.info("本地消息表消息发送成功: id={}", record.getId());
        } catch (Exception e) {
            log.error("本地消息表消息发送失败,下次重试: id={}", record.getId(), e);
            // 不更新状态,下次定时任务继续发送
        }
    }
}

// 4. 补偿:监控长时间未发送的消息
@Scheduled(cron = "0 0 2 * * ?")  // 每天凌晨 2 点
public void compensatePendingMessages() {
    // 查找超过 1 小时仍未发送的消息
    List<MessageRecord> timeout = messageDao.findByStatusAndCreateTimeBefore(
        MessageStatus.PENDING,
        LocalDateTime.now().minusHours(1)
    );
    for (MessageRecord record : timeout) {
        // 发送告警,人工介入检查
        alertService.sendAlert("消息发送超时", record);
    }
}

本地消息表 vs 事务消息选型对比

维度本地消息表RocketMQ 事务消息
依赖需要数据库 + 定时任务框架需要 RocketMQ 4.x+
复杂度需自行实现消息表管理、轮询、补偿MQ 内置,开箱即用
一致性保证数据库本地事务保证半消息 + 反查保证
实时性依赖轮询间隔(秒级延迟)半消息提交后立即投递(毫秒级)
运维成本需要监控消息表大小和堆积需要监控半消息 Topic 和反查线程
消息可靠性消息表持久化,不丢消息半消息持久化,不丢消息

选型建议

  • 如果订单量不大(日均 < 10 万)、团队已有定时任务基础设施 → 本地消息表,更可控
  • 如果订单量大、延迟敏感(秒级必须到达)、已经有 RocketMQ 运维经验 → 事务消息,更高效
  • 两种方案可以共存:核心交易链路用事务消息,内部补偿任务用本地消息表做兜底

手撕 → 框架 → 生产化,一步步把 AI Agent 工程化搞透。