【发布时间】: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