【问题标题】:Spark: Order of column arguments in repartition vs partitionBySpark:重新分区与 partitionBy 中的列参数顺序
【发布时间】:2018-06-29 14:22:55
【问题描述】:

考虑的方法(Spark 2.2.1):

  1. DataFrame.repartition(采用partitionExprs: Column*参数的两个实现)
  2. 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?

      【讨论】:

        猜你喜欢
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 2016-02-14
        • 2020-03-19
        • 2015-10-15
        • 1970-01-01
        • 1970-01-01
        相关资源
        最近更新 更多