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

深入淺出RocketMQ的事務(wù)消息

 更新時間:2023年04月09日 12:00:28   作者:小王曾是少年  
RocketMQ事務(wù)消息(Transactional?Message)是指應(yīng)用本地事務(wù)和發(fā)送消息操作可以被定義到全局事務(wù)中,要么同時成功,要么同時失敗。本文主要介紹了RocketMQ事務(wù)消息的相關(guān)知識,需要的可以參考一下

事務(wù)消息發(fā)送流程

半消息實現(xiàn)了分布式環(huán)境下的數(shù)據(jù)一致性的處理,生產(chǎn)者發(fā)送事務(wù)消息的流程如上圖所示,通過對源碼的學(xué)習(xí),我們可以弄清楚下面幾點,也是半消息機制的核心:

1.為什么prepare消息不會被Consumer消費?

2.事務(wù)消息是如何提交和回滾的?

3.定時回查本地事務(wù)狀態(tài)的實現(xiàn)細節(jié)。

發(fā)送事務(wù)消息源碼分析

發(fā)送事務(wù)消息方法TransactionMQProducer.sendMessageInTransaction:

  • msg:消息
  • tranExecuter:本地事務(wù)執(zhí)行器
  • arg:本地事務(wù)執(zhí)行器參數(shù)
public TransactionSendResult sendMessageInTransaction(final Message msg,
        final LocalTransactionExecuter localTransactionExecuter, final Object arg)
        throws MQClientException {
        TransactionListener transactionListener = getCheckListener();
        if (null == localTransactionExecuter && null == transactionListener) {
            throw new MQClientException("tranExecutor is null", null);
        }

        // 忽視消息延遲的屬性
        if (msg.getDelayTimeLevel() != 0) {
            MessageAccessor.clearProperty(msg, MessageConst.PROPERTY_DELAY_TIME_LEVEL);
        }

        Validators.checkMessage(msg, this.defaultMQProducer);
		
		// 發(fā)送半消息
        SendResult sendResult = null;
        MessageAccessor.putProperty(msg, MessageConst.PROPERTY_TRANSACTION_PREPARED, "true");
        MessageAccessor.putProperty(msg, MessageConst.PROPERTY_PRODUCER_GROUP, this.defaultMQProducer.getProducerGroup());
        try {
            sendResult = this.send(msg);
        } catch (Exception e) {
            throw new MQClientException("send message Exception", e);
        }
		
		// 處理發(fā)送半消息的結(jié)果
        LocalTransactionState localTransactionState = LocalTransactionState.UNKNOW;
        Throwable localException = null;
        switch (sendResult.getSendStatus()) {
        	// 發(fā)送半消息成功,執(zhí)行本地事務(wù)邏輯
            case SEND_OK: {
                try {
                    if (sendResult.getTransactionId() != null) {
                        msg.putUserProperty("__transactionId__", sendResult.getTransactionId());
                    }
                    String transactionId = msg.getProperty(MessageConst.PROPERTY_UNIQ_CLIENT_MESSAGE_ID_KEYIDX);
                    if (null != transactionId && !"".equals(transactionId)) {
                        msg.setTransactionId(transactionId);
                    }
                    // 執(zhí)行本地事務(wù)邏輯
                    if (null != localTransactionExecuter) {
                        localTransactionState = localTransactionExecuter.executeLocalTransactionBranch(msg, arg);
                    } else if (transactionListener != null) {
                        log.debug("Used new transaction API");
                        localTransactionState = transactionListener.executeLocalTransaction(msg, arg);
                    }
                    if (null == localTransactionState) {
                        localTransactionState = LocalTransactionState.UNKNOW;
                    }

                    if (localTransactionState != LocalTransactionState.COMMIT_MESSAGE) {
                        log.info("executeLocalTransactionBranch return {}", localTransactionState);
                        log.info(msg.toString());
                    }
                } catch (Throwable e) {
                    log.info("executeLocalTransactionBranch exception", e);
                    log.info(msg.toString());
                    localException = e;
                }
            }
            break;
            // 發(fā)送半消息失敗,標(biāo)記本地事務(wù)狀態(tài)為回滾
            case FLUSH_DISK_TIMEOUT:
            case FLUSH_SLAVE_TIMEOUT:
            case SLAVE_NOT_AVAILABLE:
                localTransactionState = LocalTransactionState.ROLLBACK_MESSAGE;
                break;
            default:
                break;
        }
		
		// 結(jié)束事務(wù),設(shè)置消息 COMMIT / ROLLBACK
        try {
            this.endTransaction(msg, sendResult, localTransactionState, localException);
        } catch (Exception e) {
            log.warn("local transaction execute " + localTransactionState + ", but end broker transaction failed", e);
        }
		
		// 返回事務(wù)發(fā)送結(jié)果
        TransactionSendResult transactionSendResult = new TransactionSendResult();
        transactionSendResult.setSendStatus(sendResult.getSendStatus());
        transactionSendResult.setMessageQueue(sendResult.getMessageQueue());
        
        // 提取Prepared消息的uniqID
        transactionSendResult.setMsgId(sendResult.getMsgId());
        transactionSendResult.setQueueOffset(sendResult.getQueueOffset());
        transactionSendResult.setTransactionId(sendResult.getTransactionId());
        transactionSendResult.setLocalTransactionState(localTransactionState);
        return transactionSendResult;
    }

該方法的入?yún)幸粋€需要用戶實現(xiàn)本地事務(wù)的LocalTransactionExecuter executer,executer中會進行事務(wù)操作以保證本地事務(wù)和消息發(fā)送這兩個操作的原子性。

由上面的源碼可知:

Producer會首先發(fā)送一個半消息到Broker中:

  • 半消息發(fā)送成功,執(zhí)行事務(wù)
  • 半消息發(fā)送失敗,不執(zhí)行事務(wù)

半消息發(fā)送到Broker后不會被Consumer消費掉的原因有以下兩點:

  • Broker在將消息寫入CommitLog時會判斷消息類型,如果是prepare或者rollback消息,ConsumeQueue的offset不變
  • Broker在構(gòu)造ConsumeQueue時會判斷是否是處于prepare或者rollback狀態(tài)的消息,如果是則不會將該消息放入ConsumeQueue里,Consumer在拉取消息時也就不會拉取到這條消息

Producer會根據(jù)半消息的發(fā)送結(jié)果和本地任務(wù)執(zhí)行結(jié)果來決定如何處理事務(wù)(commit或rollback),方法最后調(diào)用了endTransaction來處理事務(wù)的執(zhí)行結(jié)果,源碼如下:

  • sendResult:發(fā)送半消息的結(jié)果
  • localTransactionState:本地事務(wù)狀態(tài)
  • localException:執(zhí)行本地事務(wù)邏輯產(chǎn)生的異常
  • RemotingException:遠程調(diào)用異常
  • MQBrokerException:Broker異常
  • InterruptedException:當(dāng)線程中斷異常
  • UnknownHostException:未知host異常
public void endTransaction(
        final Message msg,
        final SendResult sendResult,
        final LocalTransactionState localTransactionState,
        final Throwable localException) throws RemotingException, MQBrokerException, InterruptedException, UnknownHostException {
        // 解碼消息id
        final MessageId id;
        if (sendResult.getOffsetMsgId() != null) {
            id = MessageDecoder.decodeMessageId(sendResult.getOffsetMsgId());
        } else {
            id = MessageDecoder.decodeMessageId(sendResult.getMsgId());
        }

		// 創(chuàng)建請求
        String transactionId = sendResult.getTransactionId();
        final String brokerAddr = this.mQClientFactory.findBrokerAddressInPublish(sendResult.getMessageQueue().getBrokerName());
        EndTransactionRequestHeader requestHeader = new EndTransactionRequestHeader();
        requestHeader.setTransactionId(transactionId);
        requestHeader.setCommitLogOffset(id.getOffset());
        switch (localTransactionState) {
            case COMMIT_MESSAGE:
                requestHeader.setCommitOrRollback(MessageSysFlag.TRANSACTION_COMMIT_TYPE);
                break;
            case ROLLBACK_MESSAGE:
                requestHeader.setCommitOrRollback(MessageSysFlag.TRANSACTION_ROLLBACK_TYPE);
                break;
            case UNKNOW:
                requestHeader.setCommitOrRollback(MessageSysFlag.TRANSACTION_NOT_TYPE);
                break;
            default:
                break;
        }

        doExecuteEndTransactionHook(msg, sendResult.getMsgId(), brokerAddr, localTransactionState, false);
        requestHeader.setProducerGroup(this.defaultMQProducer.getProducerGroup());
        requestHeader.setTranStateTableOffset(sendResult.getQueueOffset());
        requestHeader.setMsgId(sendResult.getMsgId());
        String remark = localException != null ? ("executeLocalTransactionBranch exception: " + localException.toString()) : null;

		// 提交 commit / rollback 消息 
        this.mQClientFactory.getMQClientAPIImpl().endTransactionOneway(brokerAddr, requestHeader, remark,
            this.defaultMQProducer.getSendMsgTimeout());
    }

該方法是將事務(wù)執(zhí)行的結(jié)果發(fā)送給Broker,再由Broker決定是否進行消息投遞,執(zhí)行步驟如下:

1.收到消息后先檢查是否是事務(wù)消息,如果不是事務(wù)消息則直接返回

2.根據(jù)請求頭里的offset查詢半消息,如果查詢結(jié)果為空則直接返回

3.根據(jù)半消息構(gòu)造新消息,新構(gòu)造的消息會被重新寫入到CommitLog里,rollback消息的消息體為空

4.如果是rollback消息,則該消息不會被投遞

具體原因上文中已經(jīng)分析過:只有commit消息才會被Broker投遞給consumer

RocketMQ會將commit消息和rollback消息都寫入到commitLog里,但rollback消息的消息體為空且不會被投遞,CommitLog在刪除過期消息時才會將其刪除。當(dāng)事務(wù)commit成功之后,RocketMQ會重新封裝半消息并將其投遞給Consumer端消費。

事務(wù)消息回查

Broker發(fā)起

相較于普通消息,事務(wù)消息主要依賴下面三個類:

1.TransactionStateService:事務(wù)狀態(tài)服務(wù),負責(zé)對事務(wù)消息進行管理,包括存儲和更新事務(wù)消息狀態(tài)、回查狀態(tài)等

2.TranStateTable:事務(wù)消息狀態(tài)存儲表,基于MappedFileQueue實現(xiàn)

3.TranRedoLog:TranStateTable的日志,每次寫入操作都會記錄日志,當(dāng)Broker宕機時,可以利用這個文件做數(shù)據(jù)恢復(fù)

存儲半消息到CommitLog時,使用offset索引到對應(yīng)的TranStateTable的位置

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

相關(guān)文章

  • maven工程打包引入本地jar包的實現(xiàn)

    maven工程打包引入本地jar包的實現(xiàn)

    我們需要將jar包發(fā)布到一些指定的第三方Maven倉庫,本文主要介紹了maven工程打包引入本地jar包的實現(xiàn),具有一定的參考價值,感興趣的可以了解一下
    2024-02-02
  • java 字符串的拼接的實現(xiàn)實例

    java 字符串的拼接的實現(xiàn)實例

    這篇文章主要介紹了java 字符串的拼接的實現(xiàn)實例的相關(guān)資料,希望通過本文大家能掌握字符拼接的實現(xiàn),需要的朋友可以參考下
    2017-09-09
  • SpringBoot定時任務(wù)詳解與案例代碼

    SpringBoot定時任務(wù)詳解與案例代碼

    SpringBoot是一個流行的Java開發(fā)框架,它提供了許多便捷的特性來簡化開發(fā)過程,其中之一就是定時任務(wù)的支持,讓開發(fā)人員可以輕松地在應(yīng)用程序中執(zhí)行定時任務(wù),本文將詳細介紹如何在Spring?Boot中使用定時任務(wù),并提供相關(guān)的代碼示例
    2023-06-06
  • 在Java中為日期增加一天的多種方法

    在Java中為日期增加一天的多種方法

    這篇文章主要給大家介紹了關(guān)于如何在Java中為日期增加一天的多種方法,在JAVA業(yè)務(wù)代碼中,經(jīng)常會遇到通過指定時間,增加指定天數(shù)的業(yè)務(wù)需求,需要的朋友可以參考下
    2023-07-07
  • Java中NIO的三大核心組件詳細解析

    Java中NIO的三大核心組件詳細解析

    這篇文章主要介紹了Java中NIO的三大核心組件詳細解析,NIO的Buffer類是一個抽象類,位于java.nio包中,提供了一組更加有效的方法,用來進行寫入和讀取的交替訪問,本質(zhì)上是一個內(nèi)存塊,既可以寫入數(shù)據(jù),也可以從中讀取數(shù)據(jù),需要的朋友可以參考下
    2023-12-12
  • java中HashMap的原理分析

    java中HashMap的原理分析

    HashMap在Java開發(fā)中有著非常重要的角色地位,每一個Java程序員都應(yīng)該了解HashMap。詳細地闡述HashMap中的幾個概念,并深入探討HashMap的內(nèi)部結(jié)構(gòu)和實現(xiàn)細節(jié),討論HashMap的性能問題
    2016-03-03
  • Java編程實現(xiàn)統(tǒng)計數(shù)組中各元素出現(xiàn)次數(shù)的方法

    Java編程實現(xiàn)統(tǒng)計數(shù)組中各元素出現(xiàn)次數(shù)的方法

    這篇文章主要介紹了Java編程實現(xiàn)統(tǒng)計數(shù)組中各元素出現(xiàn)次數(shù)的方法,涉及java針對數(shù)組的遍歷、比較、運算等相關(guān)操作技巧,需要的朋友可以參考下
    2017-07-07
  • Java?CAS機制詳解

    Java?CAS機制詳解

    這篇文章主要介紹了Java?CAS機制,CAS機制是一種數(shù)據(jù)更新的方式。在具體講什么是CAS機制之前,我們先來聊下在多線程環(huán)境下,對共享變量進行數(shù)據(jù)更新的兩種模式:悲觀鎖模式和樂觀鎖模式
    2023-01-01
  • Java序列化問題:“Serialized class has not implement Serializable interface”錯誤的解決方法

    Java序列化問題:“Serialized class has not impl

    在Java開發(fā)中,序列化(Serialization)是一個常見的操作,尤其是在分布式系統(tǒng)、網(wǎng)絡(luò)通信或數(shù)據(jù)持久化場景中,然而,序列化過程中可能會遇到各種問題,其中最常見的一個錯誤是Serialized class has not implement Serializable interface,本文給大家介紹了相關(guān)的解決方法
    2025-02-02
  • SpringBoot創(chuàng)建并簡單使用的實現(xiàn)

    SpringBoot創(chuàng)建并簡單使用的實現(xiàn)

    這篇文章主要介紹了SpringBoot創(chuàng)建并簡單使用的實現(xiàn),文中通過示例代碼介紹的非常詳細,對大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價值,需要的朋友們下面隨著小編來一起學(xué)習(xí)學(xué)習(xí)吧
    2020-10-10

最新評論

阳泉市| 广平县| 汉寿县| 马公市| 宣城市| 鱼台县| 抚远县| 广平县| 甘谷县| 韩城市| 武定县| 三原县| 惠州市| 兴化市| 平安县| 黄梅县| 宿州市| 抚远县| 麦盖提县| 合肥市| 西乡县| 江陵县| 深州市| 堆龙德庆县| 洛扎县| 同江市| 天门市| 台山市| 乌什县| 冀州市| 洞口县| 北川| 鹤岗市| 察雅县| 辽宁省| 临桂县| 阳东县| 高密市| 社旗县| 无棣县| 夏津县|