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

发布时间:2026/9/1 7:14:31
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;如果输入顺序为:abcab则产生的结果顺序为:前面的 -代表消息的属性 -代表删除 代表增加.wordcounta1b1c1-a1a2-b1b2注意:FlinkSQL对于这种普通的group by的写法 默认每来一条数据就会输出结果(如果有撤回则会产生多条先产生一条delete消息再产生一条update的消息), 如果上下游有多级的group join等逻辑则会产生较大的数据膨胀, 因此一般建议增加minibatch相关的参数。settable.exec.mini-batch.size100;settable.exec.mini-batch.allow-latency5s;FlinkSQL对于这种普通的group by的状态默认是永久保存的如果group by的key是不断增加的 那随着时间的推移 FlinkSQL保存的状态会越来越多导致作业失败或者心跳超时。对于这种作业 最好设置状态的超时时间。setstate.retention.time.min1d;setstate.retention.time.max2d;对于这种窗口下游的结果是不断更新的因此是需要下游的系统能够支持撤回和更新的 比较常见的系统有 MySQL, Iceberg, Kudu, Elasticsearch, HBase. 如果下游是Kafka我们可以保存为changelog-json格式(Kafka本身虽然不支持删除和更新 但可以通过changelog-json将这种删除和更新的行为保存起来).TIME WINDOWtime window是group window的一种特殊形式增加了时间作为窗口的范围.这里以一个简单的word count来说明:SELECTTUMBLE_START(time,INTERVAL1HOUR),word,COUNT(*)FROMwordsGROUPBYTUMBLE(time,INTERVAL1HOUR),word;原始数据:2021-10-11 00:00a2021-10-11 00:01b2021-10-11 00:01a2021-10-11 01:01b最终结果输出:timewordcount2021-10-11 00:00a22021-10-11 00:00b12021-10-11 01:00b1时间窗口相比普通的窗口:结果只有在窗口结束才会输出结果输出后状态也随之清理.由于结果是在窗口结束才会输出, 产生的数据都是insert类型因此对于下游的系统可以不用支持删除或者更新使用时间窗口字段中需要有时间属性的字段 可以是proctime (使用proctime() 函数生成的计算列)类型也可以是eventTime(当字段是timestamp类型声明watermark之后变为eventTime属性的字段)时间窗口需要使用时间窗口相关的包括: tumble, session, hopOVER WINDOWOver window与上文Group window不同的是Over window中的每一个元素都对应1个窗口每来一条数据都会进行一次窗口计算Over window可以根据数据的行或者时间戳值来确定窗口。Over window既支持event-time也支持processing-time。Over窗口分为两类Rows OVER Window和Range OVER Window。这里以word count为例:SELECTtime,wordCOUNT(amount)OVER(PARTITIONBYwordORDERBYtimeRANGEBETWEENINTERVAL1HOURPRECEDINGANDCURRENTROW)AScountFROMOrders2021-10-11 00:00a2021-10-11 00:01b2021-10-11 00:01a2021-10-11 00:02b结果输出timewordcount2021-10-11 00:00a12021-10-11 00:01b12021-10-11 00:01a22021-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,INTERVAL1DAY),id,LAST_VALUE(item1),LAST_VALUE(item2)FROMsourceTableGROUPBYTUMBLE(time,INTERVAL1DAY),id普通group窗口去重group window last_value(first_value)以下案例是基于普通窗口的去重写法数据每来一条就会输出一次结果下游的结果会不断更新.SELECTid,LAST_VALUE(item1),LAST_VALUE(item2)FROMsourceTableGROUPBYidover window窗口去重over window row_number()SELECTid,item1,item2FROM(SELECT*,ROW_NUMBER()OVER(PARTITIONBYidORDERBYproctimeASC)ASrow_numFROMsourceTable)WHERErow_num1对于普通窗口去重和 over windows的去重我们更建议使用 ROW_NUMBER()来去重 FlinkSQL内部会识别到这种写法然后优化为一个 Last_Row(取最后一条)或者 First_RowPV,UV计算字段iptime类型Stringtimestamp时间窗口计算pv, uv窗口结束之后 进行结果的输出.SELECTTUMBLE_START(time,INTERVAL1HOUR),ip,COUNT(*)ASpv,COUNT(DISTINCT(*))ASuvFROMsourceTableGROUPBYTUMBLE(time,INTERVAL1HOUR),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),ipover window计算pv, uvover 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这种写法会存在撤回)状态不会自动清理需要设置下状态清理时间本文收录于「我的数据空间」技术库——一套可私有化部署的数据平台(数据集成 / 实时计算 / 数据湖 / 湖仓查询 / 智能问数)。