「我的数据空间」实时计算实践笔记 · 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 |
则产生的结果顺序为:
前面的+, -代表消息的属性, -代表删除, + 代表增加.
| word | count | |
|---|---|---|
| + | a | 1 |
| + | b | 1 |
| + | c | 1 |
| - | a | 1 |
| + | a | 2 |
| - | b | 1 |
| + | b | 2 |
注意:
- 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:00 | a |
|---|---|
| 2021-10-11 00:01 | b |
| 2021-10-11 00:01 | a |
| 2021-10-11 01:01 | b |
最终结果输出:
| time | word | count | |
|---|---|---|---|
| + | 2021-10-11 00:00 | a | 2 |
| + | 2021-10-11 00:00 | b | 1 |
| + | 2021-10-11 01:00 | b | 1 |
时间窗口相比普通的窗口:
- 结果只有在窗口结束才会输出
- 结果输出后,状态也随之清理.
- 由于结果是在窗口结束才会输出, 产生的数据都是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:00 | a |
|---|---|
| 2021-10-11 00:01 | b |
| 2021-10-11 00:01 | a |
| 2021-10-11 00:02 | b |
结果输出
| time | word | count | |
|---|---|---|---|
| + | 2021-10-11 00:00 | a | 1 |
| + | 2021-10-11 00:01 | b | 1 |
| + | 2021-10-11 00:01 | a | 2 |
| + | 2021-10-11 00:01 | b | 2 |
这里只做简单介绍, 详细的内容可以参考Flink官方文档: https://nightlies.apache.org/flink/flink-docs-master/zh/docs/dev/table/sql/queries/window-agg/
使用案例
去重
source表结构:
| 字段名 | id | time | item1 | item2 |
|---|---|---|---|---|
| 类型 | long | timestamp | string | string |
时间窗口去重
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)FROMsourceTableGROUPBYidover 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计算
| 字段 | ip | time |
|---|---|---|
| 类型 | String | timestamp |
时间窗口计算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'),ipover 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这种写法会存在撤回) | 状态不会自动清理,需要设置下状态清理时间 |
本文收录于「我的数据空间」技术库——一套可私有化部署的数据平台(数据集成 / 实时计算 / 数据湖 / 湖仓查询 / 智能问数)。