Quatre opérations qui se ressemblent et un choix qui change la facture cloud. Pourquoi reduceByKey bat groupByKey, et les pièges du tri de clés.
Quatre méthodes du RDD, quatre noms qui se ressemblent, un seul auteur d’articles Spark sur trois qui explique correctement la différence : groupBy, groupByKey, reduceByKey, sortByKey. Elles répondent à des questions différentes, elles n’ont pas le même coût, et l’une d’elles — la plus intuitive à utiliser — est probablement la ligne la plus chère de votre pipeline si vous l’avez laissée passer une revue de code.
On prend une petite liste de prénoms et une liste de paires clé-valeur, on exécute chaque commande, on regarde ce qu’elle produit et ce qu’elle coûte, et on isole la seule différence qui compte vraiment : celle qui vit sur votre facture cloud à la fin du mois.
val rddGby1 = sc.parallelize(
List("Salim", "Martha-Patricia", "Abed", "François", "Sonia", "Abntest")
)Six chaînes de caractères, pas encore de couples clé-valeur.
groupBy — regrouper avec une fonction de son choixrddGby1.groupBy(x => x.charAt(0)).collect()groupBy prend une fonction qui calcule la clé à la volée. Ici, la clé est la première lettre du nom.
| Nom | charAt(0) |
|---|---|
| Salim | S |
| Martha-Patricia | M |
| Abed | A |
| François | F |
| Sonia | S |
| Abntest | A |
Résultat :
(S, CompactBuffer(Salim, Sonia))
(M, CompactBuffer(Martha-Patricia))
(A, CompactBuffer(Abed, Abntest))
(F, CompactBuffer(François))Le type approximatif est RDD[(Char, Iterable[String])]. L’ordre de sortie n’est pas garanti — Spark travaille en distribué et ne réordonne pas gratuitement.
Une erreur fréquente :
rddGby1.groupBy(x => x.charAt(0)).sortBy(a => a) // ← trie sur le couple completIci a est le couple (Char, Iterable[String]). Trier « sur le couple » compare aussi les valeurs, ce qui n’a aucun sens et ralentit inutilement. L’écriture correcte cible la clé :
rddGby1.groupBy(x => x.charAt(0)).sortBy(_._1).collect()_._1 signifie « prends la première composante du tuple ». Pour l’ordre décroissant :
rddGby1.groupBy(x => x.charAt(0)).sortBy(_._1, ascending = false).collect()"Al".charAt(2) lève StringIndexOutOfBoundsException. Sur un jeu de données réel avec des noms courts, un groupBy(_.charAt(2)) fait planter le job. Filtrer avant :
rddGby1
.filter(_.length > 2)
.groupBy(_.charAt(2))
.sortBy(_._1)
.collect()Petit détail, gros gain de robustesse — surtout en production, où les données ont toujours l’exception qui casse tout.
groupByKey, reduceByKey, sortByKey exigent un RDD de paires (clé, valeur). Si vous partez d’un fichier texte plat, il faut d’abord parser :
val rddPaires = sc
.parallelize(List("0,11", "1,14", "0,3", "2,19", "1,3", "5,7"))
.map { ligne =>
val parties = ligne.split(",")
(parties(0), parties(1).toInt)
}Une écriture équivalente mais plus coûteuse :
.map(x => (x.split(",")(0), x.split(",")(1).toInt))Elle appelle split deux fois par ligne. Sur un million de lignes, c’est un million d’allocations mémoire évitables. La version avec val parties fait le découpage une fois et réutilise.
Résultat :
rddPaires.collect()
// Array(("0",11), ("1",14), ("0",3), ("2",19), ("1",3), ("5",7))Type : RDD[(String, Int)].
groupByKey — regrouper sans calculerrddPaires.groupByKey().collect()Regroupe les valeurs qui partagent la même clé, sans les combiner :
("0", Iterable(11, 3))
("1", Iterable(14, 3))
("2", Iterable(19))
("5", Iterable(7))Type : RDD[(String, Iterable[Int])].
C’est utile quand on a besoin de la liste des valeurs (par exemple pour concaténer des chaînes, ou pour appliquer un traitement qui a besoin de voir toutes les valeurs d’un coup). C’est presque toujours une mauvaise idée quand on veut juste agréger — voir la section suivante.
reduceByKey — combiner par clé, tout de suiterddPaires.reduceByKey((x, y) => x + y).collect()
// équivalent :
rddPaires.reduceByKey(_ + _).collect()Regroupe les valeurs de même clé et les combine avec la fonction fournie :
("0", 14) // 11 + 3
("1", 17) // 14 + 3
("2", 19)
("5", 7)Type : RDD[(String, Int)].
Voici la ligne que vous croiserez dans tous les articles Spark, et qui mérite qu’on l’explique une fois pour toutes.
reduceByKeypeut effectuer une partie du calcul localement sur chaque partition avant le shuffle.groupByKeydoit tout déplacer sur le réseau, puis regrouper.
Concrètement, sur un cluster de 4 exécuteurs avec 1 000 lignes chacun de la clé "0" :
Le shuffle transporte 4 000 entiers contre 4. Sur des vrais volumes, c’est un facteur qui se compte en centaines. Sur une clé chaude (un user_id très fréquent), c’est la différence entre un job qui tient en mémoire et un job qui explose l’exécuteur.
Règle absolue :
// À ne jamais écrire pour une somme :
rddPaires.groupByKey().mapValues(_.sum)
// Toujours écrire :
rddPaires.reduceByKey(_ + _)Les deux formes donnent le même résultat. Elles n’ont pas le même coût. La seconde peut être 5 à 100 fois plus rapide sur des données réelles, et surtout ne fait pas exploser la mémoire des exécuteurs qui héritent des grosses clés.
groupByKey reste utile — pour la concaténation de chaînes, pour un traitement qui a besoin de la liste complète — mais dès qu’il s’agit d’une opération associative et commutative (somme, max, min, count, moyenne recalculée à partir de sommes), reduceByKey est le bon outil.
Une limite subsiste : reduceByKey impose que le résultat ait le même type que les valeurs, ce qui rend une moyenne impossible d’un seul coup. C’est le rôle de aggregateByKey et combineByKey, traités dans au-delà de reduceByKey.
sortByKey — trier sur la clé, sans se poser de questionval trie = rddPaires.reduceByKey(_ + _).sortByKey()Résultat :
("0", 14)
("1", 17)
("2", 19)
("5", 7)Ordre décroissant :
rddPaires.reduceByKey(_ + _).sortByKey(ascending = false)sortByKey ne fonctionne que sur un RDD clé-valeur. Sur un RDD non-key/value, utilisez sortBy(_._1) (voir plus haut).
Les clés de l’exemple sont des String. Le tri est donc textuel, pas numérique :
sc.parallelize(List(("1", 1), ("2", 1), ("10", 1)))
.sortByKey()
.collect()
// Array(("1",1), ("10",1), ("2",1)) ← "10" avant "2"C’est correct au sens du tri des chaînes ("1" < "10" < "2"), mais presque jamais ce qu’on veut. Deux options :
Convertir la clé en Int à la lecture (préférable) :
val rddNumerique = sc
.parallelize(List("0,11", "1,14", "10,7", "2,19"))
.map { ligne =>
val parties = ligne.split(",")
(parties(0).toInt, parties(1).toInt) // ← clé Int
}
rddNumerique.reduceByKey(_ + _).sortByKey().collect()
// Array((0,11), (1,14), (2,19), (10,7))Padding zéro à gauche si vous devez garder des String (dates au format 01, 02, 10) — ordre textuel identique à l’ordre numérique.
| Méthode | Entrée | Sortie | Ce qu’elle fait |
|---|---|---|---|
groupBy | RDD[T] + fonction | RDD[(K, Iterable[T])] | Regroupe selon une clé calculée |
groupByKey | RDD[(K, V)] | RDD[(K, Iterable[V])] | Regroupe les valeurs par clé, sans combiner |
reduceByKey | RDD[(K, V)] | RDD[(K, V)] | Combine par clé avec combine local |
sortBy | RDD[T] + fonction | RDD[T] trié | Trie selon une fonction |
sortByKey | RDD[(K, V)] | RDD[(K, V)] trié | Trie sur la clé |
val resultat = sc
.parallelize(List("0,11", "1,14", "0,3", "2,19", "1,3", "5,7"))
.map { ligne =>
val parties = ligne.split(",")
(parties(0).toInt, parties(1).toInt)
}
.reduceByKey(_ + _)
.sortByKey()
resultat.collect().foreach(println)
// (0, 14)
// (1, 17)
// (2, 19)
// (5, 7)Cinq lignes qui synthétisent tout : parsing en clé-valeur, agrégation performante par clé, tri final déterministe. Ce squelette se recycle dans 80 % des pipelines RDD que vous écrirez.
Quatre opérations, une seule décision qui compte : reduceByKey chaque fois que la combinaison est associative et commutative, groupByKey uniquement quand on a réellement besoin de voir toutes les valeurs. C’est le premier réflexe que traque une revue de code Spark, et c’est presque toujours celui qui rapporte le plus — parce qu’il coûte peu à écrire et transforme le comportement du cluster sous charge.
Pour comprendre pourquoi le combine local change tout, l’article dédié aux partitions et au shuffle explique le mécanisme complet. Et pour voir cette différence à l’œuvre sur une jointure au lieu d’un groupBy, direction broadcast, shuffle et skew — la même logique de « moins de données sur le réseau = ordre de grandeur en performance » y opère sous un autre nom.
Développement et déploiement de solutions de données — les premiers modules sont en accès libre.