Module 6 — Chaînes de traitement et validation croisée distribuée
Au module précédent, chaque étape de prétraitement était appliquée à la
main : indexer, encoder, assembler, mettre à l'échelle. Cette chaîne
tenait tant qu'elle restait courte. Dès qu'on la duplique pour la
validation croisée, elle devient une source certaine d'erreurs, dont la
plus courante — la fuite de données — reste indétectable en
production. Le Pipeline de Spark ML sert exactement à empêcher cela.
Le Pipeline de Spark ML
Un Pipeline est une séquence ordonnée de Transformer et
Estimator. Appeler fit sur le pipeline entraîne les estimateurs dans
l'ordre, chaque étape transformant le DataFrame avant la suivante. Le
résultat est un PipelineModel, lui-même un Transformer qui rejoue la
même séquence à l'inférence.
from pyspark.ml import Pipeline
from pyspark.ml.feature import StringIndexer, OneHotEncoder, VectorAssembler
from pyspark.ml.classification import LogisticRegression
indexeur_c = StringIndexer(inputCol="compagnie", outputCol="compagnie_idx", handleInvalid="keep")
indexeur_o = StringIndexer(inputCol="origine", outputCol="origine_idx", handleInvalid="keep")
encodeur = OneHotEncoder(inputCols=["compagnie_idx", "origine_idx"],
outputCols=["compagnie_ohe", "origine_ohe"])
assembleur = VectorAssembler(
inputCols=["distance_km", "heure_depart", "jour_semaine", "compagnie_ohe", "origine_ohe"],
outputCol="features", handleInvalid="skip",
)
regression = LogisticRegression(featuresCol="features", labelCol="retard_15min", maxIter=50)
pipeline = Pipeline(stages=[indexeur_c, indexeur_o, encodeur, assembleur, regression])
modele = pipeline.fit(vols_train)
predictions = modele.transform(vols_test)
Trois avantages immédiats. Reproductibilité : sauvegarder le
pipeline sauvegarde toutes les étapes ensemble, y compris les
vocabulaires appris (module 9). Clarté : la chaîne est un objet
unique qu'on relit d'un coup d'œil. Correction : le fit du
pipeline garantit que chaque estimateur apprend uniquement sur les
données que le pipeline lui-même a produites, sans que vous ayez à y
penser.
Validation croisée distribuée
CrossValidator combine un estimateur (souvent un pipeline complet), une
grille d'hyperparamètres, une métrique d'évaluation et un nombre de plis
k. Il évalue chaque combinaison sur k découpes et retourne le
meilleur modèle réentraîné sur tout le jeu d'entraînement.
from pyspark.ml.tuning import CrossValidator, ParamGridBuilder
from pyspark.ml.evaluation import BinaryClassificationEvaluator
grille = (
ParamGridBuilder()
.addGrid(regression.regParam, [0.0, 0.01, 0.1])
.addGrid(regression.elasticNetParam, [0.0, 0.5])
.build()
)
evaluateur = BinaryClassificationEvaluator(
labelCol="retard_15min", rawPredictionCol="rawPrediction", metricName="areaUnderROC",
)
cv = CrossValidator(
estimator=pipeline,
estimatorParamMaps=grille,
evaluator=evaluateur,
numFolds=5,
parallelism=4, # nombre de combinaisons ajustées en parallèle
seed=42,
)
modele_cv = cv.fit(vols_train)
meilleur_modele = modele_cv.bestModel
La distribution du calcul se fait selon deux axes. Sur les données :
chaque fit reste distribué sur les exécuteurs comme d'habitude. Sur
les combinaisons : parallelism=4 demande au pilote d'ajuster jusqu'à
quatre combinaisons en même temps, chacune bénéficiant du même pool
d'exécuteurs. Sur un petit cluster, un parallelism trop élevé écrase
les exécuteurs — commencez modeste.
Le coût d'une grille, sans naïveté
CrossValidator réalise nb_combinaisons × nb_plis entraînements
complets, plus un entraînement final sur tout le train. Six
combinaisons et cinq plis, c'est 31 entraînements. Sur un pipeline
qui dure vingt minutes, on parle de plus de dix heures. La grille doit
être choisie avec la même parcimonie qu'en scikit-learn — ce n'est pas
parce que Spark est distribué que la combinatoire disparaît.
Alternative : TrainValidationSplit fait une seule division train / val
et divise le temps par k. Moins fiable, mais souvent suffisant pour
dégrossir avant un CrossValidator final sur la grille resserrée.
La fuite quand le prétraitement sort du pipeline
Voici le piège central, avec le code fautif :
# INCORRECT
echelle = StandardScaler(inputCol="features_num", outputCol="features")
vols_std = echelle.fit(vols).transform(vols) # apprend sur tout
vols_train, vols_test = vols_std.randomSplit([0.8, 0.2], seed=42)
modele = regression.fit(vols_train) # trop optimiste
score = evaluateur.evaluate(modele.transform(vols_test))
Le StandardScaler a vu la totalité du jeu, y compris ce qui atterrira
en test. La moyenne et l'écart-type appris incorporent des informations
de test — c'est bien de la fuite, silencieuse et intégrée. La version
saine intègre l'échelle dans le pipeline :
# CORRECT
pipeline = Pipeline(stages=[assembleur, echelle, regression])
vols_train, vols_test = vols.randomSplit([0.8, 0.2], seed=42)
modele = pipeline.fit(vols_train)
score = evaluateur.evaluate(modele.transform(vols_test))
Sur notre jeu de vols, la version fautive a produit un AUC apparent de 0,86, contre 0,82 réel — quatre points de trop qui se paient à la première mise en production, quand les métriques mesurées ne correspondent plus à celles observées.
Sur des données temporelles comme les vols, randomSplit mélange le
passé et le futur. Le train contient alors des vols postérieurs à
certains vols du test, ce qui est impossible en pratique. Découper par
date (filter(F.col("annee") <= 2022) pour le train, >= 2023 pour le
test) est la seule séparation honnête pour ce type de problème.
Extraire les meilleurs hyperparamètres
Après le CrossValidator, on inspecte quels hyperparamètres ont gagné :
etape_finale = meilleur_modele.stages[-1] # LogisticRegressionModel
print("regParam :", etape_finale._java_obj.getRegParam())
print("elasticNetParam :", etape_finale._java_obj.getElasticNetParam())
for params, metrique in zip(modele_cv.getEstimatorParamMaps(), modele_cv.avgMetrics):
print(f"{{ {params} }} -> {metrique:.4f}")
C'est la même liste que dans scikit-learn : les vainqueurs guident la prochaine grille, plus resserrée autour d'eux.
En résumé
- Un
Pipelineenchaîne transformateurs et estimateurs et garantit que chaque étape apprend uniquement sur son entrée. CrossValidatordistribue les entraînements sur les exécuteurs et parallélise les combinaisons ; le coût reste combinatoire, garder la grille resserrée.- Le prétraitement hors pipeline est la source classique de fuite de données ; toujours intégrer l'échelle et l'encodage à l'intérieur.
- Sur des données temporelles,
randomSplitfuit dans le temps ; préférer une séparation par date.
Module suivant : quels algorithmes existent réellement dans MLlib, et lesquels manquent avec leurs alternatives.