一、簡介
開源流式處理系統在不斷地發展,從一開始只關注低延遲指標到現在兼顧延遲、吞吐與結果準確性,在發展程序中解決了很多問題,編程API的易用性也在不斷地提高,本文介紹一下 Flink 中的核心概念,這些概念是學習與使用 Flink 十分重要的基礎知識,在后續開發 Flink 程式程序中將會幫助開發人員更好地理解 Flink 內部的行為和機制,
這里參考一張圖來對常用的實時計算框架做個對比:

Flink 是有狀態的和容錯的,可以在維護一次應用程式狀態的同時無縫地從故障中恢復,它支持大規模計算能力,能夠在數千個節點上并發運行,它具有很好的吞吐量和延遲特性,同時,Flink 提供了多種靈活的視窗函式,Flink 在流式計算里屬于真正意義上的單條處理,每一條資料都觸發計算,而不是像 Spark 一樣的 Mini Batch 作為流式處理的妥協,Flink的容錯機制較為輕量,對吞吐量影響較小,而且擁有圖和調度上的一些優化,使得 Flink 可以達到很高的吞吐量,而 Strom 的容錯機制需要對每條資料進行ack,因此其吞吐量瓶頸也是備受詬病,
二、作業原理
Flink基本作業原理如下圖:

JobClient:負責接收程式,決議和優化程式的執行計劃,然后提交執行計劃到JobManager,這里執行的程式優化是將相鄰的Operator融合,形成Operator Chain,Operator的融合可以減少task的數量,提高TaskManager的資源利用率,
JobManagers:負責申請資源,協調以及控制整個job的執行程序,具體包括,調度任務、處理checkpoint、容錯等等,
TaskManager:TaskManager運行在不同節點上的JVM行程,負責接收并執行JobManager發送的task,并且與JobManager通信,反饋任務狀態資訊,如果說JobManager是master的話,那么TaskManager就是worker用于執行任務,每個TaskManager像是一個容器,包含一個或者多個Slot,
Slot:Slot是TaskManager資源粒度的劃分,每個Slot都有自己獨立的記憶體,所有Slot平均分配TaskManager的記憶體,值得注意的是,Slot僅劃分記憶體,不涉及CPU的劃分,即CPU是共享使用,每個Slot可以運行多個task,Slot的個數就代表了一個程式的最高并行度,
Task:Task是在operators的subtask進行鏈化之后形成的,具體Flink job中有多少task和operator的并行度和鏈化的策略有關,
SubTask:因為Flink是分布式部署的,程式中的每個算子,在實際執行中被分隔為一個或者多個subtask,運算子子任務(subtask)的數量是該特定運算子的并行度,資料流在算子之間流動,就對應到SubTask之間的資料傳輸,Flink允許同一個job中來自不同task的subtask可以共享同一個slot,每個slot可以執行一個并行的pipeline,可以將pipeline看作是多個subtask的組成的,
三、核心概念
1、Time(時間語意)
Flink 中的 Time 分為三種:事件時間、達到時間與處理時間,
1)事件時間:是事件真實發生的時間,
2)達到時間:是系統接收到事件的時間,即服務端接收到事件的時間,
3)處理時間:是系統開始處理到達事件的時間,
在某些場景下,處理時間等于達到時間,因為處理時間沒有亂序的問題,所以服務端做基于處理時間的計算是比較簡單的,無遲到與亂序資料,
Flink 中只需要通過 env 環境變數即可設定Time:
//創建環境背景關系 val env = StreamExecutionEnvironment.getExecutionEnvironment // 設定在當前程式中使用 ProcessingTime env.setStreamTimeCharacteristic(TimeCharacteristic.ProcessingTime);
2、Window(視窗)
視窗本質就是將無限資料集沿著時間(或者數量)的邊界切分成有限資料集,
1)Time Window:基于時間的,分為Tumbling Window(無資料重疊)和Sliding Window(有資料重疊) ,
2)Count Window:基于數量的,分為Tumbling Window(無資料重疊)和Sliding Window(有資料重疊),
3)Session Window:基于會話的,一個session window關閉通常是由于一段時間沒有收到元素,
4)Global Window:全域視窗,
在實際操作中,window又分為兩大型別的視窗:Keyed Window 和 Non-keyed Window,兩種型別的視窗操作API有細微的差別,
3、Trigger
1)自定義觸發器
觸發器決定了視窗何時會被觸發計算,Flink 中開發人員需要在 window 型別的操作之后才能呼叫 trigger 方法傳入觸發器定義,Flink 中的觸發器定義需要繼承并實作 Trigger 介面,該介面有以下方法:
- onElement(): 每個被添加到視窗中的元素都會被呼叫
- onEventTime(): 當事件時間定時器觸發時會被呼叫,比如watermark到達
- onProcessingTime(): 當處理時間定時器觸發時會被呼叫,比如時間周期觸發
- onMerge(): 當兩個視窗合并時兩個視窗的觸發器狀態將會被調動并合并
- clear(): 執行需要清除相關視窗的事件
以上方法會回傳決定如何觸發執行的 TriggerResult:
- CONTINUE: 什么都不做
- FIRE: 觸發計算
- PURGE: 清除視窗中的資料
- FIRE_AND_PURGE: 觸發計算后清除視窗中的資料
2)預定義觸發器
如果開發人員未指定觸發器,則 Flink 會自動根據場景使用默認的預定義好的觸發器,在基于事件時間的視窗中使用 EventTimeTrigger,該觸發器會在watermark通過視窗邊界后立即觸發(即watermark出現關閉改視窗時),在全域視窗(GlobalWindow)中使用 NeverTrigger,該觸發器永遠不會觸發,所以在使用全域視窗時用戶需要自定義觸發器,
4、State
Managed State 是由flink runtime管理來管理的,自動存盤、自動恢復,在記憶體管理上有優化機制,且Managed State 支持常見的多種資料結構,如value、list、map等,在大多數業務場景中都有適用之處,總體來說是對開發人員來說是比較友好的,因此 Managed State 是 Flink 中最常用的狀態,Managed State 又分為 Keyed State 和 Operator State 兩種,
Raw State 由用戶自己管理,需要序列化,只能使用位元組陣列的資料結構,Raw State 的使用和維度都比 Managed State 要復雜,建議在自定義的Operator場景中酌情使用,
5、狀態存盤
Flink中狀態的實作有三種:MemoryState、FsState、RocksDBState,三種狀態存盤方式與使用場景各不相同,詳細介紹如下:
1)MemoryStateBackend
建構式:MemoryStateBackend(int maxStateSize, boolean asyncSnapshot)
存盤方式:State存盤于各個 TaskManager記憶體中,Checkpoint存盤于 JobManager記憶體
容量限制:單個State最大5M、maxStateSize<=akka.framesize(10M)、總大小不超過JobManager記憶體
使用場景:無狀態或者JobManager掛掉不影響的測驗環境等,不建議在生產環境使用
2)FsStateBackend
建構式:FsStateBackend(URI checkpointUri, boolean asyncSnapshot)
存盤方式:State存盤于 TaskManager記憶體,Checkpoint存盤于 外部檔案系統(本次磁盤 or HDFS)
容量限制:State總量不超過TaskManager記憶體、Checkpoint總大小不超過外部存盤空間
使用場景:常規使用狀態的作業,分鐘級的視窗聚合等,可在生產環境使用
3)RocksDBStateBackend
建構式:RocksDBStateBackend(URI checkpointUri, boolean enableincrementCheckpoint)
存盤方式:State存盤于 TaskManager上的kv資料庫(記憶體+磁盤),Checkpoint存盤于 外部檔案系統(本次磁盤 or HDFS)
容量限制:State總量不超過TaskManager記憶體+磁盤、單key最大2g、Checkpoint總大小不超過外部存盤空間
使用場景:超大狀態的作業,天級的視窗聚合等,對讀寫性能要求不高的場景,可在生產環境使用
根據業務場景需要用戶選擇最合適的 StateBackend ,代碼中只需在相應的 env 環境中設定即可:
// flink 背景關系環境變數 val env = StreamExecutionEnvironment.getExecutionEnvironment // 設定狀態后端為 FsStateBackend,資料存盤到 hdfs /tmp/flink/checkpoint/test 中 env.setStateBackend(new FsStateBackend("hdfs://ns1/tmp/flink/checkpoint/test", false))
6、Checkpoint
Checkpoint 是分布式全域一致的,資料會被寫入hdfs等共享存盤中,且其產生是異步的,在不中斷、不影響運算的前提下產生,
用戶只需在相應的 env 環境中設定即可:
// 1000毫秒進行一次 Checkpoint 操作 env.enableCheckpointing(1000) // 模式為準確一次 env.getCheckpointConfig.setCheckpointingMode(CheckpointingMode.EXACTLY_ONCE) // 兩次 Checkpoint 之間最少間隔 500毫秒 env.getCheckpointConfig.setMinPauseBetweenCheckpoints(500) // Checkpoint 程序超時時間為 60000毫秒,即1分鐘視為超時失敗 env.getCheckpointConfig.setCheckpointTimeout(60000) // 同一時間只允許1個Checkpoint的操作在執行 env.getCheckpointConfig.setMaxConcurrentCheckpoints(1)
Asynchronous Barrier Snapshots(ABS)
異步屏障快照演算法,這個演算法基本上是Chandy-Lamport演算法的變體,針對DAG(有向無環圖)的ABS演算法執行流程如下所示:
- Barrier周期性的被注入到所有的Source中,Source節點看到Barrier后,會立即記錄自己的狀態,然后將Barrier發送到Transformation Operator,
- 當Transformation Operator從某個input channel收到Barrier后,它會立刻Block住這條通道,直到所有的input channel都收到Barrier,這個等待的程序就叫做屏障對齊(barrier alignment),此時該Operator就會記錄自身狀態,并向自己的所有output channel廣播Barrier,
- Sink接受Barrier的操作流程與Transformation Oper一樣,當所有的Barrier都到達Sink之后,并且所有的Sink也完成了Checkpoint,這一輪Snapshot就完成了,
下面這個圖展示了一個ABS演算法的執行程序:

Exactly-Once vs At-Least-Once
上面講到的屏障對齊程序是Flink exactly-once語意的基礎,因為屏障對齊能夠保證多輸入流的算子正常處理不同checkpoint區間的資料,避免它們發生交叉,即不會有資料被處理兩次,
但是對齊程序需要時間,有一些對延遲特別敏感的應用可能對準確性的要求沒有那么高,所以Flink也允許在StreamExecutionEnvironment.enableCheckpointing()方法里指定At-Least-Once語意,會取消屏障對齊,即算子收到第一個輸入的屏障之后不會阻塞,而是觸發快照,這樣一來,部分屬于檢查點n + 1的資料也會包括進檢查點n的資料里, 當恢復時,這部分交叉的資料就會被重復處理,
7、Watermark
Flink 程式并 不能自動提取資料源中哪個欄位/標識為資料的事件時間,從而也就無法自己定義 Watermark ,
開發人員需要通過 Flink 提供的 API 來 提取和定義 Timestamp/Watermark,可以在 資料源或者資料流中 定義,
1)自定義資料源設定 Timestamp/Watermark
自定義的資料源類需要繼承并實作 SourceFunction[T] 介面,其中 run 方法是定義資料生產的地方:
//自定義的資料源為自定義型別MyType class MySource extends SourceFunction[MyType]{ //重寫run方法,定義資料生產的邏輯 override def run(ctx: SourceContext[MyType]): Unit = { while (/* condition */) { val next: MyType = getNext() //設定timestamp從MyType的哪個欄位獲取(eventTimestamp) ctx.collectWithTimestamp(next, next.eventTimestamp) if (next.hasWatermarkTime) { //設定watermark從MyType的那個方法獲取(getWatermarkTime) ctx.emitWatermark(new Watermark(next.getWatermarkTime)) } } } }
2)在資料流中設定 Timestamp/Watermark
在資料流中,可以設定 stream 的 Timestamp Assigner ,該 Assigner 將會接收一個 stream,并生產一個帶 Timestamp和Watermark 的新 stream,
val withTimestampsAndWatermarks: DataStream[MyEvent] = stream .assignTimestampsAndWatermarks(new MyTimestampsAndWatermarks())
8、廣播狀態(Broadcast State)
和 Spark 中的廣播變數一樣,Flink 也支持在各個節點中各存一份小資料集,所在的計算節點實體可在本地記憶體中直接讀取被廣播的資料,可以避免Shuffle提高并行效率,
廣播狀態(Broadcast State)的引入是為了支持一些來自一個流的資料需要廣播到所有下游任務的情況,它存盤在本地,用于處理其他流上的所有傳入元素,
// key the shapes by color KeyedStream<Item, Color> colorPartitionedStream = shapeStream.keyBy(new KeySelector<Shape, Color>(){...}); // a map descriptor to store the name of the rule (string) and the rule itself. MapStateDescriptor<String, Rule> ruleStateDescriptor = new MapStateDescriptor<>("RulesBroadcastState",BasicTypeInfo.STRING_TYPE_INFO, TypeInformation.of(new TypeHint<Rule>() {})); // broadcast the rules and create the broadcast state BroadcastStream<Rule> ruleBroadcastStream = ruleStream.broadcast(ruleStateDescriptor); DataStream<Match> output = colorPartitionedStream.connect(ruleBroadcastStream).process(new KeyedBroadcastProcessFunction<Color, Item, Rule, String>(){...});
9、Operator Chain
Flink作業中,可以指定相關的chain將相關性非常強的轉換操作(operator)系結在一起,使得上下游的Task在同一個Pipeline中執行,避免因為資料在網路或者執行緒之間傳輸導致的開銷,
一般情況下Flink在Map型別的操作中默認開啟 Operator Chain 以提高整體性能,開發人員也可以根據需要創建或者禁止 Operator Chain 對任務進行細粒度的鏈條控制,
//創建 chain dataStream.filter(...).map(...).startNewChain().map(...) //禁止 chain dataStream.map(...).disableChaining()
創建的鏈條只對當前的運算子和之后的運算子有效,不不影響其他操作,如上代碼只針對兩個map操作進行鏈條系結,對前面的filter操作無效,如果需要可以在filter和map之間使用 startNewChain方法即可,
10、Side Output
除了從DataStream操作的結果中獲取主資料流之外,Flink還可以產生任意數量額外的側輸出(Side Output)結果流,側輸出結果流的資料型別不需要與主資料流的型別一致,不同側輸出流的型別也可以不同,當要拆分資料流時(通常必須復制流),從每個流過濾出不想擁有的資料時Side Output將非常有用,
DataStream<Integer> input = ...; final OutputTag<String> outputTag = new OutputTag<String>("side-output"){}; SingleOutputStreamOperator<Integer> mainDataStream = input .process(new ProcessFunction<Integer, Integer>() { @Override public void processElement( Integer value, Context ctx, Collector<Integer> out) throws Exception { // 將資料發送到常規輸出中 out.collect(value); // 將資料發送到側輸出中 ctx.output(outputTag, "sideout-" + String.valueOf(value)); } });
DataStream<String> sideOutputStream = mainDataStream.getSideOutput(outputTag);
參考:
https://zhuanlan.zhihu.com/p/93507000
https://www.jianshu.com/p/3093f6d92750
https://www.jianshu.com/p/8d6569361999
https://ci.apache.org/projects/flink/flink-docs-release-1.9/
轉載請註明出處,本文鏈接:https://www.uj5u.com/shujuku/40006.html
標籤:大數據
