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

RocketMQ特性Broker存儲事務消息實現(xiàn)

 更新時間:2022年08月17日 14:07:09   作者:奔跑的毛球  
這篇文章主要為大家介紹了RocketMQ特性Broker存儲事務消息實現(xiàn)詳解,有需要的朋友可以借鑒參考下,希望能夠有所幫助,祝大家多多進步,早日升職加薪

引言

Broker中,事務消息的初始化是通過BrokerController.initialTransaction()方法執(zhí)行的。

private void initialTransaction() {
    this.transactionalMessageService = ServiceProvider.loadClass(ServiceProvider.TRANSACTION_SERVICE_ID, TransactionalMessageService.class);
    if (null == this.transactionalMessageService) {
        this.transactionalMessageService = new TransactionalMessageServiceImpl(new TransactionalMessageBridge(this, this.getMessageStore()));
        LOG.warn("Load default transaction message hook service: {}", TransactionalMessageServiceImpl.class.getSimpleName());
    }
    this.transactionalMessageCheckListener = ServiceProvider.loadClass(ServiceProvider.TRANSACTION_LISTENER_ID, AbstractTransactionalMessageCheckListener.class);
    if (null == this.transactionalMessageCheckListener) {
        this.transactionalMessageCheckListener = new DefaultTransactionalMessageCheckListener();
        LOG.warn("Load default discard message hook service: {}", DefaultTransactionalMessageCheckListener.class.getSimpleName());
    }
    this.transactionalMessageCheckListener.setBrokerController(this);
    this.transactionalMessageCheckService = new TransactionalMessageCheckService(this);
}

這里有三個核心的初始化變量

TransactionalMessageService

事務消息主要處理服務。默認實現(xiàn)類是TransactionalMessageServiceImpl也可以自己定義事務消息處理實現(xiàn)類,通過ServiceProvider.loadClass()方法進行加載。

TransactionalMessageService類定義如下。內部屬性已加注釋標明。

public interface TransactionalMessageService {
    //用于保存Half事務消息
    PutMessageResult prepareMessage(MessageExtBrokerInner messageInner);
    CompletableFuture<PutMessageResult> asyncPrepareMessage(MessageExtBrokerInner messageInner);
    //刪除事務消息
    boolean deletePrepareMessage(MessageExt messageExt);
    //提交事務消息
    OperationResult commitMessage(EndTransactionRequestHeader requestHeader);
    //回滾事務消息
    OperationResult rollbackMessage(EndTransactionRequestHeader requestHeader);
    void check(long transactionTimeout, int transactionCheckMax, AbstractTransactionalMessageCheckListener listener);
    //打開事務消息
    boolean open();
    //關閉事務消息
    void close();
}

transactionalMessageCheckListener

事務消息回查監(jiān)聽器

transactionalMessageCheckService

事務消息回查服務,啟動一個線程定時檢查超時的Half消息是否需要回查。

處理事務消息

當初始化完成之后,Broker就可以處理事務消息了。

Broker存儲事務消息的是org.apache.rocketmq.broker.processor.SendMessageProcessor,這和普通消息其實是一樣的。

但是有兩點針對事務消息的特殊處理

第一處:

org.apache.rocketmq.broker.processor.SendMessageProcessor#sendMessage中:

//獲取擴展字段的值,若是該值為true則為事務消息
String traFlag = oriProps.get(MessageConst.PROPERTY_TRANSACTION_PREPARED);
boolean sendTransactionPrepareMessage = false;
if (Boolean.parseBoolean(traFlag)
    && !(msgInner.getReconsumeTimes() > 0 && msgInner.getDelayTimeLevel() > 0)) { 
    //判斷當前Broker配置是否支持事務消息
    if (this.brokerController.getBrokerConfig().isRejectTransactionMessage()) {
        response.setCode(ResponseCode.NO_PERMISSION);
        response.setRemark(
            "the broker[" + this.brokerController.getBrokerConfig().getBrokerIP1()
                + "] sending transaction message is forbidden");
        return response;
    }
    sendTransactionPrepareMessage = true;
}
if (sendTransactionPrepareMessage) {
    //保存Half信息
    putMessageResult = this.brokerController.getTransactionalMessageService().prepareMessage(msgInner);
} else {
    putMessageResult = this.brokerController.getMessageStore().putMessage(msgInner);
}

第二處:

存儲事務消息前的預處理,對應方法是

org.apache.rocketmq.broker.transaction.queue.TransactionalMessageBridge#parseHalfMessageInner

private MessageExtBrokerInner parseHalfMessageInner(MessageExtBrokerInner msgInner) {
    //將原消息的topic保存在擴展字段中
    MessageAccessor.putProperty(msgInner, MessageConst.PROPERTY_REAL_TOPIC, msgInner.getTopic());
    //將原消息的QueueId保存在擴展字段中
    MessageAccessor.putProperty(msgInner, MessageConst.PROPERTY_REAL_QUEUE_ID,
        String.valueOf(msgInner.getQueueId()));
    //將原消息的SysFlag保存在擴展字段中
    msgInner.setSysFlag(
        MessageSysFlag.resetTransactionValue(msgInner.getSysFlag(), MessageSysFlag.TRANSACTION_NOT_TYPE));
    //修改topic的值為RMQ_SYS_TRANS_HALF_TOPIC
    msgInner.setTopic(TransactionalMessageUtil.buildHalfTopic());
    //修改Queueid為0
    msgInner.setQueueId(0);
    msgInner.setPropertiesString(MessageDecoder.messageProperties2String(msgInner.getProperties()));
    return msgInner;
}

完成上述步驟之后,調用DefaultMessageStole.putMessage()方法將其保存到CommitLog中。

CommitLog存儲成功之后,通過org.apache.rocketmq.store.CommitLog.DefaultAppendMessageCallback#doAppend()方法對其進行處理。

final int tranType = MessageSysFlag.getTransactionValue(msgInner.getSysFlag());
switch (tranType) {
    // Prepared and Rollback message is not consumed, will not enter the consume queue
    case MessageSysFlag.TRANSACTION_PREPARED_TYPE:
    case MessageSysFlag.TRANSACTION_ROLLBACK_TYPE:
        queueOffset = 0L;
        break;
    case MessageSysFlag.TRANSACTION_NOT_TYPE:
    case MessageSysFlag.TRANSACTION_COMMIT_TYPE:
    default:
        break;
}

這里的邏輯是這樣的,當讀到的消息類型為事務消息時,設置當前消息的位點值為0,而不是設置真實的位點。這樣該位點就不會建立ConsumeQueue索引,也不會被消費。

以上就是RocketMQ特性Broker存儲事務消息實現(xiàn)的詳細內容,更多關于RocketMQ Broker存儲事務消息的資料請關注腳本之家其它相關文章!

相關文章

  • 詳解Spring Boot 部署jar和war的區(qū)別

    詳解Spring Boot 部署jar和war的區(qū)別

    本篇文章主要介紹了詳解Spring Boot 部署jar和war的區(qū)別,小編覺得挺不錯的,現(xiàn)在分享給大家,也給大家做個參考。一起跟隨小編過來看看吧
    2017-09-09
  • 在Filter中不能注入bean的問題及解決

    在Filter中不能注入bean的問題及解決

    這篇文章主要介紹了在Filter中不能注入bean的問題及解決方案,具有很好的參考價值,希望對大家有所幫助。如有錯誤或未考慮完全的地方,望不吝賜教
    2021-11-11
  • Java使用Swing實現(xiàn)一個模擬電腦計算器

    Java使用Swing實現(xiàn)一個模擬電腦計算器

    Java Swing 是一個用于創(chuàng)建 Java GUI(圖形用戶界面)的框架,它提供了一系列的 GUI 組件和工具,可以用于創(chuàng)建桌面應用程序,包括按鈕、文本框、標簽、表格等等,本文給大家介紹了Java使用Swing實現(xiàn)一個模擬計算器,感興趣的同學可以自己動手嘗試一下
    2024-05-05
  • 從零開始講解Java微信公眾號消息推送實現(xiàn)

    從零開始講解Java微信公眾號消息推送實現(xiàn)

    微信公眾號分為訂閱號和服務號,無論有沒有認證,訂閱號每天都能推送一條消息,也就是每天只能推送一次消息給粉絲,這篇文章主要給大家介紹了關于Java微信公眾號消息推送實現(xiàn)的相關資料,需要的朋友可以參考下
    2022-09-09
  • Java中的回調

    Java中的回調

    這篇文章主要介紹了Java中回調的相關資料,幫助大家更好的理解和學習java,感興趣的朋友可以了解下
    2020-08-08
  • IDEA中實體類(POJO)與JSON快速互轉問題

    IDEA中實體類(POJO)與JSON快速互轉問題

    這篇文章主要介紹了IDEA中實體類(POJO)與JSON快速互轉,本文通過圖文實例代碼相結合給大家介紹的非常詳細,需要的朋友可以參考下
    2022-08-08
  • JAVA設計模式----建造者模式詳解

    JAVA設計模式----建造者模式詳解

    這篇文章主要為大家詳細介紹了java實現(xiàn)建造者模式Builder Pattern,具有一定的參考價值,感興趣的小伙伴們可以參考一下
    2021-09-09
  • 關于mybatis一對一查詢一對多查詢遇到的問題

    關于mybatis一對一查詢一對多查詢遇到的問題

    這篇文章主要介紹了關于mybatis一對一查詢,一對多查詢遇到的錯誤,接下來是對文章進行操作,要求查詢全部文章,并關聯(lián)查詢作者,文章標簽,本文給大家介紹的非常詳細,需要的朋友可以參考下
    2022-05-05
  • 解決后端傳long類型數(shù)據(jù)到前端精度丟失問題

    解決后端傳long類型數(shù)據(jù)到前端精度丟失問題

    這篇文章主要介紹了解決后端傳long類型數(shù)據(jù)到前端精度丟失問題,具有很好的參考價值,希望對大家有所幫助,如有錯誤或未考慮完全的地方,望不吝賜教
    2024-01-01
  • SpringBoot日志進階實戰(zhàn)之Logback配置經驗和方法

    SpringBoot日志進階實戰(zhàn)之Logback配置經驗和方法

    本文給大家介紹在SpringBoot中使用Logback配置日志的經驗和方法,并提供了詳細的代碼示例和解釋,包括:滾動文件、異步日志記錄、動態(tài)指定屬性、日志級別、配置文件等常用功能,覆蓋日常Logback配置開發(fā)90%的知識點,感興趣的朋友跟隨小編一起看看吧
    2023-06-06

最新評論

鄱阳县| 集贤县| 陇西县| 峨眉山市| 北川| 东乡县| 聊城市| 册亨县| 韩城市| 耒阳市| 岑巩县| 靖安县| 满城县| 清水县| 花莲县| 鄂温| 全州县| 晋宁县| 雷波县| 宜都市| 塘沽区| 南召县| 乌鲁木齐市| 浮梁县| 达日县| 石屏县| 宁明县| 湖南省| 灵台县| 景德镇市| 阿克陶县| 巴里| 封开县| 新宾| 南靖县| 曲松县| 小金县| 彩票| 通城县| 普格县| 志丹县|