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 Window | Non-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/Session | Global + 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// 大状态(全量缓存): ForSt8.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/