读计划并算出顶点数
目标
用 EXPLAIN 获取同一个 GROUP BY 查询的执行计划,读出算子和 Exchange 的位置,并用算子链和传递方式解释真实作业的顶点数。确认 mini-batch 两阶段聚合和 distinct 拆分在计划中是怎样出现的。
为什么重要
作业慢的原因——没有下推到 Source 的过滤、集中在一个键上的聚合、断开的算子链——在吞吐量曲线上看起来一样,在计划里则各不相同。调优配置也要计划发生变化才算生效。本实验的评分器不会查询集群。它只读取你保存的 EXPLAIN 输出原文和 REST 响应,并把顶点数与 JSON 执行计划的传递方式交叉核对。
步骤
- 运行
flink-up后,在 /root/flink/plan/ddl.sql 中创建 Sourceclicks和 blackhole Sinkpage_stats,并把 /root/flink/plan/explain.sql 的EXPLAIN PLAN FOR输出保存到 /root/flink/plan/explain.out。 - 读取 explain.out 的执行计划,在 /root/flink/plan/shuffle.json 中写入
exchange、above、below、filter_pushed_into_scan。 - 把使用
EXPLAIN ESTIMATED_COST, PLAN_ADVICE的 /root/flink/plan/advice.sql 的输出保存到 /root/flink/plan/advice.out。 - 把使用
EXPLAIN JSON_EXECUTION_PLAN INSERT INTO page_stats …的 /root/flink/plan/json.sql 的输出保存到 /root/flink/plan/json.out。 - 运行以作业名称
flk-plan-chained把同一个 INSERT 跑完的 /root/flink/plan/chained.sql,并把该作业的/jobs/<jid>保存到 /root/flink/plan/chained-job.json。 - 运行关闭算子链、以作业名称
flk-plan-unchained运行的 /root/flink/plan/unchained.sql,并把/jobs/<jid>保存到 /root/flink/plan/unchained-job.json。 - 把设置 mini-batch 三个配置之后对同一个 SELECT 做 EXPLAIN 的 /root/flink/plan/twophase.sql 的输出保存到 /root/flink/plan/twophase.out。
- 把打开 distinct 拆分(桶数 64)后对
COUNT(DISTINCT user_id)做 EXPLAIN 的 /root/flink/plan/distinct.sql 的输出保存到 /root/flink/plan/distinct.out,并写出汇总各数字的 /root/flink/plan/report.json。
参考
- 源数据
/opt/lab/fixtures/data/plan_clicks.csv的列:click_id BIGINT, user_id STRING, page STRING, amount INT, click_time TIMESTAMP(3)(没有表头的 CSV)。 - 所有步骤的查询都相同:
SELECT page, COUNT(*) AS n_views, SUM(amount) AS revenue FROM clicks WHERE amount > 0 GROUP BY page。views是保留字,不能用作别名。 - 不想每次都重复表定义,可以用
sql-client.sh -i ddl.sql -f 파일.sql > 파일.out 2>&1(占位符均为文件名,-i 是先运行的初始化文件)。sql-client 每行只接受一条语句。 - EXPLAIN 的输出会以多行形式打印在表的一个单元格中。执行计划的上方靠近 Sink,下方靠近 Source。
- 常见错误:用
SET 'execution.runtime-mode' = 'batch'运行时,聚合算子的名称会不同——本实验以默认(流式)方式进行。在 INSERT 结束之前获取/jobs/<jid>,状态是 RUNNING(table.dml-sync)。 - 官方文档:EXPLAIN · Performance Tuning · Table 配置 · Flink Architecture · REST API
获取 EXPLAIN 的三个部分
运行 flink-up 后,在 /root/flink/plan/ddl.sql 中写入读取 plan_clicks.csv 的 clicks 和 CREATE TABLE page_stats (page STRING, n_views BIGINT, revenue BIGINT) WITH ('connector' = 'blackhole'),在 /root/flink/plan/explain.sql 中写入 EXPLAIN PLAN FOR + 参考中的查询,再用 sql-client.sh -i ddl.sql -f explain.sql > explain.out 2>&1 生成 /root/flink/plan/explain.out。
EXPLAIN 不会运行作业,只返回计划。输出中必须能看到 Abstract Syntax Tree · Optimized Physical Plan · Optimized Execution Plan 三个部分。初始化文件(-i)中,两条 CREATE TABLE 语句各占一行。
找到 shuffle 的位置
读取 explain.out 的执行计划部分,在 /root/flink/plan/shuffle.json 中写入 exchange(Exchange 的 distribution 值,例如 hash[…])、above(Exchange 上方的算子名称列表)、below(下方的算子名称列表,从上往下)、filter_pushed_into_scan(Source 行上是否带有 filter=[…],真/假)。
算子名称是括号前面的单词(GroupAggregate、Calc 等)。树的上方靠近 Sink。如果 WHERE 已经下推到 Source,TableSourceScan 行里会看到 filter=[…]。评分器会与你的 explain.out 核对。
附加成本估算和建议
运行使用 EXPLAIN ESTIMATED_COST, PLAN_ADVICE + 同一个查询的 /root/flink/plan/advice.sql,并把输出保存到 /root/flink/plan/advice.out。必须能看到成本(cumulative cost)和 advice[1]: [ADVICE] 行。
加上 PLAN_ADVICE 后,物理计划部分的标题会变成 'With Advice',末尾会附上建议行。请读一下它建议打开什么——第 7 步会用到。估算成本是针对没有统计信息的 Source 的假设值。
从 JSON 执行计划中查看传递方式
运行使用 EXPLAIN JSON_EXECUTION_PLAN INSERT INTO page_stats + 同一个查询的 /root/flink/plan/json.sql,并把输出保存到 /root/flink/plan/json.out。必须能看到包括 Sink(Writer)在内的节点和 ship_strategy。
要对 INSERT 而不是 SELECT 做 EXPLAIN,才会得到与实际运行的作业相同的图。== Physical Execution Plan == 之下的 nodes 数组就是算子,predecessors 的 ship_strategy 是从前一个算子接收数据的方式。请数一下 HASH 出现了几次。
数出启用算子链的作业的顶点
运行写入了 SET 'pipeline.name' = 'flk-plan-chained';、SET 'table.dml-sync' = 'true'; 和 INSERT INTO page_stats + 同一个查询的 /root/flink/plan/chained.sql,然后把该作业的 /jobs/<jid> 保存到 /root/flink/plan/chained-job.json。顶点数必须等于 json.out 中非 FORWARD 的传递数 + 1。
打开 table.dml-sync 后,sql-client 会等到 INSERT 作业结束,这样就能拿到已结束(FINISHED)作业的响应。顶点名称中用 -> 连接的算子,就是合并成一个任务的链。
关闭算子链后再数一遍
在 chained.sql 的基础上加入 SET 'pipeline.operator-chaining.enabled' = 'false';,并把作业名称改为 flk-plan-unchained,运行得到的 /root/flink/plan/unchained.sql,然后把 /jobs/<jid> 保存到 /root/flink/plan/unchained-job.json。顶点数必须等于 json.out 中的算子(节点)数。
关闭算子链后,以 FORWARD 相连的算子也会各自成为顶点(任务)。结果相同,只是任务之间的传递增多了。如果不改名称,在作业列表中会与前一步的作业混淆。
用 mini-batch 生成两阶段聚合
用 SET 设置 table.exec.mini-batch.enabled = true、table.exec.mini-batch.allow-latency = 1 s、table.exec.mini-batch.size = 1000,然后运行写入 EXPLAIN PLAN FOR + 同一个查询的 /root/flink/plan/twophase.sql,并把输出保存到 /root/flink/plan/twophase.out。执行计划中必须能看到 GlobalGroupAggregate ← Exchange ← LocalGroupAggregate,以及其下方的 MiniBatchAssigner。
第 3 步的建议所推荐的就是这些设置。只有打开 mini-batch 才会生成两阶段聚合——shuffle 之前先在每个子任务内预先合并(Local),shuffle 之后再合并(Global)。三个配置缺少任何一个,计划都不会改变。
distinct 拆分与报告
用 SET 设置 table.optimizer.distinct-agg.split.enabled = true、table.optimizer.distinct-agg.split.bucket-num = 64,并把写入 EXPLAIN PLAN FOR SELECT page, COUNT(DISTINCT user_id) AS users FROM clicks GROUP BY page; 的 /root/flink/plan/distinct.sql 的输出保存到 /root/flink/plan/distinct.out。然后在 /root/flink/plan/report.json 中写入 json_nodes、forward_edges、hash_edges(json.out),chained_vertices、unchained_vertices(两个作业的响应),distinct_buckets(distinct.out 中的桶数)。
拆分生效后,GroupAggregate 会变成 PARTIAL 和 FINAL 两层,它们之间和下方各有一个 Exchange,下方的 Calc 中会出现 MOD(HASH_CODE(user_id), 桶数)。报告中的数字都可以从你保存的文件中数出来——也请确认算子链启用时的顶点数是否等于 HASH 数 + 1。