flink中window实现的源码分析

关于Flink中window的实现分析

window

window提供了一种处理无界数据的一种手段

window的组件

首先我们看window包含了哪些组件:触发器trigger,触发器上下文triggerContext,内部状态windowState,窗口分配器windowassigner, 内部时间服务器internalTimerService,初看到这么多的组件可能会有点懵,下面的分析会一点一点介绍这些组件的作用。

今天我们从flink接收流元素进行处理的角度来分析其实现,在flink的DAG中流动的有这么几种元素:StreamRecord,LatencyRecord,WaterMark,StreamStatus,我们这里只考虑这样两种元素StreamRecordWaterMark

streamRecord

windowOperator实现了KeyContext,其实就是代表每一个windowOperator处理每一个元素就会在一个Key的上下文的环境中去做处理。

当windowOperator接收到一条StreamRecord,windowOperator会做什么呢?

1
2
3
4
5
synchronized (lock) {
numRecordsIn.inc();
streamOperator.setKeyContextElement1(record);
streamOperator.processElement(record);
}
  1. 首先设置该operator的key为当前元素
  2. 根据element所携带的时间戳(processing time或者event time)分配元素所应该属于的窗口,一个元素可能会隶属于多个窗口,比如slideWindowAssigner。
  3. 如果这个窗口是一个可merge窗口,例如session窗口,那么就会进行和原有窗口的合并和状态的更新

窗口merge原理

  1. 首先取出之前的所有未清除的窗口,和新分配到的窗口做一次merge,有重叠部分则新生成大的窗口
  2. 将各个小窗口的真实数据merge到合并后的大窗口
  3. 注册大窗口的清理时间触发器,清理原先子窗口的清理时间触发器

窗口merge完之后,则会通过triggerContext#OnElement方法去进行判断是否能触发计算,触发方式如前文分析wartermark的那样,就是通过wartermark来判断这个窗口是否可以进行计算

1
2
3
4
5
6
7
8
9
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;
}
}

如果不能触发计算,我们看到其实他是将window.maxTimestamp注册到了eventTimeQueue中, 这里我们先记一下,后文会提到他的作用。

小结

到这里windowOperator接收到一个StreamRecord元素的处理逻辑已经结束了,如果窗口不是可merge类型的除了不做窗口merge,其他的操作也是大同小异,到这里也许你会有几个疑问(其实是我自己看代码的一些疑问:)

  1. 如果我设置了allowLateness会对我们的计算结果产生什么样的影响呢?
  2. 我是用apply(),process()函数或者reduce()这种聚合函数对于window的cost代价有多大的差别?

问题一

上文中提到的在处理元素的最后会注册一个窗口cleanupTimer,那么这个时间是多少呢? window.maxTimestamp() + allowedLateness; 所以我们看到一个窗口存在时长是水位线经过window的最大时间+allowLateness的时间,因此当水位线大于窗口最大时间后就会触发计算,而计算之后状态不会清空,会保留allowedLateness的时长,而此时窗口状态还在保留,所以上游有晚到的数据来一条就会触发一次该窗口的计算,而且每次计算的数据都是该窗口的全量数据,所以业务方要慎用,或者下游要做好相应的去重或更新措施,否则可能会造成结果的不准确

问题二

不同的函数最终影响的其实是我们最终window保存数据的state的形式有ListState,也有reducingState…,最终影响了rocksdb和checkpoint的大小,不过肯定是能用聚合还是用聚合函数比较好

waterMark

当元素来的是一个watermark,window Operator又会以怎样的逻辑去处理呢?

1
2
3
4
5
6
7
8
9
public void handleWatermark(Watermark watermark) {
try {
synchronized (lock) {
lastEmittedWatermark = watermark.getTimestamp();
operator.processWatermark(watermark);
}
} catch (Exception e) {
throw new RuntimeException("Exception occurred while processing valve output watermark: ", e); }
}

这里我们将关注到开篇提到的internalTimeService,

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
public void advanceWatermark(long time) throws Exception {
currentWatermark = time;

InternalTimer<K, N> timer;

while ((timer = eventTimeTimersQueue.peek()) != null && timer.getTimestamp() <= time) {

Set<InternalTimer<K, N>> timerSet = getEventTimeTimerSetForTimer(timer);
timerSet.remove(timer);
eventTimeTimersQueue.remove();

keyContext.setCurrentKey(timer.getKey());
triggerTarget.onEventTime(timer);
}
}

watermark元素处理逻辑

  1. 当流元素是watermark时主要处理逻辑集中在internaltimerservice
  2. 如果eventTimeTimersQueue这个优先级队列中最早的时间低于了水位线,那么就会取出同一时刻的所有key的timer进行计算
  3. 处理可能意味着窗口的计算触发或者某些窗口的清理

以上便是flink对于窗口的实现逻辑。

谢谢支持