请说明在 PySpark 中提升 SQL 查询执行性能的常用手段,并列举出实践中常见的优化方法有哪些?
考察说明
考查对 PySpark SQL 执行引擎的理解及性能调优的实际经验。
回答思路
- 【回答框架 1】PySpark SQL 查询性能优化核心思路是减少数据扫描和洗牌(Shuffle)。常见方法包括:合理设置分区数(如 spark.sql.shuffle.partitions)、使用分区裁剪(Partition Pruning)和列式存储(如 Parquet)配合谓词下推(Predicate Pushdown),能大幅减少 IO。
- 【回答框架 2】优化连接(Join)操作:使用 Broadcast Join(通过 spark.sql.autoBroadcastJoinThreshold 控制)避免大表之间的 Shuffle;对倾斜数据(数据倾斜)可加盐或使用 Salting 技巧缓解;或将大表拆分为多个小表进行合并。
- 【回答框架 3】合理使用缓存(Cache)和持久化(Persist):对于反复使用的 DataFrame 或临时表,使用 cache 或 persist 到内存或磁盘,避免重复计算。注意根据数据规模选择存储级别,并注意清理缓存防止内存溢出。
- 【回答框架 4】调整执行计划:使用 Spark UI 观察物理计划,识别慢 Stage 和 Shuffle 量;通过 AQE(Adaptive Query Execution)自动优化,如动态合并 Shuffle 分区、动态调整 Join 策略;合理设计 Schema,避免读取多余列。
- 【关键点 1】分区裁剪与谓词下推是减少 IO 的基础手段。
- 【关键点 2】Broadcast Join 能显著减少 Shuffle,适合小表与大表连接。
- 【关键点 3】数据倾斜是性能杀手,加盐或两阶段聚合可有效缓解。
- 【关键点 4】Cache 适用于多次复用的中间结果,但需注意清理与存储级别。
- 【关键点 5】开启 AQE 可自动优化分区和 Join 策略,是重要的调优方向。
- 【易错点 1】过度依赖 cache 可能导致内存压力,需根据数据量选择合适存储级别。
- 【易错点 2】广播变量参与 Join 时,若表过大导致广播失败,反而增加性能开销。
- 【易错点 3】AQE 并非万能,对极不均匀的数据分布可能仍需要手动干预。