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