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

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

读计划并算出顶点数

在 TT Lab 中继续学习

目标

用 EXPLAIN 获取同一个 GROUP BY 查询的执行计划,读出算子和 Exchange 的位置,并用算子链和传递方式解释真实作业的顶点数。确认 mini-batch 两阶段聚合和 distinct 拆分在计划中是怎样出现的。

为什么重要

作业慢的原因——没有下推到 Source 的过滤、集中在一个键上的聚合、断开的算子链——在吞吐量曲线上看起来一样,在计划里则各不相同。调优配置也要计划发生变化才算生效。本实验的评分器不会查询集群。它只读取你保存的 EXPLAIN 输出原文和 REST 响应,并把顶点数与 JSON 执行计划的传递方式交叉核对。

步骤

  1. 运行 flink-up 后,在 /root/flink/plan/ddl.sql 中创建 Source clicks 和 blackhole Sink page_stats,并把 /root/flink/plan/explain.sql 的 EXPLAIN PLAN FOR 输出保存到 /root/flink/plan/explain.out。
  2. 读取 explain.out 的执行计划,在 /root/flink/plan/shuffle.json 中写入 exchange、above、below、filter_pushed_into_scan。
  3. 把使用 EXPLAIN ESTIMATED_COST, PLAN_ADVICE 的 /root/flink/plan/advice.sql 的输出保存到 /root/flink/plan/advice.out。
  4. 把使用 EXPLAIN JSON_EXECUTION_PLAN INSERT INTO page_stats … 的 /root/flink/plan/json.sql 的输出保存到 /root/flink/plan/json.out。
  5. 运行以作业名称 flk-plan-chained 把同一个 INSERT 跑完的 /root/flink/plan/chained.sql,并把该作业的 /jobs/<jid> 保存到 /root/flink/plan/chained-job.json。
  6. 运行关闭算子链、以作业名称 flk-plan-unchained 运行的 /root/flink/plan/unchained.sql,并把 /jobs/<jid> 保存到 /root/flink/plan/unchained-job.json。
  7. 把设置 mini-batch 三个配置之后对同一个 SELECT 做 EXPLAIN 的 /root/flink/plan/twophase.sql 的输出保存到 /root/flink/plan/twophase.out。
  8. 把打开 distinct 拆分(桶数 64)后对 COUNT(DISTINCT user_id) 做 EXPLAIN 的 /root/flink/plan/distinct.sql 的输出保存到 /root/flink/plan/distinct.out,并写出汇总各数字的 /root/flink/plan/report.json。

参考

获取 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。