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

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

用四种窗口 TVF 切分同一批订单

在 TT Lab 中继续学习

目标

对同一个订单流分别应用 TUMBLE、HOP、CUMULATE、SESSION 窗口,通过结果确认窗口列是如何附加的,以及一行会进入几个窗口。创建位于窗口聚合之上的窗口 Top-N,以及直接位于窗口 TVF 之上的窗口 Top-N。

为什么重要

窗口类型选错,数字会悄悄出错——把 HOP 的条数相加,会因为重叠而虚高;在 GROUP BY 中去掉窗口列,就不再是窗口聚合而变成无限聚合,会混入变更日志。窗口 TVF 只是给行附加窗口列的函数,只要准确了解这些列是如何附加的,聚合也好、排名也好,都可以用普通 SQL 叠加上去。本实验的评分器不会查询集群——它读取你保存的 sql-client 输出,并用 Python 直接从源 CSV 计算各窗口的期望值进行核对。

步骤

  1. 用 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。
  2. 把输出每个 10 分钟 TUMBLE 窗口、每个店铺(shop)的 cnt(条数)和 revenue(amount 之和)的 /root/flink/windows/tumble.sql 的输出保存到 /root/flink/windows/tumble.out。
  3. 把输出每个 HOP(slide 5 分钟,size 10 分钟)窗口的 cnt 的 /root/flink/windows/hop.sql 的输出保存到 /root/flink/windows/hop.out。
  4. 把输出每个 CUMULATE(step 10 分钟,size 1 小时)窗口的 revenue 的 /root/flink/windows/cumulate.sql 的输出保存到 /root/flink/windows/cumulate.out。
  5. 把输出每个 SESSION(PARTITION BY shop,gap 5 分钟)窗口、每个店铺的 cnt 的 /root/flink/windows/session.sql 的输出保存到 /root/flink/windows/session.out。
  6. 把输出每个 10 分钟 TUMBLE 窗口中各店铺销售额最高的 2 家(window_start, window_end, shop, revenue, rownum)的 /root/flink/windows/top-shops.sql 的输出保存到 /root/flink/windows/top-shops.out。
  7. 把不做聚合、输出每个 30 分钟 TUMBLE 窗口中金额最大的 3 笔订单(order_id, shop, amount, window_start, window_end, rownum)的 /root/flink/windows/top-orders.sql 的输出保存到 /root/flink/windows/top-orders.out。
  8. 在 /root/flink/windows/report.json 中写入 orders、hop_assignments、cumulate_windows、sessions、max_session_orders。

参考

窗口 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 倍。