TT Lab
开始
学习 学习路径 课程

Apache Flink — 用真正的引擎跑流处理

窗口 TVF——为每行附加三个窗口列的函数

在 TT Lab 中继续学习

一句话总结

Flink SQL 的窗口是接收一张表并返回一张表的函数(窗口 TVF)。TUMBLE、HOP、CUMULATE、SESSION 只是在原来的行上附加 window_start、window_end、window_time 三列后返回,聚合和 Top-N 则是按这些列分组的普通 SQL。窗口类型决定的只有一件事——一行会进入几个、哪些窗口。

为什么需要它

正如上一个模块所见,像 GROUP BY user_id 这样的无限聚合会不断改写结果,并永远为每个键持有状态。仪表板想要的通常不是这个。“每 10 分钟各店铺的销售额”“每 5 分钟刷新一次的最近 10 分钟订单数”“从今天零点到现在的累计销售额”“从进来到离开的一次会话”——全都是按时间切出来的区间上的聚合。区间关闭时,输出一次结果并丢弃状态就行了。

以前的 Flink SQL 是用 GROUP BY TUMBLE(ts, ...) 这类特殊语法(Grouped Window Functions)来做的。它可以用于聚合,却无法用于按窗口排名,或是把两个流按相同窗口连接。官方文档(Windowing TVF)指出,窗口 TVF 取代了那种语法。窗口变成了“给行附加列的函数”之后,就可以在窗口之上放任何东西——窗口聚合、窗口 Top-N、窗口连接、窗口去重。

工作原理

横轴是 09:00 到 09:30 的事件时间,上面标着六个订单点。TUMBLE 10 分钟是三个互不重叠的窗口,每个订单只进入一个。HOP 是每 5 分钟开始一个 10 分钟窗口,所以窗口一半重叠,每个订单进入两个窗口。CUMULATE 是起点同为 09:00、终点每次增加 10 分钟的三个窗口,所以越早的订单进入的窗口越多。SESSION 把间隔在 5 分钟之内连续到来的订单归为一个会话,产生两个会话,会话的末端是最后一个订单加 5 分钟的位置

窗口 TVF 写在 FROM 的位置。第一个参数是表,第二个是时间属性列,其余是大小。

SELECT window_start, window_end, shop, COUNT(*) AS cnt
FROM TUMBLE(TABLE orders, DESCRIPTOR(ts), INTERVAL '10' MINUTE)
GROUP BY window_start, window_end, shop;
TVF 参数 一行进入的窗口数 重叠
TUMBLE size 1 无
HOP slide, size size / slide 有
CUMULATE step, size 起点相同的窗口中,终点晚于该行的全部 有
SESSION (PARTITION BY 键) gap 1 无,大小各不相同

有几条规则决定了结果。

窗口是左闭右开区间。 [window_start, window_end)——恰好在 09:10:00 的订单进入的是 [09:10, 09:20),而不是 [09:00, 09:10)。window_time 如文档所说,始终是 window_end − 1ms,在流式中,这一列仍保持为时间属性,可以用于后续的窗口运算。反过来,window_start、window_end 是普通的时间戳,不是时间属性。

窗口的起点与 epoch 对齐。 10 分钟窗口以整点为基准,在 00、10、20 分开始,即使第一行在 09:01:35 到来,窗口也从 09:00 开始。要移动它,用可选参数 offset。

HOP 和 CUMULATE 的参数顺序是个陷阱。 两者都是较小的值(slide、step)在前,较大的值(size)在后。写反了,作业启动之前就会被拒绝——在实验环境中,HOP 报出“size must be an integral multiple of slide”,CUMULATE 报出“maxSize must be an integral multiple of step”。也就是说,size 必须是 slide(step)的整数倍。在 HOP(5 分钟, 10 分钟) 中,一行进入两个窗口,所以把各窗口的条数全部加起来,会是原始行数的两倍。这不是 bug,而是定义。CUMULATE 如文档所说,是“先按 size 做 TUMBLE,再把其中切成终点每个 step 递增的窗口”,所以一小时窗口、step 10 分钟时,会产生六个起点相同的窗口。

SESSION 没有大小。 同一个键的相邻行间隔小于等于 gap,就连成一个会话,会话的末端是最后一行 + gap。所以每个会话的长度都不同。按文档所述,SESSION TVF 目前还不支持批处理模式。

窗口运算只在结束时输出一次。 文档(Window Aggregation)指出,窗口聚合不输出中间结果,只在窗口结束时输出最终结果,并清除所有不再需要的状态。所以结果日志里只有 +I。窗口“结束了”的判断,由上一个模块的水位线来做。

窗口 Top-N 按窗口列划分。 像 ROW_NUMBER() OVER (PARTITION BY window_start, window_end ORDER BY ...) 这样,PARTITION BY 里有两个窗口列,优化器才会把它翻译成窗口 Top-N。这样一来,与普通 Top-N 不同,它不会在排名每次变化时输出 -U/+U,而是在窗口关闭时只输出一次前 N 名。它既可以放在窗口聚合之上,也可以直接放在窗口 TVF 之上,对行本身排名。按文档所述,直接位于 TVF 之上的窗口 Top-N 只支持 TUMBLE、HOP、CUMULATE。

在现场相遇的样子

如果用 HOP 做“每 5 分钟、最近 1 小时”的仪表板,一行会进入 12 个窗口。状态和计算量相应增加,而如果把各窗口的条数相加当作“订单总数”,就会虚高 12 倍。用 HOP 去模拟累计指标(“今天到现在”)也很常见,而这个位置正确的是起点固定的 CUMULATE。

第二种是排名的抖动。金额相同的两个订单争夺第 3 名,ROW_NUMBER 会从中选一个。它并不保证是哪一个,所以重新处理时可能会不同。需要养成习惯,在排序键里加上订单号之类的辅助键来打破并列。本实验的数据被设计成金额和各窗口销售额都没有并列。

第三种是时区。本实验使用 TIMESTAMP(3),所以窗口按所写的时刻原样切分。如果用 TIMESTAMP_LTZ 列做按天窗口,“一天”的边界会随会话时区而变。如果日窗口是在韩国时间上午 9 点切开的,就先看这一点。

下一项实验要做什么

对 119 个订单应用 10 分钟 TUMBLE,看窗口列是如何附加的,并输出按店铺的 10 分钟聚合。确认 HOP(5 分钟, 10 分钟) 的条数之和变成两倍,CUMULATE(10 分钟, 1 小时) 产生六个起点相同的窗口,SESSION(间隔 5 分钟)为每个店铺生成长度不同的会话。再用窗口 Top-N 取出每个 10 分钟窗口中销售额最高的两家店铺,以及每个 30 分钟窗口中金额最大的三笔订单,最后把数字整理成报告。