Flink水位線
1、Flink中不同的事件概念
- Processing time(處理時間): 即事件被機器處理的時間,事件流向某個算子的系統時間
- Event Time(事件時間): 事件時間是再某個生產設備上發生時間,指事件進入Flink之前嵌入的時間,通常可以從事件中獲取一個時間戳,此時間戳可以用來得出水位線,進而解決延遲,亂序,重發等情況
- Ingestion time(攝入時間): 攝入時間即是事件進入Flink的時間,是在Source Operator中設定的,
2、WaterMark(水位線)
水位線,首先水位線的主要作用是解決資料的延遲和亂序問題,水位線到底是什么?水位線其實可以理解是一個特俗的資料,用來延遲視窗的觸發(此處指的視窗每個相互獨立),具體情況,下圖說明,

假設場景,比如說上學校車,校車每天早上九點發車,但有一部分學生可能九點零二才能趕來,于是小明偷偷把司機的時間調后了兩秒,這樣當時間到了九點(延后兩秒的九點)大家都能上車了,
初學者:疑難點(誤區)
- ① 水位線如何得出?
水位線公式:watermark=當前最大事件時間-延遲時間,此公式最后的結果水位線其實代表了前多少個資料已經到齊了,每個資料進入,都會抽取資料的時間戳(事務時間)來生成一個水位線- ② 事件時間和資料容易分不清(或者混為一談)
上面有講事件時間,開始資料可能是有序的,但經過并行的處理程序中,資料難免會亂序,圖中很容易看出
③ 視窗的大小并不代表視窗處理幾條資料
視窗大小兩種要么是根據時間,要么根據處理資料個數,視情而定,不要搞混,
④ 視窗5s秒的理解
視窗五秒大小范圍是【0,5),
3、WaterMark的遲到資料
現實中很難生成一個完美的水位線,水位線就是在延遲與準確性之前做的一種權衡,那么,如果生成的水位線過于緊迫,即水位線可能會大于后來資料的時間戳,這就意味著資料有延遲,關于延遲資料的處理,Flink提供了一些機制,具體如下:
- ① 直接將遲到的資料丟棄
- ② 將遲到的資料輸出到單獨的資料流中,即使用sideOutputLateData(newOutputTag<>())實作側輸出
- ③ 根據遲到的事件更新并發出結果
其實就是水位線不一定能解決全部延遲資料,主要是處理毫秒內的資料延遲,延遲過大的資料會有剩余方案進行處理~
- 若還有別的疑問歡迎留言~🤗,有說錯之處還請指出🤨
轉載請註明出處,本文鏈接:https://www.uj5u.com/qita/436409.html
標籤:其他
上一篇:hive-SQL學習筆記11
