欢迎光临

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)