- 示例工程
- 大数据
【免费下载链接】flink-learning
flink learning blog. http://www.54tianzhisheng.cn/ 含 Flink 入门、概念、原理、实战、性能调优、源码解析等内容。涉及 Flink Connector、Metrics、Library、DataStream API、Table API & SQL 等内容的学习案例,还有 Flink 落地应用的大型项目案例(PVUV、日志存储、百亿数据实时去重、监控告警)分享。欢迎大家支持我的专栏《大数据实时计算引擎 Flink 实战与性能优化》
Apache Flink 是面向生产环境的流处理器,其 API 在数据流上提供了非常灵活的窗口定义,这使其在众多开源流处理框架中脱颖而出。本篇以无界流上的"窗口化(windowing)"思维为核心,先通过交通传感器统计汽车数量的现实场景引出 Window 的必要性,再系统讲解 Flink 内置的 Time Window、Count Window、Session Window 三种窗口的用法与实现原理,最后深入 Window 内部三大核心组件 WindowAssigner、Trigger、Evictor 的源码级机制。读完本篇,你将能理解"为什么流处理需要窗口、如何选择窗口类型、如何通过自定义 Trigger 定制窗口触发语义",并能基于仓库 flink-learning-window 模块中的完整示例代码直接上手实践。
1. Window 简介:从批处理到流处理的思维转变
目前有许多数据分析场景正从批处理演变为流处理。虽然可以将批处理作为流处理的特殊情况来处理,但是分析无穷集的流数据通常需要思维方式的转变,并且具有自己的术语,例如"windowing(窗口化)""at-least-once(至少一次)""exactly-once(只有一次)"。
对于刚刚接触流处理的人来说,这种转变和新术语可能会非常混乱。这里结合一个现实的例子来说明:统计经过某红绿灯的汽车数量之和。
假设在一个红绿灯处,我们每隔 15 秒统计一次通过此红绿灯的汽车数量。可以把汽车的经过看成一个流——一个无穷的流,不断有汽车经过此红绿灯,因此无法统计"总共"的汽车数量。但是可以换一种思路:每隔 15 秒,我们都与上一次的结果做一次 sum 操作(滑动聚合)。
这个结果依然无法回答"总共经过多少辆车"的问题,根本原因在于流是无界的。我们虽然不能限制流,但可以在一个有界的范围内处理无界的流数据。因此,需要换一个问题的提法:每分钟经过某红绿灯的汽车数量之和?
这个问题就相当于定义了一个 Window(窗口),Window 的界限是 1 分钟,且每分钟内的数据互不干扰,因此也可以称为翻滚(不重合)窗口。假设第一分钟的数量为 18,第二分钟是 28,第三分钟是 24……这样,1 个小时内会有 60 个 Window。
再考虑一种情况:每 30 秒统计一次过去 1 分钟的汽车数量之和。此时 Window 出现了重合,1 个小时内会有 120 个 Window。这就是滑动窗口的典型形态:窗口大小(size)为 1 分钟,滑动步长(slide)为 30 秒。
2. Window 有什么作用?
通常来讲,Window 就是用来对一个无限的流设置一个有限的集合,在有界的数据集上进行操作的一种机制。它的本质是将无界流切分为一个个有界的片段,让聚合、连接等批式算子得以在流上复用。
Window 又可以分为两大类:
- 基于时间(Time-based)的 Window:按时间边界切分数据,是流处理中最常用的一类;
- 基于数量(Count-based)的 Window:按元素个数切分数据,与时间无关。
3. Flink 自带的 Window 类型
Flink 在 KeyedStream(DataStream 的继承类)中提供了下面几种 Window:
- 以时间驱动的Time Window;
- 以事件数量驱动的Count Window;
- 以会话间隔驱动的Session Window。
提供上面三种 Window 机制后,由于某些特殊的需要,DataStream API 也提供了定制化的 Window 操作,供用户自定义 Window(例如自定义 WindowAssigner、自定义 Trigger 等)。
在仓库中,以上内容对应的完整学习工程位于 flink-learning-window 模块,其pom.xml依赖了公共模块 flink-learning-common(提供数据模型与执行环境工具类)。
4. 实战准备:窗口学习案例的工程结构
在深入每种窗口之前,先了解仓库中窗口模块的组成,这将帮助我们快速对照代码运行示例:
| 文件 | 作用 |
|---|---|
| Main.java | 综合演示滚动/滑动 Time Window、Count Window、Session Window 的用法 |
| Main2.java | EventTime 时间语义 + Watermark + 滚动事件时间窗口(TumblingEventTimeWindows) |
| Main3.java | 滚动处理时间窗口(TumblingProcessingTimeWindows) |
| Main4.java | 事件时间会话窗口(EventTimeSessionWindows) |
| Main5.java | 非 Keyed Stream 的全局窗口 windowAll 用法 |
| WindowAll.java | timeWindowAll 用法 |
| CustomTriggerMain.java | 自定义 Trigger 的完整示例入口 |
| CustomTrigger.java | 自定义 Trigger 实现 |
| CustomSource.java | 模拟数据源,每秒随机生成一个 WordEvent |
| LineSplitter.java | 将 socket 输入的每行文本拆分为(数量, 单词)二元组 |
| WordEvent.java | 事件模型:word、count、timestamp 三个字段 |
| WindowConstant.java | hostName、port 两个运行参数的常量定义 |
| TestWindowSize.java | 验证窗口起始时间与滑动窗口窗口集合的计算逻辑 |
其中,LineSplitter.java 是一个典型的FlatMapFunction<String, Tuple2<Long, String>>:它将终端输入的文本按空格切分,只有当第一个 token 能被解析为 long 时才输出(Long.valueOf(tokens[0]), tokens[1]),即每条记录携带一个数值与一个单词,供后续窗口做 sum 聚合。而 CustomSource.java 继承RichSourceFunction<WordEvent>,通过while (isRunning)循环每 1 秒ctx.collect一个随机的WordEvent(word, count, System.currentTimeMillis()),为测试窗口提供了稳定的数据源。
窗口模块的运行前提是执行环境工具类提供的参数(hostName、port)。以 Main.java 为例,其操作方式是:在终端执行nc -l 9000,然后输入 long text 类型的数据,程序通过env.socketTextStream(hostName, port)读取数据流。
5. Time Window 的用法及源码分析
Time Window 以时间为边界切分流数据,是使用最广泛的一类窗口。它又分为滚动时间窗口(Tumbling Window)与滑动时间窗口(Sliding Window)两种。
5.1 滚动时间窗口
Main.java 中滚动时间窗口的核心代码为:
data.flatMap(new LineSplitter()) .keyBy(1) .timeWindow(Time.seconds(30)) .sum(0) .print();即先对数据按 key(元组的第 1 个字段)分组,再开一个30 秒的滚动时间窗口,窗口内对第 0 个字段(数量)做 sum 聚合。滚动窗口的特点是:窗口之间不重合,每条数据恰好属于一个窗口;N 秒的滚动窗口在一个小时内有 3600/N 个窗口。
5.2 滑动时间窗口
Main.java 中滑动时间窗口的核心代码为:
data.flatMap(new LineSplitter()) .keyBy(1) .timeWindow(Time.seconds(60), Time.seconds(30)) .sum(0) .print();timeWindow(size, slide)第一个参数是窗口大小(60 秒),第二个参数是滑动步长(30 秒)。即每 30 秒统计一次过去 60 秒的数据,窗口之间出现重合,一条数据可能同时属于多个窗口;60 秒窗口、30 秒滑动步长的场景下,1 个小时内会产生 120 个窗口。这与前文交通传感器"每 30 秒统计过去 1 分钟汽车数量"的案例完全对应。
5.3 时间语义与底层 Assigner
从 Main.java 源码的注释可以确认一个关键事实:
//如果不指定时间的话,默认是 ProcessingTime,但是如果指定为事件事件的话,需要事件中带有时间或者添加时间水印 // env.setStreamTimeCharacteristic(TimeCharacteristic.EventTime);也就是说,Flink 默认使用ProcessingTime(处理时间)语义;一旦切换为EventTime(事件时间),则要求事件自带时间戳,并且必须为流分配 Watermark(水印)。
仓库针对两种时间语义分别给出了对照示例:
- Main2.java 设置了
env.setStreamTimeCharacteristic(TimeCharacteristic.EventTime)并setParallelism(1),随后通过assignTimestampsAndWatermarks分配时间戳与水印:extractTimestamp提取WordEvent的时间戳并维护当前最大时间戳,getCurrentWatermark返回currentTimestamp - maxTimeLag(此处 maxTimeLag 为 5000 毫秒,即允许 5 秒的乱序延迟)。之后使用TumblingEventTimeWindows.of(Time.seconds(10))开 10 秒滚动事件时间窗口,并通过WindowFunction.apply打印window.getStart()与window.getEnd()观察窗口边界。 - Main3.java 则直接使用
TumblingProcessingTimeWindows.of(Time.seconds(10))开 10 秒滚动处理时间窗口,无需 Watermark,同样在apply中打印窗口起止时间。
从源码结构看,timeWindow(Time)、timeWindow(Time, Time)这类便捷方法底层都会转换为对应的窗口分配器(WindowAssigner):滚动时间窗口对应TumblingEventTimeWindows/TumblingProcessingTimeWindows,滑动时间窗口对应SlidingEventTimeWindows/SlidingProcessingTimeWindows。
5.4 窗口起始时间的计算逻辑
TestWindowSize.java 直接验证了滚动窗口起始时间与滑动窗口窗口集合的计算公式:
//timestamp - (timestamp - offset + slide) % slide; System.out.println(l - (l + 60 * 1000) % 60000);这对应 Flink 窗口分配器计算窗口起始时间的核心公式:start = timestamp - (timestamp - offset + slide) % slide(offset 默认为 0)。对于 60 秒的窗口,任意时间戳都会被"对齐"到整分钟的边界上。
该测试还模拟了滑动窗口的窗口集合计算:
long size = Time.hours(24).toMilliseconds(); long slide = Time.hours(1).toMilliseconds(); long lastStart = (1572794063000l - (1572794063000l + slide) % slide); for (long start = lastStart; start > 1572794063000l - size; start -= slide) { System.out.println(start + " " + (start + size)); }即给定一个事件时间戳,从"最后一个不晚于该时间戳的窗口起点"开始,以 slide 为步长向前回溯,直到窗口起点早于timestamp - size为止,从而枚举出该事件所属的全部滑动窗口。这正是滑动窗口中"一条数据同时属于多个窗口"的实现基础——从源码结构与测试逻辑可以推断,SlidingWindowAssigner就是按上述公式为每个元素计算并分配多个窗口的。
6. Count Window 的用法及源码分析
Count Window 与时间无关,而是以元素个数为边界。当某个 key 的累积元素数达到设定阈值时,窗口即触发计算。
Main.java 中滚动计数窗口的核心代码为:
data.flatMap(new LineSplitter()) .keyBy(1) .countWindow(3) .sum(0) .print();countWindow(3)表示每收集满 3 个元素就触发一次计算,窗口之间不重合,即滚动计数窗口。
滑动计数窗口的写法为:
data.flatMap(new LineSplitter()) .keyBy(1) .countWindow(4, 3) .sum(0) .print();countWindow(4, 3)第一个参数是窗口长度(4 个元素),第二个参数是滑动步长(每 3 个元素滑动一次),即每来 3 个元素,就对最近 4 个元素做一次聚合,窗口之间存在重合。
从 Flink 的实现机制看,Count Window 底层由GlobalWindows(全局窗口分配器)配合计数 Trigger 实现:所有数据先被分配到一个全局窗口中,再依赖 Trigger 按元素数量决定何时触发计算。因此 Count Window 虽然 API 简单,本质上是"全局窗口 + 计数触发"的组合,这一点通过自定义 Trigger 也能自行实现。
7. Session Window 的用法及源码分析
Session Window 与固定大小的窗口不同,它以**活跃间隙(gap)**为边界:当某个 key 在 gap 时间内没有新数据到来,则认为一次会话结束,将该 gap 之前的数据合并为一个窗口进行聚合。
Main.java 中使用的是处理时间会话窗口:
data.flatMap(new LineSplitter()) .keyBy(1) .window(ProcessingTimeSessionWindows.withGap(Time.seconds(5))) .sum(0) .print();代码注释明确说明:withGap(Time.seconds(5))表示如果 5 秒内没出现数据,则认为超出会话时长,然后计算这个窗口的和。注意,Session Window 没有timeWindow便捷方法,必须显式调用.window(...)并传入会话窗口分配器。
仓库还提供了事件时间会话窗口的示例 Main4.java:它从 socket 读取"单词,时间戳"格式的数据,通过assignTimestampsAndWatermarks提取事件时间戳(当前 watermark 直接等于当前最大时间戳),然后:
keyBy(0) .window(EventTimeSessionWindows.withGap(Time.minutes(5))) .sum(1) .print("session ");即按第一个字段(单词)分组,若该单词在5 分钟事件时间内没有新数据到来,则触发一次会话窗口计算。
从实现机制看,Session Window 的两个关键点:
- gap 判定:每条数据到来时,会计算它与当前窗口边界的时间差,若超过 gap 则关闭当前窗口、开启新窗口;
- 窗口合并(merge):由于乱序与迟到数据可能使两个相邻会话窗口"相遇",Flink 的 Session Window Assigner 具备窗口合并能力,
Trigger.onMerge正是为会话窗口合并状态而设计的钩子方法(详见第 10 节)。
8. 如何自定义 Window?
当内置的窗口语义无法满足需求时,Flink 的 DataStream API 允许用户通过组合与替换 Window 的内部组件来定制窗口。自定义窗口的典型方式有两种:
- 自定义 Trigger:保留默认的 WindowAssigner(如滚动时间窗口),但替换触发逻辑,改变窗口的计算时机;
- 自定义 WindowAssigner / Evictor:完全自定义窗口的分配规则或数据清理规则。
仓库中 CustomTriggerMain.java 给出了第一种方式的完整示例:
data.keyBy(WordEvent::getWord) .timeWindow(Time.seconds(10)) .trigger(CustomTrigger.creat()) .sum("count") .print();它先为CustomSource产生的WordEvent流分配时间戳与水印(同样采用 5 秒最大乱序延迟的周期水印),按word字段分组,开 10 秒滚动时间窗口,然后通过.trigger(CustomTrigger.creat())挂载自定义触发器。
这里的CustomTrigger位于 CustomTrigger.java,继承自Trigger<WordEvent, TimeWindow>,完整覆盖了 Trigger 的五个生命周期方法:
onElement:每个元素被添加到窗口时调用,源码中打印window.getStart()与window.getEnd()并返回TriggerResult.CONTINUE;被注释掉的代码演示了更完整的实现思路——用ReducingState记录下次触发时间,registerProcessingTimeTimer注册处理时间定时器;onProcessingTime:已注册的 ProcessingTime 定时器启动时调用;onEventTime:已注册的 EventTime 定时器启动时调用;onMerge:与状态性触发器相关,当使用会话窗口、两个触发器对应的窗口合并时,合并两个触发器的状态;clear:执行任何需要清除的相应窗口(如清理定时器与状态)。
从源码结构看,CustomTrigger的类注释明确列出了TriggerResult的四种可能返回值,这是自定义 Trigger 的核心语义,详见第 10 节。通过这套机制,用户完全可以实现"每到达 N 条数据触发一次""在指定时间点触发一次"等任意定制化触发语义。
9. Window 的源码级工作原理
综合仓库代码与 Flink 的窗口抽象,一个 Window 算子的执行链路可以概括为:
- WindowAssigner(窗口分配器):决定每个元素归属于哪个/哪些窗口(见 5.4 节的起始时间公式);
- 窗口内部维护(State):每个窗口维护自己的聚合状态(如 sum 的累加值),数据按 key 分组、按窗口隔离;
- Trigger(触发器):决定窗口何时可以计算(FIRE)、何时清理(PURGE),其返回的
TriggerResult直接控制窗口生命周期; - Evictor(驱逐器):在触发计算前后,可选择性地从窗口中移除部分元素;
- 窗口函数(WindowFunction / ProcessWindowFunction / sum 等聚合):对窗口内数据执行最终的计算与输出。
其中 WindowAssigner、Trigger、Evictor 是窗口机制的三大核心组件,也是自定义窗口时最常扩展的切入点。此外,针对未分组的数据流(非 KeyedStream),Flink 还提供了全局窗口 API:如 WindowAll.java 中的timeWindowAll(Time.seconds(10)),以及 Main5.java 中通过timeWindowAll与windowAll(TumblingProcessingTimeWindows.of(Time.seconds(10)))配合ProcessAllWindowFunction处理全部数据——这类窗口的并行度受限(本质上只有一个"全量窗口"),使用时需注意。
10. Window 内部组件详解
10.1 WindowAssigner:用法与源码分析
WindowAssigner 负责为每条数据分配窗口,是窗口机制的入口。Flink 内置的 WindowAssigner 家族包括:
| Assigner | 窗口类型 | 时间语义 |
|---|---|---|
TumblingProcessingTimeWindows | 滚动时间窗口 | ProcessingTime |
TumblingEventTimeWindows | 滚动时间窗口 | EventTime |
SlidingProcessingTimeWindows | 滑动时间窗口 | ProcessingTime |
SlidingEventTimeWindows | 滑动时间窗口 | EventTime |
ProcessingTimeSessionWindows | 会话窗口 | ProcessingTime |
EventTimeSessionWindows | 会话窗口 | EventTime |
GlobalWindows | 全局窗口(需配合计数 Trigger) | 无 |
仓库中的使用示例:
- Main3.java 使用
TumblingProcessingTimeWindows.of(Time.seconds(10)); - Main2.java 使用
TumblingEventTimeWindows.of(Time.seconds(10)); - Main.java 使用
ProcessingTimeSessionWindows.withGap(Time.seconds(5)); - Main4.java 使用
EventTimeSessionWindows.withGap(Time.minutes(5))。
从源码结构与 TestWindowSize.java 的测试逻辑可以确认:时间窗口分配器通过start = timestamp - (timestamp - offset + slide) % slide将时间戳对齐到窗口边界;滑动窗口则为一个元素计算并返回多个窗口。
10.2 Trigger:用法与源码分析
Trigger 决定窗口何时触发计算,是窗口机制中"语义"最强的组件。CustomTrigger.java 的类注释给出了TriggerResult的四种可能取值,这是理解自定义 Trigger 的关键:
| TriggerResult | 含义 |
|---|---|
CONTINUE | 什么也不做,继续累积数据 |
FIRE | 触发计算(执行窗口函数并输出结果),但不清除窗口数据 |
PURGE | 清除窗口中的数据,不触发计算 |
FIRE_AND_PURGE | 触发计算,并清除窗口中的数据 |
Trigger 的四个核心方法(外加onMerge、clear)对应窗口生命周期的不同时刻:
onElement:每个元素进入窗口时调用,可在此时注册定时器、累积计数状态;onProcessingTime:ProcessingTime 定时器到点回调,配合ctx.registerProcessingTimeTimer使用;onEventTime:EventTime 定时器到点回调,配合 Watermark 推进触发;onMerge:会话窗口合并时合并两个窗口的触发状态;clear:窗口清理时执行,应在此移除注册的定时器与分区状态。
CustomTriggerMain.java 通过.trigger(CustomTrigger.creat())将自定义 Trigger 接入窗口算子,展示了一条完整的"默认时间窗口 + 定制触发语义"的实践路径。值得注意的是,CustomTrigger.java 中被注释的代码展示了利用ReducingState<Long>记录下次触发时间、周期性registerProcessingTimeTimer实现"每隔 interval 毫秒触发一次"的经典实现,从该代码结构可以推断 Flink 内置的ProcessingTimeTrigger正是基于类似的定时器机制实现的。
10.3 Evictor:用法与源码分析
Evictor(驱逐器)负责在窗口触发计算前后,从窗口中有选择地移除元素。它通常与 Trigger 配合使用,用于实现诸如"只对最近 N 个元素计算"等语义。
Evictor 的核心接口包含两个回调方法:
evictBefore:在窗口函数执行之前被调用,可先剔除部分元素再计算;evictAfter:在窗口函数执行之后被调用,可在输出结果后清理窗口元素。
Flink 内置的 Evictor 包括:
CountEvictor:保留窗口中最近 N 个元素,其余全部驱逐;TimeEvictor:保留窗口内时间戳落在最近一段时间内的元素;DeltaEvictor:基于用户自定义的 DeltaFunction 计算元素与基准值的差值,超出阈值的元素被驱逐。
Evictor 通过.evictor(...)方法挂载到窗口算子之上。需要说明的是,在仓库当前的 flink-learning-window 模块中,Evictor 主要作为窗口组件体系的一部分被介绍,尚无独立的示例类;其使用方式与 Trigger 一致,均是在窗口算子后链式调用。由于 Evictor 需要在触发前后遍历窗口内元素,使用它会带来额外的计算开销,因此只有在确实需要"部分数据参与计算"的场景下才建议使用。
11. 小结与反思
本节从生活案例出发分享了关于 Window 方面的需求,进而开始介绍 Window 相关的知识,并把 Flink 中常使用的三种窗口——Time Window、Count Window、Session Window——都一一做了介绍:它们的使用方式、时间语义差异、窗口起始时间的计算逻辑,以及基于 Main.java、Main2.java、Main3.java、Main4.java 的完整可运行代码。最后对 Window 的内部组件做了详细分析:
- WindowAssigner决定元素归属哪些窗口;
- Trigger决定窗口何时触发计算(四种
TriggerResult); - Evictor决定触发前后哪些元素可以被剔除。
这三者正是为自定义 Window 提供方法的核心切入点,仓库中的 CustomTriggerMain.java 与 CustomTrigger.java 就是"内置窗口 + 自定义触发"的最佳实践模板。
在实际选型时,可以依据以下原则判断:
- 统计固定时间周期内的指标(如每分钟车流量、每 5 分钟 PV/UV),优先考虑Time Window,并结合业务容忍的乱序程度决定使用 ProcessingTime 还是 EventTime + Watermark;
- 统计固定条数内的指标(如每 1000 条日志做一次聚合),使用Count Window;
- 面向用户行为轨迹、连续登录会话等天然以活跃间隔划分的场景,使用Session Window;
- 当内置窗口无法满足"何时计算、计算哪些数据"的语义时,通过自定义Trigger / Evictor / WindowAssigner定制。
窗口机制是 Flink 流处理的核心抽象之一,理解它,就理解了"如何在无界流上进行有界计算"这一流处理的根本命题。
- 示例工程
- 大数据
【免费下载链接】flink-learning
flink learning blog. http://www.54tianzhisheng.cn/ 含 Flink 入门、概念、原理、实战、性能调优、源码解析等内容。涉及 Flink Connector、Metrics、Library、DataStream API、Table API & SQL 等内容的学习案例,还有 Flink 落地应用的大型项目案例(PVUV、日志存储、百亿数据实时去重、监控告警)分享。欢迎大家支持我的专栏《大数据实时计算引擎 Flink 实战与性能优化》
相关推荐
Flink SQL 窗口去重(Window Deduplication)完整指南:语法、原理与实战
Flink SQL 窗口去重(Window Deduplication)完整指南:语法、原理与实战 Window Deduplication(窗口去重)是 Ap
后端大数据流处理批处理Flink 窗口去重(Window Deduplication)SQL 详解:语法、示例与实现原理
Flink 窗口去重(Window Deduplication)SQL 详解:语法、示例与实现原理 窗口去重(Window Deduplication)是 Fl
后端大数据流处理批处理flink-learning 之 Flink Window 窗口机制实战:Time Window、Count Window、Session Window 与自定义 Trigger 源码剖析
flink learning 之 Flink Window 窗口机制实战:Time Window、Count Window、Session Window 与自定
示例工程大数据
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考