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

理解Flink事件时间Over Window中的迟到数据及运行疑问

背景信息

测试数据

id1,1,2020-09-16T12:50:15
id1,2,2020-09-16T12:50:12
id1,3,2020-09-16T12:50:11
id1,4,2020-09-16T12:50:18
id1,5,2020-09-16T12:50:13
id1,6,2020-09-16T12:50:20
id1,8,2020-09-16T12:50:22
id1,9,2020-09-16T12:50:40

实现代码

import org.apache.flink.streaming.api.functions.AssignerWithPunctuatedWatermarks
import org.apache.flink.streaming.api.scala._
import org.apache.flink.streaming.api.watermark.Watermark
import org.apache.flink.table.api.bridge.scala._
import org.apache.flink.table.api.{AnyWithOperations, FieldExpression}
import org.apache.flink.types.Row
import org.example.model.Stock
import org.example.sources.StockSource

object EventTimeOverWindowSqlTest {
  def main(args: Array[String]): Unit = {
    val env = StreamExecutionEnvironment.getExecutionEnvironment
    env.setParallelism(1)
    val ds: DataStream[Stock] = env.addSource(new StockSource())
    val ds2 = ds.assignTimestampsAndWatermarks(new AssignerWithPunctuatedWatermarks[Stock] {
      var max = Long.MinValue

      override def checkAndGetNextWatermark(t: Stock, l: Long): Watermark = {
        if (t.trade_date.getTime > max) {
          max = t.trade_date.getTime
        }
        new Watermark(max - 4000) //allow 4 seconds late
      }

      override def extractTimestamp(t: Stock, l: Long): Long = t.trade_date.getTime
    })
    val tenv = StreamTableEnvironment.create(env)
    val table = tenv.fromDataStream(ds2, $"id", $"price", $"rt".rowtime())
    tenv.createTemporaryView("sourceTable", table)
    tenv.from("sourceTable").toAppendStream[Row].print()
    val sql =
      """
      select
        id,
        price,
        sum(price) OVER (PARTITION BY id ORDER BY rt rows between 2 preceding and current row) as sum_price
      from sourceTable
      """.stripMargin(' ')

    val table2 = tenv.sqlQuery(sql)

    table2.toAppendStream[Row].print()

    env.execute()

  }
}

运行输出

id1,3,3
id1,2,5
id1,5,10
id1,1,8
id1,4,10
id1,6,11
id1,8,18

结果中的price序列(3 2 5 1 4 6 8 9)完全按create_date排序,等效于执行select price from sourceTable order by create_date的结果。

疑问

  • 流式应用通常无法实现全局排序,但本次运行结果却呈现有序,想了解Flink在有限数据范围内触发排序输出的具体机制;
  • 记录id1,5,2020-09-16T12:50:13相对于已到达的id1,4,2020-09-16T12:50:18晚到5秒,按配置的4秒迟到容忍度,该记录应被判定为迟到却未被过滤,想了解原因。

解答

问题1:有限数据范围内的排序输出机制

Flink针对**Event Time Over Window(排序型窗口)**的处理逻辑,结合你使用的有限数据集场景,导致了有序输出:

  1. 状态存储逻辑:对于ORDER BY rt的Over Window,Flink会在状态中维护一个按Event Time排序的数据结构(比如有序映射),所有到达的元素都会根据rt插入到对应位置,保证状态内的数据始终有序。
  2. 有限数据集的特殊性:你的StockSource是一次性输出所有数据的有限流,Flink在处理这类场景时,会将所有数据加载到状态中完成排序,最终按顺序输出聚合结果——本质上等同于批处理模式下的排序计算。
  3. 无限流的差异:如果是持续产生数据的无限流,当Watermark推进超过某个时间点后,Flink会清理状态中过期的窗口数据,此时后续到达的迟到数据(在容忍范围内)会被单独处理并追加输出,无法保证全局有序。

问题2:迟到数据未被过滤的原因

你遇到的情况是两个关键因素共同作用的结果:

  1. Watermark与迟到数据的判定逻辑:你通过max - 4000生成Watermark,意味着当事件的rt小于当前Watermark时,才会被视为迟到数据。当id1,4(12:50:18)到达后,Watermark更新为12:50:14,而id1,5的rt是12:50:13,确实属于迟到数据,但Flink SQL的Over Window默认不会自动丢弃这类数据。
  2. 状态TTL的缺失:你没有配置状态生存时间(TTL),Flink会一直保留状态直到作业结束。在有限数据集场景下,所有数据(包括迟到数据)都会被插入到有序状态中,参与窗口聚合计算。只有当配置了table.exec.state.ttl且Watermark推进超过数据rt + TTL时,迟到数据才会被状态清理机制过滤。

另外需要注意:你使用的toAppendStream会输出所有参与计算的元素结果,即使是迟到数据,只要被状态处理就会被追加输出。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.30 00:45:04