Fichiers

Partie 1

Les RDD

La brique de base de Spark : map, filter, reduce et le célèbre word count.

Étape 1 / 12
1.1
Lecture

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 :

TransformationsActions
décrivent un calculdéclenchent le calcul
map, filter, flatMap, reduceByKeycollect, count, take, reduce
renvoient un nouveau RDDrenvoient 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.

1.2
Guidé

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é.

Tapez ce code dans une nouvelle cellule, puis exécutez-le
nombres = sc.parallelize(range(1, 11))
print("Nombre de partitions :", nombres.getNumPartitions())
Résultat attendu
Nombre de partitions : (souvent le même nombre que vos cœurs)
1.3
Guidé

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 ».

Tapez ce code dans une nouvelle cellule, puis exécutez-le
carres = nombres.map(lambda x: x * x)
print(carres)
Résultat attendu
PythonRDD[...] at RDD at PythonRDD.scala:... (aucun nombre affiché)
1.4
Guidé

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().

Tapez ce code dans une nouvelle cellule, puis exécutez-le
print(carres.collect())
Résultat attendu
[1, 4, 9, 16, 25, 36, 49, 64, 81, 100]
1.5
Exercice

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().

À vous : complétez les « ... »
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
Solution
entiers = sc.parallelize(range(1, 101))
pairs = entiers.filter(lambda x: x % 2 == 0)
print(pairs.count())
Résultat attendu
50
1.6
Exercice

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.

À vous : complétez les « ... »
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
Solution
somme = entiers.map(lambda x: x * x).reduce(lambda a, b: a + b)
print(somme)
Résultat attendu
338350
1.7
Guidé

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.

Tapez ce code dans une nouvelle cellule, puis exécutez-le
lignes = sc.textFile("data/texte.txt")
print("Nombre de lignes :", lignes.count())
lignes.take(3)
Résultat attendu
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.']
1.8
Guidé

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.

Tapez ce code dans une nouvelle cellule, puis exécutez-le
mots = lignes.flatMap(lambda ligne: ligne.split())
mots.take(8)
Résultat attendu
['Apache', 'Spark', 'est', 'un', 'moteur', 'de', 'calcul', 'distribué.']
1.9
Guidé

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 ».

Tapez ce code dans une nouvelle cellule, puis exécutez-le
paires = mots.map(lambda mot: (mot, 1))
paires.take(5)
Résultat attendu
[('Apache', 1), ('Spark', 1), ('est', 1), ('un', 1), ('moteur', 1)]
1.10
Guidé

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.

Tapez ce code dans une nouvelle cellule, puis exécutez-le
compte = paires.reduceByKey(lambda a, b: a + b)
compte.sortBy(lambda paire: paire[1], ascending=False).take(10)
Résultat attendu
Une liste de 10 paires ; les petits mots comme « de », « les » ou « est » arrivent en tête.
1.11
Exercice

Un word count plus malin

Les petits mots (« de », « les »…) ne sont pas très intéressants. Améliorez le comptage :

  1. mettez chaque mot en minuscules et retirez la ponctuation en fin de mot ;
  2. ne gardez que les mots d'au moins 5 lettres ;
  3. affichez les 10 mots les plus fréquents.
À vous : complétez les « ... »
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
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)
Résultat attendu
[('données', 11), ('spark', 9), ...]
1.12
Bonus

Les lignes qui parlent de Spark

Combien de lignes du texte contiennent le mot « spark », en majuscules ou en minuscules ?

Facultatif : à faire si vous le souhaitez
nb = lignes.filter(lambda ligne: ...).count()
print(nb)
Besoin d'un indice ?

Mettez la ligne en minuscules, puis testez "spark" in ....

Voir la solution
Solution
nb = lignes.filter(lambda ligne: "spark" in ligne.lower()).count()
print(nb)
Résultat attendu
9

Astuce : les flèches ← → du clavier permettent aussi de naviguer.