Spark SQL 是如何优化查询计划的?Explain 语句的作用是什么?
·
Spark SQL 查询优化机制
Catalyst 优化器架构
Spark SQL 使用基于规则的 Catalyst 优化器,采用函数式编程范式:
原始逻辑计划
↓
分析阶段(Analysis)
↓
逻辑优化(Logical Optimization)
↓
物理计划(Physical Planning)
↓
代码生成(Code Generation)
核心优化技术
1. 谓词下推(Predicate Pushdown)
-- 优化前:全表扫描后过滤
SELECT * FROM sales WHERE date > '2023-01-01'
-- 优化后:过滤条件下推到数据源
-- Parquet/ORC文件直接跳过不满足条件的数据块
2. 列裁剪(Column Pruning)
-- 只读取需要的列,减少I/O
SELECT name, age FROM users -- 只读取name和age列
3. 常量折叠(Constant Folding)
-- 编译时计算常量表达式
SELECT salary * 1.1 * 0.9 FROM employees
-- 优化为:SELECT salary * 0.99 FROM employees
4. 公共子表达式消除
-- 重复计算被优化
SELECT (price * quantity) as revenue,
(price * quantity) * 0.1 as tax
-- 优化为计算一次 revenue,然后派生 tax
5. 连接重排序(Join Reordering)
-- 基于统计信息重新排列join顺序
SELECT * FROM large_table l
JOIN small_table s ON l.id = s.id
-- 优化器可能选择 broadcast join 小表
6. 自动广播连接(Broadcast Join)
-- 小表自动广播到所有executor
-- 当小表 < spark.sql.autoBroadcastJoinThreshold (默认10MB)
Explain 语句的作用和用法
Explain 的三种模式
1. 简单模式(默认)
df.explain()
显示:物理执行计划
2. 扩展模式
df.explain("extended")
显示:逻辑计划 → 优化后逻辑计划 → 物理计划
3. 代码生成模式
df.explain("codegen")
显示:生成的Java代码
Explain 输出解析示例
简单模式输出:
== Physical Plan ==
*(1) Project [name#10, age#11L]
+- *(1) Filter (isnotnull(age#11L) && (age#11L > 18))
+- *(1) Scan json [name#10, age#11L]
关键信息解读:
*(1):阶段编号和并行度Project:列选择操作Filter:过滤条件Scan:数据源扫描HashAggregate:哈希聚合Exchange:数据交换(shuffle)
实际应用场景
1. 性能诊断
// 检查是否使用了广播连接
df1.join(df2, "id").explain()
// 期望看到:BroadcastHashJoin
// 如果看到:SortMergeJoin,可能需要调整配置
2. 优化验证
// 验证谓词下推是否生效
spark.conf.set("spark.sql.parquet.filterPushdown", "true")
df.filter($"date" > "2023-01-01").explain()
// 检查Filter是否在Scan之前
3. 分区策略分析
// 检查数据分布
df.repartition(100, $"dept").explain()
// 查看Exchange操作的类型和分区数
高级优化技巧
1. 统计信息收集
// 收集表统计信息(优化连接顺序)
spark.sql("ANALYZE TABLE sales COMPUTE STATISTICS")
spark.sql("ANALYZE TABLE sales COMPUTE STATISTICS FOR COLUMNS price, quantity")
2. 提示(Hints)使用
-- 强制广播连接
SELECT /*+ BROADCAST(small_table) */ *
FROM large_table JOIN small_table ON id
-- 强制合并小文件
SELECT /*+ COALESCE(10) */ * FROM table
3. 自适应查询执行(AQE)
Spark 3.0+ 特性:
- 动态合并小分区
- 动态切换连接策略
- 动态优化倾斜连接
// 启用AQE
spark.conf.set("spark.sql.adaptive.enabled", "true")
spark.conf.set("spark.sql.adaptive.coalescePartitions.enabled", "true")
性能调优实战
识别问题模式
问题1:数据倾斜
== Physical Plan ==
Exchange hashpartitioning(key#0, 200)
+- *(2) Project [key#0, value#1]
→ 某个task执行时间远长于其他task
解决方案:
// 1. 增加shuffle分区数
spark.conf.set("spark.sql.shuffle.partitions", "400")
// 2. 使用salting技术
df.withColumn("salted_key", concat($"key", lit("_"), (rand() * 10).cast("int")))
问题2:Cartesian Product
== Physical Plan ==
CartesianProduct
:- Scan table1
+- Scan table2
→ 缺少连接条件,产生笛卡尔积
问题3:低效扫描
== Physical Plan ==
Scan parquet [all_columns...]
→ 未使用列裁剪,读取了不需要的列
Explain 的最佳实践
-
开发阶段定期检查
// 在关键转换后检查执行计划 val result = complexTransformation(df) result.explain() -
对比优化效果
// 优化前后对比 originalDF.explain() optimizedDF.explain() -
监控关键指标
- 数据交换量(Exchange)
- 扫描数据量(Scan)
- 聚合操作复杂度(Aggregate)
-
结合Spark UI分析
- 查看各stage执行时间
- 识别数据倾斜
- 监控内存使用
通过Explain语句深入理解Catalyst优化器的决策过程,可以显著提升Spark应用的性能和稳定性。
更多推荐


所有评论(0)