最新国产好看的视频,伊人天堂AV在线,国产Aaaaaa视频,蜜臀视频在线观看一区,人妻av色图,密臀久久久精品影片,青青视频免费观看毛片,久草在线观看视,国产三级精品色情在线

SpringBoot集成RocketMQ發(fā)送事務(wù)消息的原理解析

 更新時(shí)間:2022年06月30日 15:48:40   作者:劍圣無(wú)痕  
RocketMQ 的事務(wù)消息提供類似 X/Open XA 的分布事務(wù)功能,通過(guò)事務(wù)消息能達(dá)到分布式事務(wù)的最終一致,這篇文章主要介紹了SpringBoot集成RocketMQ發(fā)送事務(wù)消息,需要的朋友可以參考下

簡(jiǎn)介

RocketMQ 事務(wù)消息(Transactional Message)是指應(yīng)用本地事務(wù)和發(fā)送消息操作可以被定義到全局事務(wù)中,要么同時(shí)成功,要么同時(shí)失敗。RocketMQ 的事務(wù)消息提供類似 X/Open XA 的分布事務(wù)功能,通過(guò)事務(wù)消息能達(dá)到分布式事務(wù)的最終一致。

原理

RocketMQ事務(wù)消息通過(guò)異步確保方式,保證事務(wù)的最終一致性。設(shè)計(jì)的思想可以借鑒兩個(gè)階段提交事務(wù)。其執(zhí)行流程圖如下:

圖片.png

  • 發(fā)送方向MQ服務(wù)端發(fā)送消息。
  • MQ Server將消息持久化成功之后,向發(fā)送方 ACK 確認(rèn)消息已經(jīng)發(fā)送成功,此時(shí)消息為半消息。
  • 發(fā)送方開(kāi)始執(zhí)行本地事務(wù)邏輯。
  • 發(fā)送方根據(jù)本地事務(wù)執(zhí)行結(jié)果向 MQ Server 提交二次確認(rèn)(Commit 或是 Rollback),MQ Server 收到 Commit 狀態(tài)則將半消息標(biāo)記為可投遞,訂閱方最終將收到該消息;MQ Server 收到 Rollback 狀態(tài)則刪除半消息,訂閱方將不會(huì)接受該消息。
  • 在斷網(wǎng)或者是應(yīng)用重啟的特殊情況下,上述步驟4提交的二次確認(rèn)最終未到達(dá) MQ Server,經(jīng)過(guò)固定時(shí)間后 MQ Server 將對(duì)該消息發(fā)起消息回查。
  • 發(fā)送方收到消息回查后,需要檢查對(duì)應(yīng)消息的本地事務(wù)執(zhí)行的最終結(jié)果。
  • 發(fā)送方根據(jù)檢查得到的本地事務(wù)的最終狀態(tài)再次提交二次確認(rèn),MQ Server 仍按照步驟4對(duì)半消息進(jìn)行操作。

具體實(shí)現(xiàn)

消費(fèi)者

@Component
public class TransactionProduce
{
    private Logger logger = LoggerFactory.getLogger(getClass());
    
    @Autowired
    private RocketMQTemplate rocketMQTemplate;
    
    public void sendTransactionMessage(String msg)
    {
        logger.info("start sendTransMessage hashKey:{}",msg);
       
         Message message =new Message();
         message.setBody("this is tx message".getBytes());
         TransactionSendResult result=rocketMQTemplate.sendMessageInTransaction("test-tx-rocketmq", 
                 MessageBuilder.withPayload(message).build(), msg);
         
         //發(fā)送狀態(tài)
         String sendStatus = result.getSendStatus().name();
         // 本地事務(wù)執(zhí)行狀態(tài)
         String localTxState = result.getLocalTransactionState().name();
         logger.info("send tx message sendStatus:{},localTXState:{}",sendStatus,localTxState);
    } 
}

說(shuō)明:發(fā)送事務(wù)消息采用的是sendMessageInTransaction方法,返回結(jié)果為TransactionSendResult對(duì)象,該對(duì)象中包含了事務(wù)發(fā)送的狀態(tài)、本地事務(wù)執(zhí)行的狀態(tài)等。

消費(fèi)者

@Component
@RocketMQMessageListener(consumerGroup="test-txRocketmq-group",topic="test-tx-rocketmq", messageModel = MessageModel.CLUSTERING)
public class TransactionConsumer implements RocketMQListener<String>
{
    private Logger logger =LoggerFactory.getLogger(getClass());
    @Override
    public void onMessage(String message)
    {
        logger.info("send transaction mssage parma is:{}", message);
    }
}

說(shuō)明:發(fā)送事務(wù)消息的消費(fèi)者與普通的消費(fèi)者一樣沒(méi)有太大的區(qū)別。

生產(chǎn)者消息監(jiān)聽(tīng)器

發(fā)送事務(wù)消息除了生產(chǎn)者和消費(fèi)者以外,我們還需要?jiǎng)?chuàng)建生產(chǎn)者的消息監(jiān)聽(tīng)器,來(lái)監(jiān)聽(tīng)本地事務(wù)執(zhí)行的狀態(tài)和檢查本地事務(wù)狀態(tài)。

@RocketMQTransactionListener
public class TransactionMsgListener implements RocketMQLocalTransactionListener
{
    private Logger logger = LoggerFactory.getLogger(getClass());
    /**
     * 執(zhí)行本地事務(wù)
     */
    @Override
    public RocketMQLocalTransactionState executeLocalTransaction(Message msg,
            Object obj)
    {
        logger.info("start invoke local rocketMQ transaction");
        RocketMQLocalTransactionState resultState = RocketMQLocalTransactionState.COMMIT;
        
        try
        {
            //處理業(yè)務(wù)
            String jsonStr = new String((byte[]) msg.getPayload(), StandardCharsets.UTF_8);
            logger.info("invoke msg content:{}",jsonStr);
        }
        catch (Exception e)
        {
            logger.error("invoke local mq trans error",e);
            resultState = RocketMQLocalTransactionState.UNKNOWN;
        }
        
        return resultState;
    }

    /**
     * 檢查本地事務(wù)的狀態(tài)
     */
    @Override
    public RocketMQLocalTransactionState checkLocalTransaction(Message msg)
    {
        logger.info("start check Local rocketMQ transaction");
        
        RocketMQLocalTransactionState resultState = RocketMQLocalTransactionState.COMMIT;
        
        try
        {
            String jsonStr = new String((byte[]) msg.getPayload(), StandardCharsets.UTF_8);
            logger.info("check trans msg content:{}",jsonStr);
        }
        catch (Exception e)
        {
            resultState  = RocketMQLocalTransactionState.ROLLBACK;
        }
        return resultState;
    }
}

說(shuō)明:RocketMQ本地事務(wù)狀態(tài)由如下幾種:

  • RocketMQLocalTransactionState.COMMIT:提交事務(wù),允許消費(fèi)者消費(fèi)此消息。
  • RocketMQLocalTransactionState.ROLLBACK: 回滾事務(wù),消息將被刪除,不允許被消費(fèi)。
  • RocketMQLocalTransactionState.UNKNOWN:中間狀態(tài),代表需要進(jìn)行檢查來(lái)確定狀態(tài)。

注意:Spring Boot2.0的版本之后,@RocketMQTransactionListener 已經(jīng)沒(méi)有了txProducerGroup屬性,且sendMessageInTransaction方法也將其移除。所以在同一項(xiàng)目中只能有一個(gè)@RocketMQTransactionListener,不能出現(xiàn)多個(gè),否則會(huì)報(bào)如下錯(cuò)誤:

java.lang.IllegalStateException: rocketMQTemplate already exists RocketMQLocalTransactionListener

消息事務(wù)測(cè)試

正常測(cè)試

c.s.fw.mq.produce.TransactionProduce - product start sendTransMessage msg:{"userId":"zhangsann"}
c.s.f.m.p.TransactionMsgListener - start invoke local rocketMQ transaction
c.s.f.m.p.TransactionMsgListener - invoke local transaction msg content:{"topic":null,"flag":0,"properties":null,"body":"dGhpcyBpcyB0eCBtZXNzYWdl","transactionId":null,"keys":null,"tags":null,"delayTimeLevel":0,"waitStoreMsgOK":true,"buyerId":null}
c.s.fw.mq.produce.TransactionProduce - send tx message sendStatus:SEND_OK,localTXState:COMMIT_MESSAGE
c.s.f.m.consumer.TransactionConsumer - send transaction mssage parma is:{"topic":null,"flag":0,"properties":null,"body":"dGhpcyBpcyB0eCBtZXNzYWdl","transactionId":null,"keys":null,"tags":null,"delayTimeLevel":0,"waitStoreMsgOK":true,"buyerId":null}

說(shuō)明:通過(guò)日志我們可以看出,執(zhí)行的流程與上述的一致,執(zhí)行成功后,消息執(zhí)行成功返回的結(jié)果為SEND_OK,本地事務(wù)執(zhí)行的狀態(tài)為COMMIT_MESSAGE。

異常測(cè)試

如果在執(zhí)行本地消息時(shí)出現(xiàn)異常,那么執(zhí)行結(jié)果會(huì)是怎樣?修改下本地事務(wù)執(zhí)行的方法,讓其出現(xiàn)異常。

代碼調(diào)整

  @Override
    public RocketMQLocalTransactionState executeLocalTransaction(Message msg,
            Object obj)
    {
        logger.info("start invoke local rocketMQ transaction");
        RocketMQLocalTransactionState resultState = RocketMQLocalTransactionState.COMMIT;
        
        try
        {
            //處理業(yè)務(wù)
            String jsonStr = new String((byte[]) msg.getPayload(), StandardCharsets.UTF_8);
            logger.info("invoke local transaction msg content:{}",jsonStr);
             int c=1/0;
        }
        catch (Exception e)
        {
            logger.error("invoke local mq trans error",e);
            resultState = RocketMQLocalTransactionState.UNKNOWN;
        }
        
        return resultState;
    }

執(zhí)行結(jié)果

c.s.fw.mq.produce.TransactionProduce - send tx message sendStatus:SEND_OK,localTXState:UNKNOW

從執(zhí)行的結(jié)果可以看出,消息執(zhí)行成功返回的結(jié)果為SEND_OK,本地事務(wù)執(zhí)行的狀態(tài)為:UNKNOW.所以消費(fèi)端無(wú)法消費(fèi)此消息。

總結(jié)

到此這篇關(guān)于SpringBoot集成RocketMQ發(fā)送事務(wù)消息的文章就介紹到這了,更多相關(guān)SpringBoot集成RocketMQ事務(wù)消息內(nèi)容請(qǐng)搜索腳本之家以前的文章或繼續(xù)瀏覽下面的相關(guān)文章希望大家以后多多支持腳本之家!

相關(guān)文章

  • idea使用war以及war exploded的區(qū)別說(shuō)明

    idea使用war以及war exploded的區(qū)別說(shuō)明

    本文詳細(xì)解析了war與warexploded兩種部署方式的差異及步驟,war方式是先打包成war包,再部署到服務(wù)器上;warexploded方式是直接把文件夾、class文件等移到Tomcat上部署,支持熱部署,開(kāi)發(fā)時(shí)常用,文章分別列出了warexploded模式和war包形式的具體操作步驟
    2024-10-10
  • Java數(shù)據(jù)結(jié)構(gòu)及算法實(shí)例:三角數(shù)字

    Java數(shù)據(jù)結(jié)構(gòu)及算法實(shí)例:三角數(shù)字

    這篇文章主要介紹了Java數(shù)據(jù)結(jié)構(gòu)及算法實(shí)例:三角數(shù)字,本文直接給出實(shí)現(xiàn)代碼,代碼中包含詳細(xì)注釋,需要的朋友可以參考下
    2015-06-06
  • java正則匹配讀取txt文件提取特定開(kāi)頭和結(jié)尾的字符串

    java正則匹配讀取txt文件提取特定開(kāi)頭和結(jié)尾的字符串

    通常我們可以直接通過(guò)文件流來(lái)讀取txt文件的內(nèi)容,但有時(shí)候也會(huì)遇到問(wèn)題,下面這篇文章主要給大家介紹了關(guān)于java正則匹配讀取txt文件提取特定開(kāi)頭和結(jié)尾的字符串的相關(guān)資料,需要的朋友可以參考下
    2022-11-11
  • Java仿淘寶首頁(yè)分類列表功能的示例代碼

    Java仿淘寶首頁(yè)分類列表功能的示例代碼

    這篇文章主要介紹了仿淘寶分類管理功能的示例代碼,具有很好的參考價(jià)值,希望對(duì)大家有所幫助,也給大家做個(gè)參考
    2018-05-05
  • Android應(yīng)用開(kāi)發(fā)之將SQLite和APK一起打包的方法

    Android應(yīng)用開(kāi)發(fā)之將SQLite和APK一起打包的方法

    這篇文章主要介紹了Android應(yīng)用開(kāi)發(fā)之將SQLite和APK一起打包的方法,文章時(shí)間較早,盡管現(xiàn)在開(kāi)發(fā)環(huán)境已大都遷移至Android Studio上,但打包原理依然相同,需要的朋友可以參考下
    2015-08-08
  • Java中關(guān)于Map四種取值方式

    Java中關(guān)于Map四種取值方式

    這篇文章主要介紹了Java中關(guān)于Map四種取值方式,具有很好的參考價(jià)值,希望對(duì)大家有所幫助。如有錯(cuò)誤或未考慮完全的地方,望不吝賜教
    2023-03-03
  • Maven配置多倉(cāng)庫(kù)無(wú)效的解決

    Maven配置多倉(cāng)庫(kù)無(wú)效的解決

    在項(xiàng)目中使用Maven管理jar包依賴往往會(huì)出現(xiàn)很多問(wèn)題,所以這時(shí)候就需要配置Maven多倉(cāng)庫(kù),本文介紹了如何配置以及問(wèn)題的解決
    2021-05-05
  • java 中modCount 詳解及源碼分析

    java 中modCount 詳解及源碼分析

    這篇文章主要介紹了java 中modCount 詳解及源碼分析的相關(guān)資料,需要的朋友可以參考下
    2017-02-02
  • 深入淺析hbase的優(yōu)點(diǎn)

    深入淺析hbase的優(yōu)點(diǎn)

    本文講述了HBase的特征和它的優(yōu)點(diǎn),并簡(jiǎn)要回顧了行鍵設(shè)計(jì)的重點(diǎn)之處,它還向你展示了如何在本地配置HBase環(huán)境,使用命令創(chuàng)建表、插入數(shù)據(jù)、檢索指定行以及最后如何進(jìn)行scan操作,感興趣的朋友一起看看吧
    2017-09-09
  • Java中的反射機(jī)制示例詳解

    Java中的反射機(jī)制示例詳解

    反射就是把Java類中的各個(gè)成分映射成一個(gè)個(gè)的Java對(duì)象。本文將通過(guò)示例詳細(xì)講解Java中的反射機(jī)制,感興趣的小伙伴可以跟隨小編學(xué)習(xí)一下
    2022-03-03

最新評(píng)論

娄烦县| 敦煌市| 章丘市| 遂宁市| 基隆市| 绥阳县| 遂宁市| 无锡市| 股票| 万荣县| 屏东县| 辛集市| 改则县| 江津市| 田东县| 宁波市| 宁陵县| 石林| 高州市| 塔河县| 株洲县| 左云县| 石狮市| 吉林省| 蛟河市| 滦平县| 年辖:市辖区| 德化县| 宁安市| 涞源县| 昌图县| 古田县| 太保市| 时尚| 大港区| 彭州市| 五大连池市| 五指山市| 广安市| 崇仁县| 海城市|