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

RocketMQ-延遲消息的處理流程介紹

 更新時間:2021年07月03日 09:30:01   作者:pigcoffee  
這篇文章主要介紹了RocketMQ-延遲消息的處理流程,具有很好的參考價值,希望對大家有所幫助。如有錯誤或未考慮完全的地方,望不吝賜教

概述

RocketMQ 支持發(fā)送延遲消息,但不支持任意時間的延遲消息的設置,僅支持內(nèi)置預設值的延遲時間間隔的延遲消息;

預設值的延遲時間間隔為:

1s、 5s、 10s、 30s、 1m、 2m、 3m、 4m、 5m、 6m、 7m、 8m、 9m、 10m、 20m、 30m、 1h、 2h;

在消息創(chuàng)建的時候,調(diào)用 setDelayTimeLevel(int level) 方法設置延遲時間;

broker在接收到延遲消息的時候會把對應延遲級別的消息先存儲到對應的延遲隊列中,等延遲消息時間到達時,會把消息重新存儲到對應的topic的queue里面。

Broker處理延遲消息

CommitLog.putMessage()

//獲取消息的sysflag
        final int tranType = MessageSysFlag.getTransactionValue(msg.getSysFlag());
        //非事務消息 或 已commit事務消息
        if (tranType == MessageSysFlag.TRANSACTION_NOT_TYPE
            || tranType == MessageSysFlag.TRANSACTION_COMMIT_TYPE) {
            // Delay Delivery 判斷消息是否設置延遲
            if (msg.getDelayTimeLevel() > 0) {
                //判斷延遲級別是否大于最大級別,如果大于最大值,則將延遲級別設置為最大級
                if (msg.getDelayTimeLevel() > this.defaultMessageStore.getScheduleMessageService().getMaxDelayLevel()) {
                    msg.setDelayTimeLevel(this.defaultMessageStore.getScheduleMessageService().getMaxDelayLevel());
                }
                //延遲消息的topic為 SCHEDULE_TOPIC_XXXX
                topic = ScheduleMessageService.SCHEDULE_TOPIC;
                //獲取延遲級別,一個延遲級別對應一個Queue
                queueId = ScheduleMessageService.delayLevel2QueueId(msg.getDelayTimeLevel());
 
                // Backup real topic, queueId
                //消息原始的topic,queueid保存到消息的property中
                MessageAccessor.putProperty(msg, MessageConst.PROPERTY_REAL_TOPIC, msg.getTopic());
                MessageAccessor.putProperty(msg, MessageConst.PROPERTY_REAL_QUEUE_ID, String.valueOf(msg.getQueueId()));
                msg.setPropertiesString(MessageDecoder.messageProperties2String(msg.getProperties()));
 
                msg.setTopic(topic);
                msg.setQueueId(queueId);
            }
        }

1、判斷消息類型,如果是非事務消息、已commit事務消息,才能處理延遲消息

2、判斷消息是否設置延遲級別,如果延遲級別大于0,則該消息為延遲消息

3、判斷延遲級別是否大于最大級別,如果大于最大值,則將延遲級別設置為最大級

4、延遲消息的topic為 SCHEDULE_TOPIC_XXXX

5、獲取延遲級別,一個延遲級別對應一個Queue

6、消息原始的topic,queueid保存到消息的property中

7、修改消息的topci、queueid

啟動延遲消息定時任務

ScheduleMessageService.start()

延遲消息投遞

以上為個人經(jīng)驗,希望能給大家一個參考,也希望大家多多支持腳本之家。

相關(guān)文章

最新評論

福泉市| 兰溪市| 大名县| 六安市| 邻水| 封开县| 荥经县| 丰都县| 定南县| 嘉义县| 富民县| 安庆市| 明光市| 舟曲县| 江都市| 开原市| 瑞丽市| 浮山县| 阳江市| 个旧市| 广灵县| 井冈山市| 饶河县| 尤溪县| 阳山县| 苍梧县| 晋中市| 耿马| 扎鲁特旗| 蒲江县| 兴安县| 汉中市| 临泉县| 永康市| 马边| 伊宁县| 吴旗县| 渝中区| 潮安县| 宁远县| 建德市|