【问题标题】:Spark: Order of column arguments in repartition vs partitionBySpark:重新分区与 partitionBy 中的列参数顺序
【发布时间】:2018-06-29 14:22:55
【问题描述】:
考虑的方法(Spark 2.2.1):
-
DataFrame.repartition(采用partitionExprs: Column*参数的两个实现)
DataFrameWriter.partitionBy
注意:本题不问这些方法的区别
来自docs 的partitionBy:
如果指定,则输出在文件系统上布局类似于Hive 的分区方案。例如,当我们按年和月对Dataset 进行分区时,目录布局如下所示:
- 年=2016/月=01/
- 年=2016/月=02/
由此,我推断列参数的顺序将决定目录布局;因此它是相关的。
来自docs 的repartition:
返回一个由给定分区表达式分区的新Dataset,使用spark.sql.shuffle.partitions作为分区数。生成的Dataset 被散列分区。
根据我目前的理解,repartition 决定了处理DataFrame 时的并行度。使用此定义,repartition(numPartitions: Int) 的行为很简单,但对于采用 partitionExprs: Column* 参数的 repartition 的其他两个实现,则不能这么说。
综上所述,我的疑惑如下:
- 与
partitionBy 方法一样,列顺序 输入是否也与repartition 方法相关?
- 如果上述问题的答案是
-
否:如果我们运行
SQL 查询,为并行执行提取的每个 chunk 是否包含与每个 group 中相同的数据GROUP BY 在同一列上?
-
是:请解释
repartition(columnExprs: Column*)方法的行为
- 在
repartition 的第三个实现中同时使用numPartitions: Int 和partitionExprs: Column* 参数有何相关性?
【问题讨论】:
标签:
apache-spark
dataframe
apache-spark-sql
partitioning
【解决方案1】:
在回答这个问题之前,让我先澄清一下spark中的一些概念。
块:这些物理映射到 HDFS 文件夹,并且能够存储子块和 parquet/* 文件。
parquet:数据存储压缩文件,常用于HDFS集群中存储数据。
现在来回答。
Repartition(number_of_partitions, *columns) :这将创建 parquet 文件,其中的数据会根据提供的列的不同组合值进行打乱和排序。因此列的顺序在这里没有任何区别。您可以在后台提供任何顺序 spark 将获取这些列的所有可能值,对它们进行排序并排列文件中的数据,这些文件将汇总到 number_of_partitions 。
PartionBy(*columns):这与重新分区略有不同。这将在 HDFS 中创建具有参数中提供的不同列值的块或文件夹。所以假设:
Col A = [1,2,3,4,5]
在写入 HDFS 表时,它会创建文件夹名称
colA-1
colA-2
colA-3
.
.
.
如果您提供两列,那么
colA-1/
colB-1
colB-2
colB-3
.
.
colA-2/
colA-3/
.
.
.
在其中它将存储镶木地板文件,这些文件将根据父列值对数据进行排序。此文件夹中的文件数将由 (bucketBy) 属性固定,该属性将进一步建议每个文件夹中的最大文件数。这仅在 pyspark 2.3 和 scala 1.6 及更高版本中可用。
【解决方案2】:
这两种方法的唯一相似之处是它们的名称。它们用于不同的事物并具有不同的机制,因此您根本不应该比较它们。
话虽如此,repartition 使用以下方式随机播放数据:
- 对于
partitionExprs,它在使用spark.sql.shuffle.partitions 的表达式中使用的列上使用哈希分区器。
-
partitionExprs 和 numPartitions 与前一个相同,但会覆盖 spark.sql.shuffle.partitions。
- 使用
numPartitions,它只是使用RoundRobinPartitioning 重新排列数据。
与重新分区方法相关的列输入的顺序也是?
是的。 hash((x, y)) 通常与 hash((y, x)) 不同。
df = (spark.range(5, numPartitions=4).toDF("x")
.selectExpr("cast(x as string)")
.crossJoin(spark.range(5, numPartitions=4).toDF("y")))
df.repartition(4, "y", "x").rdd.glom().map(len).collect()
[8, 6, 9, 2]
df.repartition(4, "x", "y").rdd.glom().map(len).collect()
[6, 4, 3, 12]
如果我们在相同的列上使用 GROUP BY 运行 SQL 查询,为并行执行提取的每个块是否包含与每个组中相同的数据?
取决于确切的问题。
相关How to define partitioning of DataFrame?