理解Flink事件时间Over Window中的迟到数据及运行疑问
问题与解答: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(排序型窗口)**的处理逻辑,结合你使用的有限数据集场景,导致了有序输出:
- 状态存储逻辑:对于
ORDER BY rt的Over Window,Flink会在状态中维护一个按Event Time排序的数据结构(比如有序映射),所有到达的元素都会根据rt插入到对应位置,保证状态内的数据始终有序。 - 有限数据集的特殊性:你的
StockSource是一次性输出所有数据的有限流,Flink在处理这类场景时,会将所有数据加载到状态中完成排序,最终按顺序输出聚合结果——本质上等同于批处理模式下的排序计算。 - 无限流的差异:如果是持续产生数据的无限流,当Watermark推进超过某个时间点后,Flink会清理状态中过期的窗口数据,此时后续到达的迟到数据(在容忍范围内)会被单独处理并追加输出,无法保证全局有序。
问题2:迟到数据未被过滤的原因
你遇到的情况是两个关键因素共同作用的结果:
- 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默认不会自动丢弃这类数据。 - 状态TTL的缺失:你没有配置状态生存时间(TTL),Flink会一直保留状态直到作业结束。在有限数据集场景下,所有数据(包括迟到数据)都会被插入到有序状态中,参与窗口聚合计算。只有当配置了
table.exec.state.ttl且Watermark推进超过数据rt+ TTL时,迟到数据才会被状态清理机制过滤。
另外需要注意:你使用的toAppendStream会输出所有参与计算的元素结果,即使是迟到数据,只要被状态处理就会被追加输出。
内容的提问来源于stack exchange,提问作者Tom
相关产品推荐
相关产品推荐

