Module 10 — Projet : un modèle entraîné sur un gros jeu de données
Neuf modules, neuf briques. Ce dernier module rassemble tout : lire les
vols en Parquet, construire un pipeline complet, valider par
CrossValidator, sauvegarder le modèle, écrire les prédictions. Puis on
compare avec un traitement identique sur un échantillon de 5 % en
pandas, pour savoir ce que le volume a réellement apporté.
Cadre du projet
Cible : prédire si un vol arrivera avec plus de 15 minutes de
retard (retard_15min = retard_arrivee > 15, binaire).
Données : 43 millions de vols américains sur 2015 à 2024, Parquet partitionné par année et mois, 6 Go compressés sur disque.
Métrique : AUC (areaUnderROC), avec un œil sur le rappel à
précision de 0,7 — c'est la métrique métier négociée avec la compagnie
qui exploitera les alertes.
Séparation : train jusqu'à décembre 2022, test sur 2023 et 2024. Découpe temporelle, jamais aléatoire sur des données horodatées (module 6).
Lecture et préparation
from pyspark.sql import SparkSession, functions as F
from pyspark.ml import Pipeline
from pyspark.ml.feature import StringIndexer, OneHotEncoder, VectorAssembler
from pyspark.ml.classification import RandomForestClassifier
from pyspark.ml.evaluation import BinaryClassificationEvaluator
from pyspark.ml.tuning import CrossValidator, ParamGridBuilder
spark = (
SparkSession.builder
.appName("vols-retards-projet")
.master("local[8]")
.config("spark.sql.shuffle.partitions", "128")
.config("spark.sql.adaptive.enabled", "true")
.getOrCreate()
)
vols = (
spark.read.parquet("data/vols_parquet")
.filter(F.col("retard_arrivee").isNotNull())
.filter(F.col("distance_km").between(50, 8000))
.withColumn("retard_15min", (F.col("retard_arrivee") > 15).cast("integer"))
.withColumn("heure_depart", F.hour("heure_depart_planifiee"))
.withColumn("jour_semaine", F.dayofweek("date_depart"))
)
train = vols.filter(F.col("annee") <= 2022)
test = vols.filter(F.col("annee") >= 2023)
print(f"train : {train.count():,} lignes")
print(f"test : {test.count():,} lignes")
Pipeline et validation croisée
cat = ["compagnie", "origine", "destination"]
indexeurs = [
StringIndexer(inputCol=c, outputCol=f"{c}_idx", handleInvalid="keep")
for c in cat
]
encodeur = OneHotEncoder(
inputCols=[f"{c}_idx" for c in cat],
outputCols=[f"{c}_ohe" for c in cat],
)
assembleur = VectorAssembler(
inputCols=["distance_km", "heure_depart", "jour_semaine"]
+ [f"{c}_ohe" for c in cat],
outputCol="features",
handleInvalid="skip",
)
foret = RandomForestClassifier(
featuresCol="features", labelCol="retard_15min",
numTrees=200, maxDepth=10, maxBins=400, subsamplingRate=0.7, seed=42,
)
pipeline = Pipeline(stages=indexeurs + [encodeur, assembleur, foret])
grille = (
ParamGridBuilder()
.addGrid(foret.numTrees, [100, 200])
.addGrid(foret.maxDepth, [8, 10, 12])
.build()
)
evaluateur = BinaryClassificationEvaluator(
labelCol="retard_15min", rawPredictionCol="rawPrediction",
metricName="areaUnderROC",
)
cv = CrossValidator(
estimator=pipeline, estimatorParamMaps=grille,
evaluator=evaluateur, numFolds=3, parallelism=2, seed=42,
)
modele_cv = cv.fit(train)
predictions = modele_cv.transform(test)
auc = evaluateur.evaluate(predictions)
print(f"AUC test : {auc:.4f}") # 0,761 sur nos données
Sauvegarde et écriture des scores
modele_cv.bestModel.write().overwrite().save("models/vols/retard-v1")
(
predictions
.select("annee", "mois", "compagnie", "origine", "destination",
F.col("probability").getItem(1).alias("proba_retard"))
.write.mode("overwrite")
.partitionBy("annee", "mois")
.parquet("data/predictions/vols/")
)
Comparaison avec un échantillon en pandas
Prenons 5 % du jeu, soit environ 2 millions de lignes, et refaisons la
même chaîne avec scikit-learn et xgboost :
import pandas as pd
from sklearn.compose import ColumnTransformer
from sklearn.pipeline import Pipeline as SkPipeline
from sklearn.preprocessing import OneHotEncoder as SkOHE
from sklearn.ensemble import RandomForestClassifier as SkRF
from sklearn.metrics import roc_auc_score
echantillon = train.sample(fraction=0.05, seed=42).toPandas()
X = echantillon.drop(columns=["retard_15min", "retard_arrivee"])
y = echantillon["retard_15min"]
prep = ColumnTransformer([
("cat", SkOHE(handle_unknown="ignore"), ["compagnie", "origine", "destination"]),
], remainder="passthrough")
pipe = SkPipeline([("prep", prep), ("rf", SkRF(n_estimators=200, max_depth=10, n_jobs=-1))])
pipe.fit(X, y)
Résultats mesurés sur notre banc d'essai (8 cœurs, 32 Go) :
| Chaîne | Volume train | Temps train | AUC test |
|---|---|---|---|
| Spark, forêt, tout le train | 34 M lignes | 47 min | 0,761 |
sklearn, forêt, échantillon 5 % | 1,7 M lignes | 3 min | 0,741 |
xgboost, échantillon 5 % | 1,7 M lignes | 2 min | 0,758 |
Deux enseignements. Le volume rapporte peu ici : passer de 5 % à
100 % gagne 0,003 d'AUC sur sklearn, 0,003 sur xgboost. xgboost
sur échantillon bat sklearn sur échantillon — algorithme plus fort à
volume égal. Sur ce projet, la chaîne « Spark pour préparer, xgboost
pour entraîner sur un échantillon représentatif » aurait produit un
modèle presque équivalent en une fraction du temps.
Que retirer de la comparaison ?
Le volume massif n'est pas systématiquement rentable. Il l'est quand la distribution des variables évolue (le jeu couvre trop de saisonnalités pour qu'un échantillon les capture toutes), quand la cible est très déséquilibrée (moins de 1 % de positifs, chaque exemple compte), ou quand le signal se trouve dans les combinaisons rares (routes inhabituelles). Sur un problème à signal dense et distribution stable, un bon échantillon suffit.
C'est un arbitrage à formuler explicitement, avec un chiffre. Ce cours ne vous invite pas à sortir Spark par principe, mais à le sortir en connaissance de cause — et à mesurer l'apport plutôt qu'à le supposer.
Sur un nouveau problème, commencer par un sample(fraction=0.05)
converti en pandas. Vous itérez cent fois plus vite, vous validez la
faisabilité, vous éliminez les variables inutiles. Puis vous décidez si
Spark sur tout le train apportera un gain mesuré. Cette discipline évite
d'attendre 47 minutes à chaque essai d'hyperparamètre.
En résumé
- La chaîne complète tient en une trentaine de lignes de PySpark, avec pipeline, validation croisée, sauvegarde et écriture.
- Le découpage temporel train / test est indispensable sur les vols : toute autre séparation triche.
- Le volume massif ne rapporte pas toujours : mesurer avec un échantillon avant de payer le prix du plein volume.
xgboost.sparkou monomachine sont souvent plus efficaces queRandomForestde MLlib à qualité comparable.
Module suivant : le récapitulatif complet et l'examen final de 40 questions.