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

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

同一个 GROUP BY:批处理每键一行,流式变更日志 229 行

在 TT Lab 中继续学习

目标

以批处理和流处理方式运行同一个聚合,确认最终状态相同,而流处理会输出变更日志。通过执行计划和拒绝错误,读出变更类型(+I · -U · +U · -D)是如何由算子和 Sink 决定的。

为什么重要

流式结果不是写一次就结束的值,而是会被不断改写的表。如果接收这种改写的一方无法处理撤回(-U),数字就会虚高,而只追加的存储则根本无法接收。必须学会事先从计划中读出哪些查询会产生更新、哪些 Sink 能接收它,才能在设计阶段避免事故。本实验的评分器不会查询集群——它读取你保存的 sql-client 输出和 Sink 文件,并逐行复现变更日志,从源 CSV 出发进行核对(输入只追加且并行度为 1 时,日志的顺序也是确定的)。

步骤

  1. 用 flink-up 启动集群,在 /root/flink/dynamic/batch.sql 中写一条 SQL:以批处理模式读取 /opt/lab/fixtures/data/dynamic_clicks.csv,按 user_id 输出 clicks(条数)和 dwell(dwell 之和),然后把输出保存到 /root/flink/dynamic/batch.out。
  2. 创建以流处理模式运行同一查询的 /root/flink/dynamic/stream.sql,把输出保存到 /root/flink/dynamic/stream.out。
  3. 创建以流处理方式运行按点击次数(clicks)统计用户数(users)的聚合之上的聚合的 /root/flink/dynamic/nested.sql,把输出保存到 /root/flink/dynamic/nested.out。
  4. 在 /root/flink/dynamic/explain.sql 中写两条 EXPLAIN CHANGELOG_MODE 语句(按用户的 COUNT 聚合;只选出 dwell > 60 的行的 user_id, url 的过滤),并把输出保存到 /root/flink/dynamic/explain.out。
  5. 在 /root/flink/dynamic/fs-sink.sql 中创建 filesystem 连接器表 user_clicks_fs(user_id STRING, clicks BIGINT),并写出把按用户的 COUNT INSERT INTO 进去的语句,然后把包含拒绝错误的输出保存到 /root/flink/dynamic/fs-sink.out。
  6. 在 /root/flink/dynamic/to-changelog.sql 中,把按用户的 COUNT 视图用 TO_CHANGELOG 转换后,写入 filesystem Sink(路径 /root/flink/dynamic/changelog,列 change, user_id, clicks),并把输出保存到 /root/flink/dynamic/to-changelog.out。
  7. 在 /root/flink/dynamic/upsert.sql 中创建 print 连接器表 retract_sink,以及带有 PRIMARY KEY (user_id) NOT ENFORCED 的 blackhole 连接器表 upsert_sink,再分别对两个 Sink 运行把同一聚合写入的 EXPLAIN CHANGELOG_MODE INSERT INTO ...,并把输出保存到 /root/flink/dynamic/upsert.out。
  8. 在 /root/flink/dynamic/report.json 中写入 users、retract_messages、upsert_messages、nested_deletes。

参考

用批处理一次统计

用 flink-up 启动集群,在 /root/flink/dynamic/batch.sql 中写入 SET 'execution.runtime-mode' = 'batch';、读取源 CSV 的 CREATE TABLE,以及按 user_id 输出 clicks(COUNT(*))和 dwell(SUM(dwell))的 SELECT,并把输出保存到 /root/flink/dynamic/batch.out。

filesystem 连接器的 path 是 file:///opt/lab/fixtures/data/dynamic_clicks.csv,format 是 csv。结果列用 AS clicks、AS dwell 起别名。批处理的结果没有 op 列,每个用户一行。

用流处理运行同一个查询

创建以 SET 'execution.runtime-mode' = 'streaming'; 运行与第 1 步相同查询的 /root/flink/dynamic/stream.sql,并把输出保存到 /root/flink/dynamic/stream.out。把日志应用到最后得到的最终状态必须与批处理结果相同。

每来一行,结果表都会变化。第一次见到的用户是 +I 一行,已有用户是删除旧值的 -U 和新值 +U 两行。日志行数必须是 用户数 + 2 × (点击数 − 用户数)。配置保持默认。

在聚合之上的聚合中看到 -D

创建以流处理模式运行 SELECT clicks, COUNT(*) AS users FROM (사용자별 COUNT(*) AS clicks) GROUP BY clicks(占位符为按用户统计 COUNT(*) AS clicks 的子查询)的 /root/flink/dynamic/nested.sql,并把输出保存到 /root/flink/dynamic/nested.out。

当用户从点击 1 次的栏移到 2 次的栏时,内层聚合会输出 -U(1) 和 +U(2)。外层聚合收到 -U 后,会把 1 次栏的用户数减一,而当这一栏变成 0 时,结果行应当消失,所以输出 -D。请让内层查询的列名与 clicks 保持一致。

从计划中读出变更类型

在 /root/flink/dynamic/explain.sql 中写两条以 EXPLAIN CHANGELOG_MODE 开头的语句——按用户的 COUNT 聚合,以及只选出 dwell > 60 的行的 user_id, url 的过滤——并把输出保存到 /root/flink/dynamic/explain.out。

EXPLAIN 不会运行作业,只打印计划。加上 CHANGELOG_MODE 后,Optimized Physical Plan 的每个节点都会带上 changelogMode=[...]。I 是追加,UB 是更新前,UA 是更新后,D 是删除。只有过滤的查询,其节点只会输出什么呢?

把更新放进文件 Sink 而被拒绝

在 /root/flink/dynamic/fs-sink.sql 中创建 filesystem 连接器表 user_clicks_fs(user_id STRING, clicks BIGINT)(路径为 /root/flink/dynamic/fs-out,format 为 csv),写入 INSERT INTO user_clicks_fs SELECT user_id, COUNT(*) FROM clicks GROUP BY user_id;,然后把输出保存到 /root/flink/dynamic/fs-sink.out。以错误结束是正常的。

文件无法修改已写入的行。优化器会先询问 Sink 能接收哪些变更类型,如果无法接收聚合输出的 UB/UA,就会在提交作业之前拒绝。错误语句中写明了哪个 Sink 无法接收哪个节点的什么。

把变更类型提取成列并写入文件

在 /root/flink/dynamic/to-changelog.sql 中,把按用户的 COUNT 建成视图(user_id、clicks 两列),并把 TO_CHANGELOG(input => TABLE 뷰, op => DESCRIPTOR(change))(占位符为视图名)的结果 INSERT 到 filesystem Sink(路径 /root/flink/dynamic/changelog,列 change STRING, user_id STRING, clicks BIGINT,format 为 csv)。输出保存到 /root/flink/dynamic/to-changelog.out。

TO_CHANGELOG 会把更新日志的每一行变成追加(INSERT),并把原来的变更类型放进字符串列——所以只追加的文件 Sink 也能接收。INSERT 是异步的,必须开启 table.dml-sync,作业结束之后文件才会以已确定的状态保留下来。重新运行时,请先清空 Sink 目录。

在带键的 Sink 之前,-U 消失了

在 /root/flink/dynamic/upsert.sql 中创建 print 连接器表 retract_sink(user_id STRING, clicks BIGINT),以及在相同列上带有 PRIMARY KEY (user_id) NOT ENFORCED 的 blackhole 连接器表 upsert_sink。然后运行把按用户的 COUNT 分别写入各个 Sink 的两条 EXPLAIN CHANGELOG_MODE INSERT INTO ... 语句,并把输出保存到 /root/flink/dynamic/upsert.out。

print Sink 只是按收到的内容打印,所以要删除旧行就需要撤回(UB)。声明了键的 Sink 按键查找并覆盖即可,所以不需要 UB。请比较两份计划中 GroupAggregate 的 changelogMode。NOT ENFORCED 表示 Flink 不会检查键。

报告——撤回与 upsert 的消息数

在 /root/flink/dynamic/report.json 中以整数写入 users(用户数 = stream.out 中 +I 的个数)、retract_messages(stream.out 日志的总行数)、upsert_messages(如果把同样的结果发送给 upsert Sink,会发送的消息数)、nested_deletes(nested.out 中 -D 的个数)。

upsert Sink 按键覆盖,所以不需要接收更新前的值(-U)。想一想同样的日志中只会剩下哪些 op。数字数一数保存的输出的 op 列就能得到——可以用 grep 统计行首形如“| +I |”的内容。