精品视频在线免费观看_国产精品资源网_欧美日韩亚洲综合在线_自拍视频国产精品

原創(chuàng)生活

國(guó)內(nèi) 商業(yè) 滾動(dòng)

基金 金融 股票

期貨金融

科技 行業(yè) 房產(chǎn)

銀行 公司 消費(fèi)

生活滾動(dòng)

保險(xiǎn) 海外 觀察

財(cái)經(jīng) 生活 期貨

當(dāng)前位置:原創(chuàng) >

RocketMQ 多級(jí)存儲(chǔ)設(shè)計(jì)與實(shí)現(xiàn)-熱資訊

文章來(lái)源:阿里開發(fā)者  發(fā)布時(shí)間: 2023-04-21 04:57:23  責(zé)任編輯:cfenews.com
+|-

1000w云上開發(fā)者 全棧云產(chǎn)品0元試用:點(diǎn)擊

作者:張森澤


(相關(guān)資料圖)

隨著 RocketMQ 5.1.0 的正式發(fā)布,多級(jí)存儲(chǔ)作為 RocketMQ 一個(gè)新的獨(dú)立模塊到達(dá)了 Technical Preview 里程碑:允許用戶將消息從本地磁盤卸載到其他更便宜的存儲(chǔ)介質(zhì),可以用較低的成本延長(zhǎng)消息保留時(shí)間。本文詳細(xì)介紹 RocketMQ 多級(jí)存儲(chǔ)設(shè)計(jì)與實(shí)現(xiàn)。

設(shè)計(jì)總覽

RocketMQ 多級(jí)存儲(chǔ)旨在 不影響熱數(shù)據(jù)讀寫的前提下 將數(shù)據(jù)卸載到其他存儲(chǔ)介質(zhì)中,適用于兩種場(chǎng)景:

冷熱數(shù)據(jù)分離:RocketMQ 新近產(chǎn)生的消息會(huì)緩存在 page cache 中,我們稱之為 熱數(shù)據(jù) ;當(dāng)緩存超過(guò)了內(nèi)存的容量就會(huì)有熱數(shù)據(jù)被換出成為 冷數(shù)據(jù) 。如果有少許消費(fèi)者嘗試消費(fèi)冷數(shù)據(jù)就會(huì)從硬盤中重新加載冷數(shù)據(jù)到 page cache,這會(huì)導(dǎo)致讀寫 IO 競(jìng)爭(zhēng)并擠壓 page cache 的空間。而將冷數(shù)據(jù)的讀取鏈路切換為多級(jí)存儲(chǔ)就可以避免這個(gè)問題; 延長(zhǎng)消息保留時(shí)間:將消息卸載到更大更便宜的存儲(chǔ)介質(zhì)中,可以用較低的成本實(shí)現(xiàn)更長(zhǎng)的消息保存時(shí)間。同時(shí)多級(jí)存儲(chǔ)支持為 topic 指定不同的消息保留時(shí)間,可以根據(jù)業(yè)務(wù)需要靈活配置消息 TTL。

RocketMQ 多級(jí)存儲(chǔ)對(duì)比 Kafka 和 Pulsar 的實(shí)現(xiàn)最大的不同是我們使用準(zhǔn)實(shí)時(shí)的方式上傳消息,而不是等一個(gè) CommitLog 寫滿后再上傳,主要基于以下幾點(diǎn)考慮:

均攤成本:RocketMQ 多級(jí)存儲(chǔ)需要將全局 CommitLog 轉(zhuǎn)換為 topic 維度并重新構(gòu)建消息索引,一次性處理整個(gè) CommitLog 文件會(huì)帶來(lái)性能毛刺; 對(duì)小規(guī)格實(shí)例更友好:小規(guī)格實(shí)例往往配置較小的內(nèi)存,這意味著熱數(shù)據(jù)會(huì)更快換出成為冷數(shù)據(jù),等待 CommitLog 寫滿再上傳本身就有冷讀風(fēng)險(xiǎn)。采取準(zhǔn)實(shí)時(shí)上傳的方式既能規(guī)避消息上傳時(shí)的冷讀風(fēng)險(xiǎn),又能盡快使得冷數(shù)據(jù)可以從多級(jí)存儲(chǔ)讀取。

Quick Start

多級(jí)存儲(chǔ)在設(shè)計(jì)上希望降低用戶心智負(fù)擔(dān):用戶無(wú)需變更客戶端就能實(shí)現(xiàn)無(wú)感切換冷熱數(shù)據(jù)讀寫鏈路,通過(guò)簡(jiǎn)單的修改服務(wù)端配置即可具備多級(jí)存儲(chǔ)的能力,只需以下兩步:

修改 Broker 配置,指定使用 org.apache.rocketmq.tieredstore.TieredMessageStore 作為 messageStorePlugIn 配置你想使用的儲(chǔ)存介質(zhì),以卸載消息到其他硬盤為例:配置 tieredBackendServiceProvider 為 org.apache.rocketmq.tieredstore.provider.posix.PosixFileSegment,同時(shí)指定新儲(chǔ)存的文件路徑:tieredStoreFilepath

可選項(xiàng):支持修改 tieredMetadataServiceProvider 切換元數(shù)據(jù)存儲(chǔ)的實(shí)現(xiàn),默認(rèn)是基于 json 的文件存儲(chǔ)

更多使用說(shuō)明和配置項(xiàng)可以在 GitHub 上查看多級(jí)存儲(chǔ)的 README[1]

技術(shù)架構(gòu)

architecture

接入層 :TieredMessageStore/TieredDispatcher/TieredMessageFetcher

接入層實(shí)現(xiàn) MessageStore 中的部分讀寫接口,并為他們?cè)黾恿水惒秸Z(yǔ)意。TieredDispatcher 和 TieredMessageFetcher 分別實(shí)現(xiàn)了多級(jí)存儲(chǔ)的上傳/下載邏輯,相比于底層接口這里做了較多的性能優(yōu)化:包括使用獨(dú)立的線程池,避免慢 IO 阻塞訪問熱數(shù)據(jù);使用預(yù)讀緩存優(yōu)化性能等。

容器層 :TieredCommitLog/TieredConsumeQueue/TieredIndexFile/TieredFileQueue

容器層實(shí)現(xiàn)了和 DefaultMessageStore 類似的邏輯文件抽象,同樣將文件劃分為 CommitLog、ConsumeQueue、IndexFile,并且每種邏輯文件類型都通過(guò) FileQueue 持有底層物理文件的引用。有所不同的是多級(jí)存儲(chǔ)的 CommitLog 改為 queue 維度。

驅(qū)動(dòng)層 :TieredFileSegment

驅(qū)動(dòng)層負(fù)責(zé)維護(hù)邏輯文件到物理文件的映射,通過(guò)實(shí)現(xiàn) TieredStoreProvider 對(duì)接底層文件系統(tǒng)讀寫接口(Posix、S3、OSS、MinIO 等)。目前提供了 PosixFileSegment 的實(shí)現(xiàn),可以將數(shù)據(jù)轉(zhuǎn)移到其他硬盤或通過(guò) fuse 掛載的對(duì)象存儲(chǔ)上。

消息上傳

RocketMQ 多級(jí)存儲(chǔ)的消息上傳是由 dispatch 機(jī)制觸發(fā)的:初始化多級(jí)存儲(chǔ)時(shí)會(huì)將 TieredDispatcher 注冊(cè)為 CommitLog 的 dispacher。這樣每當(dāng)有消息發(fā)送到 Broker 會(huì)調(diào)用 TieredDispatcher 進(jìn)行消息分發(fā),TieredDispatcher 將該消息寫入到 upload buffer 后立即返回成功。整個(gè) dispatch 流程中不會(huì)有任何阻塞邏輯,確保不會(huì)影響本地 ConsumeQueue 的構(gòu)建。

TieredDispatcher

TieredDispatcher 寫入 upload buffer 的內(nèi)容僅為消息的引用,不會(huì)將消息的 body 讀入內(nèi)存。因?yàn)槎嗉?jí)儲(chǔ)存以 queue 維度構(gòu)建 CommitLog,此時(shí)需要重新生成 commitLog offset 字段。

upload buffer

觸發(fā) upload buffer 上傳時(shí)讀取到每條消息的 commitLog offset 字段時(shí)采用拼接的方式將新的 offset 嵌入到原消息中。

上傳進(jìn)度控制

每個(gè)隊(duì)列都會(huì)有兩個(gè)關(guān)鍵位點(diǎn)控制上傳進(jìn)度:

dispatch offset:已經(jīng)寫入緩存但是未上傳的消息位點(diǎn) commit offset:已上傳的消息位點(diǎn)

upload progress

類比消費(fèi)者,dispatch offset 相當(dāng)于拉取消息的位點(diǎn),commit offset 相當(dāng)于確認(rèn)消費(fèi)的位點(diǎn)。commit offset 到 dispatch offset 之間的部分相當(dāng)于已拉取未消費(fèi)的消息。

消息讀取

TieredMessageStore 實(shí)現(xiàn)了 MessageStore 中的消息讀取相關(guān)接口,通過(guò)請(qǐng)求中的邏輯位點(diǎn)(queue offset)判斷是否從多級(jí)存儲(chǔ)中讀取消息,根據(jù)配置(tieredStorageLevel)有四種策略:

DISABLE:禁止從多級(jí)存儲(chǔ)中讀取消息; NOT_IN_DISK:不在 DefaultMessageStore 中的消息從多級(jí)存儲(chǔ)中讀取; NOT_IN_MEM:不在 page cache 中的消息即冷數(shù)據(jù)從多級(jí)存儲(chǔ)讀取; FORCE:強(qiáng)制所有消息從多級(jí)存儲(chǔ)中讀取,目前僅供測(cè)試使用。
/**  * Asynchronous get message  * @see #getMessage(String, String, int, long, int, MessageFilter)   getMessage  *  * @param group Consumer group that launches this query.  * @param topic Topic to query.  * @param queueId Queue ID to query.  * @param offset Logical offset to start from.  * @param maxMsgNums Maximum count of messages to query.  * @param messageFilter Message filter used to screen desired   messages.  * @return Matched messages.  */CompletableFuturegetMessageAsync(final String group, final String topic, final int queueId,    final long offset, final int maxMsgNums, final MessageFilter messageFilter);

需要從多級(jí)存儲(chǔ)中讀取的消息會(huì)交由 TieredMessageFetcher 處理:首先校驗(yàn)參數(shù)是否合法,然后按照邏輯位點(diǎn)(queue offset)發(fā)起拉取請(qǐng)求。TieredConsumeQueue/TieredCommitLog 將邏輯位點(diǎn)換算為對(duì)應(yīng)文件的物理位點(diǎn)從 TieredFileSegment 讀取消息。

// TieredMessageFetcher#getMessageAsync similar with TieredMessageStore#getMessageAsyncpublic CompletableFuturegetMessageAsync(String group, String topic, int queueId,        long queueOffset, int maxMsgNums, final MessageFilter messageFilter)

TieredFileSegment 維護(hù)每個(gè)儲(chǔ)存在文件系統(tǒng)中的物理文件位點(diǎn),并通過(guò)為不同存儲(chǔ)介質(zhì)實(shí)現(xiàn)的接口從中讀取所需的數(shù)據(jù)。

/**  * Get data from backend file system  *  * @param position the index from where the file will be read  * @param length the data size will be read  * @return data to be read  */CompletableFutureread0(long position, int length);

預(yù)讀緩存

TieredMessageFetcher 讀取消息時(shí)會(huì)預(yù)讀一部分消息供下次使用,這些消息暫存在預(yù)讀緩存中。

protected final CachereadAheadCache;

預(yù)讀緩存的設(shè)計(jì)參考了 TCP Tahoe 擁塞控制算法,每次預(yù)讀的消息量類似擁塞窗口采用加法增、乘法減的機(jī)制控制:

加法增:從最小窗口開始,每次增加等同于客戶端 batchSize 的消息量。 乘法減:當(dāng)緩存的消息超過(guò)了緩存過(guò)期時(shí)間仍未被全部拉取,在清理緩存的同時(shí)會(huì)將下次預(yù)讀消息量減半。

預(yù)讀緩存支持在讀取消息量較大時(shí)分片并發(fā)請(qǐng)求,以取得更大帶寬和更小的延遲。

某個(gè) topic 消息的預(yù)讀緩存由消費(fèi)這個(gè) topic 的所有 group 共享,緩存失效策略為:

所有訂閱這個(gè) topic 的 group 都訪問了緩存 到達(dá)緩存過(guò)期時(shí)間

故障恢復(fù)

上文中我們介紹上傳進(jìn)度由 commit offset 和 dispatch offset 控制。多級(jí)存儲(chǔ)會(huì)為每個(gè) topic、queue、fileSegment 創(chuàng)建元數(shù)據(jù)并持久化這兩種位點(diǎn)。當(dāng) Broker 重啟后會(huì)從元數(shù)據(jù)中恢復(fù),繼續(xù)從 commit offset 開始上傳消息,之前緩存的消息會(huì)重新上傳并不會(huì)丟失。

開發(fā)計(jì)劃

面向云原生的存儲(chǔ)系統(tǒng)要最大化利用云上存儲(chǔ)的價(jià)值,而對(duì)象存儲(chǔ)正是云計(jì)算紅利的體現(xiàn)。RocketMQ 多級(jí)存儲(chǔ)希望一方面利用對(duì)象存儲(chǔ)低成本的優(yōu)勢(shì)延長(zhǎng)消息存儲(chǔ)時(shí)間、拓展數(shù)據(jù)的價(jià)值;另一方面利用其共享存儲(chǔ)的特性在多副本架構(gòu)中兼得成本和數(shù)據(jù)可靠性,以及未來(lái)向 Serverless 架構(gòu)演進(jìn)。

tag 過(guò)濾

多級(jí)存儲(chǔ)拉取消息時(shí)沒有計(jì)算消息的 tag 是否匹配,tag 過(guò)濾交給客戶端處理。這樣會(huì)帶來(lái)額外的網(wǎng)絡(luò)開銷,計(jì)劃后續(xù)在服務(wù)端增加 tag 過(guò)濾能力。

廣播消費(fèi)以及多個(gè)消費(fèi)進(jìn)度不同的消費(fèi)者

預(yù)讀緩存失效需要所有訂閱這個(gè) topic 的 group 都訪問了緩存,這在多個(gè) group 消費(fèi)進(jìn)度不一致的情況下很難觸發(fā),導(dǎo)致無(wú)用的消息在緩存中堆積。

需要計(jì)算出每個(gè) group 的消費(fèi) qps 來(lái)估算某個(gè) group 能否在緩存失效前用上緩存的消息。如果緩存的消息預(yù)期在失效前都不會(huì)被再次訪問,那么它應(yīng)該被立即過(guò)期。相應(yīng)的對(duì)于廣播消費(fèi),消息的過(guò)期策略應(yīng)被優(yōu)化為所有 Client 都讀取這條消息后才失效。

和高可用架構(gòu)的融合

目前主要面臨以下三個(gè)問題:

元數(shù)據(jù)同步:如何可靠的在多個(gè)節(jié)點(diǎn)間同步元數(shù)據(jù),slave 晉升時(shí)如何校準(zhǔn)和補(bǔ)全缺失的元數(shù)據(jù); 禁止上傳超過(guò) confirm offset 的消息:為了避免消息回退,上傳的最大 offset 不能超過(guò) confirm offset; slave 晉升時(shí)快速啟動(dòng)多級(jí)存儲(chǔ):只有 master 節(jié)點(diǎn)具有寫權(quán)限,在 slave 節(jié)點(diǎn)晉升后需要快速拉起多級(jí)存儲(chǔ)斷點(diǎn)續(xù)傳。

相關(guān)鏈接:

[1] README

https://github.com/apache/rocketmq/blob/develop/tieredstore/README.md

版權(quán)聲明:本文內(nèi)容由阿里云實(shí)名注冊(cè)用戶自發(fā)貢獻(xiàn),版權(quán)歸原作者所有,阿里云開發(fā)者社區(qū)不擁有其著作權(quán),亦不承擔(dān)相應(yīng)法律責(zé)任。具體規(guī)則請(qǐng)查看《阿里云開發(fā)者社區(qū)用戶服務(wù)協(xié)議》和《阿里云開發(fā)者社區(qū)知識(shí)產(chǎn)權(quán)保護(hù)指引》。如果您發(fā)現(xiàn)本社區(qū)中有涉嫌抄襲的內(nèi)容,填寫侵權(quán)投訴表單進(jìn)行舉報(bào),一經(jīng)查實(shí),本社區(qū)將立刻刪除涉嫌侵權(quán)內(nèi)容。

關(guān)鍵詞:

專題首頁(yè)|財(cái)金網(wǎng)首頁(yè)

投資
探索

精彩
互動(dòng)

獨(dú)家
觀察

京ICP備2021034106號(hào)-38   營(yíng)業(yè)執(zhí)照公示信息  聯(lián)系我們:55 16 53 8 @qq.com  財(cái)金網(wǎng)  版權(quán)所有  cfenews.com
主站蜘蛛池模板: 亚洲爆乳无码专区| 国产精品欧美激情| 色综合天天狠天天透天天伊人 | 免费一级特黄毛片| 日韩亚洲欧美视频| 一区二区不卡在线| 亚洲欧美国产不卡| 日本一区二区视频| 久久大香伊蕉在人线观看热2| 精品人妻人人做人人爽| 欧美日韩国产免费一区二区三区 | 国产精品日韩三级| 亚洲a一级视频| 日本黄网免费一区二区精品 | 国产精品久久久久久久久久三级| 国产美女在线精品免费观看| 国产精品91一区| 日韩欧美一区三区| 精品国产免费av| 欧美国产日韩在线播放| 久久精品免费播放| 日韩专区中文字幕| 久久久久久91香蕉国产| 涩涩日韩在线| 九九热精品在线| 无码无遮挡又大又爽又黄的视频 | 国产区亚洲区欧美区| 日产精品高清视频免费| 国产精品网站免费| 欧美一区二视频在线免费观看| 国产精品美女诱惑| 欧美日韩一区在线播放| 日韩一级免费在线观看| 国产精品1234| 精品麻豆av| 亚洲熟妇av日韩熟妇在线| 国产精品日韩一区二区免费视频| 久久亚洲中文字幕无码| 欧美日韩一区二区视频在线观看| 国产欧美日韩精品在线观看| 久久精品国产精品亚洲精品色|