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

相关推荐

  • 自营机房和代理机房有什么区别,哪个更稳定可靠

    自营机房和代理机房最根本的区别在于基础设施的归属权与运营主体,自营机房由服务商自主建设、维护并持证运营,提供全程可控的IDC服务;代理机房则通过转售第三方资源,缺乏底层控制力, 选择自营还是代理,直接关系到业务的稳定性、安全边际和售后深度,近年来,随着工信部对IDC行业持证经营的要求逐步收紧,越来越多企业开始关……

    2026-07-26
    0
  • 高防服务器常见的套路猫腻如何辨别,有哪些常见的坑?

    辨别高防服务器猫腻的关键在于核实防御能力真实性、检查机房资质与服务条款,避免被低价宣传和虚假承诺误导,常见猫腻有哪些虚假防御能力不少服务商声称提供单机百G防御,实际可能只有几十G,甚至多个用户共享同一防御带宽,一旦遭遇真实攻击,防护效果远低于宣传值,部分商家还会利用“突发清洗”概念,只在攻击峰值时短暂启用清洗设……

    2026-07-26
    0
  • 选择IDC服务商千万不能忽略哪几点,有哪些注意事项?

    选择IDC服务商,资质和合规性是最容易被忽视却至关重要的环节,一家没有完整资质的服务商,无论价格多低都不值得选择, 近年来,企业数字化转型加速,对IDC服务的需求持续增长,但服务商水平参差不齐,如果你正在挑选IDC服务商,需要从资质、机房、服务、成本等几个核心维度仔细考察,才能避免业务埋雷,资质认证:合规经营的……

    2026-07-26
    0
  • 低价服务器到底为什么不能买,有什么风险?

    低价服务器看似省钱,实则隐藏性能差、不稳定、数据安全无保障、售后缺失甚至跑路等风险,最终可能让你付出更高代价,低价服务器常见的“坑”性能陷阱:超售严重低价服务器的核心套路是超售,一台物理机卖给几十甚至上百个用户,CPU、内存、带宽全面拥挤,你买到的所谓“独享”资源,实际上与邻居争抢,高峰期响应延迟直接拉满,用……

    2026-07-26
    0
  • 新手买站群服务器怎么避坑不被坑,推荐哪家服务商比较好?

    新手买站群服务器,踩坑的根源往往在于只看价格不看资质,要避开风险,必须从持牌经营、机房实地、IP池质量三个维度入手,缺一不可,站群服务器常见的“坑”有哪些IP被“污染”或“墙”了多数新手贪便宜,买到被滥用过的IP段,这类IP发出去的邮件被拒,收录慢,甚至直接被墙,你花时间搭起来的站群,还没开始跑流量就废了,IP……

    2026-07-26
    0

发表回复

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