You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何优化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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.12 09:45:06