于流批一體架構(gòu)的思考)
1. 流批一體平臺的主要作用數(shù)據(jù)采集從各種異構(gòu)數(shù)據(jù)源如數(shù)據(jù)庫、文件、消息隊(duì)列等中實(shí)時抽取數(shù)據(jù)包括增量和全量抽取。數(shù)據(jù)轉(zhuǎn)換對抽取的數(shù)據(jù)進(jìn)行清洗、轉(zhuǎn)換和處理以滿足目標(biāo)系統(tǒng)的要求。它提供了一系列的內(nèi)置轉(zhuǎn)換器和處理器也支持自定義的轉(zhuǎn)換邏輯。數(shù)據(jù)加載將轉(zhuǎn)換后的數(shù)據(jù)加載到目標(biāo)系統(tǒng)中支持多種數(shù)據(jù)目的地如搜索引擎、推薦引擎、大數(shù)據(jù)表等。實(shí)時監(jiān)控和報警提供實(shí)時監(jiān)控和報警功能可以監(jiān)控數(shù)據(jù)流的健康狀態(tài)、數(shù)據(jù)延遲和錯誤。可視化界面通過可視化界面用戶可以輕松配置和管理數(shù)據(jù)流并實(shí)時查看數(shù)據(jù)流的運(yùn)行狀態(tài)和性能指標(biāo)。2. 數(shù)據(jù)架構(gòu)演進(jìn)2.1 BI 架構(gòu)BIBusiness Intelligence商業(yè)智能的概念很早就有了正如 AI 這一概念一樣。早期它的內(nèi)涵相對模糊按照百度百科的解釋“商業(yè)智能描述了一系列的概念和方法通過應(yīng)用基于事實(shí)的支持系統(tǒng)來輔助商業(yè)決策的制定。” 隨著人們實(shí)踐不斷深入BI 系統(tǒng)的樣貌也逐漸清晰。到了上世紀(jì)九十年代BI 系統(tǒng)迎來了它的第一個輝煌時期Gartner 將各種類型的類 BI 系統(tǒng)全部統(tǒng)稱為 BIBI 產(chǎn)品也基本確定為了是一套集數(shù)據(jù)清洗、數(shù)據(jù)分析、數(shù)據(jù)挖掘、報表展示等功能于一體的完整解決方案數(shù)據(jù)倉庫也基于此建立。那時雖然沒有大數(shù)據(jù)的概念但數(shù)據(jù)分析、商業(yè)分析顯然是人們長久以來都有的需求也積累了相當(dāng)多的方法論。當(dāng)數(shù)據(jù)量不是主要矛盾時BI 系統(tǒng)能夠支持的分析方法、UI 等層面就成為了核心競爭力。當(dāng)然大多數(shù) BI 系統(tǒng)都構(gòu)建在關(guān)系型數(shù)據(jù)庫之上或者說很多 BI 系統(tǒng)本就是商業(yè)關(guān)系型數(shù)據(jù)庫的配套產(chǎn)品因此也都是支持 SQL 語言的。以 BI 系統(tǒng)為核心的數(shù)據(jù)架構(gòu)如下圖所示初代 BI 系統(tǒng)沒落的原因主要是底層構(gòu)建在傳統(tǒng)關(guān)系型數(shù)據(jù)庫之上因?yàn)榇嬖跀?shù)據(jù)一致性約束等問題支持不了大數(shù)據(jù)。不支持非結(jié)構(gòu)化數(shù)據(jù)。2.2 傳統(tǒng)大數(shù)據(jù)架構(gòu)為了解決初代 BI 系統(tǒng)存在的上述問題一些公司開始研發(fā)分布式的計(jì)算引擎和分布式的存儲平臺。其中最成功、最知名的便是 Google 研發(fā)的分布式文件系統(tǒng)與 MapReduce 計(jì)算引擎后來這套技術(shù)被開源重寫為了 Hadoop 體系的多個項(xiàng)目其生態(tài)圈也不斷擴(kuò)大。典型的傳統(tǒng)大數(shù)據(jù)架構(gòu)流程傳統(tǒng)大數(shù)據(jù)架構(gòu)它的業(yè)務(wù)系統(tǒng)數(shù)據(jù)源可能是關(guān)系型數(shù)據(jù)庫 MySQL也可能是平面文件也可能是任意未知的源數(shù)據(jù)采集和數(shù)據(jù)同步工具也是視具體的業(yè)務(wù)和上下游技術(shù)選型而定接下來數(shù)據(jù)會進(jìn)入數(shù)據(jù)倉庫大致上會依次經(jīng)過 ODS 層、DWD 層和 ADS 層最終提供給消費(fèi)方使用。集團(tuán)內(nèi)通常業(yè)務(wù)數(shù)據(jù)通過 binlog 同步到 TT或者流量日志直接上報到日志服務(wù)器再同步到 TT。TT 定期將一個時間區(qū)間內(nèi)的數(shù)據(jù)同步到 ODPSODPS 再通過每日調(diào)度的任務(wù)對這些數(shù)據(jù)進(jìn)行處理最終落到 ADS 層的表。結(jié)果表的數(shù)據(jù)再同步到 Holo 或 Lindorm 等介質(zhì)中供消費(fèi)方使用。因此單看這整個流程實(shí)際上就是典型的傳統(tǒng)大數(shù)據(jù)架構(gòu)的一種實(shí)現(xiàn)。但需要注意的是該架構(gòu)并沒有對輸入數(shù)據(jù)有結(jié)構(gòu)化的要求也沒有規(guī)定 ETL 過程使用的工具和編程語言。在這種架構(gòu)下業(yè)務(wù)系統(tǒng)和分析系統(tǒng)的隔離性做得更好了而且無論輸入數(shù)據(jù)是什么最終提供給消費(fèi)方的都是標(biāo)準(zhǔn)的結(jié)構(gòu)化數(shù)據(jù)。它的缺點(diǎn)是整個過程不再有完整的解決方案需要做大量的定制化工作。2.3 流式架構(gòu)流式架構(gòu)的思路相當(dāng)激進(jìn)。雖然傳統(tǒng)大數(shù)據(jù)架構(gòu)在技術(shù)選型上與 BI 系統(tǒng)比已經(jīng)算是脫胎換骨但其精神還是一脈相承。流式架構(gòu)干脆扔掉一整套離線的數(shù)據(jù)采集、數(shù)據(jù)同步和 ETL 工作直接讓流式計(jì)算引擎消費(fèi)業(yè)務(wù)數(shù)據(jù)庫產(chǎn)生的增量數(shù)據(jù)并直接輸出給消費(fèi)方以此提供實(shí)時的計(jì)算結(jié)果。而早期的技術(shù)儲備明顯不足以同時高質(zhì)量保證實(shí)時性和結(jié)果的準(zhǔn)確性因此只被用在了極少數(shù)對結(jié)果實(shí)時性十分敏感卻對準(zhǔn)確性要求不高的場景中。隨著技術(shù)的進(jìn)步和業(yè)務(wù)復(fù)雜度的提高這種架構(gòu)也基本銷聲匿跡了。下圖是流式架構(gòu)的典型代表2.4 lambda 架構(gòu)在早期技術(shù)無法同時支持結(jié)果的實(shí)時性和準(zhǔn)確性的情況下如何通過架構(gòu)設(shè)計(jì)同時滿足實(shí)時性與準(zhǔn)確性需求Nathan Marz 提出了 Lambda 架構(gòu)。先看lambda架構(gòu)的示意圖Lambda架構(gòu)的邏輯是流任務(wù)與批任務(wù)讀取相同的數(shù)據(jù)源實(shí)時計(jì)算結(jié)果由流任務(wù)產(chǎn)出批任務(wù)通常按天執(zhí)行計(jì)算T-1的數(shù)據(jù)并寫入到結(jié)果表中。最終數(shù)據(jù)應(yīng)用根據(jù)自己的需要對兩個結(jié)果表的結(jié)果進(jìn)行合并。其核心思路是:用流任務(wù)保證結(jié)果的實(shí)時性同時用批任務(wù)保證結(jié)果的最終一致性。Lambda 架構(gòu)包含離線和實(shí)時兩條鏈路兩條鏈路從同一個業(yè)務(wù)數(shù)據(jù)源獲取數(shù)據(jù)。離線鏈路通過定期調(diào)度周期性的同步數(shù)據(jù)源并將其持久化存儲在分布式文件系統(tǒng)隨后這些數(shù)據(jù)會被提交給 Spark、MaxCompute 等離線計(jì)算引擎進(jìn)行大規(guī)模批處理作業(yè)經(jīng)過離線處理后的結(jié)果數(shù)據(jù)通常會存入高可用、低成本且支持復(fù)雜查詢的下游存儲系統(tǒng)例如 ADBHologres 等通過暴露離線數(shù)據(jù)的數(shù)據(jù)服務(wù)進(jìn)而為各類業(yè)務(wù)場景提供基于歷史全量數(shù)據(jù)的決策支持整個鏈路提供海量數(shù)據(jù)的高吞吐、高穩(wěn)定性的處理能力但結(jié)果從采集到產(chǎn)出至用戶可見往往需要數(shù)個小時。實(shí)時鏈路則通過實(shí)時捕捉數(shù)據(jù)源的 CDC 信息能夠近乎實(shí)時地追蹤并獲取數(shù)據(jù)源的變化情況變更事件一旦產(chǎn)生就會被立即推送到消息中間件隨后經(jīng)過 FlinkSpark Streaming 等流處理引擎消費(fèi)這些消息能夠做到對數(shù)據(jù)源變更低延遲高容錯的秒級響應(yīng)能力。但 Lambda 架構(gòu)有幾個顯而易見的缺點(diǎn)需要開發(fā)、維護(hù)兩套系統(tǒng)成本太大。兩套系統(tǒng)難以保證計(jì)算口徑的一致。不同引擎間由于支持的函數(shù)不同、參數(shù)配置不同導(dǎo)致代碼不能復(fù)用容易出現(xiàn)數(shù)據(jù)一致性和質(zhì)量問題總之Lambda 架構(gòu)在滿足了部分業(yè)務(wù)需求的同時給開發(fā)和運(yùn)維同學(xué)也帶來了 “深重的災(zāi)難”。他分拆了計(jì)算鏈路和存儲鏈路導(dǎo)致一個業(yè)務(wù)邏輯需要維護(hù)兩套代碼不同引擎間由于支持的函數(shù)不同、參數(shù)配置不同導(dǎo)致代碼不能復(fù)用容易出現(xiàn)數(shù)據(jù)一致性和質(zhì)量問題導(dǎo)致開發(fā)成本大運(yùn)維成本高。個人理解Lambda 就是一種硬把批和流雜糅在一起的架構(gòu)。2.5 Kappa 架構(gòu)在流處理技術(shù)不成熟的時期主要問題之一就是吞吐量上不去。隨著 Kafka 等大數(shù)據(jù)消息隊(duì)列的出現(xiàn)吞吐量不再是瓶頸。Kappa 架構(gòu)的主要貢獻(xiàn)之一就是引入了分布式消息隊(duì)列。如下圖所示與 Lambda 架構(gòu)不同Kappa 架構(gòu)只保留了流處理層完全舍棄了批處理層。Kappa 架構(gòu)專注于事件驅(qū)動和實(shí)時處理只使用一套實(shí)時任務(wù)鏈路完成業(yè)務(wù)數(shù)據(jù)產(chǎn)出有效降低了系統(tǒng)的復(fù)雜性和維護(hù)成本同時也提高了數(shù)據(jù)處理的靈活性和響應(yīng)速度。但由于實(shí)時鏈路中普遍存在的數(shù)據(jù)延遲亂序等問題導(dǎo)致計(jì)算結(jié)果偏差需要業(yè)務(wù)對當(dāng)日數(shù)據(jù)的誤差具有一定容忍性。此外數(shù)據(jù)回刷時需要將大規(guī)模數(shù)據(jù)集重放至消息中間件中這一過程重資源消耗且對消息中間件的存儲量和吞吐量有很高要求。對于一些月度或季度匯總的統(tǒng)計(jì)周期較長的指標(biāo)純實(shí)時計(jì)算的處理方式會產(chǎn)生大量的排序消耗和中間狀態(tài)存儲。簡單來說Kappa 架構(gòu)讓其中一個流處理層正常運(yùn)行數(shù)據(jù)應(yīng)用讀取它的輸出當(dāng)數(shù)據(jù)出現(xiàn)錯誤或是業(yè)務(wù)邏輯發(fā)生變更時啟動另一個流處理層利用消息隊(duì)列的重播機(jī)制重新消費(fèi)先前的數(shù)據(jù)并輸出到另一個結(jié)果表中當(dāng)確定可以替換線上表時完成替換。總的來說Kappa 架構(gòu)雖然解決了 Lambda 架構(gòu)一套邏輯兩份代碼的問題但同時引入了資源消耗過大、數(shù)據(jù)存在誤差等問題在實(shí)踐上仍然存在很大的挑戰(zhàn)。不過Kappa 架構(gòu)的另一貢獻(xiàn)是啟發(fā)了人們用單一系統(tǒng)去實(shí)現(xiàn)曾經(jīng)需要兩套系統(tǒng)才能實(shí)現(xiàn)的需求。人們開始思考為什么流式計(jì)算引擎不能提供結(jié)果的準(zhǔn)確性是哪些環(huán)節(jié)出了問題流處理引擎是否可以批處理引擎等價的語義2.6 流批一體架構(gòu)2.6.1 SparkSpark 針對 MapReduce 數(shù)據(jù)處理過程中海量中間結(jié)果落盤影響吞吐量的問題提出了基于內(nèi)存的計(jì)算思想并且通過 DAG 執(zhí)行引擎優(yōu)化執(zhí)行流程宣稱能夠提高 Hadoop 中 MapReduce 任務(wù)效率 100 倍。并且也提供了豐富的 API 支持圖計(jì)算和機(jī)器學(xué)習(xí)計(jì)算功能。Spark 在提出時專注于批任務(wù)的處理為了迎合流任務(wù)處理的需求Spark 建設(shè)了 Spark Streaming 組件能夠支持準(zhǔn)實(shí)時的數(shù)據(jù)處理。但 Spark Streaming 組件的實(shí)現(xiàn)方式是基于微批處理將數(shù)據(jù)流劃分為一系列時間戳連續(xù)的數(shù)據(jù)塊每個數(shù)據(jù)塊經(jīng)過 DAG 執(zhí)行引擎執(zhí)行后產(chǎn)生結(jié)果。由于 Spark 將實(shí)時流切分為微批進(jìn)行處理實(shí)時任務(wù)和離線任務(wù)僅有批窗口大小的差異實(shí)時任務(wù)的窗口放大至天級別便是離線任務(wù)。因此這種實(shí)現(xiàn)方式天然的可以支持一份代碼在實(shí)時任務(wù)和離線任務(wù)運(yùn)行即剛才提到的流批一體。但這種做法只能支持小體量且延時要求不高的應(yīng)用并且不支持基于業(yè)務(wù)時間的窗口只能支持?jǐn)?shù)據(jù)到達(dá)時間的滾動窗口導(dǎo)致很多具有 session 概念的業(yè)務(wù)數(shù)據(jù)無法計(jì)算。2.6.2 FlinkSpark 將數(shù)據(jù)流視為特殊的批數(shù)據(jù)采用微批來處理數(shù)據(jù)流的做法借鑒了離線數(shù)據(jù)處理的思想能夠保證數(shù)據(jù)吞吐量但對較低響應(yīng)延遲的場景卻無能為力。Flink 專注于實(shí)時任務(wù)的處理將批處理視為特殊的有界數(shù)據(jù)流提供低延遲精確一次基于事件時間的窗口的能力來滿足秒級甚至毫秒級的實(shí)時數(shù)據(jù)處理需求。Flink 提出 State 的概念通過維護(hù)實(shí)時任務(wù)的中間狀態(tài)以及回撤流機(jī)制保證結(jié)果的準(zhǔn)確性來實(shí)現(xiàn)單記錄粒度的數(shù)據(jù)處理能力來達(dá)到實(shí)時響應(yīng)的要求并通過 checkpoint 提供容錯機(jī)制最后實(shí)現(xiàn)了 Window 和 Time 功能來實(shí)現(xiàn)基于事件時間的開窗能力。然而在實(shí)際落地過程中用戶反饋呈現(xiàn)出明顯分化在流處理場景中Flink 憑借強(qiáng)大的狀態(tài)管理、exactly-once 語義保障以及低延遲性能已成為行業(yè)事實(shí)上的標(biāo)準(zhǔn)但在批處理場景中許多用戶仍傾向于使用 Spark 或 ODPS 等系統(tǒng)因其具備更成熟的查詢優(yōu)化器、更高的吞吐效率以及更友好的開發(fā)調(diào)試體驗(yàn)。更為關(guān)鍵的是即便基于 Flink 實(shí)現(xiàn)流批統(tǒng)一開發(fā)者往往仍需為流和批分別配置執(zhí)行環(huán)境、調(diào)度策略、Checkpoint 參數(shù)及并發(fā)度等設(shè)置。SQL 層面的統(tǒng)一同樣面臨挑戰(zhàn)Flink 的流式 SQL 引入了事件時間、處理時間、窗口機(jī)制等概念而這些在傳統(tǒng)批處理場景中往往無需關(guān)注。尤其在高吞吐實(shí)時場景下系統(tǒng)性能高度依賴對狀態(tài)清理和反壓處理等機(jī)制的深入理解要求開發(fā)者掌握復(fù)雜的調(diào)優(yōu)技巧。這導(dǎo)致 Flink 應(yīng)用的整體開發(fā)門檻居高不下實(shí)際落地中仍普遍依賴專業(yè)化的實(shí)時計(jì)算團(tuán)隊(duì)支撐。即便業(yè)務(wù)邏輯可以通過同一段 SQL 表達(dá)工程實(shí)踐中卻常常需要維護(hù)兩套獨(dú)立的運(yùn)行模式。更進(jìn)一步地為了實(shí)現(xiàn) “全量 增量” 的處理流程開發(fā)者仍容易陷入經(jīng)典的 Lambda 架構(gòu)困境先編寫一個 Flink 批作業(yè)處理歷史數(shù)據(jù)以生成基線再另起一個流作業(yè)消費(fèi)實(shí)時日志進(jìn)行增量更新。盡管兩者的處理邏輯高度相似但由于執(zhí)行模式有界 vs 無界、資源配置和運(yùn)維方式的不同不得不重復(fù)開發(fā)、分別部署和獨(dú)立維護(hù)極易因版本錯配或參數(shù)差異導(dǎo)致邏輯偏差進(jìn)而引發(fā)數(shù)據(jù)不一致問題。這種困境的背后是流與批在語義層面的深層差異批處理通常基于靜態(tài)快照進(jìn)行一次性計(jì)算而流處理則依托事件時間、窗口觸發(fā)與持續(xù)更新機(jī)制具有動態(tài)性和狀態(tài)演化特性。在遲到數(shù)據(jù)處理、重復(fù)記錄去重、聚合結(jié)果修正等場景中兩者的處理邏輯天然不同。即使使用同一段 SQL若未嚴(yán)格約定時間語義、窗口策略與聚合行為仍可能產(chǎn)出不一致的結(jié)果。因此盡管 Flink 在運(yùn)行時層面實(shí)現(xiàn)了流與批執(zhí)行模型的統(tǒng)一顯著推進(jìn)了架構(gòu)收斂但距離開發(fā)者所期待的 “一套邏輯、一次開發(fā)、全域生效” 的理想目標(biāo)仍有不小的差距。2.6.3 Flink、Spark和阿里巴巴流批一體平臺SARO對比維度FlinkSparkSaro典型場景實(shí)時風(fēng)控、IoT 邊緣計(jì)算、實(shí)時大屏、金融級一致性要求場景離線數(shù)倉 ETL、機(jī)器學(xué)習(xí)訓(xùn)練、報表統(tǒng)計(jì)等高吞吐批場景中大型企業(yè)需快速構(gòu)建流批一體鏈路的業(yè)務(wù)團(tuán)隊(duì)如電商價格 / 庫存域核心思想將批視為有界流的特例把流 “切片” 成小批量再復(fù)用批處理引擎執(zhí)行本質(zhì)是 “用批模擬流”并非底層計(jì)算引擎而是基于 Flink 構(gòu)建的流批一體數(shù)據(jù)開發(fā)平臺。它不替代 Flink而是對其能力進(jìn)行封裝與增強(qiáng)通過 DSL 抽象、可視化編排、自動任務(wù)生成批 流雙作業(yè)等方式降低流批一體落地門檻。優(yōu)勢低延遲毫秒級、強(qiáng)狀態(tài)管理、事件時間語義完備天然支持事件時間Event Time、亂序處理、精確一次Exactly-once語義及低延遲狀態(tài)管理生態(tài)成熟、SQL 優(yōu)化強(qiáng)大、社區(qū)活躍、批處理性能極致開發(fā)效率高拖拽 DSL、運(yùn)維標(biāo)準(zhǔn)化、支持熱點(diǎn)隔離與成本優(yōu)化