Flink 流处理窗口类型全解析:Tumbling / Sliding / Session 窗口原理与实战(Data Engineering Zoomcamp)
2026/9/12 1:44:11 网站建设 项目流程

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' SECOND

WATERMARK是窗口结果的"发布触发器":它始终比 Flink 目前看到的最新事件时间落后 5 秒。当水印越过某个窗口的结束时刻,Flink 就发布该窗口的聚合结果。这 5 秒就是给"迟到者"的耐心——那些发生在窗口结束之前、但晚了几秒才到达的事件,仍有机会被计入原窗口。

三者的分工可以概括为:

  • 窗口(Window)= 把事件归入哪个桶(例如 1 小时);
  • 水印(Watermark)= 何时发布结果(触发器);
  • Upsert(PRIMARY KEY)= 兜底安全网:如果发布后仍有事件到达,用更新修正已发布的结果。

这正解释了为什么聚合 sink 表要声明PRIMARY KEY (window_start, PULocationID) NOT ENFORCED:滚动窗口按时间切分,晚到事件可能使 Flink 重新评估一个已经发布过结果的窗口。有了主键,Flink 通过 JDBC 连接器执行 Upsert,PostgreSQL 会原地更新已有行而不是插入重复行;若使用 append-only 的 sink,已发布的窗口无法重开,迟到事件将直接丢失。

仓库中的 producer_realtime.py 专门用来验证这套机制:它约 20% 的概率发送时间戳比当前时间晚 3-10 秒的事件(模拟网络延迟),运行输出形如:

on time -> PU=79 ts=14:23:05 on time -> PU=107 ts=14:23:05 LATE (8s) -> PU=234 ts=14:22:58 on time -> PU=48 ts=14: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(...)
典型场景按小时计行程、日营收峰值检测、移动平均、加价用户会话化、行为分析
计算成本最低高(事件复制进多个窗口)中等(需跟踪间隙状态)

选型建议可以概括为三条经验法则:

  1. 只需要固定粒度、确定性的统计报表(如按小时/按天分组计数、求和)→ 用 Tumbling,语义最接近批处理,成本最低;
  2. 关心"任意时间段内的极值或趋势",且希望时间窗口平滑移动(移动平均、激增检测)→ 用 Sliding,通过 slide 参数控制平滑程度与开销;
  3. 分析主体是"人/设备的一段连续行为",边界天然由活跃与静默定义 → 用 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),仅供参考

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

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

立即咨询