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

Java中RocketMQ的延遲消息詳解

 更新時(shí)間:2023年09月15日 09:36:59   作者:萬(wàn)貓學(xué)社  
這篇文章主要介紹了Java中RocketMQ的延遲消息詳解,RocketMQ是一款開源的分布式消息系統(tǒng),基于高可用分布式集群技術(shù),提供低延時(shí)的、高可靠、萬(wàn)億級(jí)容量、靈活可伸縮的消息發(fā)布與訂閱服務(wù),需要的朋友可以參考下

RocketMQ簡(jiǎn)介

RocketMQ是一款開源的分布式消息系統(tǒng),基于高可用分布式集群技術(shù),提供低延時(shí)的、高可靠、萬(wàn)億級(jí)容量、靈活可伸縮的消息發(fā)布與訂閱服務(wù)。

它前身是MetaQ,是阿里基于Kafka的設(shè)計(jì)使用Java進(jìn)行自主研發(fā)的。在2012年,阿里將其開源, 在2016年,阿里將其捐獻(xiàn)給Apache軟件基金會(huì)(Apache Software Foundation,簡(jiǎn)稱為ASF),正式成為孵化項(xiàng)目。2017 年,Apache軟件基金會(huì)宣布RocketMQ已孵化成為 Apache頂級(jí)項(xiàng)目(Top Level Project,簡(jiǎn)稱為TLP ),是國(guó)內(nèi)首個(gè)互聯(lián)網(wǎng)中間件在 Apache上的頂級(jí)項(xiàng)目。

延遲消息

生產(chǎn)者把消息發(fā)送到消息隊(duì)列中以后,并不期望被立即消費(fèi),而是等待指定時(shí)間后才可以被消費(fèi)者消費(fèi),這類消息通常被稱為延遲消息。

在RocketMQ中,支持延遲消息,但是不支持任意時(shí)間精度的延遲消息,只支持特定級(jí)別的延遲消息。如果要支持任意時(shí)間精度,不能避免在Broker層面做消息排序,再涉及到持久化的考量,那么消息排序就不可避免產(chǎn)生巨大的性能開銷。

消息延遲級(jí)別分別為1s 5s 10s 30s 1m 2m 3m 4m 5m 6m 7m 8m 9m 10m 20m 30m 1h 2h,共18個(gè)級(jí)別。在發(fā)送消息時(shí),設(shè)置消息延遲級(jí)別即可,設(shè)置消息延遲級(jí)別時(shí)有以下3種情況:

  1. 設(shè)置消息延遲級(jí)別等于0時(shí),則該消息為非延遲消息。
  2. 設(shè)置消息延遲級(jí)別大于等于1并且小于等于18時(shí),消息延遲特定時(shí)間,如:設(shè)置消息延遲級(jí)別等于1,則延遲1s;設(shè)置消息延遲級(jí)別等于2,則延遲5s,以此類推。
  3. 設(shè)置消息延遲級(jí)別大于18時(shí),則該消息延遲級(jí)別為18,如:設(shè)置消息延遲級(jí)別等于20,則延遲2h。

延遲消息示例

首先,寫一個(gè)消費(fèi)者,用于消費(fèi)延遲消息:

public class Consumer {
    public static void main(String[] args) throws MQClientException {
        SimpleDateFormat sdf = new SimpleDateFormat("HH:mm:ss.SSS");
        // 實(shí)例化消費(fèi)者
        DefaultMQPushConsumer consumer = new DefaultMQPushConsumer("OneMoreGroup");
        // 設(shè)置NameServer的地址
        consumer.setNamesrvAddr("localhost:9876");
        // 訂閱一個(gè)或者多個(gè)Topic,以及Tag來(lái)過(guò)濾需要消費(fèi)的消息
        consumer.subscribe("OneMoreTopic", "*");
        // 注冊(cè)回調(diào)實(shí)現(xiàn)類來(lái)處理從broker拉取回來(lái)的消息
        consumer.registerMessageListener((MessageListenerConcurrently) (msgs, context) -> {
            System.out.printf("%s %s Receive New Messages:%n"
                    , sdf.format(new Date())
                    , Thread.currentThread().getName());
            for (MessageExt msg : msgs) {
                System.out.printf("\tMsg Id: %s%n", msg.getMsgId());
                System.out.printf("\tBody: %s%n", new String(msg.getBody()));
            }
            // 標(biāo)記該消息已經(jīng)被成功消費(fèi)
            return ConsumeConcurrentlyStatus.CONSUME_SUCCESS;
        });
        // 啟動(dòng)消費(fèi)者實(shí)例
        consumer.start();
        System.out.println("Consumer Started.");
    }
}

再寫一個(gè)延遲消息的生產(chǎn)者,用于發(fā)送延遲消息:

public class DelayProducer {
    public static void main(String[] args) throws Exception {
        SimpleDateFormat sdf = new SimpleDateFormat("HH:mm:ss.SSS");
        // 實(shí)例化消息生產(chǎn)者Producer
        DefaultMQProducer producer = new DefaultMQProducer("OneMoreGroup");
        // 設(shè)置NameServer的地址
        producer.setNamesrvAddr("localhost:9876");
        // 啟動(dòng)Producer實(shí)例
        producer.start();
        Message msg = new Message("OneMoreTopic"
                , "DelayMessage", "This is a delay message.".getBytes());
        //"1s 5s 10s 30s 1m 2m 3m 4m 5m 6m 7m 8m 9m 10m 20m 30m 1h 2h"
        //設(shè)置消息延遲級(jí)別為3,也就是延遲10s。
        msg.setDelayTimeLevel(3);
        // 發(fā)送消息到一個(gè)Broker
        SendResult sendResult = producer.send(msg);
        // 通過(guò)sendResult返回消息是否成功送達(dá)
        System.out.printf("%s Send Status: %s, Msg Id: %s %n"
                , sdf.format(new Date())
                , sendResult.getSendStatus()
                , sendResult.getMsgId());
        // 如果不再發(fā)送消息,關(guān)閉Producer實(shí)例。
        producer.shutdown();
    }
}

運(yùn)行生產(chǎn)者以后,就會(huì)發(fā)送一條延遲消息:

10:37:14.992 Send Status: SEND_OK, Msg Id: C0A8006D5AB018B4AAC216E0DB690000

10秒鐘后,消費(fèi)者收到的這條延遲消息:

10:37:25.026 ConsumeMessageThread_1 Receive New Messages:
	Msg Id: C0A8006D5AB018B4AAC216E0DB690000
	Body: This is a delay message.

延遲消息的原理分析

以下分析的RocketMQ源碼的版本號(hào)是4.7.1,版本不同源碼略有差別。

CommitLog

在org.apache.rocketmq.store.CommitLog中,針對(duì)延遲消息做了一些處理:

// 延遲級(jí)別大于0,就是延時(shí)消息
if (msg.getDelayTimeLevel() > 0) {
    // 判斷當(dāng)前延遲級(jí)別,如果大于最大延遲級(jí)別,
    // 就設(shè)置當(dāng)前延遲級(jí)別為最大延遲級(jí)別。
    if (msg.getDelayTimeLevel() > this.defaultMessageStore
            .getScheduleMessageService().getMaxDelayLevel()) {
        msg.setDelayTimeLevel(this.defaultMessageStore
                .getScheduleMessageService().getMaxDelayLevel());
    }
    // 獲取延遲消息的主題,
    // 其中RMQ_SYS_SCHEDULE_TOPIC的值為SCHEDULE_TOPIC_XXXX
    topic = TopicValidator.RMQ_SYS_SCHEDULE_TOPIC;
    // 根據(jù)延遲級(jí)別獲取延遲消息的隊(duì)列Id,
    // 隊(duì)列Id其實(shí)就是延遲級(jí)別減1
    queueId = ScheduleMessageService.delayLevel2QueueId(msg.getDelayTimeLevel());
    // 備份真正的主題和隊(duì)列Id
    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()));
    // 設(shè)置延時(shí)消息的主題和隊(duì)列Id
    msg.setTopic(topic);
    msg.setQueueId(queueId);
}

可以看到,每一個(gè)延遲消息的主題都被暫時(shí)更改為SCHEDULE_TOPIC_XXXX,并且根據(jù)延遲級(jí)別延遲消息變更了新的隊(duì)列Id。接下來(lái),處理延遲消息的就是org.apache.rocketmq.store.schedule.ScheduleMessageService。

ScheduleMessageService

ScheduleMessageService是由org.apache.rocketmq.store.DefaultMessageStore進(jìn)行初始化的,初始化包括構(gòu)造對(duì)象和調(diào)用 load 方法。最后,再執(zhí)行ScheduleMessageService的 start 方法:

public void start() {
    // 使用AtomicBoolean確保start方法僅有效執(zhí)行一次
    if (started.compareAndSet(false, true)) {
        this.timer = new Timer("ScheduleMessageTimerThread", true);
        // 遍歷所有延遲級(jí)別
        for (Map.Entry<Integer, Long> entry : this.delayLevelTable.entrySet()) {
            // key為延遲級(jí)別
            Integer level = entry.getKey();
            // value為延遲級(jí)別對(duì)應(yīng)的毫秒數(shù)
            Long timeDelay = entry.getValue();
            // 根據(jù)延遲級(jí)別獲得對(duì)應(yīng)隊(duì)列的偏移量
            Long offset = this.offsetTable.get(level);
            // 如果偏移量為null,則設(shè)置為0
            if (null == offset) {
                offset = 0L;
            }
            if (timeDelay != null) {
                // 為每個(gè)延遲級(jí)別創(chuàng)建定時(shí)任務(wù),
                // 第一次啟動(dòng)任務(wù)延遲為FIRST_DELAY_TIME,也就是1秒
                this.timer.schedule(
                        new DeliverDelayedMessageTimerTask(level, offset), FIRST_DELAY_TIME);
            }
        }
        // 延遲10秒后每隔flushDelayOffsetInterval執(zhí)行一次任務(wù),
        // 其中,flushDelayOffsetInterval默認(rèn)配置也為10秒
        this.timer.scheduleAtFixedRate(new TimerTask() {
            @Override
            public void run() {
                try {
                    // 持久化每個(gè)隊(duì)列消費(fèi)的偏移量
                    if (started.get()) ScheduleMessageService.this.persist();
                } catch (Throwable e) {
                    log.error("scheduleAtFixedRate flush exception", e);
                }
            }
        }, 10000, this.defaultMessageStore
        	.getMessageStoreConfig().getFlushDelayOffsetInterval());
    }
}

遍歷所有延遲級(jí)別,根據(jù)延遲級(jí)別獲得對(duì)應(yīng)隊(duì)列的偏移量,如果偏移量不存在,則設(shè)置為0。然后為每個(gè)延遲級(jí)別創(chuàng)建定時(shí)任務(wù),第一次啟動(dòng)任務(wù)延遲為1秒,第二次及以后的啟動(dòng)任務(wù)延遲才是延遲級(jí)別相應(yīng)的延遲時(shí)間。

然后,又創(chuàng)建了一個(gè)定時(shí)任務(wù),用于持久化每個(gè)隊(duì)列消費(fèi)的偏移量。持久化的頻率由flushDelayOffsetInterval屬性進(jìn)行配置,默認(rèn)為10秒。

定時(shí)任務(wù)

ScheduleMessageService的 start 方法執(zhí)行之后,每個(gè)延遲級(jí)別都創(chuàng)建自己的定時(shí)任務(wù),這里的定時(shí)任務(wù)的具體實(shí)現(xiàn)就在DeliverDelayedMessageTimerTask類之中,它核心代碼是executeOnTimeup方法之中,我們來(lái)看一下主要部分:

// 根據(jù)主題和隊(duì)列Id獲取消息隊(duì)列
ConsumeQueue cq =
        ScheduleMessageService.this.defaultMessageStore.findConsumeQueue(
                TopicValidator.RMQ_SYS_SCHEDULE_TOPIC
                , delayLevel2QueueId(delayLevel));

如果沒(méi)有獲取到對(duì)應(yīng)的消息隊(duì)列,則在DELAY_FOR_A_WHILE(默認(rèn)為100)毫秒后再執(zhí)行任務(wù)。如果獲取到了,就繼續(xù)執(zhí)行下面操作:

// 根據(jù)消費(fèi)偏移量從消息隊(duì)列中獲取所有有效消息
SelectMappedBufferResult bufferCQ = cq.getIndexBuffer(this.offset);

如果沒(méi)有獲取到有效消息,則在DELAY_FOR_A_WHILE(默認(rèn)為100)毫秒后再執(zhí)行任務(wù)。如果獲取到了,就繼續(xù)執(zhí)行下面操作:

// 遍歷所有消息
for (; i < bufferCQ.getSize(); i += ConsumeQueue.CQ_STORE_UNIT_SIZE) {
    // 獲取消息的物理偏移量
    long offsetPy = bufferCQ.getByteBuffer().getLong();
    // 獲取消息的物理長(zhǎng)度
    int sizePy = bufferCQ.getByteBuffer().getInt();
    long tagsCode = bufferCQ.getByteBuffer().getLong();
    // 省略部分代碼...
    long now = System.currentTimeMillis();
    // 計(jì)算消息應(yīng)該被消費(fèi)的時(shí)間
    long deliverTimestamp = this.correctDeliverTimestamp(now, tagsCode);
	// 計(jì)算下一條消息的偏移量
    nextOffset = offset + (i / ConsumeQueue.CQ_STORE_UNIT_SIZE)
	long countdown = deliverTimestamp - now;
    // 省略部分代碼...
}

如果當(dāng)前消息不到消費(fèi)的時(shí)間,則在 countdown 毫秒后再執(zhí)行任務(wù)。如果到消費(fèi)的時(shí)間,就繼續(xù)執(zhí)行下面操作:

// 根據(jù)消息的物理偏移量和大小獲取消息
MessageExt msgExt =
    ScheduleMessageService.this.defaultMessageStore.lookMessageByOffset(
            offsetPy, sizePy);

如果獲取到消息,則繼續(xù)執(zhí)行下面操作:

// 重新構(gòu)建新的消息,包括:
// 1.清除消息的延遲級(jí)別
// 2.恢復(fù)真正的消息主題和隊(duì)列Id
MessageExtBrokerInner msgInner = this.messageTimeup(msgExt);
if (TopicValidator.RMQ_SYS_TRANS_HALF_TOPIC.equals(msgInner.getTopic())) {
    log.error("[BUG] the real topic of schedule msg is {},"
                    + " discard the msg. msg={}",
            msgInner.getTopic(), msgInner);
    continue;
}
// 重新把消息發(fā)送到真正的消息隊(duì)列上
PutMessageResult putMessageResult =
        ScheduleMessageService.this.writeMessageStore
                .putMessage(msgInner);

清除了消息的延遲級(jí)別,并且恢復(fù)了真正的消息主題和隊(duì)列Id,重新把消息發(fā)送到真正的消息隊(duì)列上以后,消費(fèi)者就可以立即消費(fèi)了。

總結(jié)

經(jīng)過(guò)以上對(duì)源碼的分析,可以總結(jié)出延遲消息的實(shí)現(xiàn)步驟:

如果消息的延遲級(jí)別大于0,則表示該消息為延遲消息,修改該消息的主題為SCHEDULE_TOPIC_XXXX,隊(duì)列Id為延遲級(jí)別減1。消息進(jìn)入SCHEDULE_TOPIC_XXXX的隊(duì)列中。定時(shí)任務(wù)根據(jù)上次拉取的偏移量不斷從隊(duì)列中取出所有消息。根據(jù)消息的物理偏移量和大小再次獲取消息。根據(jù)消息屬性重新創(chuàng)建消息,清除延遲級(jí)別,恢復(fù)原主題和隊(duì)列Id。重新發(fā)送消息到原主題的隊(duì)列中,供消費(fèi)者進(jìn)行消費(fèi)。

到此這篇關(guān)于Java中RocketMQ的延遲消息詳解的文章就介紹到這了,更多相關(guān)RocketMQ的延遲消息內(nèi)容請(qǐng)搜索腳本之家以前的文章或繼續(xù)瀏覽下面的相關(guān)文章希望大家以后多多支持腳本之家!

相關(guān)文章

  • SpringBoot配置文件密碼加密與解密的操作代碼

    SpringBoot配置文件密碼加密與解密的操作代碼

    我們?cè)赟pringBoot項(xiàng)目當(dāng)中,會(huì)把數(shù)據(jù)庫(kù)的用戶名密碼等配置直接放在yaml或者properties文件中,這樣維護(hù)數(shù)據(jù)庫(kù)的密碼等敏感信息顯然是有一定風(fēng)險(xiǎn)的,所以我們需要給密碼加解密,本文給大家介紹了SpringBoot配置文件密碼加密與解密的操作代碼,需要的朋友可以參考下
    2024-12-12
  • java中創(chuàng)建寫入文件的6種方式詳解與源碼實(shí)例

    java中創(chuàng)建寫入文件的6種方式詳解與源碼實(shí)例

    這篇文章主要介紹了java中創(chuàng)建寫入文件的6種方式詳解與源碼實(shí)例,Files.newBufferedWriter(Java 8),Files.write(Java 7 推薦),PrintWriter,File.createNewFile,FileOutputStream.write(byte[] b) 管道流,需要的朋友可以參考下
    2022-12-12
  • MyBatis與其使用方法示例詳解

    MyBatis與其使用方法示例詳解

    MyBatis是一個(gè)支持自定義SQL的持久層框架,通過(guò)XML文件實(shí)現(xiàn)SQL配置和數(shù)據(jù)映射,簡(jiǎn)化了JDBC代碼的編寫,本文給大家介紹MyBatis與其使用方法講解,感興趣的朋友一起看看吧
    2025-03-03
  • Kotlin開發(fā)Android應(yīng)用實(shí)例詳解

    Kotlin開發(fā)Android應(yīng)用實(shí)例詳解

    這篇文章主要介紹了Kotlin開發(fā)Android應(yīng)用實(shí)例詳解的相關(guān)資料,需要的朋友可以參考下
    2017-05-05
  • java實(shí)現(xiàn)動(dòng)態(tài)數(shù)組

    java實(shí)現(xiàn)動(dòng)態(tài)數(shù)組

    這篇文章主要為大家詳細(xì)介紹了java實(shí)現(xiàn)動(dòng)態(tài)數(shù)組,文中示例代碼介紹的非常詳細(xì),具有一定的參考價(jià)值,感興趣的小伙伴們可以參考一下
    2021-08-08
  • 部署springboot打包不打包配置文件,配置文件為外部配置文件使用詳解

    部署springboot打包不打包配置文件,配置文件為外部配置文件使用詳解

    在Spring Boot項(xiàng)目中,將配置文件排除在jar包之外,通過(guò)外部配置文件來(lái)管理不同環(huán)境的配置,可以實(shí)現(xiàn)靈活的配置管理,在pom.xml文件中添加相關(guān)配置,打包時(shí)忽略指定文件,運(yùn)行時(shí)在jar包同級(jí)目錄下創(chuàng)建config文件夾,將配置文件放入其中即可
    2025-02-02
  • 使用Nacos作為配置中心的命名空間、配置分組

    使用Nacos作為配置中心的命名空間、配置分組

    文章詳細(xì)介紹了Spring Cloud Config配置中心的命名空間、配置集、配置集ID、配置分組以及如何在微服務(wù)中加載和使用這些配置,通過(guò)配置中心,可以實(shí)現(xiàn)配置隔離和集中管理,簡(jiǎn)化微服務(wù)的配置維護(hù)
    2024-12-12
  • SWT(JFace) 簡(jiǎn)易瀏覽器 制作實(shí)現(xiàn)代碼

    SWT(JFace) 簡(jiǎn)易瀏覽器 制作實(shí)現(xiàn)代碼

    SWT(JFace) 簡(jiǎn)易瀏覽器 制作實(shí)現(xiàn)代碼
    2009-06-06
  • Maven編譯時(shí)出現(xiàn)中文亂碼的完整解決教程

    Maven編譯時(shí)出現(xiàn)中文亂碼的完整解決教程

    這篇文章主要為大家詳細(xì)介紹了Maven編譯時(shí)出現(xiàn)中文亂碼的完整解決教程,適合Windows 用戶,Maven 編譯項(xiàng)目時(shí)終端或 IDEA 出現(xiàn)中文亂碼的人群,快跟隨小編一起了解下吧
    2025-10-10
  • Java為圖片添加水印并保存實(shí)現(xiàn)方法(附帶源碼)

    Java為圖片添加水印并保存實(shí)現(xiàn)方法(附帶源碼)

    這篇文章主要介紹了如何使用Java編程語(yǔ)言在圖像上添加文字或圖片水印,并提供了一個(gè)簡(jiǎn)單的Java程序?qū)崿F(xiàn),文中給出了詳細(xì)的代碼示例,需要的朋友可以參考下
    2025-03-03

最新評(píng)論

贵南县| 泸溪县| 清流县| 常宁市| 青田县| 沾化县| 蒙阴县| 丰顺县| 华坪县| 永济市| 嘉义县| 昭苏县| 大同市| 海晏县| 沙洋县| 抚顺县| 广河县| 宜丰县| 安国市| 南部县| 舞钢市| 鹿泉市| 织金县| 油尖旺区| 甘肃省| 泸州市| 滦南县| 邹平县| 宜丰县| 常山县| 兴宁市| 天气| 丰台区| 池州市| 花垣县| 陆丰市| 英山县| 连平县| 揭西县| 顺平县| 湘西|