
示例工程大数据【免费下载链接】flink-learningflink learning blog. http://www.54tianzhisheng.cn/ 含 Flink 入门、概念、原理、实战、性能调优、源码解析等内容。涉及 Flink Connector、Metrics、Library、DataStream API、Table API SQL 等内容的学习案例还有 Flink 落地应用的大型项目案例PVUV、日志存储、百亿数据实时去重、监控告警分享。欢迎大家支持我的专栏《大数据实时计算引擎 Flink 实战与性能优化》项目地址https://gitcode.com/gh_mirrors/fl/flink-learning点击查看免费下载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 在 KeyedStreamDataStream 的继承类中提供了下面几种 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.javaEventTime 时间语义 Watermark 滚动事件时间窗口TumblingEventTimeWindowsMain3.java滚动处理时间窗口TumblingProcessingTimeWindowsMain4.java事件时间会话窗口EventTimeSessionWindowsMain5.java非 Keyed Stream 的全局窗口 windowAll 用法WindowAll.javatimeWindowAll 用法CustomTriggerMain.java自定义 Trigger 的完整示例入口CustomTrigger.java自定义 Trigger 实现CustomSource.java模拟数据源每秒随机生成一个 WordEventLineSplitter.java将 socket 输入的每行文本拆分为(数量, 单词)二元组WordEvent.java事件模型word、count、timestamp 三个字段WindowConstant.javahostName、port 两个运行参数的常量定义TestWindowSize.java验证窗口起始时间与滑动窗口窗口集合的计算逻辑其中LineSplitter.java 是一个典型的FlatMapFunctionString, Tuple2Long, String它将终端输入的文本按空格切分只有当第一个 token 能被解析为 long 时才输出(Long.valueOf(tokens[0]), tokens[1])即每条记录携带一个数值与一个单词供后续窗口做 sum 聚合。而 CustomSource.java 继承RichSourceFunctionWordEvent通过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) % slideoffset 默认为 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继承自TriggerWordEvent, 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 是窗口机制的三大核心组件也是自定义窗口时最常扩展的切入点。此外针对未分组的数据流非 KeyedStreamFlink 还提供了全局窗口 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滚动时间窗口ProcessingTimeTumblingEventTimeWindows滚动时间窗口EventTimeSlidingProcessingTimeWindows滑动时间窗口ProcessingTimeSlidingEventTimeWindows滑动时间窗口EventTimeProcessingTimeSessionWindows会话窗口ProcessingTimeEventTimeSessionWindows会话窗口EventTimeGlobalWindows全局窗口需配合计数 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每个元素进入窗口时调用可在此时注册定时器、累积计数状态onProcessingTimeProcessingTime 定时器到点回调配合ctx.registerProcessingTimeTimer使用onEventTimeEventTime 定时器到点回调配合 Watermark 推进触发onMerge会话窗口合并时合并两个窗口的触发状态clear窗口清理时执行应在此移除注册的定时器与分区状态。CustomTriggerMain.java 通过.trigger(CustomTrigger.creat())将自定义 Trigger 接入窗口算子展示了一条完整的默认时间窗口 定制触发语义的实践路径。值得注意的是CustomTrigger.java 中被注释的代码展示了利用ReducingStateLong记录下次触发时间、周期性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决定窗口何时触发计算四种TriggerResultEvictor决定触发前后哪些元素可以被剔除。这三者正是为自定义 Window 提供方法的核心切入点仓库中的 CustomTriggerMain.java 与 CustomTrigger.java 就是内置窗口 自定义触发的最佳实践模板。在实际选型时可以依据以下原则判断统计固定时间周期内的指标如每分钟车流量、每 5 分钟 PV/UV优先考虑Time Window并结合业务容忍的乱序程度决定使用 ProcessingTime 还是 EventTime Watermark统计固定条数内的指标如每 1000 条日志做一次聚合使用Count Window面向用户行为轨迹、连续登录会话等天然以活跃间隔划分的场景使用Session Window当内置窗口无法满足何时计算、计算哪些数据的语义时通过自定义Trigger / Evictor / WindowAssigner定制。窗口机制是 Flink 流处理的核心抽象之一理解它就理解了如何在无界流上进行有界计算这一流处理的根本命题。赞分享示例工程大数据【免费下载链接】flink-learningflink learning blog. http://www.54tianzhisheng.cn/ 含 Flink 入门、概念、原理、实战、性能调优、源码解析等内容。涉及 Flink Connector、Metrics、Library、DataStream API、Table API SQL 等内容的学习案例还有 Flink 落地应用的大型项目案例PVUV、日志存储、百亿数据实时去重、监控告警分享。欢迎大家支持我的专栏《大数据实时计算引擎 Flink 实战与性能优化》项目地址https://gitcode.com/gh_mirrors/fl/flink-learning点击查看免费下载相关推荐Flink SQL 窗口去重Window Deduplication完整指南语法、原理与实战Flink SQL 窗口去重Window Deduplication完整指南语法、原理与实战 Window Deduplication窗口去重是 Ap后端大数据流处理批处理Flink 窗口去重Window DeduplicationSQL 详解语法、示例与实现原理Flink 窗口去重Window DeduplicationSQL 详解语法、示例与实现原理 窗口去重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),仅供参考