SpringBoot整合RocketMq實(shí)現(xiàn)分布式事務(wù)
大家好,今天我們繼續(xù)分布式事務(wù)的學(xué)習(xí),之前我們已經(jīng)實(shí)戰(zhàn)了
來(lái)實(shí)現(xiàn)分布式事務(wù),今天我們繼續(xù)學(xué)習(xí)springboot整合RocketMq來(lái)實(shí)現(xiàn)分布式事務(wù)
MQ實(shí)現(xiàn)分布式事務(wù)原理
RocketMq提供了事務(wù)消息,要實(shí)現(xiàn)分布式事務(wù)主要還是利用它的事務(wù)消息

- 1:服務(wù)A首先會(huì)發(fā)送一條半事務(wù)的消息到MQ,此時(shí)服務(wù)接收方還是無(wú)法消費(fèi)這條消息的
- 2:半事務(wù)消息發(fā)送成功之后,服務(wù)A開(kāi)始執(zhí)行本地業(yè)務(wù)邏輯
- 3:服務(wù)A執(zhí)行完本地業(yè)務(wù)之后,提交事務(wù),事務(wù)提交成功之后,MQ這條半事務(wù)消息就會(huì)變成原始可消費(fèi)的消息
- 4:服務(wù)接收方這時(shí)候就可以消費(fèi)到這條消息,繼續(xù)執(zhí)行后續(xù)業(yè)務(wù)了
那么在這整個(gè)過(guò)程中可能會(huì)出現(xiàn)的異常有哪些呢?
1:半事務(wù)消息發(fā)送失敗
如果半事務(wù)消息發(fā)送失敗,那么服務(wù)A就不會(huì)繼續(xù)執(zhí)行接下來(lái)的業(yè)務(wù)了,整個(gè)流程會(huì)直接退出
2:本地事務(wù)提交成功,發(fā)送COMMIT消息失敗
服務(wù)A事務(wù)提交之后,需要發(fā)送一條消息告訴MQ,這條半事務(wù)消息可以消費(fèi)了,但是這時(shí)候,COMMIT消息發(fā)送失敗了,那么這條消息就還是處于半事務(wù)狀態(tài),所以MQ會(huì)進(jìn)行回查
回查服務(wù)A這個(gè)事務(wù)是否成功了,如果成功了,就會(huì)發(fā)送回查結(jié)果,如果本地事務(wù)成功了,那么回查就會(huì)發(fā)送COMMIT消息,這條消息重新設(shè)置為可消費(fèi)狀態(tài)
實(shí)戰(zhàn)
服務(wù)-A
@Transactional(rollbackFor = Exception.class)
@Override
public String mqInsert(Test test) {
//本地服務(wù)調(diào)用
testDao.insert(test);
//發(fā)送半事務(wù)消息
//Test是我們本地需要保存的一個(gè)對(duì)象
Message<String> message = MessageBuilder.withPayload(JSONObject.toJSONString(test)).build();
rocketMQTemplate.sendMessageInTransaction("test-topic", message, null);
return "success";
}
@RocketMQTransactionListener
public class TransactionMqListener implements RocketMQLocalTransactionListener {
@Resource
private TestDao testDao;
//執(zhí)行本地事務(wù)
@Override
public RocketMQLocalTransactionState executeLocalTransaction(Message message, Object o) {
//獲取半事務(wù)消息
Test test = JSONObject.parseObject(new String((byte[]) message.getPayload()), Test.class);
System.out.println("test | executeLocalTransaction | 消息是:" + JSONObject.toJSONString(test));
//根據(jù)id查詢?cè)撚涗浭欠癖4娉晒α?
Test testExist = testDao.queryById(test.getId());
if(testExist == null) {
//說(shuō)明本地事務(wù)提交失敗了,需要回滾
return RocketMQLocalTransactionState.ROLLBACK;
}
//本地事務(wù)提交成功
return RocketMQLocalTransactionState.COMMIT;
}
//事務(wù)回查
@Override
public RocketMQLocalTransactionState checkLocalTransaction(Message message) {
Test test = JSONObject.parseObject(new String((byte[]) message.getPayload()), Test.class);
System.out.println("test | checkLocalTransaction | 消息是:" + JSONObject.toJSONString(test));
//還是根據(jù)id去查詢記錄
Test testExist = testDao.queryById(test.getId());
if(testExist == null) {
//不存在,說(shuō)明本地事務(wù)提交失敗,回滾
return RocketMQLocalTransactionState.ROLLBACK;
}
//本地事務(wù)提交成功了
return RocketMQLocalTransactionState.COMMIT;
}
}
服務(wù)B
@Component
@RocketMQMessageListener(topic = "test-topic",consumerGroup = "cpy-consumer-group")
public class CpyListener implements RocketMQListener<String> {
//消息監(jiān)聽(tīng)
@Override
public void onMessage(String s) {
System.out.println("cpy服務(wù)收到消息:" + JSONObject.toJSONString(s));
}
}
測(cè)試

我們先把這里的狀態(tài)改成 UNKNOWN,來(lái)模擬本地事務(wù)提交失敗的場(chǎng)景,來(lái)驗(yàn)證事務(wù)回查的效果
發(fā)送事務(wù)消息之后,我們這里是UNKNOWN狀態(tài),所以沒(méi)有提交成功

此時(shí)服務(wù)-B也沒(méi)有消費(fèi)到消息

過(guò)了一會(huì),MQ事務(wù)消息進(jìn)行回查,此時(shí)因?yàn)閿?shù)據(jù)庫(kù)已經(jīng)存在這條記錄了,所以直接COMMIT

這時(shí)候服務(wù)消費(fèi)方也成功消費(fèi)到消息了

到此這篇關(guān)于SpringBoot整合RocketMq實(shí)現(xiàn)分布式事務(wù)的文章就介紹到這了,更多相關(guān)SpringBoot RocketMq分布式事務(wù)內(nèi)容請(qǐng)搜索腳本之家以前的文章或繼續(xù)瀏覽下面的相關(guān)文章希望大家以后多多支持腳本之家!
相關(guān)文章
Java使用正則表達(dá)式進(jìn)行匹配且對(duì)匹配結(jié)果逐個(gè)替換
這篇文章主要介紹了Java使用正則表達(dá)式進(jìn)行匹配且對(duì)匹配結(jié)果逐個(gè)替換,文章圍繞主題展開(kāi)詳細(xì)的內(nèi)容戒殺,具有一定的參考價(jià)值,需要的小伙伴可以參考一下2022-09-09
使用feign傳遞參數(shù)類型為MultipartFile的問(wèn)題
這篇文章主要介紹了使用feign傳遞參數(shù)類型為MultipartFile的問(wèn)題,具有很好的參考價(jià)值,希望對(duì)大家有所幫助。如有錯(cuò)誤或未考慮完全的地方,望不吝賜教2022-03-03
SpringBoot實(shí)現(xiàn)MQTT消息發(fā)送和接收方式
這篇文章主要介紹了SpringBoot實(shí)現(xiàn)MQTT消息發(fā)送和接收方式,具有很好的參考價(jià)值,希望對(duì)大家有所幫助。如有錯(cuò)誤或未考慮完全的地方,望不吝賜教2023-03-03
詳解SpringBoot中異步請(qǐng)求和異步調(diào)用(看完這一篇就夠了)
這篇文章主要介紹了SpringBoot中異步請(qǐng)求和異步調(diào)用問(wèn)題,非常不錯(cuò),具有一定的參考借鑒價(jià)值,需要的朋友可以參考下2019-04-04
Mybatis?MappedStatement類核心原理詳解
這篇文章主要介紹了Mybatis?MappedStatement類,mybatis的mapper文件最終會(huì)被解析器,解析成MappedStatement,其中insert|update|delete|select每一個(gè)標(biāo)簽分別對(duì)應(yīng)一個(gè)MappedStatement2022-11-11
Java基礎(chǔ)MAC系統(tǒng)下IDEA連接MYSQL數(shù)據(jù)庫(kù)JDBC過(guò)程
最近一直在學(xué)習(xí)web項(xiàng)目,當(dāng)然也會(huì)涉及與數(shù)據(jù)庫(kù)的連接這塊,這里就總結(jié)一下在IDEA中如何進(jìn)行MySQL數(shù)據(jù)庫(kù)的連接,這里提一下我的電腦是MAC系統(tǒng),使用的編碼軟件是IDEA,數(shù)據(jù)庫(kù)是MySQL2021-09-09
你所不知道的Spring的@Autowired實(shí)現(xiàn)細(xì)節(jié)分析
這篇文章主要介紹了你所不知道的Spring的@Autowired實(shí)現(xiàn)細(xì)節(jié)分析,具有很好的參考價(jià)值,希望對(duì)大家有所幫助。一起跟隨小編過(guò)來(lái)看看吧2020-08-08
Spring?Boot整合Bootstrap的超詳細(xì)步驟
之前做前端開(kāi)發(fā),在使用bootstrap的時(shí)候都是去官網(wǎng)下載,然后放到項(xiàng)目中,在頁(yè)面引用,下面這篇文章主要給大家介紹了關(guān)于Spring?Boot整合Bootstrap的超詳細(xì)步驟,需要的朋友可以參考下2023-05-05

