ISIMa AU 2025-2026
Master SD(2ème Année)
Frameworks Big Data 1
TP N◦2 : Spark & RDD(1)
(durée : 1h30 heures)
Objectifs
À la fin de cette activité, vous pourrez :
1. Découvrir Spark, un outil de calcul distribué très efficace ;
2. Savoir manipuler, de près, les RDDs sous Spark ;
3. Refaire l’exercice WordCount en utilisant les RDDs sous Spark.
1 À propos Spark & RDD
Spark Core est la base de l’ensemble du projet. Il fournit des fonctionnalités
de répartition de tâches distribuées, de planification et d’E/S de base.
Spark utilise une structure de données fondamentale spécialisée appelée
RDD (Resilient Distributed Datasets) qui est une collection logique de données
partitionnées sur plusieurs machines.
Les RDD peuvent être créés de deux manières :
1. l’une consiste à référencer des ensembles de données dans des systèmes de
stockage externes et
2. la seconde consiste à appliquer des transformations (par exemple, map,
filter, reducer, join) sur des RDD existants.
L’abstraction RDD est exposée via une API intégrée au langage. Cela sim-
plifie la complexité de la programmation car la façon dont les applications
manipulent les RDD est similaire à la manipulation de collections de données
locales.
2 Spark Shell
Spark fournit un shell interactif, un outil puissant pour analyser les données
de manière interactive. Il est disponible en langage Scala ou Python. L’abs-
traction principale de Spark est une collection distribuée d’éléments appelée
ensemble de données distribuées résilientes (RDD). Les RDD peuvent être
créés à partir de formats d’entrée Hadoop (tels que des fichiers HDFS) ou en
transformant d’autres RDD.
1. Rania MKHININI GAHAR / [Link]
1
2.1 Ouvrir le shell Spark
2.2 Créer un simple RDD
Créons un RDD simple à partir du fichier texte. Utilisez la commande sui-
vante pour créer un RDD simple.
La sortie de la commande ci-dessus est :
L’API Spark RDD introduit quelques transformations et quelques actions
pour manipuler RDD.
3 Transformations RDD
Les transformations RDD renvoient un pointeur vers un nouveau RDD
et vous permettent de créer des dépendances entre les RDD. Chaque RDD dans
la chaı̂ne de dépendances (String of Dependencies) possède une fonction pour
calculer ses données et possède un pointeur (dépendance) vers son RDD parent.
Spark est paresseux, donc rien ne sera exécuté à moins que vous n’appe-
liez une transformation ou une action qui déclenchera la création et l’exécution
d’une tâche. Regardez l’extrait suivant de l’exemple de comptage de mots(WordCount).
Par conséquent, la transformation RDD n’est pas un ensemble de données
mais une étape dans un programme (peut-être la seule étape) indiquant à Spark
comment obtenir des données et quoi en faire.
Vous trouverez ci-dessous une liste de transformations RDD.
1. map(func) : Renvoie un nouvel ensemble de données distribué, formé en
passant chaque élément de la source via une fonction func.
2. filter(func) : Renvoie un nouvel ensemble de données formé en sélectionnant
les éléments de la source sur lesquels func renvoie true.
3. flatMap(func) : Similaire à map, mais chaque élément d’entrée peut
être mappé à 0 ou plusieurs éléments de sortie (donc func doit renvoyer
une séquence plutôt qu’un seul élément).
2
4. mapPartitions(func) : Similaire à map, mais s’exécute séparément sur
chaque partition (bloc) du RDD, donc func doit être de type Iterator<T>⇒Iterator<U>
lors de l’exécution sur un RDD de type T.
5. mapPartitionsWithIndex(func) : Similaire aux mapPartitions , mais
fournit également à func une valeur entière représentant l’index de la
partition, donc func doit être de type (Int, Iterator<T>) ⇒Iterator<U>
lors de l’exécution sur un RDD de type T.
6. sample(withReplacement, fraction, seed) : Échantillonnez une frac-
tion des données, avec ou sans Replacement, à l’aide d’une graine de
générateur de nombres aléatoire donnée.
7. union(otherDataset) : Renvoie un nouvel ensemble de données qui
contient l’union des éléments de l’ensemble de données source et de l’ar-
gument.
8. intersection(otherDataset) : Renvoie un nouveau RDD qui contient
l’intersection des éléments de l’ensemble de données source et de l’argu-
ment.
9. distinct([numTasks]) : Renvoie un nouvel ensemble de données conte-
nant les éléments distincts de l’ensemble de données source.
10. groupByKey([numTasks]) : Lorsqu’il est appelé sur un ensemble de
données de paires (K, V), renvoie un ensemble de données de paires (K,
Iterable<V>).
11. reduceByKey(func, [numTasks]) : Lorsqu’il est appelé sur un en-
semble de données de paires (K, V), renvoie un ensemble de données de
paires (K, V) où les valeurs de chaque clé sont agrégées à l’aide de la fonc-
tion de rduction donnée func, qui doit être de type (V, V) ⇒V. Comme
dans groupByKey, le nombre de tâches de réduction est configurable
via un deuxième argument facultatif.
12. aggregateByKey(zeroValue)(seqOp, combOp, [numTasks]) : Lors-
qu’elle est appelée sur un ensemble de données de paires (K, V), renvoie
un ensemble de données de paires (K, U) où les valeurs de chaque clé
sont agrégées à l’aide des fonctions de combinaison données et d’une va-
leur neutre zéro . Permet un type de valeur agrégée différent du type
de valeur d’entrée, tout en évitant les allocations inutiles. Comme dans
groupByKey, le nombre de tâches de réduction est configurable via un
deuxième argument facultatif.
13. sortByKey([ascending], [numTasks]) : Lorsqu’il est appelé sur un en-
semble de données de paires (K, V) où K implémente Ordered, renvoie un
ensemble de données de paires (K, V) triées par clés dans l’ordre croissant
ou décroissant, comme spécifié dans l’argument booléen croissant.
14. join(otherDataset, [numTasks]) : Lorsqu’elle est appelée sur des en-
sembles de données de type (K, V) et (K, W), renvoie un ensemble de
données de paires (K, (V, W)) avec toutes les paires d’éléments pour
chaque clé. Les jointures externes sont prises en charge via leftOuterJoin,
rightOuterJoin et fullOuterJoin.
3
15. cogroup(otherDataset, [numTasks]) : Lorsqu’elle est appelée sur des
ensembles de données de type (K, V) et (K, W), renvoie un ensemble de
données de tuples (K, (Iterable<V>, Iterable<W>)). Cette opération est
également appelée group With.
16. cartesian(otherDataset) : Lorsqu’il est appelé sur des ensembles de
données de types T et U, renvoie un ensemble de données de paires (T,
U) (toutes les paires d’éléments).
17. pipe(command, [envVars]) : Acheminez chaque partition du RDD via
une commande shell, par exemple un script Perl ou bash. Les éléments
RDD sont écrits dans l’entrée standard du processus et les lignes de sortie
vers sa sortie standard sont renvoyées sous forme de RDD de chaı̂nes.
18. coalesce(numPartitions) : Réduisez le nombre de partitions dans le
RDD à numPartitions. Utile pour exécuter des opérations plus efficace-
ment après avoir filtré un grand ensemble de données.
19. repartition(numPartitions) : Réorganisez les données dans le RDD de
manière aléatoire pour créer plus ou moins de partitions et les équilibrer
entre elles. Cela permet de toujours mélanger toutes les données sur le
réseau.
20. repartitionAndSortWithinPartitions(partitioner) : Répartitionnez
le RDD en fonction du partitionneur donné et, dans chaque partition
résultante, triez les enregistrements par leurs clés. Cette méthode est plus
efficace que l’appel à la repartition puis au tri dans chaque partition, car
elle peut déplacer le tri vers le bas dans le mécanisme de mélange.
4 Actions
La description suivante donne une liste d’actions qui renvoient des valeurs.
1. reduce(func) : Agréger les éléments de l’ensemble de données à l’aide
d’une fonction func (qui prend deux arguments et en renvoie un). La
fonction doit être commutative et associative pour pouvoir être calculée
correctement en parallèle.
2. collect() : Renvoie tous les éléments de l’ensemble de données sous forme
de tableau au niveau du programme pilote. Cela est généralement utile
après un filtre ou une autre opération qui renvoie un sous-ensemble suffi-
samment petit des données.
3. count() : Renvoie le nombre d’éléments dans l’ensemble de données.
4. first() : Renvoie le premier élément de l’ensemble de données (similaire
à (1)).
5. take(n) : Renvoie un tableau avec les n premiers éléments de l’ensemble
de données.
6. takeSample (withReplacement,num, [seed]) : Renvoie un tableau
avec un échantillon aléatoire de num éléments de l’ensemble de données,
avec ou sans remplacement, en pré-spécifiant éventuellement une graine
de générateur de nombres aléatoires.
4
7. takeOrdered(n, [ordering]) : Renvoie les n premiers éléments du RDD
en utilisant soit leur ordre naturel, soit un comparateur personnalisé.
8. saveAsTextFile(path) : Écrit les éléments de l’ensemble de données
sous forme de fichier texte (ou d’ensemble de fichiers texte) dans un
répertoire donné du système de fichiers local, de HDFS ou de tout autre
système de fichiers pris en charge par Hadoop. Spark appelle toString sur
chaque élément pour le convertir en ligne de texte dans le fichier.
9. saveAsSequenceFile(path) (Java and Scala) : Écrit les éléments de
l’ensemble de données sous forme de fichier de squence Hadoop dans un
chemin donné du système de fichiers local, de HDFS ou de tout autre
système de fichiers pris en charge par Hadoop. Cette fonction est dispo-
nible sur les RDD de paires clé-valeur qui implémentent l’interface Wri-
table de Hadoop. Dans Scala, elle est également disponible sur les types
qui sont implicitement convertibles en Writable (Spark inclut des conver-
sions pour les types de base comme Int, Double, String, etc.).
10. saveAsObjectFile(path) (Java and Scala) : Écrit les éléments de l’en-
semble de données dans un format simple à l’aide de la sérialisation Java,
qui peut ensuite être chargé à l’aide de [Link]().
11. countByKey() : Disponible uniquement sur les RDD de type (K, V).
Renvoie une table de hachage de paires (K, Int) avec le nombre de chaque
clé.
12. foreach(func) : Exécute une fonction func sur chaque élément de l’en-
semble de données. Cette opération est généralement effectuée pour des
effets secondaires tels que la mise jour d’un accumulateur ou l’interaction
avec des systèmes de stockage externes.
5 Programmation avec RDD
Voyons les implémentations de quelques transformations et actions RDD
dans la programmation RDD à l’aide d’un exemple.
5.1 Exemple
Prenons l’exemple de WordCount : il compte chaque mot apparaissant dans
un document. Considérez le texte suivant comme une entrée et enregistrez-le
sous forme de fichier [Link] dans un répertoire personnel.
[Link], fichier d’entrée.
Suivez la procédure ci-dessous pour exécuter l’exemple donné.
5
5.2 Ouvrir le shell Spark
La commande suivante permet d’ouvrir le shell Spark. En général, Spark
est construit à l’aide de Scala. Par conséquent, un programme Spark s’exécute
dans un environnement Scala.
Si le shell Spark s’ouvre correctement, vous trouverez la sortie suivante. Re-
gardez la dernière ligne de la sortie Contexte Spark disponible en tant que sc
signifie que le conteneur Spark est automatiquement créé en tant qu’objet de
contexte Spark avec le nom sc. Avant de commencer la première étape d’un
programme, l’objet SparkContext doit être créé.
5.3 Création d’un RDD
Tout d’abord, nous devons lire le fichier d’entrée à l’aide de l’API Spark-
Scala et créer un RDD.
La commande suivante est utilisée pour lire un fichier à partir d’un emplace-
ment donné. Ici, un nouveau RDD est créé avec le nom de inputfile. La chaı̂ne
qui est donnée comme argument dans la méthode textFile() est le chemin ab-
solu pour le nom du fichier d’entrée. Cependant, si seul le nom du fichier est
donné, cela signifie que le fichier d’entrée se trouve à l’emplacement actuel.
6
5.4 Exécuter la transformation WordCount
Notre objectif est de compter les mots dans un fichier. Créez une map plate
pour diviser chaque ligne en mots (flatMap(line⇒[Link](” ”)).
Ensuite, lisez chaque mot comme une clé avec une valeur 1 (<key, value>
= <word,1>) en utilisant la fonction map (map(word ⇒(word, 1)).
Enfin, réduisez ces clés en ajoutant des valeurs de clés similaires (reduceByKey( + )).
La commande suivante est utilisée pour exécuter la logique de comptage de
mots. Après l’avoir exécutée, vous ne trouverez aucune sortie car il ne s’agit
pas d’une action, mais d’une transformation ; pointer vers un nouveau RDD ou
indiquer à Spark ce qu’il doit faire avec les données données.
5.5 RDD actuel
Lorsque vous travaillez avec le RDD, si vous souhaitez en savoir plus sur le
RDD actuel, utilisez la commande suivante. Elle vous montrera la description
du RDD actuel et de ses dépendances pour le débogage.
5.6 Mise en cache des transformations
Vous pouvez marquer un RDD comme devant être conservé à l’aide des
méthodes persist() ou cache(). La première fois qu’il est calculé dans une
action, il sera conservé en mémoire sur les nœuds. Utilisez la commande suivante
pour stocker les transformations intermédiaires en mémoire.
5.7 Application de l’action
L’application d’une action, comme le stockage de toutes les transforma-
tions, génère un fichier texte. L’argument String de la méthode saveAsText-
File(” ”) est le chemin absolu du dossier de sortie. Essayez la commande
suivante pour enregistrer la sortie dans un fichier texte. Dans l’exemple sui-
vant, le dossier output se trouve à l’emplacement actuel.
7
5.8 Vérification de la sortie
Ouvrez un autre terminal pour accéder au répertoire personnel (où Spark est
exécuté dans l’autre terminal). Utilisez les commandes suivantes pour vérifier
le répertoire de sortie.
La commande suivante est utilisée pour voir la sortie des fichiers Part-00000.
La commande suivante est utilisée pour voir la sortie des fichiers Part-00001.
8
6 Maintien du Stockage via UN-Persist
Avant de désactiver la persistance, si vous souhaitez voir l’espace de stockage
utilisé pour cette application, utilisez l’URL suivante dans votre navigateur.
Vous verrez l’écran suivant, qui montre l’espace de stockage utilisé pour les
applications, qui s’exécutent sur le shell Spark.
Si vous souhaitez annuler la persistance de l’espace de stockage d’un RDD
particulier, utilisez la commande suivante.
Vous verrez le résultat comme suit :
9
Pour vérifier l’espace de stockage dans le navigateur, utilisez l’URL suivante.
Vous verrez l’écran suivant. Il montre l’espace de stockage utilisé pour les ap-
plications excutées sur le shell Spark.
10