日韩性视频-久久久蜜桃-www中文字幕-在线中文字幕av-亚洲欧美一区二区三区四区-撸久久-香蕉视频一区-久久无码精品丰满人妻-国产高潮av-激情福利社-日韩av网址大全-国产精品久久999-日本五十路在线-性欧美在线-久久99精品波多结衣一区-男女午夜免费视频-黑人极品ⅴideos精品欧美棵-人人妻人人澡人人爽精品欧美一区-日韩一区在线看-欧美a级在线免费观看

歡迎訪問 生活随笔!

生活随笔

當前位置: 首頁 > 编程资源 > 编程问答 >内容正文

编程问答

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

發布時間:2024/4/13 编程问答 34 豆豆
生活随笔 收集整理的這篇文章主要介紹了 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.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详解之分布式事务消息(代码实现以及过程分析)的全部內容,希望文章能夠幫你解決所遇到的問題。

如果覺得生活随笔網站內容還不錯,歡迎將生活随笔推薦給好友。