用四种窗口 TVF 切分同一批订单
目标
对同一个订单流分别应用 TUMBLE、HOP、CUMULATE、SESSION 窗口,通过结果确认窗口列是如何附加的,以及一行会进入几个窗口。创建位于窗口聚合之上的窗口 Top-N,以及直接位于窗口 TVF 之上的窗口 Top-N。
为什么重要
窗口类型选错,数字会悄悄出错——把 HOP 的条数相加,会因为重叠而虚高;在 GROUP BY 中去掉窗口列,就不再是窗口聚合而变成无限聚合,会混入变更日志。窗口 TVF 只是给行附加窗口列的函数,只要准确了解这些列是如何附加的,聚合也好、排名也好,都可以用普通 SQL 叠加上去。本实验的评分器不会查询集群——它读取你保存的 sql-client 输出,并用 Python 直接从源 CSV 计算各窗口的期望值进行核对。
步骤
- 用
flink-up启动集群,运行只取order_id <= 20的订单、输出 10 分钟TUMBLE附加的order_id, ts, window_start, window_end, window_time的 /root/flink/windows/assign.sql,并把输出保存到 /root/flink/windows/assign.out。 - 把输出每个 10 分钟
TUMBLE窗口、每个店铺(shop)的cnt(条数)和revenue(amount 之和)的 /root/flink/windows/tumble.sql 的输出保存到 /root/flink/windows/tumble.out。 - 把输出每个
HOP(slide 5 分钟,size 10 分钟)窗口的cnt的 /root/flink/windows/hop.sql 的输出保存到 /root/flink/windows/hop.out。 - 把输出每个
CUMULATE(step 10 分钟,size 1 小时)窗口的revenue的 /root/flink/windows/cumulate.sql 的输出保存到 /root/flink/windows/cumulate.out。 - 把输出每个
SESSION(PARTITION BY shop,gap 5 分钟)窗口、每个店铺的cnt的 /root/flink/windows/session.sql 的输出保存到 /root/flink/windows/session.out。 - 把输出每个 10 分钟 TUMBLE 窗口中各店铺销售额最高的 2 家(
window_start, window_end, shop, revenue, rownum)的 /root/flink/windows/top-shops.sql 的输出保存到 /root/flink/windows/top-shops.out。 - 把不做聚合、输出每个 30 分钟 TUMBLE 窗口中金额最大的 3 笔订单(
order_id, shop, amount, window_start, window_end, rownum)的 /root/flink/windows/top-orders.sql 的输出保存到 /root/flink/windows/top-orders.out。 - 在 /root/flink/windows/report.json 中写入
orders、hop_assignments、cumulate_windows、sessions、max_session_orders。
参考
- 源数据:
/opt/lab/fixtures/data/windows_orders.csv,列order_id BIGINT, shop STRING, amount INT, ts TIMESTAMP(3)(没有表头的 CSV,ts 升序)。在表上设置WATERMARK FOR ts AS ts - INTERVAL '1' SECOND,并以流式模式运行。 - 形态:
FROM TUMBLE(TABLE orders, DESCRIPTOR(ts), INTERVAL '10' MINUTE)。HOP 是(TABLE, DESCRIPTOR, slide, size),CUMULATE 是(TABLE, DESCRIPTOR, step, size),SESSION 是(TABLE orders PARTITION BY shop, DESCRIPTOR(ts), gap)。 - 窗口聚合用
GROUP BY window_start, window_end, ...分组。去掉窗口列就会变成无限聚合,并混入 -U/+U。 - 常见错误:把 HOP、CUMULATE 的参数顺序(较小的值在前)弄反。会因为 size 不是 slide(step)的整数倍而被拒绝。
- 常见错误:在窗口 Top-N 的 PARTITION BY 里漏掉 window_start、window_end——那样就不是窗口 Top-N,而是普通 Top-N,排名每次变化都会产生日志。
- 官方文档:Windowing TVF · Window Aggregation · Window Top-N · Time Attributes
窗口 TVF 附加的三列
用 flink-up 启动集群,在 /root/flink/windows/assign.sql 中写入源表 orders(参考中的列和水位线),并以流式运行 SELECT order_id, ts, window_start, window_end, window_time FROM TUMBLE(TABLE orders, DESCRIPTOR(ts), INTERVAL '10' MINUTE) WHERE order_id <= 20,把输出保存到 /root/flink/windows/assign.out。
窗口 TVF 保持原有的列不变,附加三个窗口列后返回。窗口是 [起点, 终点) 左闭右开区间,所以恰好落在边界时刻的订单,会进入从该时刻开始的窗口。请看 window_time 与 window_end 的差别。
按店铺的 10 分钟聚合
在 /root/flink/windows/tumble.sql 中写出按 10 分钟 TUMBLE 窗口和 shop 分组、输出 window_start, window_end, shop, COUNT(*) AS cnt, SUM(amount) AS revenue 的窗口聚合,并把输出保存到 /root/flink/windows/tumble.out。
窗口聚合要在 GROUP BY 中放入 window_start 和 window_end。窗口不重叠,所以把 cnt 全部加起来,等于原来的订单数。结果在窗口关闭时只以 +I 输出一次。
HOP——一个订单进入两个窗口
在 /root/flink/windows/hop.sql 中写出对每个 HOP(TABLE orders, DESCRIPTOR(ts), INTERVAL '5' MINUTE, INTERVAL '10' MINUTE) 窗口输出 window_start, window_end, COUNT(*) AS cnt 的聚合,并把输出保存到 /root/flink/windows/hop.out。
HOP 的第三个参数是 slide(窗口开始的间隔),第四个是 size(窗口长度)。每 5 分钟开始一个 10 分钟窗口的话,窗口一半重叠,一个订单会进入两个窗口。请把 cnt 之和与原来的订单数比较。最前面的窗口可能从比第一个订单早 5 分钟的时刻开始。
CUMULATE——起点固定的累积窗口
在 /root/flink/windows/cumulate.sql 中写出对每个 CUMULATE(TABLE orders, DESCRIPTOR(ts), INTERVAL '10' MINUTE, INTERVAL '1' HOUR) 窗口输出 window_start, window_end, SUM(amount) AS revenue 的聚合,并把输出保存到 /root/flink/windows/cumulate.out。
CUMULATE 是先按 size(1 小时)做 TUMBLE,再把其中切成终点每个 step(10 分钟)递增的窗口。起点相同的窗口,revenue 随终点越晚,应该越大或者不变。请数一数每小时会产生几个窗口。
SESSION——每个店铺长度不同的窗口
在 /root/flink/windows/session.sql 中写出对每个 SESSION(TABLE orders PARTITION BY shop, DESCRIPTOR(ts), INTERVAL '5' MINUTE) 窗口、每个店铺输出 window_start, window_end, shop, COUNT(*) AS cnt 的聚合,并把输出保存到 /root/flink/windows/session.out。
同一店铺相邻订单的间隔小于等于 5 分钟,会话就会延续,超过就开始新的会话。会话的起点是第一个订单的时刻,末端是最后一个订单 + 5 分钟。去掉 PARTITION BY,就会把各店铺混在一起划分会话。
位于窗口聚合之上的窗口 Top-N
在 /root/flink/windows/top-shops.sql 中,先求出按 10 分钟 TUMBLE 窗口和店铺的销售额(SUM(amount) AS revenue),再用 ROW_NUMBER() OVER (PARTITION BY window_start, window_end ORDER BY revenue DESC) AS rownum 为每个窗口只保留前 2 家,输出 window_start, window_end, shop, revenue, rownum。输出保存到 /root/flink/windows/top-shops.out。
把窗口聚合作为子查询,在其上编 ROW_NUMBER,再在外层用 rownum <= 2 过滤。PARTITION BY 里要有两个窗口列,才会成为窗口 Top-N,在窗口关闭时只输出一次结果。只有一家店铺的窗口,只会有一行。
直接位于窗口 TVF 之上的窗口 Top-N
在 /root/flink/windows/top-orders.sql 中不做聚合,直接在 30 分钟 TUMBLE 窗口 TVF 之上用 ROW_NUMBER() OVER (PARTITION BY window_start, window_end ORDER BY amount DESC) 排名,输出每个窗口中金额最大的 3 笔订单的 order_id, shop, amount, window_start, window_end, rownum。输出保存到 /root/flink/windows/top-orders.out。
窗口 Top-N 不需要窗口聚合,也可以直接放在窗口 TVF 的结果之上。此时排名的对象不是聚合行,而是订单行本身。不使用 GROUP BY。数据被设计成金额没有并列,所以排名不会抖动。
报告——每个窗口被数了几次
在 /root/flink/windows/report.json 中以整数写入 orders(tumble.out 的 cnt 之和)、hop_assignments(hop.out 的 cnt 之和)、cumulate_windows(cumulate.out 的行数)、sessions(session.out 的行数)、max_session_orders(session.out 中 cnt 的最大值)。
结果行以“| +I |”开头。用 awk -F'|' 切分后,第 1 栏是空栏,第 2 栏是 op,其后按 SELECT 的列顺序依次排列。请把 hop_assignments 与 orders 比较——它是 size/slide 倍。