☰
Flink 2.3.0 从理论到实践 —— 第 8 章 窗口(Window)
2026/10/11 7:04:23 网站建设 项目流程

Flink 2.3.0 从理论到实践 —— 第 8 章 窗口(Window)

课程定位:窗口是流处理的核心抽象,将无界流切分为有限大小的"批次"以便聚合。本章深入四种窗口类型、时间/计数窗口、三种窗口函数、触发器与清理器、会话窗口动态 Gap、窗口性能优化——这些是实时大屏、报表聚合、CEP 模式匹配的基础。

版本基线:Flink 2.3.0


章节导读

  • 8.1 窗口概述
  • 8.2 四种窗口类型
  • 8.3 时间窗口 vs 计数窗口
  • 8.4 窗口函数
  • 8.5 触发器与清理器
  • 8.6 会话窗口动态 Gap
  • 8.7 窗口性能优化与状态管理
  • 8.8 本章小结与下章预告

8.1 窗口概述

8.1.1 为什么需要窗口

流是无界的,但聚合(如"每分钟订单数"、“每辆车当天累计里程”)需要有限数据集。窗口将无界流切分为有限切片:

无界流: ──数据1──数据2──数据3──数据4──数据5──数据6──► │ ▼ 按 5 分钟切窗口 窗口1: [数据1, 数据2] (10:00-10:05) 窗口2: [数据3, 数据4, 数据5] (10:05-10:10) 窗口3: [数据6] (10:10-10:15)

8.1.2 窗口的两阶段处理

┌──────────────────────────────────────────────────────────────┐ │ 窗口处理两阶段 │ └──────────────────────────────────────────────────────────────┘ 阶段 1: 累积(Window Accumulation) 数据到达 → 按 key 分区 → 分配到对应窗口 → 缓存到状态 阶段 2: 触发与计算(Trigger & Evaluation) 满足触发条件(如 Watermark >= 窗口结束) → 调用窗口函数 → 输出结果

8.1.3 Keyed vs Non-Keyed Windows

维度Keyed WindowNon-Keyed Window
前提keyBy后无keyBy
并行度多个并行(按 key)并行度 = 1(全局)
API.window(WindowAssigner).windowAll(WindowAssigner)
性能高(分布式)低(单点)

生产建议:几乎不用windowAll,优先用 Keyed Window 提升并行度。


8.2 四种窗口类型

8.2.1 类型对比

┌──────────────────────────────────────────────────────────────┐ │ 四种窗口类型 │ └──────────────────────────────────────────────────────────────┘ ① Tumbling Window (滚动窗口) ② Sliding Window (滑动窗口) ┌──┬──┬──┬──┐ ┌──┐┌──┐┌──┐ │W1│W2│W3│W4│ │W1││W2││W3│ └──┴──┴──┴──┘ └─0┘└─0┘└─0┘ 不重叠,等长 重叠,等长,滑动 ③ Session Window (会话窗口) ④ Global Window (全局窗口) ┌──┐ ┌────┐ ┌──┐ ┌──────────────┐ │W1│ │ W2 │ │W3│ │ W1(无上限) │ └──┘ └────┘ └──┘ └──────────────┘ 按活动间隔分割 需自定义触发器
类型特点适用场景
Tumbling固定大小,不重叠每分钟统计、按天聚合
Sliding固定大小,可重叠滑动平均、近 5 分钟指标
Session按 Gap 分割,长度可变用户会话、车辆行程
Global无上限,需自定义触发自定义触发逻辑

8.2.2 Tumbling Window(滚动窗口)

// 5 分钟滚动窗口(Event Time)DataStream<Stats>result=stream.keyBy(Event::getVin).window(TumblingEventTimeWindows.of(Time.minutes(5))).aggregate(newStatsAggregator());// 1 天滚动窗口(按天对齐).window(TumblingEventTimeWindows.of(Time.days(1)))

8.2.3 Sliding Window(滑动窗口)

// 5 分钟窗口,每 1 分钟滑动一次DataStream<Stats>result=stream.keyBy(Event::getVin).window(SlidingEventTimeWindows.of(Time.minutes(5),Time.minutes(1))).aggregate(newStatsAggregator());

关键点:滑动窗口的窗口大小必须是滑动步长的整数倍,否则会有"残缺"窗口。

8.2.4 Session Window(会话窗口)

// 静态 Gap: 10 分钟无活动则切窗DataStream<Stats>result=stream.keyBy(Event::getVin).window(EventTimeSessionWindows.withGap(Time.minutes(10))).aggregate(newStatsAggregator());// 动态 Gap: 每条数据决定 Gap.window(EventTimeSessionWindows.withDynamicGap(newSessionWindowTimeGapExtractor<Event>(){@Overridepubliclongextract(Eventevent){// 速度 < 5 时 Gap 长,速度 > 80 时 Gap 短returnevent.getSpeed()<5?30*60*1000L:5*60*1000L;}}));

8.2.5 Global Window

// 全局窗口: 需自定义触发器DataStream<Stats>result=stream.keyBy(Event::getVin).window(GlobalWindows.create()).trigger(newCountTrigger(100))// 每 100 条触发.aggregate(newStatsAggregator());

8.3 时间窗口 vs 计数窗口

8.3.1 对比

维度时间窗口计数窗口
触发条件时间推进元素数量
类型Tumbling/Sliding/SessionGlobal + CountTrigger
乱序处理Watermark不涉及
适用实时报表批处理模拟

项目经验:车联网行程切分用 Session Window(Event Time,动态 Gap),保证"车辆熄火 10 分钟"则行程闭合。


8.4 窗口函数

窗口触发后调用窗口函数计算结果。Flink 提供三种。

8.4.1 三种函数对比

函数增量计算访问窗口上下文性能
ReduceFunction✅❌最高
AggregateFunction✅❌高
ProcessWindowFunction❌(全量缓存)✅低

8.4.2 ReduceFunction

publicclassSumMileageReducerimplementsReduceFunction<Stats>{@OverridepublicStatsreduce(Statsa,Statsb)throwsException{a.setMileage(a.getMileage()+b.getMileage());returna;// 输入输出类型必须相同}}stream.keyBy(Event::getVin).window(TumblingEventTimeWindows.of(Time.minutes(5))).reduce(newSumMileageReducer());

特点:增量聚合,每条数据到达即合并,状态只保存一个累加值。性能最高,但输入输出类型必须相同。

8.4.3 AggregateFunction

publicclassStatsAggregatorimplementsAggregateFunction<Event,StatsAccumulator,Stats>{@OverridepublicStatsAccumulatorcreateAccumulator(){returnnewStatsAccumulator(0.0,0,Double.MIN_VALUE);}@OverridepublicStatsAccumulatoradd(Eventevent,StatsAccumulatoracc){acc.mileage+=event.getMileage();acc.count++;acc.maxSpeed=Math.max(acc.maxSpeed,event.getSpeed());returnacc;}@OverridepublicStatsgetResult(StatsAccumulatoracc){returnnewStats(acc.mileage,acc.count,acc.maxSpeed);}@OverridepublicStatsAccumulatormerge(StatsAccumulatora,StatsAccumulatorb){a.mileage+=b.mileage;a.count+=b.count;a.maxSpeed=Math.max(a.maxSpeed,b.maxSpeed);returna;}}stream.keyBy(Event::getVin).window(TumblingEventTimeWindows.of(Time.minutes(5))).aggregate(newStatsAggregator());

特点:增量聚合,输入/累加器/输出三类型可不同,灵活性强。性能高,推荐。

8.4.4 ProcessWindowFunction

publicclassWindowStatsFunctionextendsProcessWindowFunction<Event,Stats,String,TimeWindow>{@Overridepublicvoidprocess(Stringvin,Contextctx,Iterable<Event>events,Collector<Stats>out){doubletotalMileage=0;intcount=0;for(Evente:events){totalMileage+=e.getMileage();count++;}// 可访问窗口上下文longwindowStart=ctx.window().getStart();longwindowEnd=ctx.window().getEnd();longwatermark=ctx.currentWatermark();out.collect(newStats(vin,totalMileage,count,windowStart,windowEnd));}}

特点:全量缓存窗口内所有数据,可访问窗口上下文(起止时间、Watermark、状态)。性能最低,慎用。

8.4.5 组合:增量 + Process

// 最佳实践: 增量聚合 + Process 获取上下文stream.keyBy(Event::getVin).window(TumblingEventTimeWindows.of(Time.minutes(5))).aggregate(newStatsAggregator(),newWindowStatsFunction());

原理:先增量聚合(每条数据合并),触发时把聚合结果传给 ProcessWindowFunction,只处理一次。既高效又能拿到上下文。


8.5 触发器与清理器

8.5.1 触发器(Trigger)

决定何时触发窗口计算。默认触发器:

窗口类型默认触发器
EventTime 窗口Watermark >= 窗口结束
ProcessingTime 窗口处理时间到窗口结束
Global Window必须自定义

8.5.2 自定义触发器

publicclassVehicleTriggerextendsTrigger<Event,TimeWindow>{@OverridepublicTriggerResultonElement(Eventevent,longtimestamp,TimeWindowwindow,TriggerContextctx){// 1. 每收到 100 条触发一次if(ctx.getPartitionedState(valueStateDescriptor).value()>=100){returnTriggerResult.FIRE;}// 2. 注册窗口结束定时器ctx.registerEventTimeTimer(window.maxTimestamp());returnTriggerResult.CONTINUE;}@OverridepublicTriggerResultonEventTime(longtime,TimeWindowwindow,TriggerContextctx){returntime==window.maxTimestamp()?TriggerResult.FIRE:TriggerResult.CONTINUE;}@OverridepublicTriggerResultonProcessingTime(longtime,TimeWindowwindow,TriggerContextctx){returnTriggerResult.CONTINUE;}@Overridepublicvoidclear(TimeWindowwindow,TriggerContextctx){ctx.deleteEventTimeTimer(window.maxTimestamp());}}

8.5.3 清理器(Evictor)

在触发后、输出前移除某些元素:

stream.keyBy(Event::getVin).window(TumblingEventTimeWindows.of(Time.minutes(5))).evictor(newMyEvictor())// 清理某些元素.aggregate(newStatsAggregator());

慎用:Evictor 需要遍历窗口内所有元素,性能差。优先用 AggregateFunction 的过滤逻辑替代。


8.6 会话窗口动态 Gap

8.6.1 动态 Gap 应用

车联网行程切分场景,根据车辆状态动态调整 Gap:

publicclassVehicleGapExtractorimplementsSessionWindowTimeGapExtractor<Event>{@Overridepubliclongextract(Eventevent){// 1. 熄火状态: Gap 30 分钟if(event.getSpeed()==0&&event.getIgnition()==0){return30*60*1000L;// 30 分钟}// 2. 怠速状态: Gap 10 分钟if(event.getSpeed()==0){return10*60*1000L;}// 3. 行驶状态: Gap 1 分钟(短暂停车不算断开)return60*1000L;}}stream.keyBy(Event::getVin).window(EventTimeSessionWindows.withDynamicGap(newVehicleGapExtractor())).aggregate(newTripAggregator());

8.6.2 会话合并

事件: A(10:00) B(10:05) C(10:20) D(10:40) Gap = 10 分钟: A-B: 间隔 5 分钟 < Gap → 同一会话 B-C: 间隔 15 分钟 > Gap → 切窗 C-D: 间隔 20 分钟 > Gap → 切窗 → 两个会话: [A,B] [C] [D] Gap = 30 分钟: A-B-C-D 都在同会话

8.7 窗口性能优化与状态管理

8.7.1 常见性能问题

问题原因解法
状态膨胀ProcessWindowFunction 全量缓存改用 AggregateFunction
触发延迟Watermark 推进慢配合空闲超时 / 对齐
数据倾斜某 key 窗口数据多两阶段聚合
小窗口过多窗口太小增大窗口

8.7.2 状态优化

// 1. 用 AggregateFunction 替代 ProcessWindowFunction.aggregate(newStatsAggregator());// 2. 配置 TTL(防止历史窗口状态不清理)StateTtlConfigttlConfig=StateTtlConfig.newBuilder(Time.hours(48)).setUpdateType(StateTtlConfig.UpdateType.OnCreateAndWrite).build();// 3. 选择合适状态后端// 小状态(聚合结果): HashMap// 大状态(全量缓存): ForSt

8.7.3 SQL 等价写法

-- Tumbling Window (5 分钟)SELECTvin,TUMBLE_START(event_time,INTERVAL'5'MINUTE)ASwindow_start,TUMBLE_END(event_time,INTERVAL'5'MINUTE)ASwindow_end,SUM(mileage)AStotal_mileageFROMvehicle_eventsGROUPBYvin,TUMBLE(event_time,INTERVAL'5'MINUTE);-- Sliding Window (5 分钟窗口,1 分钟滑动)SELECTvin,HOP_START(event_time,INTERVAL'1'MINUTE,INTERVAL'5'MINUTE)ASwindow_start,HOP_END(event_time,INTERVAL'1'MINUTE,INTERVAL'5'MINUTE)ASwindow_end,SUM(mileage)AStotal_mileageFROMvehicle_eventsGROUPBYvin,HOP(event_time,INTERVAL'1'MINUTE,INTERVAL'5'MINUTE);-- Session Window (动态 Gap)SELECTvin,SESSION_START(event_time,INTERVAL'10'MINUTE)ASwindow_start,SESSION_END(event_time,INTERVAL'10'MINUTE)ASwindow_end,SUM(mileage)AStotal_mileageFROMvehicle_eventsGROUPBYvin,SESSION(event_time,INTERVAL'10'MINUTE);

生产建议:实时大屏优先用 SQL TVF(Table-Valued Function)写法,更简洁:

-- TVF 写法(Flink 2.x 推荐)SELECT*FROMTABLE(TUMBLE(TABLEvehicle_events,DESCRIPTOR(event_time),INTERVAL'5'MINUTE))GROUPBYwindow_start,window_end,vin;

8.8 本章小结与下章预告

本章小结

┌────────────────────────────────────────────────────────────────┐ │ 第 8 章 要点回顾 │ └────────────────────────────────────────────────────────────────┘ ✓ 窗口 = 把无界流切有限批次,用于聚合 两阶段: 累积 → 触发计算 Keyed Window(并行) > Non-Keyed(单点) ✓ 四种窗口类型: Tumbling(滚动,不重叠) → 每分钟统计 Sliding(滑动,可重叠) → 滑动平均 Session(会话,动态) → 用户会话/车辆行程 Global(全局,自定义触发) ✓ 三种窗口函数: ReduceFunction(增量,同类型,最快) AggregateFunction(增量,三类型,推荐) ProcessWindowFunction(全量,上下文,慢) 组合: 增量 + Process(最佳) ✓ 触发器: 决定何时触发 EventTime: Watermark >= 窗口结束 自定义: 计数 + 时间 ✓ 会话窗口动态 Gap: 根据事件属性决定 Gap(熄火/怠速/行驶) ✓ 性能优化: AggregateFunction > ProcessWindowFunction TTL 防状态膨胀 SQL TVF 写法(Flink 2.x 推荐)

下章预告

第 9 章 多流 Join 与异步 IO:讲解双流 Join(窗口 Join、Interval Join)、维表 Join(Lookup Join 同步/异步)、流批 Join 选型、AsyncDataStream 异步 IO(Ordered/Unordered)、Broadcast State 模式、Join 性能优化与数据倾斜处理——这是实时数仓 Join 与维表关联的核心。

官方参考资料

  • Flink Windows:https://nightlies.apache.org/flink/flink-docs-stable/docs/dev/datastream/windows/
  • Window Functions:https://nightlies.apache.org/flink/flink-docs-stable/docs/dev/datastream/windows/#window-functions
  • SQL Windows:https://nightlies.apache.org/flink/flink-docs-stable/docs/dev/table/sql/queries/window/
  • Window TVF:https://nightlies.apache.org/flink/flink-docs-stable/docs/dev/table/sql/queries/window-tvf/

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询