【发布时间】:2018-11-29 22:53:57
【问题描述】:
我在协调 sqlContext.sql("set spark.sql.shuffle.partitions=n") 和使用 df.repartition(n) 重新分区 Spark DataFrame 之间的差异(如果存在的话)有点困难。
Spark 文档指出set spark.sql.shuffle.partitions=n 配置了在洗牌数据时使用的分区数,而df.repartition 似乎返回了一个按指定键数分区的新 DataFrame。
为了让这个问题更清楚,这里是一个玩具示例,说明我相信df.reparition 和spark.sql.shuffle.partitions 是如何工作的:
假设我们有一个 DataFrame,如下所示:
ID | Val
--------
A | 1
A | 2
A | 5
A | 7
B | 9
B | 3
C | 2
-
场景一:3 Shuffle Partitions,Reparition DF by ID:
如果我要设置
sqlContext.sql("set spark.sql.shuffle.partitions=3"),然后设置df.repartition($"ID"),我希望我的数据被重新分区为 3 个分区,其中一个分区保存 ID 为“A”的所有行的 3 个 val,另一个保存所有行的 2 个 val ID 为“B”的行,最后一个分区保存所有 ID 为“C”的行中的 1 个值。 - 场景 2:5 个 shuffle partitions,Reparititon DF by ID:在这种场景下,我仍然希望每个分区只保存带有相同 ID 标记的数据。 也就是说,则不会在同一分区内混合具有不同 ID 的行。
我的理解是否在这里?一般来说,我的问题是:
我正在尝试优化数据帧的分区以避免 skew,但要让每个分区拥有尽可能多的相同键 尽可能的信息。我如何使用
set spark.sql.shuffle.partitions和df.repartiton实现这一目标?有链接吗 在
set spark.sql.shuffle.partitions和df.repartition之间?如果 那么,那个链接是什么?
谢谢!
【问题讨论】:
标签: apache-spark pyspark apache-spark-sql