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 的最佳实践

  1. 开发阶段定期检查

    // 在关键转换后检查执行计划
    val result = complexTransformation(df)
    result.explain()
    
  2. 对比优化效果

    // 优化前后对比
    originalDF.explain()
    optimizedDF.explain()
    
  3. 监控关键指标

    • 数据交换量(Exchange)
    • 扫描数据量(Scan)
    • 聚合操作复杂度(Aggregate)
  4. 结合Spark UI分析

    • 查看各stage执行时间
    • 识别数据倾斜
    • 监控内存使用

通过Explain语句深入理解Catalyst优化器的决策过程,可以显著提升Spark应用的性能和稳定性。

Logo

码道开发者社区,聚焦华为云码道 CodeArts 代码智能体,沉淀 Agent、Skill、鸿蒙开发实战内容,供开发者查阅资料、交流技术、分享工程实践

更多推荐