
1. 引言Apache Flink 作为一款分布式流处理引擎窗口Window是其最核心的抽象之一。在实际业务中我们经常需要统计一段时间内的数据例如「最近 5 分钟的订单金额」「最近 1 小时的 UV」等。Flink 提供了多种窗口类型其中Sliding Window滑动窗口是最常用、也最灵活的一种。本文将深入讲解 Flink Sliding Window 的原理、触发机制、与 Tumbling Window 的区别并通过完整的代码实战Java Flink 1.17带你从零跑通一个滑动窗口统计任务。2. 什么是 Sliding Window2.1 基本概念Sliding Window 由两个参数定义窗口大小Window Size窗口覆盖的时间范围例如 60 秒。滑动步长Slide Size窗口向前移动的间隔例如 10 秒。当窗口大小 滑动步长时相邻窗口之间会有重叠一条数据会同时属于多个窗口。2.2 直观理解假设窗口大小为 60 秒滑动步长为 10 秒窗口 1[00:00, 00:60)窗口 2[00:10, 01:10)窗口 3[00:20, 01:20)可以看到窗口 1 和窗口 2 在[00:10, 00:60)区间重叠落在该区间的数据会被两个窗口同时统计。2.3 与 Tumbling Window 的区别对比项Tumbling Window滚动窗口Sliding Window滑动窗口窗口大小固定固定滑动步长等于窗口大小小于窗口大小窗口重叠无有数据归属每条数据只属于一个窗口一条数据可能属于多个窗口触发频率每个窗口结束触发一次每个滑动步长触发一次3. Sliding Window 的触发机制3.1 时间语义Flink 支持三种时间语义Event Time事件时间数据本身携带的时间戳最真实推荐使用。Ingestion Time摄入时间数据进入 Flink 的时间。Processing Time处理时间算子本地处理数据的系统时间。Sliding Window 的触发与时间语义密切相关。使用 Event Time 时窗口的触发依赖于 Watermark水位线的推进。3.2 触发条件一个滑动窗口在满足以下条件时触发计算当前 Watermark 超过窗口的end_time。窗口内至少有一条数据默认情况下空窗口不触发。3.3 窗口生命周期以窗口大小 60s、滑动步长 10s为例Flink 内部会维护 6 个活跃窗口。每当 Watermark 前进 10 秒最老的窗口关闭并触发计算同时新建一个窗口。是否数据流进入WindowAssigner 分配窗口数据同时加入多个重叠窗口Watermark 推进是否超过窗口 end_time触发窗口计算输出聚合结果4. 代码实战基于 Processing Time 的 Sliding Window4.1 项目依赖首先在pom.xml中引入 Flink 依赖dependenciesdependencygroupIdorg.apache.flink/groupIdartifactIdflink-streaming-java/artifactIdversion1.17.2/version/dependencydependencygroupIdorg.apache.flink/groupIdartifactIdflink-clients/artifactIdversion1.17.2/version/dependency/dependencies4.2 完整代码下面实现一个「每 10 秒统计最近 60 秒内每个用户的订单金额总和」的滑动窗口任务importorg.apache.flink.api.common.functions.MapFunction;importorg.apache.flink.api.java.tuple.Tuple2;importorg.apache.flink.streaming.api.datastream.DataStream;importorg.apache.flink.streaming.api.environment.StreamExecutionEnvironment;importorg.apache.flink.streaming.api.windowing.assigners.SlidingProcessingTimeWindows;importorg.apache.flink.streaming.api.windowing.time.Time;publicclassSlidingWindowDemo{publicstaticvoidmain(String[]args)throwsException{// 1. 创建执行环境StreamExecutionEnvironmentenvStreamExecutionEnvironment.getExecutionEnvironment();// 2. 模拟数据源每 1 秒发送一条订单数据DataStreamStringsourceenv.socketTextStream(localhost,9999);// 3. 解析数据格式为 用户ID,订单金额DataStreamTuple2String,Doubleorderssource.map(newMapFunctionString,Tuple2String,Double(){OverridepublicTuple2String,Doublemap(Stringline)throwsException{String[]fieldsline.split(,);StringuserIdfields[0];doubleamountDouble.parseDouble(fields[1]);returnTuple2.of(userId,amount);}});// 4. 按用户 ID 分组// 5. 应用滑动窗口窗口大小 60 秒滑动步长 10 秒DataStreamTuple2String,Doubleresultorders.keyBy(order-order.f0).window(SlidingProcessingTimeWindows.of(Time.seconds(60),Time.seconds(10))).sum(1);// 6. 输出结果result.print();// 7. 执行任务env.execute(Flink Sliding Window Demo);}}4.3 运行与测试启动一个 Socket 数据源nc-lk9999输入测试数据user1,100 user1,200 user2,50 user1,150观察输出你会看到user1的金额在多个重叠窗口中累加。5. 代码实战基于 Event Time 的 Sliding Window5.1 为什么需要 Event TimeProcessing Time 无法处理乱序数据。在真实业务中数据可能延迟到达此时应使用 Event Time Watermark。5.2 完整代码importorg.apache.flink.api.common.eventtime.SerializableTimestampAssigner;importorg.apache.flink.api.common.eventtime.WatermarkStrategy;importorg.apache.flink.api.common.functions.MapFunction;importorg.apache.flink.api.java.tuple.Tuple3;importorg.apache.flink.streaming.api.datastream.DataStream;importorg.apache.flink.streaming.api.environment.StreamExecutionEnvironment;importorg.apache.flink.streaming.api.windowing.assigners.SlidingEventTimeWindows;importorg.apache.flink.streaming.api.windowing.time.Time;importjava.time.Duration;publicclassEventTimeSlidingWindowDemo{publicstaticvoidmain(String[]args)throwsException{StreamExecutionEnvironmentenvStreamExecutionEnvironment.getExecutionEnvironment();// 数据格式时间戳,用户ID,订单金额DataStreamStringsourceenv.socketTextStream(localhost,9999);DataStreamTuple3Long,String,Doubleorderssource.map(newMapFunctionString,Tuple3Long,String,Double(){OverridepublicTuple3Long,String,Doublemap(Stringline)throwsException{String[]fieldsline.split(,);longtimestampLong.parseLong(fields[0]);StringuserIdfields[1];doubleamountDouble.parseDouble(fields[2]);returnTuple3.of(timestamp,userId,amount);}});// 提取事件时间并设置 Watermark允许 5 秒乱序DataStreamTuple3Long,String,DoublewithWatermarkorders.assignTimestampsAndWatermarks(WatermarkStrategy.Tuple3Long,String,DoubleforBoundedOutOfOrderness(Duration.ofSeconds(5)).withTimestampAssigner(newSerializableTimestampAssignerTuple3Long,String,Double(){OverridepubliclongextractTimestamp(Tuple3Long,String,Doubleelement,longrecordTimestamp){returnelement.f0;}}));// 滑动窗口窗口大小 60 秒滑动步长 10 秒DataStreamTuple3Long,String,DoubleresultwithWatermark.keyBy(order-order.f1).window(SlidingEventTimeWindows.of(Time.seconds(60),Time.seconds(10))).sum(2);result.print();env.execute(Flink Event Time Sliding Window Demo);}}5.3 测试数据示例1700000000000,user1,100 1700000005000,user1,200 1700000010000,user2,50 1700000015000,user1,1506. 进阶使用 ProcessWindowFunction 获取窗口上下文有时我们不仅需要聚合结果还需要知道窗口的起止时间。此时可以使用ProcessWindowFunctionimportorg.apache.flink.api.java.tuple.Tuple2;importorg.apache.flink.streaming.api.functions.windowing.ProcessWindowFunction;importorg.apache.flink.streaming.api.windowing.windows.TimeWindow;importorg.apache.flink.util.Collector;publicclassWindowWithTimeFunctionextendsProcessWindowFunctionTuple2String,Double,String,String,TimeWindow{Overridepublicvoidprocess(Stringkey,Contextcontext,IterableTuple2String,Doubleelements,CollectorStringout)throwsException{doublesum0.0;for(Tuple2String,Doubleelement:elements){sumelement.f1;}longwindowStartcontext.window().getStart();longwindowEndcontext.window().getEnd();out.collect(窗口 [windowStart, windowEnd) 用户 key 订单总额sum);}}使用方式DataStreamStringresultorders.keyBy(order-order.f0).window(SlidingProcessingTimeWindows.of(Time.seconds(60),Time.seconds(10))).process(newWindowWithTimeFunction());7. 常见问题与调优建议7.1 窗口重叠导致的数据重复计算这是 Sliding Window 的固有特性。如果业务上不允许重复统计请改用 Tumbling Window 或使用Session Window。7.2 窗口数量过多当窗口大小 / 滑动步长的比值过大时Flink 需要维护大量活跃窗口内存压力较大。建议合理设置步长或使用增量聚合 全量聚合结合的方式优化。7.3 延迟数据处理使用 Event Time 时可以设置allowedLateness允许迟到数据.window(SlidingEventTimeWindows.of(Time.seconds(60),Time.seconds(10))).allowedLateness(Time.seconds(30))7.4 增量聚合优化对于求和、计数等场景推荐使用reduce或aggregate进行增量聚合避免全量遍历窗口数据DataStreamTuple2String,Doubleresultorders.keyBy(order-order.f0).window(SlidingProcessingTimeWindows.of(Time.seconds(60),Time.seconds(10))).reduce(newReduceFunctionTuple2String,Double(){OverridepublicTuple2String,Doublereduce(Tuple2String,Doublev1,Tuple2String,Doublev2)throwsException{returnTuple2.of(v1.f0,v1.f1v2.f1);}});8. 总结本文详细介绍了 Flink Sliding Window 的核心概念、触发机制并给出了 Processing Time 与 Event Time 两种时间语义下的完整代码实战。滑动窗口适合「最近 N 时间内的统计」类业务场景但要注意窗口重叠带来的重复计算问题以及窗口数量对内存的影响。希望本文能帮助你彻底掌握 Flink Sliding Window并在实际项目中灵活运用。