知识卡片
多输入算子水位线取最小值的正确性原理
内容
当一个算子同时接收多个上游输入管道的数据时(比如合并了两个分区),它自身对外发出的水位线必须取所有输入管道中最小的那个水位线,而不是最大值或平均值。原因是水位线代表”这个时间点之前的数据都到齐了”的承诺,只要还有一个上游管道进度慢,就不能假装所有数据都到齐,否则会把那条慢管道后续到达的数据误判为迟到并计算错误。这个机制有一个副作用:一条数据稀疏的分区会拖慢整体水位线推进,导致窗口迟迟不触发、内存中缓存越积越多。Flink为此提供了空闲检测(withIdleness):超过设定时间没有新数据的分区会被标记为空闲,取最小值时不再考虑它,直到它重新变活跃。这体现了一个通用原则:正确性保证必须以最慢的那个源为准,而工程上再单独处理”慢源等于不存在”这种退化情况。
参考来源
《Flink入门与实战》第3章《时间和窗口》