¿Hay alguna forma de obtener el número actual de particiones de un DataFrame? Revisé el javadoc de DataFrame (spark 1.6) y no encontré un método para eso, ¿o simplemente me lo perdí? (En el caso de JavaRDD, hay un método getNumPartitions()).
Debe llamar a getNumPartitions() en el RDD subyacente de DataFrame, por ejemplo, df.rdd.getNumPartitions() . En el caso de Scala, este es un método sin parámetros: df.rdd.getNumPartitions .
dataframe.rdd.partitions.size es otra alternativa además de df.rdd.getNumPartitions() o df.rdd.length .
déjame explicarte esto con un ejemplo completo...
val x = (1 to 10).toList val numberDF = x.toDF(“number”) numberDF.rdd.partitions.size // => 4Para probar la cantidad de particiones que obtuvimos arriba ... guarde ese marco de datos como csv
numberDF.write.csv(“/Users/Ram.Ghadiyaram/output/numbers”)Así es como se separan los datos en las diferentes particiones.
Partition 00000: 1, 2 Partition 00001: 3, 4, 5 Partition 00002: 6, 7 Partition 00003: 8, 9, 10@Hemanth hizo una buena pregunta en el comentario... básicamente por qué el número de particiones es 4 en el caso anterior
Respuesta corta: depende de los casos en los que esté ejecutando. desde que usé local[4], obtuve 4 particiones.
Respuesta larga :
Estaba ejecutando el programa anterior en mi máquina local y usé el maestro como local [4] en función de que estaba tomando 4 particiones.
val spark = SparkSession.builder() .appName(this.getClass.getName) .config("spark.master", "local[4]").getOrCreate()Si es una cáscara de chispa en hilo maestro, obtuve el número de particiones como 2
ejemplo: spark-shell --master yarn y volvió a escribir los mismos comandos
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 y en función de su Runtime.getRuntime.availableProcessors() es decir local[Runtime.getRuntime.availableProcessors()] , intentará asignar esa cantidad de particiones. si la cantidad de procesadores disponibles es 12 (es decir local[Runtime.getRuntime.availableProcessors()]) y tiene una lista de 1 a 10, solo se crearán 10 particiones.NOTA:
Si está en una computadora portátil de 12 núcleos donde estoy ejecutando el programa Spark y, de forma predeterminada, la cantidad de particiones/tareas es la cantidad de todos los núcleos disponibles, es decir, 12. Eso significa
local[*]os"local[${Runtime.getRuntime.availableProcessors()}]")pero en este caso solo hay 10 números, por lo que se limitará a 10
teniendo en cuenta todos estos consejos, le sugiero que pruebe por su cuenta
convertir a RDD y luego obtener la longitud de las particiones
DF.rdd.partitions.length