【问题标题】:Get current number of partitions of a DataFrame获取 DataFrame 的当前分区数
【发布时间】:2017-06-29 12:24:13
【问题描述】:

有什么方法可以获取 DataFrame 的当前分区数? 我检查了 DataFrame javadoc (spark 1.6) 并没有找到一种方法,或者我只是错过了它? (对于 JavaRDD,有一个 getNumPartitions() 方法。)

【问题讨论】:

    标签: python scala dataframe apache-spark apache-spark-sql


    【解决方案1】:

    您需要在 DataFrame 的底层 RDD 上调用 getNumPartitions(),例如,df.rdd.getNumPartitions()。对于 Scala,这是一个无参数方法:df.rdd.getNumPartitions

    【讨论】:

    • 减去 (),所以不完全正确 - 至少在 SCALA 模式下不正确
    • 这会导致从DFRDD转换昂贵)吗?
    • 这很贵
    • @javadba 你有一个不适合 RDD API 的答案吗?
    • 不,我不这样做:不幸的是,spark 没有像 hive 那样更好地管理元数据。您的回答是正确的,但我也观察到这是昂贵的。
    【解决方案2】:

    dataframe.rdd.partitions.size 是除df.rdd.getNumPartitions()df.rdd.length 之外的另一种选择。

    让我用完整的例子来解释一下......

    val x = (1 to 10).toList
    val numberDF = x.toDF(“number”)
    numberDF.rdd.partitions.size // => 4
    

    为了证明我们在上面得到了多少个分区...将该数据帧保存为 csv

    numberDF.write.csv(“/Users/Ram.Ghadiyaram/output/numbers”)
    

    这是数据在不同分区上的分离方式。

    Partition 00000: 1, 2
    Partition 00001: 3, 4, 5
    Partition 00002: 6, 7
    Partition 00003: 8, 9, 10
    

    更新:

    @Hemanth 在评论中问了一个很好的问题......基本上为什么数字 在上述情况下,分区数为 4

    简短回答:取决于您执行的情况。自从我使用 local[4] 以来,我得到了 4 个分区。

    长答案:

    我在本地机器上运行上面的程序,并使用 master 作为本地 [4],因为它被用作 4 分区。

    val spark = SparkSession.builder()
        .appName(this.getClass.getName)
        .config("spark.master", "local[4]").getOrCreate()
    

    如果它在 master yarn 中的 spark-shell 我得到的分区数为 2

    示例:spark-shell --master yarn 并再次键入相同的命令

    scala> val x = (1 to 10).toList
    x: List[Int] = List(1, 2, 3, 4, 5, 6, 7, 8, 9, 10)
    
    
    scala> val numberDF = x.toDF("number")
    numberDF: org.apache.spark.sql.DataFrame = [number: int]
    
    scala> numberDF.rdd.partitions.size
    res0: Int = 2
    
    • 这里2是default parllelism of spark
    • spark 将根据 hashpartitioner 决定分配多少个分区。如果您在--master local 中运行并基于您的Runtime.getRuntime.availableProcessors()local[Runtime.getRuntime.availableProcessors()]它将尝试分配 这些分区数。如果您的可用处理器数量为 12(即local[Runtime.getRuntime.availableProcessors()]),并且您有 1 到 10 个列表,则只会创建 10 个分区。

    注意:

    如果您使用的是 12 核笔记本电脑,我正在执行 spark 程序,并且默认情况下,分区/任务的数量是所有可用内核的数量,即 12。 表示local[*]s"local[${Runtime.getRuntime.availableProcessors()}]") 但在这 如果只有 10 个数字,所以会限制为 10

    记住所有这些建议,我建议你自己尝试

    【讨论】:

    • 感谢您的出色回答。我很好奇为什么在转换为 DF 时将 10 个数字的列表分为 4 个分区。请您解释一下好吗?
    • 这个since local[4] I used, I got 4 partitions. 对 3.x 仍然有效吗?我有 200 个带有本地 [4] 的分区。
    • @Sergey Bushmanov : see here 还有spark docs
    • 您提供的 2 个链接确实确认当前分区数与 local[n] 不同。实际上,由于 map/reduce 并行性,预计分区数与 local[n] 几乎没有关系。
    • 我们可以在map函数中获取分区号吗?比如 rdd.map{ r => this.partitionNum } ?
    【解决方案3】:

    转换为RDD然后得到分区长度

    DF.rdd.partitions.length
    

    【讨论】:

    • 我们可以在map函数中获取分区号吗?比如 rdd.map{ r => this.partitionNum } ?
    【解决方案4】:
     val df = Seq(
      ("A", 1), ("B", 2), ("A", 3), ("C", 1)
    ).toDF("k", "v")
    
    df.rdd.getNumPartitions
    

    【讨论】:

    • 请阅读此how-to-answer 以提供高质量的答案。
    • 我们可以在map函数中获取分区号吗?比如 rdd.map{ r => this.partitionNum } ?
    【解决方案5】:

    另一种获取分区数量的有趣方法是“使用 mapPartitions”转换。 示例代码 -

    val x = (1 to 10).toList
    val numberDF = x.toDF()
    numberDF.rdd.mapPartitions(x => Iterator[Int](1)).sum()
    

    欢迎 Spark 专家对其性能发表评论。

    【讨论】:

    • 我们可以在map函数中获取分区号吗?比如 rdd.map{ r => this.partitionNum } ?
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2010-12-09
    • 2020-04-28
    • 2017-03-13
    • 1970-01-01
    • 1970-01-01
    • 2017-01-15
    相关资源
    最近更新 更多