【问题标题】:Minimize shuffle spill and shuffle write最小化随机溢出和随机写入
【发布时间】:2021-06-11 22:56:49
【问题描述】:

我在 HDFS 中有以下管道,我正在 spark 中处理这些管道

输入表: batch, team, user, metric1, metric2

此表可以按小时批量包含用户级别的指标。在同一小时内,一个用户可以有多个条目。

1 级聚合:此聚合用于获取每个用户每批次的最新条目

agg(metric1) as user_metric1, agg(metric2) as user_metric2 (group by batch, team, user)

2 级聚合: 获取团队级别指标

agg(user_metric1) as team_metric1, agg(user_metric2) as team_metric2 (group by batch, team)

HDFS 中的输入表大小为 8gb(快速拼花格式)。我的 spark 工作显示 shuffle 写入 40gb 并且每个执行器 shuffle 溢出至少 1 gb。

为了尽量减少这种情况,如果我在执行聚合之前在用户级别重新分区输入表,

df = df.repartition('user')

它会提高性能吗?如果我想减少洗牌,我应该如何解决这个问题?

使用以下资源运行时

spark.executor.cores=6
spark.cores.max=48
spark.sql.shuffle.partitions=200

【问题讨论】:

    标签: apache-spark pyspark shuffle


    【解决方案1】:

    Spark 将数据从一个节点转移到另一个节点,因为资源分布在集群上(输入数据...),这可能会使计算速度变慢,并且可能会在集群上呈现大量网络流量,对于您的情况,数字shuffle 的数量是由于 group by ,如果你根据 goup by 的三列进行重新分区,它将减少 shuffle 的数量,对于 spark 配置,默认 spark.sql.shuffle.partitions 是 200,假设我们将让 spark 配置保持原样,重新分区需要一些时间,一旦完成计算会更快:

    new_df = df.repartition("batch","team", "user")
    
    

    【讨论】:

      猜你喜欢
      • 2015-11-19
      • 2017-05-30
      • 1970-01-01
      • 1970-01-01
      • 2016-06-07
      • 1970-01-01
      • 1970-01-01
      • 2018-01-14
      • 1970-01-01
      相关资源
      最近更新 更多