在流上做join是比较常见的情形,Flink的双流Join是需要基于两个前提条件:
两个流的时间类型一致;
被关联的流数据都是在一个窗口内的.
之前构造了一个带watermark特性的双流Join的场景,但是发现结果总是与自己的预期不一样.
比如,通过Kafka构造两个流stream1和stream2作为输入,Flink的时间特性定义为EventTime,以3s作为一个滚动窗口计算.如果只是单纯做流上的Join,那么只需满足上面提到的两点即可.但是实际场景中,流数据都会有迟到的情形,所以我在两个流上都加上了watermark特性,有了watermark后窗口触发的条件就不一样了.




从kafak生产了几条数据,格式如下:


Flink这端应用起来后,输出结果如下:

按照3s一个滚动窗口的计算预期,watermark结合窗口的特性,watermark_time>=window_endtime,才会触发窗口的计算,看起来应该是1000000057000这条数据到达后,开始计算.但是结果却在1000000056000便触发了计算.
最后我把这个疑问抛给了Flink社区,从社区得到了回复,触发的计算条件都没问题,问题是理解窗口划分的定义.Flink的窗口界定源码如下:
public static long getWindowStartWithOffset(long timestamp, long offset, long windowSize) {
return timestamp - (timestamp - offset + windowSize) % windowSize;
};
也就是说,窗口的开始与闭合是遵循上面的计算公式的,比如第一条数据时间戳是1000000055000,那么
WindowStart_Time=1000000055000-(1000000055000-0+3000)%3000=1000000053000
WindowEnd_Time=1000000053000+3000
所以,窗口开端并不是传入首条记录作为开端,这个例子的窗口范围是[1000000053000,1000000056000).这样,才会出现在1000000056000到达的时候便触发了计算.




