基于Event Time和Watermark的Flink窗口输出结果存疑求助
Flink Event Time窗口与Watermark问题分析
问题场景
编写了基于Event Time和Watermark的Flink滚动窗口代码,窗口长度8秒,允许4秒乱序时间,对输出结果中的某个事件存在疑问。
测试代码
import java.text.SimpleDateFormat import java.util.Date import java.util.concurrent.TimeUnit import org.apache.flink.streaming.api.TimeCharacteristic import org.apache.flink.streaming.api.functions.timestamps.BoundedOutOfOrdernessTimestampExtractor import org.apache.flink.streaming.api.scala.function.WindowFunction import org.apache.flink.streaming.api.scala.{StreamExecutionEnvironment, _} import org.apache.flink.streaming.api.windowing.assigners.TumblingEventTimeWindows import org.apache.flink.streaming.api.windowing.time.Time import org.apache.flink.streaming.api.windowing.windows.TimeWindow import org.apache.flink.util.Collector object EventTimeWindowTest{ def to_milli(str: String) = new SimpleDateFormat("yyyy-MM-dd HH:mm:ss").parse(str).getTime def to_char(milli: Long) = { val date = if (milli <= 0) new Date(0) else new Date(milli) new SimpleDateFormat("yyyy-MM-dd HH:mm:ss").format(date) } def main(args: Array[String]): Unit = { val env = StreamExecutionEnvironment.getExecutionEnvironment env.setStreamTimeCharacteristic(TimeCharacteristic.EventTime) env.setParallelism(1) val data = Seq( ("A", "2020-08-30 10:50:15"), ("A", "2020-08-30 10:50:11"), ("B", "2020-08-30 10:50:14"), ("B", "2020-08-30 10:50:09"), ("A", "2020-08-30 10:50:21"), ("A", "2020-08-30 10:50:10") ) val stream = env.fromCollection(data).setParallelism(1).assignTimestampsAndWatermarks(new BoundedOutOfOrdernessTimestampExtractor[(String, String)](Time.seconds(4)) { override def extractTimestamp(evt: (String, String)): Long = to_milli(evt._2) }).keyBy(_._1) .window(TumblingEventTimeWindows.of(Time.of(8, TimeUnit.SECONDS))) .apply(new WindowFunction[(String, String), String, String, TimeWindow] { override def apply(key: String, window: TimeWindow, windowData: Iterable[(String, String)], out: Collector[String]): Unit = { val start = to_char(window.getStart) val end = to_char(window.getEnd) val sb = new StringBuilder sb.append(s"$key($start, $end):") windowData.foreach { e => sb.append(e._2 + ",") } out.collect(sb.toString().substring(0, sb.length - 1)) } }) stream.print().setParallelism(1) env.execute() } }
代码输出
A(2020-08-30 10:50:08, 2020-08-30 10:50:16):2020-08-30 10:50:15,2020-08-30 10:50:11,2020-08-30 10:50:10 B(2020-08-30 10:50:08, 2020-08-30 10:50:16):2020-08-30 10:50:14,2020-08-30 10:50:09 A(2020-08-30 10:50:16, 2020-08-30 10:50:24):2020-08-30 10:50:21
用户疑问
无法理解为什么Key为A、事件时间为2020-08-30 10:50:10的事件会被输出。按照理解,此前的事件("A", "2020-08-30 10:50:21")会将Watermark推进到2020-08-30 10:50:17,此时窗口(2020-08-30 10:50:08, 2020-08-30 10:50:16)应该已经关闭,该事件不应被包含在内。
核心原因解析
1. 事件处理与Watermark生成的顺序
Flink处理单个事件的固定流程是:
- 先分配事件到窗口:提取事件时间戳,将事件添加到对应时间窗口
- 再更新Watermark:更新当前观察到的最大事件时间,生成并传递新的Watermark
2. 有限数据流的特殊触发逻辑
你使用fromCollection创建的是有限数据流,Flink会先按顺序处理完所有事件,再发送一个终极Watermark(Long.MAX_VALUE),统一触发所有未计算的窗口。
3. 你的事件处理全流程拆解
按事件顺序逐一分析:
- 前5个事件处理完毕后,Watermark被推进到
10:50:17(10:50:21 - 4秒),此时窗口[10:50:08, 10:50:16)的结束时间已小于Watermark,但窗口并未立即触发计算——因为还有最后一个事件未处理。 - 处理第6个事件
("A", "2020-08-30 10:50:10")时,先将其分配到窗口[10:50:08, 10:50:16),之后更新Watermark(因该事件时间戳小于当前最大事件时间,Watermark保持10:50:17)。 - 所有事件处理完成后,终极Watermark触发窗口计算,此时窗口
[10:50:08, 10:50:16)已经包含了第6个事件,所以输出中会出现它。
4. 真实无限流场景的差异
如果是真实的无限流环境,若第6个事件在Watermark推进到10:50:17之后才到达,该事件会被默认丢弃(因为未设置窗口的allowedLateness),不会被包含进窗口。
内容的提问来源于stack exchange,提问作者Tom
相关产品推荐
相关产品推荐

