RocketMQ事務消息圖文示例講解
更新時間:2022年12月27日 16:10:11 作者:一個雙子座的Java攻城獅
RocketMQ事務消息(Transactional Message)是指應用本地事務和發(fā)送消息操作可以被定義到全局事務中,要么同時成功,要么同時失敗。RocketMQ的事務消息提供類似 X/Open XA 的分布式事務功能,通過事務消息能達到分布式事務的最終一致
RocketMQ 也允許我們像mysql 一樣發(fā)送具有事務特征的消息
MQ 的事務流程(本地代碼正常執(zhí)行)

MQ 的消息補償過程(當本地代碼執(zhí)行失敗時)

MQ 消息的三種狀態(tài)
- 提交狀態(tài):允許進入隊列,此消息與非事務消息無區(qū)別
- 回滾狀態(tài):不允許進入隊列,此消息等同于未發(fā)送過
- 中間狀態(tài):完成了 half 消息的發(fā)送,未對 MQ 進行二次狀態(tài)確認(未知狀態(tài))
注意:事務消息僅與生產者有關,與消費者無關
生產者代碼(提交狀態(tài)、回滾狀態(tài)):
public class Producer {
public static void main(String[] args) throws Exception{
//事務消息使用的生產者是TransactionMQProducer
TransactionMQProducer producer = new TransactionMQProducer("group1");
producer.setNamesrvAddr("192.168.23.127:9876");
//添加本地事務對應的監(jiān)聽
producer.setTransactionListener(new TransactionListener() {
//正常事務過程
@Override
public LocalTransactionState executeLocalTransaction(Message message, Object o) {
// 此處寫本地事務處理業(yè)務
// 如果成功,消息改為提交,如果失敗改為 回滾,如果是多線程處理狀態(tài)未知,就提交為未知等待事務補償過程
//事務提交狀態(tài)
return LocalTransactionState.COMMIT_MESSAGE;// 類似于msql 的 commit
//return LocalTransactionState.ROLLBACK_MESSAGE;回滾狀態(tài)
}
//事務補償過程
@Override
public LocalTransactionState checkLocalTransaction(MessageExt messageExt) {
return null;
}
});
producer.start();
Message msg = new Message("topic8",("事務消息:hello rocketmq ").getBytes("UTF-8"));
SendResult result = producer.sendMessageInTransaction(msg,null);
System.out.println("返回結果:"+result);
producer.shutdown();
}
}
生產者(中間狀態(tài)):
public class Producer {
public static void main(String[] args) throws Exception{
//事務消息使用的生產者是TransactionMQProducer
TransactionMQProducer producer = new TransactionMQProducer("group1");
producer.setNamesrvAddr("192.168.23.127:9876");
//添加本地事務對應的監(jiān)聽
producer.setTransactionListener(new TransactionListener() {
//正常事務過程
@Override
public LocalTransactionState executeLocalTransaction(Message message, Object o) {
return LocalTransactionState.UNKNOW;
}
//事務補償過程
@Override
public LocalTransactionState checkLocalTransaction(MessageExt messageExt) {
System.out.println("事務補償過程執(zhí)行");
return LocalTransactionState.COMMIT_MESSAGE;
}
});
producer.start();
Message msg = new Message("topic8",("事務消息:hello rocketmq ").getBytes("UTF-8"));
SendResult result = producer.sendMessageInTransaction(msg,null);
System.out.println("返回結果:"+result);
//事務補償過程必須保障服務器在運行過程中,否則將無法進行正常的事務補償
//producer.shutdown();
}
}
到此這篇關于RocketMQ事務消息圖文示例講解的文章就介紹到這了,更多相關RocketMQ事務消息內容請搜索腳本之家以前的文章或繼續(xù)瀏覽下面的相關文章希望大家以后多多支持腳本之家!
相關文章
java發(fā)送post請求使用multipart/form-data格式文件數據到接口代碼示例
這篇文章主要介紹了java發(fā)送post請求使用multipart/form-data格式文件數據到接口的相關資料,文中指定了數據編碼格式為UTF-8,并強調了所需依賴工具類,需要的朋友可以參考下2024-12-12
Spring Boot配置application.yml及根據application.yml選擇啟動配置的操作
Spring Boot中可以選擇applicant.properties 作為配置文件,也可以通過在application.yml中進行配置,讓Spring Boot根據你的選擇進行加載啟動配置文件,本文給大家介紹Spring Boot配置application.yml及根據application.yml選擇啟動配置的操作方法,感興趣的朋友一起看看吧2023-10-10

