【Flink SQL 学习笔记 二】实时聚合:用窗口和 Watermark 把订单按时段统计
上一篇用
GROUP BY city做了全局聚合——从任务启动到现在,每个城市一共多少订单。但业务场景不是这么用的:大盘要看"每分钟订单量",运营要看"近 5 分钟趋势",日报要看"当天累计 GMV"。这些需求都有时间边界——按时段切分,而不是无限累加。这就是窗口解决的问题:把无限的数据流切成一段段有限的,每段单独聚合。
本文要点
- 为什么需要窗口:流是无限的,聚合必须有边界
- 三种窗口类型:滚动、滑动、累积
- TVF(Table-Valued Function):Flink SQL 写窗口聚合的标配语法
- Watermark:让迟到的订单归到正确的窗口
- 迟到数据处理:窗口关窗后的数据怎么办
- 完整任务:按分钟统计各城市订单量,写入 Paimon ADS
零 为什么需要窗口
1 全局聚合的问题
上一篇的完整任务里写了这样的 SQL:
INSERT INTO ads_city_stats
SELECT city, COUNT(*), SUM(amount)
FROM orders_source
GROUP BY city;
这个聚合没有时间边界——从任务启动开始,北京的订单量一直在累加,1 分钟后是 50 条,1 小时后是 3000 条,永远不会"结束"。这在业务上没有意义——大盘看的是"当前每分钟多少订单",不是"从今天任务启动到现在一共多少"。
2 窗口做了什么
窗口把无限的数据流切成一段段有限的:
无限流: e1 → e2 → e3 → e4 → e5 → e6 → e7 → ...(永远不停)
↕
窗口切分: [e1, e2, e3] | [e4, e5] | [e6, e7] | ...
窗口1 窗口2 窗口3
(00:00-01:00) (01:00-02:00) (02:00-03:00)
每个窗口是一段有限数据,可以独立聚合,聚合完就输出结果。窗口之间互不依赖——窗口 1 不需要等窗口 2,窗口 2 也不需要等窗口 1 关闭。数据来了一条 00:20 的,进窗口 1 累加;紧接着来了一条 01:10 的,进窗口 2 累加。两个窗口同时在攒数据,各自独立关闭、各自输出结果。
一 窗口类型
Flink SQL 支持三种主流窗口。窗口既是统计分组的维度——把数据按时间切成段,每段独立聚合;也决定了结果什么时候输出——窗口关闭时才把聚合结果写出去。上一篇的全局聚合(GROUP BY city)和本篇的窗口聚合不冲突,反而互补:全局聚合给的是"今天一共多少单",窗口聚合多了"每分钟多少单"这个时间维度。用订单场景说明三种窗口的区别:

1.1 滚动窗口(Tumbling Window)
电商大盘需要看"每分钟各城市的订单量"——00:00 到 00:59 的订单算一个数,01:00 到 01:59 的订单算另一个数,分钟之间互不影响。滚动窗口每 1 分钟切一段,段与段不重叠、不遗漏。每段结束时输出一次结果,下一段从零开始重新计数。
1.2 滑动窗口(Sliding Window)
运营看每分钟订单量时发现一个问题:单分钟的数字跳动太大,看不出趋势。需要的是"近 2 分钟"的数据,每 30 秒刷新一次。滚动窗口做不到——它每段不重叠,看不到跨段的"近 2 分钟"。滑动窗口让窗口之间重叠:窗口长度 2 分钟、步长 30 秒,00:00-01:59 是一个窗口,00:30-02:29 又是一个窗口。同一条订单会出现在多个窗口里,每个窗口各自聚合、各自输出。代价是 State 更大,因为一条数据要在多个窗口里各存一份。
1.3 累积窗口(Cumulative Window)
财务要出"当日渐进式销售额"——每 15 分钟看一次从 0 点到现在的累计 GMV。滚动窗口只能给每 15 分钟独立的数字,给不了"从 0 点累计到现在"的结果。累积窗口在周期内逐步扩大:第一次输出 00:00-00:15 的聚合值,第二次输出 00:00-00:30 的聚合值,到 01:00 周期结束归零,下一轮重新开始。每次输出的数据范围都在扩大,包含之前所有数据。
二 TVF:用 SQL 写窗口聚合
1 TVF 是什么
TVF(Table-Valued Function,表值函数)是 Flink SQL 写窗口聚合的标配方式。核心思路:把窗口定义放在 FROM TABLE(...) 里,当成一张"窗口表"来查,GROUP BY 只管聚合逻辑:
SELECT window_start, window_end, city, COUNT(*) AS order_cnt, SUM(amount) AS total_amount
FROM TABLE(
TUMBLE(TABLE orders_source, DESCRIPTOR(order_time), INTERVAL '1' MINUTE)
)
GROUP BY window_start, window_end, city;
窗口定义(TUMBLE(...))和数据聚合(GROUP BY city)分开了,读起来更清晰:先开窗口,再按窗口+城市分组聚合。
TUMBLE 函数三个参数的含义:
TABLE orders_source:在哪张表上开窗口——就是 Source 表的名字DESCRIPTOR(order_time):用哪个时间字段切窗口——必须是 DDL 里定义了 WATERMARK 的字段INTERVAL '1' MINUTE:窗口多长——这里是 1 分钟
执行后,每行数据会被加上 window_start 和 window_end 两个隐藏字段,表示这条数据属于哪个窗口。GROUP BY window_start, window_end, city 就是按窗口+城市分组聚合。
2 滑动窗口 TVF
SELECT
window_start,
window_end,
city,
COUNT(*) AS order_cnt
FROM TABLE(
HOP(TABLE orders_source, DESCRIPTOR(order_time),
INTERVAL '30' SECOND, -- 滑动步长
INTERVAL '1' MINUTE) -- 窗口长度
)
GROUP BY window_start, window_end, city;
HOP 函数四个参数的含义:
TABLE orders_source:在哪张表上开窗口DESCRIPTOR(order_time):用哪个时间字段切窗口INTERVAL '30' SECOND:滑动步长——每隔 30 秒开一个新窗口INTERVAL '1' MINUTE:窗口长度——每个窗口覆盖 1 分钟的数据
步长 < 窗口长度时,窗口之间就会重叠。步长 30 秒、窗口 1 分钟,意味着同一条数据会同时落在两个窗口里。步长 = 窗口长度时,滑动窗口就退化成了滚动窗口。
3 累积窗口 TVF
SELECT
window_start,
window_end,
city,
SUM(amount) AS cumulative_amount
FROM TABLE(
CUMULATE(TABLE orders_source, DESCRIPTOR(order_time),
INTERVAL '15' SECOND, -- 步长
INTERVAL '1' MINUTE) -- 总周期
)
GROUP BY window_start, window_end, city;
CUMULATE 函数四个参数的含义:
TABLE orders_source:在哪张表上开窗口DESCRIPTOR(order_time):用哪个时间字段切窗口INTERVAL '15' SECOND:步长——每 15 秒输出一次累计结果INTERVAL '1' MINUTE:总周期——1 分钟后归零重新开始
和滑动窗口的区别:滑动窗口的每个窗口长度固定(1 分钟),窗口之间平行移动;累积窗口的窗口从周期起点开始,逐步扩大到周期结束——第一次覆盖 15 秒,第二次覆盖 30 秒,直到覆盖完整的 1 分钟后归零重置。
4 三种 TVF 对比
三种窗口的 SQL 结构完全一样——SELECT ... FROM TABLE(窗口函数) GROUP BY window_start, window_end, ...。区别只在 FROM TABLE(...) 里的函数名和参数:
| 函数 | 窗口类型 | 参数 | 参数含义 |
|---|---|---|---|
TUMBLE |
滚动 | 表, 时间字段, 窗口长度 | 3 个参数,窗口等长不重叠 |
HOP |
滑动 | 表, 时间字段, 步长, 窗口长度 | 4 个参数,多一个"步长"——每隔多久开一个新窗口 |
CUMULATE |
累积 | 表, 时间字段, 步长, 总周期 | 4 个参数,多一个"总周期"——周期结束时归零重置 |
记住函数名和参数差异,其余部分套同一个模板就行。
三 Watermark:让迟到的订单归到正确的窗口
1 窗口什么时候关
上一节的 TVF 任务跑起来了——Kafka 里的订单一条条进来,按 order_time 分配到 1 分钟窗口,按城市聚合。窗口 1 [00:00-01:00) 在攒数据,窗口 2 [01:00-02:00) 也在攒数据,互不干扰。
但马上遇到一个问题:窗口 1 什么时候把结果输出出来?
01:00(处理时间)到了就关窗?可能不行——一条 00:30 的北京订单,因为网络延迟 01:02 才到 Flink。01:00 关窗,这条就丢了,北京的订单量少算了一条。
一直等?等 01:10?01:30?等越久结果越延迟,等太短又丢数据。到底等多久才"够"?
需要一个机制告诉 Flink:事件时间推进到什么程度,就可以认为之前的窗口数据到齐了。这就是 Watermark。
2 订单为什么会乱序到达
要搞清楚"等多久才够",得先看数据为什么会乱序。
订单从下单到进入 Flink,经过 App → API → Kafka → Flink 这条链路。网络延迟、重试、Kafka 分区不均都会导致数据不按下单顺序到达。
涉及两个时间维度:
- 事件时间(Event Time):订单实际下单的时间,印在数据里(
order_time字段) - 处理时间(Processing Time):Flink 收到这条数据的时间,看机器时钟
处理时间一定 >= 事件时间——数据必须先产生才能到达,所以只有"晚到"的数据,没有"早到"的数据。但晚的程度不同:
| 订单 | 事件时间(下单时间) | 处理时间(到达 Flink) | 延迟 |
|---|---|---|---|
| 订单A(北京100元) | 00:30 | 00:42 | 12 秒 |
| 订单B(上海200元) | 00:40 | 00:41 | 1 秒 |
订单B 下单更晚(00:40),但延迟更小(1 秒),所以 00:41 就到了。订单A 下单更早(00:30),但延迟更大(12 秒),00:42 才到。结果:事件时间更晚的订单B 反而先到——这就是乱序。
如果按处理时间分窗口:订单A 到达时是 00:42,会被分到 00:42 所在的 [00:00-01:00) 窗口——碰巧对了。但换个场景:如果订单A 的下单时间是 00:58,到达时间是 01:02,按处理时间它会被分到 [01:00-02:00) 窗口——错了,它应该在前一分钟。只有按事件时间分窗口才是对的。
3 Watermark:给延迟一个容忍上限
既然数据都会晚到,但晚的程度不同,Watermark 的思路是:给延迟定一个容忍上限。假设最多晚 5 秒,超过 5 秒的数据认为不会来了,就可以安全地关窗。
定义方式(在 Source 表的 DDL 里):
CREATE TABLE orders_source (
...
order_time TIMESTAMP(3),
-- 允许 5 秒延迟
WATERMARK FOR order_time AS order_time - INTERVAL '5' SECOND
) WITH ( ... );
order_time - INTERVAL '5' SECOND 的含义:Flink 不断追踪已到达数据中最大的事件时间,然后减去 5 秒得到 Watermark。比如已经收到事件时间最大为 00:30 的数据,Watermark 就是 00:25——告诉 Flink,00:25 之前的数据应该都到齐了,可以安全地关闭 00:25 之前的窗口。
为什么要减 5 秒,而不是直接用最大事件时间? 因为最大事件时间只代表"这个时间点的数据到了",不代表"比它更早的数据都到了"。如果收到 00:30 的订单就直接把 WM 推到 00:30,等于告诉 Flink “00:30 之前全到了”——但一条 00:28 的订单可能延迟 2 秒还没到,窗口此时关闭,这条就丢了。减 5 秒就是留一个缓冲:收到 00:30 时 WM 只推到 00:25,00:25 到 00:30 之间的数据不急着关窗,再等等。等后续收到 01:05 的订单时 WM 才到 01:00,窗口 [00:00-01:00) 才真正关闭。这 5 秒不是浪费时间,是给迟到数据留的等待期——不减就是"看到一条就认为前面全到了",乱序数据全丢。
5 秒不是拍脑袋定的——要观察线上数据的延迟分布,延迟 99 分位在 3 秒以内,设 5 秒就够了。但延迟容忍值要和窗口大小匹配:1 分钟的窗口设 5 秒延迟,关窗推迟 5 秒,影响不大;同样的窗口设 30 秒延迟,等于半个窗口的数据都可能迟到,结果要等 1 分 30 秒才出来,延迟太高。一般经验是延迟容忍不超过窗口长度的 10%——1 分钟窗口最多容忍 6 秒,5 分钟窗口可以容忍 30 秒。如果数据延迟确实大到和窗口一个量级,说明该把窗口调大,而不是硬设一个大延迟值。
4 一条条订单看 Watermark 怎么推进窗口关门
我们可以按照事件推进过程来看:
回到第一节的问题:窗口 1 [00:00-01:00) 什么时候关?答案是:Watermark >= 01:00 时关。Watermark 到 01:00 意味着 Flink 认为 01:00 之前的数据都到齐了(允许 5 秒延迟)。
用 5 条订单走一遍(图中从左到右是到达顺序):
- 北京100元(下单 00:20)到达 → 最大事件时间 00:20,WM = 00:15。窗口 1 没关(00:15 < 01:00)
- 上海200元(下单 00:50)到达 → 最大事件时间 00:50,WM = 00:45。窗口 1 没关
- 深圳300元(下单 01:00)到达 → 最大事件时间 01:00,WM = 00:55。窗口 1 仍然没关(00:55 < 01:00)
- 北京150元(下单 00:40)到达 → 事件时间 00:40 < 当前最大 01:00,WM 保持 00:55 不变。窗口 1 还没关(00:55 < 01:00),这条数据还能进窗口 1
- 上海80元(下单 01:50)到达 → 最大事件时间 01:50,WM = 01:45 >= 01:00,窗口 1 关闭,输出结果:北京 250(100+150),上海 200
关键规则:Watermark 只增不减。第 4 步来了一条事件时间更小的数据(下单 00:40),Watermark 不会倒退,保证窗口只会关一次。
5 窗口关了之后又来了迟到数据怎么办
窗口 1 在第 5 步关闭后,输出了结果。如果此时又来了一条 深圳50元(下单 00:35)——它的窗口是 [00:00-01:00),但 Watermark 已经是 01:45,远超窗口结束时间 01:00。这条数据属于"迟到",Flink 默认丢弃。
对于每分钟订单量这种大盘场景,偶尔丢几条不影响整体准确性。如果是支付金额等不能丢的场景,可以配置允许延迟(Allowed Lateness),让窗口"关了又开"——迟到数据到达后重新计算并输出更新后的结果。具体配置在第四篇生产实战中讲。
四 完整任务:按分钟统计各城市订单量
把窗口、TVF、Watermark 组合到一起,写一个完整的 Flink 任务。
Paimon ADS 表(提前用 Spark SQL 建好):
CREATE TABLE ads.city_minute_stats (
window_start TIMESTAMP(3),
window_end TIMESTAMP(3),
city STRING,
order_cnt BIGINT,
total_amount DECIMAL(10, 2),
PRIMARY KEY (window_start, city) NOT ENFORCED
) USING paimon
TBLPROPERTIES ('bucket' = '2');
Flink 任务(提交后持续运行):
-- ============================================
-- Flink 任务:Kafka 订单流 → 1 分钟滚动窗口聚合 → Paimon ADS
-- ============================================
-- 1. Source 表:对接 Kafka 订单 topic
CREATE TABLE orders_source (
order_id BIGINT,
user_id BIGINT,
product_id BIGINT,
city STRING,
amount DECIMAL(10, 2),
order_time TIMESTAMP(3),
-- 用订单时间作为事件时间,允许 5 秒延迟
WATERMARK FOR order_time AS order_time - INTERVAL '5' SECOND
) WITH (
'connector' = 'kafka',
'topic' = 'orders',
'properties.bootstrap.servers' = 'localhost:9092',
'properties.group.id' = 'flink-ads-city-minute',
'scan.startup.mode' = 'latest-offset',
'format' = 'json'
);
-- 2. 执行:1 分钟滚动窗口 + 按城市聚合 + 写入 Paimon
INSERT INTO ads.city_minute_stats
SELECT
window_start,
window_end,
city,
COUNT(*) AS order_cnt,
SUM(amount) AS total_amount
FROM TABLE(
TUMBLE(TABLE orders_source, DESCRIPTOR(order_time), INTERVAL '1' MINUTE)
)
GROUP BY window_start, window_end, city;

任务运行时的行为:
- 订单从 Kafka 不断流入,按
order_time分配到对应的 1 分钟窗口 - 每个窗口内按城市聚合,结果暂存在 Flink State 里
- Watermark 到达窗口结束时间时,窗口关闭,输出该窗口的聚合结果到 Paimon
- Paimon 主键表收到结果,自动 upsert——窗口重新输出时,覆盖旧值
- 窗口关闭后,对应的 State 被清理,释放资源
这里说的 State 是 Flink 存放中间计算结果的地方。生产环境通常用 RocksDB 作为 State Backend——数据存在 TaskManager 所在机器的本地磁盘上,通过嵌入式 RocksDB 引擎管理,容量只受磁盘限制。单个 TaskManager 的 State 存储量通常在几十 GB 到数百 GB 级别。State 不是永久占着不放的:窗口关闭后数据清理,也可以设置 TTL 让过期状态自动回收。容错靠 Checkpoint——定期把 State 快照上传到远程存储(S3/HDFS),TaskManager 挂了就从快照恢复。
和上一篇全局聚合的区别,就是多了一个时间维度:
| 对比项 | 全局聚合(上一篇) | 窗口聚合(本篇) |
|---|---|---|
| GROUP BY | city |
window_start, window_end, city |
| 聚合范围 | 从启动到现在的全部数据 | 每个 1 分钟窗口内的数据 |
| 结果行数 | 每个城市 1 行(不断更新) | 每个城市 × 每个窗口 1 行 |
| 输出方式 | 每来一条就撤回旧值发新值 | 窗口内累积,关窗时输出一次 |
这里说的"累积"不是攒批——每来一条订单,Flink 立刻在 State 里更新聚合值(SUM 加一下、COUNT 加一),计算是逐条增量的。攒的不是数据,是结果:结果暂存在 State 里不往外发,等窗口关闭了才输出。这和 Spark 微批不同,Spark 微批是攒一批数据最后一起算;Flink 窗口聚合是边来边算,只是延迟输出。
两种不是优劣关系,是适用场景不同:全局聚合适合看总量(“今天一共多少单”),窗口聚合适合看趋势(“每分钟多少单”)。窗口只是把 GROUP BY 多加了一个时间维度。
五 小结
- 流是无限的,聚合需要边界——窗口把无限流切成有限段,每段独立聚合
- 三种窗口覆盖主流场景——滚动(每分钟统计)、滑动(近 N 分钟趋势)、累积(渐进式日报)
- TVF 是写窗口聚合的标配——
TUMBLE(TABLE ..., DESCRIPTOR(...), INTERVAL ...)先开窗口,再 GROUP BY - Watermark 解决数据乱序——用事件时间减去延迟容忍值,表示"这个时间之前的数据都到齐了"
- 窗口关闭 = Watermark 过了窗口结束时间——关了之后迟到的数据默认丢弃
- 窗口聚合和全局聚合不是优劣关系——窗口只是多加了一个时间维度,看趋势用窗口聚合,看总量用全局聚合
下一篇讲 Join:订单流怎么和城市维表、商品维表关联,把 ID 转成可读的名称。
更多推荐
所有评论(0)