1. 引言
Apache Flink 作为一款分布式流处理引擎,窗口(Window)是其最核心的抽象之一。在实际业务中,我们经常需要统计一段时间内的数据,例如「最近 5 分钟的订单金额」「最近 1 小时的 UV」等。Flink 提供了多种窗口类型,其中Sliding Window(滑动窗口)是最常用、也最灵活的一种。
本文将深入讲解 Flink Sliding Window 的原理、触发机制、与 Tumbling Window 的区别,并通过完整的代码实战(Java + Flink 1.17)带你从零跑通一个滑动窗口统计任务。
2. 什么是 Sliding Window
2.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 秒,最老的窗口关闭并触发计算,同时新建一个窗口。
4. 代码实战:基于 Processing Time 的 Sliding Window
4.1 项目依赖
首先在pom.xml中引入 Flink 依赖:
<dependencies><dependency><groupId>org.apache.flink</groupId><artifactId>flink-streaming-java</artifactId><version>1.17.2</version></dependency><dependency><groupId>org.apache.flink</groupId><artifactId>flink-clients</artifactId><version>1.17.2</version></dependency></dependencies>4.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. 创建执行环境StreamExecutionEnvironmentenv=StreamExecutionEnvironment.getExecutionEnvironment();// 2. 模拟数据源:每 1 秒发送一条订单数据DataStream<String>source=env.socketTextStream("localhost",9999);// 3. 解析数据:格式为 "用户ID,订单金额"DataStream<Tuple2<String,Double>>orders=source.map(newMapFunction<String,Tuple2<String,Double>>(){@OverridepublicTuple2<String,Double>map(Stringline)throwsException{String[]fields=line.split(",");StringuserId=fields[0];doubleamount=Double.parseDouble(fields[1]);returnTuple2.of(userId,amount);}});// 4. 按用户 ID 分组// 5. 应用滑动窗口:窗口大小 60 秒,滑动步长 10 秒DataStream<Tuple2<String,Double>>result=orders.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 Window
5.1 为什么需要 Event Time
Processing 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{StreamExecutionEnvironmentenv=StreamExecutionEnvironment.getExecutionEnvironment();// 数据格式:时间戳,用户ID,订单金额DataStream<String>source=env.socketTextStream("localhost",9999);DataStream<Tuple3<Long,String,Double>>orders=source.map(newMapFunction<String,Tuple3<Long,String,Double>>(){@OverridepublicTuple3<Long,String,Double>map(Stringline)throwsException{String[]fields=line.split(",");longtimestamp=Long.parseLong(fields[0]);StringuserId=fields[1];doubleamount=Double.parseDouble(fields[2]);returnTuple3.of(timestamp,userId,amount);}});// 提取事件时间,并设置 Watermark(允许 5 秒乱序)DataStream<Tuple3<Long,String,Double>>withWatermark=orders.assignTimestampsAndWatermarks(WatermarkStrategy.<Tuple3<Long,String,Double>>forBoundedOutOfOrderness(Duration.ofSeconds(5)).withTimestampAssigner(newSerializableTimestampAssigner<Tuple3<Long,String,Double>>(){@OverridepubliclongextractTimestamp(Tuple3<Long,String,Double>element,longrecordTimestamp){returnelement.f0;}}));// 滑动窗口:窗口大小 60 秒,滑动步长 10 秒DataStream<Tuple3<Long,String,Double>>result=withWatermark.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 获取窗口上下文
有时我们不仅需要聚合结果,还需要知道窗口的起止时间。此时可以使用ProcessWindowFunction:
importorg.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;publicclassWindowWithTimeFunctionextendsProcessWindowFunction<Tuple2<String,Double>,String,String,TimeWindow>{@Overridepublicvoidprocess(Stringkey,Contextcontext,Iterable<Tuple2<String,Double>>elements,Collector<String>out)throwsException{doublesum=0.0;for(Tuple2<String,Double>element:elements){sum+=element.f1;}longwindowStart=context.window().getStart();longwindowEnd=context.window().getEnd();out.collect("窗口 ["+windowStart+", "+windowEnd+") 用户 "+key+" 订单总额:"+sum);}}使用方式:
DataStream<String>result=orders.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进行增量聚合,避免全量遍历窗口数据:
DataStream<Tuple2<String,Double>>result=orders.keyBy(order->order.f0).window(SlidingProcessingTimeWindows.of(Time.seconds(60),Time.seconds(10))).reduce(newReduceFunction<Tuple2<String,Double>>(){@OverridepublicTuple2<String,Double>reduce(Tuple2<String,Double>v1,Tuple2<String,Double>v2)throwsException{returnTuple2.of(v1.f0,v1.f1+v2.f1);}});8. 总结
本文详细介绍了 Flink Sliding Window 的核心概念、触发机制,并给出了 Processing Time 与 Event Time 两种时间语义下的完整代码实战。滑动窗口适合「最近 N 时间内的统计」类业务场景,但要注意窗口重叠带来的重复计算问题,以及窗口数量对内存的影响。
希望本文能帮助你彻底掌握 Flink Sliding Window,并在实际项目中灵活运用。