入門:流批一體架構(gòu)與DataStream實戰(zhàn))
很多準(zhǔn)備上手 Flink 的朋友一開始都會問同一個問題我已經(jīng)會寫 Spark也處理過離線數(shù)倉的 ETL為什么還要再學(xué)一個流處理框架我第一次接觸 Flink 時也有類似的困惑甚至一度覺得它不過是把 Spark Streaming 里的微批窗口又做了一遍。直到真的把線上風(fēng)控、實時指標(biāo)這類場景跑在 Flink 上我才意識到Flink 對“流”這件事的底層理解和 Spark 是完全不同的。這篇是 Flink 零基礎(chǔ)入門系列的第一篇我不堆概念直接把核心基礎(chǔ)、架構(gòu)原理、分層 API、核心組件這四塊講透再帶你把第一個 WordCount 程序跑起來。不管你之前是做后端、數(shù)倉還是數(shù)據(jù)分析只要想系統(tǒng)入門 Flink這篇文章都值得認真讀一遍。讀完之后你會發(fā)現(xiàn)Flink 的入門門檻不在于語法而在于你有沒有把“流”的思維方式建立起來。本文會盡量用大白話解釋那些看起來高深的名詞同時保留足夠的技術(shù)細節(jié)方便你隨時對照官方文檔繼續(xù)深挖。1. 流批一體這個概念為什么 Flink 要反復(fù)講1.1 從批與流的差別說起批處理面對的是已經(jīng)落地的歷史數(shù)據(jù)數(shù)據(jù)是“靜止”的。你可以反復(fù)掃描、排序、聚合算多少遍都行今天跑不完明天接著跑。流處理面對的是源源不斷產(chǎn)生的事件數(shù)據(jù)就像水管里的水你不可能把整條河攔下來再算只有在它經(jīng)過你面前的那一瞬間抓住它立刻做出反應(yīng)。這兩種模式對應(yīng)的是完全不同的業(yè)務(wù)訴求。離線報表關(guān)心的是昨天全天的 GMV晚兩個小時出結(jié)果問題不大但支付風(fēng)控如果不能在幾十毫秒內(nèi)識別出異常交易等報表跑完錢早就被轉(zhuǎn)走了。所以流處理的核心價值不是“快”而是“及時”。很多框架想同時兼顧這兩件事但實現(xiàn)路徑完全不同。Spark Streaming 的做法是把流拆成一個個小批次用微批調(diào)度去模擬流的效果好處是復(fù)用 Spark 的批處理生態(tài)壞處是它本質(zhì)上還是批延遲很難壓到真正的毫秒級。而 Flink 從誕生那天起就選擇了真正的流式處理模型每一條數(shù)據(jù)都會立即被處理這也是它敢說自己“低延遲”的底氣。1.2 Flink 的答案批是有界流流是廣袤宇宙Flink 里有一句非常經(jīng)典的話一切數(shù)據(jù)皆流Everything is a stream。離線批處理的數(shù)據(jù)其實只是有開始、有結(jié)束的“有界流”而實時數(shù)據(jù)是不知道什么時候結(jié)束的“無界流”。這個哲學(xué)直接決定了架構(gòu)設(shè)計。Spark 是“批為體、流為用”Flink 是“流為體、批為用”。在 Flink 中同一套執(zhí)行引擎、同一套 Runtime 既能跑流任務(wù)也能跑批任務(wù)。你寫一套 DataStream API 的代碼把輸入源換成有界數(shù)據(jù)集它就是一個批任務(wù)換成 Kafka 這種無界數(shù)據(jù)源它就變成了實時任務(wù)。換句話說你不需要為批處理和流處理維護兩套代碼也不需要引入兩個不同的集群。這一點在 Flink 1.12 之后被推到了前臺社區(qū)明確提出“流批一體”的概念之后 Flink SQL 也在逐步統(tǒng)一流批的行為。比如同樣是GROUP BY在有界輸入上就是離線聚合在無界輸入上就是持續(xù)更新的結(jié)果表上層用戶的開發(fā)體驗是一致的。對于很多既有離線數(shù)倉、又要建設(shè)實時數(shù)倉的團隊來說Flink 是他們能在同一套技術(shù)棧里同時解決兩類需求的現(xiàn)實選擇。1.3 適合 Flink 的典型場景Flink 最常見的落地場景我列幾個自己接觸過的你對照看看自己屬于哪一類實時數(shù)倉從 Kafka 接入業(yè)務(wù)日志和變更數(shù)據(jù)在 Flink 里做清洗、關(guān)聯(lián)、聚合結(jié)果寫入 ClickHouse、Doris、StarRocks 或 Hive供 BI 報表、實時大屏使用。實時風(fēng)控交易、登錄、領(lǐng)取優(yōu)惠券等事件發(fā)生后Flink 在毫秒級內(nèi)完成規(guī)則匹配。比如同一設(shè)備短時間內(nèi)多次登錄不同賬號這樣的模式用傳統(tǒng)批處理很難攔截但在 Flink 里做一個滾動窗口判斷即可。用戶行為分析埋點日志進入 Flink 后實時計算 PV、UV、轉(zhuǎn)化率、留存率運營側(cè)可以即時看到活動效果并調(diào)整策略。CDC 數(shù)據(jù)同步通過 Flink CDC 捕獲數(shù)據(jù)庫 binlog 的變更把數(shù)據(jù)實時同步到數(shù)倉或下游系統(tǒng)替代傳統(tǒng)基于定時任務(wù)的增量同步。監(jiān)控與告警業(yè)務(wù)指標(biāo)異常檢測、日志錯誤率突增、訂單量暴跌這類場景都需要對數(shù)據(jù)流做持續(xù)計算并及時觸發(fā)告警。這些場景有一個共同特點數(shù)據(jù)的價值隨時間快速衰減。等幾個小時后再算出來結(jié)論可能已經(jīng)沒用了。Flink 的價值就是把這個衰減周期壓縮到秒級甚至毫秒級。2. 從提交一個 Jar 開始看 Flink 內(nèi)部發(fā)生了什么2.1 三個大角色Client、JobManager、TaskManagerFlink 的集群架構(gòu)并不復(fù)雜核心角色就三個Client、JobManager、TaskManager。你可以把它想象成一個軟件開發(fā)團隊Client 是需求方負責(zé)把項目方案提交上來JobManager 是項目經(jīng)理負責(zé)拆任務(wù)、排計劃、協(xié)調(diào)資源TaskManager 是干活的一線程序員真正執(zhí)行每一個計算單元。具體到數(shù)據(jù)流上任務(wù)運行的完整鏈路是這樣的你寫好 Flink 代碼打包成 Jar通過flink run命令提交。Client 負責(zé)解析代碼、生成作業(yè)圖JobGraph并把作業(yè)圖提交給 JobManager。JobManager 收到作業(yè)后會生成可執(zhí)行的執(zhí)行圖ExecutionGraph并向 ResourceManager 申請所需的資源。TaskManager 有空閑的 Task Slot 時JobManager 把任務(wù)部署到對應(yīng)的 Slot 上。TaskManager 開始消費數(shù)據(jù)、執(zhí)行計算并把狀態(tài)信息匯報給 JobManager。這里有一個容易混淆的點JobManager 并不直接參與數(shù)據(jù)處理。它只負責(zé)調(diào)度、協(xié)調(diào)、故障恢復(fù)真正干活的是 TaskManager。很多人剛開始以為 JobManager 是“主節(jié)點負責(zé)大部分計算”這是完全錯誤的。類比成項目經(jīng)理也很貼切項目經(jīng)理不會自己寫代碼但他知道進度、知道誰在做什么、出了問題能立刻調(diào)整方案。2.2 一張圖是怎么從代碼變成可執(zhí)行計劃的Flink 里還有一組容易讓人暈頭的名詞StreamGraph、JobGraph、ExecutionGraph、物理執(zhí)行圖。其實它們只是同一個作業(yè)在不同階段的“表現(xiàn)形式”。StreamGraph由用戶代碼直接生成的邏輯圖。你在代碼里寫了多少個算子它就對應(yīng)多少個節(jié)點此時還沒有做任何優(yōu)化。JobGraph經(jīng)過合并優(yōu)化后的圖。Flink 會把符合條件的相鄰算子“串”在一起減少線程切換和網(wǎng)絡(luò)傳輸開銷。ExecutionGraphJobManager 拿到 JobGraph 后根據(jù)并行度把它拆分成可調(diào)度的任務(wù)并加上并行實例形成執(zhí)行圖。物理執(zhí)行圖任務(wù)真正部署到 TaskManager 上之后每個并行子任務(wù)實際運行的形態(tài)。這四個階段里算子鏈Operator Chain是最值得理解的一個優(yōu)化手段。舉個例子一個任務(wù)里有map和filter兩個相鄰算子它們之間沒有重新分區(qū)即不是keyBy而是直接串聯(lián)Flink 就會把它們放進同一個線程里執(zhí)行。你少了一次線程切換和一次網(wǎng)絡(luò)序列化性能提升是實打?qū)嵉摹T?Web UI 上看到的任務(wù)鏈很多就是這種算子鏈合并后的結(jié)果。2.3 Task Slot 和并行度資源和并發(fā)不是一回事接著說說最容易混淆的兩個概念Task Slot 和并行度。Task Slot 是 TaskManager 內(nèi)部的資源劃分單元。比如一個 TaskManager 配置了 4 個 Slot那它最多可以同時運行 4 個并行任務(wù)。Slot 限制的是“能同時跑多少任務(wù)”是靜態(tài)的資源池。并行度是算子被拆分成多少份同時執(zhí)行。比如一個map算子的并行度是 4那就表示這個map會在 4 個 Slot 里各跑一個實例每個實例處理一部分數(shù)據(jù)。它們之間的關(guān)系可以用一句話概括并行度決定了需要多少 SlotSlot 決定了能不能跑得下。假設(shè)你的 Flink 集群有 3 個 TaskManager每個 8 個 Slot那一共是 24 個 Slot如果你的作業(yè)全局并行度設(shè)置成了 30那在提交階段就會因為資源不足而失敗日志會明確提示沒有足夠的 Slot。這個坑我見過太多次。很多人把并行度無限調(diào)大以為能提升性能結(jié)果瓶頸反而出現(xiàn)在資源不夠或者單 Slot 上的任務(wù)爭搶。合理經(jīng)驗是并行度要結(jié)合數(shù)據(jù)量和資源量一起看不是越大越好。比如 Kafka 分區(qū)數(shù)是 12你給 Source 算子設(shè)置并行度 12 是最理想的設(shè)置 24 反而有大量空閑任務(wù)在空轉(zhuǎn)。3. 分層 API 不是四層樓梯而是一條決策鏈3.1 四層 API 各自適合什么Flink 提供了從高到低四種 API很多人喜歡把它們畫成四層金字塔但我覺得更準(zhǔn)確的比喻是一條“決策鏈”從最省事的 SQL 開始當(dāng)你發(fā)現(xiàn)它表達不了你的邏輯時就往下走一層。這四層分別是Flink SQL最高層聲明式 API。你只需要告訴它“我要什么”不需要關(guān)心“怎么做”。適合絕大多數(shù)聚合統(tǒng)計、關(guān)聯(lián)查詢、過濾清洗場景。Table APISQL 的編程化形式或者說“嵌入 Java/Scala 里的 DSL”。它比 SQL 表達能力強一點可以寫一些復(fù)雜的 UDF同時還能享受優(yōu)化器的自動優(yōu)化。DataStream APIFlink 的核心 API。它提供細粒度的算子map、flatMap、keyBy、window、process 等你可以精確控制數(shù)據(jù)轉(zhuǎn)換、狀態(tài)、時間語義是大多數(shù)實時業(yè)務(wù)的主戰(zhàn)場。ProcessFunction最底層、最靈活的接口。它允許你訪問事件的時間戳、Watermark、定時器、狀態(tài)實現(xiàn)任意復(fù)雜的業(yè)務(wù)邏輯。代價是代碼量最大優(yōu)化全靠自己。我可以給一個非常直觀的對比表API 層級開發(fā)效率表達能力適合人群典型場景Flink SQL最高中數(shù)據(jù)分析師、BI 開發(fā)者實時報表、指標(biāo)統(tǒng)計、簡單 ETLTable API高中高希望用聲明式但又需要編程能力的開發(fā)者有 UDF 的流批任務(wù)DataStream API中高Java/Scala 后端工程師復(fù)雜流處理、狀態(tài)計算、窗口邏輯ProcessFunction低最高框架開發(fā)者、極致定制場景精確時間控制、自定義狀態(tài)、定時器3.2 最常見的組合玩法Table 和 DataStream 互轉(zhuǎn)很多人誤以為四層 API 是互相隔離的選了一層就不能用另一層。實際上 Flink 允許你在同一個作業(yè)里混合使用它們。最常見的組合就是 DataStream 和 Table 互相轉(zhuǎn)換。比如你的數(shù)據(jù)源是 Kafka里面是一條條 JSON 日志。你可以先用 DataStream API 讀取并用 POJO 解析然后把它注冊成一張臨時表后續(xù)用 SQL 去做聚合聚合結(jié)果如果只是簡單的指標(biāo)直接execute().print()輸出即可但如果你想對結(jié)果做更精細的加工又可以轉(zhuǎn)回 DataStream。代碼看起來是這樣StreamExecutionEnvironment env StreamExecutionEnvironment.getExecutionEnvironment(); StreamTableEnvironment tableEnv StreamTableEnvironment.create(env); DataStreamEvent eventStream env.addSource(...); // DataStream 轉(zhuǎn) Table Table eventTable tableEnv.fromDataStream(eventStream); tableEnv.createTemporaryView(events, eventTable); // 用 SQL 聚合 Table result tableEnv.sqlQuery( SELECT user_id, COUNT(*) AS cnt FROM events GROUP BY user_id ); // Table 轉(zhuǎn)回 DataStream追加流 DataStreamRow resultStream tableEnv.toDataStream(result);這種組合方式在實際項目中非常普遍。因為純 SQL 在處理復(fù)雜事件時間、多流 order 的時候很別扭而純 DataStream 寫聚合又太啰嗦。兩頭各取所長是生產(chǎn)環(huán)境最常見的寫法。3.3 實際項目里我推薦怎么選從我自己的經(jīng)驗出發(fā)我給零基礎(chǔ)同學(xué)一個非常實際的選型建議如果你的需求是“把一份數(shù)據(jù)做過濾、字段映射、聚合然后寫出去”直接上 Flink SQL。它內(nèi)置了大量連接器Kafka、JDBC、Elasticsearch、Hive 都有現(xiàn)成的配上 CREATE TABLE 的 DDL 就能跑學(xué)習(xí)成本極低。如果需求里出現(xiàn)了“多流關(guān)聯(lián)”“狀態(tài)去重”“自定義窗口觸發(fā)”這類需要精細控制的邏輯用 DataStream API。它比 SQL 更靈活多寫一些代碼但可掌控性高得多。如果需求里涉及“在某個精確時間點做某事”“根據(jù)業(yè)務(wù)規(guī)則動態(tài)生成定時器”等那就必須上 ProcessFunction。它是唯一能直接訪問定時器和原始事件的 API。很多人有個誤區(qū)覺得“既然 Flink SQL 這么方便是不是只學(xué) SQL 就夠了”。我的看法是如果你永遠只做報表類應(yīng)用SQL 確實夠用但只要你需要排查線上問題、定位某個算子為什么延遲高、理解狀態(tài)為什么會膨脹你就必須懂 DataStream 和底層架構(gòu)。SQL 給你的是效率底層知識給你的是護城河。4. 核心組件逐個過JobManager、TaskManager 和它們背后的機制4.1 JobManager 不是一個人是一個治理團隊表面上我們只說“JobManager”但它的內(nèi)部實際上是三個角色在協(xié)同工作Dispatcher、ResourceManager、JobMaster。Dispatcher負責(zé)接收用戶提交的作業(yè)。你可以把它理解成前臺接待員拿到 Jar 之后它會把作業(yè)交給對應(yīng)的 JobMaster 去管理。ResourceManager負責(zé)資源分配。當(dāng) JobMaster 需要 Slot 來執(zhí)行任務(wù)時它會向 ResourceManager 申請如果資源不夠它還會嘗試向 Kubernetes 或 YARN 請求擴容。JobMaster每個作業(yè)都有一個 JobMaster它是這個作業(yè)的總負責(zé)人管理作業(yè)的整個生命周期包括調(diào)度任務(wù)、協(xié)調(diào) Checkpoint、處理故障恢復(fù)。這樣拆開之后很多現(xiàn)象就解釋得通了。比如你同時提交多個作業(yè)每個作業(yè)會有一個獨立的 JobMaster它們共享同一個 Dispatcher 和 ResourceManager。如果其中一個作業(yè)因為邏輯問題導(dǎo)致頻繁失敗理論上并不會拖垮其他作業(yè)只要 Slot 資源充足。在部署層面JobManager 進程通常是獨立于 TaskManager 的。Flink 官方建議 JobManager 和 TaskManager 分開部署不能放在同一臺物理機上否則大數(shù)據(jù)量下 TaskManager 的 CPU 和內(nèi)存壓力會直接影響調(diào)度穩(wěn)定性。這一點在 Standalone 集群、YARN 或 K8s 部署時都是同樣的原則。4.2 TaskManager、狀態(tài)后端和容錯TaskManager 是真正干活的進程。一個 TaskManager 內(nèi)部會劃分出多個 Task Slot每個 Slot 可以跑一個并行子任務(wù)。前面說過Slot 是靜態(tài)資源隔離單元它限制的是并發(fā)任務(wù)數(shù)不是 CPU 配額。如果你想讓某個任務(wù)獨占 CPU那需要借助更底層的資源調(diào)度能力而不是單純靠 Slot 配置。理解了 TaskManager 之后就必須聊一聊狀態(tài)和容錯因為它們密不可分。Flink 是我們說的“有狀態(tài)流計算”狀態(tài)就是算子保存的中間結(jié)果一個sum算子需要記住當(dāng)前累加值一個窗口算子需要保存窗口里收集到的數(shù)據(jù)一個去重算子需要記錄已經(jīng)見過的 key。沒有狀態(tài)流處理就退化成了純粹的“一條條數(shù)據(jù)獨立處理”很多業(yè)務(wù)邏輯根本做不了。狀態(tài)存放的位置由狀態(tài)后端決定。Flink 提供了兩種主流選擇HashMapStateBackend狀態(tài)全部存在內(nèi)存里讀寫極快適合狀態(tài)量不大的場景。缺點是狀態(tài)一大就容易撐爆內(nèi)存。EmbeddedRocksDBStateBackend狀態(tài)存在本地磁盤RocksDB同時利用內(nèi)存做緩存適合 GB 甚至 TB 級別的超大狀態(tài)。代價是序列化和磁盤 IO 有一定開銷。很多人在最開始都默認用內(nèi)存狀態(tài)后端直到線上狀態(tài)漲到幾個 GB 才發(fā)現(xiàn)問題。我的經(jīng)驗是如果狀態(tài)小于幾百 MB內(nèi)存后端沒問題一旦你預(yù)估狀態(tài)可能超過 1 GB直接上 RocksDB別抱有僥幸心理。容錯機制的核心是Checkpoint。Flink 會定期把每個算子的狀態(tài)做一次快照保存到外部存儲比如 HDFS 或 S3。當(dāng)任務(wù)因機器宕機、網(wǎng)絡(luò)抖動失敗時JobManager 會從最近一次成功的 Checkpoint 恢復(fù)狀態(tài)然后重新拉起任務(wù)實現(xiàn)“精確一次”的語義。這里插一句Checkpoint 的間隔很有講究太頻繁會增加存儲和網(wǎng)絡(luò)開銷太稀疏又會拉長故障恢復(fù)時間生產(chǎn)中一般從 1 分鐘起步觀察后再調(diào)整。4.3 反壓Backpressure的簡要理解反壓是 Flink 里一個重要又容易被忽視的機制。簡單說當(dāng)下游算子的處理速度跟不上上游數(shù)據(jù)產(chǎn)生的速度時壓力會逐級向上傳導(dǎo)最終讓數(shù)據(jù)源放慢消費或者暫停消費。這就像一條水管出口接了一個很細的噴嘴水流速度超過噴嘴的排出能力水就會在管道里積壓最終倒灌回源頭。反壓不是一種錯誤而是系統(tǒng)自我保護的方式。如果沒有反壓機制數(shù)據(jù)會在內(nèi)存里無限積壓最終導(dǎo)致 OOM。在 Flink Web UI 的 Job 頁面有一個“Backpressure”標(biāo)簽可以直接看到哪些算子處于高反壓狀態(tài)。遇到反壓時優(yōu)先檢查下游算子的處理邏輯是不是出現(xiàn)了瓶頸比如 SQL 里某個 UDF 太慢、某個 keyBy 熱點導(dǎo)致數(shù)據(jù)傾斜而不是盲目加并行度。加并行度只能解決資源不夠解決不了單個算子邏輯太慢的問題。5. 本地跑通第一個 Flink 程序環(huán)境準(zhǔn)備與完整實操5.1 環(huán)境準(zhǔn)備與版本選擇動手之前先把環(huán)境裝好。基礎(chǔ)要求并不高JDK 1.8 或 11、17我推薦 11兼容性和性能比較均衡Maven 3.6IDEA 或 Eclipse推薦 IDEAFlink 插件體驗更好Flink 發(fā)行版去官網(wǎng)下載一個穩(wěn)定版壓縮包即可比如 1.18.x 或 1.19.x版本選擇這里多說一句。Flink 的版本迭代節(jié)奏不算快但不同小版本之間 API 有細微差異。別一上來就追最新版盡量選當(dāng)前社區(qū)主推的穩(wěn)定版本比如 1.18/1.19 這條線。后續(xù)如果想升級再對照官方遷移指南改。還有一點很重要Flink 版本和連接器版本尤其是 Kafka、JDBC、CDC 連接器必須匹配。連接器版本不對是最常見的“NoSuchMethodError”“ClassNotFoundException”來源。5.2 創(chuàng)建工程依賴該怎么配在 IDEA 里新建一個普通的 Maven 項目Java 版本選 8 或 11。pom.xml最小依賴如下properties flink.version1.18.1/flink.version maven.compiler.source1.8/maven.compiler.source maven.compiler.target1.8/maven.compiler.target /properties dependencies dependency groupIdorg.apache.flink/groupId artifactIdflink-streaming-java/artifactId version${flink.version}/version /dependency dependency groupIdorg.apache.flink/groupId artifactIdflink-clients/artifactId version${flink.version}/version /dependency dependency groupIdorg.apache.flink/groupId artifactIdflink-table-api-java-bridge/artifactId version${flink.version}/version /dependency /dependencies注意我在這里沒有把 scope 設(shè)置為 provided原因是為了方便在 IDEA 里直接運行。如果你要打包提交到服務(wù)器上的 Flink 集群再把 scope 改為 provided并加上 maven-shade-plugin 做打包避免和集群自帶的依賴沖突。build plugins plugin groupIdorg.apache.maven.plugins/groupId artifactIdmaven-shade-plugin/artifactId version3.4.1/version executions execution phasepackage/phase goalsgoalshade/goal/goals /execution /executions /plugin /plugins /build這個打包插件是必須的否則你提交到集群后會選不到主類或找不到依賴。5.3 用 DataStream API 寫一個流式 WordCount接下來寫一個經(jīng)典的流式 WordCount。我們從一個 Socket 端口讀取文本按空格拆分單詞統(tǒng)計每個單詞出現(xiàn)的次數(shù)。import org.apache.flink.api.common.functions.FlatMapFunction; import org.apache.flink.api.java.tuple.Tuple2; import org.apache.flink.streaming.api.datastream.DataStream; import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment; import org.apache.flink.util.Collector; public class SocketWordCount { public static void main(String[] args) throws Exception { // 1. 創(chuàng)建流執(zhí)行環(huán)境 StreamExecutionEnvironment env StreamExecutionEnvironment.getExecutionEnvironment(); // 2. 從本機 9999 端口讀取文本流 DataStreamString text env.socketTextStream(localhost, 9999); // 3. 拆詞、計數(shù) DataStreamTuple2String, Integer counts text .flatMap(new Tokenizer()) .keyBy(value - value.f0) .sum(1); // 4. 打印結(jié)果 counts.print(); // 5. 提交作業(yè) env.execute(Socket WordCount); } public static final class Tokenizer implements FlatMapFunctionString, Tuple2String, Integer { Override public void flatMap(String value, CollectorTuple2String, Integer out) { String[] tokens value.toLowerCase().split(\\W); for (String token : tokens) { if (token.length() 0) { out.collect(Tuple2.of(token, 1)); } } } } }這段代碼有幾個關(guān)鍵點值得注意socketTextStream是一個本地調(diào)試用的 Source生產(chǎn)環(huán)境基本不會用但作為入門演示非常直觀。flatMap里我把每個單詞拆出來并輸出(word, 1)。keyBy按單詞分組相當(dāng)于 SQL 里的GROUP BY。注意 Flink 1.12 之后keyBy必須返回KeySelector而不能再直接用字段名做字符串參數(shù)除非用 Table API。sum(1)是對 Tuple 的第 2 個字段做累加。在 IDEA 里直接運行main方法之前需要先在本地啟動一個 Socket 服務(wù)。在 macOS/Linux 終端執(zhí)行nc -lk 9999運行后在終端輸入任意英文句子回車你就能在 IDEA 控制臺看到每個單詞的累計次數(shù)。輸入flink flink learning再輸入flink spark輸出會清楚地展示flink的計數(shù)如何累加。這就是流式計數(shù)的基本形態(tài)。5.4 用 Flink SQL 寫同一個 WordCount同樣的邏輯用 Flink SQL 寫非常簡單。因為本地不好演示 Socket 的 SQL 連接器我用 Datagen 連接器生成模擬數(shù)據(jù)import org.apache.flink.table.api.EnvironmentSettings; import org.apache.flink.table.api.Table; import org.apache.flink.table.api.TableEnvironment; public class SqlWordCount { public static void main(String[] args) { // 創(chuàng)建 SQL 執(zhí)行環(huán)境 TableEnvironment tableEnv TableEnvironment.create(EnvironmentSettings.inStreamingMode()); // 定義一個數(shù)據(jù)生成表 tableEnv.executeSql( CREATE TABLE source_table ( word STRING ) WITH ( connector datagen, fields.word.length 5 )); // 查詢統(tǒng)計 Table result tableEnv.sqlQuery( SELECT word, COUNT(*) AS cnt FROM source_table GROUP BY word); // 打印結(jié)果 result.execute().print(); } }這個作業(yè)的意義不在于輸出多少結(jié)果而在于讓你感受 Flink SQL 的開發(fā)模式建表、寫查詢、拿結(jié)果。表結(jié)構(gòu)、連接器參數(shù)都在 DDL 里聲明邏輯和運行分離這是它和 DataStream API 最大的不同。5.5 提交運行從 IDE 到獨立集群在 IDEA 里跑通是第一步但真實項目一定要學(xué)會把作業(yè)提交到集群。先把 Flink 發(fā)行版解壓直接啟動本地集群tar -xzf flink-1.18.1-bin-scala_2.12.tgz cd flink-1.18.1 ./bin/start-cluster.sh啟動之后打開http://localhost:8081這就是 Flink Web UI可以看到 JobManager 狀態(tài)、TaskManager 數(shù)量和 Slot 數(shù)。然后打包你的 Maven 工程mvn clean package最后提交作業(yè)./bin/flink run -d -c com.example.SocketWordCount ./target/flink-quickstart-1.0-SNAPSHOT.jar-d表示后臺運行-c指定主類。提交成功后回到 Web UI 的 Jobs 頁面你能看到一個 RUNNING 狀態(tài)的作業(yè)。點進去可以看到執(zhí)行圖、并行度、每個算子的輸入輸出條數(shù)、Checkpoint 狀態(tài)、反壓情況。強烈建議每個初學(xué)者提交一個作業(yè)后把 Web UI 的每個頁面都點一遍這是理解 Flink 運行時最快的方式比看十篇文章都有用。6. 入門期最容易卡住的四個問題都是真實踩過的坑6.1 Watermark為什么我的窗口遲遲不觸發(fā)Watermark 是 Flink 入門階段最難啃、也是問得最多的問題。網(wǎng)上關(guān)于 Flink SQL 中 Watermark 的搜索量一直很高原因就是大家在寫窗口聚合時經(jīng)常發(fā)現(xiàn)“窗口怎么不輸出結(jié)果”。先理解概念流處理里有兩個時間事件時間Event Time和處理時間Processing Time。處理時間是數(shù)據(jù)到達 Flink 機器的那一刻簡單但不可靠事件時間是業(yè)務(wù)真正發(fā)生的時間才是大多數(shù)業(yè)務(wù)需要的。由于網(wǎng)絡(luò)延遲、數(shù)據(jù)亂序事件時間晚到的數(shù)據(jù)會打破窗口的邊界所以 Flink 引入了 Watermark用來表示“在這個時間戳之前的數(shù)據(jù)都應(yīng)該到了”。Watermark 就像一個“遲到截止線”它等于當(dāng)前觀察到的事件時間減去一個允許亂序的閾值。一個典型的 SQL 建表語句是這樣CREATE TABLE orders ( order_id BIGINT, order_time TIMESTAMP(3), amount DOUBLE, WATERMARK FOR order_time AS order_time - INTERVAL 5 SECOND ) WITH (...)這里WATERMARK允許最多 5 秒的亂序數(shù)據(jù)。如果你設(shè)置的窗口是 1 分鐘滾動窗口那只有當(dāng) Watermark 超過窗口結(jié)束時間窗口才會觸發(fā)計算并輸出結(jié)果。窗口不觸發(fā)的原因八成是 Watermark 一直沒有推進比如數(shù)據(jù)源里根本沒有事件時間字段、或者事件時間字段格式不對、或者數(shù)據(jù)本身長時間斷流。遇到這類問題先查數(shù)據(jù)源的時間戳解析再查 Watermark 閾值是否設(shè)置得太大。6.2 JDBC 連接器報錯大多逃不過這幾種原因Flink SQL 里的 JDBC 連接器比如 MySQL、PostgreSQL、TiDB是高頻使用組件但報錯率也極高。我見過最多的幾類異常異常現(xiàn)象根本原因解決辦法ClassNotFoundException: com.mysql.cj.jdbc.Driver缺少 JDBC 驅(qū)動把 mysql-connector-java 顯式加進依賴或放入 Flink lib 目錄Communications link failure網(wǎng)絡(luò)不通、URL 配錯、驅(qū)動版本太老確認數(shù)據(jù)庫地址端口可達URL 加上 useSSLfalseserverTimezoneAsia/ShanghaiConnection is not available, request timed out連接池被占滿下游寫入太慢調(diào)大 sink 并行度或者優(yōu)化批量寫入?yún)?shù)Data truncation / 字段類型不匹配表結(jié)構(gòu)與 DDL 字段不一致核對 Flink 側(cè)聲明的類型與數(shù)據(jù)庫側(cè)字段類型一一對應(yīng)JDBC 連接器和 Flink 的其它連接器不太一樣它本身是一個“查詢/寫入型”連接器不適合做超大吞吐的維表。如果維表數(shù)據(jù)量大、需要頻繁關(guān)聯(lián)記得加上緩存在 DDL 里配置lookup.cache.max-rows和lookup.cache.ttl能顯著降低連接數(shù)和延遲。6.3 SQL 夠用了為啥還要學(xué) DataStream這是我被問到最多次的問題。本質(zhì)上Flink SQL 是給“大部分場景”準(zhǔn)備的它足夠簡單、足夠快但有一類問題是它很難覆蓋的復(fù)雜的自定義邏輯和細粒度的狀態(tài)控制。舉個例子你要判斷一個用戶在 10 分鐘內(nèi)是否連續(xù)點擊了 5 次同一商品并且點擊間隔不超過 2 秒。這個邏輯用 SQL 也能寫但窗口定義、狀態(tài)清理、時間判斷會讓 SQL 變得極其抽象可讀性也很差。用 DataStream API 加一個 ProcessFunction你可以輕松維護一個用戶維度的狀態(tài)注冊定時器在自定義事件里精確控制判斷邏輯。另外SQL 的執(zhí)行計劃是優(yōu)化器生成的你無法直接控制每個算子的行為而 DataStream API 讓你清楚地知道數(shù)據(jù)從哪里來、經(jīng)過哪些算子、最終到哪里去。排查線上問題時DataStream 的直觀性優(yōu)勢非常明顯。我的建議是學(xué)的時候不要把 SQL 和 DataStream 對立起來兩者是互補關(guān)系。先用 SQL 快速實現(xiàn)遇到瓶頸再下沉到 DataStream 精確控制。6.4 下一步往哪走CDC、SQL Gateway、數(shù)據(jù)血緣等你把基礎(chǔ)概念都吃透就可以開始探索 Flink 生態(tài)里更有意思的方向了。Flink CDC是目前最熱的方向它可以直接捕獲 MySQL、PostgreSQL 等數(shù)據(jù)庫的 binlog 變更把增量數(shù)據(jù)流式同步到數(shù)據(jù)倉庫或消息隊列。相比傳統(tǒng)基于時間戳的增量同步CDC 更實時、更完整而且不會對源庫造成額外查詢壓力。很多團隊甚至直接用 Flink CDC 替代了曾經(jīng)的 CanalKafka 鏈路。Flink SQL Gateway解決的是 SQL 作業(yè)管理的問題。以前每個人在本地寫 SQL、在控制臺提交腳本散落各處。有了 SQL Gateway你可以通過 REST 接口提交 SQL 作業(yè)、查詢結(jié)果甚至把它集成到自己公司的調(diào)度平臺里統(tǒng)一管理 SQL 任務(wù)的生命周期。數(shù)據(jù)血緣是肉眼可見的方向。Flink 在 SQL 執(zhí)行計劃中天然攜帶了數(shù)據(jù)來源、字段映射關(guān)系、輸出去向信息通過解析這些元數(shù)據(jù)可以做字段級血緣追蹤。對于數(shù)據(jù)治理要求嚴格的公司這是剛需能力。還有一個很常見的組合是 Flink TiDB用 Flink SQL 的 JDBC 連接器或 TiDB Connector 直接做實時數(shù)據(jù)寫入和維表查詢用 TiDB 的分布式事務(wù)能力承接 Flink 的下游存儲。這類“Flink 數(shù)據(jù)庫”的整合方案在實時數(shù)倉架構(gòu)里越來越常見值得花時間研究。最后分享一點我在項目里的個人體會Flink 入門最忌諱“只看不跑”。你把本文的 WordCount 跑通再把 Web UI 每個頁面點一遍接著親手改幾個并行度、把狀態(tài)后端從內(nèi)存換成 RocksDB、在 SQL 里加一個 Watermark這些動作比讀十遍概念都有用。等這些基本功都扎實了再回頭看 Kafka 連接器、Flink CDC、SQL Gateway 這些生態(tài)組件你會有一種“原來如此”的通透感。希望這篇系列第一篇能幫你在 Flink 這條路上開個好頭。