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

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

读懂执行计划 — 一条 SQL 会变成多少个任务

在 TT Lab 中继续学习

一句话总结

EXPLAIN 会显示 SQL 被转换成哪些算子,而算子之间的 Exchange 就是通过网络 shuffle 数据的位置。不经过 shuffle 直接相连(FORWARD)的算子,会通过算子链(chaining)合并进同一个任务。因此,数出计划中的 Exchange 个数,就能提前知道作业会有几个顶点;两阶段聚合、distinct 拆分这类调优是否真正生效,也要通过计划来确认。

为什么需要它

收到“作业很慢”的反馈时,通常会先调高并行度。但如果瓶颈是集中在某一个键上的聚合,即使调高并行度,负责那个键的一个子任务仍然要承接全部数据。也可能是过滤条件没有下推到 Source,以致读取了所有行;或者算子链断开,任务之间的传递增多了。这些原因在吞吐量曲线上看起来一模一样,而在计划里看起来则完全不同。

调优选项也是如此。并不是打开 table.exec.mini-batch.enabled 就会出现两阶段聚合,只有当优化器判断条件满足时,计划才会改变。本模块的要点就是:不要相信配置文件,要看计划。

工作原理

左边是把 EXPLAIN 执行计划从下(Source)向上(Sink)画出来的样子。TableSourceScan、Calc 之后是 Exchange hash[page],其上是 GroupAggregate、Calc、Writer。中间是启用算子链的作业,以 Exchange 为界,合并成两个顶点——Source→Calc 和 GroupAggregate→Calc→Writer。右边是关闭算子链的作业,五个算子各自成为一个顶点。下方的色带表示:启用 mini-batch 后,Exchange 下方会出现 LocalGroupAggregate,上方的聚合变成 GlobalGroupAggregate

EXPLAIN 的三个部分。 EXPLAIN PLAN FOR <질의>(占位符为查询)会输出三个部分(官方文档中的 EXPLAIN Statements)。

部分 内容
== Abstract Syntax Tree == 对 SQL 原样转换的逻辑计划。LogicalFilter、LogicalAggregate
== Optimized Physical Plan == 应用规则之后的物理计划。过滤下推到 Source,聚合变成 GroupAggregate
== Optimized Execution Plan == 实际会生成的算子树

树中上方靠近 Sink,下方靠近 Source。对本实验的查询以流式方式运行 EXPLAIN,会得到 GroupAggregate ← Exchange(distribution=[hash[page]]) ← Calc ← TableSourceScan,并且 Source 那一行带有 filter=[>(amount, 0)]——这表示 WHERE 已经下推到了 Source。Exchange 是为了把相同键的行发送到同一个子任务而按哈希 shuffle 的位置,聚合一定位于它的上方。

EXPLAIN 还可以附加一些细节。ESTIMATED_COST 为每个节点附上估算行数和累计成本,CHANGELOG_MODE 附上变更日志的类型(第 2 个模块),PLAN_ADVICE 附上风险警告和调优建议,JSON_EXECUTION_PLAN 则以 JSON 形式附上算子图。在这个 Pod 上给 GROUP BY 加上 PLAN_ADVICE,会出现提示“试试打开 local-global two-phase”的 [ADVICE]。估算成本是针对没有统计信息的文件 Source 的假设值,不能当作绝对量来读。

算子链与顶点。 JSON 执行计划中的节点是一个个算子,节点之间写着 ship_strategy(FORWARD · HASH 等)。以 FORWARD 相连且并行度相同的算子会成为链,在同一个线程中运行(第 1 个模块)。所以对于线性管道,顶点数 = 非 FORWARD 的传递数 + 1。用 pipeline.operator-chaining.enabled = false 关闭后,每个算子都会成为一个顶点——这是调试时为了查看各算子指标而临时关闭的开关,生产环境中不应该关着。

两阶段聚合。 mini-batch 是先把输入短暂攒起来,使每个键对状态的访问减少为一次的机制(Performance Tuning)。启用 mini-batch 后,优化器可以把聚合拆成 LocalGroupAggregate(shuffle 之前,在每个子任务内预先合并)和 GlobalGroupAggregate(shuffle 之后再合并)。文档把它比作 MapReduce 的 Combine + Reduce。即使有集中在一个键上的数据,穿过 Exchange 的也只是预先合并好的累加器。在这个 Pod 上,只要设置 mini-batch 的三个配置(agg-phase-strategy 保持默认值 AUTO),计划中就会出现 MiniBatchAssigner 和 Local/Global。

distinct 拆分。 COUNT(DISTINCT user_id) 即使预先合并也很难缩小——因为累加器最终就是一份用户列表。打开 table.optimizer.distinct-agg.split.enabled 后,优化器会把查询改写成两层。第一层把 MOD(HASH_CODE(user_id), 버킷 수)(占位符为桶数)加到键上再 shuffle(PARTIAL),第二层再按原来的键 shuffle 并合计(FINAL)。桶数的默认值是 1024。

在现场相遇的样子

最常用的场景是“过滤条件在哪里生效”。如果 Source 支持过滤下推,TableSourceScan 那一行会带上 filter=[…]。如果没有带上,所有行都会从 Source 出来,到 Calc 才被丢弃。

第二种是热点键。如果仪表板上只有一个子任务很忙,就去看计划里 Exchange(distribution=[hash[…]]) 的键。如果该键的分布是倾斜的,答案就不是并行度,而是两阶段聚合或 distinct 拆分。打开之后,一定要用 EXPLAIN 确认是否出现了 Local/Global 或 PARTIAL/FINAL。

第三种是顶点数与预期不符。并行度不同的算子之间、shuffle 的位置之间,算子链会断开。/jobs/<jid> 的顶点名称中会出现用 -> 连接的算子列表,把它与 JSON 执行计划的传递方式并排对照,就能立刻看出在哪里断开了。

下一项实验要做什么

对点击文件按页面聚合的查询做 EXPLAIN,得到三个部分,并把 Exchange 的上方和下方写进 JSON 答案。拿到附带成本、建议的 EXPLAIN 和 JSON 执行计划之后,把同一个 INSERT 分别在启用和关闭算子链的情况下各运行一次,使顶点数与 JSON 中的传递方式对上。最后拿到 mini-batch 两阶段聚合和 COUNT(DISTINCT) 拆分的计划,把数字汇总成报告。