客服熱線
186-8811-5347、186-7086-0265
官方郵箱
contactus@mingting.cn
添加微信
立即線上溝通
客服微信
詳情請(qǐng)咨詢客服
客服熱線
186-8811-5347、186-7086-0265
官方郵箱
contactus@mingting.cn
2022-05-15 來(lái)源:金山毒霸電腦優(yōu)化作者:電腦技巧&問(wèn)題
BIGO?于?2014?年成立,是一家高速發(fā)展的科技公司。基于強(qiáng)大的音視頻處理技術(shù)、全球音視頻實(shí)時(shí)傳輸技術(shù)、人工智能技術(shù)、CDN?技術(shù),BIGO?推出了一系列音視頻類社交及內(nèi)容產(chǎn)品,包括?Bigo?Live(直播)和?Likee(短視頻)等,在全球已擁有近?1?億用戶,產(chǎn)品及服務(wù)已覆蓋超過(guò)?150?個(gè)國(guó)家和地區(qū)。
1挑戰(zhàn)
最初,BIGO?的消息流平臺(tái)主要采用開源?Kafka?作為數(shù)據(jù)支撐。隨著數(shù)據(jù)規(guī)模日益增長(zhǎng),產(chǎn)品不斷迭代,BIGO?消息流平臺(tái)承載的數(shù)據(jù)規(guī)模出現(xiàn)了成倍增長(zhǎng),下游的在線模型訓(xùn)練、在線推薦、實(shí)時(shí)數(shù)據(jù)分析、實(shí)時(shí)數(shù)倉(cāng)等業(yè)務(wù)對(duì)消息流平臺(tái)的實(shí)時(shí)性和穩(wěn)定性提出了更高的要求。開源的?Kafka?集群難以支撐海量數(shù)據(jù)處理場(chǎng)景,我們需要投入更多的人力去維護(hù)多個(gè)?Kafka?集群,這樣成本會(huì)越來(lái)越高,主要體現(xiàn)在以下幾個(gè)方面:
數(shù)據(jù)存儲(chǔ)和消息隊(duì)列服務(wù)綁定,集群擴(kuò)縮容?/?分區(qū)均衡需要大量拷貝數(shù)據(jù),造成集群性能下降。
當(dāng)分區(qū)副本不處于?ISR(同步)狀態(tài)時(shí),一旦有?broker?發(fā)生故障,可能會(huì)造成數(shù)據(jù)丟失或該分區(qū)無(wú)法提供讀寫服務(wù)。
當(dāng)?Kafka?broker?磁盤故障?/?空間占用率過(guò)高時(shí),需要進(jìn)行人工干預(yù)。
集群跨區(qū)域同步使用?KMM(Kafka?Mirror?Maker),性能和穩(wěn)定性難以達(dá)到預(yù)期。
在?catch-up?讀場(chǎng)景下,容易出現(xiàn)?PageCache?污染,造成讀寫性能下降。
Kafka?broker?上存儲(chǔ)的?topic?分區(qū)數(shù)量有限,分區(qū)數(shù)越多,磁盤讀寫順序性越差,讀寫性能越低。
Kafka?集群規(guī)模增長(zhǎng)導(dǎo)致運(yùn)維成本急劇增長(zhǎng),需要投入大量的人力進(jìn)行日常運(yùn)維;在?BIGO,擴(kuò)容一臺(tái)機(jī)器到?Kafka?集群并進(jìn)行分區(qū)均衡,需要?0.5?人?/?天;縮容一臺(tái)機(jī)器需要?1?人?/?天。
如果繼續(xù)使用?Kafka,成本會(huì)不斷上升:擴(kuò)縮容機(jī)器、增加運(yùn)維人力。同時(shí),隨著業(yè)務(wù)規(guī)模增長(zhǎng),我們對(duì)消息系統(tǒng)有了更高的要求:系統(tǒng)要更穩(wěn)定可靠、便于水平擴(kuò)展、延遲低。為了提高消息隊(duì)列的實(shí)時(shí)性、穩(wěn)定性和可靠性,降低運(yùn)維成本,我們開始考慮是否要基于開源?Kafka?做本地化二次開發(fā),或者看看社區(qū)中有沒(méi)有更好的解決方案,來(lái)解決我們?cè)诰S護(hù)?Kafka?集群時(shí)遇到的問(wèn)題。
2為什么選擇?Pulsar
2019?年?11?月,我們開始調(diào)研消息隊(duì)列,對(duì)比當(dāng)前主流消息流平臺(tái)的優(yōu)缺點(diǎn),并跟我們的需求對(duì)接。在調(diào)研過(guò)程中,我們發(fā)現(xiàn)?Apache?Pulsar?是下一代云原生分布式消息流平臺(tái),集消息、存儲(chǔ)、輕量化函數(shù)式計(jì)算為一體。Pulsar?能夠無(wú)縫擴(kuò)容、延遲低、吞吐高,支持多租戶和跨地域復(fù)制。最重要的是,Pulsar?存儲(chǔ)、計(jì)算分離的架構(gòu)能夠完美解決?Kafka?擴(kuò)縮容的問(wèn)題。Pulsar?producer?把消息發(fā)送給?broker,broker?通過(guò)?bookie?client?寫到第二層的存儲(chǔ)?BookKeeper?上。
Pulsar?采用存儲(chǔ)、計(jì)算分離的分層架構(gòu)設(shè)計(jì),支持多租戶、持久化存儲(chǔ)、多機(jī)房跨區(qū)域數(shù)據(jù)復(fù)制,具有強(qiáng)一致性、高吞吐以及低延時(shí)的高可擴(kuò)展流數(shù)據(jù)存儲(chǔ)特性。
水平擴(kuò)容:能夠無(wú)縫擴(kuò)容到成百上千個(gè)節(jié)點(diǎn)。
低延遲:在大規(guī)模的消息量下依然能夠保持低延遲(小于?5?ms)。
持久化機(jī)制:Pulsar?的持久化機(jī)制構(gòu)建在?Apache?BookKeeper?上,實(shí)現(xiàn)了讀寫分離。
讀寫分離:BookKeeper?的讀寫分離?IO?模型極大發(fā)揮了磁盤順序?qū)懶阅埽瑢?duì)機(jī)械硬盤相對(duì)比較友好,單臺(tái)?bookie?節(jié)點(diǎn)支撐的?topic?數(shù)不受限制。
為了進(jìn)一步加深對(duì)?Apache?Pulsar?的理解,衡量?Pulsar?能否真正滿足我們生產(chǎn)環(huán)境大規(guī)模消息?Pub-Sub?的需求,我們從?2019?年?12?月開始進(jìn)行了一系列壓測(cè)工作。由于我們使用的是機(jī)械硬盤,沒(méi)有?SSD,在壓測(cè)過(guò)程中遇到了一些性能問(wèn)題,在?StreamNative?的協(xié)助下,我們分別對(duì)?Broker?和?BookKeeper?進(jìn)行了一系列的性能調(diào)優(yōu),Pulsar?的吞吐和穩(wěn)定性均有所提高。
經(jīng)過(guò)?3~4?個(gè)月的壓測(cè)和調(diào)優(yōu),我們認(rèn)為?Pulsar?完全能夠解決我們使用?Kafka?時(shí)遇到的各種問(wèn)題,并于?2020?年?4?月在測(cè)試環(huán)境上線?Pulsar。
3Apache?Pulsar?at?BIGO:Pub-Sub?消費(fèi)模式
2020?年?5?月,我們正式在生產(chǎn)環(huán)境中使用?Pulsar?集群。Pulsar?在?BIGO?的場(chǎng)景主要是?Pub-Sub?的經(jīng)典生產(chǎn)消費(fèi)模式,前端有?Baina?服務(wù)(用?C++?實(shí)現(xiàn)的數(shù)據(jù)接收服務(wù)),Kafka?的?Mirror?Maker?和?Flink,以及其他語(yǔ)言如?Java、
在下游,我們對(duì)接的業(yè)務(wù)場(chǎng)景有實(shí)時(shí)數(shù)倉(cāng)、實(shí)時(shí)?ETL(Extract-Transform-Load,將數(shù)據(jù)從來(lái)源端經(jīng)過(guò)抽取(extract)、轉(zhuǎn)換(transform)、加載(load)至目的端的過(guò)程)、實(shí)時(shí)數(shù)據(jù)分析和實(shí)時(shí)推薦。大部分業(yè)務(wù)場(chǎng)景使用?Flink?消費(fèi)?Pulsar?topic?中的數(shù)據(jù),并進(jìn)行業(yè)務(wù)邏輯處理;其他業(yè)務(wù)場(chǎng)景消費(fèi)使用的客戶端語(yǔ)言主要分布在?C++、Go、Python?等。數(shù)據(jù)經(jīng)過(guò)各自業(yè)務(wù)邏輯處理后,最終會(huì)寫入?Hive、Pulsar?topic?以及?ClickHouse、HDFS、Redis?等第三方存儲(chǔ)服務(wù)。
4Pulsar?+?Flink?實(shí)時(shí)流平臺(tái)
在?BIGO,我們借助?Flink?和?Pulsar?打造了實(shí)時(shí)流平臺(tái)。在介紹這個(gè)平臺(tái)之前,我們先了解下?Pulsar?Flink?Connector?的內(nèi)部運(yùn)行機(jī)理。在?Pulsar?Flink?Source/Sink?API?中,上游有一個(gè)?Pulsar?topic,中間是?Flink?job,下游有一個(gè)?Pulsar?topic。我們?cè)趺聪M(fèi)這個(gè)?topic,又怎樣處理數(shù)據(jù)并寫入?Pulsar?topic?呢?
按照上圖左側(cè)代碼示例,初始化一個(gè)?StreamExecutionEnvironment,進(jìn)行相關(guān)配置,比如修改?property、topic?值。然后創(chuàng)建一個(gè)?FlinkPulsarSource?對(duì)象,這個(gè)?Source?里面填上?serviceUrl(brokerlist)、adminUrl(admin?地址)以及?topic?數(shù)據(jù)的序列化方式,最終會(huì)把?property?傳進(jìn)去,這樣就能夠讀取?Pulsar?topic?中的數(shù)據(jù)。Sink?的使用方法非常簡(jiǎn)單,首先創(chuàng)建一個(gè)?FlinkPulsarSink,Sink?里面指定?target?topic,再指定?TopicKeyExtractor?作為?key,并調(diào)用?addsink,把數(shù)據(jù)寫入?Sink。這個(gè)生產(chǎn)消費(fèi)模型很簡(jiǎn)單,和?Kafka?很像。
Pulsar?topic?和?Flink?的消費(fèi)如何聯(lián)動(dòng)呢?如下圖所示,新建?FlinkPulsarSource?時(shí),會(huì)為?topic?的每一個(gè)分區(qū)新創(chuàng)建一個(gè)?reader?對(duì)象。要注意的是?Pulsar?Flink?Connector?底層使用?reader?API?消費(fèi),會(huì)先創(chuàng)建一個(gè)?reader,這個(gè)?reader?使用?Pulsar?Non-Durable?Cursor。Reader?消費(fèi)的特點(diǎn)是讀取一條數(shù)據(jù)后馬上提交(commit),所以在監(jiān)控上可能會(huì)看到?reader?對(duì)應(yīng)的?subion?沒(méi)有?backlog?信息。
Offset?Commit?完成后,Pulsar?broker?會(huì)將?Offset?信息(在?Pulsar?中以?Cursor?表示)存儲(chǔ)到底層的分布式存儲(chǔ)系統(tǒng)?BookKeeper?中,這樣做的好處是當(dāng)?Flink?任務(wù)重啟后,會(huì)有兩層恢復(fù)保障。第一種情況是從?checkpoint?恢復(fù):可以直接從?checkpoint?里獲得上一次消費(fèi)的?message?id,通過(guò)這個(gè)?message?id?獲取數(shù)據(jù),這個(gè)數(shù)據(jù)流就能繼續(xù)消費(fèi)。如果沒(méi)有從?checkpoint?恢復(fù),F(xiàn)link?任務(wù)重啟后,會(huì)根據(jù)?SubionName?從?Pulsar?中獲取上一次?Commit?對(duì)應(yīng)的?Offset?位置開始消費(fèi)。這樣就能有效防止?checkpoint?損壞導(dǎo)致整個(gè)?Flink?任務(wù)無(wú)法成功啟動(dòng)的問(wèn)題。
為了降低?Flink?消費(fèi)?Pulsar?topic?的門檻,讓?Pulsar?Flink?Connector?支持更加豐富的?Flink?新特性,BIGO?消息隊(duì)列團(tuán)隊(duì)為?Pulsar?Flink?Connector?增加了?Pulsar?Flink?SQL?DDL(Data?Definition?Language,數(shù)據(jù)定義語(yǔ)言)?和?Flink?1.11?支持。此前官方提供的?Pulsar?Flink?SQL?只支持?Catalog,要想通過(guò)?DDL?形式消費(fèi)、處理?Pulsar?topic?中的數(shù)據(jù)不太方便。在?BIGO?場(chǎng)景中,大部分?topic?數(shù)據(jù)都以?JSON?格式存儲(chǔ),而?JSON?的?schema?沒(méi)有提前注冊(cè),所以只能在?Flink?SQL?中指定?topic?的?DDL?后才可以消費(fèi)。針對(duì)這種場(chǎng)景,BIGO?基于?Pulsar?Flink?Connector?做了二次開發(fā),提供了通過(guò)?Pulsar?Flink?SQL?DDL?形式消費(fèi)、解析、處理?Pulsar?topic?數(shù)據(jù)的代碼框架(如下圖所示)。
左邊的代碼中,第一步是配置?Pulsar?topic?的消費(fèi),首先指定?topic?的?DDL?形式,比如?rip、rtime、uid?等,下面是消費(fèi)?Pulsar?topic?的基礎(chǔ)配置,比如?topic?名稱、service-url、admin-url?等。底層?reader?讀到消息后,會(huì)根據(jù)?DDL?解出消息,將數(shù)據(jù)存儲(chǔ)在?test_flink_sql?表中。第二步是常規(guī)邏輯處理(如對(duì)表進(jìn)行字段抽取、做?join?等),得出相關(guān)統(tǒng)計(jì)信息或其他相關(guān)結(jié)果后,返回這些結(jié)果,寫到?HDFS?或其他系統(tǒng)上等。第三步,提取相應(yīng)字段,將其插入一張?hive?表。由于?Flink?1.11?對(duì)?hive?的寫入支持比?1.9.1?更加優(yōu)秀,所以?BIGO?又做了一次?API?兼容和版本升級(jí),使?Pulsar?Flink?Connector?支持?Flink?1.11。BIGO?基于?Pulsar?和?Flink?構(gòu)建的實(shí)時(shí)流平臺(tái)主要用于實(shí)時(shí)?ETL?處理場(chǎng)景和?AB-test?場(chǎng)景。
實(shí)時(shí)?ETL?處理場(chǎng)景
實(shí)時(shí)?ETL?處理場(chǎng)景主要運(yùn)用?Pulsar?Flink?Source?及?Pulsar?Flink?Sink。這個(gè)場(chǎng)景中,Pulsar?topic?實(shí)現(xiàn)幾百甚至上千個(gè)?topic,每個(gè)?topic?都有獨(dú)立的?schema。我們需要對(duì)成百上千個(gè)?topic?進(jìn)行常規(guī)處理,如字段轉(zhuǎn)換、容錯(cuò)處理、寫入?HDFS?等。每個(gè)?topic?都對(duì)應(yīng)?HDFS?上的一張表,成百上千個(gè)?topic?會(huì)在?HDFS?上映射成百上千張表,每張表的字段都不一樣,這就是我們遇到的實(shí)時(shí)?ETL?場(chǎng)景。
隨著程序運(yùn)行,我們發(fā)現(xiàn)這種方案也存在問(wèn)題:算子之間壓力不均衡。因?yàn)橛行?topic?流量大,有些流量小,如果完全通過(guò)隨機(jī)哈希的方式映射到對(duì)應(yīng)的?task?manager?上去,有些?task?manager?處理的流量會(huì)很高,而有些?task?manager?處理的流量很低,導(dǎo)致有些?task?機(jī)器上積塞非常嚴(yán)重,拖慢?Flink?流的處理。所以我們引入了?slot?group?概念,根據(jù)每個(gè)?topic?的流量情況進(jìn)行分組,流量會(huì)映射到?topic?的分區(qū)數(shù),在創(chuàng)建?topic?分區(qū)時(shí)也以流量為依據(jù),如果流量很高,就多為?topic?創(chuàng)建分區(qū),反之少一些。分組時(shí),把流量小的?topic?分到一個(gè)?group?中,把流量大的?topic?單獨(dú)放在一個(gè)?group?中,很好地隔離了資源,保證?task?manager?總體上流量均衡。
AB-test?場(chǎng)景
實(shí)時(shí)數(shù)倉(cāng)需要提供小時(shí)表或天表為數(shù)據(jù)分析師及推薦算法工程師提供數(shù)據(jù)查詢服務(wù),簡(jiǎn)單來(lái)講就是?app?應(yīng)用中會(huì)有很多打點(diǎn),各種類型的打點(diǎn)會(huì)上報(bào)到服務(wù)端。如果直接暴露原始打點(diǎn)給業(yè)務(wù)方,不同的業(yè)務(wù)使用方就需要訪問(wèn)各種不同的原始表從不同維度進(jìn)行數(shù)據(jù)抽取,并在表之間進(jìn)行關(guān)聯(lián)計(jì)算。頻繁對(duì)底層基礎(chǔ)表進(jìn)行數(shù)據(jù)抽取和關(guān)聯(lián)操作會(huì)嚴(yán)重浪費(fèi)計(jì)算資源,所以我們提前從基礎(chǔ)表中抽取用戶關(guān)心的維度,將多個(gè)打點(diǎn)合并在一起,構(gòu)成一張或多張寬表,覆蓋上面推薦相關(guān)的或數(shù)據(jù)分析相關(guān)的?80%?~?90%?場(chǎng)景任務(wù)。
在實(shí)時(shí)數(shù)倉(cāng)場(chǎng)景下還需實(shí)時(shí)中間表,我們的解決方案是,針對(duì)?topic?A?到?topic?K?,我們使用?Pulsar?Flink?SQL?將消費(fèi)到的數(shù)據(jù)解析成相應(yīng)的表。通常情況下,將多張表聚合成一張表的常用做法是使用?join,如把表?A?到?K?按照?uid?進(jìn)行?join?操作,形成非常寬的寬表;但在?Flink?SQL?中?join?多張寬表效率較低。所以?BIGO?使用?union?來(lái)替代?join,做成很寬的視圖,以小時(shí)為單位返回視圖,寫入?ClickHouse,提供給下游的業(yè)務(wù)方實(shí)時(shí)查詢。使用?union?來(lái)替代?join?加速表的聚合,能夠把小時(shí)級(jí)別的中間表產(chǎn)出控制在分鐘級(jí)別。
輸出天表可能還需要?join?存放在?hive?上的表或其他存儲(chǔ)介質(zhì)上的離線表,即流表和離線表之間?join?的問(wèn)題。如果直接?join,checkpoint?中需要存儲(chǔ)的中間狀態(tài)會(huì)比較大,所以我們?cè)诹硗庖粋€(gè)維度上做了優(yōu)化。
左側(cè)部分類似于小時(shí)表,每個(gè)?topic?使用?Pulsar?Flink?SQL?消費(fèi)并轉(zhuǎn)換成對(duì)應(yīng)的表,表之間進(jìn)行?union?操作,將?union?得到的表以天為單位輸入到?HBase(此處引入?HBase?是為了做替代它的?join)。
右側(cè)需要?join?離線數(shù)據(jù),使用?Spark?聚合離線的?Hive?表(如表?a1、a2、a3),聚合后的數(shù)據(jù)會(huì)通過(guò)精心設(shè)計(jì)的?row-key?寫入?HBase?中。數(shù)據(jù)聚合后狀態(tài)如下:假設(shè)左邊數(shù)據(jù)的?key?填了寬表的前?80?列,后面?Spark?任務(wù)算出的數(shù)據(jù)對(duì)應(yīng)同樣一個(gè)?key,填上寬表的后?20?列,在?HBase?中組成一張很大的寬表,把最終數(shù)據(jù)再次從?HBase?抽出,寫入?ClickHouse,供上層用戶查詢,這就是?AB-test?的主體架構(gòu)。
5業(yè)務(wù)收益
從?2020?年?5?月上線至今,Pulsar?運(yùn)行穩(wěn)定,日均處理消息數(shù)百億,字節(jié)入流量為?2~3?GB/s。Apache?Pulsar?提供的高吞吐、低延遲、高可靠性等特性極大提高了?BIGO?消息處理能力,降低了消息隊(duì)列運(yùn)維成本,節(jié)約了近?50%?的硬件成本。目前,我們?cè)趲资_(tái)物理主機(jī)上部署了上百個(gè)?Pulsar?broker?和?bookie?進(jìn)程,采用?bookie?和?broker?在同一個(gè)節(jié)點(diǎn)的混部模式,已經(jīng)把?ETL?從?Kafka?遷移到?Pulsar,并逐步將生產(chǎn)環(huán)境中消費(fèi)?Kafka?集群的業(yè)務(wù)(比如?Flink、Flink?SQL、ClickHouse?等)遷移到?Pulsar?上。隨著更多業(yè)務(wù)的遷移,Pulsar?上的流量會(huì)持續(xù)上漲。
我們的?ETL?任務(wù)有一萬(wàn)多個(gè)?topic,每個(gè)?topic?平均有?3?個(gè)分區(qū),使用?3?副本的存儲(chǔ)策略。之前使用?Kafka,隨著分區(qū)數(shù)增加,磁盤由順序讀寫逐漸退化為隨機(jī)讀寫,讀寫性能退化嚴(yán)重。Apache?Pulsar?的存儲(chǔ)分層設(shè)計(jì)能夠輕松支持百萬(wàn)?topic,為我們的?ETL?場(chǎng)景提供了優(yōu)雅支持。
6未來(lái)展望
BIGO?在?Pulsar?broker?負(fù)載均衡、broker?cache?命中率優(yōu)化、broker?相關(guān)監(jiān)控、BookKeeper?讀寫性能優(yōu)、BookKeeper?磁盤?IO?性能優(yōu)化、Pulsar?與?Flink、Pulsar?與?Flink?SQL?結(jié)合等方面做了大量工作,提升了?Pulsar?的穩(wěn)定性和吞吐,也降低了?Flink?與?Pulsar?結(jié)合的門檻,為?Pulsar?的推廣打下了堅(jiān)實(shí)基礎(chǔ)。
未來(lái),我們會(huì)增加?Pulsar?在?BIGO?的場(chǎng)景應(yīng)用,幫助社區(qū)進(jìn)一步優(yōu)化、完善?Pulsar?功能,具體如下:
為?Apache?Pulsar?研發(fā)新特性,比如支持?topic?policy?相關(guān)特性。
遷移更多任務(wù)到?Pulsar。這項(xiàng)工作涉及兩方面,一是遷移之前使用?Kafka?的任務(wù)到?Pulsar。二是新業(yè)務(wù)直接接入?Pulsar。
BIGO?準(zhǔn)備使用?KoP?來(lái)保證數(shù)據(jù)遷移平滑過(guò)渡。因?yàn)?BIGO?有大量消費(fèi)?Kafka?集群的?Flink?任務(wù),我們希望能夠直接在?Pulsar?中做一層?KoP,簡(jiǎn)化遷移流程。
對(duì)?Pulsar?及?BookKeeper?持續(xù)進(jìn)行性能優(yōu)化。由于生產(chǎn)環(huán)境中流量較高,BIGO?對(duì)系統(tǒng)的可靠性和穩(wěn)定性要求較高。
持續(xù)優(yōu)化?BookKeeper?的?IO?協(xié)議棧。Pulsar?的底層存儲(chǔ)本身是?IO?密集型系統(tǒng),保證底層?IO?高吞吐,才能夠提升上層吞吐,保證性能穩(wěn)定。
最后,小編給您推薦,金山毒霸“手機(jī)數(shù)據(jù)恢復(fù)”,如您遇到誤刪手機(jī)圖片、文件、通訊錄、通話記錄和短信等,金山手機(jī)數(shù)據(jù)恢復(fù)統(tǒng)統(tǒng)都能幫您恢復(fù)。