如何缩短Databricks中大数据PySpark脚本的运行时长?
Databricks 6GiB数据处理优化方案
缓存清理的必要性
- 必须定期清理缓存!Notebook中反复测试转换操作时,Spark会自动缓存大量中间DataFrame,占用集群内存后会迫使后续操作读写磁盘,直接拖慢速度。
- 清理方式:
- 清空全局缓存:
spark.catalog.clearCache() - 释放单个DataFrame缓存:
df.unpersist()
- 清空全局缓存:
核心代码优化技巧
6GiB的数据完全没必要接受长运行时长,以下是无需扩容集群的优化方法:
1. 杜绝无差别全表扫描
- 替换小数据常用的
df[df[col1] == 5]为df.filter(col("col1") == 5),配合分区表使用:- 如果表按
col1(或关联字段)分区,Spark只会读取目标分区数据,过滤速度会提升数倍。 - 若未分区,优先对高频过滤字段做分区(适合基数适中的字段)或分桶(适合高基数字段)。
- 如果表按
2. 提前裁剪数据量级
- 测试阶段用采样/限制行数缩小数据集,避免加载全表:
# 随机采样1%数据(保证测试随机性) df_sample = df.sample(fraction=0.01, seed=42) # 或直接取前1万行(适合快速验证逻辑) df_small = df.limit(10000) - 加载数据时只选择需要的列,禁用
select("*"):df = spark.table("your_table").select("col1", "col2", "target_col")
3. 优化DataFrame操作逻辑
- 重复使用的中间结果先缓存,用完立即释放:
filtered_df = df.filter(col("col1") == 5).cache() # 执行多次基于filtered_df的转换操作 filtered_df.unpersist() # 用完及时释放内存 - 用
explain()排查低效执行计划:
重点看是否存在全表扫描、不必要的shuffle(比如笛卡尔积、未优化的join),针对性调整逻辑。df.filter(col("col1") ==5).explain(mode="extended")
4. 利用Delta Lake加速(若使用Delta表)
- 对高频过滤字段做Z-Order排序,开启数据跳过:
执行一次后,后续过滤OPTIMIZE your_delta_table ZORDER BY (col1)col1的操作会直接跳过无关数据块,大幅缩短扫描时间。
5. 替换Python UDF为内置函数
Python UDF会触发跨进程序列化,性能极差,优先用Spark内置函数:
- 比如用
when/otherwise替代自定义判断UDF,用concat替代自定义字符串拼接UDF。
其他实用调整
- Notebook分段运行,每次测试完成后清理中间变量和缓存;
- 调整
spark.sql.shuffle.partitions参数:默认200,6GiB数据建议设为50-100,减少shuffle开销。
内容的提问来源于stack exchange,提问作者Joakim Torsvik
相关产品推荐
相关产品推荐

