暂无图片
暂无图片
暂无图片
暂无图片
暂无图片

Flink的双流Join

Nathan的笔记 2020-04-04
840

在流上做join是比较常见的情形,Flink的双流Join是需要基于两个前提条件:

  1. 两个流的时间类型一致;

  2. 被关联的流数据都是在一个窗口内的.

之前构造了一个带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到达的时候便触发了计算.


文章转载自Nathan的笔记,如果涉嫌侵权,请发送邮件至:contact@modb.pro进行举报,并提供相关证据,一经查实,墨天轮将立刻删除相关内容。

评论