EventTimeTrigger
EventTimeTrigger 的觸發完全依賴 watermark,換言之,如果 stream 中沒有 watermark,就不會觸發 EventTimeTrigger,
watermark 之于事件時間就是如此重要,來看一下 watermark 的定義先~
Watermarks 是某個 event time 視窗中所有資料都到齊的標志,
Watermarks 作為資料流的一部分流動并攜帶時間戳 t,Watermark(t) 斷言資料流中不會再有小于時間戳 t 的事件出現,
換言之,當 Watermark(t) 到達時,標志著所有小于時間戳 t 的事件都已到齊,可放心地對時間戳 t 之前的所有事件執行聚合、視窗關閉等動作,
來看一下 《Streaming Systems》書中關于 Watermark 原汁原味的定義:
Watermarks are temporal notions of input completeness in the event-time domain. Worded differently, they are the way the system measures progress and completeness relative to the event times of the records being processed in a stream of events.
That point in event time, E, is the point up to which the system believes all inputs with event times less than E have been observed. In other words, it’s an assertion that no more data with event times less than E will ever be seen again.
如下圖所示,watermark 和 event 一樣在 pipeline 中流動,并且都攜帶時間戳,W(4) 表示在此之后不會再收到事件時間小于 4 的事件,W(9) 表示在此之后不會再接收到事件時間小于 9 的事件,
假設我們準備對這個資料流做視窗聚合操作,時間視窗大小為 4 個時間單位,視窗內元素做求和聚合操作,示例代碼如下:
input.window(TumblingEventTimeWindows.of(Time.minutes(4L)))
.trigger(EventTimeTrigger.create())
.reduce(new SumReduceFunction());
我們跟蹤 Flink 原始碼看看 EventTimeTrigger 的實作邏輯:
/**
* A Trigger that fires once the watermark passes the end of the window to which a pane belongs.
*/
public class EventTimeTrigger extends Trigger<Object, TimeWindow> {
private static final long serialVersionUID = 1L;
private EventTimeTrigger() {}
@Override
public TriggerResult onElement(
Object element, long timestamp, TimeWindow window, TriggerContext ctx)
throws Exception {
if (window.maxTimestamp() <= ctx.getCurrentWatermark()) {
// if the watermark is already past the window fire immediately
return TriggerResult.FIRE;
} else {
ctx.registerEventTimeTimer(window.maxTimestamp());
return TriggerResult.CONTINUE;
}
}
@Override
public TriggerResult onEventTime(long time, TimeWindow window, TriggerContext ctx) {
return time == window.maxTimestamp() ? TriggerResult.FIRE : TriggerResult.CONTINUE;
}
public static EventTimeTrigger create() {
return new EventTimeTrigger();
}
}
onElement 方法在每次資料進入該 window 時都會觸發:首先判斷當前的 watermark 是否已經超過了 window 的最大時間(視窗右界時間),如果已經超過,回傳觸發結果 FIRE;如果尚未超過,則根據視窗右界時間注冊一個事件時間定時器,標記觸發結果為 CONTINUE,
繼續跟蹤 registerEventTimeTimer 的動作,發現原來是將定時器放到了一個優先佇列 eventTimeTimersQueue 中,優先佇列的權重為定時器的時間戳,時間越小越希望先被觸發,自然排在佇列的前面,
onEventTime 方法是在什么時候被呼叫呢,一路跟蹤來到了 InternalTimerServiceImpl 類,發現只有在 advanceWatermark 的時候會觸發 onEventTime,具體邏輯是:當 watermark 到來時,根據 watermark 攜帶的時間戳 t,從事件時間定時器佇列中出隊所有時間戳小于 t 的定時器,然后觸發 onEventTime,onEventTime 方法的具體實作中,首先比較觸發時間是否是自己當前視窗的結束時間,是則 FIRE,否則繼續 CONTINUE,
public class InternalTimerServiceImpl<K, N> implements InternalTimerService<N> {
/** Event time timers that are currently in-flight. */
private final KeyGroupedInternalPriorityQueue<TimerHeapInternalTimer<K, N>>
eventTimeTimersQueue;
@Override
public void registerEventTimeTimer(N namespace, long time) {
eventTimeTimersQueue.add(
new TimerHeapInternalTimer<>(time, (K) keyContext.getCurrentKey(), namespace));
}
public void advanceWatermark(long time) throws Exception {
currentWatermark = time;
InternalTimer<K, N> timer;
while ((timer = eventTimeTimersQueue.peek()) != null && timer.getTimestamp() <= time) {
eventTimeTimersQueue.poll();
keyContext.setCurrentKey(timer.getKey());
triggerTarget.onEventTime(timer);
}
}
}
我們以上述示例資料為例,看一下具體的執行程序:
- 當事件到達時,根據事件時間將事件分配到相應的視窗中,事件 '2' 被分配到時間視窗 T1-T4 中,如下圖所示:
- 后續時間陸續到達,事件 '2','3','1','3' 都陸續被分配到時間視窗 T1-T4 中,由于 currentWatermark 未超過視窗結束時間 4 ,因此注冊事件時間定時器 Timer(4);事件 '7' 被分配到視窗 T5-T8 中,也因為 currentWatermark 未超過視窗結束時間 8,因此注冊事件時間定時器 Timer(8),目前為止事件時間定時器佇列中有 2 個定時器,
- 重點來了,接下來到達的是 watermark(4),那么就去事件時間定時器佇列中找到所有定時時間小于等于 4 的定時器,Timer(4) 出隊,觸發 onEventTime 方法,視窗 T1-T4 的結束時間和 Timer(4) 的時間戳相等,FIRE 視窗 T1-T4(圖中標記為綠色表示 FIRE),此時定時器佇列中只剩下 1 個定時器,
- 重復同樣程序,事件 '5','9','6' 接踵而至,'5', '6' 被分配到視窗 T5-T8 中,'9' 則被分配到新視窗 T9-T12 中,'5','6' 對應的注冊定時器為 Timer(8),'9' 注冊了一個新定時器 Timer(12),此時定時器佇列中剩下 2 個定時器,
watermark(9) 的到達會促使定時器 Timer(8) 出隊,進而 FIRE 視窗 T5-T8,以此類推,
ProcessingTimeTrigger
與 EventTimeTrigger 相比,ProcessingTimeTrigger 就相對簡單了,它的觸發只依賴系統時間,我們跟蹤 Flink 原始碼看看 ProcessingTimeTrigger 的實作邏輯:
/**
* A Trigger that fires once the current system time passes the end of the window to which a pane belongs.
*/
public class ProcessingTimeTrigger extends Trigger<Object, TimeWindow> {
private static final long serialVersionUID = 1L;
private ProcessingTimeTrigger() {}
@Override
public TriggerResult onElement(
Object element, long timestamp, TimeWindow window, TriggerContext ctx) {
ctx.registerProcessingTimeTimer(window.maxTimestamp());
return TriggerResult.CONTINUE;
}
@Override
public TriggerResult onProcessingTime(long time, TimeWindow window, TriggerContext ctx) {
return TriggerResult.FIRE;
}
/** Creates a new trigger that fires once system time passes the end of the window. */
public static ProcessingTimeTrigger create() {
return new ProcessingTimeTrigger();
}
}
onElement 方法在每次資料進入該 window 時都會觸發:直接根據視窗結束時間注冊一個處理時間定時器,同樣是放到一個處理時間定時器優先佇列中,
系統根據系統時間定時地從處理時間定時器佇列中取出小于當前系統時間的定時器,然后呼叫 onProcessingTime 方法,FIRE 視窗,
public class InternalTimerServiceImpl<K, N> implements InternalTimerService<N> {
/** Processing time timers that are currently in-flight. */
private final KeyGroupedInternalPriorityQueue<TimerHeapInternalTimer<K, N>>
@Override
public void registerProcessingTimeTimer(N namespace, long time) {
InternalTimer<K, N> oldHead = processingTimeTimersQueue.peek();
if (processingTimeTimersQueue.add(
new TimerHeapInternalTimer<>(time, (K) keyContext.getCurrentKey(), namespace))) {
long nextTriggerTime = oldHead != null ? oldHead.getTimestamp() : Long.MAX_VALUE;
// check if we need to re-schedule our timer to earlier
if (time < nextTriggerTime) {
if (nextTimer != null) {
nextTimer.cancel(false);
}
nextTimer = processingTimeService.registerTimer(time, this::onProcessingTime);
}
}
}
private void onProcessingTime(long time) throws Exception {
// null out the timer in case the Triggerable calls registerProcessingTimeTimer()
// inside the callback.
nextTimer = null;
InternalTimer<K, N> timer;
while ((timer = processingTimeTimersQueue.peek()) != null && timer.getTimestamp() <= time) {
processingTimeTimersQueue.poll();
keyContext.setCurrentKey(timer.getKey());
triggerTarget.onProcessingTime(timer);
}
if (timer != null && nextTimer == null) {
nextTimer =
processingTimeService.registerTimer(
timer.getTimestamp(), this::onProcessingTime);
}
}
}
歡迎留言、討論,轉載請注明出處,
轉載請註明出處,本文鏈接:https://www.uj5u.com/shujuku/469863.html
標籤:大數據
