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