RocketMQ 的事务事务消息机制是其核心特性之一,专门用于解决分布式系统中 “本地事务与远程操作一致性” 问题(即分布式事务问题)。它通过 “两阶段提交” 思想,确保本地业务操作与消息发送的最终一致性,避免出现 “本地事务成功但消息未发出” 或 “消息发出但本地事务失败” 的情况。

先帮助大家快速理解一下:

RocketMQ 的事务消息,简单说就是保证 “本地办事” 和 “发消息通知” 要么都成,要么都不成的机制。

比如你网购下单:

  1. 系统先扣库存(本地办事),再发消息让支付系统生成支付单(发通知)。
  2. 怕就怕扣了库存,发消息时服务器崩了 —— 库存没了,支付单没生成,订单废了。

事务消息怎么解决?分三步:

  • 先 “占个坑”:发一条 “半消息” 给 MQ,这消息暂时藏着,别人看不到。
  • 办正事:扣库存。成了就告诉 MQ “消息可以发出去了”;败了就告诉 MQ “把那条消息删了”。
  • 怕失忆?回头查:如果扣完库存没来得及通知 MQ(比如服务器崩了),MQ 过一会儿会主动来问 “你刚才那事办没办成?”,系统查下库存记录,再告诉 MQ 结果。

这样就确保了:要么库存扣了、消息发了;要么库存没扣、消息也没发,两边一致,不添乱。

详::

核心场景:为什么需要事务消息?

举个典型例子:电商系统中 “下单减库存” 与 “发送支付通知” 的流程。

  • 正确逻辑:先扣减库存(本地事务),再发送 “待支付” 消息通知用户付款。
  • 问题:若扣减库存成功后,发送消息时服务器宕机,会导致用户收不到通知,订单异常。
  • 事务消息作用:确保 “扣减库存” 和 “发送消息” 要么都成功,要么都失败(或回滚到一致状态)。

事务消息机制原理:两阶段提交 + 回查补偿

RocketMQ 事务消息的核心是将消息发送分为 “预处理” 和 “确认 / 回滚” 两个阶段,并通过 “消息回查” 机制解决本地事务结果未知的问题。整体流程如下:

1. 阶段一:发送半事务消息(Half Message)
  • 半事务消息:一种特殊消息,发送到 RocketMQ 后,会被标记为 “暂不能投递” 状态,消费者无法消费。
  • 作用:验证消息发送通道是否通畅,同时预留消息存储空间,为后续确认 / 回滚做准备。

流程

  1. 生产者(业务系统)向 RocketMQ 发送 “半事务消息”(包含业务数据,如订单 ID、库存信息)。
  2. RocketMQ 收到消息后,存储到本地,并返回 “发送成功” 给生产者(此时消息处于 “半事务状态”,消费者不可见)。
2. 阶段二:执行本地事务 + 提交 / 回滚消息

生产者收到半事务消息发送成功的响应后,执行本地业务逻辑,并根据结果决定消息的最终状态。

流程

  1. 生产者执行本地事务(如扣减库存、创建订单)。
  2. 根据本地事务结果,向 RocketMQ 发送 “提交” 或 “回滚” 指令
    • 若本地事务成功:发送 “提交” 指令,RocketMQ 将半事务消息标记为 “可投递”,消费者可正常消费。
    • 若本地事务失败:发送 “回滚” 指令,RocketMQ 直接删除半事务消息,消费者不会收到。
3. 关键补偿:事务状态回查(解决未知状态)

若生产者在执行本地事务后,因网络中断、服务器宕机等原因,未向 RocketMQ 发送 “提交 / 回滚” 指令,RocketMQ 会主动发起 “事务回查”,确认本地事务的实际结果。

流程

  1. RocketMQ 定期(默认 60 秒,可配置)扫描处于 “半事务状态” 的消息,若超过超时时间未收到确认指令,则触发回查。
  2. RocketMQ 调用生产者预设的 “回查接口”(由业务系统实现),查询本地事务的实际状态(成功 / 失败)。
  3. 生产者根据本地事务日志(需业务系统记录,如订单表中的 “事务状态” 字段),返回明确结果:
    • 若本地事务成功:RocketMQ 提交消息,消费者可见。
    • 若本地事务失败:RocketMQ 回滚消息,删除该消息。
    • 若仍未知(如本地事务未执行完):RocketMQ 会在下次周期继续回查(默认最多回查 15 次,超过则丢弃消息)。

事务消息核心组件

  1. TransactionMQProducer:生产者类,专门用于发送事务消息,需配置 “本地事务执行器” 和 “回查处理器”。
  2. LocalTransactionExecuter:本地事务执行器,定义本地业务逻辑(如扣库存),返回 “提交 / 回滚 / 未知” 状态。
  3. TransactionCheckListener:事务回查监听器,定义回查逻辑(查询本地事务状态),用于 RocketMQ 回查时调用。

代码示例:事务消息发送与处理

以下是 Java 代码示例,展示事务消息的核心流程(基于 RocketMQ 4.x):

1. 定义生产者与事务处理器
import org.apache.rocketmq.client.producer.TransactionMQProducer;
import org.apache.rocketmq.client.producer.LocalTransactionExecuter;
import org.apache.rocketmq.client.producer.LocalTransactionState;
import org.apache.rocketmq.common.message.Message;
import org.apache.rocketmq.common.message.MessageExt;

public class TransactionProducer {
    public static void main(String[] args) throws Exception {
        // 1. 创建事务消息生产者
        TransactionMQProducer producer = new TransactionMQProducer("transaction_producer_group");
        producer.setNamesrvAddr("localhost:9876"); // NameServer 地址

        // 2. 设置本地事务执行器(执行本地业务 + 返回状态)
        producer.setTransactionExecuter(new LocalTransactionExecuter() {
            @Override
            public LocalTransactionState executeLocalTransactionBranch(Message msg, Object arg) {
                try {
                    // 解析消息中的业务数据(如订单ID)
                    String orderId = new String(msg.getBody());
                    // 执行本地事务:扣减库存(示例逻辑)
                    boolean stockDeductSuccess = deductStock(orderId);
                    
                    if (stockDeductSuccess) {
                        // 本地事务成功,返回“提交”
                        return LocalTransactionState.COMMIT_MESSAGE;
                    } else {
                        // 本地事务失败,返回“回滚”
                        return LocalTransactionState.ROLLBACK_MESSAGE;
                    }
                } catch (Exception e) {
                    // 异常情况(如网络问题),返回“未知”,等待回查
                    return LocalTransactionState.UNKNOW;
                }
            }
        });

        // 3. 设置事务回查监听器(RocketMQ 回查时调用)
        producer.setTransactionCheckListener((MessageExt msg) -> {
            // 解析消息,查询本地事务状态(需从数据库中查询订单/库存状态)
            String orderId = new String(msg.getBody());
            boolean isTransactionSuccess = checkLocalTransactionStatus(orderId);
            
            if (isTransactionSuccess) {
                return LocalTransactionState.COMMIT_MESSAGE;
            } else {
                return LocalTransactionState.ROLLBACK_MESSAGE;
            }
        });

        // 启动生产者
        producer.start();

        // 4. 发送半事务消息
        Message msg = new Message(
            "transaction_topic", // 主题
            "order_tag", // 标签
            "ORDER_12345".getBytes() // 消息体(订单ID)
        );
        producer.sendMessageInTransaction(msg, null); // 第二个参数为自定义业务参数

        // 关闭生产者(实际生产环境不轻易关闭)
        // producer.shutdown();
    }

    // 模拟扣减库存(本地事务)
    private static boolean deductStock(String orderId) {
        // 实际业务中操作数据库扣减库存,返回是否成功
        System.out.println("扣减订单 " + orderId + " 的库存...");
        return true; // 假设扣减成功
    }

    // 模拟查询本地事务状态(回查时调用)
    private static boolean checkLocalTransactionStatus(String orderId) {
        // 从数据库查询订单的事务状态(如订单表中“stock_deduct_status”字段)
        System.out.println("回查订单 " + orderId + " 的库存扣减状态...");
        return true; // 假设实际扣减成功
    }
}
2. 消费者接收消息(常规消费逻辑)
import org.apache.rocketmq.client.consumer.DefaultMQPushConsumer;
import org.apache.rocketmq.client.consumer.listener.ConsumeConcurrentlyContext;
import org.apache.rocketmq.client.consumer.listener.ConsumeConcurrentlyStatus;
import org.apache.rocketmq.client.consumer.listener.MessageListenerConcurrently;
import org.apache.rocketmq.common.message.MessageExt;

import java.util.List;

public class TransactionConsumer {
    public static void main(String[] args) throws Exception {
        DefaultMQPushConsumer consumer = new DefaultMQPushConsumer("transaction_consumer_group");
        consumer.setNamesrvAddr("localhost:9876");
        consumer.subscribe("transaction_topic", "*"); // 订阅事务消息主题

        // 注册消息监听器,接收并处理已提交的消息
        consumer.registerMessageListener((MessageListenerConcurrently) (msgs, context) -> {
            for (MessageExt msg : msgs) {
                String orderId = new String(msg.getBody());
                System.out.println("收到订单 " + orderId + " 的支付通知消息,开始通知用户...");
            }
            return ConsumeConcurrentlyStatus.CONSUME_SUCCESS; // 消费成功
        });

        consumer.start();
        System.out.println("消费者启动成功,等待接收消息...");
    }
}

关键特性与限制

  1. 最终一致性:事务消息不保证 “实时一致性”,但通过回查机制确保 “最终一致性”(最多回查 15 次,超时后消息被丢弃,需业务兜底)。
  2. 性能损耗:相比普通消息,多了 “半事务消息” 和 “回查” 步骤,性能略低,适合核心业务(如交易、支付)。
  3. 消息不可篡改:半事务消息一旦发送,内容不可修改,确保回查时数据一致。
  4. 依赖本地日志:回查机制依赖业务系统的本地事务日志(如订单状态表),需确保日志准确、可查询。

总结

RocketMQ 事务消息通过 “半事务消息预处理→本地事务执行→提交 / 回滚→回查补偿” 四步流程,解决了分布式系统中 “本地操作与消息发送” 的一致性问题。其核心是用 “两阶段提交”+“定时回查” 确保消息最终状态与本地事务一致,是分布式事务的轻量级、高可用解决方案。

Logo

北京人形旗下天工造物具身智能开源社区,聚焦具身天工与慧思开物两大平台

更多推荐