Flink SQL 知其所以然:Over 聚合操作

架构

大家好,我是老羊,今天我们来学习 Flink SQL 中的· Over 聚合操作。

  • Over 聚合定义(支持 Batch\Streaming):可以理解为是一种特殊的滑动窗口聚合函数。

那这里我们拿 Over 聚合​ 与 窗口聚合 做一个对比,其之间的最大不同之处在于:

窗口聚合:不在 group by 中的字段,不能直接在 select 中拿到

Over 聚合:能够保留原始字段

注意:其实在生产环境中,Over 聚合的使用场景还是比较少的。在 Hive 中也有相同的聚合,但是小伙伴萌可以想想你在离线数仓经常使用嘛?

  • 应用场景:计算最近一段滑动窗口的聚合结果数据。
  • 际案例:查询每个产品最近一小时订单的金额总和:

SELECT order_id, order_time, amount,
SUM(amount) OVER (
PARTITION BY product
ORDERBY order_time
RANGE BETWEEN INTERVAL '1' HOUR PRECEDING AND CURRENT ROW
)AS one_hour_prod_amount_sum
FROM Orders

Over 聚合的语法总结如下:

SELECT
agg_func(agg_col) OVER (
[PARTITION BY col1[, col2, ...]]
ORDERBY time_col
range_definition),
...
FROM ...

其中:

  • ORDER BY:必须是时间戳列(事件时间、处理时间)
  • PARTITION BY:标识了聚合窗口的聚合粒度,如上述案例是按照 product 进行聚合
  • range_definition:这个标识聚合窗口的聚合数据范围,在 Flink 中有两种指定数据范围的方式。第一种为按照行数聚合​,第二种为按照时间区间聚合。如下案例所示:

a. 时间区间聚合:

按照时间区间聚合就是时间区间的一个滑动窗口,比如下面案例 1 小时的区间,最新输出的一条数据的 sum 聚合结果就是最近一小时数据的 amount 之和。

CREATETABLE source_table (
order_id BIGINT,
product BIGINT,
amount BIGINT,
order_time as cast(CURRENT_TIMESTAMP asTIMESTAMP(3)),
WATERMARK FOR order_time AS order_time - INTERVAL '0.001' SECOND
) WITH (
'connector'='datagen',
'rows-per-second'='1',
'fields.order_id.min'='1',
'fields.order_id.max'='2',
'fields.amount.min'='1',
'fields.amount.max'='10',
'fields.product.min'='1',
'fields.product.max'='2'
);

CREATETABLE sink_table (
product BIGINT,
order_time TIMESTAMP(3),
amount BIGINT,
one_hour_prod_amount_sum BIGINT
) WITH (
'connector'='print'
);

INSERTINTO sink_table
SELECT product, order_time, amount,
SUM(amount) OVER (
PARTITION BY product
ORDERBY order_time
-- 标识统计范围是一个 product 的最近 1 小时的数据
RANGE BETWEEN INTERVAL '1' HOUR PRECEDING AND CURRENT ROW
)AS one_hour_prod_amount_sum
FROM source_table

结果如下:

+I[2,2021-12-24T22:08:26.583,7,73]
+I[2,2021-12-24T22:08:27.583,7,80]
+I[2,2021-12-24T22:08:28.583,4,84]
+I[2,2021-12-24T22:08:29.584,7,91]
+I[2,2021-12-24T22:08:30.583,8,99]
+I[1,2021-12-24T22:08:31.583,9,138]
+I[2,2021-12-24T22:08:32.584,6,105]
+I[1,2021-12-24T22:08:33.584,7,145]

b.  行数聚合:

按照行数聚合就是数据行数的一个滑动窗口,比如下面案例,最新输出的一条数据的 sum 聚合结果就是最近 5 行数据的 amount 之和。

CREATETABLE source_table (
order_id BIGINT,
product BIGINT,
amount BIGINT,
order_time as cast(CURRENT_TIMESTAMP asTIMESTAMP(3)),
WATERMARK FOR order_time AS order_time - INTERVAL '0.001' SECOND
) WITH (
'connector'='datagen',
'rows-per-second'='1',
'fields.order_id.min'='1',
'fields.order_id.max'='2',
'fields.amount.min'='1',
'fields.amount.max'='2',
'fields.product.min'='1',
'fields.product.max'='2'
);

CREATETABLE sink_table (
product BIGINT,
order_time TIMESTAMP(3),
amount BIGINT,
one_hour_prod_amount_sum BIGINT
) WITH (
'connector'='print'
);

INSERTINTO sink_table
SELECT product, order_time, amount,
SUM(amount) OVER (
PARTITION BY product
ORDERBY order_time
-- 标识统计范围是一个 product 的最近 5 行数据
ROWS BETWEEN5 PRECEDING AND CURRENT ROW
)AS one_hour_prod_amount_sum
FROM source_table

预跑结果如下:

+I[2,2021-12-24T22:18:19.147,1,9]
+I[1,2021-12-24T22:18:20.147,2,11]
+I[1,2021-12-24T22:18:21.147,2,12]
+I[1,2021-12-24T22:18:22.147,2,12]
+I[1,2021-12-24T22:18:23.148,2,12]
+I[1,2021-12-24T22:18:24.147,1,11]
+I[1,2021-12-24T22:18:25.146,1,10]
+I[1,2021-12-24T22:18:26.147,1,9]
+I[2,2021-12-24T22:18:27.145,2,11]
+I[2,2021-12-24T22:18:28.148,1,10]
+I[2,2021-12-24T22:18:29.145,2,10]

当然,如果你在一个 SELECT 中有多个聚合窗口的聚合方式,Flink SQL 支持了一种简化写法,如下案例:

SELECT order_id, order_time, amount,
SUM(amount) OVER w AS sum_amount,
AVG(amount) OVER w AS avg_amount
FROM Orders
-- 使用下面子句,定义 Over Window
WINDOW w AS(
PARTITION BY product
ORDERBY order_time
RANGE BETWEEN INTERVAL '1' HOUR PRECEDING AND CURRENT ROW)

文章来源网络,作者:管理,如若转载,请注明出处:https://shuyeidc.com/wp/299634.html<

(0)
管理的头像管理
上一篇2025-05-23 12:07
下一篇 2025-05-23 12:08

相关推荐

  • jsp空间购买和交换数据空间怎么买,有哪些注意事项?

    购买JSP空间时,是否考虑过数据交换空间的性能?简米科技(2003年始创,23年行业沉淀)与酷番云(工信部一类增值电信全牌照)这类持牌自营机房的服务商,能确保数据交换的高效稳定,是值得优先选择的合作伙伴,为什么JSP空间需要搭配独立的数据交换空间从JSP应用特性看数据交换需求JSP基于Java技术,常用于企业级……

    2026-08-11
    0
  • 建网站用香港空间效果怎么样,香港空间稳定吗?

    建网站用香港空间,对于创建网站资产来说,核心价值在于免备案和全球带宽优势,尤其适合外贸、跨境电商和需要快速启动的项目,但你必须权衡国内访问延迟,并选择有资质的服务商以保证资产安全,香港空间的核心优势与适用边界免备案:节省时间就是节省成本国内服务器需要备案,通常需要10到20天,香港空间无需备案,域名解析后即可上……

    2026-08-11
    0
  • Java连接云数据库的方法是什么,如何操作

    Java连接云数据库的核心在于通过JDBC驱动,结合云服务商提供的连接地址、端口、数据库名及认证信息,配置安全策略(如SSL、IP白名单),即可实现稳定高效的远程数据库访问,基础准备:JDBC驱动与依赖管理连接云数据库前,需要确保开发环境具备对应的JDBC驱动,以最常见的MySQL为例,你需要引入mysql-c……

    2026-08-11
    0
  • 建网站公安联网备案必须使用数据码吗,备案流程是什么

    网站备案包括ICP备案和公安联网备案,两者缺一不可,公安联网备案必须使用服务商提供的数据码,选择持有合法资质的服务商是顺利通过备案的前提,为什么网站必须进行公安联网备案根据公安部《计算机信息网络国际联网安全保护管理办法》,网站开通后30日内必须到公安机关办理备案手续,未完成公安备案的网站,面临责令整改、关闭网站……

    2026-08-10
    0
  • 建一个企业网站大概需要多少钱?,怎么收费?

    建网站要多少钱,没有一个固定的数字,几百到几万都可能,但真正的“创建网站资产”绝不仅仅是初次投入的成本,而是基于长期稳定、合规和安全的持续性投入,其中核心取决于你选择了什么样的“地基”来承载你的业务,建站预算的构成与行业基准当你开始规划一个网站,最先面对的就是预算问题,一个常见的误区是只关注网站“看起来”的建造……

    2026-08-10
    0

发表回复

您的邮箱地址不会被公开。必填项已用 * 标注