Skip to content

事务消息

事务消息为Apache RocketMQ中的高级特性消息。

应用场景

分布式系统调用的特点为一个核心业务逻辑的执行,同时需要调用多个下游业务进行处理。因此,如何保证核心业务和多个下游业务的执行结果完全一致,是分布式事务需要解决的主要问题。 事务场景图 以电商交易场景为例,用户支付订单这一核心操作的同时会涉及到下游物流发货、积分变更、购物车状态清空等多个子系统的变更。当前业务的处理分支包括:

  • 主分支订单系统状态更新:由未支付变更为支付成功。
  • 物流系统状态新增:新增待发货物流记录,创建订单物流记录。
  • 积分系统状态变更:变更用户积分,更新用户积分表。
  • 购物车系统状态变更:清空购物车,更新用户购物车记录。

解决方案

基于传统XA事务方案

关键点:性能不足

为了保证上述四个分支的执行结果一致性,典型方案是基于XA协议的分布式事务系统来实现。 将四个调用分支封装成包含四个独立事务分支的大事务。 基于XA分布式事务的方案可以满足业务处理结果的正确性,但最大的缺点是多分支环境下资源锁定范围大,并发度低,随着下游分支的增加,系统性能会越来越差。

XA事务

XA事务是一种分布式事务协议,用于确保多个数据库或资源管理器之间的事务一致性。以下是关于XA事务的详细解释:

定义与组成:

XA协议由Tuxedo首先提出,并由X/Open组织标准化,作为资源管理器(如数据库)与事务管理器之间的接口标准。它包括事务管理器(TM)和资源管理器(RM),其中资源管理器通常由数据库实现,如Oracle、DB2、MySQL等。

工作原理:

XA事务基于两阶段提交(2PC)协议实现,确保数据的强一致性。第一阶段是准备阶段,所有参与者准备执行事务并锁定所需资源;第二阶段是提交阶段,当事务管理器确认所有参与者都准备好后,发送提交命令。

优点与缺点:

XA事务的优点在于能够确保分布式系统中多个数据库或资源之间的事务一致性。然而,它也有性能上的缺点,通常不适合高并发场景,因为其性能可能远低于单个数据库的事务处理。

实际应用:

在需要跨多个数据库或资源执行事务一致性的场景中,如金融交易系统,XA事务被广泛应用。

总体来说:

XA事务是一种强大的技术,用于确保分布式系统中数据的一致性和完整性,但在选择使用时应考虑其性能影响。

基于普通消息方案

关键点:一致性保障困难

将上述基于XA事务的方案进行简化,将订单系统变更作为本地事务,剩下的系统变更作为普通消息的下游来执行,事务分支简化成普通消息+订单表事务,充分利用消息异步化的能力缩短链路,提高并发度。 事务和普通消息 该方案中消息下游分支和订单系统变更的主分支很容易出现不一致的现象,例如:

  • 消息发送成功,订单没有执行成功,需要回滚整个事务。
  • 订单执行成功,消息没有发送成功,需要额外补偿才能发现不一致。
  • 消息发送超时未知,此时无法判断需要回滚订单还是提交订单变更。

基于事务消息方案

关键点:支持最终一致性

上述普通消息方案中,普通消息和订单事务无法保证一致的原因,本质上是由于普通消息无法像单机数据库事务一样,具备提交、回滚和统一协调的能力。

而基于Apache RocketMQ实现的分布式事务消息功能,在普通消息基础上,支持二阶段的提交能力。将二阶段提交和本地事务绑定,实现全局提交结果的一致性。 最终一致性 Apache RocketMQ事务消息的方案,具备高性能、可扩展、业务开发简单的优势。具体事务消息的原理和流程,请参见下文的功能原理。

功能原理

什么是事务消息?

事务消息是 Apache RocketMQ 提供的一种高级消息类型,支持在分布式场景下保障消息生产和本地事务的最终一致性。

事务消息处理流程

事务流程

mermaid
sequenceDiagram
    participant P as Producer
    participant B as RocketMQ Broker
    participant DB as 本地数据库
    participant C as Consumer

    P->>B: 发送半事务消息
    B-->>P: 半消息持久化成功 Ack
    P->>DB: 执行本地事务
    alt 本地事务成功
        P->>B: Commit
        B->>C: 投递消息
    else 本地事务失败
        P->>B: Rollback
        B-->>B: 丢弃半消息
    else 状态未知或生产者宕机
        B->>P: 回查本地事务状态
        P->>DB: 查询订单/事务表最终状态
        P-->>B: Commit 或 Rollback
    end
  1. 生产者将消息发送至Apache RocketMQ服务端。
  2. Apache RocketMQ服务端将消息持久化成功之后,向生产者返回Ack确认消息已经发送成功,此时消息被标记为"暂不能投递",这种状态下的消息即为半事务消息。
  3. 生产者开始执行本地事务逻辑。
  4. 生产者根据本地事务执行结果向服务端提交二次确认结果(Commit或是Rollback),服务端收到确认结果后处理逻辑如下:
    • 二次确认结果为Commit:服务端将半事务消息标记为可投递,并投递给消费者。
    • 二次确认结果为Rollback:服务端将回滚事务,不会将半事务消息投递给消费者。
  5. 在断网或者是生产者应用重启的特殊情况下,若服务端未收到发送者提交的二次确认结果,或服务端收到的二次确认结果为Unknown未知状态,经过固定时间后,服务端将对消息生产者即生产者集群中任一生产者实例发起消息回查。服务端回查的间隔时间和最大回查次数,请参见参数限制。
  6. 生产者收到消息回查后,需要检查对应消息的本地事务执行的最终结果。
  7. 生产者根据检查到的本地事务的最终状态再次提交二次确认,服务端仍按照步骤4对半事务消息进行处理。

生命周期

生命周期

  • 初始化:半事务消息被生产者构建并完成初始化,待发送到服务端的状态。
  • 事务待提交:半事务消息被发送到服务端,和普通消息不同,并不会直接被服务端持久化,而是会被单独存储到事务存储系统中,等待第二阶段本地事务返回执行结果后再提交。此时消息对下游消费者不可见。
  • 消息回滚:第二阶段如果事务执行结果明确为回滚,服务端会将半事务消息回滚,该事务消息流程终止。
  • 提交待消费:第二阶段如果事务执行结果明确为提交,服务端会将半事务消息重新存储到普通存储系统中,此时消息对下游消费者可见,等待被消费者获取并消费。
  • 消费中:消息被消费者获取,并按照消费者本地的业务逻辑进行处理的过程。 此时服务端会等待消费者完成消费并提交消费结果,如果一定时间后没有收到消费者的响应,Apache RocketMQ会对消息进行重试处理。具体信息,请参见消费重试
  • 消费提交:消费者完成消费处理,并向服务端提交消费结果,服务端标记当前消息已经被处理(包括消费成功和失败)。 Apache RocketMQ默认支持保留所有消息,此时消息数据并不会立即被删除,只是逻辑标记已消费。消息在保存时间到期或存储空间不足被删除前,消费者仍然可以回溯消息重新消费。
  • 消息删除:Apache RocketMQ按照消息保存机制滚动清理最早的消息数据,将消息从物理文件中删除。更多信息,请参见消息存储和清理机制。

使用限制

消息类型一致性

事务消息仅支持在 MessageType 为 Transaction 的主题内使用,即事务消息只能发送至类型为事务消息的主题中,发送的消息的类型必须和主题的类型一致。

消费事务性

Apache RocketMQ事务消息保证本地主分支事务和下游消息发送事务的一致性,但不保证消息消费结果和上游事务的一致性。因此需要下游业务分支自行保证消息正确处理,建议消费端做好消费重试,如果有短暂失败可以利用重试机制保证最终处理成功。

中间状态可见性

Apache RocketMQ事务消息为最终一致性,即在消息提交到下游消费端处理完成之前,下游分支和上游事务之间的状态会不一致。因此,事务消息仅适合接受异步执行的事务场景。

事务超时机制

Apache RocketMQ事务消息的生命周期存在超时机制,即半事务消息被生产者发送服务端后,如果在指定时间内服务端无法确认提交或者回滚状态,则消息默认会被回滚。事务超时时间,请参见参数限制。

代码 Demo:事务消息生产者

下面示例演示“先发送半消息,再执行本地事务,最后提交或回滚消息”的核心流程。

java
TransactionMQProducer producer = new TransactionMQProducer("ORDER_TX_PRODUCER_GROUP");
producer.setNamesrvAddr("127.0.0.1:9876");
producer.setTransactionListener(new TransactionListener() {
    @Override
    public LocalTransactionState executeLocalTransaction(Message msg, Object arg) {
        String orderId = (String) arg;
        try {
            // 这里执行本地事务,例如更新订单为已支付。
            boolean success = payOrder(orderId);
            return success ? LocalTransactionState.COMMIT_MESSAGE
                           : LocalTransactionState.ROLLBACK_MESSAGE;
        } catch (Exception e) {
            return LocalTransactionState.UNKNOW;
        }
    }

    @Override
    public LocalTransactionState checkLocalTransaction(MessageExt msg) {
        String orderId = msg.getKeys();
        // Broker 回查时,根据订单表最终状态决定提交还是回滚。
        return isOrderPaid(orderId)
            ? LocalTransactionState.COMMIT_MESSAGE
            : LocalTransactionState.ROLLBACK_MESSAGE;
    }
});

producer.start();

Message message = new Message(
    "ORDER_EVENT",
    "ORDER_PAID",
    "order-1001",
    "{\"orderId\":\"order-1001\"}".getBytes(StandardCharsets.UTF_8)
);
producer.sendMessageInTransaction(message, "order-1001");

事务消息只能保证“本地事务成功后消息最终可见”。消费端仍然要处理重复消费、消费失败和业务补偿。