Flink中TumblingEventTimeWindows无输出问题求助
Flink事件时间滚动窗口无输出,处理时间窗口正常的问题排查
我在Flink中使用TumblingEventTimeWindows时没有任何输出,但切换为TumblingProcessingTimeWindows后运行正常。查了官方文档,没找到必须额外配置的项(比如触发器、驱逐器、允许延迟时间)。
输入数据示例:
{"userId":1,"count":11,"dt":"2023-04-11T09:29:12.244"}
系统时间与输入时间一致。
初始代码
WatermarkStrategy<UserModel> strategy = WatermarkStrategy.<UserModel>forBoundedOutOfOrderness(Duration.ofSeconds(20)) .withTimestampAssigner((i, timestamp) -> Timestamp.valueOf(i.dt).getTime()); ds.assignTimestampsAndWatermarks(strategy) .windowAll(TumblingEventTimeWindows.of(Time.seconds(10))).reduce((acc, i) -> { acc.count += i.count; acc.dt = i.dt; return acc; }).addSink(new PrintSinkFunction());
排查步骤(更新2)
- 在
withTimestampAssigner中添加打印,确认每个事件都会触发该方法; - 添加
OutputTag捕获延迟事件,无延迟数据输出; - 在reduce函数内添加调试打印,确认每个事件都会执行reduce逻辑,但窗口关闭后sink仍无输出。
完整代码:
private static void m4(DataStream<UserModel> ds) { WatermarkStrategy<UserModel> strategy = WatermarkStrategy.<UserModel>forBoundedOutOfOrderness(Duration.ofSeconds(20)) .withTimestampAssigner((i, timestamp) -> { long time = i.dt.toInstant(ZoneOffset.UTC).toEpochMilli(); System.out.println(i.dt + " is: " + time + " dont know: " + timestamp); return time; }); OutputTag<UserModel> lateTag = new OutputTag<UserModel>("late"){}; SingleOutputStreamOperator<UserModel> reduce = ds.assignTimestampsAndWatermarks(strategy) .windowAll(TumblingEventTimeWindows.of(Time.seconds(10))) .sideOutputLateData(lateTag) .reduce((acc, i) -> { System.out.println(i.dt + " reDUCE:"); acc.count += i.count; acc.dt = i.dt; return acc; }); reduce.getSideOutput(lateTag).print(); reduce.addSink(new PrintSinkFunction()); }
进一步尝试(更新3)
尝试添加ProcessAllWindowFunction,但该函数并未被调用,仍无输出;而切换为TumblingProcessingTimeWindows时,无需该函数也能正常关闭窗口并输出到sink。
相关代码:
public class Rich extends ProcessAllWindowFunction<UserModel, UserModelEx, TimeWindow> { @Override public void process(ProcessAllWindowFunction<UserModel, UserModelEx, TimeWindow>.Context context, Iterable<UserModel> iterable, Collector<UserModelEx> collector) throws Exception { UserModel um = iterable.iterator().next(); System.out.println(um.count + " rich:" + um.dt); collector.collect(new UserModelEx() {{ userId = um.userId; count = um.count; wStart = LocalDateTime.ofInstant(Instant.ofEpochMilli(context.window().getStart()), ZoneOffset.UTC); wEnd = LocalDateTime.ofInstant(Instant.ofEpochMilli(context.window().getEnd()), ZoneOffset.UTC); }}); } } private static void m4(DataStream<UserModel> ds) { WatermarkStrategy<UserModel> strategy = WatermarkStrategy.<UserModel>forBoundedOutOfOrderness(Duration.ofSeconds(20)) .withTimestampAssigner((i, timestamp) -> { long time = i.dt.toInstant(ZoneOffset.UTC).toEpochMilli(); System.out.println(i.dt + " assignEvent: " + time + " : " + timestamp); return time; }); SingleOutputStreamOperator<UserModelEx> reduce = ds.assignTimestampsAndWatermarks(strategy) .windowAll(TumblingEventTimeWindows.of(Time.seconds(10))) .reduce((acc, i) -> { acc.count += i.count; acc.dt = i.dt; System.out.println(acc.dt + " reduce:" + acc.count); return acc; }, new Rich()); reduce.print(); //reduce.addSink(new PrintSinkFunction<UserModelEx>()); }
问题原因及解决思路
事件时间窗口的触发完全依赖水印(Watermark)推进到窗口结束时间之后。当前代码中使用forBoundedOutOfOrderness(Duration.ofSeconds(20)),水印计算逻辑为最大事件时间减去20秒。
你的输入事件时间是2023-04-11T09:29:12.244,所属的10秒滚动窗口结束时间是09:29:20。由于只有这一条数据,最大事件时间停留在09:29:12.244,水印则为09:29:12.244 - 20秒 = 09:28:52.244,远未达到窗口结束时间09:29:20,因此窗口永远不会触发输出。
而处理时间窗口基于系统时间触发,到点就会关闭窗口输出,所以能正常运行。同时reduce函数是增量计算,你能看到打印是因为每条事件都会触发增量聚合,但只有窗口触发时才会把最终结果输出到sink。
解决办法
- 添加后续事件:生产环境中后续事件的到来会自然推进最大事件时间,进而推进水印;测试时可构造事件时间晚于
09:29:40的事件(此时水印为09:29:40 - 20秒 = 09:29:20,刚好达到窗口结束时间),触发窗口输出。 - 调整水印策略:如果测试时不想依赖后续事件,可使用
WatermarkStrategy.forMonotonousTimestamps()(适用于事件时间单调递增的场景),此时水印等于最大事件时间,当事件所属窗口的结束时间小于等于水印时,窗口会触发。 - 自定义水印生成器:如果需要更灵活的水印控制,可实现
WatermarkGenerator接口,手动控制水印推进逻辑。
内容的提问来源于stack exchange,提问作者padavan
相关产品推荐
相关产品推荐

