首先 flink 的窗口分配是发生在 StreamTask 初始化的过程中。
核心方法是 TumblingProcessingTimeWindows.assignWindows(…)
// TODO : 在初始化StreamTask的时候需要分配好窗口
@Override
public Collection assignWindows(
Object element, long timestamp, WindowAssignerContext context) {
final long now = context.getCurrentProcessingTime();
// TODO : 默认情况下 staggerOffset = 0
if (staggerOffset == null) {
staggerOffset =
windowStagger.getStaggerOffset(context.getCurrentProcessingTime(), size);
}
// TODO : 获取窗口起始时间
long start =
TimeWindow.getWindowStartWithOffset(
now, (globalOffset + staggerOffset) % size, size);
return Collections.singletonList(new TimeWindow(start, start + size));
}
可以通过这个方法往上点,会发现他是在StreamTask初始化的时候触发的。
这个方法调用了一个很重要的方法来计算窗口开始时间:TimeWindow.getWindowStartWithOffset(…)
// TODO : 默认 offset = 0
public static long getWindowStartWithOffset(long timestamp, long offset, long windowSize) {
return timestamp - (timestamp - offset + windowSize) % windowSize;
// TODO : 如果 offset = 0,当前时间 - 当前时间除去windowSize的余数
// TODO : 如果 offset != 0且为正数, 由于 offset 不会大于 windowSize,所以会导致余数变小了,最终得到的窗口 startTime 变大了
}
注释里有我的简单总结,供参考。



