Catalyst 优化器深入
概述
Catalyst 是 Spark SQL 的查询优化框架,采用 树 + 规则(Rule) 的模式,把一条 SQL 从字符串一步步变成可高效执行的物理计划。本文深入四个阶段:Parser、Analyzer、Optimizer、SparkPlanner,讲清逻辑计划到物理计划的完整过程。
一、核心设计:树与规则
1.1 树(TreeNode)
所有计划(Plan)、表达式(Expression)都是不可变的树结构:
Project(user_id, amount) ← 节点
└─ Filter(dt = '2026-08-01') ← 节点
└─ Relation(orders) ← 叶子| 类型 | 说明 |
|---|---|
| LogicalPlan | 逻辑计划树 |
| Expression | 表达式树(条件、运算) |
| PhysicalPlan | 物理计划树 |
1.2 规则(Rule)
规则是 输入树 → 输出新树 的变换函数,Catalyst 通过规则批量优化:
Rule = 模式匹配 + 替换
例:谓词下推规则
Filter(cond, Join(a, b)) → Join(Filter(cond_a, a), Filter(cond_b, b))规则反复应用,直到树不再变化(固定点)。
二、阶段一:Parser
2.1 作用
把 SQL 字符串解析为 未解析的逻辑计划(Unresolved Logical Plan)。
输入:SELECT user_id, SUM(amount) FROM orders WHERE dt = '2026-08-01'
输出(语法树):
Project(user_id, SUM(amount))
└─ Filter(dt = '2026-08-01')
└─ UnresolvedRelation(orders)2.2 特点
| 特点 | 说明 |
|---|---|
| 只做语法检查 | 不校验表/列是否存在 |
| 输出 Unresolved 树 | 待 Analyzer 绑定 |
| 支持方言 | HiveQL、标准 SQL、ANSI 模式 |
三、阶段二:Analyzer
3.1 作用
用 Catalog 把未解析的标识符绑定到真实 Schema,进行语义分析。
输入:UnresolvedRelation(orders) + Filter(dt = ...)
Catalog 查表 orders 的 Schema
输出:
Relation(orders, schema=[user_id: Long, amount: Double, dt: String])
Filter(dt = '2026-08-01') ← 校验 dt 列存在且类型匹配3.2 绑定内容
| 绑定 | 说明 |
|---|---|
| 表名 | 关联到 Catalog 中的表 |
| 列名 | 绑定真实列与类型 |
| 函数 | 解析内置/UDF |
| 类型检查 | 表达式类型匹配校验 |
| 隐式转换 | 类型自动提升(如 Int→Long) |
四、阶段三:Optimizer(逻辑优化)
4.1 作用
对逻辑计划应用规则优化,输出 Optimized Logical Plan。
4.2 核心规则
| 规则 | 示例 |
|---|---|
| 谓词下推 | Filter 下推到数据源/Join 前 |
| 列裁剪 | 只投影需要的列 |
| 常量折叠 | 1 + 2 → 3 |
| 常量过滤 | WHERE 1 = 0 → 空结果 |
| 布尔简化 | a AND true → a |
| 消除重复 | 去重冗余 Project |
| 子查询消除 | 扁平化相关子查询 |
4.3 示例
sql
SELECT user_id, amount FROM (
SELECT * FROM orders WHERE amount > 100
) t WHERE dt = '2026-08-01'优化后:
原始:Project(user_id, amount, dt)
└─ Filter(dt = '2026-08-01')
└─ Filter(amount > 100)
└─ Relation(orders)
优化:Project(user_id, amount) ← 列裁剪
└─ Filter(amount > 100 AND dt = '2026-08-01')
└─ Relation(orders) ← 谓词下推合并五、阶段四:SparkPlanner(物理计划)
5.1 作用
把逻辑计划转换为 Physical Plan,选择具体执行策略。
5.2 物理策略选择
| 逻辑算子 | 物理实现选择 |
|---|---|
| Join | BroadcastHashJoin / SortMergeJoin / ShuffledHashJoin |
| Scan | FileSourceScan / JDBCScan / InMemoryScan |
| Aggregate | HashAggregate(代码生成) |
| Sort | 内存排序 / 外部排序 |
5.3 Join 策略决策
| 条件 | 选择 |
|---|---|
| 小表可广播(默认 10MB) | BroadcastHashJoin(无 Shuffle) |
| 大表互 join | SortMergeJoin(Shuffle + 排序归并) |
| 有桶表 | BucketedSortMergeJoin(免 Shuffle) |
5.4 执行准备
物理计划再经 WholeStageCodegen 生成 Java 代码,形成 RDD 执行。
六、查看执行计划
6.1 explain 方法
scala
df.explain() // 物理计划
df.explain(true) // 逻辑 + 优化 + 物理全计划
df.explain("extended") // 同上
df.explain("codegen") // 生成的代码6.2 计划解读示例
text
== Optimized Logical Plan ==
Aggregate [user_id], [sum(amount) AS total]
+- Filter (dt = 2026-08-01)
+- Relation[user_id, amount, dt]
== Physical Plan ==
*(2) HashAggregate(keys=[user_id], functions=[sum(amount)])
+- Exchange hashpartitioning(user_id, 200)
+- *(1) FileScan parquet ... ← 谓词已下推* 号表示启用了 WholeStageCodegen。
七、自定义规则与扩展
7.1 自定义优化规则
scala
class MyRule(spark: SparkSession)
extends Rule[LogicalPlan] {
def apply(plan: LogicalPlan): LogicalPlan = plan.transform {
// 模式匹配并替换
}
}
spark.experimental.extraOptimizations += new MyRule(spark)7.2 自定义函数(UDF)
| 方式 | 说明 |
|---|---|
spark.udf.register("func", fn) | SQL 中调用 |
UDFRegistration | DataFrame 注册 |
| 自定义类型 UDT | 复杂类型支持 |
常见问题速查
| 问题 | 原因与处理 |
|---|---|
| 谓词下推未生效 | 自定义 UDF 阻断,改 Catalyst 可识别写法 |
| Broadcast 未触发 | 表超过阈值,用 hint 或调大 spark.sql.autoBroadcastJoinThreshold |
| 执行计划看不懂 | 从 Optimized 开始读,关注 Exchange(Shuffle)位置 |
| SQL 慢但看不出问题 | explain(true) 对比逻辑与物理,找意外 Shuffle |
| 规则不生效 | 检查表是否 cached,缓存表绕过部分优化 |