PyFlink Table API 窗口操作实战:Tumble / Slide / Session 窗口示例与源码解析

发布时间:2026/9/25 3:36:59
PyFlink Table API 窗口操作实战:Tumble / Slide / Session 窗口示例与源码解析 大数据流处理批处理数据工程【免费下载链接】flink项目地址https://gitcode.com/gh_mirrors/fli/flink点击查看免费下载导读本文以 Apache Flink 仓库中 flink-python/docs/examples/table/window.rst 关联的 PyFlink 窗口示例为核心完整讲解 Table API 中三种分组窗口Tumble 滚动窗口、Slide 滑动窗口、Session 会话窗口的 Python 实现从环境搭建、时间属性与 Watermark 定义到窗口 DSL 的链式调用、聚合查询与结果输出。读完本文你将掌握用 PyFlink 编写可运行的窗口聚合作业的完整套路并理解窗口在事件时间语义下的触发原理。一、文档与示例概览原文档 window.rst 通过literalinclude引用了三个可直接运行的示例脚本它们位于仓库的 flink-python/pyflink/examples/table/windowing/ 目录窗口类型示例脚本核心 DSLTumble滚动窗口tumble_window.pyTumble.over(size).on(time_field).alias(w)Slide滑动窗口sliding_window.pySlide.over(size).every(slide).on(time_field).alias(w)Session会话窗口session_window.pySession.with_gap(gap).on(time_field).alias(w)此外同一目录还提供了 over_window.pyOver 窗口示例可作为延伸阅读。三个示例共享同一套数据源 Schema/Watermark 定义 print Sink 窗口聚合骨架差异仅在窗口 DSL 一行。下面先讲清楚公共骨架再逐个深入。二、运行环境准备这些示例依赖 PyFlink 的 Table API 模块。在当前仓库中PyFlink 的完整实现位于 flink-python/ 模块Python API 源码在 flink-python/pyflink/ 下Table 相关代码见 flink-python/pyflink/table/。运行示例前需要安装 PyFlink可通过pip install pyflink安装发布版或基于本仓库flink-python目录构建安装准备 Python 环境示例使用了pyflink.common、pyflink.datastream、pyflink.table等模块需要与当前 Flink 版本匹配的 PyFlink 版本执行方式直接在本地执行python tumble_window.py或 sliding/session 脚本即可。示例内部通过StreamTableEnvironment在本地 mini cluster 中运行作业末尾的.wait()会阻塞等待作业完成并打印输出若提交到远程集群需要移除.wait()原示例注释已明确说明。三个示例都以if __name__ __main__:作为入口先配置logging输出到标准输出再调用对应的*_window_demo()函数if __name__ __main__: logging.basicConfig(streamsys.stdout, levellogging.INFO, format%(message)s) tumble_window_demo()三、公共骨架环境、时间属性与 Watermark三个示例前半段完全一致先建立流式表环境env StreamExecutionEnvironment.get_execution_environment() env.set_parallelism(1) t_env StreamTableEnvironment.create(stream_execution_environmentenv)set_parallelism(1)将并行度固定为 1保证集合数据源的处理顺序和输出结果可预期便于演示。3.1 构造带时间戳的集合数据源ds env.from_collection( collection[ (Instant.of_epoch_milli(1000), Alice, 110.1), (Instant.of_epoch_milli(4000), Bob, 30.2), (Instant.of_epoch_milli(3000), Alice, 20.0), (Instant.of_epoch_milli(2000), Bob, 53.1), (Instant.of_epoch_milli(5000), Alice, 13.1), (Instant.of_epoch_milli(3000), Bob, 3.1), (Instant.of_epoch_milli(7000), Bob, 16.1), (Instant.of_epoch_milli(10000), Alice, 20.1) ], type_infoTypes.ROW([Types.INSTANT(), Types.STRING(), Types.FLOAT()]))每条记录三元组为(事件时间, 用户名, 价格)时间用pyflink.common.time.InstantInstant.of_epoch_milli从毫秒时间戳构造数据类型通过Types.ROW声明为INSTANT / STRING / FLOAT注意数据本身是乱序到达的事件时间并非递增这正是为了演示 Watermark 在乱序事件流中的作用。3.2 从 DataStream 转 Table 并定义 Watermarktable t_env.from_data_stream( ds, Schema.new_builder() .column_by_expression(ts, CAST(f0 AS TIMESTAMP(3))) .column(f1, DataTypes.STRING()) .column(f2, DataTypes.FLOAT()) .watermark(ts, ts - INTERVAL 3 SECOND) .build() ).alias(ts, name, price)这段代码是窗口作业的关键前置column_by_expression(ts, CAST(f0 AS TIMESTAMP(3)))把 DataStream 的第 0 列f0INSTANT通过表达式转换为TIMESTAMP(3)毫秒精度并命名为ts作为事件时间列column(f1, DataTypes.STRING())/column(f2, DataTypes.FLOAT())映射剩余两列类型.watermark(ts, ts - INTERVAL 3 SECOND)在ts上声明 Watermark 策略允许最多 3 秒的乱序延迟。按此策略某时刻的 Watermark 已观测到的最大事件时间 − 3 秒最后.alias(ts, name, price)将三列重命名为语义化名称后续窗口表达式统一使用col(ts)、col(name)、col(price)。为什么必须有 Watermark事件时间窗口的关闭与触发依赖 Watermark 推进当 Watermark 越过窗口结束时间严格来说是window_end - 1ms时窗口才允许计算并输出。示例中最大事件时间为 10s因此最终 Watermark 为 7s10s − 3s这决定了哪些窗口会被触发详见第五节的推演。3.3 定义 print 结果表Sinkt_env.create_temporary_table( sink, TableDescriptor.for_connector(print) .schema(Schema.new_builder() .column(name, DataTypes.STRING()) .column(total_price, DataTypes.FLOAT()) .column(w_start, DataTypes.TIMESTAMP_LTZ()) .column(w_end, DataTypes.TIMESTAMP_LTZ()) .build()) .build())使用TableDescriptor.for_connector(print)创建临时表sink输出四列用户名、总价FLOAT、窗口开始时间、窗口结束时间TIMESTAMP_LTZ。print连接器会把每条结果以I(...)形式打印到标准输出。四、三种分组窗口的 DSL 与完整示例窗口分组本质上是一种特殊的group_by见 window.py 中GroupWindow的 docstring。三种窗口的聚合写法完全一致table table.window(窗口定义) .group_by(col(name), col(w)) .select(col(name), col(price).sum, col(w).start, col(w).end)col(w)是窗口别名group_by同时按用户名和窗口分组col(price).sum对窗口内价格求和col(w).start/col(w).end访问窗口的开始/结束时间属性。4.1 Tumble 滚动窗口固定长度、连续不重叠滚动窗口是长度固定、连续、互不重叠的窗口。例如 5 分钟滚动窗口将数据按 5 分钟间隔切分0–5 分钟、5–10 分钟……。完整示例 tumble_window.py 如下import logging import sys from pyflink.common.time import Instant from pyflink.common import Types from pyflink.datastream import StreamExecutionEnvironment from pyflink.table import (DataTypes, TableDescriptor, Schema, StreamTableEnvironment) from pyflink.table.expressions import lit, col from pyflink.table.window import Tumble def tumble_window_demo(): env StreamExecutionEnvironment.get_execution_environment() env.set_parallelism(1) t_env StreamTableEnvironment.create(stream_execution_environmentenv) # define the source with watermark definition ds env.from_collection( collection[ (Instant.of_epoch_milli(1000), Alice, 110.1), (Instant.of_epoch_milli(4000), Bob, 30.2), (Instant.of_epoch_milli(3000), Alice, 20.0), (Instant.of_epoch_milli(2000), Bob, 53.1), (Instant.of_epoch_milli(5000), Alice, 13.1), (Instant.of_epoch_milli(3000), Bob, 3.1), (Instant.of_epoch_milli(7000), Bob, 16.1), (Instant.of_epoch_milli(10000), Alice, 20.1) ], type_infoTypes.ROW([Types.INSTANT(), Types.STRING(), Types.FLOAT()])) table t_env.from_data_stream( ds, Schema.new_builder() .column_by_expression(ts, CAST(f0 AS TIMESTAMP(3))) .column(f1, DataTypes.STRING()) .column(f2, DataTypes.FLOAT()) .watermark(ts, ts - INTERVAL 3 SECOND) .build() ).alias(ts, name, price) # define the sink t_env.create_temporary_table( sink, TableDescriptor.for_connector(print) .schema(Schema.new_builder() .column(name, DataTypes.STRING()) .column(total_price, DataTypes.FLOAT()) .column(w_start, DataTypes.TIMESTAMP_LTZ()) .column(w_end, DataTypes.TIMESTAMP_LTZ()) .build()) .build()) # define the tumble window operation table table.window(Tumble.over(lit(5).seconds).on(col(ts)).alias(w)) \ .group_by(col(name), col(w)) \ .select(col(name), col(price).sum, col(w).start, col(w).end) # submit for execution table.execute_insert(sink) \ .wait() # remove .wait if submitting to a remote cluster if __name__ __main__: logging.basicConfig(streamsys.stdout, levellogging.INFO, format%(message)s) tumble_window_demo()窗口定义Tumble.over(lit(5).seconds).on(col(ts)).alias(w)的含义Tumble.over(lit(5).seconds)窗口大小为 5 秒lit(5).seconds构造 5 秒时间间隔表达式.on(col(ts))指定按事件时间列ts分组.alias(w)为窗口起别名w供group_by与select中的col(w).start/end引用。4.2 Slide 滑动窗口固定大小 滑动步长可重叠滑动窗口具有固定大小 size 和滑动步长 slide。当 slide 小于 size 时窗口彼此重叠一条元素可能同时属于多个窗口。例如 size15 分钟、slide5 分钟时每 5 分钟评估一次 15 分钟的数据每条数据落入 3 个连续窗口。完整示例 sliding_window.py 与滚动窗口的差异仅在窗口定义一行# define the sliding window operation table table.window(Slide.over(lit(5).seconds).every(lit(2).seconds).on(col(ts)).alias(w))\ .group_by(col(name), col(w)) \ .select(col(name), col(price).sum, col(w).start, col(w).end)窗口定义Slide.over(lit(5).seconds).every(lit(2).seconds).on(col(ts)).alias(w)的含义Slide.over(lit(5).seconds)窗口大小 5 秒.every(lit(2).seconds)每 2 秒滑动一次即每 2 秒启动一个新窗口因此窗口会重叠.on(col(ts)).alias(w)按事件时间列分组并起别名。4.3 Session 会话窗口以不活动间隔为界会话窗口以不活动间隔gap为边界窗口在超过 gap 时间没有新事件到达时闭合反之若新事件在 gap 内到达则并入当前会话并顺延其结束时间。这非常适合一次用户会话、一次点击行为序列等场景。完整示例 session_window.py 与前述差异数据源略有调整(Instant.of_epoch_milli(8000), Bob, 16.1)替换了滚动/滑动示例中的 5000ms 记录窗口定义使用Session.with_gap(...)# define the session window operation table table.window(Session.with_gap(lit(5).seconds).on(col(ts)).alias(w)) \ .group_by(col(name), col(w)) \ .select(col(name), col(price).sum, col(w).start, col(w).end)Session.with_gap(lit(5).seconds)表示 gap 为 5 秒——任何两条事件时间相差超过 5 秒的记录会被分到不同会话窗口。五、事件时间语义下的窗口触发推演结合 Watermark 策略ts - INTERVAL 3 SECOND最大延迟 3 秒与单并行度集合源可以推演三个示例的实际触发行为基于 Flink 事件时间窗口标准触发语义Watermark 达到window_end - 1ms时窗口触发。5.1 数据归属将 8 条记录按窗口区间归类tumble/sliding 数据事件时间(ms)用户价格Tumble[0,5)Slide[0,5)Slide[2,7)1000Alice110.1✓✓2000Bob53.1✓✓✓3000Alice20.0✓✓✓3000Bob3.1✓✓✓4000Bob30.2✓✓✓5000Alice13.1✓7000Bob16.1✓10000Alice20.15.2 Tumble5 秒滚动窗口触发结果最终 Watermark 10000 − 3000 7000ms。窗口[0, 5)window_end - 1 4999Watermark 7000 ≥ 4999 →触发。Alice 合计110.1 20.0 130.1Bob 合计53.1 3.1 30.2 86.4窗口[5, 10)window_end - 1 9999Watermark 7000 9999 →不触发延迟事件可能落入该窗口但 Watermark 尚未推进到 10s。因此滚动窗口示例预期的核心输出为 Alice130.1 与 Bob86.4对应[0,5)窗口同时[5,10)窗口因 Watermark 不足而保持未触发——这直观演示了Watermark 决定窗口何时关闭这一核心机制。5.3 Slide5 秒大小、2 秒滑动触发结果窗口起点依次为 0、2、4、6……秒窗口[0, 5)Watermark 7000 ≥ 4999 → 触发Alice130.1Bob86.4窗口[2, 7)window_end - 1 6999Watermark 7000 ≥ 6999 → 触发Alice 为20.0 13.1 33.1Bob 为53.1 3.1 30.2 16.1 102.5窗口[4, 9)需要 Watermark ≥ 8999 → 不触发。可见同一条事件如 3000ms 的 Alice 记录同时参与了[0,5)与[2,7)两个窗口的聚合这正是滑动窗口重叠性的体现。5.4 Session5 秒 gap的行为特征会话窗口数据Bob 的记录在 8000ms按事件时间排布Alice1000、3000相差 2000 5000同一会话10000 与上一事件 3000 相差 7000 5000严格按 gap 语义应开启新会话。但需要注意最终 Watermark 7000ms而会话闭合要求 Watermark 越过最后事件时间 gap。本地运行该示例时可以观察到会话窗口的输出由 Watermark 推进决定——若 Watermark 未能越过会话结束边界会话窗口将保持打开等待后续数据。这提醒我们会话窗口的输出时机高度依赖 Watermark 的推进速度在真实数据源如 Kafka上会话窗口的闭合会随 Watermark 持续前移而自然发生。说明上述推演基于集合源在单并行度下顺序处理的假设结合事件时间窗口触发语义得出用于理解机制实际输出以运行环境为准。三个示例的核心价值在于演示窗口 API 的完整用法输出行本身会随数据与 Watermark 变化。六、源码级解析窗口 DSL 的链式实现窗口类的 Python 实现在 flink-python/pyflink/table/window.py它通过 Py4J 网关把 Python 表达式翻译到 Java 侧的org.apache.flink.table.expressions.Tumble / Slide / Session / Over。6.1 链式 API 的部分定义设计每种窗口都按 Builder 模式拆分为多个中间状态类逐步补全窗口定义窗口链式步骤类TumbleTumble.over(size)→TumbleWithSize→.on(time)→TumbleWithSizeOnTime→.alias(name)→GroupWindowSlideSlide.over(size)→SlideWithSize→.every(slide)→SlideWithSizeAndSlide→.on(time)→SlideWithSizeAndSlideOnTime→.alias(name)→GroupWindowSessionSession.with_gap(gap)→SessionWithGap→.on(time)→SessionWithGapOnTime→.alias(name)→GroupWindowOverOver.partition_by(...)/Over.order_by(...)→ 中间类 →.preceding(...)/.following(...)→.alias(name)→OverWindow例如Tumble.over的实现window.pyclassmethod def over(cls, size: Expression) - TumbleWithSize: Creates a tumbling window. Tumbling windows are fixed-size, consecutive, non-overlapping windows of a specified fixed length. return TumbleWithSize(get_gateway().jvm.Tumble.over(_get_java_expression(size)))get_gateway().jvm.Tumble.over(...)直接把 Python 表达式对象转换为 Java 表达式并调用 Java 侧工厂方法_get_java_expression负责翻译。随后.on(...)指定时间属性window.py.alias(...)通过get_method(self._java_window, as)(alias)调用 Java 的as()方法返回GroupWindowwindow.py。源码 docstring 还明确了适用边界window.py流式表可按事件时间event-time或处理时间processing-time属性分组批式表可按 timestamp 或 long 类型字段分组分组窗口为基于时间的 groupBy 提供快捷方式。6.2 滑动窗口的重叠语义Slide.every(slide)的 docstringwindow.py明确说明slide 决定窗口的启动间隔当 slide 小于窗口 size 时窗口重叠一条记录可贡献给多个窗口例如 size15 分钟、slide3 分钟时每 3 分钟对 15 分钟数据做一次分组每条记录属于 5 个窗口。6.3 Over 窗口延伸同一目录的 over_window.py 演示了 Over 窗口类似于 SQL 的 OVER 聚合对每行在其相邻行范围内做聚合仅支持流式表。其窗口定义为table table.over_window( Over.partition_by(col(name)) .order_by(col(ts)) .preceding(row_interval(2)) .following(CURRENT_ROW) .alias(w)) \ .select(col(name), col(price).max.over(col(w)))partition_by(col(name))按用户名分区order_by(col(ts))按事件时间排序Over 窗口第二参数必须是时间属性preceding(row_interval(2))上边界为当前行之前的 2 行行数间隔following(CURRENT_ROW)下边界为当前行聚合col(price).max.over(col(w))计算每个分区内截至当前行的滑动最大值。与之配套的偏移常量定义在 flink-python/pyflink/table/expressions.pyUNBOUNDED_ROW无界行数上边界、UNBOUNDED_RANGE无界时间范围上边界、CURRENT_ROW当前行、CURRENT_RANGE当前排序键所在范围。七、测试验证窗口查询计划的关键字断言仓库测试 flink-python/pyflink/table/tests/test_window.py 直接验证了窗口 DSL 生成的查询计划字符串是理解窗口语义的可靠依据测试方法断言的关键内容test_tumble_windowTumbleWindow(field: [a], size: [2])test_slide_windowSlideWindow(field: [a], slide: [1000], size: [2000])test_session_windowSessionWindow(field: [a], gap: [1000])test_over_window断言 Over 窗口第二参数必须为时间属性Second argument in OVER window must be a TIME ATTRIBUTE其中test_slide_window验证了 slide1000ms、size2000ms 的配置被正确编码进SlideWindow查询操作批式环境下的test_tumble_window/test_slide_window/test_session_window分别调用Tumble.over(row_interval(2))、Slide.over(lit(2).seconds).every(lit(1).seconds)、Session.with_gap(lit(1).seconds)确认了批表也支持按时间间隔定义的分组窗口。八、实战要点与常见误区事件时间列与 Watermark 缺一不可窗口按事件时间聚合前必须先在 Schema 中声明时间属性column_by_expressionwatermark否则窗口作业无法正确处理乱序数据与延迟结果col(w).start / .end是窗口属性访问只有在.alias(w)之后select中才能访问窗口起止时间输出列类型应使用TIMESTAMP_LTZ承载窗口边界本地执行用.wait()远程提交需移除示例末尾注释明确提示提交到远程集群时应删除.wait()避免客户端阻塞并行度影响输出顺序示例通过env.set_parallelism(1)保证单并行度下结果可预期真实作业请按吞吐与负载设置并行度乱序与延迟Watermark 的延迟值如 3 秒决定了窗口等待迟到数据的容忍度过小会丢弃迟到数据过大会延迟窗口输出会话窗口的闭合依赖 Watermark 推进gap 只是不活动间隔的上限语义实际触发仍受 Watermark 控制参见第五节推演。结语本文以 window.rst 关联的三个示例为骨架完整还原了 PyFlink Table API 中 Tumble / Slide / Session 三种分组窗口的编写方式并结合 window.py 的链式实现、test_window.py 的查询计划断言以及事件时间 Watermark 的触发推演把照着写提升为理解着写。进一步可阅读同目录的 over_window.py 与 expressions.py 中UNBOUNDED_ROW、CURRENT_ROW等常量掌握 Over 窗口的进阶用法。赞分享大数据流处理批处理数据工程【免费下载链接】flink项目地址https://gitcode.com/gh_mirrors/fli/flink点击查看免费下载相关推荐PyFlink Table Window 窗口 API 完全指南Tumble、Slide、Session 与 Over 窗口PyFlink Table Window 窗口 API 完全指南Tumble、Slide、Session 与 Over 窗口 窗口Window是流式数据处大数据流处理批处理数据工程Flink SQL 窗口表值函数Windowing TVF完全指南TUMBLE、HOP、CUMULATE 与 SESSION 实战Flink SQL 窗口表值函数Windowing TVF完全指南TUMBLE、HOP、CUMULATE 与 SESSION 实战 窗口Window是大数据流处理批处理数据工程Reactor Core 窗口操作符完全指南时间窗口与计数窗口实战解析Reactor Core 窗口操作符完全指南时间窗口与计数窗口实战解析 Reactor Core 作为 JVM 上最强大的响应式编程框架其窗口操作符是处理实后端异步编程上一篇Midway Swagger多文档管理终极指南轻松实现API版本控制下一篇Adapt Intent Parser架构设计解析理解引擎、解析器与标记器的协作机制创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考