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

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

去除重复并排名

在 TT Lab 中继续学习

目标

用同一种 ROW_NUMBER 模式分别创建保留第一行、保留最后一行和 Top-N,并通过输出的条数和执行计划,确认各自产生的变更日志的形态和数量。

为什么重要

清除重发的事件与保留当前状态,在 SQL 里只差 ASC 与 DESC 一个词,但前者产生 append-only 结果,后者产生撤回与更新日志。这个差别决定了可以使用哪种 Sink,以及下游聚合会收到什么。Top-N 则取决于是否输出行号,同样的结果,日志量会相差很大。本实验的评分器不会查询集群,而是把你保存的 sql-client 输出与直接从源 CSV 计算出的值进行核对。

步骤

  1. 用 flink-up 启动集群,在 /root/flink/dedup/ddl.sql 中定义 events 和 sales(在 events 的 ts 上设置水位线)。再把以批处理方式输出 n_rows(总行数)和 n_orders(不同的订单数)的 /root/flink/dedup/count.sql 的输出保存到 /root/flink/dedup/count.out。
  2. 把为每个订单保留最早事件的保留第一行查询 /root/flink/dedup/first.sql(order_id、status、amount)的输出保存到 /root/flink/dedup/first.out。
  3. 把为每个订单保留最晚事件的保留最后一行查询 /root/flink/dedup/last.sql 的输出保存到 /root/flink/dedup/last.out。
  4. 把对第 2、3 步的查询执行 EXPLAIN CHANGELOG_MODE 的 /root/flink/dedup/explain.sql 的输出保存到 /root/flink/dedup/explain.out。
  5. 把在保留最后一行之上按 status 统计 orders(订单数)和 amount(合计)的 /root/flink/dedup/board.sql 的输出保存到 /root/flink/dedup/board.out。
  6. 把用 sales 按类别取累计 qty 合计最高的 3 个、并带上行号的 /root/flink/dedup/topn.sql(category、product、total、rn)的输出保存到 /root/flink/dedup/topn.out。
  7. 把同样的 Top-3 只在外层 SELECT 中去掉 rn 的 /root/flink/dedup/topn-norank.sql 的输出保存到 /root/flink/dedup/topn-norank.out。
  8. 在 /root/flink/dedup/report.json 中写入 n_rows、n_orders、update_pairs、shipped_now、topn_log_rows、norank_log_rows、top_books。

参考

定义两个 Source 并统计行数

运行 flink-up 后,在 /root/flink/dedup/ddl.sql 中定义 events 和 sales(在 events 的 ts 上设置 WATERMARK FOR ts AS ts)。用 sql-client.sh -i ddl.sql -f count.sql 运行以批处理模式输出 n_rows 和 n_orders 的 /root/flink/dedup/count.sql,并保存到 /root/flink/dedup/count.out。

订单数是 COUNT(DISTINCT order_id)。两个值的差,就是同一订单第二条以后到来的事件数——在第 3 步会再次遇到这个数字。

保留第一行——输出一次就结束

在 /root/flink/dedup/first.sql 中以流式模式写出去重:为每个订单(order_id)保留 ts 最早的一个事件,并输出 order_id、status、amount。输出保存到 /root/flink/dedup/first.out。

在内层 SELECT 中用 ROW_NUMBER() OVER (PARTITION BY order_id ORDER BY ts ASC) 编行号,在外层只选行号为 1 的行。看看输出的 op 是否全是 +I——输出过一次的第一行不会再改变。

保留最后一行——撤回并修正

在 /root/flink/dedup/last.sql 中写出为每个订单保留 ts 最晚事件的去重(order_id、status、amount),并把输出保存到 /root/flink/dedup/last.out。看看日志条数(+I · -U · +U)与第 1 步的两个数字是什么关系。

只需改变排序方向。同一订单的事件每次新到来,就用 -U 收回之前输出的行,并用 +U 输出新行。即使重发带来完全相同的行,也会产生一对。按处理时间排序的话,这个数量可能会波动,所以请使用 ts。

从计划中区分两种去重

在 /root/flink/dedup/explain.sql 中写两条语句,分别在第 2 步和第 3 步的查询前加上 EXPLAIN CHANGELOG_MODE,并把输出保存到 /root/flink/dedup/explain.out。

执行计划(Optimized Execution Plan)中有 Deduplicate(keep=[...]) 节点,物理计划(Optimized Physical Plan)的每一行都带有 changelogMode=[...]。比较两个查询的 keep、outputInsertOnly 和 changelogMode 有何不同。

按当前状态统计——去重之上的聚合

在 /root/flink/dedup/board.sql 中写出这样的查询:内层使用第 3 步的保留最后一行,外层用 GROUP BY status 输出 orders(订单数)和 amount(amount 之和);并把输出保存到 /root/flink/dedup/board.out。

订单从 created 变为 paid 时,去重会输出 -U created · +U paid,聚合则在 created 栏减去、在 paid 栏加上。最终 orders 的合计必须等于订单数。如果不去重而直接统计 events,这个合计会增大到行数那么多。

Top-3——每当排名变化就修正

在 /root/flink/dedup/topn.sql 中,先把 sales 按(category, product)汇总为 SUM(qty) AS total,然后写出查询:每个类别按 total 降序取前 3 个,并带上行号 rn(category、product、total、rn)。输出保存到 /root/flink/dedup/topn.out。

把 GROUP BY 的合计作为子查询,在其上用 ROW_NUMBER() OVER (PARTITION BY category ORDER BY total DESC) 编行号,再在外层用 rn <= 3 选取。合计每次变化,排名都会移动,所以输出中会混有 -U · -D。只看最终状态的话,每个类别是三行。

去掉行号,日志就减少

运行只把第 6 步查询外层 SELECT 中的 rn 去掉的 /root/flink/dedup/topn-norank.sql(category、product、total),把输出保存到 /root/flink/dedup/topn-norank.out。最终的前 3 名相同,日志行数应比第 6 步少。

行号列是结果唯一键的一部分,所以输出它的话,某个商品排名上升时,它下面各名次的行会全部重新发出。去掉行号,就只需发送变化的商品的行。用 grep -c 数出两个文件的日志行数并进行比较。

报告——用数字表示日志的量

在 /root/flink/dedup/report.json 中写入 n_rows 和 n_orders(第 1 步)、update_pairs(last.out 中 -U 的个数)、shipped_now(board.out 最终状态中 shipped 的订单数)、topn_log_rows 和 norank_log_rows(两个 Top-3 输出的日志行数)、top_books(books 类别的第 1 名商品)。

全都从前面保存的输出中抄过来。日志行是像“| +I |”这样以 op 开头的行。带有变更日志的表,其最终值是该键最后一个 +I · +U。也请用第 1 步的两个数字确认 update_pairs 是否正确。