关于Flink窗口延迟触发配置的理解确认
关于Flink窗口触发逻辑的理解确认
示例1代码
Window.<KV<String, DeviceData>>into(FixedWindows.of(Duration.standardSeconds(options.getWindowSize()))) .triggering( AfterWatermark.pastEndOfWindow() .withLateFirings(AfterProcessingTime.pastFirstElementInPane().plusDelayOf(Duration.standardMinutes(5))) ) .withAllowedLateness(Duration.standardHours(3)) .accumulatingFiredPanes();
示例2代码
Window.<KV<String, DeviceData>>into(FixedWindows.of(Duration.standardSeconds(options.getWindowSize()))) .triggering( AfterWatermark.pastEndOfWindow() .withEarlyFirings(AfterProcessingTime.pastFirstElementInPane().plusDelayOf(Duration.standardMinutes(1))) .withLateFirings(AfterProcessingTime.pastFirstElementInPane().plusDelayOf(Duration.standardHours(1)))) .withAllowedLateness(Duration.standardHours(6)) .accumulatingFiredPanes();
理解纠正与确认
示例1逻辑说明:
- 窗口会在水印超过窗口结束时间时触发一次,这是
AfterWatermark.pastEndOfWindow()的默认行为。 - 窗口结束后的3小时内,若有延迟数据到达,会生成新的pane,并在该pane第一个元素到达后的5分钟(处理时间)触发一次,并非延迟数据一到就触发。
- 最大允许延迟时长确实为3小时,超出该时长的延迟数据会被直接丢弃。
- 窗口会在水印超过窗口结束时间时触发一次,这是
示例2逻辑说明:
- 窗口结束前:每当有新元素进入窗口生成新的pane,会在该pane第一个元素到达后的1分钟(处理时间)触发一次——若窗口内多次有新元素进入,可能触发多次提前计算,并非仅触发一次。
- 窗口结束时:水印超过窗口结束时间时,会触发一次最终计算。
- 窗口结束后:6小时内到达的延迟数据,每生成一个新pane,会在该pane第一个元素到达后的1小时(处理时间)触发一次。
- 最大允许延迟时长为6小时,超时时长的延迟数据会被丢弃。
内容的提问来源于stack exchange,提问作者Alex Tbk
相关产品推荐
相关产品推荐

