)
Flink 流处理窗口类型全解析Tumbling / Sliding / Session 窗口原理与实战Data Engineering Zoomcamp【免费下载链接】data-engineering-zoomcampData Engineering Zoomcamp is a free 9-week course on building production-ready data pipelines. Join the course here 项目地址: https://gitcode.com/GitHub_Trending/da/data-engineering-zoomcamp本文是 Data Engineering Zoomcamp 流处理Streaming模块的窗口类型专题。在完成 10-aggregation-with-tumbling-windows.md 中的滚动窗口聚合之后本文将系统讲解 Apache Flink 支持的全部三种窗口类型——滚动窗口Tumbling、滑动窗口Sliding与会话窗口Session并结合作品仓库中真实的 PyFlink 作业代码说明每种窗口的边界语义、Flink SQL 表值函数写法、典型业务场景以及它们与水印Watermark、迟到事件Late Events、Upsert 之间的关系。读完本文你将能根据业务需求准确选型窗口类型并能在本仓库的出租车实时数据流水线中直接落地实现。一、为什么需要理解窗口类型窗口决定了事件如何被分组在流处理中数据是无穷无尽、持续到达的无法像批处理那样等待全部数据到齐再统一计算。窗口Window就是把无界数据流切分成有界片段的手段它决定了一个事件到底属于哪个计算桶bucket。窗口类型决定了事件归属于一个固定桶、多个重叠桶还是一个由不活动间隙界定的突发burst集合。在 12-understanding-window-types.md 中明确指出前面的单元已经使用过滚动窗口而 Flink 支持三种窗口类型窗口类型大小是否重叠事件归属窗口何时关闭Tumbling滚动固定否每个事件只属于恰好一个窗口到达固定时间点Sliding滑动固定是一个事件可属于多个窗口到达固定时间点多个起点Session会话动态否按不活动间隙切分超过指定时间的静默下面逐一深入。二、Tumbling 滚动窗口固定大小、互不重叠滚动窗口是三种窗口中最直观的一种固定大小、互不重叠fixed-size, non-overlapping每个事件恰好属于一个窗口。| Window 1 | Window 2 | Window 3 | | 1 hour | 1 hour | 1 hour |如果你来自批处理世界滚动窗口一定是最熟悉的概念——它只是把数据切分成固定片段本质上是加速版批处理。以 1 小时为例窗口就是 00:00-01:00、01:00-02:00、02:00-03:00……时间线被整齐地切成等长的桶互不交叉。典型应用场景按小时统计出租车行程数counting trips per hour、每日营收汇总daily revenue summaries。在本仓库中滚动窗口正是前一个单元 10-aggregation-with-tumbling-windows.md 的核心。参考实现位于 aggregation_job.py核心 SQL 为INSERT INTO processed_events_aggregated SELECT window_start, PULocationID, COUNT(*) AS num_trips, SUM(total_amount) AS total_revenue FROM TABLE( TUMBLE(TABLE events, DESCRIPTOR(event_timestamp), INTERVAL 1 HOUR) ) GROUP BY window_start, PULocationID;要点拆解TUMBLE(TABLE events, DESCRIPTOR(event_timestamp), INTERVAL 1 HOUR)是 Flink SQL 的**表值函数Table-Valued Function**形式三个参数分别是数据源表、事件时间列描述符、窗口大小DESCRIPTOR(event_timestamp)必须指向定义了 WATERMARK 的那一列本作业中是计算列event_timestamp AS TO_TIMESTAMP_LTZ(tpep_pickup_datetime, 3)配合WATERMARK否则无法进行基于事件时间的窗口切分分组键GROUP BY window_start, PULocationID同时按时间窗口与上车地点两个维度聚合这与 sink 表PRIMARY KEY (window_start, PULocationID)一一对应。仓库中还提供了一个便于观察窗口关闭行为的演示作业 aggregation_job_demo.py它把窗口从 1 小时改为 10 秒并特意保留了latest-offset与注释Use with producer_realtime.py to observe watermark behavior: - Watermark event_timestamp - 5 seconds - Late events (5s) arrive before the watermark closes the window - included - Late events (5s) may arrive after the watermark closes the window - dropped这说明滚动窗口虽然语义简单但在真实流中窗口何时关闭、迟到事件算不算完全取决于水印策略详见下文第五节。三、Sliding 滑动窗口固定大小、互相重叠滑动窗口同样是固定大小但窗口之间互相重叠overlapping一个事件可以同时属于多个窗口。提到1 小时窗口大多数人想到的是 00:00-01:00。但其实 00:15-01:15、00:30-01:30 也都是 1 小时窗口只是起点不同。滑动窗口把这些不同起点的窗口全部纳入计算|--- Window 1 (1 hour) ---| |--- Window 2 (1 hour) ---| |--- Window 3 (1 hour) ---| - 15 min slide -滑动窗口有两个关键参数窗口大小size与滑动步长slide。上图中窗口大小 1 小时、每 15 分钟滑动一次因此任意时刻同时有 4 个窗口处于打开状态一个新事件会被计入这 4 个窗口。对应的 Flink SQL 表值函数是HOP即原文档中给出的示例HOP(TABLE events, DESCRIPTOR(event_timestamp), INTERVAL 15 MINUTE, INTERVAL 1 HOUR)HOP的三个时间参数依次为滑动步长15 分钟、窗口大小1 小时。同样的示例也可以在仓库的 workshop/README.md 中查到这是 Flink SQL 对滑动窗口的标准写法在部分教材中也称为 Sliding Window / HOP Window。典型应用场景寻找峰值与低谷finding peaks and valleys——任意一个 1 小时窗口内我们的峰值流量是多少这种重叠窗口让你能精确锁定时间线上取值最高或最低的时刻非常适合求最值min-maxing、移动平均moving averages以及激增检测surge detection例如网约车平台的动态加价ride-share surge pricing要判断过去任意 1 小时内的需求是否异常飙升非重叠的滚动窗口可能会恰好把爆发点切分到两个桶里而错过峰值滑动窗口则不会。滑动窗口的代价是计算开销每个事件会被复制到多个窗口参与聚合重叠度越高slide 越小冗余计算越大。四、Session 会话窗口基于不活动间隙的动态窗口会话窗口与前两类有本质区别窗口大小不是固定的。窗口不会在某个指定时间点关闭而是在一段指定时长的不活动inactivity之后才关闭。|--events--| gap |--events------| gap |--events--| | Session 1| | Session 2 | | Session 3|上图中每段连续事件流构成一个会话一旦事件流中出现超过阈值session gap的静默期当前会话就结束下一个事件开启新会话。典型应用场景把用户行为聚合在一起grouping user behavior together。设想一个用户登录 App、连续点击若干按钮、离开 2 分钟、然后又回来——从行为分析的角度看这仍然是同一个会话。你可以设置一个会话间隙session gap例如 30 分钟不活动Flink 会把该间隙内的所有事件归入同一个会话窗口。会话化Sessionization在行为分析behavioral analytics中非常强大例如一次购物旅程、一次游戏对局、一段网页浏览序列等。在 Flink SQL 中对应的是SESSION表值函数其时间参数为会话间隙例如SESSION(TABLE events, DESCRIPTOR(event_timestamp), INTERVAL 30 MINUTE)。会话窗口的独特之处在于窗口长度动态变化取决于事件的实际到达模式事件之间的间隙大于阈值时窗口才关闭并触发计算结果输出它天然适合以用户/设备为粒度 以活跃期为单位的分析是滚动、滑动窗口无法直接表达的语义。五、从仓库实战看窗口与水印、迟到事件、Upsert 的协同理解三种窗口类型之后必须回到本仓库的真实流水线语境窗口只是定义了把事件归入哪个桶真正驱动结果发布的是水印Watermark。在 10-aggregation-with-tumbling-windows.md 与 aggregation_job.py 中Kafka 源表的关键声明是event_timestamp AS TO_TIMESTAMP_LTZ(tpep_pickup_datetime, 3), WATERMARK FOR event_timestamp AS event_timestamp - INTERVAL 5 SECONDWATERMARK是窗口结果的发布触发器它始终比 Flink 目前看到的最新事件时间落后 5 秒。当水印越过某个窗口的结束时刻Flink 就发布该窗口的聚合结果。这 5 秒就是给迟到者的耐心——那些发生在窗口结束之前、但晚了几秒才到达的事件仍有机会被计入原窗口。三者的分工可以概括为窗口Window 把事件归入哪个桶例如 1 小时水印Watermark 何时发布结果触发器UpsertPRIMARY KEY 兜底安全网如果发布后仍有事件到达用更新修正已发布的结果。这正解释了为什么聚合 sink 表要声明PRIMARY KEY (window_start, PULocationID) NOT ENFORCED滚动窗口按时间切分晚到事件可能使 Flink 重新评估一个已经发布过结果的窗口。有了主键Flink 通过 JDBC 连接器执行 UpsertPostgreSQL 会原地更新已有行而不是插入重复行若使用 append-only 的 sink已发布的窗口无法重开迟到事件将直接丢失。仓库中的 producer_realtime.py 专门用来验证这套机制它约 20% 的概率发送时间戳比当前时间晚 3-10 秒的事件模拟网络延迟运行输出形如on time - PU79 ts14:23:05 on time - PU107 ts14:23:05 LATE (8s) - PU234 ts14:22:58 on time - PU48 ts14:23:06配合 11-late-events-and-upserts.md 中的双终端实验一端跑实时生产者、另一端watch聚合结果可以直观看到随着迟到事件到达较早窗口的num_trips计数会通过 Upsert 持续增大。这一点对窗口类型选型有直接含义无论选择滚动、滑动还是会话窗口都需要同时设计水印容忍多少迟到、sink 的 Upsert 能力如何修正已发布结果以及检查点如何在故障后恢复未关闭的窗口状态。本仓库作业通过env.enable_checkpointing(10 * 1000)每 10 秒快照一次作业状态包括尚未关闭的窗口内容。六、三种窗口的对比与选型建议维度Tumbling 滚动Sliding 滑动Session 会话窗口大小固定固定动态由数据决定重叠无有无事件归属恰好一个窗口多个窗口一个会话关闭条件固定时间点固定时间点多起点超过间隙时长的静默Flink SQL 函数TUMBLE(...)HOP(...)SESSION(...)典型场景按小时计行程、日营收峰值检测、移动平均、加价用户会话化、行为分析计算成本最低高事件复制进多个窗口中等需跟踪间隙状态选型建议可以概括为三条经验法则只需要固定粒度、确定性的统计报表如按小时/按天分组计数、求和→ 用 Tumbling语义最接近批处理成本最低关心任意时间段内的极值或趋势且希望时间窗口平滑移动移动平均、激增检测→ 用 Sliding通过 slide 参数控制平滑程度与开销分析主体是人/设备的一段连续行为边界天然由活跃与静默定义 → 用 Session它把业务上的一次会话直接映射为计算窗口。七、继续深入相关文档与源码索引想在本仓库中继续验证窗口机制可以按以下路径深入窗口类型讲义原文12-understanding-window-types.md滚动窗口聚合完整作业DDL、Watermark、Upsert、提交命令10-aggregation-with-tumbling-windows.md 及源码 aggregation_job.py观察窗口关闭与迟到事件丢弃的 10 秒窗口演示aggregation_job_demo.py迟到事件与 Upsert 行为验证11-late-events-and-upserts.md 与 producer_realtime.py消费起始位置earliest-offset/latest-offset对窗口作业重放的影响09-offsets-earliest-vs-latest.md流处理模块整体结构与作业清单07-streaming/README.md掌握滚动、滑动、会话三种窗口并理解它们与水印、迟到事件、Upsert 的协同关系是构建健壮流处理管线的关键一步——这也是本仓库从Python 消费者逐条处理升级到Flink 一行 SQL 完成窗口聚合的核心价值所在。【免费下载链接】data-engineering-zoomcampData Engineering Zoomcamp is a free 9-week course on building production-ready data pipelines. Join the course here 项目地址: https://gitcode.com/GitHub_Trending/da/data-engineering-zoomcamp创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考