Flink 窗口聚合(Window Aggregation)完全指南:TVF 语法、GROUPING SETS 与多级聚合实战
Flink 窗口聚合Window Aggregation完全指南TVF 语法、GROUPING SETS 与多级聚合实战【免费下载链接】flink项目地址: https://gitcode.com/gh_mirrors/fli/flink窗口聚合是 Apache Flink SQL 处理无限流数据时最核心的计算模式之一它借助窗口表值函数Windowing TVF把无界流切分为有界的桶再通过GROUP BY对每个桶内的数据进行聚合从而产出按时间维度统计的实时指标。本文基于 Flink 官方文档中的窗口聚合章节完整讲解滚动TUMBLE、滑动HOP、累积CUMULATE、会话SESSION四种窗口的 SQL 写法深入剖析GROUPING SETS/ROLLUP/CUBE多维度分组、window_time时间属性传递与多级窗口聚合同时结合仓库源码揭示窗口聚合在计划器与运行时中的底层实现帮助你在批流两种模式下写出正确、高效的窗口聚合查询。窗口表值函数TVF聚合{{ label Batch }} {{ label Streaming }}窗口聚合是通过GROUP BY子句定义的其特征是分组键中包含窗口表值函数产生的window_start和window_end列。与普通GROUP BY一样窗口聚合对每个组计算并产出一行数据SELECT ... FROM windowed_table -- relation applied windowing TVF GROUP BY window_start, window_end, ...这里windowed_table是在FROM子句中调用窗口表值函数得到的开窗后的表。窗口 TVF 在返回原始列的同时还会附加三列window_start、window_end和window_time后文会分别说明它们的语义与用途。窗口聚合与普通连续表聚合有两个关键差异不产生中间结果普通分组聚合group aggregation会随每条数据的到来持续更新输出而窗口聚合只在窗口结束时一次性输出该窗口的聚合结果因此下游观察到的是一条条追加式的完整窗口统计。自动清理中间状态窗口一旦结束并输出其对应的聚合中间状态即可被清除不会像无界聚合那样长期累积状态这也是窗口聚合在流处理中状态开销可控的根本原因。四种窗口表值函数Flink 支持在TUMBLE、HOP、CUMULATE和SESSION上进行窗口聚合对时间属性的要求如下流模式下窗口表值函数的时间属性字段必须是事件时间或处理时间属性批模式下时间属性字段必须是TIMESTAMP或TIMESTAMP_LTZ类型注意SESSION窗口聚合目前不支持批模式。关于窗口 TVF 的完整语法与参数含offset窗口偏移、命名参数调用方式参见 Windowing TVF。下文以一张带事件时间与 watermark 的Bid出价表为例逐一演示四类窗口聚合。示例数据如下-- tables must have time attribute, e.g. bidtime in this table Flink SQL desc Bid; ----------------------------------------------------------------------------------------- | name | type | null | key | extras | watermark | ----------------------------------------------------------------------------------------- | bidtime | TIMESTAMP(3) *ROWTIME* | true | | | bidtime - INTERVAL 1 SECOND | | price | DECIMAL(10, 2) | true | | | | | item | STRING | true | | | | | supplier_id | STRING | true | | | | ----------------------------------------------------------------------------------------- Flink SQL SELECT * FROM Bid; -------------------------------------------- | bidtime | price | item | supplier_id | -------------------------------------------- | 2020-04-15 08:05 | 4.00 | C | supplier1 | | 2020-04-15 08:07 | 2.00 | A | supplier1 | | 2020-04-15 08:09 | 5.00 | D | supplier2 | | 2020-04-15 08:11 | 3.00 | B | supplier2 | | 2020-04-15 08:13 | 1.00 | E | supplier1 | | 2020-04-15 08:17 | 6.00 | F | supplier2 | --------------------------------------------说明为便于理解窗口行为上表省略了时间戳尾部的 0。在 Flink SQL Client 中TIMESTAMP(3)类型的2020-04-15 08:05实际显示为2020-04-15 08:05:00.000。滚动窗口聚合TUMBLE滚动窗口把数据分配到固定大小、连续且不重叠的时间区间内。例如下面查询以 10 分钟为桶统计每个窗口内的总出价金额-- tumbling window aggregation Flink SQL SELECT window_start, window_end, SUM(price) AS total_price FROM TABLE( TUMBLE(TABLE Bid, DESCRIPTOR(bidtime), INTERVAL 10 MINUTES)) GROUP BY window_start, window_end; ------------------------------------------------- | window_start | window_end | total_price | ------------------------------------------------- | 2020-04-15 08:00 | 2020-04-15 08:10 | 11.00 | | 2020-04-15 08:10 | 2020-04-15 08:20 | 10.00 | -------------------------------------------------TUMBLE(TABLE data, DESCRIPTOR(timecol), size [, offset])中size为窗口时长。08:05、08:07、08:09 三条数据被分入[08:00, 08:10)合计 11.00其余三条被分入[08:10, 08:20)合计 10.00。滑动窗口聚合HOP滑动窗口同时具备窗口大小size与滑动步长slide两个参数当slide size时窗口之间产生重叠一条数据可能同时属于多个窗口。下面的查询使用 5 分钟滑动、10 分钟大小的窗口-- hopping window aggregation Flink SQL SELECT window_start, window_end, SUM(price) AS total_price FROM TABLE( HOP(TABLE Bid, DESCRIPTOR(bidtime), INTERVAL 5 MINUTES, INTERVAL 10 MINUTES)) GROUP BY window_start, window_end; ------------------------------------------------- | window_start | window_end | total_price | ------------------------------------------------- | 2020-04-15 08:00 | 2020-04-15 08:10 | 11.00 | | 2020-04-15 08:05 | 2020-04-15 08:15 | 15.00 | | 2020-04-15 08:10 | 2020-04-15 08:20 | 10.00 | | 2020-04-15 08:15 | 2020-04-15 08:25 | 6.00 | -------------------------------------------------可以看到[08:05, 08:15)的 15.00 是由 08:05、08:07、08:09属于上一个窗口与 08:11、08:13属于下一个窗口共同贡献的这正是滑动窗口重叠分配的结果。HOP的完整签名为HOP(TABLE data, DESCRIPTOR(timecol), slide, size [, offset])。累积窗口聚合CUMULATE累积窗口适用于提前触发的滚动窗口场景比如每日仪表盘从 00:00 开始每分钟绘制一次当日累计 UV10:00 时的值就是从 00:00 到 10:00 的累计 UV。可以把它理解为先按最大窗口大小做滚动窗口再把每个滚动窗口按步长拆分成多个起始时间相同、结束时间递增的子窗口因此累积窗口起始固定、会产生重叠、大小不固定。-- cumulative window aggregation Flink SQL SELECT window_start, window_end, SUM(price) AS total_price FROM TABLE( CUMULATE(TABLE Bid, DESCRIPTOR(bidtime), INTERVAL 2 MINUTES, INTERVAL 10 MINUTES)) GROUP BY window_start, window_end; ------------------------------------------------- | window_start | window_end | total_price | ------------------------------------------------- | 2020-04-15 08:00 | 2020-04-15 08:06 | 4.00 | | 2020-04-15 08:00 | 2020-04-15 08:08 | 6.00 | | 2020-04-15 08:00 | 2020-04-15 08:10 | 11.00 | | 2020-04-15 08:10 | 2020-04-15 08:12 | 3.00 | | 2020-04-15 08:10 | 2020-04-15 08:14 | 4.00 | | 2020-04-15 08:10 | 2020-04-15 08:16 | 4.00 | | 2020-04-15 08:10 | 2020-04-15 08:18 | 10.00 | | 2020-04-15 08:10 | 2020-04-15 08:20 | 10.00 | -------------------------------------------------上例以 2 分钟为步长step、10 分钟为最大窗口size数据进入后会同时被分配到从[08:00, 08:06)直到[08:00, 08:10)的所有以 08:00 为起始的累积窗口。注意size必须是step的整数倍。CUMULATE完整签名为CUMULATE(TABLE data, DESCRIPTOR(timecol), step, size [, offset])。会话窗口聚合SESSION会话窗口没有固定的时间区间其边界由不活动间隙gap决定当超过gap时长没有新数据到达时当前会话窗口关闭后续数据开启新会话。会话窗口天然适用于按用户/供应商等维度切分一次连续交互的场景因此通常搭配PARTITION BY使用。-- session window aggregation with partition keys Flink SQL SELECT window_start, window_end, supplier_id, SUM(price) AS total_price FROM TABLE( SESSION(TABLE Bid PARTITION BY supplier_id, DESCRIPTOR(bidtime), INTERVAL 2 MINUTES)) GROUP BY window_start, window_end, supplier_id; -------------------------------------------------------------- | window_start | window_end | supplier_id | total_price | -------------------------------------------------------------- | 2020-04-15 08:05 | 2020-04-15 08:09 | supplier1 | 6.00 | | 2020-04-15 08:09 | 2020-04-15 08:13 | supplier2 | 8.00 | | 2020-04-15 08:13 | 2020-04-15 08:15 | supplier1 | 1.00 | | 2020-04-15 08:17 | 2020-04-15 08:19 | supplier2 | 6.00 | --------------------------------------------------------------在PARTITION BY supplier_id下会话窗口按供应商分别闭合supplier1在 08:05 和 08:07 有出价间隔 2 分钟以内合并为[08:05, 08:09)合计 6.0008:13 的出价与上一会话间隔超过 2 分钟另开新会话[08:13, 08:15)。若去掉分区键所有数据进入同一个会话维度结果完全不同-- session window aggregation without partition keys Flink SQL SELECT window_start, window_end, SUM(price) AS total_price FROM TABLE( SESSION(TABLE Bid, DESCRIPTOR(bidtime), INTERVAL 2 MINUTES)) GROUP BY window_start, window_end; ------------------------------------------------- | window_start | window_end | total_price | ------------------------------------------------- | 2020-04-15 08:05 | 2020-04-15 08:15 | 15.00 | | 2020-04-15 08:17 | 2020-04-15 08:19 | 6.00 | -------------------------------------------------SESSION完整签名为SESSION(TABLE data [PARTITION BY(keycols, ...)], DESCRIPTOR(timecol), gap)其中keycols决定会话窗口按哪些列分区gap是两个事件被认为属于同一会话的最大时间间隔。GROUPING SETS窗口聚合同样支持GROUPING SETS语法可以在一条 SQL 中表达更复杂的分组操作数据会按每一个指定的 Grouping Sets 分别分组并像普通GROUP BY一样对每组进行聚合。使用窗口聚合 GROUPING SETS时有两条硬性约束GROUP BY子句必须包含window_start和window_end列GROUPING SETS子句中不能包含这两个字段。示例按供应商分组与全量分组同时统计Flink SQL SELECT window_start, window_end, supplier_id, SUM(price) AS total_price FROM TABLE( TUMBLE(TABLE Bid, DESCRIPTOR(bidtime), INTERVAL 10 MINUTES)) GROUP BY window_start, window_end, GROUPING SETS ((supplier_id), ()); -------------------------------------------------------------- | window_start | window_end | supplier_id | total_price | -------------------------------------------------------------- | 2020-04-15 08:00 | 2020-04-15 08:10 | (NULL) | 11.00 | | 2020-04-15 08:00 | 2020-04-15 08:10 | supplier2 | 5.00 | | 2020-04-15 08:00 | 2020-04-15 08:10 | supplier1 | 6.00 | | 2020-04-15 08:10 | 2020-04-15 08:20 | (NULL) | 10.00 | | 2020-04-15 08:10 | 2020-04-15 08:20 | supplier2 | 9.00 | | 2020-04-15 08:10 | 2020-04-15 08:20 | supplier1 | 1.00 | --------------------------------------------------------------GROUPING SETS的每个子列表可以是空的、多列或表达式其解释方式与直接使用GROUP BY相同。其中空 Grouping Sets上例中的()表示把所有行聚合到一个分组下即使没有数据也会输出结果。对于空子列表结果数据中对应的分组/表达式列会用NULL代替——例如上例()对应的supplier_id列显示为(NULL)。ROLLUPROLLUP是 Grouping Sets 的简写语法表示指定表达式及其所有前缀外加空列表。例如ROLLUP (one, two)等价于GROUPING SETS ((one, two), (one), ())。ROLLUP窗口聚合同样要求GROUP BY包含window_start和window_end且ROLLUP子句中不能包含这两个字段。下面查询与上面的GROUPING SETS ((supplier_id), ())例子效果完全相同SELECT window_start, window_end, supplier_id, SUM(price) AS total_price FROM TABLE( TUMBLE(TABLE Bid, DESCRIPTOR(bidtime), INTERVAL 10 MINUTES)) GROUP BY window_start, window_end, ROLLUP (supplier_id);CUBECUBE表示指定列表及其所有可能的子集幂集。同样地CUBE窗口聚合要求GROUP BY包含window_start和window_end且CUBE子句中不能包含这两个字段。下面两个查询完全等效SELECT window_start, window_end, item, supplier_id, SUM(price) AS total_price FROM TABLE( TUMBLE(TABLE Bid, DESCRIPTOR(bidtime), INTERVAL 10 MINUTES)) GROUP BY window_start, window_end, CUBE (supplier_id, item); SELECT window_start, window_end, item, supplier_id, SUM(price) AS total_price FROM TABLE( TUMBLE(TABLE Bid, DESCRIPTOR(bidtime), INTERVAL 10 MINUTES)) GROUP BY window_start, window_end, GROUPING SETS ( (supplier_id, item), (supplier_id ), ( item), ( ) )选取分组窗口的开始和结束时间戳在GROUP BY中纳入window_start与window_end后这两列会作为普通分组列进入结果供你在SELECT中直接选定用来标识每条聚合结果所属窗口的时间边界。多级窗口聚合window_start和window_end是普通的时间戳字段并不是时间属性因此不能在后续操作中直接用于基于时间的操作例如再次开窗、interval join、over 聚合等。为了传递时间属性需要在GROUP BY子句中额外添加window_time列——它是窗口 TVF 产生的三列之一本质是窗口的时间属性其值恒等于window_end - 1ms参见 Windowing TVFs 中的窗口函数说明。把window_time加入GROUP BY后即可被SELECT选定并用于后续基于时间的操作典型场景就是多级窗口聚合和窗口 TopN。下面的例子先按supplier_id做 5 分钟滚动聚合并把window_time作为新的时间属性传出随后在第一步结果上再做 10 分钟滚动聚合-- tumbling 5 minutes for each supplier_id CREATE VIEW window1 AS -- Note: The window start and window end fields of inner Window TVF are optional in the select clause. However, if they appear in the clause, they need to be aliased to prevent name conflicting with the window start and window end of the outer Window TVF. SELECT window_start AS window_5mintumble_start, window_end AS window_5mintumble_end, window_time AS rowtime, SUM(price) AS partial_price FROM TABLE( TUMBLE(TABLE Bid, DESCRIPTOR(bidtime), INTERVAL 5 MINUTES)) GROUP BY supplier_id, window_start, window_end, window_time; -- tumbling 10 minutes on the first window SELECT window_start, window_end, SUM(partial_price) AS total_price FROM TABLE( TUMBLE(TABLE window1, DESCRIPTOR(rowtime), INTERVAL 10 MINUTES)) GROUP BY window_start, window_end;这里有一个容易踩坑的细节SQL 注释中已强调内层窗口 TVF 的window_start、window_end字段在 select 子句中是可选的但如果出现必须为其起别名如window_5mintumble_start以避免与外层窗口 TVF 产生的同名window_start/window_end冲突。分组窗口聚合已过时{{ label Batch }} {{ label Streaming }}{{ hint warning }}警告分组窗口聚合group window aggregation已经过时推荐使用更加强大和高效的窗口表值函数聚合。窗口表值函数聚合相对分组窗口聚合的优势包括包含性能调优中提到的所有性能优化如切片复用、提前聚合等支持标准的GROUPING SETS语法可以在窗口聚合结果上使用窗口 TopN等等。 {{ /hint }}分组窗口聚合同样定义在 SQL 的GROUP BY子句中包含分组窗口函数的GROUP BY查询会对各组分别计算、各产生一个结果行。批处理表和流表上的 SQL 都支持以下分组窗口函数。分组窗口函数Group Window FunctionDescriptionTUMBLE(time_attr, interval)定义一个滚动时间窗口。它把数据分配到连续且不重叠的固定时间区间interval例如一个 5 分钟的滚动窗口以 5 分钟为间隔对数据进行分组。滚动窗口可以被定义在事件时间流 批或处理时间流上。HOP(time_attr, interval, interval)定义一个滑动时间窗口它有窗口大小第二个interval参数和滑动间隔第一个interval参数两个参数。如果滑动间隔小于窗口大小窗口会产生重叠数据可以被分配到多个窗口。例如一个 15 分钟大小、5 分钟滑动间隔的滑动窗口会把每一行分配给 3 个不同的 15 分钟窗口。滑动窗口可以被定义在事件时间流 批或处理时间流上。SESSION(time_attr, interval)定义一个会话时间窗口。会话窗口没有固定的时间区间其边界通过不活动时间interval定义一个会话窗口会在指定时长内没有事件出现时关闭。例如一个 30 分钟间隔的会话窗口收到一条数据时如果此前已经 30 分钟不活动就会开启一个新窗口否则该数据被分配到已存在的窗口中如果 30 分钟之内没有新数据到来窗口就会关闭。会话窗口可以被定义在事件时间流 批或处理时间流上。时间属性流处理模式下分组窗口函数的time_attr必须是有效的处理时间或事件时间属性如何定义参见时间属性文档批处理模式下time_attr必须是TIMESTAMP类型的字段。选取分组窗口开始和结束时间戳分组窗口的开始、结束时间戳以及时间属性可以通过以下辅助函数获取辅助函数描述TUMBLE_START(time_attr, interval)HOP_START(time_attr, interval, interval)SESSION_START(time_attr, interval)返回相应滚动、滑动或会话窗口的下限时间戳inclusive即窗口开始时间。TUMBLE_END(time_attr, interval)HOP_END(time_attr, interval, interval)SESSION_END(time_attr, interval)返回相应滚动、滑动或会话窗口的上限时间戳exclusive即窗口结束时间。注意上限时间戳exclusive不能作为 rowtime attribute 用于后续基于时间的操作例如 interval joins、group window 或 over window aggregations。TUMBLE_ROWTIME(time_attr, interval)HOP_ROWTIME(time_attr, interval, interval)SESSION_ROWTIME(time_attr, interval)返回相应滚动、滑动或会话窗口的上限时间戳inclusive即窗口事件时间或窗口处理时间。返回的值是 rowtime attribute可以用于后续基于时间的操作比如 interval joins、group window 或 over window aggregations。TUMBLE_PROCTIME(time_attr, interval)HOP_PROCTIME(time_attr, interval, interval)SESSION_PROCTIME(time_attr, interval)返回的值是 proctime attribute可以用于后续基于时间的操作比如 interval joins、group window 或 over window aggregations。注意辅助函数的参数必须与GROUP BY子句中的分组窗口函数保持一致。下面是一个在流式表上使用分组窗口 SQL 查询的完整示例含 watermark 定义CREATE TABLE Orders ( user BIGINT, product STRING, amount INT, order_time TIMESTAMP(3), WATERMARK FOR order_time AS order_time - INTERVAL 1 MINUTE ) WITH (...); SELECT user, TUMBLE_START(order_time, INTERVAL 1 DAY) AS wStart, SUM(amount) FROM Orders GROUP BY TUMBLE(order_time, INTERVAL 1 DAY), user该查询以 1 天为滚动窗口、user为分组键通过TUMBLE_START取出每个窗口的开始时间输出每个用户每天的订单金额合计。源码视角窗口聚合在 Flink 内部如何执行窗口聚合的 SQL 能力背后是计划器planner与运行时runtime的紧密配合。以下从当前仓库的源码出发梳理关键实现环节。执行节点StreamExecWindowAggregate在流式计划器中由 TVF 语法翻译而来的窗口聚合对应执行节点 StreamExecWindowAggregate。该类的类注释明确指出它与旧的StreamExecGroupWindowAggregate的区别在于——前者由窗口 TVF 语法翻译而来后者来自已过时的GROUP WINDOW FUNCTION语法且在将来StreamExecGroupWindowAggregate会被移除。从节点结构看StreamExecWindowAggregate内部封装了grouping分组键、aggCalls聚合调用、windowingWindowingStrategy 窗口策略与namedWindowProperties命名窗口属性即window_start/window_end/window_time的产出配置。该节点通过ExecNodeMetadata标注为stream-exec-window-aggregate其consumedOptions包含table.local-time-zone——这印证了文档中时间属性为TIMESTAMP_LTZ时窗口分配受本地时区影响的语义。类中还定义了WINDOW_AGG_MEMORY_RATIO 100与WINDOW_AGGREGATE_TRANSFORMATION window-aggregate分别用于限制窗口聚合的托管内存占比和标记其 Transformation 名称。运行时算子SliceAssigner 与 WindowAggOperatorBuilder窗口 TVF 聚合的运行时构建入口是 WindowAggOperatorBuilder。它的构建式 API 表明窗口聚合算子的核心由窗口分配器assigner与聚合函数aggregate组合而成例如WindowAggOperatorBuilder.builder() .inputType(inputType) .keyTypes(keyFieldTypes) .assigner(SliceAssigners.tumbling(rowtimeIndex, Duration.ofSeconds(5))) .aggregate(genAggsFunction, accTypes) .build();窗口分配器由 SliceAssigners 工厂类创建分别对应四种窗口SliceAssigners.tumbling(...)创建TumblingSliceAssigner滚动窗口切片分配器SliceAssigners.hopping(...)创建HoppingSliceAssigner滑动窗口SliceAssigners.cumulative(...)创建CumulativeSliceAssigner累积窗口会话窗口则通过UnsliceAssigners.session(...)走非切片unslicing路径因为会话窗口边界动态、不适用固定切片。从 SliceAssigners 的工厂方法签名可以看到两个值得注意的实现细节rowtimeIndex输入行中时间字段的下标-1表示基于处理时间shiftTimeZone窗口的时区偏移——当时间属性类型为TIMESTAMP_LTZ时使用用户在TableConfig中配置的时区否则使用 UTC即不偏移。同时TumblingSliceAssigner的构造函数用checkArgument校验了size 0且|offset| size从源码层面印证了窗口参数大小、偏移量的合法性约束。这也解释了文档中反复强调的规则window_start/window_end是普通时间戳字段而window_time才是可继续传递的时间属性——因为在运行时层面切片分配器只负责把数据映射到窗口区间并产出边界列真正被下游当作时间属性使用的是经过专门设计的window_timewindow_end - 1ms。测试与计划验证仓库中为窗口聚合提供了大量可验证依据计划节点级测试flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/nodes/exec/stream/下的WindowAggregateTestPrograms.java、WindowAggregateEventTimeRestoreTest.java等覆盖事件时间/处理时间下的滚动、滑动、累积、会话窗口聚合的 JSON 计划生成与恢复计划资源文件flink-table/flink-table-planner/src/test/resources/restore-tests/stream-exec-window-aggregate_1/目录中包含window-aggregate-tumble-event-time-two-phase、window-aggregate-session-partition-event-time、window-aggregate-hop-event-time-two-phase-with-offset等计划其中two-phase体现了窗口聚合的**两阶段局部聚合 全局聚合**优化with-offset体现了窗口偏移参数的传递运行时算子测试flink-table/flink-table-runtime/src/test/java/org/apache/flink/table/runtime/operators/aggregate/window/WindowAggOperatorTestBase.java直接验证WindowAggOperator在切片/非切片两种处理器路径下的窗口聚合行为。小结与选型建议优先使用窗口 TVF 聚合它是当前及未来的推荐语法天然支持GROUPING SETS/ROLLUP/CUBE、窗口 TopN并继承了性能调优中针对窗口的所有优化手段如切片提前聚合、两阶段聚合旧的分组窗口聚合TUMBLE(time_attr, interval)语法已过时仅建议用于存量作业迁移场景。按业务形态选窗口固定周期统计用TUMBLE需要重叠/滑动统计用HOP需要当日累计、逐步触发用CUMULATE需要按活跃度切分用户会话用SESSION注意其暂不支持批模式与性能调优优化。关注时间属性的正确传递window_start/window_end只是普通字段跨窗口的二次开窗、interval join、over 聚合必须通过GROUP BY中的window_time传递时间属性。批流模式差异流模式下窗口 TVF 的时间字段必须是事件时间或处理时间属性批模式下必须是TIMESTAMP/TIMESTAMP_LTZ字段SESSION窗口仅支持流模式。【免费下载链接】flink项目地址: https://gitcode.com/gh_mirrors/fli/flink创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考