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

RocketMQ設(shè)計(jì)之同步刷盤

 更新時(shí)間:2022年03月21日 10:17:40   作者:周杰倫本人  
這篇文章主要介紹了RocketMQ設(shè)計(jì)之同步刷盤,文章主要通過CommitLog的handleDiskFlush方法展開全文內(nèi)容,實(shí)現(xiàn)同步刷盤,下面文章詳細(xì)介紹,需要的小伙伴可以參考一下

同步刷盤方式:在返回寫成功狀態(tài)時(shí),消息已經(jīng)被寫入磁盤。具體流程是,消息寫入內(nèi)存的PAGECACHE后,立刻通知刷盤線程刷盤,然后等待刷盤完成,刷盤線程執(zhí)行完成后喚醒等待的線程,返回消息寫成功的狀態(tài)。

在同步刷盤模式下,當(dāng)消息寫到內(nèi)存后,會(huì)等待數(shù)據(jù)寫到磁盤的CommitLog文件。

CommitLog的handleDiskFlush方法:

public void handleDiskFlush(AppendMessageResult result, PutMessageResult putMessageResult, MessageExt messageExt) {
? ? // Synchronization flush
? ? if (FlushDiskType.SYNC_FLUSH == this.defaultMessageStore.getMessageStoreConfig().getFlushDiskType()) {
? ? ? ? final GroupCommitService service = (GroupCommitService) this.flushCommitLogService;
? ? ? ? if (messageExt.isWaitStoreMsgOK()) {
? ? ? ? ? ? GroupCommitRequest request = new GroupCommitRequest(result.getWroteOffset() + result.getWroteBytes());
? ? ? ? ? ? service.putRequest(request);
? ? ? ? ? ? boolean flushOK = request.waitForFlush(this.defaultMessageStore.getMessageStoreConfig().getSyncFlushTimeout());
? ? ? ? ? ? if (!flushOK) {
? ? ? ? ? ? ? ? log.error("do groupcommit, wait for flush failed, topic: " + messageExt.getTopic() + " tags: " + messageExt.getTags()
? ? ? ? ? ? ? ? ? ? + " client address: " + messageExt.getBornHostString());
? ? ? ? ? ? ? ? putMessageResult.setPutMessageStatus(PutMessageStatus.FLUSH_DISK_TIMEOUT);
? ? ? ? ? ? }
? ? ? ? } else {
? ? ? ? ? ? service.wakeup();
? ? ? ? }
? ? }
? ? // Asynchronous flush
? ? else {
? ? ? ? if (!this.defaultMessageStore.getMessageStoreConfig().isTransientStorePoolEnable()) {
? ? ? ? ? ? flushCommitLogService.wakeup();
? ? ? ? } else {
? ? ? ? ? ? commitLogService.wakeup();
? ? ? ? }
? ? }
}


class GroupCommitService extends FlushCommitLogService {
? ? ? ? private volatile List<GroupCommitRequest> requestsWrite = new ArrayList<GroupCommitRequest>();
? ? ? ? private volatile List<GroupCommitRequest> requestsRead = new ArrayList<GroupCommitRequest>();

? ? ?? ?//提交刷盤任務(wù)到任務(wù)列表
? ? ? ? public synchronized void putRequest(final GroupCommitRequest request) {
? ? ? ? ? ? synchronized (this.requestsWrite) {
? ? ? ? ? ? ? ? this.requestsWrite.add(request);
? ? ? ? ? ? }
? ? ? ? ? ? if (hasNotified.compareAndSet(false, true)) {
? ? ? ? ? ? ? ? waitPoint.countDown(); // notify
? ? ? ? ? ? }
? ? ? ? }

? ? ? ? private void swapRequests() {
? ? ? ? ? ? List<GroupCommitRequest> tmp = this.requestsWrite;
? ? ? ? ? ? this.requestsWrite = this.requestsRead;
? ? ? ? ? ? this.requestsRead = tmp;
? ? ? ? }

? ? ? ? private void doCommit() {
? ? ? ? ? ? synchronized (this.requestsRead) {
? ? ? ? ? ? ? ? if (!this.requestsRead.isEmpty()) {
? ? ? ? ? ? ? ? ? ? for (GroupCommitRequest req : this.requestsRead) {
? ? ? ? ? ? ? ? ? ? ? ? // There may be a message in the next file, so a maximum of
? ? ? ? ? ? ? ? ? ? ? ? // two times the flush
? ? ? ? ? ? ? ? ? ? ? ? boolean flushOK = false;
? ? ? ? ? ? ? ? ? ? ? ? for (int i = 0; i < 2 && !flushOK; i++) {
? ? ? ? ? ? ? ? ? ? ? ? ? ? flushOK = CommitLog.this.mappedFileQueue.getFlushedWhere() >= req.getNextOffset();

? ? ? ? ? ? ? ? ? ? ? ? ? ? if (!flushOK) {
? ? ? ? ? ? ? ? ? ? ? ? ? ? ? ? CommitLog.this.mappedFileQueue.flush(0);
? ? ? ? ? ? ? ? ? ? ? ? ? ? }
? ? ? ? ? ? ? ? ? ? ? ? }

? ? ? ? ? ? ? ? ? ? ? ? req.wakeupCustomer(flushOK);
? ? ? ? ? ? ? ? ? ? }

? ? ? ? ? ? ? ? ? ? long storeTimestamp = CommitLog.this.mappedFileQueue.getStoreTimestamp();
? ? ? ? ? ? ? ? ? ? if (storeTimestamp > 0) {
? ? ? ? ? ? ? ? ? ? ? ? CommitLog.this.defaultMessageStore.getStoreCheckpoint().setPhysicMsgTimestamp(storeTimestamp);
? ? ? ? ? ? ? ? ? ? }

? ? ? ? ? ? ? ? ? ? this.requestsRead.clear();
? ? ? ? ? ? ? ? } else {
? ? ? ? ? ? ? ? ? ? // Because of individual messages is set to not sync flush, it
? ? ? ? ? ? ? ? ? ? // will come to this process
? ? ? ? ? ? ? ? ? ? CommitLog.this.mappedFileQueue.flush(0);
? ? ? ? ? ? ? ? }
? ? ? ? ? ? }
? ? ? ? }

? ? ? ? public void run() {
? ? ? ? ? ? CommitLog.log.info(this.getServiceName() + " service started");

? ? ? ? ? ? while (!this.isStopped()) {
? ? ? ? ? ? ? ? try {
? ? ? ? ? ? ? ? ? ? this.waitForRunning(10);
? ? ? ? ? ? ? ? ? ? this.doCommit();
? ? ? ? ? ? ? ? } catch (Exception e) {
? ? ? ? ? ? ? ? ? ? CommitLog.log.warn(this.getServiceName() + " service has exception. ", e);
? ? ? ? ? ? ? ? }
? ? ? ? ? ? }

? ? ? ? ? ? // Under normal circumstances shutdown, wait for the arrival of the
? ? ? ? ? ? // request, and then flush
? ? ? ? ? ? try {
? ? ? ? ? ? ? ? Thread.sleep(10);
? ? ? ? ? ? } catch (InterruptedException e) {
? ? ? ? ? ? ? ? CommitLog.log.warn("GroupCommitService Exception, ", e);
? ? ? ? ? ? }

? ? ? ? ? ? synchronized (this) {
? ? ? ? ? ? ? ? this.swapRequests();
? ? ? ? ? ? }

? ? ? ? ? ? this.doCommit();

? ? ? ? ? ? CommitLog.log.info(this.getServiceName() + " service end");
? ? ? ? }

? ? ? ? @Override
? ? ? ? protected void onWaitEnd() {
? ? ? ? ? ? this.swapRequests();
? ? ? ? }

? ? ? ? @Override
? ? ? ? public String getServiceName() {
? ? ? ? ? ? return GroupCommitService.class.getSimpleName();
? ? ? ? }

? ? ? ? @Override
? ? ? ? public long getJointime() {
? ? ? ? ? ? return 1000 * 60 * 5;
? ? ? ? }
? ? }

GroupCommitRequest是刷盤任務(wù),提交刷盤任務(wù)后,會(huì)在刷盤隊(duì)列中等待刷盤,而刷盤線程

GroupCommitService每隔10毫秒寫一批數(shù)據(jù)到磁盤。之所以不直接寫是磁盤io壓力大,寫入性能低,每隔10毫秒寫一次可以提升磁盤io效率和寫入性能。

  • putRequest(request) 提交刷盤任務(wù)到任務(wù)列表
  • request.waitForFlush同步等待GroupCommitService將任務(wù)列表中的任務(wù)刷盤完成。

兩個(gè)隊(duì)列讀寫分離,requestsWrite是寫隊(duì)列,用戶保存添加進(jìn)來的刷盤任務(wù),requestsRead是讀隊(duì)列,在刷盤之前會(huì)把寫隊(duì)列的數(shù)據(jù)放入讀隊(duì)列。

CommitLog的doCommit方法:

private void doCommit() {
? ? ? ? ? ? synchronized (this.requestsRead) {
? ? ? ? ? ? ? ? if (!this.requestsRead.isEmpty()) {
? ? ? ? ? ? ? ? ? ? for (GroupCommitRequest req : this.requestsRead) {
? ? ? ? ? ? ? ? ? ? ? ? // There may be a message in the next file, so a maximum of
? ? ? ? ? ? ? ? ? ? ? ? // two times the flush
? ? ? ? ? ? ? ? ? ? ? ? boolean flushOK = false;
? ? ? ? ? ? ? ? ? ? ? ? for (int i = 0; i < 2 && !flushOK; i++) {
? ? ? ? ? ? ? ? ? ? ? ? ? ? //根據(jù)offset確定是否已經(jīng)刷盤
? ? ? ? ? ? ? ? ? ? ? ? ? ? flushOK = CommitLog.this.mappedFileQueue.getFlushedWhere() >= req.getNextOffset();

? ? ? ? ? ? ? ? ? ? ? ? ? ? if (!flushOK) {
? ? ? ? ? ? ? ? ? ? ? ? ? ? ? ? CommitLog.this.mappedFileQueue.flush(0);
? ? ? ? ? ? ? ? ? ? ? ? ? ? }
? ? ? ? ? ? ? ? ? ? ? ? }

? ? ? ? ? ? ? ? ? ? ? ? req.wakeupCustomer(flushOK);
? ? ? ? ? ? ? ? ? ? }

? ? ? ? ? ? ? ? ? ? long storeTimestamp = CommitLog.this.mappedFileQueue.getStoreTimestamp();
? ? ? ? ? ? ? ? ? ? if (storeTimestamp > 0) {
? ? ? ? ? ? ? ? ? ? ? ? CommitLog.this.defaultMessageStore.getStoreCheckpoint().setPhysicMsgTimestamp(storeTimestamp);
? ? ? ? ? ? ? ? ? ? }
?? ??? ??? ??? ??? ?//清空已刷盤的列表
? ? ? ? ? ? ? ? ? ? this.requestsRead.clear();
? ? ? ? ? ? ? ? } else {
? ? ? ? ? ? ? ? ? ? // Because of individual messages is set to not sync flush, it
? ? ? ? ? ? ? ? ? ? // will come to this process
? ? ? ? ? ? ? ? ? ? CommitLog.this.mappedFileQueue.flush(0);
? ? ? ? ? ? ? ? }
? ? ? ? ? ? }
? ? ? ? }
  • 刷盤的時(shí)候依次讀取requestsRead中的數(shù)據(jù)寫入磁盤,
  • 寫入完成后清空requestsRead。

讀寫分離設(shè)計(jì)的目的是在刷盤時(shí)不影響任務(wù)提交到列表。

CommitLog.this.mappedFileQueue.flush(0);是刷盤操作:

public boolean flush(final int flushLeastPages) {
? ? boolean result = true;
? ? MappedFile mappedFile = this.findMappedFileByOffset(this.flushedWhere, this.flushedWhere == 0);
? ? if (mappedFile != null) {
? ? ? ? long tmpTimeStamp = mappedFile.getStoreTimestamp();
? ? ? ? int offset = mappedFile.flush(flushLeastPages);
? ? ? ? long where = mappedFile.getFileFromOffset() + offset;
? ? ? ? result = where == this.flushedWhere;
? ? ? ? this.flushedWhere = where;
? ? ? ? if (0 == flushLeastPages) {
? ? ? ? ? ? this.storeTimestamp = tmpTimeStamp;
? ? ? ? }
? ? }

? ? return result;
}

通過MappedFile映射的CommitLog文件寫入磁盤

這就是RocketMQ高可用設(shè)計(jì)之同步刷盤的基本情況了,大體思路就是一個(gè)讀寫分離的隊(duì)列來刷盤,同步刷盤任務(wù)提交后會(huì)在刷盤隊(duì)列中等待刷盤完成后再返回,而GroupCommitService每隔10毫秒寫一批數(shù)據(jù)到磁盤。

到此這篇關(guān)于RocketMQ設(shè)計(jì)之同步刷盤的文章就介紹到這了,更多相關(guān)RocketMQ同步刷盤內(nèi)容請(qǐng)搜索腳本之家以前的文章或繼續(xù)瀏覽下面的相關(guān)文章希望大家以后多多支持腳本之家!

相關(guān)文章

  • Java如何定義Long類型

    Java如何定義Long類型

    這篇文章主要介紹了Java如何定義Long類型,具有很好的參考價(jià)值,希望對(duì)大家有所幫助。如有錯(cuò)誤或未考慮完全的地方,望不吝賜教
    2022-07-07
  • maven鏡像倉(cāng)庫(kù)的配置過程

    maven鏡像倉(cāng)庫(kù)的配置過程

    本文詳細(xì)介紹了MAVEN_HOME的配置步驟、Path環(huán)境變量的設(shè)置、檢測(cè)配置是否成功的方法、修改默認(rèn)的maven依賴包下載路徑以及配置阿里鏡像倉(cāng)庫(kù)的路徑,同時(shí)分享了作者在配置過程中遇到的問題,如命令不識(shí)別、版本不匹配等,并提供了解決方案
    2024-09-09
  • 一文帶你了解SpringBoot的啟動(dòng)原理

    一文帶你了解SpringBoot的啟動(dòng)原理

    大家通常只需要給一個(gè)類添加一個(gè)@SpringBootApplication 注解,然后再加一個(gè)main 方法里面固定的寫法 SpringApplication.run(Application.class, args);那么spring boot 到底是如何啟動(dòng)服務(wù)的呢,接下來咱們通過源碼解析,需要的朋友可以參考下
    2023-05-05
  • Java集合的組內(nèi)平均值的計(jì)算方法總結(jié)

    Java集合的組內(nèi)平均值的計(jì)算方法總結(jié)

    在Java中,經(jīng)常需要對(duì)集合進(jìn)行各種操作,其中之一就是計(jì)算集合的組內(nèi)平均值,本文將介紹如何使用Java集合來計(jì)算組內(nèi)平均值,并提供一些示例代碼和實(shí)用技巧
    2024-08-08
  • 關(guān)于IDEA2020.1新建項(xiàng)目maven PKIX 報(bào)錯(cuò)問題解決方法

    關(guān)于IDEA2020.1新建項(xiàng)目maven PKIX 報(bào)錯(cuò)問題解決方法

    這篇文章主要介紹了關(guān)于IDEA2020.1新建項(xiàng)目maven PKIX 報(bào)錯(cuò)問題解決方法,文中通過示例代碼介紹的非常詳細(xì),對(duì)大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價(jià)值,需要的朋友們下面隨著小編來一起學(xué)習(xí)學(xué)習(xí)吧
    2020-06-06
  • java中rss解析器(rome.jar和jdom.jar)示例

    java中rss解析器(rome.jar和jdom.jar)示例

    這篇文章主要介紹了java中rss解析器(rome.jar和jdom.jar)示例,需要的朋友可以參考下
    2014-03-03
  • java WebSocket實(shí)現(xiàn)聊天消息推送功能

    java WebSocket實(shí)現(xiàn)聊天消息推送功能

    這篇文章主要為大家詳細(xì)介紹了java WebSocket實(shí)現(xiàn)聊天消息推送功能,具有一定的參考價(jià)值,感興趣的小伙伴們可以參考一下
    2018-07-07
  • 詳解Java如何實(shí)現(xiàn)與JS相同的Des加解密算法

    詳解Java如何實(shí)現(xiàn)與JS相同的Des加解密算法

    這篇文章主要介紹了如何在Java中實(shí)現(xiàn)與JavaScript相同的DES(Data Encryption Standard)加解密算法,確保在兩個(gè)平臺(tái)之間可以無縫地傳遞加密信息,希望對(duì)大家有一定的幫助
    2025-04-04
  • Java中Integer和int的區(qū)別解讀

    Java中Integer和int的區(qū)別解讀

    這篇文章主要介紹了Java中Integer和int的區(qū)別解讀,大家都知道他可以表示一個(gè)整數(shù),而且也知道可以表示整數(shù)的還有int,只是使用Integer的次數(shù)要比int多得多,今天我們就來好好探究一下Integer與int的區(qū)別以及更深處的知識(shí),需要的朋友可以參考下
    2023-12-12
  • 深入淺出重構(gòu)Mybatis與Spring集成的SqlSessionFactoryBean(上)

    深入淺出重構(gòu)Mybatis與Spring集成的SqlSessionFactoryBean(上)

    通常來講,重構(gòu)是指不改變功能的情況下優(yōu)化代碼,但本文所說的重構(gòu)也包括了添加功能。這篇文章主要介紹了重構(gòu)Mybatis與Spring集成的SqlSessionFactoryBean(上)的相關(guān)資料,需要的朋友可以參考下
    2016-11-11

最新評(píng)論

长治县| 昌乐县| 年辖:市辖区| 康马县| 崇州市| 垣曲县| 汤阴县| 乌审旗| 西充县| 徐州市| 项城市| 东明县| 古丈县| 荣昌县| 麻栗坡县| 临潭县| 万安县| 广饶县| 安多县| 嵩明县| 正安县| 金川县| 烟台市| 上高县| 金坛市| 临西县| 天峻县| 高邑县| 和硕县| 贡觉县| 邹平县| 法库县| 额尔古纳市| 朝阳区| 鹤庆县| 夏津县| 辉县市| 大同县| 益阳市| 蓬莱市| 大方县|