同一个 GROUP BY:批处理每键一行,流式变更日志 229 行
目标
以批处理和流处理方式运行同一个聚合,确认最终状态相同,而流处理会输出变更日志。通过执行计划和拒绝错误,读出变更类型(+I · -U · +U · -D)是如何由算子和 Sink 决定的。
为什么重要
流式结果不是写一次就结束的值,而是会被不断改写的表。如果接收这种改写的一方无法处理撤回(-U),数字就会虚高,而只追加的存储则根本无法接收。必须学会事先从计划中读出哪些查询会产生更新、哪些 Sink 能接收它,才能在设计阶段避免事故。本实验的评分器不会查询集群——它读取你保存的 sql-client 输出和 Sink 文件,并逐行复现变更日志,从源 CSV 出发进行核对(输入只追加且并行度为 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。 - 创建以流处理模式运行同一查询的 /root/flink/dynamic/stream.sql,把输出保存到 /root/flink/dynamic/stream.out。
- 创建以流处理方式运行按点击次数(
clicks)统计用户数(users)的聚合之上的聚合的 /root/flink/dynamic/nested.sql,把输出保存到 /root/flink/dynamic/nested.out。 - 在 /root/flink/dynamic/explain.sql 中写两条
EXPLAIN CHANGELOG_MODE语句(按用户的 COUNT 聚合;只选出dwell > 60的行的user_id, url的过滤),并把输出保存到 /root/flink/dynamic/explain.out。 - 在 /root/flink/dynamic/fs-sink.sql 中创建 filesystem 连接器表
user_clicks_fs(user_id STRING, clicks BIGINT),并写出把按用户的 COUNTINSERT INTO进去的语句,然后把包含拒绝错误的输出保存到 /root/flink/dynamic/fs-sink.out。 - 在 /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。 - 在 /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。 - 在 /root/flink/dynamic/report.json 中写入
users、retract_messages、upsert_messages、nested_deletes。
参考
- 源列:
click_id BIGINT, user_id STRING, url STRING, dwell INT, ts TIMESTAMP(3)(没有表头的 CSV,文件顺序 = 到达顺序)。 - 运行 SQL 文件:
sql-client.sh -f 파일.sql > 파일.out 2>&1(占位符均为文件名)。遇到出错的语句就会停止,输出中会留下[ERROR]。 - 流式结果表的最前面会多出
op列。因为读取的是有终点的文件,流处理作业读完文件后也会结束。 - 常见错误:第 6 步运行两次,Sink 目录中的文件会累积。重新运行时,请先清空目录。
INSERT默认是异步提交,所以作业结束之前 sql-client 就退出了——设置SET 'table.dml-sync' = 'true';就会等到结束。 - 常见错误:一旦打开
SET 'table.exec.mini-batch.enabled'这类配置,中间日志会被合并,行数就会不同。本实验保持默认配置。 - 官方文档:Dynamic Tables · EXPLAIN · Changelog Conversion · FileSystem · Print · BlackHole
用批处理一次统计
用 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 |”的内容。