
「我的數據空間」實時計算實踐筆記 · Flink SQL 系列StreamingSQL和BatchSQL的流程對比以下是一個嵌套的FlinkSQL代碼:SELECTcnt,count(cnt)ASfreqFROM(SELECTword,COUNT(num)AScntFROMTableGROUPBYword)GROUPBYcnt;計算不同的單詞出現的頻率(task 1)計算不同頻率下的單詞的個數(task 2)原始的數據:wordnumHello1Bob1Word1Hello1batch的結果輸出:計算不同單詞的出現頻率, task1的輸出如下所示:wordnumHello2Bob1World1計算不同頻率下的個數, task2的輸出如下所示:cntfreq1221streaming的結果輸出:對于Streaming作業來說數據是源源不斷的 因此寫到下游的結果也是源源不斷的. 因此對于StreamingSQL無法像batchSQL那樣只產生一次的結果輸出.sourcetask 1task 2消息編號wordnumwordnumcntfreq1hello1hello1112word1word1123bob1bob1134hello1hello22112如上表格所示為不同的消息進入系統之后各個task產生的輸出情況.前三條消息到來之后 task1 和task2的輸出比較容易理解是正常的累加邏輯.重點看第四條消息到來時的各個task的輸出:task 1的輸出為 hello, 2 這里比較好理解 由于是按字段進行累計和count task1緩存了 hello, 1的狀態 當 hello, 1消息進入時進行累計計算 輸出為 hello, 2.task 2 在接受到 (hello,2)之后 輸出 ( 2 1) , 并同時產生(1, 2) 用來覆蓋掉之前的( 1, 3).retract機制實現原理解析通過上面的結果輸出 我們大致能明白 Streaming作業和batch作業兩種作業的差異, Streaming作業的結果會根據當前的實時數據不斷的去修正最終的結果.其中的關鍵問題就是: task 2 怎樣才能輸出 (1 , 2) 這條結果 可能會存在兩種方案:task2接受到(hello, 2)這條消息之后, 通過內部的狀態信息得出需要減去(hello, 1)這條消息 將(1, 3) 減去 (hello, 1)得到(1, 2).task2需要接受到上游發送過來的-(hello, 1)的消息 將 (1, 3)減去(hello, 1), 得到 (1, 2)所以以上的問題就變成了: 是由task1 還是task2來產生 -(hello, 1) 這條消息 ?由task2來產生減(hello, 1)的消息由task1來產生減(hello, 1)的消息task2產生 -(hello, 1)的消息task2中保存的狀態MapString, Integer 保存所有 word以及對應的countMapInteger, Integer 保存單詞出現個數以及對應的頻率task1保存的狀態MapString, Integer 保存所有word對應的count.task1產生-(hello, 1)的消息task2中保存的狀態MapInteger, Integer 保存單詞出現的個數以及對應的頻率task1中保存的狀態MapString, Integer 保存所有的word以及對應的count很明顯 task1中已經保存了所有的word對應的count 則task2也不需要進行保存 由task1產生 減(hello, 1)的所需要保存的狀態較少.因此 task1在接收到第四條消息時需要產生兩條消息:6. - (hello, 1)7. (hello, 2)什么場景下需要retract簡單SQLSELECTword,num%10asnumAScntFROMTable;簡單的字段轉換或者映射 不涉及到task之間task和外部系統的數據更新操作 則不需要retract.聚合SQLSELECTcnt,count(word)ASfreqFROM(SELECTword,COUNT(num)AScntFROMTableGROUPBYword);聚合操作(非時間窗口)涉及到task和task之間task和外部系統之間的數據更新 則需要retract機制.Flink框架本身支持處理和產生add, update, delete類型的消息 同時也需要最終的sink也需要能夠處理這些類型的消息 如果不能支持 則有些場景下就可能無法支持.類型特點相關系統append只能接受append消息 無法處理updatedelete消息消息隊列(Kafka), druid, opentsdbupsert可以處理 add, update, delete消息mysql, hbase, kv, es, kudu, esretract只能處理 add, delete消息 無法處理update消息print(測試用)tips:在實際的應用中 Kafka并不僅僅只能作為append表 雖然Kafka系統本身無法處理delete或者update消息 但是在實現上 可以將 append, delete, update等消息的類型也一并寫入到消息體中 由下游再去處理不同類型的消息類型即可 實現細節可以參考FlinkKafka Retract-Table支持寫retract信息到下游Kafka很少有系統真正是retract表 一般支持刪除的系統都支持處理update消息.retract無法處理update消息 如果下游是retract表 那么Flink框架會將update的消息轉化為 delete add 消息.~~如果使用了Upsert類型的sink表 一定要使用 insert into SinkTable select xxx, sum(xxx) group by xxx的寫法 讓框架能夠識別到sink表的主鍵用于優化生成的DAG圖. ~~在Flink1.12中 聲明主鍵即可本文收錄于「我的數據空間」技術庫——一套可私有化部署的數據平臺(數據集成 / 實時計算 / 數據湖 / 湖倉查詢 / 智能問數)。產品介紹見我的數據空間官網,支持私有化部署與 OEM 合作。