repartitionByRange
Renvoie un nouveau DataFrame partitionné par les expressions de partitionnement données. Le DataFrame résultant est partitionné par plage.
Syntaxe
repartitionByRange(numPartitions: Union[int, "ColumnOrName"], *cols: "ColumnOrName")
parameter
parameter | Type | Description |
|---|---|---|
| int | peut être un entier pour spécifier le nombre cible de partitions ou une colonne. S'il s'agit d'une colonne, elle sera utilisée comme première colonne de partitionnement. S'il n'est pas spécifié, le nombre de partitions default est utilisé. |
| str or Column | colonnes de partitionnement. |
Renvoie
DataFrame: DataFrame repartitionné.
Notes
Au moins une expression de partitionnement doit être spécifiée. Lorsqu'aucun ordre de tri explicite n'est spécifié, « ascendant, nulls en premier » est assumé.
Pour des raisons de performance, cette méthode utilise l'échantillonnage pour estimer les plages. Par conséquent, le résultat peut ne pas être cohérent, car l'échantillonnage peut renvoyer des valeurs différentes. La taille de l'échantillon peut être contrôlée par la configuration spark.sql.execution.rangeExchange.sampleSizePerPartition.
Exemples
from pyspark.sql import functions as sf
spark.createDataFrame(
[(14, "Tom"), (23, "Alice"), (16, "Bob")], ["age", "name"]
).repartitionByRange(2, "age").select(
"age", "name", sf.spark_partition_id()
).show()
# +---+-----+--------------------+
# |age| name|SPARK_PARTITION_ID()|
# +---+-----+--------------------+
# | 14| Tom| 0|
# | 16| Bob| 0|
# | 23|Alice| 1|
# +---+-----+--------------------+