Aller au contenu principal

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.

La fuite temporelle

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 Pipeline enchaîne transformateurs et estimateurs et garantit que chaque étape apprend uniquement sur son entrée.
  • CrossValidator distribue 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, randomSplit fuit 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.