news 2026/10/8 13:20:59

Flink Sliding Window 详解及代码实现

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
Flink Sliding Window 详解及代码实现

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 触发条件

一个滑动窗口在满足以下条件时触发计算:

  1. 当前 Watermark 超过窗口的end_time。
  2. 窗口内至少有一条数据(默认情况下,空窗口不触发)。

3.3 窗口生命周期

以窗口大小 60s、滑动步长 10s为例,Flink 内部会维护 6 个活跃窗口。每当 Watermark 前进 10 秒,最老的窗口关闭并触发计算,同时新建一个窗口。

是

否

数据流进入

WindowAssigner 分配窗口

数据同时加入多个重叠窗口

Watermark 推进

是否超过窗口 end_time

触发窗口计算

输出聚合结果

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 运行与测试

  1. 启动一个 Socket 数据源:
nc-lk9999
  1. 输入测试数据:
user1,100 user1,200 user2,50 user1,150
  1. 观察输出,你会看到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,150

6. 进阶:使用 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,并在实际项目中灵活运用。

版权声明: 本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若内容造成侵权/违法违规/事实不符,请联系邮箱:809451989@qq.com进行投诉反馈,一经查实,立即删除!
网站建设 2026/10/8 13:20:47

PA Agent 演示模式使用教程:零API成本回放历史K线分析记录

PA Agent 演示模式使用教程&#xff1a;零API成本回放历史K线分析记录 【免费下载链接】PA_Agent 项目地址: https://gitcode.com/gh_mirrors/pa/PA_Agent PA Agent 是一款基于价格行为学&#xff08;Price Action&#xff09;的 AI K 线分析工具&#xff0c;而它的演示…

作者头像 李华
网站建设 2026/10/8 13:19:55

国产CAD二选一:从天正兼容到信创适配逐项看清

"这两款都是国产CAD&#xff0c;价格差不太多&#xff0c;到底选哪一个&#xff1f;"在帮助企业推进软件国产化的这些年里&#xff0c;这句话几乎是被问得最多的。决策者手里往往摆着几份产品介绍&#xff0c;功能清单看起来都挺齐全&#xff0c;可真要落笔签字&…

作者头像 李华
网站建设 2026/10/8 13:19:50

2026硬核实测|9款主流AI论文工具排名(算力+风控双重测评)

2026年高校AI检测全面升级&#xff0c;只看写作能力、不看风控合规的工具测评全部过时。今年大量毕业生翻车核心原因&#xff1a;用新版大模型写论文&#xff0c;逻辑更丝滑、语句更规整&#xff0c;但机器特征直接拉满&#xff0c;AI检测百分百爆红。本次排名摒弃虚夸宣传&…

作者头像 李华