背景
在一個(gè)微服務(wù)架構(gòu)的項(xiàng)目中,一個(gè)業(yè)務(wù)操作可能涉及到多個(gè)服務(wù),這些服務(wù)往往是獨(dú)立部署,構(gòu)成一個(gè)個(gè)獨(dú)立的系統(tǒng)。這種分布式的系統(tǒng)架構(gòu)往往面臨著分布式事務(wù)的問題。為了保證系統(tǒng)數(shù)據(jù)的一致性,我們需要確保這些服務(wù)中的操作要么全部成功,要么全部失敗。通過使用RocketMQ實(shí)現(xiàn)分布式事務(wù),我們可以協(xié)調(diào)這些服務(wù)的操作,保證數(shù)據(jù)的一致性。
功能原理
RocketMQ的分布式事務(wù)消息功能,在普通消息基礎(chǔ)上,支持二階段的提交。將二階段提交和本地事務(wù)綁定,實(shí)現(xiàn)全局提交結(jié)果的一致性。
整個(gè)事務(wù)消息的詳細(xì)交互流程如下圖所示:
1、生產(chǎn)者將消息發(fā)送至RocketMQ服務(wù)端。
2、RocketMQ服務(wù)端將消息持久化成功之后,向生產(chǎn)者返回Ack確認(rèn)消息已經(jīng)發(fā)送成功,此時(shí)消息被標(biāo)記為"暫不能投遞",這種狀態(tài)下的消息即為半事務(wù)消息。
3、生產(chǎn)者開始執(zhí)行本地事務(wù)邏輯。
4、生產(chǎn)者根據(jù)本地事務(wù)執(zhí)行結(jié)果向服務(wù)端提交二次確認(rèn)結(jié)果(Commit或是Rollback),服務(wù)端收到確認(rèn)結(jié)果后處理邏輯如下:
-
二次確認(rèn)結(jié)果為Commit:服務(wù)端將半事務(wù)消息標(biāo)記為可投遞,并投遞給消費(fèi)者。
-
二次確認(rèn)結(jié)果為Rollback:服務(wù)端將回滾事務(wù),不會(huì)將半事務(wù)消息投遞給消費(fèi)者。
5、在斷網(wǎng)或者是生產(chǎn)者應(yīng)用重啟的特殊情況下,若服務(wù)端未收到生產(chǎn)者提交的二次確認(rèn)結(jié)果,或服務(wù)端收到的二次確認(rèn)結(jié)果為Unknown未知狀態(tài),經(jīng)過固定時(shí)間后,服務(wù)端將對(duì)消息生產(chǎn)者集群中任一生產(chǎn)者實(shí)例發(fā)起消息回查。
6、生產(chǎn)者收到消息回查后,需要檢查對(duì)應(yīng)消息的本地事務(wù)執(zhí)行的最終結(jié)果。
7、生產(chǎn)者根據(jù)檢查到的本地事務(wù)的最終狀態(tài)再次提交二次確認(rèn),服務(wù)端仍按照步驟4對(duì)半事務(wù)消息進(jìn)行處理。
注意問題
消息類型
事務(wù)消息僅支持在MessageType為Transaction的主題使用,即事務(wù)消息只能發(fā)送至類型為事務(wù)消息的主題中。
消息消費(fèi)
RocketMQ事務(wù)消息保證生產(chǎn)者本地事務(wù)和下游消息發(fā)送事務(wù)的一致性,但不保證消息消費(fèi)結(jié)果和上游事務(wù)的一致性。因此需要下游業(yè)務(wù)自行保證消息正確處理,建議消費(fèi)端做好消費(fèi)重試。
中間狀態(tài)
RocketMQ事務(wù)消息一致性為最終一致性,即在消息提交到下游消費(fèi)端處理完成之前,下游和上游事務(wù)之間的狀態(tài)會(huì)不一致。因此,事務(wù)消息僅適合能接受異步執(zhí)行的場(chǎng)景。
事務(wù)超時(shí)
RocketMQ事務(wù)消息的生命周期存在超時(shí)機(jī)制,即半事務(wù)消息被生產(chǎn)者發(fā)送服務(wù)端后,如果在指定時(shí)間內(nèi)服務(wù)端無法確認(rèn)提交或者回滾狀態(tài),則消息默認(rèn)會(huì)被回滾。
示例代碼
以下為RocketMQ 4.x版本事務(wù)消息示例代碼,
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)建事務(wù)消息生產(chǎn)者
TransactionMQProducer producer = new TransactionMQProducer("transaction_producer");
producer.setNamesrvAddr("127.0.0.1:9876");
// 設(shè)置事務(wù)監(jiān)聽器
TransactionListener transactionListener = new MyTransactionListener();
producer.setTransactionListener(transactionListener);
// 設(shè)置事務(wù)回查的線程池,可以不必設(shè)置,如果不設(shè)置也會(huì)默認(rèn)生成一個(gè)
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);
// 啟動(dòng)生產(chǎn)者
producer.start();
// 發(fā)送事務(wù)消息
Message message = new Message("transaction_topic", "test_tag", "test_key", "Hello RocketMQ".getBytes());
producer.sendMessageInTransaction(message, null);
// 關(guān)閉生產(chǎn)者
producer.shutdown();
}
}
/**
* 事務(wù)監(jiān)聽器
*/
class MyTransactionListener implements TransactionListener {
@Override
public LocalTransactionState executeLocalTransaction(Message msg, Object arg) {
// 執(zhí)行本地事務(wù)操作
System.out.println("執(zhí)行本地事務(wù)操作,消息內(nèi)容:" + new String(msg.getBody()));
return LocalTransactionState.COMMIT_MESSAGE; // 提交事務(wù),允許消費(fèi)者消費(fèi)該消息
// return LocalTransactionState.ROLLBACK_MESSAGE;// 回滾事務(wù),消息將被丟棄不允許消費(fèi)。
// return LocalTransactionState.UNKNOW;// 暫時(shí)無法判斷狀態(tài),等待固定時(shí)間以后Broker端根據(jù)回查規(guī)則向生產(chǎn)者進(jìn)行消息回查。
}
@Override
public LocalTransactionState checkLocalTransaction(MessageExt msg) {
// 檢查本地事務(wù)狀態(tài)
System.out.println("檢查本地事務(wù)狀態(tài),消息內(nèi)容:" + new String(msg.getBody()));
return LocalTransactionState.COMMIT_MESSAGE;
}
}
代碼解釋:
1、事務(wù)消息的生產(chǎn)者使用TransactionMQProducer
創(chuàng)建。
2、MyTransactionListener
作為事務(wù)監(jiān)聽器,實(shí)現(xiàn)了接口TransactionListener
,該接口有兩個(gè)方法,分別是:
-
executeLocalTransaction
:
半事務(wù)消息發(fā)送成功后,執(zhí)行本地事務(wù)的方法,具體執(zhí)行完本地事務(wù)后,可以在該方法中返回以下三種狀態(tài):
LocalTransactionState.COMMIT_MESSAGE: 提交事務(wù),允許消費(fèi)者消費(fèi)該消息。
LocalTransactionState.ROLLBACK_MESSAGE: 回滾事務(wù),消息將被丟棄不允許消費(fèi)。
LocalTransactionState.UNKNOW: 暫時(shí)無法判斷狀態(tài),等待固定時(shí)間以后RocketMQ服務(wù)端根據(jù)回查規(guī)則向生產(chǎn)者進(jìn)行消息回查。文章來源:http://www.zghlxwxcb.cn/news/detail-838864.html -
checkLocalTransaction
:
二次確認(rèn)消息沒有收到,RocketMQ服務(wù)端回查生產(chǎn)者端事務(wù)結(jié)果的方法?;夭橐?guī)則:本地事務(wù)執(zhí)行完成后,若RocketMQ服務(wù)端收到的本地事務(wù)返回狀態(tài)為L(zhǎng)ocalTransactionState.UNKNOW,或生產(chǎn)者應(yīng)用退出導(dǎo)致本地事務(wù)未提交任何狀態(tài)。則RocketMQ服務(wù)端會(huì)向消息生產(chǎn)者發(fā)起事務(wù)回查,第一次回查后仍未獲取到事務(wù)狀態(tài),則之后每隔一段時(shí)間會(huì)再次回查。文章來源地址http://www.zghlxwxcb.cn/news/detail-838864.html
到了這里,關(guān)于基于RocketMQ實(shí)現(xiàn)分布式事務(wù)的文章就介紹完了。如果您還想了解更多內(nèi)容,請(qǐng)?jiān)谟疑辖撬阉鱐OY模板網(wǎng)以前的文章或繼續(xù)瀏覽下面的相關(guān)文章,希望大家以后多多支持TOY模板網(wǎng)!