欢迎光临

Spark Catalyst 优化器深度解析:从逻辑计划到物理计划的查询优化全流程

引言:为什么理解 Catalyst 对 Spark 开发者至关重要

在日常的 Spark SQL 开发中,我们写下一条 SQL 语句或一段 DataFrame API 调用,Spark 便能在海量数据上高效执行。这种”写起来简单、跑起来快”的体验,背后核心功臣便是 Catalyst 优化器。Catalyst 是 Spark SQL 的心脏,它将用户编写的查询转化为经过多重优化的物理执行计划,直接影响着作业的性能与资源消耗。

很多 Spark 开发者在遇到性能问题时,习惯于调整并行度、增加内存或修改 Shuffle 参数,却忽略了从查询计划层面寻找根因。理解 Catalyst 的工作流程,不仅能帮助你读懂

1
explain()

输出中的每个节点,更能让你写出”对优化器友好”的查询,从源头上避免性能陷阱。

本文将从 Catalyst 的整体架构出发,逐步拆解从 SQL 解析到物理计划生成的完整流程,结合大量代码示例和执行计划对比,让你真正掌握 Spark 查询优化的核心原理。

Spark Catalyst 优化器

Catalyst 总体架构:一棵树的旅程

Catalyst 的核心设计基于规则驱动(Rule-Based)和代价驱动(Cost-Based)的双重优化体系。整个优化流程可以看作一棵逻辑树经过多轮变换,最终生成物理执行树的过程。

Catalyst 的处理流程分为四个主要阶段:

  • 解析(Analysis):将未解析的逻辑计划转化为已解析的逻辑计划
  • 逻辑优化(Logical Optimization):基于规则的逻辑计划变换(谓词下推、列裁剪等)
  • 物理计划(Physical Planning):生成多个物理计划候选,基于代价选择最优
  • 代码生成(Code Generation):将物理计划编译为 Java 字节码

每个阶段的输入输出都是树形结构(TreeNode),这是 Catalyst 的基础数据模型。理解 TreeNode 的不可变(Immutable)和转换(Transform)特性,是理解 Catalyst 的第一步。

TreeNode:一切皆树

Spark SQL 中的表达式、查询计划、逻辑节点都是 TreeNode 的子类。TreeNode 的核心特征:


1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
// TreeNode 的核心操作
// 1. 不可变:所有变换返回新节点
val newPlan = oldPlan.transformDown {
  case Filter(condition, child) => Filter(optimizedCondition, child)
}

// 2. 模式匹配:利用 Scala 样例类做结构化匹配
plan match {
  case Filter(Literal(true, BooleanType), child) => child
  case Project(attrs, Filter(cond, child)) => ???
}

// 3. 递归变换:transformDown / transformUp
plan.transformDown { rule }
plan.transformUp { rule }

不可变性保证了优化规则可以安全地并行应用,而模式匹配让规则定义极其简洁。这种设计使得 Catalyst 的扩展性极强——添加一个新的优化规则只需编写一个 PartialFunction。

第一阶段:解析——从字符串到逻辑计划

当用户提交一条 SQL 或一段 DataFrame 代码时,Catalyst 首先将其转化为未解析的逻辑计划(Unresolved Logical Plan)。此时的计划只是一棵语法树,其中的表名、列名都是未验证的字符串。


1
2
3
4
5
6
7
8
-- 用户 SQL
SELECT name, age FROM users WHERE age > 30 JOIN orders ON users.id = orders.user_id

// 解析后生成的未解析逻辑计划(示意)
// 'UnresolvedRelation [users]
// 'UnresolvedRelation [orders]
// 'Filter ('age > 30)
// 'UnresolvedAttribute [name, age]

Analyzer 的职责是利用Catalog(元数据目录)将未解析的计划转化为已解析的逻辑计划(Resolved Logical Plan)。解析规则包括:

  • ResolveRelations:将表名绑定到实际的 Catalog 表
  • ResolveReferences:将列名绑定到对应表的 schema
  • ResolveFunctions:将函数名绑定到注册的 UDF 或内置函数
  • TypeCasts:插入必要的类型转换

解析失败的常见场景:


1
2
3
4
5
6
7
8
9
10
11
12
// 常见解析错误及原因
// 1. 列名歧义:JOIN 两表有同名列
spark.sql("SELECT id FROM t1 JOIN t2 ON t1.id = t2.id")
// Reference 'id' is ambiguous

// 2. 列不存在
spark.sql("SELECT non_exist FROM users")
// cannot resolve 'non_exist' given input columns

// 3. 类型不匹配
spark.sql("SELECT * FROM t WHERE string_col > 100")
// 需要隐式类型转换,或报错

开发者常遇到的”cannot resolve”错误,绝大多数源自这一阶段。排查方法是先打印未解析计划:


1
2
3
val df = spark.sql("SELECT name FROM users WHERE age > 30")
df.queryPlan.logical // 查看逻辑计划
df.queryPlan.analyzed // 查看解析后的计划

第二阶段:逻辑优化——规则驱动的计划变换

解析完成后,Catalyst 进入逻辑优化阶段。这是 Catalyst 最核心的部分,通过一系列优化规则(Optimization Rules)对逻辑计划进行等价变换,使查询在语义不变的前提下更高效。

逻辑优化规则分为两大类:基于规则的优化(RBO)和基于代价的优化(CBO)。其中 RBO 是 Catalyst 的传统强项,CBO 在 Spark 2.x 后逐步引入并增强。

逻辑优化流程

核心 RBO 规则详解

1. 谓词下推(Predicate Pushdown)

将 Filter 操作尽可能下推到数据源层,减少上游需要处理的数据量。这是最经典也最有效的优化之一。


1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
// 优化前
SELECT name FROM (SELECT * FROM users JOIN orders
  ON users.id = orders.user_id) WHERE age > 30

// 逻辑计划(优化前)
Filter (age > 30)
  |- Join (users.id = orders.user_id)
       |- Scan users [all columns]
       |- Scan orders [all columns]

// 优化后:Filter 下推到 users 表
Join (users.id = orders.user_id)
  |- Filter (age > 30)
       |- Scan users [all columns]
  |- Scan orders [all columns]

对于支持谓词下推的数据源(Parquet、ORC、Delta Lake 等),Filter 甚至可以下推到文件读取层,直接跳过不满足条件的文件块,效果极为显著。

2. 列裁剪(Column Pruning)

只读取查询真正需要的列,避免全列扫描。


1
2
3
4
5
6
7
8
9
// 优化前:读取 users 表所有列,但只用 name 和 age
Project [name, age]
  |- Filter (age > 30)
       |- Relation users [id, name, age, email, phone, address, ...]

// 优化后:只读取需要的列
Project [name, age]
  |- Filter (age > 30)
       |- Relation users [name, age]  // 列裁剪后

对于列式存储格式(Parquet),列裁剪可以减少 80% 以上的 I/O,因为实际场景中”SELECT *”查询远少于按需查询。

3. 常量折叠(Constant Folding)

在编译期计算常量表达式,避免运行时重复计算。


1
2
3
4
5
6
// 优化前
SELECT * FROM t WHERE 1 = 1 AND age > 30 + 5

// 优化后
SELECT * FROM t WHERE age > 35
// 1=1 被消除,30+5 被预计算为 35

4. 空值传播(Null Propagation)


1
2
3
4
5
6
// 优化前
SELECT NULL + age, NULL * price FROM users

// 优化后
SELECT NULL, NULL FROM users
// NULL 与任何值运算结果均为 NULL,直接替换

5. Filter/Project 合并与重排


1
2
3
4
5
6
7
8
9
10
11
12
// 优化前:两个连续的 Filter
Filter (age > 30)
  Filter (city = 'Beijing')
    Scan users

// 优化后:合并为一个 Filter
Filter (age > 30 AND city = 'Beijing')
  Scan users

// 更重要:选择性高的条件优先执行
Filter (city = 'Beijing' AND age > 30)
  Scan users

CBO:代价驱动的优化

Spark 2.x 引入了基于代价的优化器,通过表统计信息(行数、列的 NDV、数据大小等)做更智能的决策。CBO 的核心场景是Join 重排序。


1
2
3
4
5
6
7
8
9
10
// 开启 CBO
spark.conf.set("spark.sql.cbo.enabled", true)
spark.conf.set("spark.sql.statistics.histogram.enabled", true)

// 收集统计信息
spark.sql("ANALYZE TABLE large_table COMPUTE STATISTICS FOR ALL COLUMNS")

// CBO 生效的 Join 重排序示例
// 大表 JOIN 小表 → 优化为 小表 JOIN 大表
// 确保小表作为 Build 侧,减少内存消耗

CBO 的统计信息还可用于过滤条件的选择率估算,从而更准确地预估中间结果大小,指导 Join 策略选择。

第三阶段:物理计划——从逻辑到物理的跨越

逻辑优化完成后,Catalyst 需要将逻辑算子转化为可执行的物理算子。这一阶段最大的特点是一个逻辑算子可能对应多种物理实现,Catalyst 需要从中选择最优方案。

物理计划选择

Join 策略选择

Join 是物理计划选择中最关键的决策。Spark 提供了五种 Join 策略:

Join 策略 适用场景 内存要求 是否需要 Shuffle
Broadcast Hash Join 小表(< 10MB 默认) 低 不需要
Sort Merge Join 大表 JOIN 大表 中 需要
Shuffle Hash Join 中表 JOIN 大表 较高 需要
Broadcast Nested Loop 交叉 Join 或非等值 Join 低 不需要
Cartesian Product 无条件交叉连接 极高 不需要

1
2
3
4
5
6
7
8
9
// 手动指定 Join 策略的 Hint
spark.sql("""
  SELECT /*+ BROADCAST(small_table) */ *
  FROM large_table JOIN small_table ON large_table.id = small_table.id
""")

// 或通过 DataFrame API
import org.apache.spark.sql.functions.broadcast
large_df.join(broadcast(small_df), "id")

Join 策略的选择逻辑(简化版):

  • 如果有 Broadcast Hint 或表大小小于
    1
    spark.sql.autoBroadcastJoinThreshold

    (默认 10MB)→ Broadcast Hash Join

  • 如果没有等值条件 → Broadcast Nested Loop 或 Cartesian Product
  • 否则 → Sort Merge Join(默认大表场景)或 Shuffle Hash Join(如果一侧足够小)

分区策略选择

物理计划还需要决定数据的分区方式,这直接影响 Shuffle 的必要性。


1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
// 查看物理计划中的分区信息
df.explain(true)

// 常见的分区优化
// 1. 避免 Shuffle:利用已有的分区
spark.sql("SELECT * FROM t PARTITION BY (col)")

// 2. 控制分区数
spark.conf.set("spark.sql.shuffle.partitions", 200)

// 3. 利用 Bucket 避免重复 Shuffle
spark.sql("""
  CREATE TABLE orders (
    user_id BIGINT, order_id BIGINT, amount DOUBLE
  ) USING PARQUET
  CLUSTERED BY (user_id) INTO 32 BUCKETS
""")
// 当 JOIN 键与 bucket 键一致时,可以跳过 Shuffle

第四阶段:代码生成——Volcano 模型的终结者

传统的数据库执行引擎采用 Volcano 模型(迭代器模型),每个算子通过

1
next()

方法逐行传递数据。这种模型虽然优雅,但存在大量虚函数调用开销,无法被 JVM JIT 有效优化。

Spark 从 2.0 开始引入全阶段代码生成(Whole-Stage Code Generation),将多个物理算子融合为一个 Java 函数,大幅减少虚函数调用,提升 CPU 效率。


1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
// Volcano 模型(传统)
while (filter.hasNext()) {
  Row row = filter.next();  // 虚函数调用
  if (project.hasNext()) {
    Object val = project.next();  // 虚函数调用
  }
}

// Whole-Stage Code Generation(Spark)
while (scan.hasNext()) {
  InternalRow row = scan.next();
  // Filter + Project 内联在同一循环
  if (row.getInt(0) &gt; 30) {
    result.setValue(0, row.getString(1));
    return result;
  }
}

代码生成对性能的影响极大,在简单查询(如过滤、聚合)场景下可以带来 10 倍以上的性能提升。你可以通过执行计划中的

1
*

标记确认代码生成是否生效:


1
2
3
4
5
6
7
8
9
10
df.explain()
// *(1) Filter (age#1 &gt; 30)     ← 星号表示代码生成生效
// +- *(1) ColumnarToRow
//    +- Scan Parquet
//
// Exchange hashpartitioning(id#2, 200)   ← 无星号
// +- *(2) HashAggregate

// 关闭代码生成对比
spark.conf.set("spark.sql.codegen.wholeStage", false)

实战:读懂 explain() 输出

理解了 Catalyst 的四个阶段后,我们用实际案例来解读

1
explain()

的输出。


1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
val df = spark.table("orders")
  .filter("amount &gt; 100")
  .join(spark.table("users"), "user_id")
  .select("name", "amount")

df.explain(true)

// == Parsed Logical Plan ==
// Project [name#5, amount#3]
//   Join Inner, (user_id#2 = user_id#6)
//     Filter (amount#3 &gt; 100)
//       UnresolvedRelation [orders]
//     UnresolvedRelation [users]
//
// == Analyzed Logical Plan ==
// Project [name#5, amount#3]
//   Join Inner, (user_id#2 = user_id#6)
//     Filter (amount#3 &gt; cast(100 as double))
//       SubqueryAlias orders
//         Relation[order_id#1, user_id#2, amount#3]
//     SubqueryAlias users
//       Relation[user_id#6, name#5, age#7]
//
// == Optimized Logical Plan ==
// Project [name#5, amount#3]
//   Join Inner, (user_id#2 = user_id#6)
//     Filter (amount#3 &gt; 100.0)
//       Relation[order_id#1, user_id#2, amount#3]
//     Relation[user_id#6, name#5]  ← 列裁剪生效
//
// == Physical Plan ==
// *(1) Project [name#5, amount#3]
// +- *(1) SortMergeJoin [user_id#2 = user_id#6]
//    :- *(1) Sort [user_id#2 ASC]
//    :  +- Exchange hashpartitioning(user_id#2, 200)
//    :     +- *(1) Filter (amount#3 &gt; 100.0)
//    :        +- *(1) ColumnarToRow
//    :           +- FileScan parquet [user_id#2, amount#3]
//    +- *(2) Sort [user_id#6 ASC]
//       +- Exchange hashpartitioning(user_id#6, 200)
//          +- *(2) ColumnarToRow
//             +- FileScan parquet [user_id#6, name#5]

从上面的输出可以清晰看到 Catalyst 的完整工作流:从未解析计划到已解析计划,再到优化后的逻辑计划(注意列裁剪移除了

1
age

),最后到包含 SortMergeJoin + Exchange 的物理计划。

常见性能陷阱与 Catalyst 优化建议

1. 避免破坏谓词下推的操作

某些操作会阻止 Catalyst 进行谓词下推,导致全表扫描:


1
2
3
4
5
6
7
8
9
10
11
12
// 错误:Filter 在新增列上
val df = spark.table("logs")
  .select("timestamp", "level", "message")
  .withColumn("date", to_date("timestamp"))
  .filter("date &gt; '2024-01-01'")
// 无法下推到数据源

// 正确:优先过滤再转换
val df = spark.table("logs")
  .filter("timestamp &gt; '2024-01-01 00:00:00'")
  .withColumn("date", to_date("timestamp"))
  .select("date", "level", "message")

2. 注意 UDF 对优化的阻断

Python UDF 是 Catalyst 优化的”黑洞”——优化器无法穿透 UDF 内部逻辑进行下推或列裁剪。


1
2
3
4
5
6
7
8
9
10
11
12
13
14
// Python UDF 阻断优化
from pyspark.sql.functions import udf
from pyspark.sql.types import StringType

@udf(StringType)
def complex_logic(col):
    return processed_value

// 替代方案:使用内置函数组合
from pyspark.sql import functions as F
df.withColumn("result",
  F.when(F.col("status") == "active", F.upper(F.col("name")))
   .otherwise(F.lit("INACTIVE"))
)

3. 利用 AQE 弥补静态计划的不足

Catalyst 在计划阶段是基于静态信息做决策的,运行时实际情况可能不同。AQE(Adaptive Query Execution)在运行时根据 Shuffle Map 阶段的实际统计信息动态调整计划:


1
2
3
4
5
6
7
8
9
10
11
// 开启 AQE
spark.conf.set("spark.sql.adaptive.enabled", true)

// AQE 的三大优化:
// 1. 动态合并 Shuffle 分区
// 2. 动态切换 Join 策略(运行时发现小表 → Broadcast)
// 3. 动态优化数据倾斜(skewJoin)

spark.conf.set("spark.sql.adaptive.skewJoin.enabled", true)
spark.conf.set("spark.sql.adaptive.skewJoin.skewedPartitionFactor", 5)
spark.conf.set("spark.sql.adaptive.skewJoin.skewedPartitionThresholdInBytes", "256MB")

4. 善用 Hint 引导优化器

当 Catalyst 的自动选择不理想时,可以通过 Hint 显式指导优化器:


1
2
3
4
5
6
7
8
9
10
11
12
13
// 常用 Hint
spark.sql("SELECT /*+ BROADCAST(t2) */ * FROM t1 JOIN t2 ON t1.id = t2.id")
spark.sql("SELECT /*+ MERGE(t1, t2) */ * FROM t1 JOIN t2 ON t1.id = t2.id")
spark.sql("SELECT /*+ REPARTITION(100) */ * FROM t")
spark.sql("SELECT /*+ COALESCE(10) */ * FROM t")

// 多 Hint 组合
spark.sql("""
  SELECT /*+ BROADCAST(small_dim), MERGE(fact1, fact2) */
  * FROM fact1
    JOIN fact2 ON fact1.id = fact2.id
    JOIN small_dim ON fact1.dim_id = small_dim.id
""")

自定义 Catalyst 扩展:编写自己的优化规则

Catalyst 的扩展性是其最强大的设计之一。你可以通过

1
SparkSessionExtensions

注入自定义规则。


1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
// 自定义规则:将 Filter(1=1) 完全消除
import org.apache.spark.sql.SparkSession
import org.apache.spark.sql.catalyst.rules.Rule
import org.apache.spark.sql.catalyst.plans.logical.{Filter, LogicalPlan}
import org.apache.spark.sql.catalyst.expressions.Literal

object EliminateAlwaysTrueFilter extends Rule[LogicalPlan] {
  def apply(plan: LogicalPlan): LogicalPlan = plan.transformDown {
    case Filter(Literal(true, BooleanType), child) =&gt; child
    case Filter(Literal(false, BooleanType), _) =&gt;
      LocalRelation(plan.output)
  }
}

// 注册扩展
val spark = SparkSession.builder()
  .appName("CustomCatalyst")
  .withExtensions(_.injectOptimizerRule(EliminateAlwaysTrueFilter))
  .getOrCreate()

更实际的场景是针对业务特有的数据分布编写规则,例如针对特定列的字典编码压缩、针对时间分区表的早期过滤等。

总结与最佳实践

Catalyst 优化器是 Spark SQL 性能的基石,理解其工作流程是成为高效 Spark 开发者的必经之路。以下是要点回顾和实践建议:

  • 善用 explain():始终在优化性能前查看执行计划,定位瓶颈在哪个阶段
  • 写对优化器友好的代码:先 Filter 再 Transform,优先使用内置函数而非 UDF
  • 利用统计信息:对关键表执行
    1
    ANALYZE TABLE

    ,让 CBO 做出更好的 Join 决策

  • 开启 AQE:让 Catalyst 在运行时根据真实数据调整计划
  • 必要时使用 Hint:当你比优化器更了解数据时,显式 Hint 比调参数更直接
  • 监控代码生成:确认 Whole-Stage Code Generation 生效,星号标记是你的朋友

Catalyst 的设计哲学——规则驱动的树变换、可扩展的优化框架、运行时自适应——值得每一位数据工程师深入学习。掌握了 Catalyst,你就掌握了 Spark SQL 性能调优的钥匙。

【本站文章皆为原创,未经允许不得转载】:汤不热吧 » Spark Catalyst 优化器深度解析:从逻辑计划到物理计划的查询优化全流程
分享到: 更多 (0)