Files
ENI-ml-mlops/pipeline.py

71 lines
2.9 KiB
Python

"""Pipeline ML automatise : split -> entrainement (--register) -> promotion.
Orchestre en une seule commande (`python pipeline.py`) les etapes du fil rouge
menant a la mise a disposition du modele dans le MLflow Model Registry :
1. split : (re)cree data/{train,validation,test}.parquet depuis /data/modelling
(strategie de split figee dans lab/constants.py) ;
2. train : entraine la strategie de service (`full`) et enregistre une
nouvelle version du modele dans le Registry (log_model + --register) ;
3. promote : (re)pointe l'alias `champion` vers cette nouvelle version.
Chaque etape est un sous-processus lance avec `check=True` : si une etape echoue,
le pipeline s'arrete immediatement avec un code de sortie non nul et n'execute
aucune etape suivante (Partie 2, Q2.2). Le pipeline est reexecutable sans effet
de bord indesirable (idempotence, Q2.3) : le split est deterministe, `--register`
empile proprement une nouvelle version, et l'alias `champion` est repointe
(jamais duplique).
Les etapes reutilisent telles quelles les CLI du package `lab/` (memes commandes
qu'une execution manuelle des TP precedents) : le pipeline ne fait que les
enchainer et calculer le numero de la version a promouvoir.
"""
import logging
import subprocess
import sys
from mlflow import MlflowClient
from lab import constants
logging.basicConfig(level=logging.INFO, format="%(asctime)s [pipeline] %(message)s")
logger = logging.getLogger("pipeline")
def run_step(title: str, module_args: list[str]) -> None:
"""Execute une etape (module `lab`) dans un sous-processus ; stoppe si elle echoue."""
logger.info("=== ETAPE %s : python -m %s", title, " ".join(module_args))
subprocess.run([sys.executable, "-m", *module_args], check=True)
def latest_model_version(name: str) -> int:
"""Numero de la derniere version enregistree du modele (celle qui vient d'etre creee)."""
client = MlflowClient()
versions = client.search_model_versions(f"name = '{name}'")
if not versions:
raise RuntimeError(f"Aucune version enregistree pour le modele '{name}'.")
return max(int(mv.version) for mv in versions)
def main() -> None:
strategy = constants.SERVING_STRATEGY
name = constants.REGISTERED_MODEL_NAME
alias = constants.MODEL_ALIAS
# 1. Preparation des jeux de donnees (deterministe).
run_step("split", ["lab.split.cli"])
# 2. Entrainement + enregistrement d'une nouvelle version dans le Model Registry.
run_step("train", ["lab.modeling.cli", strategy.value, "--register"])
# 3. Promotion : repointer l'alias vers la version qui vient d'etre enregistree.
version = latest_model_version(name)
run_step("promote", ["lab.registry.cli", "promote", "--version", str(version), "--alias", alias])
logger.info("Pipeline termine avec succes : %s v%s -> alias '%s'", name, version, alias)
if __name__ == "__main__":
main()