200字范文,内容丰富有趣,生活中的好帮手!
200字范文 > RocketMQ的Producer详解之分布式事务消息(代码实现以及过程分析)

RocketMQ的Producer详解之分布式事务消息(代码实现以及过程分析)

时间:2021-01-12 17:50:07

相关推荐

RocketMQ的Producer详解之分布式事务消息(代码实现以及过程分析)

执行流程

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对半消息进行操作。

package cn.learn.rocketmq.transaction;import org.apache.rocketmq.client.producer.TransactionMQProducer;import org.mon.message.Message;public class TransactionProducer {public static void main(String[] args) throws Exception {TransactionMQProducer producer = newTransactionMQProducer("transaction_producer");producer.setNamesrvAddr("localhost:9876");// 设置事务监听器producer.setTransactionListener(new TransactionListenerImpl());producer.start();// 发送消息Message message = new Message("pay_topic", "用户A给用户B转账500元".getBytes("UTF-8"));producer.sendMessageInTransaction(message, null);Thread.sleep(999999);producer.shutdown();}}

package cn.learn.rocketmq.transaction;import org.apache.rocketmq.client.producer.LocalTransactionState;import org.apache.rocketmq.client.producer.TransactionListener;import org.mon.message.Message;import org.mon.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*/@Overridepublic LocalTransactionState executeLocalTransaction(Message msg, Object arg) {try {System.out.println("用户A账户减500元.");Thread.sleep(500); //模拟调用服务// System.out.println(1/0);System.out.println("用户B账户加500元.");Thread.sleep(800);STATE_MAP.put(msg.getTransactionId(), MIT_MESSAGE);// 二次提交确认// return LocalTransactionState.UNKNOW;return MIT_MESSAGE;} catch (Exception e) {e.printStackTrace();}STATE_MAP.put(msg.getTransactionId(), LocalTransactionState.ROLLBACK_MESSAGE);// 回滚return LocalTransactionState.ROLLBACK_MESSAGE;}/*** 消息回查** @param msg* @return*/@Overridepublic LocalTransactionState checkLocalTransaction(MessageExt msg) {System.out.println("状态回查 ---> " + msg.getTransactionId() +" " +STATE_MAP.get(msg.getTransactionId()) );return STATE_MAP.get(msg.getTransactionId());}}

package cn.learn.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.mon.message.MessageExt;import java.io.UnsupportedEncodingException;import java.util.List;public class TransactionConsumer {public static void main(String[] args) throws Exception {DefaultMQPushConsumer consumer = newDefaultMQPushConsumer("LEARN_CONSUMER");consumer.setNamesrvAddr("localhost:9876");// 订阅topic,接收此Topic下的所有消息consumer.subscribe("pay_topic", "*");consumer.registerMessageListener(new MessageListenerConcurrently() {@Overridepublic 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();}}

本内容不代表本网观点和政治立场,如有侵犯你的权益请联系我们处理。
网友评论
网友评论仅供其表达个人看法,并不表明网站立场。