Skip to content

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

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

基本概念介绍:

三类窗口的语义差异:GROUP WINDOW / TIME WINDOW / OVER WINDOW

对于Flink的窗口,我们大概可以分为三类:

GROUP WINDOW(普通窗口)

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

sql
-- words表只有一个字段, 每一行是不同的单词
SELECT word, COUNT(*) AS cnt
FROM words
GROUP BY word;

如果输入顺序为:

a
b
c
a
b

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

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

注意:

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

TIME WINDOW

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

sql
SELECT TUMBLE_START(`time`, INTERVAL '1' HOUR), word, COUNT(*)
FROM words
GROUP BY TUMBLE(`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为例:

sql
SELECT `time`, word
  COUNT(amount) OVER (
    PARTITION BY word
    ORDER BY `time`
    RANGE BETWEEN INTERVAL '1' HOUR PRECEDING AND CURRENT ROW
  ) AS `count`
FROM Orders
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天的时间窗口, 所以最终结果会在凌晨进行输出.

sql
SELECT TUMBLE_START(`time`, INTERVAL '1' DAY), id, LAST_VALUE(item1), LAST_VALUE(item2)
FROM sourceTable
GROUP BY TUMBLE(`time`, INTERVAL '1' DAY), id

普通group窗口去重

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

sql
SELECT id, LAST_VALUE(item1), LAST_VALUE(item2)
FROM sourceTable
GROUP BY id

over window窗口去重

over window + row_number()

sql
SELECT id, item1, item2
FROM (
  SELECT *,
    ROW_NUMBER() OVER (PARTITION BY id ORDER BY proctime ASC) AS row_num
  FROM sourceTable)
WHERE row_num = 1

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

PV,UV计算

字段iptime
类型Stringtimestamp

时间窗口计算pv, uv

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

sql
SELECT TUMBLE_START(`time`, INTERVAL '1' HOUR), ip, COUNT(*) AS pv, COUNT(DISTINCT(*)) AS uv
FROM sourceTable
GROUP BY TUMBLE(`time`, INTERVAL '1' HOUR), ip

普通窗口计算pv, uv

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

sql
SELECT DATE_FORMAT(`time`, 'yyyy-MM-dd hh:00:00'), ip, COUNT(*) as pv, COUNT(DISTINCT(*)) AS uv
FROM sourceTable
GROUP BY DATE_FORMAT(`time`, 'yyyy-MM-dd hh:00:00'), ip

over window计算pv, uv

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

sql
-- 使用over window计算每条数据到来时累计的数据
CREATE VIEW pv_uv_per_10min AS
SELECT 
  MAX(SUBSTR(DATE_FORMAT(ts, 'HH:mm'),1,4) || '0') OVER w AS time_str, 
  COUNT(ip) OVER w AS pv,
  COUNT(DISTINCT ip) OVER w AS uv
FROM sourceTable
WINDOW w AS (ORDER BY proctime ROWS BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW);

-- 使用groupBy过滤出10分钟内的最大值
SELECT time_str, MAX(pv), MAX(uv)
FROM pv_uv_per_10min
GROUP BY time_str;

总结

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

本文是「我的数据空间」实时计算实践笔记的一篇。平台把这类链路做成了托管作业(Flink SQL 流作业、实时集成、Iceberg 表管理与维护),支持私有化部署与 OEM 合作 —— 见核心能力总览技术博客,交流合作 QQ:1559851993。