news 2026/9/1 7:14:20

FlinkSQL 窗口使用:三类窗口的语义差异与去重、PV/UV 的三种写法

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
FlinkSQL 窗口使用:三类窗口的语义差异与去重、PV/UV 的三种写法

「我的数据空间」实时计算实践笔记 · Flink SQL 系列

平台已经支持编写FlinkSQL作业, 包括Streaming, Batch。
本文主要介绍一下FlinkSQL比较常见的窗口计算 包括over window, group window (普通窗口, 时间窗口).

基本概念介绍:

GROUP WINDOW(普通窗口)

group window比较容易理解, 按固定的字段进行分组, 通过聚合函数(sum, min, max, count)等函数进行计算。和batchSQL不同的是, FlinkSQL产生的结果是不断更新的, 它采用了一种回撤机制, 如果SQL中包含多级的group by操作, 每一层都会将结果不断更新并传递给下游,最终也会将传递到结果表. 通过不断的回撤和更新, 可以保证和batchSQL的最终结果一致.
这里以一个简单的word count来说明:

-- words表只有一个字段, 每一行是不同的单词SELECTword,COUNT(*)AScntFROMwordsGROUPBYword;

如果输入顺序为:

a
b
c
a
b

则产生的结果顺序为:
前面的+, -代表消息的属性, -代表删除, + 代表增加.

wordcount
+a1
+b1
+c1
-a1
+a2
-b1
+b2

注意:

  • FlinkSQL对于这种普通的group by的写法, 默认每来一条数据,就会输出结果(如果有撤回,则会产生多条,先产生一条delete消息,再产生一条update的消息), 如果上下游有多级的group , join等逻辑,则会产生较大的数据膨胀, 因此一般建议增加minibatch相关的参数。
settable.exec.mini-batch.size=100;settable.exec.mini-batch.allow-latency=5s;
  • FlinkSQL对于这种普通的group by的状态默认是永久保存的,如果group by的key是不断增加的, 那随着时间的推移, FlinkSQL保存的状态会越来越多,导致作业失败或者心跳超时。对于这种作业, 最好设置状态的超时时间。
setstate.retention.time.min=1d;setstate.retention.time.max=2d;
  • 对于这种窗口,下游的结果是不断更新的,因此是需要下游的系统能够支持撤回和更新的, 比较常见的系统有 MySQL, Iceberg, Kudu, Elasticsearch, HBase. 如果下游是Kafka,我们可以保存为changelog-json格式(Kafka本身虽然不支持删除和更新, 但可以通过changelog-json将这种删除和更新的行为保存起来).

TIME WINDOW

time window是group window的一种特殊形式,增加了时间作为窗口的范围.
这里以一个简单的word count来说明:

SELECTTUMBLE_START(`time`,INTERVAL'1'HOUR),word,COUNT(*)FROMwordsGROUPBYTUMBLE(`time`,INTERVAL'1'HOUR),word;

原始数据:

2021-10-11 00:00a
2021-10-11 00:01b
2021-10-11 00:01a
2021-10-11 01:01b

最终结果输出:

timewordcount
+2021-10-11 00:00a2
+2021-10-11 00:00b1
+2021-10-11 01:00b1

时间窗口相比普通的窗口:

  • 结果只有在窗口结束才会输出
  • 结果输出后,状态也随之清理.
  • 由于结果是在窗口结束才会输出, 产生的数据都是insert类型,因此对于下游的系统可以不用支持删除或者更新
  • 使用时间窗口,字段中需要有时间属性的字段, 可以是proctime (使用proctime() 函数生成的计算列)类型,也可以是eventTime(当字段是timestamp类型,声明watermark之后变为eventTime属性的字段)
  • 时间窗口需要使用时间窗口相关的包括: tumble, session, hop

OVER WINDOW

Over window与上文Group window不同的是,Over window中的每一个元素都对应1个窗口,每来一条数据都会进行一次窗口计算,Over window可以根据数据的行或者时间戳值来确定窗口。Over window既支持event-time,也支持processing-time。Over窗口分为两类,Rows OVER Window和Range OVER Window。
这里以word count为例:

SELECT`time`,wordCOUNT(amount)OVER(PARTITIONBYwordORDERBY`time`RANGEBETWEENINTERVAL'1'HOURPRECEDINGANDCURRENTROW)AS`count`FROMOrders
2021-10-11 00:00a
2021-10-11 00:01b
2021-10-11 00:01a
2021-10-11 00:02b

结果输出

timewordcount
+2021-10-11 00:00a1
+2021-10-11 00:01b1
+2021-10-11 00:01a2
+2021-10-11 00:01b2

这里只做简单介绍, 详细的内容可以参考Flink官方文档: https://nightlies.apache.org/flink/flink-docs-master/zh/docs/dev/table/sql/queries/window-agg/

使用案例

去重

source表结构:

字段名idtimeitem1item2
类型longtimestampstringstring

时间窗口去重

tumble_window + last_value(first_value)
以下案例是基于时间窗口的去重写法, 最终结果在时间窗口结束后输出. 该案例是使用1天的时间窗口, 所以最终结果会在凌晨进行输出.

SELECTTUMBLE_START(`time`,INTERVAL'1'DAY),id,LAST_VALUE(item1),LAST_VALUE(item2)FROMsourceTableGROUPBYTUMBLE(`time`,INTERVAL'1'DAY),id

普通group窗口去重

group window + last_value(first_value)
以下案例是基于普通窗口的去重写法,数据每来一条就会输出一次结果,下游的结果会不断更新.

SELECTid,LAST_VALUE(item1),LAST_VALUE(item2)FROMsourceTableGROUPBYid

over window窗口去重

over window + row_number()

SELECTid,item1,item2FROM(SELECT*,ROW_NUMBER()OVER(PARTITIONBYidORDERBYproctimeASC)ASrow_numFROMsourceTable)WHERErow_num=1

对于普通窗口去重,和 over windows的去重,我们更建议使用 ROW_NUMBER()来去重, FlinkSQL内部会识别到这种写法,然后优化为一个 Last_Row(取最后一条)或者 First_Row

PV,UV计算

字段iptime
类型Stringtimestamp

时间窗口计算pv, uv

窗口结束之后, 进行结果的输出.

SELECTTUMBLE_START(`time`,INTERVAL'1'HOUR),ip,COUNT(*)ASpv,COUNT(DISTINCT(*))ASuvFROMsourceTableGROUPBYTUMBLE(`time`,INTERVAL'1'HOUR),ip

普通窗口计算pv, uv

每来一条数据就会输出一次结果.

SELECTDATE_FORMAT(`time`,'yyyy-MM-dd hh:00:00'),ip,COUNT(*)aspv,COUNT(DISTINCT(*))ASuvFROMsourceTableGROUPBYDATE_FORMAT(`time`,'yyyy-MM-dd hh:00:00'),ip

over window计算pv, uv

over window一般用来统计递增的数据, 最终可以输出一个递增的结果.

-- 使用over window计算每条数据到来时累计的数据CREATEVIEWpv_uv_per_10minASSELECTMAX(SUBSTR(DATE_FORMAT(ts,'HH:mm'),1,4)||'0')OVERwAStime_str,COUNT(ip)OVERwASpv,COUNT(DISTINCTip)OVERwASuvFROMsourceTable WINDOW wAS(ORDERBYproctimeROWSBETWEENUNBOUNDEDPRECEDINGANDCURRENTROW);-- 使用groupBy过滤出10分钟内的最大值SELECTtime_str,MAX(pv),MAX(uv)FROMpv_uv_per_10minGROUPBYtime_str;

总结

结果输出最终结果个数结果类型状态
时间窗口窗口结束后输出, 结果有一个窗口的延迟每个窗口只产生一条结果只会产生append 数据,不会存在撤回状态在窗口结束后清理
普通窗口每来一条数据就进行输出, 实时输出,实时更新,不断修正最终结果每个窗口只产生一条结果(将不断更新和删除的数据合并,最终只会有一条结果)会有删除和更新操作状态默认不会清理,需要自己设置状态清理时间
over window每来一条数据就会输出, 实时产生新结果每个窗口只产生一条结果(over window每条数据都会开一个新窗口,因此相比普通窗口最终结果会多一些)产生append 消息, 不会存在撤回 (topN, top1这种写法会存在撤回)状态不会自动清理,需要设置下状态清理时间

本文收录于「我的数据空间」技术库——一套可私有化部署的数据平台(数据集成 / 实时计算 / 数据湖 / 湖仓查询 / 智能问数)。

版权声明: 本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若内容造成侵权/违法违规/事实不符,请联系邮箱:809451989@qq.com进行投诉反馈,一经查实,立即删除!
网站建设 2026/9/1 7:11:26

基于SpringBoot的民间艺术传承管理系统(源码+文档+部署+讲解)

温馨提示:本人主页置顶文章(点我)开头有 CSDN 平台官方提供的学长联系方式的名片! 温馨提示:本人主页置顶文章(点我)开头有 CSDN 平台官方提供的学长联系方式的名片! 温馨提示:本人主页置顶文章(点我)开头有 CSDN 平台…

作者头像 李华
网站建设 2026/9/1 7:10:03

数据中心IBMS精细化改造路径

在数据中心基础设施日趋高密度化、集约化的当下,传统运维模式的局限性愈发凸显。长期以来,机房运维依赖二维组态、静态报表与纸质图纸,数据参数与物理空间相互割裂,设备状态、空间容量、能耗分布等关键信息无法直观联动&#xff0…

作者头像 李华
网站建设 2026/9/1 7:09:32

四年级孩子学C++需要哪些基础

结合小学生CSP-J备赛、适配四年级孩子的C入门规划相关背景,四年级孩子学C不需要提前掌握复杂的编程知识,只需要具备几类基础能力就可以顺利入门,完全适配低龄孩子的认知节奏。 一、必备的基础能力 1、基础数学能力‌: 熟练掌握加…

作者头像 李华
网站建设 2026/9/1 7:09:00

基于SpringBoot的大学生志愿服务管理系统(源码+lw+部署文档+讲解等)

联系博主 温馨提示:本人主页置顶文章(点我)开头有 CSDN 平台官方提供的学长联系方式的名片! 温馨提示:本人主页置顶文章(点我)开头有 CSDN 平台官方提供的学长联系方式的名片! 温馨提示:本人主页置顶文章(点我)开头有 …

作者头像 李华
网站建设 2026/9/1 7:06:40

用Python构建虚拟主播直播带货控制中枢:架构设计与实战踩坑总结

简介:面向直播间带货、语音助手、数字导游等场景,这套Python虚拟数字人控制中枢可对接UE4等引擎形象,内置WebSocket协议与UE4对接示例,也能独立作为语音助理。资源包共136个文件、10.5MB,以Python脚本、XML/conf配置、…

作者头像 李华
网站建设 2026/9/1 7:03:45

数据库索引优化与慢查询分析实战:原型怎样变成可用功能

数据库索引优化与慢查询分析实战:原型怎样变成可用功能演示 Demo 的陷阱:AI 生成的“完美索引”在生产环境引发写放大 在实验室里测试 AI 数据库 Agent 时,演示效果好得令人惊讶。只需将 slow_query_log 输入给 Agent,它就能迅速给…

作者头像 李华