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

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。

解决办法

  1. 添加后续事件:生产环境中后续事件的到来会自然推进最大事件时间,进而推进水印;测试时可构造事件时间晚于09:29:40的事件(此时水印为09:29:40 - 20秒 = 09:29:20,刚好达到窗口结束时间),触发窗口输出。
  2. 调整水印策略:如果测试时不想依赖后续事件,可使用WatermarkStrategy.forMonotonousTimestamps()(适用于事件时间单调递增的场景),此时水印等于最大事件时间,当事件所属窗口的结束时间小于等于水印时,窗口会触发。
  3. 自定义水印生成器:如果需要更灵活的水印控制,可实现WatermarkGenerator接口,手动控制水印推进逻辑。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.25 04:25:01