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

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

去重与 Top-N——同一个 ROW_NUMBER,不同的变更日志

在 TT Lab 中继续学习

一句话总结

Flink SQL 的去重和 Top-N 都是给 ROW_NUMBER() OVER (PARTITION BY ... ORDER BY ...) 加上行号条件的同一种模式。然而,保留第一行时,结果是 append-only;保留最后一行或按排名输出时,得到的是会撤回并修正已输出行的变更日志。究竟是哪一种,决定了下游 Sink 和聚合的行为。

为什么需要它

官方文档举的例子就是现实中的样子。如果上游 ETL 无法保证端到端的精确一次,故障恢复时同一条记录会两次进入 Sink。在这种状态下做 SUM、COUNT,数字就会虚高。所以在分析之前必须先把重复清除掉。

不过,“重复”混着两种含义。一种是同一个事件来了两次(重发)——只保留最先来的那一条即可。另一种是同一对象的状态变了很多次(订单从创建 → 支付 → 发货)——要保留当前状态,也就是最后一条。这两种需求在 SQL 里只差 ASC 与 DESC 一个词,但在引擎内部做的事完全不同。因为第一行输出一次就结束了,而最后一行则是“到目前为止的最后一行”,会不断变化。

工作原理

左边是订单 o0018 的三条事件(created、paid、同一条 paid 的重发)进入时,两个算子输出的日志。保留第一行时,把 created 作为 +I 输出一次就结束。保留最后一行时,先把 created 作为 +I 输出,paid 到来后输出 -U created 和 +U paid 一对,重发的 paid 到来时也会再输出 -U paid 和 +U paid 一对。右边是文档中的 Top-N 示例:排名第 9 的商品升到第 1 时,输出行号的查询会把第 1 名到第 9 名共九行全部重新发送,而去掉行号后只发送发生变化的那一个商品

必须严格遵守模式。 文档把去重定义为 ROW_NUMBER()、PARTITION BY 키、ORDER BY 시간 속성(占位符依次为分区键、时间属性)以及外层的 WHERE rownum = 1,并指出必须严格遵循这个形态,优化器才能识别。ORDER BY 必须是时间属性(处理时间或事件时间),ASC 保留第一行,DESC 保留最后一行。从理论上说,去重是 N 为 1 且按时间排序的 Top-N 的特例。

从执行计划上看,两者是这样分开的。EXPLAIN CHANGELOG_MODE 的物理计划把这个算子写成 Rank(strategy=[AppendFastStrategy], rankRange=[rankStart=1, rankEnd=1], ...),并在行尾附上输出的变更类型。在执行计划中,同一位置会变成 Deduplicate 节点。

물리 계획  Rank(... orderBy=[ROWTIME ts ASC] ...)   changelogMode=[I]
실행 계획  Deduplicate(keep=[FirstRow], key=[order_id], order=[ROWTIME], outputInsertOnly=[true])

물리 계획  Rank(... orderBy=[ROWTIME ts DESC] ...)  changelogMode=[I,UA,D]
실행 계획  Deduplicate(keep=[LastRow], key=[order_id], order=[ROWTIME], outputInsertOnly=[false])

保留第一行只在状态里为每个键记下“已经见过”的标记,从第二条开始就丢弃。结果是 insert-only,所以可以直接写入文件这类 append-only Sink。保留最后一行则在状态里为每个键保存到目前为止的最后一行,新的行到来时,先撤回之前输出的行,再输出新的行。实测中输入了 622 条事件、240 个订单,得到了 240 个 +I 和 382 对 -U/+U。382 = 622 − 240,也就是同一个键从第二条起每一条产生一对。即使是完全相同的行被重发,也会输出一对。按处理时间(PROCTIME())排序时,同样的输入在 41 条重发中有 18 条没有产生成对的输出——由于这种排序依赖到达时刻,这类数量会随运行环境而变化,所以实验以事件时间为准。

Top-N 是行号条件为 <= N 的同一种模式。文档指出,Top-N 属于结果更新型,前 N 名发生变化时,会以撤回和更新的形式发送变化的行。如果输入本身就是更新的(销量 SUM 之上的排名),执行计划中会出现 Rank(strategy=[RetractStrategy], ...)。这里重要的选择是是否输出行号。行号列会成为结果唯一键的一部分,所以正如文档的例子,第 9 名升到第 1 名时,第 1–9 名的行会全部重新发出。如果在外层 SELECT 中去掉行号,就只需要发送变化的那一个商品。实测中,同一份销售文件按类别取 Top-3,输出行号时日志为 3,419 行,去掉行号时为 1,523 行,最终结果相同。

在保留最后一行之上再叠加 GROUP BY status,两个算子就咬合在一起了。当某个订单从 created 变为 paid,去重会输出 -U created 和 +U paid,聚合接收后在 created 栏减一、在 paid 栏加一。如果不去重而直接统计事件,同一个订单会同时被计入多个栏。

上一个模块在 temporal join 中创建版本视图,用的正是这种保留最后一行。窗口级 Top-N 已在窗口模块中讲过——那里只在窗口关闭时输出一次,所以没有撤回。

在现场相遇的样子

最常见的错误是“把去重结果写入文件或 append-only 主题”。保留第一行没有问题,但一旦改成保留最后一行,结果就变成变更日志,无法写入不接受更新的 Sink(这正是在动态表模块中见过的拒绝)。首先要从需求上分清到底需要哪一种。如果是清除重发,用第一行;如果需要当前状态,用最后一行。

第二种是 Top-N 的日志暴增。把实时排行榜写入键值存储时,如果写入量是预期的好几倍,首先要看是否输出了行号。如果前端可以自己给出名次,仅仅去掉行号就能大幅减少写入。

第三种是排序依据。按处理时间去重,结果取决于到达顺序,每次重新运行都可能不同。如果是需要重新处理后得到相同答案的管道,就使用事件时间。

下一项实验要做什么

定义订单状态事件和销售记录,统计行数和订单数。运行保留第一行和保留最后一行,确认日志数量有何不同,并用 EXPLAIN 比较两份计划的变更日志模式。在保留最后一行之上统计各状态的订单情况,再分别带行号和不带行号运行按类别的 Top-3,比较日志量,最后整理成报告。