defaultParallelism : combien de cœurs Spark utilise-t-il ?

Spark déduit le nombre de partitions des cœurs qu’il croit avoir. Le code pour afficher ce chiffre, d’où il vient, et le piège des deux réglages.

scala
val rdd = sc.parallelize(1 to 100)
rdd.getNumPartitions

Cette ligne ne renvoie pas le même nombre sur votre portable, sur celui de votre collègue, et sur le cluster de production. Sur l’un ce sera 8, sur l’autre 20, en production 480. Aucun de ces chiffres n’apparaît dans votre code — Spark le déduit du nombre de cœurs qu’il croit avoir à disposition.

Ce chiffre invisible détermine le parallélisme de départ de tous vos jobs RDD. Le connaître, savoir d’où il vient et comment le forcer, c’est la base de tout réglage de performance Spark. Voici le code pour l’afficher, et l’explication de ce qu’il raconte.

Le code de diagnostic, à garder sous la main

C’est le snippet à coller en premier dans chaque nouveau notebook ou spark-shell :

scala
println(s"master                 : ${sc.master}")
println(s"defaultParallelism     : ${sc.defaultParallelism}")
println(s"cœurs JVM locaux       : ${Runtime.getRuntime.availableProcessors}")
println(s"exécuteurs actifs      : ${sc.statusTracker.getExecutorInfos.length}")
println(s"spark.sql.shuffle.part : ${spark.conf.get("spark.sql.shuffle.partitions")}")

Sortie typique sur un poste de développement à 10 cœurs physiques / 20 logiques :

master                 : local[*]
defaultParallelism     : 20
cœurs JVM locaux       : 20
exécuteurs actifs      : 1
spark.sql.shuffle.part : 200

Cinq lignes qui répondent d’un coup à « pourquoi mon job a N tâches », « pourquoi il ne va pas plus vite », et « pourquoi le chiffre change entre dev et prod ».

defaultParallelism : d’où vient le chiffre

sc.defaultParallelism est la valeur que Spark utilise quand vous ne précisez pas le nombre de partitions. Sa provenance dépend entièrement du mode d’exécution :

Mode (master)defaultParallelism vaut
local1
local[4]4
local[*]Tous les cœurs logiques vus par la JVM
YARN / Kubernetes / Standalonenb exécuteurs × cœurs par exécuteur, avec un minimum de 2

En mode local[*], Spark appelle littéralement Runtime.getRuntime.availableProcessors — c’est-à-dire les cœurs logiques, hyperthreading compris. Un processeur 10 cœurs / 20 threads donne 20, pas 10.

Vous pouvez forcer la valeur explicitement :

scala
spark.conf.set("spark.default.parallelism", "40")
// ou au lancement :
// spark-submit --conf spark.default.parallelism=40

À noter : spark.default.parallelism doit être posé avant la création du SparkContext pour être pris en compte partout. Dans un spark-shell déjà lancé, certains chemins de code l’ignorent.

parallelize(1 to 100) : la mise en pratique

scala
val rdd = sc.parallelize(1 to 100)
rdd.getNumPartitions
// res0: Int = 20    ← égal à defaultParallelism

Cent éléments, vingt partitions, cinq éléments chacune. Vérifions avec glom(), qui transforme chaque partition en tableau :

scala
rdd.glom().collect().zipWithIndex.foreach { case (elements, index) =>
  println(s"Partition $index : ${elements.mkString(", ")}")
}

Sortie :

Partition 0  : 1, 2, 3, 4, 5
Partition 1  : 6, 7, 8, 9, 10
Partition 2  : 11, 12, 13, 14, 15
...
Partition 19 : 96, 97, 98, 99, 100

Le calcul est trivial — 100 ÷ 20 = 5 — mais l’affichage confirme deux choses utiles : la répartition est contiguë (pas aléatoire) et équilibrée. C’est le comportement de parallelize sur une Range ; sur d’autres sources, la répartition peut être bien moins régulière.

Une version plus riche, qui donne aussi les tailles

Sur de vrais volumes, on veut le compte par partition, pas le contenu :

scala
def profilPartitions[T](rdd: org.apache.spark.rdd.RDD[T]): Unit = {
  val tailles = rdd.mapPartitions(it => Iterator(it.size)).collect()

  println(s"partitions : ${tailles.length}")
  println(s"total      : ${tailles.sum}")
  println(s"min / max  : ${tailles.min} / ${tailles.max}")
  println(s"médiane    : ${tailles.sorted.apply(tailles.length / 2)}")
}

profilPartitions(rdd)

mapPartitions(it => Iterator(it.size)) compte sans rapatrier les données — contrairement à glom().collect(), qui ramène tout sur le driver et fait exploser la mémoire sur un gros RDD. C’est la version utilisable en production, et le meilleur détecteur de skew : quand max vaut vingt fois la médiane, vous avez trouvé pourquoi une tâche traîne.

Forcer le nombre de partitions

Le second argument de parallelize court-circuite defaultParallelism :

scala
val rdd20 = sc.parallelize(1 to 100, 20)   // exactement 20, partout
rdd20.getNumPartitions                      // 20

C’est la bonne pratique dans tout code destiné à être partagé : un notebook qui repose sur defaultParallelism donne des résultats différents chez chaque personne, et des benchmarks incomparables. Fixer le nombre rend le comportement reproductible.

Le piège : defaultParallelismspark.sql.shuffle.partitions

Deux réglages, deux mondes, et une confusion très répandue :

PropriétéS’applique àValeur par défaut
spark.default.parallelismAPI RDDparallelize, textFile, shuffles RDDNombre de cœurs total
spark.sql.shuffle.partitionsAPI DataFrame / SQL — après groupBy, join, distinct200, en dur

Conséquence concrète : régler spark.default.parallelism à 400 ne change rien aux partitions produites par un groupBy sur un DataFrame. Celui-là écoute spark.sql.shuffle.partitions, dont le défaut historique de 200 ne correspond à presque aucun cluster moderne. C’est le sujet du tutoriel sur les partitions et le shuffle.

scala
// Un job DataFrame : c'est ce réglage qui compte
spark.conf.set("spark.sql.shuffle.partitions", "800")

// Un job RDD : c'est celui-là
spark.conf.set("spark.default.parallelism", "800")

Si vous travaillez en DataFrame — ce qui est presque toujours le bon choix — le second n’aura quasiment aucun effet sur vos temps d’exécution.

Sur un vrai cluster : le calcul complet

En mode YARN ou Kubernetes, defaultParallelism dérive de votre demande de ressources :

bash
spark-submit \
  --num-executors 30 \
  --executor-cores 4 \
  ...
# defaultParallelism = 30 × 4 = 120

Et la recommandation universelle en matière de dimensionnement :

Viser 2 à 3 fois le nombre total de cœurs en partitions.

Sur 120 cœurs, cela donne 240 à 360 partitions. Pourquoi plus que le nombre de cœurs ? Pour que l’ordonnanceur ait toujours du travail en réserve : quand une tâche finit tôt, le cœur libéré en reprend une autre immédiatement au lieu d’attendre la fin du stage. Avec exactement 120 partitions sur 120 cœurs, la moindre tâche lente laisse 119 cœurs inactifs.

La checklist de démarrage

Avant tout réglage de performance, ces quatre questions — auxquelles le snippet du début répond en une exécution :

  1. Quel master ? local[*] en dev, YARN/K8s en prod : le parallélisme n’a rien à voir.
  2. Combien de cœurs réellement disponibles ? defaultParallelism vous le dit, availableProcessors confirme en local.
  3. Combien de partitions dans mon RDD/DataFrame ? getNumPartitions, à chaque étape critique.
  4. Sont-elles équilibrées ? profilPartitions ci-dessus, et comparez max à la médiane.

Ce qu’il faut retenir

Spark ne devine pas votre intention : il compte les cœurs qu’il voit et découpe en conséquence. sc.defaultParallelism est ce chiffre, getNumPartitions en est le reflet dans un RDD donné, et glom() permet de vérifier de ses yeux. Le réflexe qui fait gagner le plus de temps : fixer explicitement le nombre de partitions dans tout code partagé, et ne jamais confondre spark.default.parallelism (RDD) avec spark.sql.shuffle.partitions (DataFrame) — ce sont deux leviers distincts pour deux APIs distinctes.

Pour changer le nombre de partitions après coup, l’article dédié à coalesce vs repartition, démontré avec glom() prend le relais exactement là où celui-ci s’arrête.

Ce sujet fait partie d’un cours complet

Développement et déploiement de solutions de données — les premiers modules sont en accès libre.

Voir le plan du cours

Continuer sur le même sujet