Partie 1
Les RDD
La brique de base de Spark : map, filter, reduce et le célèbre word count.
C'est quoi un RDD ?
Un RDD (Resilient Distributed Dataset) est une collection d'éléments (nombres, lignes de texte…) découpée en morceaux, les partitions. Chaque partition peut être traitée par un cœur différent, ou une machine différente : c'est ce qui rend Spark rapide sur de gros volumes.
On manipule un RDD avec deux familles d'opérations :
| Transformations | Actions |
|---|---|
| décrivent un calcul | déclenchent le calcul |
map, filter, flatMap, reduceByKey | collect, count, take, reduce |
| renvoient un nouveau RDD | renvoient un résultat Python |
Dans ces opérations, on passe souvent de petites fonctions appelées lambdas : lambda x: x * 2 est une fonction qui prend x et renvoie x * 2.
Créer un RDD
sc.parallelize transforme une liste Python en RDD. Ici, les entiers de 1 à 10. getNumPartitions() indique en combien de morceaux Spark l'a découpé.
nombres = sc.parallelize(range(1, 11))
print("Nombre de partitions :", nombres.getNumPartitions())Nombre de partitions : (souvent le même nombre que vos cœurs)
Une transformation : map
map applique une fonction à chaque élément. Ici, on met chaque nombre au carré.
Regardez bien ce qui s'affiche : pas de résultat, seulement une description du RDD ! Spark est paresseux (lazy) : une transformation ne calcule rien, elle note seulement la « recette ».
carres = nombres.map(lambda x: x * x)
print(carres)PythonRDD[...] at RDD at PythonRDD.scala:... (aucun nombre affiché)
Une action : collect
collect() est une action : elle force Spark à faire le calcul et rapatrie tous les résultats dans une liste Python. Si la Spark UI est ouverte, un nouveau job vient d'y apparaître.
Attention : collect() ramène tout sur votre machine. Parfait pour 10 nombres, dangereux pour des milliards de lignes. Sur de gros volumes, on préfère take(n) (les n premiers) ou count().
print(carres.collect())[1, 4, 9, 16, 25, 36, 49, 64, 81, 100]
Garder les nombres pairs
Créez un RDD avec les entiers de 1 à 100, gardez seulement les nombres pairs avec filter, puis comptez-les avec count().
entiers = sc.parallelize(range(1, 101))
pairs = entiers.filter(lambda x: ...)
print(pairs.count())Besoin d'un indice ?
Un nombre est pair quand le reste de sa division par 2 vaut 0 : x % 2 == 0.
Voir la solution
entiers = sc.parallelize(range(1, 101))
pairs = entiers.filter(lambda x: x % 2 == 0)
print(pairs.count())50
La somme des carrés
Calculez 1² + 2² + … + 100² en enchaînant un map (mettre au carré) puis un reduce (tout additionner).
reduce(lambda a, b: a + b) combine les éléments deux par deux jusqu'à n'en garder qu'un.
somme = entiers.map(...).reduce(...)
print(somme)Besoin d'un indice ?
map(lambda x: x * x) puis reduce(lambda a, b: a + b)
Voir la solution
somme = entiers.map(lambda x: x * x).reduce(lambda a, b: a + b)
print(somme)338350
Lire un fichier texte
On attaque le word count : compter combien de fois chaque mot apparaît dans un texte. C'est le « Hello World » du Big Data.
sc.textFile lit un fichier : chaque ligne devient un élément du RDD.
lignes = sc.textFile("data/texte.txt")
print("Nombre de lignes :", lignes.count())
lignes.take(3)Nombre de lignes : 20 ['Apache Spark est un moteur de calcul distribué.', 'Spark permet de traiter de grandes quantités de données sur plusieurs machines.', 'Avec Spark, les données sont découpées en partitions.']
Découper en mots : flatMap
ligne.split() découpe une ligne en liste de mots. Avec map, on obtiendrait un RDD de listes. flatMap fait la même chose puis met tout à plat : on obtient un RDD de mots.
mots = lignes.flatMap(lambda ligne: ligne.split())
mots.take(8)['Apache', 'Spark', 'est', 'un', 'moteur', 'de', 'calcul', 'distribué.']
Former des paires (mot, 1)
On transforme chaque mot en une paire (mot, 1). Le mot est la clé, le 1 est la valeur : chaque apparition du mot « compte pour 1 ».
paires = mots.map(lambda mot: (mot, 1))
paires.take(5)[('Apache', 1), ('Spark', 1), ('est', 1), ('un', 1), ('moteur', 1)]Additionner par mot : reduceByKey
reduceByKey regroupe les paires qui ont la même clé et additionne leurs valeurs. On trie ensuite par nombre d'apparitions, du plus grand au plus petit, et on affiche les 10 premiers.
compte = paires.reduceByKey(lambda a, b: a + b)
compte.sortBy(lambda paire: paire[1], ascending=False).take(10)Une liste de 10 paires ; les petits mots comme « de », « les » ou « est » arrivent en tête.
Un word count plus malin
Les petits mots (« de », « les »…) ne sont pas très intéressants. Améliorez le comptage :
- mettez chaque mot en minuscules et retirez la ponctuation en fin de mot ;
- ne gardez que les mots d'au moins 5 lettres ;
- affichez les 10 mots les plus fréquents.
top_mots = (lignes
.flatMap(lambda ligne: ligne.split())
.map(lambda mot: ...)
.filter(lambda mot: ...)
.map(lambda mot: (mot, 1))
.reduceByKey(lambda a, b: a + b)
.sortBy(lambda paire: paire[1], ascending=False))
top_mots.take(10)Besoin d'un indice ?
mot.lower().strip(".,:;!?") pour nettoyer, et len(mot) >= 5 pour filtrer.
Voir la solution
top_mots = (lignes
.flatMap(lambda ligne: ligne.split())
.map(lambda mot: mot.lower().strip(".,:;!?"))
.filter(lambda mot: len(mot) >= 5)
.map(lambda mot: (mot, 1))
.reduceByKey(lambda a, b: a + b)
.sortBy(lambda paire: paire[1], ascending=False))
top_mots.take(10)[('données', 11), ('spark', 9), ...]Les lignes qui parlent de Spark
Combien de lignes du texte contiennent le mot « spark », en majuscules ou en minuscules ?
nb = lignes.filter(lambda ligne: ...).count()
print(nb)Besoin d'un indice ?
Mettez la ligne en minuscules, puis testez "spark" in ....
Voir la solution
nb = lignes.filter(lambda ligne: "spark" in ligne.lower()).count()
print(nb)9
Astuce : les flèches ← → du clavier permettent aussi de naviguer.