如何优化Databricks中joined_df.display()的执行性能?
问题分析
生成joined_df仅耗时1.8秒是因为Spark只是构建了执行计划,实际计算要等到display()这类action操作触发才会执行。你遇到的多作业、长耗时问题,核心是执行计划里存在冗余shuffle、不必要排序,以及分区策略不合理导致的重复计算。
优化方案
1. 清理冗余的分区与广播操作
all_combinations_df是crossJoin后的结果,数据量通常较大,绝对不能对大表用broadcast——这会把整张表拷贝到所有节点,徒增网络开销。直接去掉F.broadcast():all_combinations_df_partitioned = all_combinations_df.repartition(4, "goods_no")ranked_goods_df如果数据量不小,repartition(8, "goods_no")后再broadcast完全冗余,直接去掉hint("broadcast")和不必要的重分区(如果Redshift读取时已经做了合理分区):ranked_goods_partitioned = ranked_goods_df # 保留原分区,或根据实际数据量调整- 仅保留小表的
broadcast:gl_goods_df和goods_df这类维度表用broadcast是合理的,继续保留。
2. 移除无意义的中间排序
filtered_date_df里的orderBy("char_date")完全多余,后续crossJoin不需要日期有序,去掉这个排序能减少一次shuffle作业:
filtered_date_df = date_df.select( F.date_format(F.col("date"), "yyyy-MM").alias("char_date") ).dropDuplicates() # 删除orderBy("char_date")
3. 对齐窗口函数的分区策略
窗口函数Window.partitionBy("goods_no")需要上游DataFrame的分区键和它一致,避免额外shuffle。调整all_combinations_df_partitioned的分区数,和Spark默认并行度或数据量匹配:
# 比如设置为8,和后续窗口分区、Spark并行度对齐 all_combinations_df_partitioned = all_combinations_df.repartition(8, "goods_no")
4. 提前缓存中间结果
如果需要多次查看joined_df,先缓存避免重复计算:
joined_df.cache() joined_df.count() # 触发缓存执行 joined_df.display()
若数据量极大,改用persist指定存储级别(如磁盘+内存):
from pyspark.storagelevel import StorageLevel joined_df.persist(StorageLevel.MEMORY_AND_DISK)
5. 限制display的数据量
display()默认加载全量数据,若只需验证结果,取前N行即可大幅缩短耗时:
joined_df.limit(1000).display()
6. 优化Redshift读取并行度
读取Redshift时指定合理的并行度,减少读取阶段的作业数:
# 示例:读取时设置numPartitions参数 all_goods_no_df = read_from_redshift(all_goods_no_query, numPartitions=8)
内容的提问来源于stack exchange,提问作者DOHEE KIM
相关产品推荐
相关产品推荐

