关于Flink中window的实现分析
window
window提供了一种处理无界数据的一种手段
window的组件
首先我们看window包含了哪些组件:触发器trigger,触发器上下文triggerContext,内部状态windowState,窗口分配器windowassigner, 内部时间服务器internalTimerService,初看到这么多的组件可能会有点懵,下面的分析会一点一点介绍这些组件的作用。
今天我们从flink接收流元素进行处理的角度来分析其实现,在flink的DAG中流动的有这么几种元素:StreamRecord,LatencyRecord,WaterMark,StreamStatus,我们这里只考虑这样两种元素StreamRecord和WaterMark
streamRecord
windowOperator实现了KeyContext,其实就是代表每一个windowOperator处理每一个元素就会在一个Key的上下文的环境中去做处理。
当windowOperator接收到一条StreamRecord,windowOperator会做什么呢?
1 | synchronized (lock) { |
- 首先设置该operator的key为当前元素
- 根据element所携带的时间戳(processing time或者event time)分配元素所应该属于的窗口,一个元素可能会隶属于多个窗口,比如slideWindowAssigner。
- 如果这个窗口是一个可merge窗口,例如session窗口,那么就会进行和原有窗口的合并和状态的更新
窗口merge原理

- 首先取出之前的所有未清除的窗口,和新分配到的窗口做一次merge,有重叠部分则新生成大的窗口
- 将各个小窗口的真实数据merge到合并后的大窗口
- 注册大窗口的清理时间触发器,清理原先子窗口的清理时间触发器
窗口merge完之后,则会通过triggerContext#OnElement方法去进行判断是否能触发计算,触发方式如前文分析wartermark的那样,就是通过wartermark来判断这个窗口是否可以进行计算
1 | public TriggerResult onElement(Object element, long timestamp, TimeWindow window, TriggerContext ctx) throws Exception { |
如果不能触发计算,我们看到其实他是将window.maxTimestamp注册到了eventTimeQueue中, 这里我们先记一下,后文会提到他的作用。
小结
到这里windowOperator接收到一个StreamRecord元素的处理逻辑已经结束了,如果窗口不是可merge类型的除了不做窗口merge,其他的操作也是大同小异,到这里也许你会有几个疑问(其实是我自己看代码的一些疑问:)
- 如果我设置了allowLateness会对我们的计算结果产生什么样的影响呢?
- 我是用
apply(),process()函数或者reduce()这种聚合函数对于window的cost代价有多大的差别?
问题一
上文中提到的在处理元素的最后会注册一个窗口cleanupTimer,那么这个时间是多少呢? window.maxTimestamp() + allowedLateness; 所以我们看到一个窗口存在时长是水位线经过window的最大时间+allowLateness的时间,因此当水位线大于窗口最大时间后就会触发计算,而计算之后状态不会清空,会保留allowedLateness的时长,而此时窗口状态还在保留,所以上游有晚到的数据来一条就会触发一次该窗口的计算,而且每次计算的数据都是该窗口的全量数据,所以业务方要慎用,或者下游要做好相应的去重或更新措施,否则可能会造成结果的不准确
问题二
不同的函数最终影响的其实是我们最终window保存数据的state的形式有ListState,也有reducingState…,最终影响了rocksdb和checkpoint的大小,不过肯定是能用聚合还是用聚合函数比较好
waterMark
当元素来的是一个watermark,window Operator又会以怎样的逻辑去处理呢?
1 | public void handleWatermark(Watermark watermark) { |
这里我们将关注到开篇提到的internalTimeService,
1 | public void advanceWatermark(long time) throws Exception { |
watermark元素处理逻辑

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