Module 9 — Écriture des résultats et intégration en aval
Un modèle qui ne sort jamais du carnet n'a aucune valeur métier. Ce module traite des trois artefacts que produit un travail Spark ML : les prédictions écrites, le modèle sauvegardé, et la chaîne d'inférence qui les rejoint en production. Chacun a des règles propres qu'il faut respecter à la lettre pour éviter les mauvaises surprises.
Écrire un DataFrame
L'écriture par défaut est en Parquet, partitionnée, avec un mode qui gère la présence préalable du dossier :
(
predictions
.select("vol_id", "annee", "mois", "compagnie", "probabilite_retard")
.write
.mode("overwrite") # remplace le contenu existant
.partitionBy("annee", "mois")
.parquet("s3://scoring/vols/predictions/")
)
Les quatre modes d'écriture ont des conséquences très différentes qu'il faut savoir énoncer :
"errorifexists"(défaut) : échoue si le chemin existe. Prudent pour du calcul manuel, gênant pour de la production automatique."overwrite": supprime puis récrit. Attention,partitionBycombiné àoverwriteen mode dynamique nécessitespark.sql.sources.partitionOverwriteMode=dynamicpour n'écraser que les partitions écrites, pas tout le dossier."append": ajoute de nouveaux fichiers dans le dossier existant. Simple, mais crée rapidement des collections de milliers de petits fichiers si utilisé quotidiennement — préférer une compaction périodique."ignore": ne fait rien si le chemin existe. Utile pour idempotence.
Formats et cible aval
Choisir la cible en fonction de la consommation :
- Parquet partitionné : format neutre, lu par presque tout (Spark, DuckDB, Trino, Athena, BigQuery). Défaut sain.
- Delta Lake / Iceberg / Hudi : Parquet avec un journal
transactionnel. Permet
INSERT,MERGE,UPDATE, versionnement et voyage dans le temps (versionAsOf). C'est le défaut sur Databricks (Delta) et devient standard dans les entrepôts modernes. - JDBC (
.jdbc(url, table, mode, properties)) : écriture directe dans PostgreSQL, MySQL, SQL Server. Attention à la parallélisation : 200 exécuteurs qui insèrent en parallèle saturent la base.numPartitions=8ou moins, batching activé. - Kafka : sortie en flux via
.writeStream.format("kafka"), sujet du dernier cours du parcours.
Sauvegarder le pipeline entraîné
Le PipelineModel sauvegarde toutes les étapes ajustées, y compris les
vocabulaires appris par StringIndexer et les moyennes de
StandardScaler. Le format est un dossier structuré, portable entre
versions mineures de Spark :
meilleur_modele.write().overwrite().save("s3://modeles/vols/retard-v3")
Le rechargement se fait par PipelineModel.load :
from pyspark.ml import PipelineModel
modele = PipelineModel.load("s3://modeles/vols/retard-v3")
predictions = modele.transform(vols_a_scorer)
Une seule commande relit tout : les indexeurs, l'encodeur, l'assembleur, la régression. C'est cette portabilité qui rend le pipeline indispensable — sans lui, il faudrait resérialiser chaque étape à la main, et le premier oubli casse silencieusement l'inférence.
Un modèle doit se nommer avec une version explicite dans son chemin :
retard-v3, jamais retard-final. La production peut à tout moment
avoir besoin de revenir à une version antérieure — un modèle écrasé est
un modèle perdu. Coupler avec un fichier metadata.json qui contient
la date d'entraînement, la métrique obtenue, la période des données.
Le scoring en aval
Le scoring quotidien lit le pipeline, lit les nouvelles données, transforme, écrit :
def scorer_quotidiennement(date_str: str) -> None:
spark = SparkSession.builder.appName("scoring-vols").getOrCreate()
modele = PipelineModel.load("s3://modeles/vols/retard-v3")
a_scorer = (
spark.read.parquet("s3://vols/planifies/")
.filter(F.col("date_depart") == date_str)
)
predictions = modele.transform(a_scorer)
(
predictions
.select("vol_id", "date_depart", "compagnie",
F.col("probability").getItem(1).alias("proba_retard"))
.write.mode("overwrite")
.parquet(f"s3://scoring/vols/date={date_str}/")
)
if __name__ == "__main__":
import sys
scorer_quotidiennement(sys.argv[1])
Trois précautions récurrentes. filter(F.col("date_depart") == date_str)
avant transform limite les données lues à la journée traitée, sinon
on rescore tout le passé chaque nuit. Le vecteur probability de
Spark ML est un Vector à deux composantes pour la classification
binaire ; extraire l'indice 1 pour la probabilité de la classe positive.
Réutiliser le même chemin d'écriture par date rend le travail
idempotent : le rejouer donne le même résultat sans doublonner.
Planification
Le travail se lance par un ordonnanceur : Airflow (SparkSubmitOperator),
Databricks Jobs, cron, GitHub Actions pour du petit volume. Trois
propriétés à garantir :
- Idempotence : rejouer le travail deux fois pour la même date produit exactement le même résultat.
- Journalisation : capturer le journal complet, y compris l'ID d'application Spark, pour retrouver l'exécution dans l'interface.
- Alerte : le travail doit remonter clairement une erreur ; en
cluster, un
exit(1)du script principal est ce que l'ordonnanceur détecte.
Réentraîner : quand et comment
Un modèle vieillit. Deux stratégies :
- Fenêtre glissante : chaque mois, réentraîner sur les 24 derniers mois, remplacer la version en production.
- Détection de dérive (drift) : surveiller la distribution des variables d'entrée et la métrique aval ; réentraîner quand un seuil est franchi.
Dans les deux cas, le pipeline complet, du prétraitement à l'algorithme, doit être réentraîné ensemble : réentraîner uniquement l'algorithme sur des variables produites par un ancien prétraitement est un contrat rompu, qui produit des erreurs silencieuses.
En résumé
- Le mode d'écriture (
overwrite,append,errorifexists,ignore) gouverne le comportement face à un dossier existant ; le connaître avant de lancer un travail de production. - Un
PipelineModelsauvegarde toutes les étapes ajustées ;.savepuis.loadsont l'API portable. - Le scoring quotidien doit être idempotent, filtré à la journée utile, et versionné par chemin.
- Le pipeline complet se réentraîne ensemble ; réentraîner seulement l'algorithme casse silencieusement la chaîne.
Module suivant : le projet complet de bout en bout, de la lecture à
l'écriture, en comparant à un traitement pandas équivalent.