欢迎您访问程序员文章站本站旨在为大家提供分享程序员计算机编程知识!
您现在的位置是: 首页

RocketMQ之分布式事务消息

程序员文章站 2022-07-14 23:34:52
...

 

目录

什么是分布式事务消息

HalfMessage(半消息)

Message Status Check(消息回查)

分布式消息的执行原理

模拟事务正常执行的功能

模拟事务失败执行的功能

模拟消息回查执行的功能


 

什么是分布式事务消息

 

聊什么是事务,最经典的例子就是转账操作,用户A转账给用户B1000元的过程如下:

  • 用户A发起转账请求,用户A账户减去1000元
  • 用户B的账户增加1000元

如果,用户A账户减去1000元后,出现了故障(如网络故障),那么需要将该操作回滚,用户A账户增加1000元。这就是事务。

 

随着项目越来越复杂,越来越服务化,就会导致系统间的事务问题,这个就是分布式事务问题。分布式事务分类有这几种:
基于单个JVM,数据库分库分表了(跨多个数据库)。
基于多JVM,服务拆分了(不跨数据库)。
基于多JVM,服务拆分了 并且数据库分库分表了。

 

 

HalfMessage(半消息)

 

指的是暂不能投递的消息,发送方已经将消息成功发送到了 MQ 服务端,但是服务端未收到生产者对该消息的二次确认,此时该消息被标记成“暂不能投递”状态,处于该种状态下的消息即半消息

 

Message Status Check(消息回查)

 

由于网络闪断、生产者应用重启等原因,导致某条事务消息的二次确认丢失,MQ 服务端通过扫描发现某条消息长期处于“半消息”时,需要主动向消息生产者询问该消息的最终状态(Commit 或是 Rollback),该过程即消息回查。

 

 

分布式消息的执行原理

 

RocketMQ之分布式事务消息

  • 1. 发送方向 MQ 服务端发送消息。
  • 2. MQ Server 将消息持久化成功之后,向发送方 ACK 确认消息已经发送成功,此时消息为半消息。
  • 3. 发送方开始执行本地事务逻辑。
  • 4. 发送方根据本地事务执行结果向 MQ Server 提交二次确认(Commit 或是 Rollback),MQ Server 收到Commit 状态则将半消息标记为可投递,订阅方最终将收到该消息;MQ Server 收到 Rollback 状态则删除半消息,订阅方将不会接受该消息。
  • 5. 在断网或者是应用重启的特殊情况下,上述步骤4提交的二次确认最终未到达 MQ Server,经过固定时间后MQ Server 将对该消息发起消息回查。
  • 6. 发送方收到消息回查后,需要检查对应消息的本地事务执行的最终结果。
  • 7. 发送方根据检查得到的本地事务的最终状态再次提交二次确认,MQ Server 仍按照步骤4对半消息进行操作。 

 

 

模拟事务正常执行的功能

 

新建Producer

package cn.itcast.rocketmq.transaction;

import org.apache.rocketmq.client.producer.TransactionMQProducer;
import org.apache.rocketmq.common.message.Message;

public class TransactionProducer {
    public static void main(String[] args) throws Exception {

        TransactionMQProducer producer = new
                TransactionMQProducer("transaction_producer");
        producer.setNamesrvAddr("192.168.62.132:9876");

        // 设置事务监听器
        producer.setTransactionListener(new TransactionListenerImpl());
        producer.start();

        // 发送消息
        Message message = new Message("pay_topic", "用户A给用户B转账1000元".getBytes("UTF-8"));
        producer.sendMessageInTransaction(message, null);

        Thread.sleep(999999);
        producer.shutdown();
    }
}

 

新建TransactionListenerImpl实现类

这个实现类模拟转账的过程,事务正常运行情况下,只需要把转账的代码用事务的try..catch捕获即可。

package cn.itcast.rocketmq.transaction;

import org.apache.rocketmq.client.producer.LocalTransactionState;
import org.apache.rocketmq.client.producer.TransactionListener;
import org.apache.rocketmq.common.message.Message;
import org.apache.rocketmq.common.message.MessageExt;

import java.util.HashMap;
import java.util.Map;

public class TransactionListenerImpl implements TransactionListener {

    private static Map<String, LocalTransactionState> STATE_MAP = new HashMap<>();

    /**
     * 执行具体的业务逻辑
     *
     * @param msg 发送的消息对象
     * @param arg
     * @return
     */
    @Override
    public LocalTransactionState executeLocalTransaction(Message msg, Object arg) {
        try {
            System.out.println("用户A账户减1000元.");
            Thread.sleep(500); //模拟调用服务

//             System.out.println(1/0);

            System.out.println("用户B账户加1000元.");
            Thread.sleep(800);


            //STATE_MAP.put(msg.getTransactionId(), LocalTransactionState.COMMIT_MESSAGE);

            // 二次提交确认
//            return LocalTransactionState.UNKNOW;
            return LocalTransactionState.COMMIT_MESSAGE;
        } catch (Exception e) {
            e.printStackTrace();
        }

        STATE_MAP.put(msg.getTransactionId(), LocalTransactionState.ROLLBACK_MESSAGE);
        // 回滚
        return LocalTransactionState.ROLLBACK_MESSAGE;
    }

    /**
     * 消息回查
     *
     * @param msg
     * @return
     */
    @Override
    public LocalTransactionState checkLocalTransaction(MessageExt msg) {
        System.out.println("状态回查 ---> " + msg.getTransactionId() +" " +STATE_MAP.get(msg.getTransactionId()) );
        return STATE_MAP.get(msg.getTransactionId());
    }
}

 

新建TransactionConsumer

事务中最省事的就是消费者,因为事务无论成功或失败只需提供给消费者一个状态码就行。

package cn.itcast.rocketmq.transaction;

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.io.UnsupportedEncodingException;
import java.util.List;

public class TransactionConsumer {
    public static void main(String[] args) throws Exception {
        DefaultMQPushConsumer consumer = new
                DefaultMQPushConsumer("HAOKE_CONSUMER");
        consumer.setNamesrvAddr("192.168.62.132:9876");

        // 订阅topic,接收此Topic下的所有消息
        consumer.subscribe("pay_topic", "*");
        consumer.registerMessageListener(new MessageListenerConcurrently() {
            @Override
            public ConsumeConcurrentlyStatus consumeMessage(List<MessageExt> msgs,
                                                            ConsumeConcurrentlyContext context) {
                for (MessageExt msg : msgs) {
                    try {
                        System.out.println(new String(msg.getBody(), "UTF-8"));
                    } catch (UnsupportedEncodingException e) {
                        e.printStackTrace();
                    }
                }
                return ConsumeConcurrentlyStatus.CONSUME_SUCCESS;
            }
        });
        consumer.start();
    }
}

 

测试

测试时我们先运行consumer再运行producter。正常转账,没有问题。

RocketMQ之分布式事务消息

 

 

 

模拟事务失败执行的功能

 

我们改一下TransactionListenerImpl中的代码,在生产者扣钱后,我们模拟一个异常看看消费者能不能收到1000元

RocketMQ之分布式事务消息

运行Producer,可以看到由于事务回滚了,所以消费者并没有收到任何消息

RocketMQ之分布式事务消息

 

 

模拟消息回查执行的功能

 

还记得为什么要回查吗?如果生产者正常执行成功,但是由于未知原因,执行成功并没有通知消息队列,这个时候消息队列和消费者都是处于懵逼状态的。我们通过消息队列的回查机制,当消息队列没有收到生产者的消息时,主动回查生产者的事务是否执行成功,如果成功把commit的状态码发给消费者,如果失败把rollback的状态码发给消费者。并分别对事务的成功和失败做出处理。

 

接下来我们在事务监听实现类中修改代码,注释异常代码(1/0)。事务执行成功时加上一句返回成功的消息码,失败时加上一句返回回滚的消息码。新增消息回查方法,消息码用来查询事务的。

package cn.itcast.rocketmq.transaction;

import org.apache.rocketmq.client.producer.LocalTransactionState;
import org.apache.rocketmq.client.producer.TransactionListener;
import org.apache.rocketmq.common.message.Message;
import org.apache.rocketmq.common.message.MessageExt;

import java.util.HashMap;
import java.util.Map;

public class TransactionListenerImpl implements TransactionListener {

    private static Map<String, LocalTransactionState> STATE_MAP = new HashMap<>();

    /**
     * 执行具体的业务逻辑
     *
     * @param msg 发送的消息对象
     * @param arg
     * @return
     */
    @Override
    public LocalTransactionState executeLocalTransaction(Message msg, Object arg) {
        try {
            System.out.println("用户A账户减1000元.");
            Thread.sleep(500); //模拟调用服务

            //System.out.println(1/0);

            System.out.println("用户B账户加1000元.");
            Thread.sleep(800);

            STATE_MAP.put(msg.getTransactionId(), LocalTransactionState.COMMIT_MESSAGE);

            // 二次提交确认
            return LocalTransactionState.UNKNOW;
            //return LocalTransactionState.COMMIT_MESSAGE;
        } catch (Exception e) {
            e.printStackTrace();
        }

        STATE_MAP.put(msg.getTransactionId(), LocalTransactionState.ROLLBACK_MESSAGE);
        // 回滚
        return LocalTransactionState.ROLLBACK_MESSAGE;
    }

    /**
     * 消息回查
     *
     * @param msg
     * @return
     */
    @Override
    public LocalTransactionState checkLocalTransaction(MessageExt msg) {
        System.out.println("状态回查 ---> " + msg.getTransactionId() +" " +STATE_MAP.get(msg.getTransactionId()) );
        return STATE_MAP.get(msg.getTransactionId());
    }
}

测试

当我们执行生产者代码时,出现的情况是

RocketMQ之分布式事务消息

当状态回查成功后,消费者那里就会出现转账成功信息

RocketMQ之分布式事务消息

 

相关标签: [RocketMQ]