Travaux pratiques - Introduction à PySpark et aux DataFrames Spark

Références externes utiles :

Objectif de cette séance

Cette séance introduit PySpark, l’interface Python d’Apache Spark, avec pour objectifs :

  1. Comprendre comment démarrer une session Spark et s’y connecter depuis Jupyter.

  2. Créer, lire, écrire et inspecter des DataFrames Spark.

  3. Maîtriser les opérations fondamentales : sélection, filtrage, colonnes calculées, agrégation, jointure.

  4. Comprendre la distinction entre transformations et actions, ainsi que l’évaluation paresseuse (lazy).

  5. Utiliser Spark SQL pour interroger des DataFrames.

  6. Mettre en œuvre l”échantillonnage dans Spark.

  7. Utiliser cache() et persist() à bon escient.

Nous travaillons sur des données peu volumineuses mais les mécanismes mis en œuvre sont les mêmes que sur un vrai cluster à grande échelle.

Mise en place de l’environnement

Accédez au serveur JupyterHub via Moodle (Mes enseignements > RCP216 > Accès à JupyterHub). PySpark est déjà installé. Importez le cahier Jupyter de ce TP (bouton de téléversement).

Dans Jupyter, initialisez la session Spark ainsi :

import os, sys
os.environ['PYSPARK_PYTHON']        = sys.executable
os.environ['PYSPARK_DRIVER_PYTHON'] = sys.executable

from pyspark.sql import SparkSession
spark = SparkSession.builder \
            .master("local[*]") \
            .appName("RCP216-TP3") \
            .getOrCreate()

print(spark.version)

local[*] indique à Spark d’utiliser tous les cœurs disponibles sur la machine locale. Sur un vrai cluster on remplacerait cette URL par l’adresse du cluster manager.

Premiers pas : texte et évaluation paresseuse

Préparation des données

Téléchargeons d’abord quelques fichiers de travail.

mkdir -p tpPySpark/data
wget -nc https://cedric.cnam.fr/vertigo/Cours/RCP216/docs/sophie.txt -P tpPySpark/data/
wget -nc http://cedric.cnam.fr/vertigo/Cours/RCP216/docs/geyser.csv -P tpPySpark/data/
wget -nc https://files.grouplens.org/datasets/movielens/ml-latest-small.zip -P tpPySpark/
unzip -n tpPySpark/ml-latest-small.zip -d tpPySpark/
cp tpPySpark/ml-latest-small/*.csv tpPySpark/data/

DataFrame de texte : actions et transformations

Créons un premier DataFrame à partir du fichier texte sophie.txt (première partie des « Malheurs de Sophie » de la comtesse de Ségur) :

texteSophie = spark.read.text("tpPySpark/data/sophie.txt")

# Inspecter la structure
texteSophie.printSchema()   # une seule colonne "value" de type string
texteSophie.show(5)

Spark crée un DataFrame dont chaque ligne correspond à une ligne du fichier. La lecture est paresseuse : le fichier n’est pas encore chargé en mémoire.

Actions : déclenchent l’exécution (et la lecture) :

texteSophie.count()          # nombre de lignes (déclenche la lecture + comptage)
texteSophie.first()          # première ligne
texteSophie.take(3)          # 3 premières lignes sous forme de liste Python

Attention à collect()

texteSophie.collect() ramène toutes les données sur le nœud driver. Sur un vrai cluster avec des téraoctets de données, cette opération provoquerait un débordement mémoire. À éviter donc, sauf pour de très petits DataFrames de résultats finaux.

Transformations : paresseuses, ne déclenchent pas de calcul immédiat :

from pyspark.sql.functions import length, col

# Filtrer les lignes contenant "poupée"
lignesPoupee = texteSophie.filter(texteSophie.value.contains("poupée"))

# Calculer la longueur de chaque ligne
longueursLignes = texteSophie.select(length(texteSophie.value).alias("longueur"))

À ce stade, aucun calcul n’a encore lieu. Ce n’est qu’à l’appel d’une action que Spark planifie et exécute l’ensemble du graphe de transformations :

longueursLignes.show(5)      # action : déclenche lecture + calcul des longueurs

Question

  1. Combien de lignes du fichier contiennent le mot « poupée » ?

  2. Enchaînez les transformations et l’action en une seule expression pour calculer la longueur totale du texte (somme des longueurs de toutes les lignes). Utilisez from pyspark.sql.functions import sum.

Évaluation paresseuse et optimisation

La puissance de l’évaluation paresseuse est visible lorsqu’on enchaîne plusieurs transformations : Spark les optimise globalement avant exécution.

# Ces trois transformations ne font rien immédiatement
etape1 = texteSophie.filter(length(texteSophie.value) > 20)
etape2 = etape1.select(texteSophie.value, length(texteSophie.value).alias("longueur"))
etape3 = etape2.filter(col("longueur") < 80)

# Le plan (optimisé par Catalyst) est exécuté seulement ici
etape3.show(5)

# Afficher le plan d'exécution physique optimisé
etape3.explain()

En lisant la sortie de explain() on constate que Catalyst a fusionné les deux filtres en un seul passage sur les données (filter pushdown).

Persistance : cache() et persist()

Lorsqu’un DataFrame est utilisé plusieurs fois dans des calculs différents, Spark le recalcule entièrement à chaque action (en rejouant le lineage). Pour éviter ce coût, on peut demander à Spark de conserver le DataFrame en mémoire :

# Sans cache : le fichier est relu deux fois
nb_lignes   = texteSophie.count()
nb_poupee   = texteSophie.filter(texteSophie.value.contains("poupée")).count()

# Avec cache : le fichier est lu une seule fois, le DataFrame est conservé en mémoire
texteSophie = spark.read.text("tpPySpark/data/sophie.txt").cache()
nb_lignes   = texteSophie.count()      # premier appel : lecture + mise en cache
nb_poupee   = texteSophie.filter(      # deuxième appel : lecture depuis le cache
    texteSophie.value.contains("poupée")).count()

# Libérer le cache quand on n'en a plus besoin
texteSophie.unpersist()

persist() offre plus de contrôle sur le niveau de stockage :

from pyspark import StorageLevel
texteSophie.persist(StorageLevel.MEMORY_AND_DISK)
# MEMORY_ONLY, MEMORY_AND_DISK, DISK_ONLY, MEMORY_ONLY_2 (réplication x2), ...

Question

Reprenez les calculs nb_lignes et nb_poupee en activant cache(). Observez dans l’interface Spark UI (accessible sur http://localhost:4040 si Spark tourne en local sur votre ordinateur) la différence entre le premier appel (lecture depuis disque) et le second (lecture depuis le cache). Quel onglet de Spark UI montre les DataFrames mis en cache ?

DataFrames numériques : lecture, écriture, manipulation

Lecture de fichiers CSV

# Lecture basique : colonnes nommées _c0, _c1, type string par défaut
geyser = spark.read.csv("tpPySpark/data/geyser.csv")
geyser.printSchema()
geyser.show(5)
# Lecture avec en-tête et inférence de types
geyser = (spark.read
          .option("header", True)
          .option("inferSchema", True)
          .csv("tpPySpark/data/geyser.csv"))
geyser.printSchema()    # colonnes duration et interval de type double
geyser.show(5)
geyser.describe().show()

Note sur inferSchema

L’inférence de types nécessite une seconde passe sur les données pour déterminer le type de chaque colonne. Sur un très grand fichier il est préférable de spécifier le schéma explicitement :

from pyspark.sql.types import StructType, StructField, DoubleType

schema = StructType([
    StructField("duration", DoubleType(), True),
    StructField("interval", DoubleType(), True),
])
geyser = spark.read.schema(schema).csv("tpPySpark/data/geyser.csv",
                                       header=True)
geyser.printSchema()

Lecture des fichiers CSV MovieLens

Chargeons dans Spark les données MovieLens utilisées dans la séance précédente :

ratings = (spark.read
           .option("header", True)
           .option("inferSchema", True)
           .csv("tpPySpark/data/ratings.csv"))

movies = (spark.read
          .option("header", True)
          .option("inferSchema", True)
          .csv("tpPySpark/data/movies.csv"))

ratings.printSchema()
ratings.show(5)
print(f"Nombre de notes : {ratings.count()}")
print(f"Nombre de films : {movies.count()}")

Écriture de DataFrames

# Écriture en CSV (produit un répertoire contenant un fichier par partition)
geyser.write.mode("overwrite").csv("tpPySpark/data/geyser_out", header=True)

# Pour obtenir un seul fichier CSV : forcer une seule partition
geyser.repartition(1).write.mode("overwrite").csv("tpPySpark/data/geyser_1fichier",
                                                   header=True)

# Écriture en Parquet (format binaire compressé, recommandé pour la production)
ratings.write.mode("overwrite").parquet("tpPySpark/data/ratings_parquet")

# Relecture depuis Parquet
ratings_parquet = spark.read.parquet("tpPySpark/data/ratings_parquet")
ratings_parquet.printSchema()

Note sur les formats

En production, Parquet est préféré à CSV : il est binaire (plus compact), supporte la compression, permet le predicate pushdown (Spark ne lit que les colonnes et lignes nécessaires) et préserve les types. Delta Lake et Iceberg sont des évolutions de Parquet qui ajoutent, entre autres, le support des transactions (ACID : atomicity, consistency, isolation, and durability).

Sélection, filtrage et colonnes calculées

Sélection de colonnes

from pyspark.sql.functions import col

# Plusieurs syntaxes équivalentes pour sélectionner des colonnes
ratings.select("userId", "rating").show(5)
ratings.select(col("userId"), col("rating")).show(5)
ratings.select(ratings.userId, ratings.rating).show(5)

# Sélection dynamique à partir d'une liste
colonnes = ["userId", "movieId", "rating"]
ratings.select(*colonnes).show(5)

Filtrage

# Notes >= 4.5
bonnes_notes = ratings.filter(col("rating") >= 4.5)
bonnes_notes.count()
# Conditions multiples (& pour ET, | pour OU, ~ pour NON)
bons_recents = ratings.filter(
    (col("rating") >= 4.0) & (col("timestamp") > 1_000_000_000)
)
bons_recents.show(5)
# Filtrage sur une chaîne de caractères
dramas = movies.filter(col("genres").contains("Drama"))
dramas.count()

Colonnes calculées avec withColumn

from pyspark.sql.functions import from_unixtime, year, round as spark_round

# Convertir le timestamp Unix en date lisible, extraire l'année
ratings = (ratings
           .withColumn("date", from_unixtime(col("timestamp")))
           .withColumn("annee", year(from_unixtime(col("timestamp")))))

ratings.select("userId", "movieId", "rating", "date", "annee").show(5)

# Arrondir la note à l'entier le plus proche
ratings = ratings.withColumn("rating_arrondi", spark_round(col("rating")))

# Renommer et supprimer une colonne
ratings = ratings.withColumnRenamed("rating", "note")
ratings = ratings.drop("timestamp")
ratings.printSchema()

Question

  1. Créez une nouvelle colonne note_normalisee contenant la note ramenée entre 0 et 1 (divisez par 5). Vérifiez avec describe().

  2. Combien de films du genre « Animation » ont reçu au moins une note dans ratings ? (filtrez movies sur genres, ensuite faites une jointure avec .join())

Agrégation et groupBy

Agrégations de base

from pyspark.sql.functions import avg, count, min as spark_min, max as spark_max, stddev

# Note moyenne globale
ratings.select(avg("note")).show()

# Plusieurs agrégations en une passe
ratings.agg(
    avg("note").alias("note_moy"),
    count("*").alias("nb_notes"),
    spark_min("note").alias("note_min"),
    spark_max("note").alias("note_max"),
    stddev("note").alias("note_std"),
).show()

Agrégation par groupe

# Statistiques par film
stats_par_film = (ratings
                  .groupBy("movieId")
                  .agg(
                      avg("note").alias("note_moy"),
                      count("*").alias("nb_notes"),
                      stddev("note").alias("note_std"),
                  ))
stats_par_film.show(10)
stats_par_film.orderBy(col("nb_notes").desc()).show(10)
# Statistiques par année
stats_par_annee = (ratings
                   .groupBy("annee")
                   .agg(
                       avg("note").alias("note_moy"),
                       count("*").alias("nb_notes"),
                   )
                   .orderBy("annee"))
stats_par_annee.show()

Question

  1. Calculez la note moyenne par utilisateur et identifiez les 5 utilisateurs les plus sévères (note moyenne la plus basse, parmi ceux ayant noté au moins 20 films).

  2. Calculez, par genre principal (premier genre de la colonne genres), la note moyenne et le nombre de films distincts. Triez par note moyenne décroissante. Indice : utilisez split(col("genres"), "\\|").getItem(0) pour extraire le premier genre.

Jointures

# Jointure gauche : enrichir les notes avec les titres et genres
notes_enrichies = ratings.join(movies, "movieId", "left")
notes_enrichies.show(5)

# Vérification : aucune ligne perdue ?
print(f"Avant jointure : {ratings.count()}")
print(f"Après jointure : {notes_enrichies.count()}")
# Top 10 films les plus populaires avec leur note moyenne
top_films = (notes_enrichies
             .groupBy("movieId", "title")
             .agg(
                 avg("note").alias("note_moy"),
                 count("*").alias("nb_notes"),
             )
             .filter(col("nb_notes") >= 50)
             .orderBy(col("note_moy").desc()))
top_films.show(10, truncate=False)

Note sur les jointures distribuées

Dans Spark, deux stratégies principales existent pour les jointures :

  • Sort-merge join (par défaut pour les grands DataFrames) : les deux DataFrames sont triés sur la clé de jointure puis fusionnés. Exige un shuffle (transfert de données entre nœuds) coûteux.

  • Broadcast join : si un des DataFrames est petit, Spark peut l’envoyer (broadcast) en entier à chaque nœud worker, évitant le shuffle coûteux. Spark le fait automatiquement en-dessous d’un seuil configurable (spark.sql.autoBroadcastJoinThreshold, 10 Mo par défaut).

from pyspark.sql.functions import broadcast

# Forcer un broadcast join sur le petit DataFrame movies
notes_enrichies = ratings.join(broadcast(movies), "movieId", "left")
notes_enrichies.explain()   # vérifier que "BroadcastHashJoin" apparaît

Spark SQL

Spark permet d’interroger des DataFrames via des requêtes SQL standard, il suffit pour cela d’enregistrer le DataFrame comme une vue temporaire :

# Enregistrement des DataFrames comme vues SQL temporaires
ratings.createOrReplaceTempView("notes")
movies.createOrReplaceTempView("films")

# Requête SQL : top 10 films les mieux notés (≥ 50 notes)
top_sql = spark.sql("""
    SELECT f.title,
           AVG(n.note)  AS note_moy,
           COUNT(*)     AS nb_notes
    FROM   notes  n
    JOIN   films  f ON n.movieId = f.movieId
    GROUP  BY f.movieId, f.title
    HAVING COUNT(*) >= 50
    ORDER  BY note_moy DESC
    LIMIT  10
""")
top_sql.show(truncate=False)
# Mélanger SQL et API DataFrame : le résultat de spark.sql() est un DataFrame
films_action = spark.sql(
    "SELECT movieId, title FROM films WHERE genres LIKE '%Action%'"
)
films_action.count()

# Continuer avec l'API DataFrame
films_action.join(ratings, "movieId").groupBy("title") \
            .agg(avg("note").alias("note_moy")).orderBy(col("note_moy").desc()) \
            .show(10, truncate=False)

Question

Écrivez en Spark SQL la requête suivante : pour chaque utilisateur ayant noté plus de 50 films, calculez la moyenne de ses notes. Retournez les 5 utilisateurs avec les notes moyennes les plus élevées. Comparez avec la version API DataFrame équivalente.

Échantillonnage dans Spark

Comme vu en cours, l’échantillonnage est une première approche de réduction du volume de données. Spark propose deux méthodes directement sur les DataFrames.

Échantillonnage simple

# Échantillon d'environ 10 % des notes, sans remise, reproductible (seed=42)
echantillon = ratings.sample(withReplacement=False, fraction=0.1, seed=42)
print(f"Taille originale  : {ratings.count()}")
print(f"Taille échantillon: {echantillon.count()}")
print(f"Taux réel         : {echantillon.count() / ratings.count():.3f}")

# Vérification : la distribution des notes est-elle préservée ?
ratings.groupBy("note").count().orderBy("note").show()
echantillon.groupBy("note").count().orderBy("note").show()

Question

Comparez la note moyenne sur l’ensemble des données et sur l’échantillon à 10 %. Répétez avec des fractions de 1 %, 5 %, 10 % et 50 %. À partir de quelle fraction l’estimation de la note moyenne est-elle stable (écart < 0.01) ?

Échantillonnage stratifié

L’échantillonnage stratifié garantit une représentation contrôlée de chaque strate. Ici nous l’appliquons sur les valeurs de note (strate = valeur de note) pour garantir que toutes les valeurs de note sont représentées dans l’échantillon dans des proportions choisies :

# Fractions par valeur de note : même taux pour toutes les strates
valeurs_notes = [row["note"] for row in
                 ratings.select("note").distinct().collect()]
fractions = {note: 0.1 for note in valeurs_notes}

echantillon_strat = ratings.stat.sampleBy("note", fractions=fractions, seed=42)
print(f"Taille : {echantillon_strat.count()}")

# Comparaison des distributions
print("Distribution originale :")
ratings.groupBy("note").count() \
       .withColumn("pct", col("count") / ratings.count() * 100) \
       .orderBy("note").show()

print("Distribution échantillon stratifié :")
echantillon_strat.groupBy("note").count() \
                 .withColumn("pct", col("count") / echantillon_strat.count() * 100) \
                 .orderBy("note").show()

Question

Créez un échantillon stratifié sur la valeur de note en sur-représentant les notes extrêmes (0.5 et 5.0 à 50 %) et en sous-représentant les notes médianes (2.5 et 3.0 à 5 %). Pour quelle application de fouille de données ce type d’échantillonnage serait-il utile ?

Exercice de synthèse optionnel sur MovieLens

Question 1 : Exploration

En utilisant l’API DataFrame (pas SQL) :

  1. Chargez ratings.csv, movies.csv et tags.csv avec les bons types.

  2. Combien y a-t-il d’utilisateurs distincts, de films distincts, et de tags distincts ?

  3. Quelle est la plage temporelle des notes (dates min et max) ?

Question 2 : Passage à l’échelle

Comparez les performances de Spark et Pandas pour l’agrégation groupBy("movieId") avec calcul de la note moyenne, sur le jeu de données MovieLens complet. Mesurez avec %%time. Quel outil est le plus rapide ici, et pourquoi le résultat peut-il surprendre par rapport à la comparaison avec Polars la séance précédente ?