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

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

动态表与变更日志——流式结果是不断被改写的表

在 TT Lab 中继续学习

一句话总结

在 Flink SQL 中,流是不断追加行的表(动态表),其上的查询结果也是一张表。每当结果表发生变化,引擎就会输出变更日志(+I 追加 · -U 撤回旧值 · +U 新值 · -D 删除)。哪个算子会产生哪种变更,决定了结果可以写入哪种 Sink。

为什么需要它

批处理 SQL 在输入全部到齐后运行一次就结束,结果写一次就完事了。流里没有“全部到齐之后”。对流运行按用户统计点击数的查询,第一次点击到来时回答 u04 → 1,第二次点击到来时就必须修正这个答案。已经输出的答案该如何修正——这就是流式 SQL 的核心问题。

Flink 给出的答案是借用数据库的物化视图(materialized view)。官方文档(Dynamic Tables)是这样总结的:数据库中的表由 INSERT、UPDATE、DELETE 的流(changelog stream)构成,物化视图接收这个流,并持续改写结果。把流看作表,把连续查询看作视图维护,就能把 SQL 的语义原样搬到流上。而且文档有一个承诺——连续查询的结果,在任何时刻,都与对该时刻的输入快照运行同一个批查询所得的结果语义相同。本模块的实验就从亲自验证这个承诺开始。

工作原理

点击流的每一行都以追加(INSERT)的方式进入动态表。GroupAggregate 按用户统计点击数并改写结果表,再把这些变化以日志形式输出。第一次见到的用户是 +I,已有的用户是 -U(旧值)和 +U(新值)两行。撤回流两行都发送,带键的 upsert Sink 只接收 +U 一行。像文件 Sink 这样只能追加的 Sink 无法接收 -U,会在计划阶段被拒绝

流 → 表。 Source 的每条记录都被解释为对结果表的 INSERT。从文件或日志读取的点击,是只追加(append-only)的表。

表 → 表。 每当输入表发生变化,查询就会改写结果表。过滤(WHERE dwell > 60)或投影是一行输入对应一行结果,到此为止,所以结果也只追加。相比之下,GROUP BY user_id 的 COUNT(*) 每来一行相同键的行,就必须修改已输出的结果行。文档把前者称为 append 查询,后者称为 update 查询,并指出 update 查询为了修正已输出的结果,需要持有更多状态。

表 → 流。 把结果表的变化输出到外部时,必须选择编码方式。文档列举了三种。

编码 发送什么 接收方需要什么
append-only 只发送追加的行 无——只有结果只追加时才可行
retract INSERT 是追加,DELETE 是撤回,UPDATE 是先撤回(旧行)再追加(新行)共两条消息 无——拿到旧行的值就原样删除
upsert INSERT、UPDATE 是一条 upsert 消息,DELETE 是删除 唯一键——按键查找并覆盖

sql-client 表模式输出中看到的 op 列,正是这个变更类型。+I 是追加,-U 是更新前值的撤回,+U 是更新后的值,-D 是删除。如果是按用户 COUNT,新用户输出 +I 一行,已有用户输出 -U、+U 两行。如果输入只追加且并行度为 1,这份日志的顺序和条数完全由输入顺序决定——这也是实验评分器用 Python 逐行复现日志并进行核对的依据。

在聚合之上再叠加聚合,还会产生删除。统计“点击了 n 次的用户有多少”时,某个用户从 1 号栏移到 2 号栏,内层聚合的 -U 会让外层的 1 号栏减一。当这一栏变成 0,结果行本身就应该消失,于是会出现 -D。

每个算子会输出哪种变更,EXPLAIN CHANGELOG_MODE 会在计划里标出来。

GroupAggregate(groupBy=[user_id], ..., changelogMode=[I,UB,UA])   ← 갱신 전·후를 다 낸다
GroupAggregate(groupBy=[user_id], ..., changelogMode=[I,UA])      ← 같은 집계, 업서트 싱크 앞
Calc(select=[user_id, url], where=[>(dwell, 60)], changelogMode=[I])

同样的聚合却有两种模式,原因是优化器会从 Sink 能接收什么出发,反向向上推导。没有键的 print Sink 需要撤回才能删除旧行,所以要求 UB;声明了 PRIMARY KEY 的 Sink 按键覆盖即可,所以去掉 UB。消息数几乎减少了一半。反过来,如果把 update 查询放进像 filesystem 这样无法修改已写入行的 Sink,在作业启动之前就会被以“doesn't support consuming update changes”拒绝。如果仍想把日志保存到文件,可以用 TO_CHANGELOG 把变更类型提取成列,把所有行都变成追加(文档中的 Changelog Conversion)。

在现场相遇的样子

最常见的事故是仪表板上的数字翻了一倍。原因要么是接收撤回流的一方忽略了 -U 而只累加 +U,要么是把日志原样贴到了只追加的存储里。如果不看 op 而只数行,更新就变成了重复。首先要成对确认结果表的定义(键是什么)与 Sink 的写入方式(覆盖还是追加)。

第二种是“流式结果与批处理不同”的报告。比较最终状态,通常是相同的。不同的是中间日志被谁以何种方式接收。不过,结果相同的保证只在输入相同时成立。如果查询里有会按时间丢弃行的部分(下一个模块的水位线),或者设置了按时间清理状态,就可能与批处理不同。

第三种是 upsert Sink 的键与结果表的键不一致。聚合键是 user_id,却把 Sink 的 PRIMARY KEY 定成了别的列,那么 Sink 收到的变更按 Sink 键来看,顺序可能会被打乱。文档(Configuration 中的 table.exec.sink.upsert-materialize)指出,出现这种打乱时,优化器会在 Sink 之前插入 upsert materialize 算子——它按键把收到的行保存为状态,再以正确的 upsert 顺序重新输出。相当于又多了一份状态,所以原则是一开始就让 Sink 的键与结果表的唯一键保持一致。

下一项实验要做什么

分别以批处理和流处理方式运行按用户统计 120 次点击的查询,确认最终状态相同,并确认流式日志有 229 行。观察聚合之上的聚合中出现 -D,并用 EXPLAIN CHANGELOG_MODE 比较聚合与过滤的变更类型。把聚合放进文件 Sink 遭到拒绝后,改用 TO_CHANGELOG 写入文件,再看计划在没有键的 Sink 和有键的 Sink 之前有何不同。最后在报告中写出:如果把同样的日志以 upsert 方式发送,消息会有多少条。