Une partition est une tâche potentielle, pas un cœur. L’exécution par vagues, et comment inspecter les partitions avec mapPartitionsWithIndex.
val rdd = sc.parallelize(1 to 100, 100)
rdd.sum()Cent partitions, cent tâches. Sur une machine à vingt cœurs, combien de tâches tournent en même temps ? Vingt. Les quatre-vingts autres attendent leur tour. C’est le point qui manque à la plupart des explications sur les partitions Spark, et c’est celui qui rend la suite intuitive : une partition n’est pas un processeur, c’est une tâche potentielle.
Cet article traite du modèle d’exécution — combien de tâches, dans quel ordre, sur quelles ressources — et donne les deux outils pour voir quelle partition traite quel élément. C’est le complément direct du tutoriel sur le dimensionnement des partitions : ici, on regarde ce qui se passe à l’exécution.
La chaîne est simple, et il faut la garder en tête intégralement :
Données → Partitions → Tâches → Slots d'exécution → Cœursnb exécuteurs × cœurs par exécuteur.Le nombre de partitions détermine donc la quantité de travail découpé, jamais la capacité de calcul. Confondre les deux mène directement à l’erreur classique : ajouter des partitions en espérant aller plus vite, alors que les cœurs sont déjà tous occupés.
Cent tâches, vingt slots. Spark n’attend pas d’avoir cent cœurs — il exécute en vagues :
Le calcul est direct :
nombre de vagues ≈ nombre de partitions ÷ nombre de slots
100 ÷ 20 = 5 vaguesCe n’est pas un ordonnancement strict par paquets de vingt — dès qu’une tâche finit, la suivante démarre immédiatement sur le slot libéré. Mais l’ordre de grandeur est celui-là, et il explique deux comportements courants :
Pourquoi viser 2 à 3 fois le nombre de cœurs. Avec exactement vingt partitions sur vingt slots, une seule tâche lente laisse dix-neuf cœurs inactifs jusqu’à la fin du stage. Avec soixante partitions, il reste toujours du travail en réserve pour occuper les cœurs qui se libèrent — l’ordonnanceur lisse naturellement les écarts de durée.
Pourquoi trop de partitions coûte cher. Chaque tâche a un coût fixe : sérialisation du code, planification, comptabilité, lancement. Compter en millisecondes par tâche paraît négligeable — jusqu’à ce que vous ayez cinquante mille tâches qui traitent chacune trois lignes.
mapPartitionsWithIndexC’est l’outil de référence pour inspecter la répartition physique. Il donne accès au numéro de partition et à son contenu :
val rdd = sc.parallelize(1 to 20, 4)
rdd.mapPartitionsWithIndex { case (index, elements) =>
Iterator(s"Partition $index : ${elements.mkString(", ")}")
}.collect().foreach(println)Sortie :
Partition 0 : 1, 2, 3, 4, 5
Partition 1 : 6, 7, 8, 9, 10
Partition 2 : 11, 12, 13, 14, 15
Partition 3 : 16, 17, 18, 19, 20Deux avantages décisifs sur glom().collect() :
zipWithIndex côté driver.// Compter sans rapatrier les données — utilisable en production
rdd.mapPartitionsWithIndex { case (index, elements) =>
Iterator((index, elements.size))
}.collect().foreach { case (index, taille) =>
println(f"Partition $index%3d : $taille%,d éléments")
}glom().collect() ramène toutes les données sur le driver et le fait exploser dès que le RDD devient sérieux. mapPartitionsWithIndex qui renvoie une taille ne transporte qu’un entier par partition.
def diagnostiquer[T](rdd: org.apache.spark.rdd.RDD[T], nom: String = "rdd"): Unit = {
val tailles = rdd
.mapPartitionsWithIndex { case (i, it) => Iterator((i, it.size)) }
.collect()
val compte = tailles.map(_._2)
val mediane = compte.sorted.apply(compte.length / 2)
val (pireIndex, pireTaille) = tailles.maxBy(_._2)
println(s"[$nom] ${tailles.length} partitions, ${compte.sum} éléments")
println(s"[$nom] médiane $mediane, max $pireTaille (partition $pireIndex)")
if (mediane > 0 && pireTaille > mediane * 5)
println(s"[$nom] ⚠ skew : la pire partition fait ${pireTaille / mediane}× la médiane")
}Un ratio max/médiane supérieur à 5 explique presque toujours la tâche qui traîne à 99 % pendant que les autres sont finies. C’est le premier diagnostic à lancer sur un stage anormalement lent.
TaskContext.getPartitionId() : le numéro depuis l’intérieurPour connaître la partition au niveau de chaque élément, sans changer la structure du RDD :
import org.apache.spark.TaskContext
val rdd = sc.parallelize(1 to 20, 4)
rdd.map { nombre =>
(TaskContext.getPartitionId(), nombre)
}.collect().foreach { case (partition, nombre) =>
println(s"Partition $partition traite $nombre")
}Sortie :
Partition 0 traite 1
Partition 0 traite 2
...
Partition 1 traite 6TaskContext est un objet accessible depuis le code exécuté sur les exécuteurs. Il donne aussi attemptNumber() (utile pour tracer les tâches réessayées après échec) et stageId().
Son usage principal en production n’est pas l’affichage : c’est le logging contextuel. Quand une exception remonte d’une tâche, savoir quelle partition a échoué fait gagner des heures.
rdd.map { valeur =>
try traiter(valeur)
catch {
case e: Exception =>
val ctx = TaskContext.get()
throw new RuntimeException(
s"échec partition ${ctx.partitionId()}, tentative ${ctx.attemptNumber()}", e
)
}
}Chaque partition produit un fichier lors d’une écriture distribuée. C’est mécanique, et c’est la source du problème de petits fichiers.
sc.parallelize(1 to 100, 10).saveAsTextFile("sortie-10")Produit :
sortie-10/
├── part-00000
├── part-00001
...
├── part-00009
└── _SUCCESSDix partitions, dix fichiers part-*. Le fichier _SUCCESS, vide, signale que l’écriture s’est terminée sans erreur — plusieurs orchestrateurs l’attendent avant de démarrer l’étape suivante.
Poussons le raisonnement à l’absurde :
sc.parallelize(1 to 100, 100).saveAsTextFile("sortie-100")
// 100 fichiers, contenant chacun UN nombreCent fichiers pour cent entiers. Sur un stockage objet facturé à la requête, avec des métadonnées à gérer par fichier, c’est exactement ce qu’il ne faut pas produire. La correction tient en un appel :
sc.parallelize(1 to 100, 100).coalesce(2).saveAsTextFile("sortie-2")
// 2 fichiersLe mécanisme complet — et pourquoi coalesce plutôt que repartition ici — est détaillé dans l’article coalesce vs repartition démontré avec glom(). Et les conséquences côté lecture sont traitées dans écrire proprement en Parquet et Delta Lake.
| Situation | Code | Ce qui se passe |
|---|---|---|
| Trop de partitions | parallelize(1 to 100, 1000) | 1 000 tâches pour 100 éléments, 900 partitions vides, coût de planification > coût du calcul |
| Trop peu de partitions | parallelize(1 to 100000000, 1) | 1 seule tâche, 1 cœur occupé sur 20, partition énorme, risque d’OutOfMemoryError |
| Équilibré | parallelize(1 to 100000000, 60) | 60 tâches sur 20 slots, 3 vagues, ordonnanceur toujours alimenté |
La règle qui découle du modèle par vagues : 2 à 3 fois le nombre de slots d’exécution. Pour vingt cœurs, quarante à soixante partitions. Pour un cluster à 120 cœurs, 240 à 360.
1. Créer et observer. Créez un RDD des nombres 1 à 30 avec trois partitions, puis affichez le contenu de chacune avec mapPartitionsWithIndex. Combien d’éléments par partition ?
2. Une partition par élément. Créez un RDD de 1 à 20 avec vingt partitions. Vérifiez avec getNumPartitions, puis affichez le contenu. Que se passerait-il si vous demandiez trente partitions pour vingt éléments ?
3. Réduction. Créez un RDD de cinq partitions, réduisez-le à deux avec coalesce, et comparez la répartition avant/après. Quelles partitions ont fusionné ?
4. Augmentation. Créez un RDD de deux partitions et portez-le à huit. coalesce ou repartition ? Vérifiez votre réponse avec getNumPartitions — l’une des deux méthodes va vous ignorer silencieusement.
5. Fichiers de sortie. Sauvegardez un RDD de dix partitions, comptez les fichiers part-*. Réduisez à deux partitions, sauvegardez ailleurs, comparez. Combien de fichiers dans chaque dossier, et pourquoi ?
Un rappel pour l’exercice 5 : relancer saveAsTextFile sur un chemin existant lève FileAlreadyExistsException. Spark refuse d’écraser pour ne pas détruire un résultat par accident — changez de nom ou supprimez le dossier. Le détail est dans l’article sur le word count RDD.
Une partition est un morceau de données et une tâche potentielle — jamais un processeur. Le nombre de partitions décide de la granularité du découpage ; le nombre de slots décide de la vitesse à laquelle ce découpage est consommé, par vagues successives. De cette seule distinction découlent la règle des 2-3× cœurs, le coût réel d’un nombre excessif de tâches, et la relation directe entre partitions et fichiers de sortie.
Deux outils à retenir pour ne plus travailler à l’aveugle : mapPartitionsWithIndex pour voir et mesurer la répartition sans faire tomber le driver, et TaskContext pour savoir quelle partition a échoué quand une exception remonte. Ensemble, ils transforment le débogage de performance Spark d’une intuition en une mesure.
Développement et déploiement de solutions de données — les premiers modules sont en accès libre.