Fiche de révision : Introduction à Spark et Big Data

Plan du Cours

  1. Fondamentaux et architecture de Spark
  2. Big Data et écosystème Spark
  3. Architecture interne et évaluation paresseuse
  4. RDD, transformations et actions
  5. DataFrames et Spark SQL
  6. Spark Streaming et modes de sortie
  7. Fenêtrage, watermark et checkpoints
  8. MLlib : préparation des données
  9. Algorithmes de machine learning
  10. Déploiement et deep learning avec Spark

1. Fondamentaux et architecture de Spark

Notions clés & Définitions

  • Big Data : Ensemble de données caractérisées notamment par de grands volumes, une vélocité élevée, une grande variété, une qualité variable (véracité) et une valeur métier extraite.
  • Écosystème Apache Spark : Plateforme unifiée qui fournit des briques pour traiter des données structurées, des flux en temps réel et des tâches de machine learning distribuées.
  • Architecture Master/Slave : Modèle distribué de Spark où un driver prépare le plan, un gestionnaire orchestre les ressources et des executors exécutent les tâches sur les nœuds.
  • Lazy Evaluation : Évaluation paresseuse dans laquelle les transformations sont enregistrées dans un DAG sans calculer tout de suite, et les actions déclenchent l’exécution.
  • SparkSession : Point d’entrée unique introduit à partir de Spark 2.0 pour accéder aux fonctionnalités de Spark et initialiser l’environnement.

Points essentiels

  • Spark vise un changement de paradigme par rapport à MapReduce : les calculs sont faits en mémoire (In-Memory) plutôt que d’écrire sur disque à chaque étape, ce qui accélère fortement les algorithmes itératifs jusqu’à 100×.
  • Dans l’architecture distribuée, le Driver Program contient main(), crée la SparkSession et transforme le code en plan d’exécution (DAG) avant d’ordonnancer le travail.
  • Le Cluster Manager (YARN, Kubernetes ou Standalone) alloue les ressources CPU et RAM, tandis que les Executors exécutent les tâches et stockent les données en mémoire ou sur disque.
  • Un Job est découpé en Stages, puis en Tasks, et une Task correspond à la plus petite unité envoyée à un executor.
  • En Spark, les Transformations (ex. map, filter, groupBy) ne s’exécutent pas immédiatement, alors que les Actions (ex. count, collect, save) déclenchent le calcul réel et renvoient un résultat au driver.
  • Avec Spark 2.0, la SparkSession devient l’unique point d’entrée pour utiliser les fonctionnalités de Spark.

Astuce mémo

5V de Big Data : V O Vélocité Variété Véracité Valeur (V=Volume, V=Vélocité, V=Variété, V=Véracité, V=Valeur).

2. Big Data et écosystème Spark

Notions clés & Définitions

  • Spark SQL : Brique de Spark pour manipuler des données structurées via des DataFrames et exécuter des requêtes de type SQL.
  • Spark Streaming : Brique de Spark pour traiter des flux de données en temps réel plutôt que des jeux figés.

Points essentiels

  • Le Big Data se caractérise par 5V : Volume, Vélocité, Variété, Véracité et Valeur.
  • Spark est souvent plus rapide que Hadoop MapReduce sur les algorithmes itératifs car il calcule en mémoire plutôt que d’écrire sur disque après chaque étape, jusqu’à 100 fois plus vite.
  • Spark SQL sert de passerelle pour interroger et transformer des données structurées via SQL et DataFrames, avec un traitement optimisé par l’architecture Spark.
  • Spark Streaming permet de traiter des données générées en continu, en temps réel ou quasi temps réel.
  • L’écosystème Spark inclut aussi MLlib pour le machine learning distribué et GraphX pour le traitement de graphes.

Astuce mémo

5V = Volume Vélocité Variété Véracité Valeur : Volume+Vélocité arrivent vite, variété change, véracité manque parfois, valeur justifie le projet.

3. Architecture interne et évaluation paresseuse

Notions clés & Définitions

  • Lineage Spark : Ligne d’historique des transformations d’un RDD qui permet à Spark de reconstruire les données si une partition est perdue.
  • Partitionnement distribué : Découpage des données en partitions réparties sur plusieurs machines pour exécuter le calcul en parallèle.
  • RDD immuable : Objet RDD qui ne se modifie jamais : toute opération produit un nouveau RDD via une transformation.
  • Transformations paresseuses : Transformations qui définissent une logique de calcul sans l’exécuter immédiatement sur le cluster.
  • Actions immédiates : Actions qui déclenchent l’exécution et renvoient un résultat au Driver côté Spark.

Points essentiels

  • En cas de perte d’un nœud, Spark reconstruit les partitions perdues grâce au Lineage, c’est-à-dire l’historique des transformations.
  • Les transformations Lazy ne lancent aucun calcul tant qu’aucune action Eager n’est exécutée.
  • Une action (comme collect, count ou take) déclenche le calcul et renvoie le résultat au Driver.
  • Par défaut, un RDD est recalculé à chaque action, sauf si on le persiste via cache() ou persist(level).
  • persist(level) permet de choisir où stocker les résultats (RAM, disque, ou sérialisés) pour éviter des recalculs.

Astuce mémo

Lazy = j’écris la recette ; Action = je cuisine et je ramène le résultat au Driver.

4. RDD, transformations et actions

Notions clés & Définitions

  • Transformation Spark : Une transformation Spark produit un nouveau jeu de données à partir d’un jeu existant sans exécuter le calcul tout de suite.
  • Action Spark : Une action Spark déclenche l’exécution du traitement distribué et renvoie un résultat exploitable (par exemple pour écrire une sortie).
  • Variable Broadcast : Une variable broadcast permet de distribuer efficacement un objet (comme un dictionnaire) à tous les nœuds pour éviter de le recopier à chaque tâche.
  • Accumulateur Spark : Un accumulateur sert à compter ou agréger des événements globaux pendant un traitement distribué avec un mécanisme fiable.

Points essentiels

  • La fonction de traitement principal s’écrit comme une transformation, puis on la valide avec un test unitaire nommé test_log_processor.py qui affiche OK si tout passe.
  • Pour traduire des codes de sévérité, le dictionnaire n’est pas codé en dur par tâche mais transformé en variable Broadcast avec les clés INF, WRN et ERR.
  • L’accumulateur sert à compter des lignes malformées, y compris les cas sans séparateur :, les lignes vides successives après une ligne marquée, et les lignes sans correspondance dans le dictionnaire de broadcast.
  • Un test unitaire de log doit être exécuté via docker exec -it spark-master python3 /opt/spark/work-dir/test_log_processor.py et réussir avant le déploiement du job Spark.

Astuce mémo

Broadcast = dictionnaire distribué une fois ; Accumulateur = compteur global fiable en distribué.

5. DataFrames et Spark SQL

Notions clés & Définitions

  • DataFrame Spark : Un DataFrame est une table distribuée structurée où les colonnes ont un schéma défini et où on applique des transformations SQL-like.
  • explode : explode est une fonction qui transforme une colonne contenant une collection en plusieurs lignes, une par élément.
  • UDF Python : Une UDF Python est une fonction définie par l’utilisateur qui s’exécute côté Python et n’est pas analysée par Catalyst pour optimiser le plan.
  • Fonctions natives Spark SQL : Les fonctions natives de pyspark.sql.functions sont fournies par Spark SQL et s’exécutent directement dans le moteur JVM.
  • Pandas UDF : Une Pandas UDF est une variante de fonction utilisateur qui traite des lots via Apache Arrow pour réduire les coûts d’échange Python↔Spark.

Points essentiels

  • Spark 3.5+ formalise les Python UDTF, mais explode reste l’approche standard pédagogique pour créer une ligne par élément.
  • Catalyst ne peut pas optimiser le code à l’intérieur d’une UDF Python car il y voit une boîte noire, ce qui limite les optimisations de requête.
  • Priorité aux fonctions natives pyspark.sql.functions (ex: upper(), date_format()) car elles s’exécutent directement en JVM.
  • Les Pandas UDF sont souvent plus rapides sur gros volumes car Apache Arrow transfère des données par blocs au lieu de ligne par ligne.
  • L’exemple de classification utilise INCONNU si montant vaut NULL, PETIT si montant < 50, MOYEN si montant < 200, sinon GROS.
  • Dans le TP CSV -> Spark, définir le schéma manuellement évite inferSchema=True qui nécessite 2 passes sur le fichier et devient lent à grande échelle.

Astuce mémo

Native en JVM > Pandas UDF en blocs (Arrow) > UDF Python boîte noire (pas d’optimisation Catalyst).

6. Spark Streaming et modes de sortie

Notions clés & Définitions

  • Streaming DataFrameReader : Un lecteur de flux sert à créer un DataFrame vivant à partir d’une source continue en utilisant spark.readStream pour ingérer les données au fil de l’eau.
  • writeStream : Un DataFrameWriter de streaming décrit le flux de sortie en utilisant df.writeStream pour publier en continu les résultats produits par Spark.
  • Mode de sortie : Un paramètre de sortie indique si Spark écrit uniquement les nouveautés, tout le résultat, ou seulement les lignes modifiées à chaque déclenchement.

Points essentiels

  • Le passage Batch→Streaming se fait surtout par le changement de fonctions : spark.read devient spark.readStream et df.write devient df.writeStream.
  • Le mode Append (défaut) n’écrit que les nouvelles lignes apparues depuis le dernier déclenchement.
  • Le mode Complete réécrit toute la table de résultat à chaque micro-batch, ce qui est requis pour des agrégations.
  • Le mode Update écrit uniquement les lignes dont la valeur a changé depuis le dernier déclenchement.
  • Dans l’exemple socket, l’écriture vers la console utilise un outputMode("complete") avec format("console") et start() pour afficher chaque mise à jour globale.

Astuce mémo

Append=ajoute, Complete=réécrit tout, Update=met à jour seulement les lignes changées.

7. Fenêtrage, watermark et checkpoints

Notions clés & Définitions

  • Fenêtrage glissant : Un fenêtrage glissant regroupe des événements dans une fenêtre de temps fixe qui avance régulièrement, ce qui permet de cumuler des montants sur des périodes qui se chevauchent.
  • Watermark : Une watermark limite l’attente des événements tardifs en précisant le délai toléré par Spark, afin de déclencher l’agrégation sans attendre indéfiniment.
  • Checkpointing : Le checkpointing enregistre l’avancement et l’état du traitement du flux sur le disque pour reprendre automatiquement au bon endroit après un arrêt ou un crash.
  • Mode update : Le mode update n’envoie à la sortie (et donc à MongoDB) que les vendeurs dont le total a effectivement changé dans chaque micro-lot.

Points essentiels

  • La fenêtre testée a une taille de 30 secondes avec un glissement de 10 secondes, et l’alerte se déclenche quand la somme dépasse 1000.
  • Le code d’alerte utilise un watermark sur la colonne ts avec un délai de 1 minute avant de grouper par window(30 seconds, 10 seconds) et de filtrer sum(montant) > 1000.
  • Le checkpointing sauvegarde l’état de progression sur disque, et si le script plante, Spark redémarre là où il s’était arrêté.
  • Si le code change le schéma (par exemple des noms de colonnes) et que vous relancez avec l’ancien checkpoint, Spark peut planter, ce qui impose de supprimer le dossier checkpoints/ pour repartir à zéro.
  • En mode update, MongoDB ne reçoit que les mises à jour des vendeurs dont le total a changé, tandis qu’en mode complete Spark renvoie toute la table à chaque micro-batch (coûteux si elle est grande).

Astuce mémo

Watermark (1 minute) = “retard accepté” avant de fermer les fenêtres.

8. MLlib : préparation des données

Notions clés & Définitions

  • Feature Engineering : La préparation des données consiste à convertir les colonnes d’entrée en représentations compatibles avec les modèles ML de Spark.
  • StringIndexer : Un algorithme qui remplace une colonne catégorielle par des indices numériques en apprenant la table de correspondance lors du fit.
  • OneHotEncoder : Un transformeur qui transforme les indices de catégories en vecteurs binaires à présenter comme variables d’entrée.
  • VectorAssembler : Un transformeur qui regroupe plusieurs colonnes numériques (dont celles issues des encodages) dans une seule colonne de type vecteur, généralement nommée features.

Points essentiels

  • Dans Spark MLlib, les modèles attendent en entrée une colonne unique de type vecteur (souvent features) contenant toutes les variables explicatives.
  • StringIndexer apprend les identifiants numériques des catégories en scannant toute la colonne lors du fit avant de produire la colonne indexée.
  • OneHotEncoder convertit les indices de catégories en vecteurs binaires (un vecteur par ligne).
  • Dans le TP 5, le label est créé ainsi : 1 si montant > 150, sinon 0.
  • Le pipeline de préparation pour le TP 5 assemble les variables sous la forme d’un vecteur features à partir de la sortie d’encodage vendeur et de montant.

Astuce mémo

Indexer → OneHot → Assembler : Catégories → indices → vecteur features.

9. Algorithmes de machine learning

Notions clés & Définitions

  • Régression Linéaire : L’algorithme de régression linéaire apprend une relation entre des variables explicatives features et une cible numérique continue à prédire.
  • RandomForestClassifier : Le classifieur Random Forest combine les prédictions de nombreux arbres entraînés sur des sous-ensembles de données pour produire une prédiction plus stable.
  • KMeans : K-means est un algorithme de clustering non supervisé qui regroupe des données similaires en k clusters sans variable label à prédire.
  • LogisticRegression : La régression logistique réalise une classification binaire en estimant la probabilité d’appartenir à la classe 1 via une fonction sigmoïde.

Points essentiels

  • LinearRegression s’entraîne avec fit(train_data) et fournit ensuite des prédictions via transform (Estimator devient Transformer après fit).
  • RandomForestClassifier entraîne numTrees=20 arbres et peut afficher l’importance des variables via featureImportances.
  • KMeans ne possède pas de label à prédire : il cherche des groupes de données cohérents, par exemple avec k=3 et des centres calculés par le modèle.
  • LogisticRegression prédit une probabilité entre 0 et 1 puis convertit en classe selon un seuil par défaut de 0.5.
  • Si la probabilité prédite dépasse le seuil (0.5), Spark affecte la prédiction à 1, sinon à 0.

Astuce mémo

Continue→Régression linéaire, Beaucoup d’arbres→Random Forest, Groupes→K-means, Probabilité→Logistic(sigmoïde, seuil 0.5).

10. Déploiement et deep learning avec Spark

Notions clés & Définitions

  • Broadcast Spark : Le broadcast Spark diffuse un objet Python depuis le driver vers tous les workers pour éviter de le renvoyer à chaque exécution.
  • pandas_udf : Une pandas_udf décrit une UDF exécutée par lots avec des structures Pandas, ce qui facilite le passage et le calcul côté Python.
  • PipelineModel : Un PipelineModel stocke et charge un pipeline ML entraîné pour appliquer la même préparation et l’inférence lors du scoring.
  • Dockerfile personnalisé Spark : Un Dockerfile personnalisé définit une image Spark incluant les bibliothèques nécessaires au deep learning, assurant un environnement identique sur master et workers.

Points essentiels

  • Le scoring temps réel charge le modèle avec PipelineModel depuis HDFS, applique une préparation identique aux données, puis exécute l’inférence avant d’écrire vers MongoDB.
  • En production Spark-Keras, un entraînement est fait sur une machine puissante (avec GPU), puis le modèle est diffusé aux workers via broadcast pour prédire en parallèle sur des millions de lignes.
  • Le mode le plus courant pour l’inférence Spark/Keras utilise une pandas_udf et un modèle diffusé, chaque worker prédisant sur son lot via NumPy.
  • Le Dockerfile doit partir de apache/spark:3.5.0, installer tensorflow pandas pyarrow numpy, puis repasser à USER 1001 pour installer les dépendances sur l’image de cluster.
  • Pour que master et workers aient les mêmes versions, docker-compose doit remplacer la clé image par build: . et reconstruire avec docker-compose up -d --build.
  • La validation rapide consiste à exécuter sur spark-master: import tensorflow puis afficher la version, afin de vérifier que le worker importera bien TensorFlow.

Astuce mémo

Broadcast = “je l’envoie une fois” : le modèle (ou ses poids) part du driver vers tous les workers, puis chaque worker prédit son lot avec une pandas_udf.

Repères chronologiques

DateÉvénement
2026-03-01Date de vente dans le fichier ventes_v2.csv du TP 5
2026-03-02Date de vente dans le fichier ventes_v2.csv du TP 5
2026-03-03Date de vente dans le fichier ventes_v2.csv du TP 5
2026-03-04Date de vente dans le fichier ventes_v2.csv du TP 5
2026-03-05Date de vente dans le fichier ventes_v2.csv du TP 5

Tableaux de synthèse

Modes de sortie Structured Streaming

ModeÉcritureCas typique
AppendUniquement les nouvelles lignes apparues depuis le dernier déclenchementFlux sans besoin de réécrire tout le résultat
CompleteRéécrit toute la table de résultat à chaque micro-batchNécessaire pour les agrégations
UpdateÉcrit uniquement les lignes dont la valeur a changé depuis le dernier déclenchementIdéal pour mettre à jour les totaux en temps réel

Transformer, Estimator, Pipeline (MLlib)

TermeRôleMéthode clé
TransformerTransforme un DataFrame en un autre DataFrame.transform()
EstimatorApprend en ajustant sur les données puis fournit un Model.fit()
PipelineEnchaîne plusieurs étapes (pour éviter la divergence entre train/test)pipeline.fit(train_df)

Pièges & confusions fréquents

  1. Confondre Transformations (Lazy, enregistrées dans un DAG) et Actions (Eager, déclenchent le calcul et renvoient un résultat au Driver).
  2. Croire qu’un RDD est modifiable : un RDD est immuable, toute opération crée un nouveau RDD via transformation.
  3. Oublier que, par défaut, un RDD est recalculé à chaque action : sans cache()/persist(level), les temps explosent en enchaînant les actions.
  4. Mal utiliser Broadcast : coder en dur (par tâche) une grosse variable dans une fonction map surcharge le réseau au lieu d’être envoyée une seule fois.
  5. Penser que Catalyst optimise le contenu d’une UDF Python : Catalyst ne peut pas optimiser la boîte noire d’une UDF Python.
  6. Utiliser le mauvais outputMode en streaming : Complete réécrit toute la table (coûteux) alors que Update ne pousse que les lignes dont le total a changé.
  7. Redémarrer un streaming avec un schéma modifié et un ancien checkpoint : cela peut provoquer un Schema Mismatch, nécessitant de supprimer checkpoints/ pour repartir à zéro.

Checklist Examen

  1. Lister les 5V du Big Data et expliquer pourquoi Spark peut être jusqu’à 100 fois plus rapide que MapReduce pour les algorithmes itératifs (in-memory).
  2. Décrire le modèle Master/Slave de Spark : Driver Program, Cluster Manager (YARN/Kubernetes/Standalone) et Executors, ainsi que la décomposition Job→Stages→Tasks.
  3. Expliquer le mécanisme Lazy Evaluation : pourquoi les transformations (map/filter/groupBy) ne s’exécutent pas et comment une action (count/collect/save) déclenche réellement le calcul.
  4. Citer et distinguer cache() et persist(level), et associer un Storage Level (MEMORY_ONLY, MEMORY_AND_DISK, DISK_ONLY, MEMORY_ONLY_SER) à sa logique de stockage.
  5. Expliquer quand utiliser Broadcast et comment l’accès se fait via .value, puis expliquer l’usage d’un Accumulateur pour compter des événements globaux.
  6. Savoir utiliser DataFrames/Spark SQL avec Catalyst : création d’une vue temporaire et exécution via spark.sql, et rappeler pourquoi les fonctions natives pyspark.sql.functions sont prioritaires aux UDF Python.
  7. Comparer Spark SQL natif, Pandas UDF (Arrow, traitement par blocs) et UDF Python (boîte noire) en termes d’optimisation et de performance attendue.
  8. Maîtriser les modes Structured Streaming (Append/Complete/Update) et relier les agrégations au mode Complete et aux mises à jour au mode Update.
  9. Expliquer Fenêtrage glissant (window(…, …, …)) et Watermark (délai toléré des événements tardifs), puis rappeler le rôle des checkpoints (/offsets, /commits, /metadata, /sources).
  10. Décrire l’enchaînement MLlib en termes de Pipeline, Feature Engineering (StringIndexer→OneHotEncoder→VectorAssembler) et de la séparation Estimator/Transformer, puis identifier l’algorithme selon le type de problème (Linear/Logistic/RandomForest/KMeans).
  11. Pour le scoring temps réel, expliquer le rôle de PipelineModel (chargement) + model.transform() sur le flux (Kafka) et l’écriture vers MongoDB, puis relier la robustesse à checkpointLocation et aux offsets.
  12. Pour Deep Learning avec Spark, décrire le principe de diffusion (broadcast du modèle) + pandas_udf et la nécessité d’un Dockerfile (apache/spark:3.5.0 + tensorflow/pandas/pyarrow/numpy) validée par un import tensorflow sur spark-master.

Teste tes connaissances

Teste tes connaissances sur Introduction à Spark et Big Data avec 20 questions à choix multiples et corrections détaillées.

1. Quel est le rôle principal du Driver Program dans l’architecture distribuée de Spark ?

2. Quel composant convertit une colonne catégorielle en indices numériques appris lors du fit ?

Faire le QCM →

Révisez avec les flashcards

Mémorisez les concepts clés de Introduction à Spark et Big Data avec 20 flashcards interactives.

Big Data — caractéristiques principales ?

Volume, Vélocité, Variété, Véracité, Valeur

Écosystème Spark — composantes clés ?

Spark SQL, Spark Streaming, MLlib, GraphX

Architecture Spark — modèle ?

Master/Slave avec Driver, Cluster Manager, Executors

Voir les flashcards →

Cours similaires

Crée tes propres fiches de révision

Importe ton cours et l'IA génère fiches, QCM et flashcards en 30 secondes.

Générateur de fiches