基于RocketMQ實現(xiàn)分布式事務的方法
背景
在一個微服務架構的項目中,一個業(yè)務操作可能涉及到多個服務,這些服務往往是獨立部署,構成一個個獨立的系統(tǒng)。這種分布式的系統(tǒng)架構往往面臨著分布式事務的問題。為了保證系統(tǒng)數(shù)據(jù)的一致性,我們需要確保這些服務中的操作要么全部成功,要么全部失敗。通過使用RocketMQ實現(xiàn)分布式事務,我們可以協(xié)調這些服務的操作,保證數(shù)據(jù)的一致性。
功能原理
RocketMQ的分布式事務消息功能,在普通消息基礎上,支持二階段的提交。將二階段提交和本地事務綁定,實現(xiàn)全局提交結果的一致性。
整個事務消息的詳細交互流程如下圖所示:
1、生產(chǎn)者將消息發(fā)送至RocketMQ服務端。
2、RocketMQ服務端將消息持久化成功之后,向生產(chǎn)者返回Ack確認消息已經(jīng)發(fā)送成功,此時消息被標記為"暫不能投遞",這種狀態(tài)下的消息即為半事務消息。
3、生產(chǎn)者開始執(zhí)行本地事務邏輯。
4、生產(chǎn)者根據(jù)本地事務執(zhí)行結果向服務端提交二次確認結果(Commit或是Rollback),服務端收到確認結果后處理邏輯如下:
二次確認結果為Commit:服務端將半事務消息標記為可投遞,并投遞給消費者。
二次確認結果為Rollback:服務端將回滾事務,不會將半事務消息投遞給消費者。
5、在斷網(wǎng)或者是生產(chǎn)者應用重啟的特殊情況下,若服務端未收到生產(chǎn)者提交的二次確認結果,或服務端收到的二次確認結果為Unknown未知狀態(tài),經(jīng)過固定時間后,服務端將對消息生產(chǎn)者集群中任一生產(chǎn)者實例發(fā)起消息回查。
6、生產(chǎn)者收到消息回查后,需要檢查對應消息的本地事務執(zhí)行的最終結果。
7、生產(chǎn)者根據(jù)檢查到的本地事務的最終狀態(tài)再次提交二次確認,服務端仍按照步驟4對半事務消息進行處理。
注意問題
消息類型事務消息僅支持在MessageType為Transaction的主題使用,即事務消息只能發(fā)送至類型為事務消息的主題中。
消息消費RocketMQ事務消息保證生產(chǎn)者本地事務和下游消息發(fā)送事務的一致性,但不保證消息消費結果和上游事務的一致性。因此需要下游業(yè)務自行保證消息正確處理,建議消費端做好消費重試。
中間狀態(tài)RocketMQ事務消息一致性為最終一致性,即在消息提交到下游消費端處理完成之前,下游和上游事務之間的狀態(tài)會不一致。因此,事務消息僅適合能接受異步執(zhí)行的場景。
事務超時RocketMQ事務消息的生命周期存在超時機制,即半事務消息被生產(chǎn)者發(fā)送服務端后,如果在指定時間內(nèi)服務端無法確認提交或者回滾狀態(tài),則消息默認會被回滾。
示例代碼
以下為RocketMQ 4.x版本事務消息示例代碼,
import org.apache.rocketmq.client.producer.LocalTransactionState; import org.apache.rocketmq.client.producer.TransactionListener; import org.apache.rocketmq.client.producer.TransactionMQProducer; import org.apache.rocketmq.common.message.Message; import org.apache.rocketmq.common.message.MessageExt; import java.util.concurrent.*; public class RocketMqTransactionDemo { public static void main(String[] args) throws Exception { // 創(chuàng)建事務消息生產(chǎn)者 TransactionMQProducer producer = new TransactionMQProducer("transaction_producer"); producer.setNamesrvAddr("127.0.0.1:9876"); // 設置事務監(jiān)聽器 TransactionListener transactionListener = new MyTransactionListener(); producer.setTransactionListener(transactionListener); // 設置事務回查的線程池,可以不必設置,如果不設置也會默認生成一個 ExecutorService executorService = new ThreadPoolExecutor(2, 5, 100, TimeUnit.SECONDS, new ArrayBlockingQueue <Runnable> (2000), new ThreadFactory() { @Override public Thread newThread(Runnable r) { Thread thread = new Thread(r); thread.setName("client-transaction-msg-check-thread"); return thread; } }); producer.setExecutorService(executorService); // 啟動生產(chǎn)者 producer.start(); // 發(fā)送事務消息 Message message = new Message("transaction_topic", "test_tag", "test_key", "Hello RocketMQ".getBytes()); producer.sendMessageInTransaction(message, null); // 關閉生產(chǎn)者 producer.shutdown(); } } /** * 事務監(jiān)聽器 */ class MyTransactionListener implements TransactionListener { @Override public LocalTransactionState executeLocalTransaction(Message msg, Object arg) { // 執(zhí)行本地事務操作 System.out.println("執(zhí)行本地事務操作,消息內(nèi)容:" + new String(msg.getBody())); return LocalTransactionState.COMMIT_MESSAGE; // 提交事務,允許消費者消費該消息 // return LocalTransactionState.ROLLBACK_MESSAGE;// 回滾事務,消息將被丟棄不允許消費。 // return LocalTransactionState.UNKNOW;// 暫時無法判斷狀態(tài),等待固定時間以后Broker端根據(jù)回查規(guī)則向生產(chǎn)者進行消息回查。 } @Override public LocalTransactionState checkLocalTransaction(MessageExt msg) { // 檢查本地事務狀態(tài) System.out.println("檢查本地事務狀態(tài),消息內(nèi)容:" + new String(msg.getBody())); return LocalTransactionState.COMMIT_MESSAGE; } }
代碼解釋:
1、事務消息的生產(chǎn)者使用TransactionMQProducer
創(chuàng)建。
2、MyTransactionListener
作為事務監(jiān)聽器,實現(xiàn)了接口TransactionListener
,該接口有兩個方法,分別是:
executeLocalTransaction
:
半事務消息發(fā)送成功后,執(zhí)行本地事務的方法,具體執(zhí)行完本地事務后,可以在該方法中返回以下三種狀態(tài):
LocalTransactionState.COMMIT_MESSAGE: 提交事務,允許消費者消費該消息。
LocalTransactionState.ROLLBACK_MESSAGE: 回滾事務,消息將被丟棄不允許消費。
LocalTransactionState.UNKNOW: 暫時無法判斷狀態(tài),等待固定時間以后RocketMQ服務端根據(jù)回查規(guī)則向生產(chǎn)者進行消息回查。checkLocalTransaction
:
二次確認消息沒有收到,RocketMQ服務端回查生產(chǎn)者端事務結果的方法?;夭橐?guī)則:本地事務執(zhí)行完成后,若RocketMQ服務端收到的本地事務返回狀態(tài)為LocalTransactionState.UNKNOW,或生產(chǎn)者應用退出導致本地事務未提交任何狀態(tài)。則RocketMQ服務端會向消息生產(chǎn)者發(fā)起事務回查,第一次回查后仍未獲取到事務狀態(tài),則之后每隔一段時間會再次回查。
到此這篇關于基于RocketMQ實現(xiàn)分布式事務的文章就介紹到這了,更多相關RocketMQ分布式事務內(nèi)容請搜索腳本之家以前的文章或繼續(xù)瀏覽下面的相關文章希望大家以后多多支持腳本之家!
相關文章
分析java并發(fā)中的wait notify notifyAll
一個線程修改一個對象的值,而另一個線程則感知到了變化,然后進行相應的操作,這就是wait()、notify()和notifyAll()方法的本質。本文將詳細來介紹它們概念實現(xiàn)以及區(qū)別2021-06-06如何使用Spring?integration在Springboot中集成Mqtt詳解
MQTT是多個客戶端通過一個中央服務器傳遞信息的多對多協(xié)議,能高效地將信息分發(fā)給一個或多個訂閱者,下面這篇文章主要給大家介紹了關于如何使用Spring?integration在Springboot中集成Mqtt的相關資料,需要的朋友可以參考下2023-02-02解決新版 Idea 中 SpringBoot 熱部署不生效的問題
這篇文章主要介紹了解決新版 Idea 中 SpringBoot 熱部署不生效的問題,本文通過圖文并茂的形式給大家介紹的非常詳細,對大家的學習或工作具有一定的參考借鑒價值,需要的朋友可以參考下2023-08-08springboot中validator數(shù)據(jù)校驗功能的實現(xiàn)
這篇文章主要介紹了springboot中validator數(shù)據(jù)校驗功能,校驗分為普通校驗和分組校驗,每種校驗方式通過實例代碼給大家介紹的非常詳細,需要的朋友可以參考下2021-10-10Intellij IDEA中啟動多個微服務(開啟Run Dashboard管理)
這篇文章主要介紹了Intellij IDEA中啟動多個微服務(開啟Run Dashboard管理),文中通過示例代碼介紹的非常詳細,對大家的學習或者工作具有一定的參考學習價值,需要的朋友們下面隨著小編來一起學習學習吧2020-07-07