Aller au contenu principal

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îneVolume trainTemps trainAUC test
Spark, forêt, tout le train34 M lignes47 min0,761
sklearn, forêt, échantillon 5 %1,7 M lignes3 min0,741
xgboost, échantillon 5 %1,7 M lignes2 min0,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.

La règle « échantillonner d'abord »

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.spark ou monomachine sont souvent plus efficaces que RandomForest de MLlib à qualité comparable.

Module suivant : le récapitulatif complet et l'examen final de 40 questions.