Spark RDD par l’exemple : le word count en Scala

Le word count de Spark décortiqué ligne par ligne en Scala : flatMap, map, reduceByKey, saveAsTextFile sur un cas concret, avec les pièges classiques.

5 min de lecturesparkscalarddbig-data

Tout le monde ouvre son apprentissage de Spark par le même exercice : le word count. Ce n’est pas une coquetterie pédagogique. Il tient dans dix lignes de Scala et pourtant il vous fait manipuler, dans l’ordre, tout ce qui compte au quotidien : un RDD, une transformation à un vers plusieurs, une clé-valeur, une agrégation par clé, une écriture distribuée. Une fois ces cinq gestes lus une fois pour de vrai, le reste de l’API Spark se lit sans effort.

On prend un fichier etudiants.txt court, on compte combien de fois chaque prénom apparaît, et on discute à chaque étape ce que Spark vient vraiment de faire.

Le fichier de départ

Trois lignes, sept mots au total :

Ali Sara Ali
Sara Yassine
Ali Yassine

Résultat attendu :

(Ali, 3)
(Sara, 2)
(Yassine, 2)

1. Lire le fichier — un élément par ligne

scala
val rdd1 = sc.textFile("etudiants.txt")

Ce sc est le SparkContext, créé automatiquement par le spark-shellce qu’il est et en quoi il diffère de spark est expliqué à part.

Spark crée un RDD dont chaque élément est une ligne. Rien n’est encore lu : textFile construit un plan, il n’ouvre pas le fichier. Si la nature d’un RDD — collection distribuée, pas table, pas List un peu plus grosse — n’est pas encore claire, le modèle mental est posé ici.

Pour voir le contenu :

scala
rdd1.collect()
// Array("Ali Sara Ali", "Sara Yassine", "Ali Yassine")

collect() ramène tous les éléments distribués sur la machine qui pilote le job — le driver. C’est très pratique en développement, dangereux en production : sur un RDD de plusieurs millions de lignes, on transforme le driver en goulot d’étranglement, souvent jusqu’au OutOfMemoryError. Règle simple : collect() sur les exemples, take(n) ou show() sur les vrais volumes.

2. Découper chaque ligne en mots

scala
val rdd2 = rdd1.flatMap(line => line.split(" "))

split(" ") renvoie un tableau. On veut un seul RDD plat de mots, pas un RDD de tableaux — d’où flatMap au lieu de map :

OpérationRésultat
map(_.split(" "))Array(Array("Ali","Sara","Ali"), Array("Sara","Yassine"), Array("Ali","Yassine"))
flatMap(_.split(" "))Array("Ali","Sara","Ali","Sara","Yassine","Ali","Yassine")
scala
rdd2.collect()
// Array("Ali", "Sara", "Ali", "Sara", "Yassine", "Ali", "Yassine")

Retenez le contrat : map = une entrée, une sortie ; flatMap = une entrée, N sorties concaténées. Découper, tokeniser, exploser un JSON en événements : c’est toujours flatMap.

3. Étiqueter chaque mot avec un 1

scala
val rdd3 = rdd2.map(word => (word, 1))

On passe d’un RDD de String à un RDD de paires (clé, valeur). La clé sera le mot, la valeur 1 est la contribution de cette occurrence au décompte final.

scala
rdd3.collect()
// Array((Ali,1), (Sara,1), (Ali,1), (Sara,1), (Yassine,1), (Ali,1), (Yassine,1))

Cette étape a l’air décorative — elle ne l’est pas : les agrégateurs de Spark (reduceByKey, aggregateByKey, groupByKey) opèrent tous sur des paires. Sans clé, pas d’agrégation.

4. Additionner les valeurs par clé

scala
val rdd4 = rdd3.reduceByKey(_ + _)

reduceByKey regroupe les paires ayant la même clé et combine leurs valeurs deux à deux avec la fonction fournie. Pour Ali :

(Ali, 1) + (Ali, 1) → (Ali, 2)
(Ali, 2) + (Ali, 1) → (Ali, 3)

La syntaxe _ + _ est un raccourci Scala pour (a, b) => a + b. Les deux formes suivantes sont strictement équivalentes :

scala
rdd3.reduceByKey(_ + _)
rdd3.reduceByKey((a, b) => a + b)

Résultat :

scala
rdd4.collect()
// Array((Ali, 3), (Sara, 2), (Yassine, 2))

L’ordre peut varier d’une exécution à l’autre — Spark travaille en parallèle sur des partitions et ne garantit pas l’ordre à la sortie.

Pourquoi `reduceByKey` et pas `groupByKey` ?

groupByKey amène toutes les valeurs d’une clé sur un même exécuteur avant de combiner. Sur un dataset volumineux, un mot très fréquent (les articles, les mots vides) fait exploser la mémoire de l’exécuteur qui hérite de cette clé.

reduceByKey fait la somme localement dans chaque partition avant le shuffle. Le réseau ne voit passer qu’un total partiel par partition, pas la liste brute des 1. C’est presque toujours ce qu’on veut.

5. Écrire le résultat

scala
rdd4.saveAsTextFile("resultatEtudiants")

Spark ne crée pas un fichier resultatEtudiants.txt — il crée un dossier :

resultatEtudiants/
├── _SUCCESS
├── part-00000
└── part-00001

Chaque part-* correspond à une partition. La présence de _SUCCESS (vide) est la façon dont Spark signale que l’écriture s’est terminée sans erreur — plusieurs jobs en aval (Airflow, downstream Spark) l’attendent avant de démarrer.

À l’intérieur d’un part-* :

(Ali,3)
(Sara,2)
(Yassine,2)

Le piège : relancer saveAsTextFile sur un dossier existant

scala
rdd4.saveAsTextFile("resultatEtudiants") // OK la première fois
rdd4.saveAsTextFile("resultatEtudiants") // FileAlreadyExistsException

Spark refuse d’écraser silencieusement — c’est une protection, pas un bug. En développement, on résout ça de trois façons, par ordre de préférence :

scala
// 1. Changer le nom (traçabilité maximale)
rdd4.saveAsTextFile("resultatEtudiants-v2")

// 2. Supprimer le dossier avant
import scala.reflect.io.Directory
import java.io.File
new Directory(new File("resultatEtudiants")).deleteRecursively()
rdd4.saveAsTextFile("resultatEtudiants")

// 3. Passer au DataFrame API et son mode overwrite
rdd4.toDF("prenom", "n").write.mode("overwrite").csv("resultatEtudiants")

En production, on ne réécrit jamais au même chemin : on horodate le dossier de sortie (s3://bucket/wordcount/dt=2026-08-10/), on écrit à côté, on bascule le pointeur (une table Hive, un symlink S3, un manifeste) une fois que _SUCCESS est là. C’est ce qui permet de rejouer un job sans perdre l’ancien résultat.

Le pipeline complet, d’un coup d’œil

Et le programme dans sa forme minimale :

scala
val rdd1 = sc.textFile("etudiants.txt")
val rdd2 = rdd1.flatMap(line => line.split(" "))
val rdd3 = rdd2.map(word => (word, 1))
val rdd4 = rdd3.reduceByKey(_ + _)

rdd4.collect().foreach(println)
rdd4.saveAsTextFile("resultatEtudiants")

Ce qu’il faut retenir

Cinq étapes, cinq gestes qu’on retrouvera partout ailleurs : lire, transformer un vers plusieurs (flatMap), mettre en clé-valeur, agréger par clé, écrire en distribué. Chaque piège de cet exemple — collect naïf, oubli de flatMap, groupByKey ruineux, écriture qui refuse d’écraser — se rejouera à l’identique sur des datasets qui font mille fois la taille de ce fichier. C’est pour ça que le word count survit à chaque nouvelle version de Spark : il condense en dix lignes ce qui vous prendra dix ans à ne plus jamais rater.

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