【发布时间】:2017-06-29 12:24:13
【问题描述】:
有什么方法可以获取 DataFrame 的当前分区数? 我检查了 DataFrame javadoc (spark 1.6) 并没有找到一种方法,或者我只是错过了它? (对于 JavaRDD,有一个 getNumPartitions() 方法。)
【问题讨论】:
标签: python scala dataframe apache-spark apache-spark-sql
有什么方法可以获取 DataFrame 的当前分区数? 我检查了 DataFrame javadoc (spark 1.6) 并没有找到一种方法,或者我只是错过了它? (对于 JavaRDD,有一个 getNumPartitions() 方法。)
【问题讨论】:
标签: python scala dataframe apache-spark apache-spark-sql
您需要在 DataFrame 的底层 RDD 上调用 getNumPartitions(),例如,df.rdd.getNumPartitions()。对于 Scala,这是一个无参数方法:df.rdd.getNumPartitions。
【讨论】:
DF 到RDD 的转换(昂贵)吗?
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
--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
记住所有这些建议,我建议你自己尝试
【讨论】:
since local[4] I used, I got 4 partitions. 对 3.x 仍然有效吗?我有 200 个带有本地 [4] 的分区。
local[n] 不同。实际上,由于 map/reduce 并行性,预计分区数与 local[n] 几乎没有关系。
转换为RDD然后得到分区长度
DF.rdd.partitions.length
【讨论】:
val df = Seq(
("A", 1), ("B", 2), ("A", 3), ("C", 1)
).toDF("k", "v")
df.rdd.getNumPartitions
【讨论】:
另一种获取分区数量的有趣方法是“使用 mapPartitions”转换。 示例代码 -
val x = (1 to 10).toList
val numberDF = x.toDF()
numberDF.rdd.mapPartitions(x => Iterator[Int](1)).sum()
欢迎 Spark 专家对其性能发表评论。
【讨论】: