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 :
Comprendre comment démarrer une session Spark et s’y connecter depuis Jupyter.
Créer, lire, écrire et inspecter des
DataFramesSpark.Maîtriser les opérations fondamentales : sélection, filtrage, colonnes calculées, agrégation, jointure.
Comprendre la distinction entre transformations et actions, ainsi que l’évaluation paresseuse (lazy).
Utiliser Spark SQL pour interroger des DataFrames.
Mettre en œuvre l”échantillonnage dans Spark.
Utiliser
cache()etpersist()à 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).
Installez PySpark avec pip install pyspark. Une installation Java (JDK 21
ou 25) est nécessaire. Suivez les
instructions d’installation
en cas de difficulté.
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
Combien de lignes du fichier contiennent le mot « poupée » ?
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
Créez une nouvelle colonne
note_normaliseecontenant la note ramenée entre 0 et 1 (divisez par 5). Vérifiez avecdescribe().Combien de films du genre « Animation » ont reçu au moins une note dans
ratings? (filtrezmoviessurgenres, 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
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).
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 : utilisezsplit(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) :
Chargez
ratings.csv,movies.csvettags.csvavec les bons types.Combien y a-t-il d’utilisateurs distincts, de films distincts, et de tags distincts ?
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 ?