Travaux pratiques - Recommandation avec Spark MLlib et GraphFrames¶
Références externes utiles :
Objectif du TP¶
Ce TP met en œuvre un système de recommandation de films complet sur le jeu de données MovieLens avec Spark MLlib. Le travail est décomposé en quatre parties :
Exploration des données et structure de la matrice d’utilités.
Factorisation matricielle avec ALS : entraînement, évaluation, recommandations.
Réglage des hyperparamètres : impact du rang et de la régularisation.
Recommandation par graphe : PPR et filtrage collaboratif via GraphFrames.
Les données MovieLens ont été vues dans le TP de préparation Python. Nous tentons d’utiliser ici la version complète MovieLens 25M (25 millions de notes, 162 000 utilisateurs, 62 000 films) pour que la parallélisation de Spark soit significative.
Partie 1 : Données et exploration¶
Vérification des bibliothèques¶
La communication entre Spark et Polars exige la version 2 de sparkpl et la version 4.x de Pyspark, une vérification est nécessaire :
import sys
!{sys.executable} -m pip install "pyspark>=4,<5" sparkpl==2.0.1 polars altair graphframes-py==0.12.1
Si la réponse vous indique que pyspark était déjà installé alors vous pouvez poursuivre. En revanche, si cette instruction vient d’installer la version 4 à la place d’une plus ancienne alors il faut redémarrer le noyau (icône circulaire du menu en haut de la page) avant de poursuivre !
Récupération des données MovieLens 25M¶
import os, sys
os.environ['PYSPARK_PYTHON'] = sys.executable
os.environ['PYSPARK_DRIVER_PYTHON'] = sys.executable
dossier = "tpRecommandation/movielens/"
os.makedirs(dossier, exist_ok=True)
archive = dossier + "ml-25m.zip"
if not os.path.exists(dossier + "ml-25m/ratings.csv"):
if not os.path.exists(archive):
print("Téléchargement MovieLens 25M (~250 Mo)...")
os.system(f"wget -q https://files.grouplens.org/datasets/movielens/ml-25m.zip -P {dossier}")
print("Décompression...")
os.system(f"unzip -q {archive} -d {dossier}")
print("OK")
else:
print("Données déjà présentes.")
Initialisation de Spark¶
from importlib.metadata import version
from pyspark.sql import SparkSession
from pyspark.sql import functions as F
gf_version = version("graphframes-py") # garde le JAR aligné sur le paquet pip
spark = (SparkSession.builder
.appName("tp-recommandation")
.config("spark.driver.memory", "4g")
.config("spark.jars.packages",
f"io.graphframes:graphframes-spark4_2.13:{gf_version}")
.getOrCreate())
Chargement des données¶
ratings = spark.read.csv(
dossier + "ml-25m/ratings.csv",
header=True, inferSchema=True)
movies = spark.read.csv(
dossier + "ml-25m/movies.csv",
header=True, inferSchema=True)
print("Schéma ratings :"), ratings.printSchema()
print("Schéma movies :"), movies.printSchema()
n_ratings = ratings.count()
n_users = ratings.select("userId").distinct().count()
n_movies = ratings.select("movieId").distinct().count()
print(f"Notes : {n_ratings:,}")
print(f"Utilisateurs: {n_users:,}")
print(f"Films : {n_movies:,}")
print(f"Taux de remplissage : {n_ratings / (n_users * n_movies):.6f}")
Question
Calculez le nombre de notes par film et ensuite par utilisateur. Calculez la médiane, la moyenne et la valeur maximale du nombre de notes par utilisateur. Comment caractériser la distribution ?
Distribution des notes¶
import polars as pl
from sparkpl.converter import spark_to_polars
import altair as alt
alt.data_transformers.enable("default", max_rows=None)
# Distribution des notes (agrégation Spark, visualisation Altair)
distrib = (ratings
.groupBy("rating")
.agg(F.count("*").alias("nb"))
.orderBy("rating"))
distrib_pl = spark_to_polars(distrib)
alt.Chart(distrib_pl).mark_bar().encode(
x=alt.X("rating:O", title="Note"),
y=alt.Y("nb:Q", title="Nombre de notes"),
tooltip=["rating:O", alt.Tooltip("nb:Q", format=",")]
).properties(title="Distribution des notes pour MovieLens 25M", width=400, height=280)
Question
Observez-vous un biais vers les notes élevées (positivity bias) ? Comment cela peut affecter un système de recommandation ?
Films les plus notés et les mieux notés¶
stats_films = (ratings
.groupBy("movieId")
.agg(F.count("rating").alias("nb_notes"),
F.mean("rating").alias("note_moy"))
.join(movies.select("movieId", "title", "genres"), on="movieId")
)
# Top 10 films par nombre de notes
print("Films les plus notés :")
stats_films.orderBy(F.desc("nb_notes")).show(10, truncate=False)
# Top 10 films les mieux notés (parmi ceux ayant au moins 1000 notes)
print("Films les mieux notés (≥ 1000 notes) :")
(stats_films
.filter(F.col("nb_notes") >= 1000)
.orderBy(F.desc("note_moy"))
.show(10, truncate=False))
Question
Pourquoi c’est important de filtrer les films ayant peu de notes avant de calculer un classement par note moyenne ? Que se passerait-il sans ce filtre ?
Partie 2 : Factorisation matricielle avec ALS¶
Rappel du principe¶
ALS cherche des vecteurs latents \(\mathbf{u}_i\) (utilisateurs) et \(\mathbf{a}_j\) (films) de dimension \(m\) minimisant :
La prédiction de note pour un couple [utilisateur, film] non observé est \(\hat{x}_{kl} = \mathbf{u}_k^T \cdot \mathbf{a}_l\).
Séparation en apprentissage et test¶
# Partitionnement 80/20
train, test = ratings.randomSplit([0.8, 0.2], seed=42)
train.cache()
print(f"Train : {train.count():,} | Test : {test.count():,}")
Entraînement du modèle¶
Pour un premier essai nous choisissons a priori les valeurs des paramètres :
from pyspark.ml.recommendation import ALS
from pyspark.ml.evaluation import RegressionEvaluator
als = ALS(
rank=20, # dimension de l'espace latent
maxIter=10, # nombre d'itérations ALS
regParam=0.1, # paramètre de régularisation lambda
userCol="userId",
itemCol="movieId",
ratingCol="rating",
coldStartStrategy="drop" # ignorer les utilisateurs et films hors du train
)
model = als.fit(train)
print("Entraînement terminé.")
Note
L’option coldStartStrategy="drop" supprime du jeu de test les lignes
correspondant à des utilisateurs et à des films non vus à l’entraînement
(problème de démarrage à froid). L’alternative "nan" conserve ces lignes
avec une prédiction NaN, ce qui est utile pour les analyser ultérieurement.
Évaluation sur le jeu de test¶
Evaluation d’un modèle naïf (baseline) qui prédit toujours la note moyenne globale :
evaluator_rmse = RegressionEvaluator(
metricName="rmse", labelCol="rating", predictionCol="prediction")
evaluator_mae = RegressionEvaluator(
metricName="mae", labelCol="rating", predictionCol="prediction")
mu = ratings.select(F.mean("rating")).collect()[0][0]
naive_preds = test.withColumn("prediction", F.lit(mu))
rmse_naive = evaluator_rmse.evaluate(naive_preds)
mae_naive = evaluator_mae.evaluate(naive_preds)
print(f"RMSE naïf : {rmse_naive:.4f}")
print(f"MAE naïf : {mae_naive:.4f}")
Evaluation du modèle ALS obtenu :
predictions = model.transform(test)
rmse = evaluator_rmse.evaluate(predictions)
mae = evaluator_mae.evaluate(predictions)
print(f"RMSE : {rmse:.4f}")
print(f"MAE : {mae:.4f}")
Question
L’erreur RMSE (Root Mean Square Error) obtenue vous semble satisfaisante ? Quelle amélioration apporte ALS par rapport au prédicteur naïf ?
Inspection des facteurs latents¶
# Facteurs des films : chaque film est un vecteur de dimension rank
item_factors = model.itemFactors
item_factors.show(5, truncate=False)
# Facteurs des utilisateurs
user_factors = model.userFactors
print(f"Utilisateurs : {user_factors.count()}, Films : {item_factors.count()}")
Exercice
Calculez la norme \(L_2\) des vecteurs de films et triez par norme décroissante. A quoi correspondent les films de forte norme (films très populaires, très polarisants, etc.) ?
Génération de recommandations¶
Obtenons d’abord les 10 meilleures recommandations pour chaque utilisateur :
user_recs = model.recommendForAllUsers(numItems=10)
user_recs.show(5, truncate=False)
Les films déjà notés par un utilisateur spécifique, identifié par son user_id :
user_id = 1
# Films déjà notés par cet utilisateur
deja_notes = (ratings
.filter(F.col("userId") == user_id)
.join(movies, on="movieId")
.orderBy(F.desc("rating"))
.select("title", "rating", "genres"))
print(f"Films notés par l'utilisateur {user_id} :")
deja_notes.show(10, truncate=False)
Les meilleures recommandations obtenues avec ALS pour cet utilisateur :
recs_user = (user_recs
.filter(F.col("userId") == user_id)
.select(F.explode("recommendations").alias("rec"))
.select(
F.col("rec.movieId").alias("movieId"),
F.col("rec.rating").alias("score_predit"))
.join(movies, on="movieId")
.select("title", "genres", "score_predit")
.orderBy(F.desc("score_predit")))
print(f"Top 10 recommandations pour l'utilisateur {user_id} :")
recs_user.show(10, truncate=False)
Question
Les recommandations sont-elles cohérentes avec les films déjà notés par cet utilisateur ? Trouvez-vous une cohérence de genre ou de style ?
Trouvons maintenant les 10 utilisateurs les plus susceptibles d’apprécier un film spécifique, identifié par son titre :
film_titre = "Toy Story (1995)"
film_id = movies.filter(F.col("title") == film_titre).first()["movieId"]
film_recs = model.recommendForAllItems(numUsers=10)
(film_recs
.filter(F.col("movieId") == film_id)
.select(F.explode("recommendations").alias("rec"))
.select(F.col("rec.userId").alias("userId"),
F.col("rec.rating").alias("score_predit"))
.orderBy(F.desc("score_predit"))
.show(10))
Partie 3 : Réglage des hyperparamètres¶
Nous explorons ici l’espace des hyperparamètres afin de voir si nos choix de départ étaient raisonnables. Cette partie peut prendre plus de temps, en attendant que les calculs soient faits regardez les questions qui suivent.
Effet du rang¶
Le rang \(m\) est le nombre de facteurs latents. Un rang trop faible sous-ajuste alors qu’un rang trop élevé sur-ajuste (et ralentit les calculs).
resultats = []
for rank in [5, 10, 20, 50]:
m = ALS(rank=rank, maxIter=5, regParam=0.1,
userCol="userId", itemCol="movieId", ratingCol="rating",
coldStartStrategy="drop")
mod = m.fit(train)
rmse = evaluator_rmse.evaluate(mod.transform(test))
resultats.append({"rank": rank, "RMSE": rmse})
print(f"rank={rank:3d} RMSE={rmse:.4f}")
res_pl = pl.DataFrame(resultats)
alt.Chart(res_pl).mark_line(point=True).encode(
x=alt.X("rank:O", title="Rang (dimension latente)"),
y=alt.Y("RMSE:Q", scale=alt.Scale(zero=False), title="RMSE test"),
tooltip=["rank:O", alt.Tooltip("RMSE:Q", format=".4f")]
).properties(title="RMSE suivant le rang ALS", width=400, height=280)
Question
La RMSE diminue systématiquement avec l’augmentation du rang ? À partir de quel rang on observe une stagnation ou une dégradation (sur-apprentissage) ?
Effet de la régularisation¶
resultats_reg = []
for reg in [0.01, 0.05, 0.1, 0.2]:
m = ALS(rank=10, maxIter=5, regParam=reg,
userCol="userId", itemCol="movieId", ratingCol="rating",
coldStartStrategy="drop")
mod = m.fit(train)
rmse = evaluator_rmse.evaluate(mod.transform(test))
resultats_reg.append({"regParam": reg, "RMSE": rmse})
print(f"regParam={reg:.2f} RMSE={rmse:.4f}")
Question
Quel est l’effet d’une régularisation trop faible ? Trop forte ? Quelle
valeur de regParam donne le meilleur RMSE ?
Recherche en grille avec validation croisée¶
Une recherche en grille permet de couvrir de façon plus systématique l’espace de variation des valeurs des hyper-paramètres. Attention, l’exécution de cette cellule peut prendre jusqu’à 30 minutes.
from pyspark.ml.tuning import CrossValidator, ParamGridBuilder
als_cv = ALS(userCol="userId", itemCol="movieId", ratingCol="rating",
coldStartStrategy="drop")
param_grid = (ParamGridBuilder()
.addGrid(als_cv.rank, [20, 50])
.addGrid(als_cv.regParam, [0.05, 0.1])
.build())
cv = CrossValidator(
estimator=als_cv,
estimatorParamMaps=param_grid,
evaluator=evaluator_rmse,
numFolds=3,
parallelism=2) # évaluer 2 configurations en parallèle
cv_model = cv.fit(train)
print("Meilleure configuration :")
print(f" rank = {cv_model.bestModel.rank}")
print(f" regParam = {cv_model.bestModel._java_obj.parent().getRegParam()}")
print(f" RMSE = {evaluator_rmse.evaluate(cv_model.transform(test)):.4f}")
Question
La validation croisée donne-t-elle une configuration différente de celle trouvée manuellement ? Comment expliquer d’éventuelles différences ?
Partie 4 : Recommandation par graphe avec GraphFrames¶
Le filtrage collaboratif avec factorisation matricielle par ALS exploite la matrice d’utilités mais ignore la structure de réseau entre utilisateurs et articles. Nous illustrons maintenant la recommandation par Personalized PageRank (PPR) sur le graphe biparti qui représente les liens entre utilisateurs et films, en travaillant sur un sous-ensemble MovieLens 100K (pour une exécution plus rapide).
# Sous-ensemble : MovieLens 100K, déjà employé dans tpPandasPolarsPlotly
dossier100k = "tpPandasPolarsPlotly/data/"
ratings_100k = spark.read.csv(
dossier100k + "ratings.csv",
header=True, inferSchema=True)
movies_100k = spark.read.csv(
dossier100k + "movies.csv",
header=True, inferSchema=True)
print(f"Ratings 100K : {ratings_100k.count():,}")
Construction du graphe biparti¶
Dans le graphe biparti, les nœuds sont soit des utilisateurs soit des films.
Nous préfixons les identifiants pour éviter les collisions. Aussi, pour que PageRank fonctionne il est nécessaire de rendre les liens symétriques, sinon la marche aléatoire s’arrête après le premier lien utilisateur → film. Pour cela on crée une copie de l’ensemble des arêtes, en renommant src en dest et dest en src, et ensuite on fusionne les deux ensembles. Il est critique d’utiliser unionByName et non union car unionByName fusionne en tenant compte du nom des colonnes alors que union se contente de leur position (première colonne avec première colonne, deuxième avec deuxième, etc.).
from graphframes import GraphFrame
from pyspark.sql.types import StringType
# Nœuds utilisateurs
verts_users = (ratings_100k
.select("userId")
.distinct()
.withColumn("id", F.concat(F.lit("u_"), F.col("userId").cast(StringType())))
.withColumn("type", F.lit("user"))
.select("id", "type"))
# Nœuds films
verts_items = (movies_100k
.select("movieId", "title")
.withColumn("id", F.concat(F.lit("m_"), F.col("movieId").cast(StringType())))
.withColumn("type", F.lit("movie"))
.select("id", F.col("title").alias("type")))
vertices = verts_users.union(
verts_items.withColumn("type", F.lit("movie"))
.select("id", "type"))
# Arêtes utilisateur → film avec poids = note
edges = (ratings_100k
.withColumn("src", F.concat(F.lit("u_"), F.col("userId").cast(StringType())))
.withColumn("dst", F.concat(F.lit("m_"), F.col("movieId").cast(StringType())))
.withColumn("weight", F.col("rating"))
.select("src", "dst", "weight"))
# Arêtes symétriques (graphe non orienté)
edges_rev = edges.select(F.col("dst").alias("src"),
F.col("src").alias("dst"),
"weight")
edges_bipartite = edges.unionByName(edges_rev)
bipartite_g = GraphFrame(vertices, edges_bipartite)
print(f"Nœuds : {bipartite_g.vertices.count()}, "
f"Arêtes : {bipartite_g.edges.count()}")
spark.sparkContext.setCheckpointDir("checkpoints")
PageRank personnalisé pour un utilisateur¶
Obtenons les recommandations PPR pour le même utilisateur (que nous avions sélectionné plus haut). D’abord, quels sont les films qu’il a déjà vus, pour pouvoir les exclure de la sélection :
target_user_id = 1
source_id = f"u_{target_user_id}"
# Films déjà notés par cet utilisateur (à exclure des recommandations)
deja_vus = set(
ratings_100k
.filter(F.col("userId") == target_user_id)
.select(F.concat(F.lit("m_"), F.col("movieId").cast(StringType())).alias("mid"))
.rdd.flatMap(lambda x: x)
.collect())
print(f"Films déjà notés par l'utilisateur {target_user_id} : {len(deja_vus)}")
Ensuite les recommandations, obtenues par PPR et excluant les films déjà vus.
# PPR depuis l'utilisateur cible
ppr = bipartite_g.pageRank(
resetProbability=0.15,
maxIter=10,
sourceId=source_id)
# Recommandations : films à fort score PPR non encore vus
recs_ppr = (ppr.vertices
.filter(F.col("id").startswith("m_")) # garder uniquement les films
.filter(~F.col("id").isin(deja_vus)) # exclure les films déjà vus
.withColumn("movieId",
F.col("id").substr(3, 10).cast("int"))
.join(movies_100k.select("movieId", "title", "genres"), on="movieId")
.select("title", "genres", "pagerank")
.orderBy(F.desc("pagerank"))
.limit(10))
print(f"Top 10 recommandations PPR pour l'utilisateur {target_user_id} :")
recs_ppr.show(truncate=False)
Question
Comparez les recommandations PPR avec celles d’ALS pour le même utilisateur (sachant toutefois que ALS a été entraîné sur un ensemble bien plus grand). Quelles différences observez-vous ?
Question
Le PPR exploite les chemins de longueur supérieure à 2 dans le graphe biparti. Qu’est-ce que cela signifie concrètement ? Quel type d’information est ignoré par le filtrage collaboratif user-based au premier degré mais exploité par PPR ?
Recommandation basée sur les communautés¶
# Détection de communautés d'utilisateurs sur le graphe biparti
lp = bipartite_g.labelPropagation(maxIter=5)
n_communities_lp = lp.select("label").distinct().count()
print(f"Nombre de communautés (Label Propagation) : {n_communities_lp}")
# A quelle communauté appartient l'utilisateur cible ?
comm_user = lp.filter(F.col("id") == source_id).first()["label"]
print(f"Communauté de l'utilisateur {target_user_id} : {comm_user}")
# Films les plus populaires dans cette communauté, non encore vus par la cible
users_same_comm = (lp
.filter(F.col("label") == comm_user)
.filter(F.col("id").startswith("u_"))
.withColumn("userId",
F.col("id").substr(3, 10).cast("int"))
.select("userId"))
recs_comm = (ratings_100k
.join(users_same_comm, on="userId")
.groupBy("movieId")
.agg(F.mean("rating").alias("note_moy_comm"),
F.count("*").alias("nb_notes_comm"))
.filter(F.col("nb_notes_comm") >= 5)
.filter(~F.concat(F.lit("m_"),
F.col("movieId").cast(StringType())).isin(deja_vus))
.join(movies_100k.select("movieId", "title", "genres"), on="movieId")
.orderBy(F.desc("note_moy_comm"))
.select("title", "genres", "note_moy_comm", "nb_notes_comm")
.limit(10))
print("Top 10 films populaires dans la communauté :")
recs_comm.show(truncate=False)
Question
La recommandation par communautés vous semble pertinente ? Calculez et affichez les tailles des communautés. Que contatez-vous ?
Synthèse¶
Dans ce TP nous avons illustré trois approches complémentaires de recommandation :
1. Le filtrage collaboratif par factorisation matricielle avec ALS (Alternating Least Squares) est le plus performant en termes de RMSE sur les notes explicites. Il est scalable (Spark parallélise les étapes ALS) mais ne tient pas compte de la structure de réseau et souffre du problème de démarrage à froid.
2. La recommandation par graphe avec PPR (Personalized PageRank) exploite les chemins de longueur arbitraire dans le graphe biparti. Il est interprétable (on peut tracer le chemin menant à une recommandation) et gère naturellement les utilisateurs sans beaucoup d’historique (quelques interactions suffisent à amorcer la diffusion du score). Il est plus lent qu’ALS sur de très grands graphes.
3. Recommandation par communauté : la plus simple et la plus interprétable, mais aussi la plus grossière car tous les utilisateurs d’une même communauté reçoivent les mêmes recommandations. Elle est utile comme baseline ou en complément des deux approches précédentes.
En production, les systèmes contemporains sont hybrides : ALS (ou LightGCN) pour la précision, PPR ou similarité de contenu pour la gestion du démarrage à froid et filtrage par communauté pour la diversité.
Références¶
[HK15] Harper, F. M. and Konstan, J. A. The MovieLens Datasets: History and Context. ACM Transactions on Interactive Intelligent Systems, 5(4), 2015. https://grouplens.org/datasets/movielens/
[KBV09] Koren, Y., Bell, R., Volinsky, C. Matrix factorization techniques for recommender systems. Computer, 42(8):30–37, 2009.
[HDW20] He, X., Deng, K., Wang, X., Li, Y., Zhang, Y., Wang, M. LightGCN: Simplifying and Powering Graph Convolution Network for Recommendation. SIGIR 2020.
Documentation Spark MLlib ALS : https://spark.apache.org/docs/latest/ml-collaborative-filtering.html