【问题标题】:Why is `getNumPartitions()` not giving me the correct number of partitions specified by `repartition`?为什么 `getNumPartitions()` 没有给我 `repartition` 指定的正确分区数?
【发布时间】:2016-03-21 23:42:21
【问题描述】:

我有一个textFile 和 RDD,如下所示:sc.textFile(<file_name>)

我尝试重新分区 RDD 以加快处理速度:

sc.repartition(<n>)

无论我为<n> 输入什么,它似乎都没有改变,如下所示:

RDD.getNumPartitions() 总是打印相同的数字(3) 无论如何。

如何更改分区数以提高性能?

【问题讨论】:

  • 你能检查你的配置并确保你有足够的执行器吗?

标签: apache-spark pyspark partition hadoop-partitioning


【解决方案1】:

这是因为 RDD 是不可变的。 您不能更改 RDD 的分区,但可以创建一个具有所需分区数的新分区。

scala> val a = sc.parallelize( 1 to 1000)
a: org.apache.spark.rdd.RDD[Int] = ParallelCollectionRDD[0] at  parallelize at <console>:21
scala> a.partitions.size
res2: Int = 4
scala> val b = a.repartition(6)
b: org.apache.spark.rdd.RDD[Int] = MapPartitionsRDD[4] at repartition at <console>:23
scala> a.partitions.size
res3: Int = 4
scala> b.partitions.size
res4: Int = 6

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2018-10-28
    • 1970-01-01
    • 2013-01-29
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多