You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

基于Event Time和Watermark的Flink窗口输出结果存疑求助

问题场景

编写了基于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. 你的事件处理全流程拆解

按事件顺序逐一分析:

  1. 前5个事件处理完毕后,Watermark被推进到10:50:17(10:50:21 - 4秒),此时窗口[10:50:08, 10:50:16)的结束时间已小于Watermark,但窗口并未立即触发计算——因为还有最后一个事件未处理。
  2. 处理第6个事件("A", "2020-08-30 10:50:10")时,先将其分配到窗口[10:50:08, 10:50:16),之后更新Watermark(因该事件时间戳小于当前最大事件时间,Watermark保持10:50:17)。
  3. 所有事件处理完成后,终极Watermark触发窗口计算,此时窗口[10:50:08, 10:50:16)已经包含了第6个事件,所以输出中会出现它。

4. 真实无限流场景的差异

如果是真实的无限流环境,若第6个事件在Watermark推进到10:50:17之后才到达,该事件会被默认丢弃(因为未设置窗口的allowedLateness),不会被包含进窗口。


内容的提问来源于stack exchange,提问作者Tom

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.31 00:12:23