执行流程

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.apache.rocketmq.common.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.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*/@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(), 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*/@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.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 = 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();}
}

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

  1. RocketMQ的Producer详解之分布式事务消息(回顾事务)

    回顾什么事务 聊什么是事务,最经典的例子就是转账操作,用户A转账给用户B1000元的过程如下: 用户A发起转账请求,用户A账户减去1000元 用户B的账户增加1000元 如果,用户A账户减去1000元 ...

  2. RocketMQ的Producer详解之分布式事务消息(原理分析)

    Half(Prepare) Message 指的是暂不能投递的消息,发送方已经将消息成功发送到了 MQ 服务端,但是服务端未收到生产者对该消息的二次确认,此时该消息被标记成"暂不能投递&qu ...

  3. 分布式事务详解【分布式事务的几种解决方案】彻底搞懂分布式事务

    文章目录 一.基本概念 什么是事务 本地事务 分布式事务 分布式事务产生的场景 二.分布式事务基础理论 CAP理论 CP - Consistency/Partition Tolerance AP - ...

  4. 详解Mysql分布式事务XA

    在开发中,为了降低单点压力,通常会根据业务情况进行分表分库,将表分布在不同的库中(库可能分布在不同的机器上).在这种场景下,事务的提交会变得相对复杂,因为多个节点(库)的存在,可能存在部分节点提交失败 ...

  5. 一文详解,分布式事务Seata

    事务ACID原则 原子性:事务中的所有操作,要么全部成功,要么全部失败一致性:要保证数据库内部完整性约束.声明性约束隔离性:对同一资源操作的事务不能同时发生持久性:对数据库做的一切修改将永久保存,不管 ...

  6. RocketMQ的Producer详解之顺序消息(原理)

    顺序消息 在某些业务中,consumer在消费消息时,是需要按照生产者发送消息的顺序进行消费的,比如在电商系统中,订单的消息,会有创建订单.订单支付.订单完成,如果消息的顺序发生改变,那么这样的消息就 ...

  7. Apache RocketMQ 正式开源分布式事务消息

    摘要: 近日,Apache RocketMQ 社区正式发布4.3版本.此次发布不仅包括提升性能,减少内存使用等原有特性增强,还修复了部分社区提出的若干问题,更重要的是该版本开源了社区最为关心的分布式事 ...

  8. 消息中间件学习总结(15)——Apache RocketMQ 正式开源分布式事务消息

    近日,Apache RocketMQ 社区正式发布4.3版本.此次发布不仅包括提升性能,减少内存使用等原有特性增强,还修复了部分社区提出的若干问题,更重要的是该版本开源了社区最为关心的分布式事务消息, ...

  9. Session机制详解及分布式中Session共享解决方案

    Session机制详解及分布式中Session共享解决方案 参考文章: (1)Session机制详解及分布式中Session共享解决方案 (2)https://www.cnblogs.com/jing ...

最新文章

  1. python 归一化_只需 45 秒,Python 给故宫画一组手绘图!
  2. 还款每个月90.85元, 到 2012年10月,2012 11月 2256元,共 5799.25元
  3. SPY++ 学习总结
  4. pcm 采样率转换_44.1KHz够用吗?我们是否需要更高的采样率?
  5. [云炬创业基础笔记]第四章测试21
  6. 维沃手机有没有智能机器人_权威发布!2019世界智能移动终端产业高峰会议获奖名单...
  7. 采用计算机辅助电话调查,计算机辅助电话调查(CATI)-实验.pdf
  8. 手机modem开发(28)---开发电信VoLTE开关默认值设置
  9. 利用nginx重写url参数并跳转
  10. Tcpip详解卷一第3章(2)
  11. Revel敏捷后台开发框架
  12. 使用plupload压缩图片
  13. iOS14:AirPods Auto Switching
  14. 如何求地球上两点之间的最短距离_初中数学求线段之和最小的问题,知识点题型汇总...
  15. mysql int 11 最大多少_MySQL中int(11)最大长度是多少?
  16. 【新媒体 | 自媒体 运营】虚拟素材(图片,字体,音频,视频)商用及CC版权相关问题
  17. 历年全国计算机技术与软件专业资格(水平)考试真题及答案汇总
  18. 妖怪屋 服务器维护中,《阴阳师:妖怪屋》3月3日维护更新公告
  19. sofa框架server-client搭建
  20. android 华为 多语言,原来华为手机自带翻译神器!这3个方法,一键实现多国语言翻译...

热门文章

  1. sqlserver数据库事务
  2. bzoj1059: [ZJOI2007]矩阵游戏
  3. 类与类集合的基本使用
  4. 《大道至简》第二章 读后感
  5. JQuery:JQuery添加元素
  6. 二:java语法基础:
  7. 判断线段相交(hdu1558 Segment set 线段相交+并查集)
  8. POJ 1273 (基础最大流) Drainage Ditches
  9. 网址收藏 plc实现
  10. MySQL数据库引擎详解