Compare commits
10 Commits
f8549972d5
...
main
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
ead28bf663 | ||
|
|
7c5660c455 | ||
|
|
67566ef076 | ||
|
|
47109efad9 | ||
|
|
e063effc5d | ||
|
|
b0bd6cdb07 | ||
|
|
9a52395c91 | ||
|
|
e6a05fbbc7 | ||
|
|
d8c2e85083 | ||
|
|
0c0bd77156 |
2
.dvc/.gitignore
vendored
Normal file
2
.dvc/.gitignore
vendored
Normal file
@@ -0,0 +1,2 @@
|
||||
/tmp
|
||||
/cache
|
||||
8
.dvc/config
Normal file
8
.dvc/config
Normal file
@@ -0,0 +1,8 @@
|
||||
[core]
|
||||
analytics = false
|
||||
remote = garage
|
||||
['remote "garage"']
|
||||
url = s3://dvc-store
|
||||
endpointurl = https://garage.192-168-122-143.nip.io
|
||||
region = garage
|
||||
ssl_verify = /etc/ssl/certs/ca-certificates.crt
|
||||
3
.dvcignore
Normal file
3
.dvcignore
Normal file
@@ -0,0 +1,3 @@
|
||||
# Add patterns of files dvc should ignore, which could improve
|
||||
# the performance. Learn more at
|
||||
# https://dvc.org/doc/user-guide/dvcignore
|
||||
22
.env.example
Normal file
22
.env.example
Normal file
@@ -0,0 +1,22 @@
|
||||
# Copier en .env (gitignore) et renseigner les secrets.
|
||||
# A sourcer avant de lancer les scripts : set -a; source .env; set +a
|
||||
|
||||
# --- MLflow (serveur de la VM, basic auth) ---
|
||||
MLFLOW_TRACKING_URI=https://mlflow.192-168-122-143.nip.io
|
||||
MLFLOW_TRACKING_USERNAME=admin
|
||||
MLFLOW_TRACKING_PASSWORD=change-me
|
||||
MLFLOW_EXPERIMENT_NAME=tp02_electricity_consumption
|
||||
|
||||
# --- CA interne ENI MLOps (les clients Python n'utilisent pas le store systeme par defaut) ---
|
||||
REQUESTS_CA_BUNDLE=/etc/ssl/certs/ca-certificates.crt
|
||||
AWS_CA_BUNDLE=/etc/ssl/certs/ca-certificates.crt
|
||||
|
||||
# --- S3 Garage : artefacts MLflow (TP03) ---
|
||||
# Necessaire pour log_model (upload de l'artefact) ET pour le chargement du modele
|
||||
# par l'API (download depuis s3://mlflow-artifacts). Cle S3 "mlflow" (RWO sur le bucket).
|
||||
AWS_ACCESS_KEY_ID=GKxxxxxxxxxxxxxxxxxxxxxxxx
|
||||
AWS_SECRET_ACCESS_KEY=change-me
|
||||
MLFLOW_S3_ENDPOINT_URL=https://garage.192-168-122-143.nip.io
|
||||
|
||||
# --- Import du package lab ---
|
||||
PYTHONPATH=/home/user/tp
|
||||
50
.forgejo/workflows/pipeline.yml
Normal file
50
.forgejo/workflows/pipeline.yml
Normal file
@@ -0,0 +1,50 @@
|
||||
name: pipeline
|
||||
|
||||
# Partie 4 : le pipeline ML est rejoue automatiquement a chaque Pull Request
|
||||
# vers la branche principale (Q4.1). Autres declencheurs possibles (Q4.4) :
|
||||
# push sur main, tag de release, planification (schedule/cron), declenchement
|
||||
# manuel (workflow_dispatch).
|
||||
on:
|
||||
pull_request:
|
||||
branches: [main]
|
||||
|
||||
jobs:
|
||||
train-and-register:
|
||||
runs-on: docker
|
||||
# Le runner (mode docker, reseau mlops-net) lance le job dans ce conteneur.
|
||||
# On y monte en lecture seule les donnees source du fil rouge et le magasin
|
||||
# de certificats du host (pour faire confiance a la CA interne ENI MLOps).
|
||||
container:
|
||||
image: node:20-bookworm
|
||||
volumes:
|
||||
- /data/modelling:/data/modelling:ro
|
||||
- /etc/ssl/certs/ca-certificates.crt:/etc/ssl/certs/ca-certificates.crt:ro
|
||||
steps:
|
||||
- name: Recuperation du code
|
||||
uses: actions/checkout@v4
|
||||
|
||||
- name: Installation de Python et des dependances
|
||||
run: |
|
||||
apt-get update
|
||||
apt-get install -y --no-install-recommends python3-venv
|
||||
python3 -m venv .venv
|
||||
.venv/bin/pip install --upgrade pip
|
||||
.venv/bin/pip install -r requirements.txt
|
||||
|
||||
- name: Execution du pipeline (split -> train --register -> promote)
|
||||
run: .venv/bin/python pipeline.py
|
||||
env:
|
||||
# Serveur MLflow (basic auth) : URLs HTTPS via Caddy, resolues sur mlops-net.
|
||||
MLFLOW_TRACKING_URI: ${{ secrets.MLFLOW_TRACKING_URI }}
|
||||
MLFLOW_TRACKING_USERNAME: ${{ secrets.MLFLOW_TRACKING_USERNAME }}
|
||||
MLFLOW_TRACKING_PASSWORD: ${{ secrets.MLFLOW_TRACKING_PASSWORD }}
|
||||
MLFLOW_EXPERIMENT_NAME: ${{ secrets.MLFLOW_EXPERIMENT_NAME }}
|
||||
# Artefacts sur S3 Garage (log_model).
|
||||
AWS_ACCESS_KEY_ID: ${{ secrets.AWS_ACCESS_KEY_ID }}
|
||||
AWS_SECRET_ACCESS_KEY: ${{ secrets.AWS_SECRET_ACCESS_KEY }}
|
||||
MLFLOW_S3_ENDPOINT_URL: ${{ secrets.MLFLOW_S3_ENDPOINT_URL }}
|
||||
# CA interne ENI MLOps (requests/boto3 n'utilisent pas le store systeme par defaut).
|
||||
REQUESTS_CA_BUNDLE: /etc/ssl/certs/ca-certificates.crt
|
||||
AWS_CA_BUNDLE: /etc/ssl/certs/ca-certificates.crt
|
||||
# Import du package lab depuis la racine du depot.
|
||||
PYTHONPATH: ${{ github.workspace }}
|
||||
23
.gitignore
vendored
23
.gitignore
vendored
@@ -1,7 +1,17 @@
|
||||
# Donnees du fil rouge (volumineuses, gerees hors git / DVC plus tard)
|
||||
data/
|
||||
*.parquet
|
||||
*.csv
|
||||
# Donnees : gerees par DVC. Seuls les pointeurs .dvc et le .gitignore
|
||||
# genere par DVC sont versionnes ; les .parquet reels vont sur le remote S3 (Garage).
|
||||
/data/*
|
||||
!/data/*.dvc
|
||||
!/data/.gitignore
|
||||
|
||||
# Secrets / environnement
|
||||
.env
|
||||
*.key
|
||||
*.pem
|
||||
.dvc/config.local
|
||||
|
||||
# MLflow local eventuel (on utilise le serveur distant)
|
||||
mlruns/
|
||||
|
||||
# Jupyter
|
||||
.ipynb_checkpoints/
|
||||
@@ -14,8 +24,3 @@ __pycache__/
|
||||
*.pyc
|
||||
.venv/
|
||||
venv/
|
||||
|
||||
# Secrets / environnement
|
||||
.env
|
||||
*.key
|
||||
*.pem
|
||||
|
||||
135
README.md
Normal file
135
README.md
Normal file
@@ -0,0 +1,135 @@
|
||||
# TP02 - Comparer et tracer des experimentations ML (DVC + MLflow)
|
||||
|
||||
Fil rouge : prediction de la consommation electrique. Ce module compare des strategies
|
||||
de features et de split, en versionnant les datasets avec **DVC** (remote S3 = Garage de la VM)
|
||||
et en tracant les experiences avec **MLflow** (serveur de la VM).
|
||||
|
||||
## Prerequis (sur la VM)
|
||||
|
||||
- venv : `/opt/venvs/mlops`
|
||||
- donnees source : `/data/modelling/features.parquet` + `target.parquet`
|
||||
- `.env` renseigne (voir `.env.example`), puis :
|
||||
|
||||
```bash
|
||||
cd /home/user/tp
|
||||
set -a; source .env; set +a
|
||||
alias py=/opt/venvs/mlops/bin/python
|
||||
alias dvc=/opt/venvs/mlops/bin/dvc
|
||||
```
|
||||
|
||||
## Package
|
||||
|
||||
- `lab/constants.py` : strategies de split (`full_history`, `recent_history`), strategies de
|
||||
features (`short_memory`, `seasonality`, `tendency`, `mixed`, `full`), `RIDGE_ALPHAS`.
|
||||
- `lab/split/cli.py` : lit `/data/modelling`, ecrit `data/{train,validation,test}.parquet`
|
||||
selon `CHOSEN_SPLIT_STRATEGY`.
|
||||
- `lab/modeling/cli.py` : `py -m lab.modeling.cli <strategy>` -> regression lineaire -> MLflow.
|
||||
- `lab/modeling_ridge/cli.py` : `py -m lab.modeling_ridge.cli [--strategy <s>]` -> Ridge sur `RIDGE_ALPHAS`.
|
||||
|
||||
## Versionner un dataset avec DVC
|
||||
|
||||
```bash
|
||||
py -m lab.split.cli
|
||||
dvc add data/train.parquet data/validation.parquet data/test.parquet
|
||||
git add data/*.dvc data/.gitignore lab .gitignore
|
||||
git commit -m "split <strategie>"
|
||||
git tag dataset-<version>
|
||||
dvc push
|
||||
```
|
||||
|
||||
Restaurer une version anterieure du dataset :
|
||||
|
||||
```bash
|
||||
git checkout dataset-v1-full-history -- data/train.parquet.dvc data/validation.parquet.dvc data/test.parquet.dvc
|
||||
dvc checkout
|
||||
```
|
||||
|
||||
## Entrainements
|
||||
|
||||
```bash
|
||||
for s in short_memory seasonality tendency mixed; do py -m lab.modeling.cli $s; done # Parties 1 et 2
|
||||
py -m lab.modeling_ridge.cli --strategy mixed # Partie 3
|
||||
```
|
||||
|
||||
Resultats et comparaisons : https://mlflow.192-168-122-143.nip.io (experience `tp02_electricity_consumption`).
|
||||
|
||||
## Livrable
|
||||
|
||||
Synthese des resultats et reponses aux questions : `SYNTHESE.md`.
|
||||
|
||||
---
|
||||
|
||||
# TP03 - Exposer un modele via une API REST (Model Registry + FastAPI)
|
||||
|
||||
Prolonge le TP02 : on enregistre le meilleur modele dans le **MLflow Model Registry**, on le
|
||||
**promeut via un alias**, puis on l'expose par une **API REST FastAPI**.
|
||||
|
||||
## Prerequis (en plus du TP02)
|
||||
|
||||
`.env` complete avec les creds S3 Garage (voir `.env.example`) : `AWS_ACCESS_KEY_ID`,
|
||||
`AWS_SECRET_ACCESS_KEY`, `MLFLOW_S3_ENDPOINT_URL`. Necessaires pour `log_model` (upload de
|
||||
l'artefact) et pour le chargement du modele par l'API (download depuis `s3://mlflow-artifacts`).
|
||||
|
||||
## Enregistrer et promouvoir (Partie 1)
|
||||
|
||||
```bash
|
||||
set -a; source .env; set +a
|
||||
py -m lab.modeling.cli full --register # v1 -> Registry (champion vise)
|
||||
py -m lab.modeling.cli mixed --register # v2 -> Registry (comparaison)
|
||||
py -m lab.registry.cli versions # lister versions + alias
|
||||
py -m lab.registry.cli promote --version 1 --alias champion
|
||||
```
|
||||
|
||||
`log_model` est appele a **chaque** run (artefact sauvegarde) ; `--register` empile en plus une
|
||||
version dans le Registry sous le nom `electricity-consumption`.
|
||||
|
||||
## Servir l'API (Parties 2 a 4)
|
||||
|
||||
```bash
|
||||
./serve.sh # uvicorn 0.0.0.0:8000, charge models:/...@champion
|
||||
curl -s localhost:8000/health # {"status":"ok"}
|
||||
curl -s -X POST localhost:8000/predict \
|
||||
-H 'content-type: application/json' -d '{"client_id":"MT_124"}'
|
||||
curl -s -X POST localhost:8000/predict/batch \
|
||||
-H 'content-type: application/json' -d '{"client_ids":["MT_124","MT_158"]}'
|
||||
```
|
||||
|
||||
- Endpoints : `GET /health`, `POST /predict`, `POST /predict/batch`, Swagger `GET /docs`.
|
||||
- Feature store **simule** par un dictionnaire (`lab/serving/features.py`) ; client inconnu -> **404**.
|
||||
- Depuis le poste (cert ENI de confiance) : **https://api.192-168-122-143.nip.io/docs**
|
||||
(reverse-proxy Caddy vers uvicorn). Unite systemd transitoire : `sudo systemctl status tp03-api`.
|
||||
|
||||
## Livrable TP03
|
||||
|
||||
Reponses aux questions et recap : `SYNTHESE_TP03.md`.
|
||||
|
||||
---
|
||||
|
||||
# TP04 - Automatiser le pipeline avec la CI/CD (pipeline.py + Forgejo Actions)
|
||||
|
||||
Prolonge les TP precedents : les etapes manuelles (split -> entrainement -> enregistrement/promotion
|
||||
Registry) sont orchestrees en une seule commande, puis rejouees automatiquement en CI.
|
||||
|
||||
## Pipeline (Parties 2-3)
|
||||
|
||||
`pipeline.py` (racine) enchaine, via des sous-processus (`check=True`, arret au premier echec) :
|
||||
|
||||
```bash
|
||||
set -a; source .env; set +a
|
||||
python pipeline.py # split -> train full --register -> promote champion
|
||||
```
|
||||
|
||||
- 1. `lab.split.cli` : (re)cree `data/{train,validation,test}.parquet` (deterministe) ;
|
||||
- 2. `lab.modeling.cli full --register` : entraine + `log_model` + nouvelle version au Registry ;
|
||||
- 3. `lab.registry.cli promote` : repointe l'alias `champion` sur cette version.
|
||||
|
||||
`requirements.txt` fige les dependances runtime (env de reference de la VM, Python 3.14).
|
||||
|
||||
## CI/CD (Partie 4)
|
||||
|
||||
`.forgejo/workflows/pipeline.yml` : declenche sur **pull_request vers `main`**, installe les
|
||||
dependances puis lance `pipeline.py`. Execute par les runners du **Forgejo de la VM**
|
||||
(`forgejo.192-168-122-143.nip.io`), mode docker sur `mlops-net` (les donnees `/data/modelling` et la
|
||||
CA du host sont montees en lecture seule). Secrets (MLflow + S3 Garage) dans les secrets Actions du repo.
|
||||
|
||||
Livrable : `SYNTHESE_TP04.md`.
|
||||
200
SYNTHESE.md
Normal file
200
SYNTHESE.md
Normal file
@@ -0,0 +1,200 @@
|
||||
# TP02 - Synthèse : comparer et tracer des expérimentations ML
|
||||
|
||||
Fil rouge : prédiction de la consommation électrique (kWh). Datasets versionnés avec **DVC**
|
||||
(remote S3 = Garage), expériences tracées avec **MLflow** (expérience `tp02_electricity_consumption`,
|
||||
13 runs : 10 régressions linéaires + 3 Ridge).
|
||||
|
||||
Deux versions de dataset produites et versionnées (tags git + DVC) :
|
||||
|
||||
| Version DVC | Tag git | Train | Validation | Test |
|
||||
|---|---|---|---|---|
|
||||
| v1 `full_history` | `dataset-v1-full-history` | 2011-2012 | 2013 | 2014 |
|
||||
| v2 `recent_history` | `dataset-v2-recent-history` | 2013 | jan-mai 2014 | juin-déc 2014 |
|
||||
|
||||
Métriques comparées : **RMSE** et **MAE** sur la **validation** (le test reste réservé).
|
||||
|
||||
---
|
||||
|
||||
## Partie 1 - Comparaison de stratégies de features (split `full_history`)
|
||||
|
||||
Stratégies testées (hypothèse -> features) :
|
||||
|
||||
| Stratégie | Hypothèse | Features |
|
||||
|---|---|---|
|
||||
| `short_memory` | la conso dépend surtout de la veille | lag_1d |
|
||||
| `seasonality` | conso plus saisonnière que journalière | lag_7d, lag_30d |
|
||||
| `tendency` | conso surtout tendancielle | rolling_mean_7d, rolling_mean_30d |
|
||||
| `mixed` | mémoire courte + tendance | lag_1d, lag_7d, lag_30d, rolling_mean_30d |
|
||||
| `full` | toutes les features | les 6 |
|
||||
|
||||
Résultats (validation, `full_history`) :
|
||||
|
||||
| Stratégie | RMSE | MAE |
|
||||
|---|---|---|
|
||||
| `full` | **7.41** | **4.07** |
|
||||
| `mixed` | 7.53 | 4.11 |
|
||||
| `short_memory` | 8.82 | 4.58 |
|
||||
| `seasonality` | 9.14 | 5.04 |
|
||||
| `tendency` | 20.79 | 14.66 |
|
||||
|
||||
**1.1 - Quelle stratégie obtient les meilleures performances ?**
|
||||
`full` (6 features) puis `mixed` (4 features), quasi à égalité. Toutes deux combinent lags courts
|
||||
et tendance. `tendency` (moyennes glissantes seules) est de loin la pire.
|
||||
|
||||
**1.2 - Ajouter davantage de features améliore-t-il systématiquement le modèle ?**
|
||||
Non. `short_memory` (1 feature) fait 8.82 alors que `tendency` (2 features) fait 20.79 : c'est la
|
||||
**pertinence** des features, pas leur nombre, qui compte. Et `full` (6) n'améliore `mixed` (4) que
|
||||
marginalement (7.41 vs 7.53) : rendements décroissants, les features supplémentaires apportent peu.
|
||||
|
||||
**1.3 - Quelles hypothèses semblent validées ?**
|
||||
- « La conso dépend de la veille » : validée, `lag_1d` est le prédicteur dominant (voir coefficients).
|
||||
- « Plutôt tendancielle » (moyennes glissantes seules) : **invalidée** (pire modèle).
|
||||
- La conso combine **mémoire court terme + tendance** : validée (mixed/full gagnants).
|
||||
|
||||
---
|
||||
|
||||
## Partie 2 - Comparaison des stratégies d'entraînement (split)
|
||||
|
||||
Chaque stratégie ré-entraînée sur `recent_history`. Matrice RMSE (validation) :
|
||||
|
||||
| Stratégie | `full_history` | `recent_history` |
|
||||
|---|---|---|
|
||||
| `short_memory` | 8.82 | 8.30 |
|
||||
| `seasonality` | 9.14 | 7.60 |
|
||||
| `tendency` | 20.79 | 18.56 |
|
||||
| `mixed` | 7.53 | 6.67 |
|
||||
| `full` | 7.41 | **6.58** |
|
||||
|
||||
La stratégie `full` a été relancée sur **la première version du dataset (v1) restaurée via DVC**
|
||||
(`git checkout dataset-v1-full-history -- ... && dvc checkout`), illustrant la reproductibilité.
|
||||
|
||||
**2.1 - Quel split obtient les meilleures performances ?**
|
||||
`recent_history` : RMSE plus bas pour **les 5 stratégies** (env. -12 %).
|
||||
|
||||
**2.2 - Davantage d'historique ou données plus récentes ?**
|
||||
Ici, **données plus récentes**. Nuance importante : la comparaison n'est pas strictement iso car les
|
||||
fenêtres de validation diffèrent (2013 pour `full_history` vs début 2014 pour `recent_history`).
|
||||
Prédire début 2014 à partir de 2013 (année adjacente) est plus « facile » que prédire 2013 à partir
|
||||
de 2011-2012. Enseignement : la **proximité temporelle train/validation** (moins de dérive) prime sur
|
||||
le simple volume d'historique ancien.
|
||||
|
||||
**2.3 - Avantages / inconvénients de chaque approche ?**
|
||||
- `full_history` : + capte les cycles longs (saisonnalité annuelle), robustesse au bruit ponctuel ;
|
||||
- inclut des données anciennes possiblement obsolètes (data/concept drift), plus coûteux, validation lointaine.
|
||||
- `recent_history` : + colle au régime actuel (moins de drift), moins coûteux ;
|
||||
- moins de données (variance accrue), saisonnalités longues moins fiables (`lag_365d`), sensible aux
|
||||
événements récents atypiques.
|
||||
|
||||
**2.4 - Bonus : coefficients du modèle** (linéaire `full`, `recent_history`), par influence décroissante :
|
||||
|
||||
| Feature | Coefficient |
|
||||
|---|---|
|
||||
| lag_1d | +0.460 |
|
||||
| lag_7d | +0.371 |
|
||||
| rolling_mean_7d | +0.288 |
|
||||
| rolling_mean_30d | -0.268 |
|
||||
| lag_365d | +0.090 |
|
||||
| lag_30d | +0.058 |
|
||||
|
||||
- Les plus influentes : `lag_1d` (la veille) et `lag_7d` (semaine dernière), puis la paire
|
||||
`rolling_mean_7d` (+) / `rolling_mean_30d` (-) qui agit en **différentiel de tendance** court/moyen terme.
|
||||
- Les peu utilisées : `lag_30d` et `lag_365d` (mensuel/annuel apportent peu une fois les autres présentes).
|
||||
- Cela **confirme les hypothèses de la Partie 1** : mémoire court terme + tendance courte dominent.
|
||||
|
||||
---
|
||||
|
||||
## Partie 3 - Influence des hyperparamètres avec Ridge
|
||||
|
||||
Sur le meilleur couple (features `full`, split `recent_history`), alphas `[1, 1e3, 1e9]` :
|
||||
|
||||
| alpha | RMSE (val) | MAE (val) | RMSE (train) |
|
||||
|---|---|---|---|
|
||||
| 1 | **6.576** | **3.677** | 7.397 |
|
||||
| 1e3 | **6.576** | **3.677** | 7.397 |
|
||||
| 1e9 | 7.059 | 4.185 | 8.115 |
|
||||
|
||||
Coefficients selon alpha :
|
||||
|
||||
| Feature | alpha=1 | alpha=1e3 | alpha=1e9 |
|
||||
|---|---|---|---|
|
||||
| lag_1d | 0.460 | 0.460 | 0.284 |
|
||||
| lag_7d | 0.371 | 0.371 | 0.263 |
|
||||
| lag_30d | 0.058 | 0.058 | 0.158 |
|
||||
| lag_365d | 0.090 | 0.090 | 0.174 |
|
||||
| rolling_mean_7d | 0.288 | 0.287 | 0.069 |
|
||||
| rolling_mean_30d | -0.268 | -0.268 | 0.043 |
|
||||
|
||||
**3.1 - Quelle valeur d'alpha obtient les meilleures performances ?**
|
||||
`alpha=1` (identique à `alpha=1e3`). `alpha=1e9` dégrade nettement (RMSE 7.06).
|
||||
|
||||
**3.2 - Que se passe-t-il sur les coefficients quand alpha augmente ?**
|
||||
Ils sont **contraints vers zéro** (régularisation L2). À `alpha=1` et `1e3` : quasi identiques, car la
|
||||
pénalité reste négligeable devant ~4,5 M d'échantillons. À `alpha=1e9` : forte contraction, les
|
||||
coefficients dominants s'écrasent (`lag_1d` 0.46 -> 0.28, `rolling_mean_7d` 0.29 -> 0.07,
|
||||
`rolling_mean_30d` -0.27 -> +0.04) ; le modèle devient plus « plat » (biais accru) -> RMSE augmente.
|
||||
|
||||
**3.3 - Pourquoi est-il indispensable de logger les hyperparamètres ?**
|
||||
Pour la **reproductibilité et la comparabilité** : une métrique n'a de sens qu'associée à ses HP (alpha,
|
||||
features, split). Sans, impossible de reproduire un run, d'expliquer un écart de performance ou de
|
||||
comparer objectivement. MLflow lie chaque métrique à ses paramètres -> traçabilité complète.
|
||||
|
||||
**3.4 - Bonus : que représente le bruit ? Faut-il l'apprendre ?**
|
||||
Le bruit = variations aléatoires non reproductibles (erreurs de mesure, aléas ponctuels). Il **ne faut
|
||||
pas** l'apprendre : le modèle mémoriserait des accidents particuliers (overfitting) au lieu de la relation
|
||||
générale, et généraliserait mal. Exemple : un pic de conso dû à une canicule exceptionnelle un 15 août
|
||||
ne doit pas devenir une règle.
|
||||
|
||||
**3.5 - L'ordre de grandeur des coefficients a-t-il un effet ? Et avec du bruit ?**
|
||||
Oui : de grands coefficients rendent la prédiction très sensible aux petites variations des features
|
||||
(amplification). En présence de bruit, ils **amplifient ce bruit** -> prédictions instables, variance
|
||||
élevée (overfitting). Des coefficients plus petits lissent la réponse.
|
||||
|
||||
**3.6 - Rôle de l'hyperparamètre alpha ?**
|
||||
alpha règle le **compromis biais/variance** en pénalisant la magnitude des coefficients : alpha faible ->
|
||||
modèle libre (variance élevée, risque d'overfitting) ; alpha élevé -> coefficients contraints (biais
|
||||
élevé, risque d'underfitting). C'est le levier de régularisation pour améliorer la généralisation face
|
||||
au bruit.
|
||||
|
||||
**3.7 - Comment utiliser un dataset de test dans le choix d'un hyperparamètre ?**
|
||||
On ne choisit **jamais** alpha sur le test. On règle alpha sur la **validation** (comparaison des alphas),
|
||||
puis on mesure **une seule fois** le modèle retenu sur le **test** (jamais vu) pour une estimation non
|
||||
biaisée de la généralisation. Utiliser le test pour régler alpha revient à le « fuiter » -> estimation
|
||||
trop optimiste.
|
||||
|
||||
---
|
||||
|
||||
## Partie 4 - Réflexion en production
|
||||
|
||||
**4.1 - Sur quelles données entraîner ?**
|
||||
Sur une **fenêtre glissante de données récentes** représentatives du régime actuel (ici `recent_history`
|
||||
l'emporte), en réintégrant régulièrement les dernières observations, tout en gardant assez d'historique
|
||||
pour capter les saisonnalités utiles (semaine, éventuellement année si stable). Compromis récence/volume.
|
||||
|
||||
**4.2 - À quelle fréquence réentraîner ?**
|
||||
Réentraînement **périodique** (hebdomadaire à mensuel selon le coût et la vitesse de dérive) **et**
|
||||
déclenché **par événement** (dégradation des métriques ou détection de drift). La forte saisonnalité de
|
||||
la conso justifie au minimum un rythme régulier accompagné d'une surveillance.
|
||||
|
||||
**4.3 - Quels indicateurs signalent un nouvel entraînement ?**
|
||||
- **Dégradation des métriques** en production (RMSE/MAE qui remontent vs baseline).
|
||||
- **Data drift** : la distribution des features d'entrée s'éloigne de celle d'entraînement.
|
||||
- **Concept drift** : la relation features -> cible change (nouveaux usages, réglementation, météo extrême).
|
||||
- Écart croissant entre distributions **train vs production** (surveillance type Evidently/Grafana).
|
||||
|
||||
---
|
||||
|
||||
## Livrables - synthèse finale
|
||||
|
||||
- **Stratégies de features testées** : `short_memory`, `seasonality`, `tendency`, `mixed`, `full`.
|
||||
- **Résultats** : voir matrices ci-dessus (validation RMSE/MAE).
|
||||
- **Meilleur run MLflow** : features `full` sur le split `recent_history` (RMSE **6.576**, MAE **3.677**) ;
|
||||
en linéaire comme en Ridge `alpha` ≤ 1e3 (résultats indistinguables ; run Ridge `alpha=1e3`
|
||||
`id 57cac1f5...`).
|
||||
- **Impact du changement de split** : `recent_history` améliore **toutes** les stratégies (~ -12 % RMSE),
|
||||
résultat à nuancer (fenêtres de validation différentes -> proximité temporelle avantageuse).
|
||||
- **Stratégie recommandée en production** : features **`mixed`** (parcimonie : quasi identique à `full`
|
||||
mais 4 features au lieu de 6, `lag_30d`/`lag_365d` contribuant peu), entraînée sur une **fenêtre de
|
||||
données récentes**, avec **Ridge `alpha` modéré (1 à 1e3)**.
|
||||
- **Raisons** : meilleure performance en validation, modèle plus simple donc plus robuste et
|
||||
maintenable, données récentes = moins de dérive, régularisation légère = filet de sécurité contre le
|
||||
bruit sans coût de performance.
|
||||
181
SYNTHESE_TP03.md
Normal file
181
SYNTHESE_TP03.md
Normal file
@@ -0,0 +1,181 @@
|
||||
# TP03 - Synthèse : exposer un modèle ML via une API REST
|
||||
|
||||
Fil rouge : prédiction de la consommation électrique (kWh). On sélectionne le meilleur modèle
|
||||
du TP02, on l'enregistre dans le **MLflow Model Registry**, on le **promeut via un alias**, puis
|
||||
on l'expose par une **API REST FastAPI**.
|
||||
|
||||
## Ce qui a été construit
|
||||
|
||||
- **`log_model`** ajouté à chaque run d'entraînement (`lab/modeling/cli.py`, `lab/modeling_ridge/cli.py`) :
|
||||
sauvegarde l'artefact complet du modèle sur S3 (`s3://mlflow-artifacts`).
|
||||
- **Modèle enregistré** : `electricity-consumption` (Registry MLflow). Versions créées avec le flag
|
||||
`--register` : v1 = `full` (6 features, RMSE val 6.576), v2 = `mixed` (4 features, RMSE 6.671).
|
||||
- **Promotion** : alias `champion` -> **v1** (`full`), via `lab/registry/cli.py promote`.
|
||||
- **API FastAPI** (`lab/serving/`) : `GET /health`, `POST /predict`, `POST /predict/batch`, Swagger `/docs`.
|
||||
Le modèle est chargé au démarrage par `models:/electricity-consumption@champion`.
|
||||
- **Exposition** : `https://api.192-168-122-143.nip.io/docs` (reverse-proxy Caddy -> uvicorn `:8000`).
|
||||
|
||||
---
|
||||
|
||||
## Partie 1 - Sélection et enregistrement du modèle
|
||||
|
||||
**1.1 - Que contiennent les artefacts du modèle sauvegardé par MLflow, pourquoi sont-ils utiles ?**
|
||||
Sauvegarder un modèle ne se limite pas à ses coefficients : MLflow enregistre tout l'environnement
|
||||
d'exécution. Contenu observé (`s3://mlflow-artifacts/1/models/<id>/artifacts/`) :
|
||||
- `model.pkl` : le modèle sérialisé (poids/coefficients) ;
|
||||
- `MLmodel` : métadonnées (flavors `sklearn`/`pyfunc`, **signature** = schéma entrées/sorties) ;
|
||||
- `requirements.txt`, `conda.yaml`, `python_env.yaml` : versions exactes des dépendances ;
|
||||
- `input_example.json`, `serving_input_example.json` : exemple d'entrée.
|
||||
Utiles pour **recharger le modèle partout** (`load_model`), **reproduire l'environnement** (mêmes
|
||||
versions -> mêmes prédictions), et connaître le **contrat d'E/S** (signature).
|
||||
|
||||
**1.2 - Quel mécanisme vous permet de promouvoir un modèle ?**
|
||||
L'**alias** du Model Registry : `MlflowClient().set_registered_model_alias(name, "champion", version)`
|
||||
fait pointer un alias mobile vers une version précise. (Les anciens *stages* Staging/Production sont
|
||||
dépréciés en MLflow 3.x au profit des alias + tags.)
|
||||
|
||||
**1.3 - Plusieurs environnements de production (un par région), plusieurs modèles : comment les identifier ?**
|
||||
Avec des **alias** et **tags** spécifiques : par ex. des alias `production-eu`, `production-us`,
|
||||
`champion-north`... sur un même modèle enregistré, et/ou un modèle enregistré par région, complétés par
|
||||
des **tags** (région, environnement) portés par le modèle ou la version. Le Registry gère plusieurs
|
||||
alias par modèle et des tags arbitraires : l'API cible alors `models:/<nom>@<alias-région>`.
|
||||
|
||||
---
|
||||
|
||||
## Partie 2 - Service de prédiction
|
||||
|
||||
### Étape 1 - Création de l'API
|
||||
|
||||
**2.1 - Quel endpoint de vérification ? Quel code de statut attendu ?**
|
||||
`GET /health`, qui renvoie **200 OK** avec `{"status":"ok"}`.
|
||||
|
||||
**2.2 - À quoi sert un endpoint de santé ?**
|
||||
À vérifier que le service est vivant et prêt (*liveness/readiness*), sans exécuter de vraie prédiction.
|
||||
Il est sondé par l'orchestrateur / load-balancer / monitoring pour router le trafic, redémarrer un
|
||||
conteneur en échec, ou alerter.
|
||||
|
||||
### Étape 2 - Requête de prédiction
|
||||
|
||||
**2.3 - Quelles informations le client doit-il fournir ? Pourquoi ?**
|
||||
L'**identifiant client** (`client_id`) et, en paramètre, la **date** de prédiction. Ce sont les seules
|
||||
informations du **contexte métier** que le client connaît : elles disent *qui* et *quand*. Elles servent
|
||||
de clé pour retrouver le reste côté serveur.
|
||||
|
||||
**2.4 - Quelles informations le client ne peut-il pas fournir ? Comment le service les récupère-t-il ?**
|
||||
Les **features calculées** (lags, moyennes glissantes) issues de l'historique de consommation : le
|
||||
client ne les possède/calcule pas. Le service les **récupère lui-même** depuis un feature store / une
|
||||
base, à partir du `client_id` (+ date). Ici, c'est **simulé par un dictionnaire Python** (`lab/serving/features.py`).
|
||||
|
||||
**2.5 - Quelle méthode HTTP pour le endpoint de prédiction ? Pourquoi ?**
|
||||
**POST** : la requête transporte un **corps JSON structuré** (et potentiellement volumineux en batch),
|
||||
et déclenche un **calcul** (action, non une simple lecture cacheable de ressource comme le ferait GET).
|
||||
|
||||
### Étape 3 - Récupération des features
|
||||
|
||||
**2.6 - Rappel des features nécessaires à la prédiction.**
|
||||
Modèle promu = stratégie `full`, soit les **6 features** :
|
||||
`lag_1d`, `lag_7d`, `lag_30d`, `lag_365d`, `rolling_mean_7d`, `rolling_mean_30d`.
|
||||
|
||||
**2.7 - Dans un système réel, d'où proviennent ces features ?**
|
||||
D'un **feature store / pipeline de features** : un job (batch ou streaming) calcule lags et moyennes
|
||||
glissantes depuis la série temporelle brute, les stocke dans une base (ex. Feast), et les sert à
|
||||
l'inférence. Elles doivent être calculées **de façon identique à l'entraînement** (éviter le
|
||||
*training/serving skew*).
|
||||
|
||||
**2.8 - Pourquoi séparer récupération des features et calcul de la prédiction ?**
|
||||
**Séparation des responsabilités** : la source des features peut évoluer (BDD, cache, feature store)
|
||||
sans toucher au modèle ; le modèle reste une **fonction pure** `features -> prédiction`, testable et
|
||||
réutilisable. « L'application du modèle n'est qu'une étape de la chaîne de prédiction. »
|
||||
|
||||
**2.9 (Bonus) - Si les features ne peuvent pas être calculées/récupérées ?**
|
||||
L'API ne doit pas planter : elle renvoie une erreur explicite. Client/features introuvables -> **404** ;
|
||||
payload invalide -> 422 ; feature store indisponible (panne transitoire) -> **503**. Toujours un JSON
|
||||
d'erreur clair.
|
||||
|
||||
### Étape 4 - Chargement du modèle depuis le Registry
|
||||
|
||||
**2.10 - Pourquoi un alias plutôt qu'un numéro de version ?**
|
||||
L'alias (`champion`) est **stable et mobile** : l'API charge toujours `models:/electricity-consumption@champion`
|
||||
et l'on **re-pointe l'alias** vers une nouvelle version sans modifier ni redéployer le code. Un numéro
|
||||
de version est **figé** : chaque changement de modèle imposerait d'éditer la config et de redéployer.
|
||||
L'alias **découple** « quel modèle est en prod » (décision côté MLflow) du code de service.
|
||||
|
||||
**2.11 - Avantages du Registry par rapport à un simple fichier modèle ?**
|
||||
Versioning centralisé (historique de toutes les versions), **alias/tags**, **lignée** vers le run
|
||||
d'entraînement (params/métriques), chargement par URI depuis n'importe où, workflow de promotion,
|
||||
traçabilité/audit, environnement (requirements) attaché. Un simple fichier n'offre rien de tout cela
|
||||
(pas d'historique, pas de métadonnées, distribution manuelle et fragile).
|
||||
|
||||
**2.12 (Bonus) - Si les requirements du modèle sont incohérents avec l'environnement de l'API ?**
|
||||
Risque : incompatibilité de versions (scikit-learn, numpy) -> erreur de désérialisation ou écarts
|
||||
numériques silencieux. **Architecture** : **isoler le modèle dans son propre runtime** construit à
|
||||
partir de son `requirements.txt` (image conteneur dédiée par modèle, ex. `mlflow models build-docker` /
|
||||
MLflow serving), l'API l'appelant via HTTP ; ou figer l'environnement de l'API depuis les requirements
|
||||
du modèle. On **découple** l'API des dépendances du modèle.
|
||||
|
||||
### Étape 5 - Endpoint de prédiction
|
||||
|
||||
**2.13 - Étapes lorsqu'une requête de prédiction arrive.**
|
||||
1. Valider le payload (schéma Pydantic). 2. Récupérer les features du client (404 si inconnu).
|
||||
3. Assembler le vecteur de features **dans l'ordre attendu** par le modèle. 4. `model.predict`.
|
||||
5. Formater et renvoyer la réponse JSON (200).
|
||||
|
||||
**2.14 - Quel format de réponse ? Quel statut HTTP ?**
|
||||
JSON : `{client_id, prediction_kwh, model_name, model_version}`, statut **200 OK**.
|
||||
|
||||
---
|
||||
|
||||
## Partie 3 (Bonus) - Gestion des erreurs
|
||||
|
||||
**3.1 - L'application doit-elle échouer ou intercepter cette erreur ?**
|
||||
**Intercepter.** Un `client_id` inconnu est une **erreur cliente** (mauvaise entrée), pas un bug
|
||||
serveur : le service reste debout et renvoie une réponse d'erreur propre.
|
||||
|
||||
**3.2 - Quel code HTTP est adapté ?**
|
||||
**404 Not Found** (la ressource/le client demandé n'existe pas). (422 si le payload lui-même est
|
||||
malformé.) Implémenté via `HTTPException(status_code=404, ...)`.
|
||||
|
||||
---
|
||||
|
||||
## Partie 4 (Bonus) - Prédictions en batch
|
||||
|
||||
**4.1 - Un modèle doit-il forcément être exposé par API ? Dans quel cas ?**
|
||||
Non. L'API (temps réel) convient aux prédictions **à la demande, individuelles, à faible latence**
|
||||
(appli interactive). Pour de **gros volumes calculés périodiquement** (ex. tous les clients chaque
|
||||
nuit), l'**inférence batch** (job planifié qui écrit les résultats en base) est plus adaptée et moins
|
||||
coûteuse. On expose par API quand on a besoin de prédictions fraîches, unitaires et synchrones.
|
||||
|
||||
**4.2 - Quelles briques restent identiques entre batch et temps réel ?**
|
||||
Le **modèle** (même artefact du Registry), la **logique/définition des features**, le **préprocessing**,
|
||||
le **code de prédiction** (`features -> prédiction`). Ce qui diffère : le **déclencheur/orchestration**
|
||||
(requête HTTP vs job planifié), les **E/S** (un JSON unitaire vs une table en masse) et le profil
|
||||
latence/débit.
|
||||
|
||||
**4.3 - Comment récupérer les informations de prédiction ?**
|
||||
En **masse** : lire les features de tous les clients pour la date depuis le feature store / une table
|
||||
(BDD ou parquet), appeler `model.predict` sur le **lot entier** (vectorisé), puis écrire les résultats
|
||||
en base/fichier. Ici, `POST /predict/batch` prend une liste de `client_ids` et renvoie la liste des
|
||||
prédictions (+ `unknown_client_ids` pour les clients ignorés).
|
||||
|
||||
---
|
||||
|
||||
## Reproduire
|
||||
|
||||
Sur la VM, depuis `/home/user/tp` :
|
||||
|
||||
```bash
|
||||
set -a; source .env; set +a # MLflow + creds S3 (Garage)
|
||||
|
||||
# Partie 1 : enregistrer + promouvoir
|
||||
python -m lab.modeling.cli full --register # v1 (champion)
|
||||
python -m lab.modeling.cli mixed --register # v2 (comparaison)
|
||||
python -m lab.registry.cli versions # consulter les versions/alias
|
||||
python -m lab.registry.cli promote --version 1 --alias champion
|
||||
|
||||
# Partie 2-4 : servir l'API
|
||||
./serve.sh # uvicorn :8000 (Swagger /docs)
|
||||
curl -s localhost:8000/health
|
||||
curl -s -X POST localhost:8000/predict -H 'content-type: application/json' -d '{"client_id":"MT_124"}'
|
||||
```
|
||||
|
||||
Depuis le poste (certificat ENI de confiance) : **https://api.192-168-122-143.nip.io/docs**.
|
||||
160
SYNTHESE_TP04.md
Normal file
160
SYNTHESE_TP04.md
Normal file
@@ -0,0 +1,160 @@
|
||||
# TP04 - Synthèse : automatisation d'un pipeline ML avec CI/CD
|
||||
|
||||
Fil rouge : prédiction de la consommation électrique (kWh). Aux TP02/TP03, chaque étape (split,
|
||||
entraînement, `log_model`, enregistrement + promotion au Model Registry) était lancée **à la main**.
|
||||
Le TP04 **automatise** cet enchaînement en un pipeline reproductible et idempotent (`pipeline.py`),
|
||||
puis le **déclenche automatiquement via Forgejo Actions** à chaque Pull Request vers `main`.
|
||||
|
||||
## Ce qui a été construit
|
||||
|
||||
- **`pipeline.py`** (racine du dépôt) : orchestre en une seule commande `split -> train --register ->
|
||||
promote`. Chaque étape est un sous-processus `subprocess.run(..., check=True)` : un échec stoppe le
|
||||
pipeline avec un code de sortie non nul. Réutilise **telles quelles** les CLI du package `lab/`
|
||||
(mêmes commandes qu'en manuel) ; ne fait qu'enchaîner et calculer le numéro de version à promouvoir.
|
||||
- **`requirements.txt`** : dépendances runtime figées sur l'environnement de référence de la VM
|
||||
(venv `/opt/venvs/mlops`, Python 3.14 ; mlflow 3.13.0, scikit-learn 1.9.0, pandas, pyarrow, boto3, s3fs).
|
||||
- **`.forgejo/workflows/pipeline.yml`** : workflow CI déclenché sur `pull_request` vers `main`, qui
|
||||
installe les dépendances puis exécute `pipeline.py`.
|
||||
- **Repo + CI sur le Forgejo de la VM** (`forgejo.192-168-122-143.nip.io/trainer-admin/ENI-ml-mlops`) :
|
||||
les 2 runners existants (mode docker, réseau `mlops-net`) exécutent le job. Le code de référence reste
|
||||
archivé publiquement sur `git.lidge.fr` (comme aux TP02/TP03).
|
||||
|
||||
---
|
||||
|
||||
## Partie 1 - Analyse de la pipeline existante
|
||||
|
||||
**1.1 - Quelles étapes sont nécessaires pour produire un nouveau modèle ?**
|
||||
1. **Préparation des données** (`lab/split/cli.py`) : lire le dataset source `/data/modelling/{features,
|
||||
target}.parquet` et le découper en `train` / `validation` / `test` selon la stratégie de split figée.
|
||||
2. **Entraînement** (`lab/modeling/cli.py`) : entraîner sur `train`, évaluer sur `validation`, logguer
|
||||
paramètres et métriques dans MLflow, sauvegarder l'artefact du modèle (`log_model`).
|
||||
3. **Enregistrement au Registry** (option `--register`) : empiler une **nouvelle version** du modèle
|
||||
`electricity-consumption`.
|
||||
4. **Promotion** (`lab/registry/cli.py`) : (re)pointer l'alias `champion` vers cette nouvelle version
|
||||
(mise à disposition pour l'API de service).
|
||||
|
||||
**1.2 - Quelle étape dépend directement du résultat de la précédente ?**
|
||||
Le pipeline est une **chaîne linéaire** : `train` dépend de `split` (il lit `train.parquet` et
|
||||
`validation.parquet` produits par le split) ; l'`enregistrement`/`log_model` dépend de `train` (il
|
||||
sauvegarde le modèle entraîné) ; la `promotion` dépend de l'enregistrement (elle a besoin du **numéro
|
||||
de version** qui vient d'être créé). Chaque étape consomme la sortie de la précédente
|
||||
(parquets -> modèle -> version -> alias).
|
||||
|
||||
**1.3 - Quels artefacts sont produits tout au long du pipeline ?**
|
||||
- Les jeux de données `data/{train,validation,test}.parquet` (versionnables par DVC) ;
|
||||
- un **run MLflow** avec ses paramètres (`strategy`, `split_strategy`, `features`) et métriques
|
||||
(`train/validation_rmse`, `train/validation_mae`, coefficients) ;
|
||||
- l'**artefact du modèle** sur `s3://mlflow-artifacts` (`model.pkl`, `MLmodel` + signature,
|
||||
`requirements.txt`/`conda.yaml`/`python_env.yaml`, `input_example`) ;
|
||||
- une **version** au Model Registry (`electricity-consumption` vN) ;
|
||||
- l'**alias** `champion` pointant vers cette version.
|
||||
|
||||
---
|
||||
|
||||
## Partie 2 - Construction d'une pipeline automatisé
|
||||
|
||||
**2.1 - Quels avantages par rapport à une exécution manuelle ?**
|
||||
Reproductibilité (mêmes étapes, même ordre, mêmes paramètres à chaque exécution), suppression des
|
||||
erreurs humaines (étape oubliée, mauvais ordre, paramètre incohérent), rapidité (**une seule
|
||||
commande**), traçabilité, et surtout **exécutabilité par une machine** : le pipeline devient
|
||||
déclenchable par une plateforme de CI/CD. Le script documente aussi le processus de façon vivante.
|
||||
|
||||
**2.2 - Que se passe-t-il si une étape échoue ?**
|
||||
Le pipeline **s'arrête immédiatement** : chaque étape est lancée avec `check=True`, donc un échec lève
|
||||
une exception et le script sort avec un **code non nul**. Les étapes suivantes ne sont pas exécutées
|
||||
(on n'enregistre pas un modèle à partir d'un split incomplet). En CI, ce code non nul fait **échouer le
|
||||
job** (signal rouge visible). On corrige la cause, puis on relance.
|
||||
|
||||
**2.3 - Pourquoi est-il important qu'un pipeline soit ré-exécutable sans effets de bord (idempotence) ?**
|
||||
Pour pouvoir le **rejouer en confiance** (reprise après échec, ré-entraînement périodique, exécution en
|
||||
CI) sans corrompre l'état ni accumuler d'effets parasites. Ici : le **split est déterministe** (mêmes
|
||||
dates -> mêmes fichiers, réécrits proprement) ; `--register` **empile une nouvelle version** proprement
|
||||
(le versioning est le comportement attendu, pas un doublon anarchique) ; l'alias `champion` est
|
||||
**repointé** (mobile), jamais dupliqué. Rejouer le pipeline redonne donc un état cohérent et prévisible.
|
||||
|
||||
---
|
||||
|
||||
## Partie 3 - Exécution de la pipeline
|
||||
|
||||
Exécution de référence sur la VM (`set -a; source .env; set +a; python pipeline.py`).
|
||||
|
||||
**3.1 - Quels artefacts ont été produits ?**
|
||||
- `data/train.parquet` (~115 Mo, 4 485 120 lignes), `validation.parquet` (~48 Mo, 1 855 488),
|
||||
`test.parquet` (~68 Mo, 2 629 632) ;
|
||||
- un nouveau run dans l'expérience MLflow `tp02_electricity_consumption` ;
|
||||
- une **nouvelle version** de `electricity-consumption` (stratégie `full`, `validation_rmse = 6.576`) ;
|
||||
- l'alias `champion` repointé vers cette nouvelle version.
|
||||
|
||||
**3.2 - Comment vérifier qu'un nouveau modèle a bien été enregistré dans MLflow ?**
|
||||
- En CLI : `python -m lab.registry.cli versions` (liste les versions, leurs alias et le
|
||||
`validation_rmse`) ;
|
||||
- dans l'**UI MLflow** (`https://mlflow.192-168-122-143.nip.io` -> Models -> `electricity-consumption`) :
|
||||
la nouvelle version apparaît, l'alias `champion` pointe dessus ;
|
||||
- par l'API : `MlflowClient().search_model_versions(...)` / `get_model_version_by_alias(name, "champion")`.
|
||||
|
||||
**3.3 - Quels éléments garantissent la reproductibilité ?**
|
||||
- Le **code** versionné (git) et le **workflow** figé ;
|
||||
- la **configuration** figée dans `lab/constants.py` (`CHOSEN_SPLIT_STRATEGY`, listes de features,
|
||||
`SERVING_STRATEGY`) et le **split déterministe** par dates ;
|
||||
- les **données** versionnées par DVC (remote S3 Garage) ;
|
||||
- les **dépendances épinglées** (`requirements.txt`) + l'environnement capturé par MLflow dans
|
||||
l'artefact (`requirements.txt`/`conda.yaml`) ;
|
||||
- le **tracking MLflow** (paramètres, métriques, artefacts) qui permet de retrouver exactement chaque run.
|
||||
|
||||
---
|
||||
|
||||
## Partie 4 - Automatisation avec Forgejo Actions
|
||||
|
||||
**4.1 - Quel événement déclenche le workflow ?**
|
||||
Une **Pull Request vers la branche principale** : `on: pull_request: branches: [main]`. L'ouverture
|
||||
(et chaque mise à jour) d'une PR ciblant `main` lance le workflow. (Démonstration réalisée : PR #1
|
||||
`tp04-ci-pipeline -> main`, workflow **vert**, une nouvelle version du modèle enregistrée par la CI.)
|
||||
|
||||
**4.2 - Quel est le rôle de Forgejo Actions dans cette architecture ?**
|
||||
C'est la **plateforme de CI/CD intégrée au forge**. Elle détecte l'événement (PR), planifie un job et
|
||||
l'assigne à un **runner**, fournit l'environnement d'exécution (conteneur), **injecte les secrets**,
|
||||
exécute les étapes (récupération du code, installation des dépendances, exécution de `pipeline.py`) et
|
||||
**rapporte le statut** (vert/rouge) rattaché à la PR. C'est l'orchestrateur qui automatise l'exécution
|
||||
du pipeline à chaque changement proposé, sert de garde-fou (une PR qui casse le pipeline est visible
|
||||
avant le merge) et centralise la traçabilité.
|
||||
|
||||
**4.3 - Différence entre le pipeline Python (partie 2) et le workflow Forgejo Actions ?**
|
||||
`pipeline.py` porte le **quoi** : la logique métier (l'enchaînement `split -> train -> register ->
|
||||
promote`), exécutable partout (poste, VM, CI). Le workflow porte le **quand** et le **où** :
|
||||
l'orchestration CI (quel événement déclenche, sur quel runner, dans quel environnement/conteneur, avec
|
||||
quels secrets) et le reporting. Le workflow ne contient **aucune logique ML** : il prépare
|
||||
l'environnement et appelle `pipeline.py`. Cette séparation des responsabilités permet de changer le
|
||||
déclencheur ou l'infrastructure sans toucher à la logique, et inversement.
|
||||
|
||||
**4.4 - Quels autres événements pourraient déclencher automatiquement le pipeline ?**
|
||||
Un **push** sur `main` (après merge), la création d'un **tag/release**, une **planification**
|
||||
(`schedule`/cron, pour un ré-entraînement périodique), un **déclenchement manuel** (`workflow_dispatch`),
|
||||
ou un **événement externe** (webhook signalant l'arrivée de nouvelles données, appel depuis un autre
|
||||
workflow).
|
||||
|
||||
---
|
||||
|
||||
## Détails d'implémentation (Forgejo Actions sur la VM)
|
||||
|
||||
- Le job tourne dans un conteneur `node:20-bookworm` sur le réseau `mlops-net` : il résout les services
|
||||
par leur nom et joint MLflow/Garage via Caddy en HTTPS (`https://mlflow...`, `https://garage...`),
|
||||
hôtes autorisés par `MLFLOW_SERVER_ALLOWED_HOSTS`.
|
||||
- Deux volumes montés en lecture seule : les **données source** `/data/modelling` (le split en a besoin)
|
||||
et le **magasin de certificats** du host `/etc/ssl/certs/ca-certificates.crt` (confiance envers la CA
|
||||
interne ENI MLOps ; les runners ont été autorisés à monter ces chemins via `valid_volumes`).
|
||||
- Les identifiants sont fournis en **secrets Actions** du repo (MLflow basic auth + clés S3 Garage),
|
||||
jamais committés (`.env` reste gitignoré). Python et les dépendances sont installés dans le job
|
||||
(`python3-venv` + `pip install -r requirements.txt`).
|
||||
|
||||
## Reproduire
|
||||
|
||||
Sur la VM, `~/tp` :
|
||||
|
||||
```bash
|
||||
set -a; source .env; set +a
|
||||
python pipeline.py # split -> train full --register -> promote champion
|
||||
python -m lab.registry.cli versions # verifier la nouvelle version + alias champion
|
||||
```
|
||||
|
||||
En CI : ouvrir une Pull Request vers `main` sur `forgejo.192-168-122-143.nip.io/trainer-admin/ENI-ml-mlops`
|
||||
-> le workflow `pipeline` s'exécute et enregistre une nouvelle version du modèle (onglet **Actions**).
|
||||
178
SYNTHESE_TP05.md
Normal file
178
SYNTHESE_TP05.md
Normal file
@@ -0,0 +1,178 @@
|
||||
# TP05 - Synthèse : monitoring, dérive et redéploiement
|
||||
|
||||
Fil rouge : prédiction de la consommation électrique (kWh, 128 clients portugais, pas de 15 min).
|
||||
|
||||
Scénario : un modèle a été **entraîné sur 2011** et tourne toujours en production plusieurs années plus
|
||||
tard. En ingénieur MLOps, on vérifie si les données **2014** observées en production ont **dérivé** par
|
||||
rapport à l'entraînement, on analyse cette dérive avec **Evidently**, puis on mesure son **impact sur les
|
||||
performances** (2012 -> 2014) pour décider d'une stratégie de maintenance.
|
||||
|
||||
## Ce qui a été construit
|
||||
|
||||
- **`tp_module5_monitoring_derive.ipynb`** (notebook exécuté) : Partie 1 (moyennes + histogrammes de
|
||||
`lag_30d` 2011 vs 2014), Partie 2 (rapport de dérive Evidently `DataDriftPreset`), Partie 3 (modèle
|
||||
linéaire entraîné sur 2011, évalué sur 2011-2014). Réutilise `lab/constants.py` (chemins, cible, features).
|
||||
- **`tp05_evidently_drift_2011_vs_2014.html`** : rapport Evidently interactif complet (à ouvrir dans un
|
||||
navigateur).
|
||||
- **Choix technique** : la feature `lag_365d` (consommation 365 jours plus tôt) est **exclue** car elle
|
||||
est indéfinie (NaN) sur **toute** l'année 2011 (première année, pas d'historique 2010). On travaille donc
|
||||
sur les 5 features définies sur les deux périodes : `lag_1d, lag_7d, lag_30d, rolling_mean_7d,
|
||||
rolling_mean_30d`. Cet ensemble sert à la fois au rapport de dérive et au modèle entraîné sur 2011.
|
||||
|
||||
---
|
||||
|
||||
## Partie 1 - Observer une dérive
|
||||
|
||||
Feature `lag_30d`, sur les années complètes 2011 (train) et 2014 (production) :
|
||||
|
||||
| Année | Moyenne `lag_30d` | Effectif |
|
||||
|--------------|-------------------|------------|
|
||||
| 2011 (train) | **62.041 kWh** | 4 116 352 |
|
||||
| 2014 (prod) | **54.054 kWh** | 4 485 120 |
|
||||
| Écart | **-7.986 kWh** | **-12.9 %** |
|
||||
|
||||
**1.1 - Les deux distributions semblent-elles similaires ?**
|
||||
Non. Elles gardent la même **forme** générale (distribution asymétrique étalée vers la droite : beaucoup de
|
||||
petites consommations, une longue queue de fortes valeurs), mais l'histogramme 2014 est **décalé vers les
|
||||
valeurs plus faibles** et sa moyenne est nettement plus basse. Visuellement, le décalage est net.
|
||||
|
||||
**1.2 - La moyenne a-t-elle évolué entre 2011 et 2014 ?**
|
||||
Oui, franchement : de **62.04 kWh** à **54.05 kWh**, soit **-7.99 kWh (-12.9 %)**. La consommation moyenne
|
||||
(observée via `lag_30d`) a **baissé d'environ 13 %** entre l'entraînement et la production.
|
||||
|
||||
**1.3 - Cette évolution paraît-elle suffisamment importante pour parler de dérive ?**
|
||||
Une baisse de ~13 % de la moyenne est un signal **fort et cohérent** d'un changement de distribution : c'est
|
||||
un indice sérieux de dérive des données. Attention toutefois : la moyenne seule ne **prouve** pas une dérive
|
||||
statistique (deux distributions peuvent avoir des moyennes proches et des formes différentes, ou l'inverse).
|
||||
Il faut confirmer avec une analyse sur toute la distribution, avec un test et un seuil objectifs (Partie 2).
|
||||
|
||||
**1.4 - Cette première analyse est-elle suffisante pour conclure sur l'état du modèle ? Pourquoi ?**
|
||||
Non. Elle ne porte que sur **une** feature (`lag_30d`), via un **seul** indicateur (moyenne + comparaison
|
||||
visuelle), **sans test statistique ni seuil**, et elle ne dit **rien** des autres features ni de l'**impact
|
||||
réel sur la performance** du modèle. C'est un premier indice, pas une conclusion. Il faut (a) un outil de
|
||||
monitoring qui teste toutes les features (Partie 2) et (b) mesurer la performance dans le temps (Partie 3).
|
||||
|
||||
---
|
||||
|
||||
## Partie 2 - Analyse avec un outil de monitoring (Evidently)
|
||||
|
||||
Rapport `DataDriftPreset` comparant 2011 (référence) et 2014 (courant), échantillon de 200 000 lignes par
|
||||
période (seed fixe). Test choisi automatiquement par Evidently pour de grands échantillons numériques :
|
||||
**distance de Wasserstein normalisée**, seuil **0.1**.
|
||||
|
||||
| Feature | Test | Drift score | Seuil | Dérive ? |
|
||||
|--------------------|-------------------------------|-------------|-------|----------|
|
||||
| lag_1d | Wasserstein distance (normed) | 0.165 | 0.1 | Oui |
|
||||
| lag_7d | Wasserstein distance (normed) | 0.167 | 0.1 | Oui |
|
||||
| lag_30d | Wasserstein distance (normed) | 0.167 | 0.1 | Oui |
|
||||
| rolling_mean_7d | Wasserstein distance (normed) | 0.188 | 0.1 | Oui |
|
||||
| rolling_mean_30d | Wasserstein distance (normed) | 0.191 | 0.1 | Oui |
|
||||
|
||||
Verdict global : **`dataset_drift = True`**, **5 / 5 features en dérive (100 %)**.
|
||||
|
||||
**2.1 - Combien de features présentent une dérive ?**
|
||||
**Les 5** features analysées présentent une dérive (5 / 5, soit 100 %). Toutes dépassent le seuil de 0.1.
|
||||
|
||||
**2.2 - Toutes les features évoluent-elles de la même manière ?**
|
||||
Toutes dérivent, mais **pas avec la même intensité**. Les distances de Wasserstein vont de **0.165**
|
||||
(`lag_1d`) à **0.191** (`rolling_mean_30d`). Les **moyennes glissantes** (`rolling_mean_7d` 0.188,
|
||||
`rolling_mean_30d` 0.191) dérivent un peu **plus** que les **lags bruts** (0.165-0.167) : les indicateurs
|
||||
lissés/tendanciels captent davantage la baisse durable du niveau de consommation.
|
||||
|
||||
**2.3 - Identifiez deux métriques ou indicateurs du rapport. À quoi servent-ils ?**
|
||||
- **Le drift score par feature** (ici la distance de Wasserstein normalisée) : il **quantifie** l'écart
|
||||
entre la distribution de référence (2011) et la distribution courante (2014) pour chaque variable, et le
|
||||
compare à un **seuil** (0.1). Il sert à **détecter et localiser** la dérive, feature par feature.
|
||||
- **Le verdict global "Dataset Drift" + la part de colonnes en dérive** : il **agrège** les décisions par
|
||||
feature (nombre / pourcentage de colonnes dérivées, ici 5/5 = 100 %) en une conclusion **au niveau du
|
||||
dataset entier**. Il sert à **trancher globalement** (le dataset a-t-il dérivé, oui/non).
|
||||
(Le rapport affiche aussi, par feature, les **histogrammes de distribution référence vs courant**, utiles
|
||||
pour visualiser la nature du décalage.)
|
||||
|
||||
**2.4 - Le rapport conclut-il à une dérive globale du dataset ? Justifiez.**
|
||||
**Oui.** `dataset_drift = True`. Le critère par défaut d'Evidently (déclarer une dérive du dataset quand la
|
||||
**part de colonnes en dérive dépasse 50 %**) est **largement** franchi : **100 %** des colonnes dérivent.
|
||||
Le rapport conclut donc sans ambiguïté à une dérive globale des données d'entrée.
|
||||
|
||||
---
|
||||
|
||||
## Partie 3 - Mesurer l'impact sur les performances
|
||||
|
||||
Modèle (régression linéaire) entraîné **sur 2011 uniquement**, évalué sur chaque année complète :
|
||||
|
||||
| Année | Rôle | RMSE (kWh) | MAE (kWh) |
|
||||
|-------|-------------------|------------|-----------|
|
||||
| 2011 | train (référence) | 7.904 | 4.450 |
|
||||
| 2012 | évaluation (prod) | 7.595 | 4.178 |
|
||||
| 2013 | évaluation (prod) | 7.486 | 4.095 |
|
||||
| 2014 | évaluation (prod) | **7.101** | **3.911** |
|
||||
|
||||
**Résultat marquant : l'erreur ne se dégrade pas, elle diminue légèrement d'année en année.**
|
||||
|
||||
**3.1 - Peut-on comparer 2011 vs 2012 au même titre que 2012 vs 2013 ?**
|
||||
Non. **2011 est l'année d'entraînement** : l'erreur y est mesurée **in-sample** (le modèle a déjà vu ces
|
||||
données), donc **optimiste**. Comparer 2011 (in-sample) à 2012 (hors échantillon) mélange deux régimes
|
||||
différents, alors que 2012 vs 2013 compare **deux années hors échantillon** entre elles (comparaison
|
||||
homogène). De plus, le **niveau de consommation change** d'une année à l'autre : comparer des **RMSE bruts**
|
||||
est biaisé par l'échelle de la cible (voir 3.3). 2011 doit donc être traité comme **référence de train**,
|
||||
pas comme un point de comparaison équivalent aux années de production.
|
||||
|
||||
**3.2 - Observe-t-on une dégradation progressive ?**
|
||||
**Non.** En valeur absolue, la performance **s'améliore** légèrement : RMSE 7.90 -> 7.60 -> 7.49 -> 7.10 et
|
||||
MAE 4.45 -> 4.18 -> 4.10 -> 3.91 de 2011 à 2014. Aucune dégradation observée sur ces métriques.
|
||||
|
||||
**3.3 - Cette (absence de) dégradation est-elle cohérente avec la dérive détectée ?**
|
||||
À première vue c'est surprenant (dérive nette en Partie 2, mais pas de perte de performance), et c'est
|
||||
**instructif** : **une dérive des données ne signifie pas automatiquement une dégradation du modèle**.
|
||||
Explication : la consommation a **baissé** (~-13 % sur `lag_30d`, Partie 1), donc la **cible** (kWh) est plus
|
||||
petite en 2014 ; comme RMSE et MAE sont des erreurs **absolues** en kWh, elles **diminuent mécaniquement**
|
||||
quand l'échelle de la cible diminue. Pendant ce temps, la **relation features -> cible** (les lags prédisent
|
||||
la consommation courante) est restée **stable**. On a donc une dérive des données **sans** dérive de la
|
||||
relation. Pour comparer rigoureusement les années, il faudrait une métrique **indépendante de l'échelle**
|
||||
(MAPE, R², RMSE normalisé).
|
||||
|
||||
**3.4 - Peut-on conclure que le modèle 2011 est encore fiable en 2014 ?**
|
||||
Sur la base des erreurs absolues, il **ne s'est pas dégradé** (il fait même un peu mieux en kWh). Mais on ne
|
||||
peut pas **conclure** à sa pleine fiabilité sur ces seuls chiffres : (1) les métriques absolues sont
|
||||
**trompeuses** quand l'échelle de la cible change ; (2) on n'a mesuré que l'**amplitude** de l'erreur, pas un
|
||||
éventuel **biais systématique** (sur/sous-estimation) ni la performance par segment (client, saison).
|
||||
Verdict prudent : **pas de signe de dégradation**, mais à **confirmer** avec des métriques relatives et une
|
||||
analyse des résidus avant de déclarer le modèle fiable.
|
||||
|
||||
**3.5 - Que faudrait-il faire ensuite : surveiller, réentraîner, ou redéployer ?**
|
||||
**Surveiller.** La dérive des données est réelle, mais **sans impact mesuré** sur la performance : rien
|
||||
n'impose un réentraînement ou un redéploiement immédiat (ce serait du travail et du risque pour un gain non
|
||||
démontré). On met en place un **monitoring continu** (dérive des features + performance suivie avec une
|
||||
métrique relative + dérive des prédictions et de la cible dès que la vérité terrain arrive) et on **déclenche
|
||||
un réentraînement seulement si/quand** la performance se dégrade réellement ou qu'un concept drift apparaît.
|
||||
|
||||
**3.6 - Quel type de dérive semble en jeu ? Quelles analyses préconisez-vous ?**
|
||||
Type : **data drift** (dérive des covariables / *covariate shift*). Les distributions des features d'entrée
|
||||
ont changé (baisse durable du niveau de consommation), tandis que la **relation** features -> cible semble
|
||||
**stable** (pas de perte de performance) : **pas de concept drift évident**. Analyses préconisées :
|
||||
- suivre une **métrique de performance indépendante de l'échelle** (MAPE, R², nRMSE) dans le temps, plutôt
|
||||
que le RMSE brut ;
|
||||
- monitorer aussi la **dérive des prédictions** et de la **cible** (*target drift*), pas seulement des
|
||||
features en entrée ;
|
||||
- **analyser les résidus** (biais moyen, hétéroscédasticité) et **segmenter** (par client, par saison) ;
|
||||
- surtout, dès que la **vérité terrain** est disponible (avec délai), **recalculer la performance réelle en
|
||||
production** pour détecter précocement un éventuel **concept drift** : c'est la fermeture de la *feedback
|
||||
loop* MLOps (surveiller -> alerter -> réentraîner -> redéployer).
|
||||
|
||||
---
|
||||
|
||||
## Reproduire
|
||||
|
||||
Sur la VM, dans `~/tp` (venv `/opt/venvs/mlops`). Le TP05 est **100 % local** (données `/data/modelling`) :
|
||||
aucun secret réseau (MLflow, S3) n'est nécessaire.
|
||||
|
||||
```bash
|
||||
# Ré-exécuter le notebook (régénère les sorties + le rapport HTML)
|
||||
/opt/venvs/mlops/bin/jupyter nbconvert --to notebook --execute --inplace tp_module5_monitoring_derive.ipynb
|
||||
```
|
||||
|
||||
Consultation interactive (service `jupyter-tp`, token `mlops`) :
|
||||
|
||||
```
|
||||
https://jupyter.192-168-122-143.nip.io/lab/tree/tp_module5_monitoring_derive.ipynb?token=mlops
|
||||
```
|
||||
5
data/test.parquet.dvc
Normal file
5
data/test.parquet.dvc
Normal file
@@ -0,0 +1,5 @@
|
||||
outs:
|
||||
- md5: 6ecb52ceb1d8322a91454aca4d902646
|
||||
size: 68171008
|
||||
hash: md5
|
||||
path: test.parquet
|
||||
5
data/train.parquet.dvc
Normal file
5
data/train.parquet.dvc
Normal file
@@ -0,0 +1,5 @@
|
||||
outs:
|
||||
- md5: 24bf315461302605f8fd229eb263f499
|
||||
size: 115206905
|
||||
hash: md5
|
||||
path: train.parquet
|
||||
5
data/validation.parquet.dvc
Normal file
5
data/validation.parquet.dvc
Normal file
@@ -0,0 +1,5 @@
|
||||
outs:
|
||||
- md5: 5b3f656605f48a8aabeb5b21a6ae27ca
|
||||
size: 48074392
|
||||
hash: md5
|
||||
path: validation.parquet
|
||||
0
lab/__init__.py
Normal file
0
lab/__init__.py
Normal file
87
lab/constants.py
Normal file
87
lab/constants.py
Normal file
@@ -0,0 +1,87 @@
|
||||
import datetime
|
||||
from enum import StrEnum
|
||||
from pathlib import Path
|
||||
from typing import Literal
|
||||
|
||||
# Racine du depot de travail (/home/user/tp sur la VM) : lab/constants.py -> parents[1]
|
||||
REPO_ROOT = Path(__file__).resolve().parents[1]
|
||||
|
||||
# Donnees source, deja preparees (hors git, volumineuses) : voir /data sur la VM
|
||||
SOURCE_DIR = Path("/data/modelling")
|
||||
|
||||
# Sorties de split versionnees par DVC dans le depot
|
||||
DATASET_DIR = REPO_ROOT / "data"
|
||||
|
||||
FEATURE_FILENAME = "features.parquet"
|
||||
TARGET_FILENAME = "target.parquet"
|
||||
|
||||
|
||||
class SplitStrategy(StrEnum):
|
||||
FULL_HISTORY = "full_history"
|
||||
RECENT_HISTORY = "recent_history"
|
||||
|
||||
|
||||
# Strategie de split active : pilote a la fois le decoupage produit par split/cli.py
|
||||
# et le parametre "split_strategy" logge dans MLflow. On la modifie (et on committe)
|
||||
# a chaque changement de version de dataset pour synchroniser DVC et Git.
|
||||
CHOSEN_SPLIT_STRATEGY = SplitStrategy.RECENT_HISTORY
|
||||
|
||||
DatasetPart = Literal["train", "test", "validation"]
|
||||
|
||||
DATASET_SPLIT_DATES: dict[SplitStrategy, dict[DatasetPart, tuple[datetime.date, datetime.date]]] = {
|
||||
# Partie 1 : tout l'historique disponible pour l'entrainement
|
||||
SplitStrategy.FULL_HISTORY: {
|
||||
"train": (datetime.date(2011, 1, 1), datetime.date(2012, 12, 31)),
|
||||
"validation": (datetime.date(2013, 1, 1), datetime.date(2013, 12, 31)),
|
||||
"test": (datetime.date(2014, 1, 1), datetime.date(2014, 12, 31)),
|
||||
},
|
||||
# Partie 2 : donnees plus recentes uniquement
|
||||
SplitStrategy.RECENT_HISTORY: {
|
||||
"train": (datetime.date(2013, 1, 1), datetime.date(2013, 12, 31)),
|
||||
"validation": (datetime.date(2014, 1, 1), datetime.date(2014, 5, 31)),
|
||||
"test": (datetime.date(2014, 6, 1), datetime.date(2014, 12, 31)),
|
||||
},
|
||||
}
|
||||
|
||||
|
||||
class ModellingStrategy(StrEnum):
|
||||
SHORT_MEMORY = "short_memory"
|
||||
SEASONALITY = "seasonality"
|
||||
TENDENCY = "tendency"
|
||||
MIXED = "mixed"
|
||||
FULL = "full"
|
||||
|
||||
|
||||
Features = Literal["lag_1d", "lag_7d", "lag_30d", "lag_365d", "rolling_mean_7d", "rolling_mean_30d"]
|
||||
TARGET = "consumption_kwh"
|
||||
|
||||
MODELLING_FEATURES: dict[ModellingStrategy, list] = {
|
||||
# La conso depend surtout de la veille
|
||||
ModellingStrategy.SHORT_MEMORY: ["lag_1d"],
|
||||
# La conso est plus saisonniere que journaliere
|
||||
ModellingStrategy.SEASONALITY: ["lag_7d", "lag_30d"],
|
||||
# La conso suit surtout une tendance
|
||||
ModellingStrategy.TENDENCY: ["rolling_mean_7d", "rolling_mean_30d"],
|
||||
# Melange lags + tendance
|
||||
ModellingStrategy.MIXED: ["lag_1d", "lag_7d", "lag_30d", "rolling_mean_30d"],
|
||||
# Toutes les features disponibles (ajoutee en Partie 2 Etape 3)
|
||||
ModellingStrategy.FULL: ["lag_1d", "lag_7d", "lag_30d", "lag_365d", "rolling_mean_7d", "rolling_mean_30d"],
|
||||
}
|
||||
|
||||
# Valeurs d'alpha demandees par l'enonce (Partie 3)
|
||||
RIDGE_ALPHAS = [1, 1e3, 1e9]
|
||||
|
||||
|
||||
# --- TP03 : Model Registry + service de prediction ---
|
||||
# Nom sous lequel les modeles sont enregistres dans le MLflow Model Registry.
|
||||
# On garde un nom stable pour retrouver le modele et empiler ses versions.
|
||||
REGISTERED_MODEL_NAME = "electricity-consumption"
|
||||
|
||||
# Alias pointant vers la version promue (chargee par l'API). On identifie le modele
|
||||
# a servir par son alias (mobile) plutot que par un numero de version (fige).
|
||||
MODEL_ALIAS = "champion"
|
||||
|
||||
# Strategie de features du modele expose par l'API : "full" (les 6 features).
|
||||
# L'ordre des colonnes servies doit correspondre a celui de l'entrainement.
|
||||
SERVING_STRATEGY = ModellingStrategy.FULL
|
||||
SERVING_FEATURES = MODELLING_FEATURES[SERVING_STRATEGY]
|
||||
0
lab/modeling/__init__.py
Normal file
0
lab/modeling/__init__.py
Normal file
101
lab/modeling/cli.py
Normal file
101
lab/modeling/cli.py
Normal file
@@ -0,0 +1,101 @@
|
||||
import logging
|
||||
|
||||
import mlflow
|
||||
import pandas as pd
|
||||
import typer
|
||||
from mlflow.models import infer_signature
|
||||
from sklearn import linear_model
|
||||
from sklearn import metrics
|
||||
|
||||
from .. import constants
|
||||
|
||||
app = typer.Typer()
|
||||
|
||||
logging.basicConfig(level=logging.INFO)
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
|
||||
@app.command()
|
||||
def main(
|
||||
strategy: constants.ModellingStrategy,
|
||||
register: bool = typer.Option(
|
||||
False,
|
||||
"--register/--no-register",
|
||||
help="Enregistrer le modele dans le Model Registry (cree une nouvelle version).",
|
||||
),
|
||||
):
|
||||
training_file_path = constants.DATASET_DIR / "train.parquet"
|
||||
validation_file_path = constants.DATASET_DIR / "validation.parquet"
|
||||
|
||||
features = constants.MODELLING_FEATURES[strategy]
|
||||
|
||||
train_df = pd.read_parquet(training_file_path)
|
||||
validation_df = pd.read_parquet(validation_file_path)
|
||||
|
||||
train_df = train_df.dropna(
|
||||
subset=features + [constants.TARGET]
|
||||
)
|
||||
validation_df = validation_df.dropna(
|
||||
subset=features + [constants.TARGET]
|
||||
)
|
||||
|
||||
X_train = train_df[features]
|
||||
y_train = train_df[constants.TARGET]
|
||||
|
||||
X_validation = validation_df[features]
|
||||
y_validation = validation_df[constants.TARGET]
|
||||
|
||||
logger.info(f"Training model with strategy '{strategy}' and features {features}")
|
||||
|
||||
with mlflow.start_run(run_name=f"modelling_{strategy.value}"):
|
||||
mlflow.log_param("model_type", "linear")
|
||||
mlflow.log_param("strategy", strategy.value)
|
||||
mlflow.log_param("split_strategy", constants.CHOSEN_SPLIT_STRATEGY.value)
|
||||
|
||||
mlflow.log_param("features", ",".join(features))
|
||||
model = linear_model.LinearRegression()
|
||||
model.fit(X_train, y_train)
|
||||
|
||||
train_predictions = model.predict(X_train)
|
||||
validation_predictions = model.predict(X_validation)
|
||||
|
||||
train_rmse = metrics.root_mean_squared_error(y_train, train_predictions)
|
||||
validation_rmse = metrics.root_mean_squared_error(y_validation, validation_predictions)
|
||||
|
||||
train_mae = metrics.mean_absolute_error(y_train, train_predictions)
|
||||
validation_mae = metrics.mean_absolute_error(y_validation, validation_predictions)
|
||||
|
||||
mlflow.log_metric("train_rmse", train_rmse)
|
||||
mlflow.log_metric("validation_rmse", validation_rmse)
|
||||
|
||||
mlflow.log_metric("train_mae", train_mae)
|
||||
mlflow.log_metric("validation_mae", validation_mae)
|
||||
|
||||
for feature_name, coefficient in zip(
|
||||
features,
|
||||
model.coef_,
|
||||
strict=True,
|
||||
):
|
||||
mlflow.log_metric(f"coef_{feature_name}", float(coefficient))
|
||||
|
||||
# Sauvegarde de l'artefact du modele (poids + signature + environnement
|
||||
# d'execution : requirements.txt, conda.yaml, MLmodel). --register empile
|
||||
# une nouvelle version dans le Model Registry pour les meilleures experiences.
|
||||
signature = infer_signature(X_train, train_predictions)
|
||||
mlflow.sklearn.log_model(
|
||||
sk_model=model,
|
||||
name="model",
|
||||
signature=signature,
|
||||
input_example=X_train.iloc[:5],
|
||||
registered_model_name=(
|
||||
constants.REGISTERED_MODEL_NAME if register else None
|
||||
),
|
||||
)
|
||||
if register:
|
||||
logger.info(
|
||||
f"Model registered as '{constants.REGISTERED_MODEL_NAME}' (nouvelle version)"
|
||||
)
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
app()
|
||||
0
lab/modeling_ridge/__init__.py
Normal file
0
lab/modeling_ridge/__init__.py
Normal file
97
lab/modeling_ridge/cli.py
Normal file
97
lab/modeling_ridge/cli.py
Normal file
@@ -0,0 +1,97 @@
|
||||
import logging
|
||||
|
||||
import mlflow
|
||||
import pandas as pd
|
||||
import typer
|
||||
from mlflow.models import infer_signature
|
||||
from sklearn import linear_model
|
||||
from sklearn import metrics
|
||||
|
||||
from .. import constants
|
||||
|
||||
app = typer.Typer()
|
||||
|
||||
logging.basicConfig(level=logging.INFO)
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
|
||||
@app.command()
|
||||
def main(
|
||||
strategy: constants.ModellingStrategy = constants.ModellingStrategy.MIXED,
|
||||
register: bool = typer.Option(
|
||||
False,
|
||||
"--register/--no-register",
|
||||
help="Enregistrer chaque modele (par alpha) dans le Model Registry.",
|
||||
),
|
||||
):
|
||||
training_file_path = constants.DATASET_DIR / "train.parquet"
|
||||
validation_file_path = constants.DATASET_DIR / "validation.parquet"
|
||||
|
||||
features = constants.MODELLING_FEATURES[strategy]
|
||||
|
||||
train_df = pd.read_parquet(training_file_path)
|
||||
validation_df = pd.read_parquet(validation_file_path)
|
||||
|
||||
train_df = train_df.dropna(
|
||||
subset=features + [constants.TARGET]
|
||||
)
|
||||
validation_df = validation_df.dropna(
|
||||
subset=features + [constants.TARGET]
|
||||
)
|
||||
|
||||
X_train = train_df[features]
|
||||
y_train = train_df[constants.TARGET]
|
||||
|
||||
X_validation = validation_df[features]
|
||||
y_validation = validation_df[constants.TARGET]
|
||||
|
||||
logger.info(f"Training Ridge with strategy '{strategy}' and features {features}")
|
||||
|
||||
for alpha in constants.RIDGE_ALPHAS:
|
||||
with mlflow.start_run(run_name=f"modelling_{strategy.value}_ridge_alpha_{alpha:g}"):
|
||||
mlflow.log_param("model_type", "ridge")
|
||||
mlflow.log_param("strategy", strategy.value)
|
||||
mlflow.log_param("split_strategy", constants.CHOSEN_SPLIT_STRATEGY.value)
|
||||
|
||||
mlflow.log_param("features", ",".join(features))
|
||||
mlflow.log_param("alpha", alpha)
|
||||
model = linear_model.Ridge(alpha=alpha)
|
||||
model.fit(X_train, y_train)
|
||||
|
||||
train_predictions = model.predict(X_train)
|
||||
validation_predictions = model.predict(X_validation)
|
||||
|
||||
train_rmse = metrics.root_mean_squared_error(y_train, train_predictions)
|
||||
validation_rmse = metrics.root_mean_squared_error(y_validation, validation_predictions)
|
||||
|
||||
train_mae = metrics.mean_absolute_error(y_train, train_predictions)
|
||||
validation_mae = metrics.mean_absolute_error(y_validation, validation_predictions)
|
||||
|
||||
mlflow.log_metric("train_rmse", train_rmse)
|
||||
mlflow.log_metric("validation_rmse", validation_rmse)
|
||||
|
||||
mlflow.log_metric("train_mae", train_mae)
|
||||
mlflow.log_metric("validation_mae", validation_mae)
|
||||
|
||||
for feature_name, coefficient in zip(
|
||||
features,
|
||||
model.coef_,
|
||||
strict=True,
|
||||
):
|
||||
mlflow.log_metric(f"coef_{feature_name}", float(coefficient))
|
||||
|
||||
# Sauvegarde de l'artefact du modele (voir lab/modeling/cli.py).
|
||||
signature = infer_signature(X_train, train_predictions)
|
||||
mlflow.sklearn.log_model(
|
||||
sk_model=model,
|
||||
name="model",
|
||||
signature=signature,
|
||||
input_example=X_train.iloc[:5],
|
||||
registered_model_name=(
|
||||
constants.REGISTERED_MODEL_NAME if register else None
|
||||
),
|
||||
)
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
app()
|
||||
0
lab/registry/__init__.py
Normal file
0
lab/registry/__init__.py
Normal file
63
lab/registry/cli.py
Normal file
63
lab/registry/cli.py
Normal file
@@ -0,0 +1,63 @@
|
||||
"""Gestion du MLflow Model Registry : consulter les versions et promouvoir via alias.
|
||||
|
||||
La promotion = pointer un alias (mobile) vers une version precise (figee) du modele.
|
||||
L'API charge ensuite le modele par `models:/<nom>@<alias>` sans connaitre le numero.
|
||||
"""
|
||||
|
||||
import logging
|
||||
|
||||
import typer
|
||||
from mlflow import MlflowClient
|
||||
|
||||
from .. import constants
|
||||
|
||||
app = typer.Typer(help="MLflow Model Registry (versions + promotion par alias).")
|
||||
|
||||
logging.basicConfig(level=logging.INFO)
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
|
||||
@app.command()
|
||||
def versions(
|
||||
name: str = constants.REGISTERED_MODEL_NAME,
|
||||
):
|
||||
"""Lister les versions du modele enregistre, avec leurs alias."""
|
||||
client = MlflowClient()
|
||||
results = client.search_model_versions(f"name = '{name}'")
|
||||
if not results:
|
||||
logger.warning(f"Aucune version pour le modele '{name}'.")
|
||||
raise typer.Exit(code=1)
|
||||
|
||||
# Les alias sont portes par le modele enregistre (dict alias -> version),
|
||||
# pas par les objets renvoyes par search_model_versions.
|
||||
alias_by_version: dict[str, list[str]] = {}
|
||||
for alias, version in client.get_registered_model(name).aliases.items():
|
||||
alias_by_version.setdefault(str(version), []).append(alias)
|
||||
|
||||
for mv in sorted(results, key=lambda v: int(v.version)):
|
||||
aliases = ", ".join(alias_by_version.get(str(mv.version), [])) or "-"
|
||||
run = client.get_run(mv.run_id) if mv.run_id else None
|
||||
strategy = run.data.params.get("strategy", "?") if run else "?"
|
||||
val_rmse = run.data.metrics.get("validation_rmse") if run else None
|
||||
rmse_txt = f"{val_rmse:.3f}" if val_rmse is not None else "?"
|
||||
typer.echo(
|
||||
f"v{mv.version:<3} | alias: {aliases:<12} | strategy={strategy:<12}"
|
||||
f" | validation_rmse={rmse_txt} | run={mv.run_id}"
|
||||
)
|
||||
|
||||
|
||||
@app.command()
|
||||
def promote(
|
||||
version: int = typer.Option(..., help="Numero de version a promouvoir."),
|
||||
alias: str = typer.Option(constants.MODEL_ALIAS, help="Alias a (re)pointer."),
|
||||
name: str = constants.REGISTERED_MODEL_NAME,
|
||||
):
|
||||
"""Promouvoir une version : (re)pointer l'alias vers cette version."""
|
||||
client = MlflowClient()
|
||||
client.set_registered_model_alias(name=name, alias=alias, version=str(version))
|
||||
mv = client.get_model_version_by_alias(name=name, alias=alias)
|
||||
logger.info(f"Alias '{alias}' -> {name} v{mv.version} (source: {mv.source})")
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
app()
|
||||
0
lab/serving/__init__.py
Normal file
0
lab/serving/__init__.py
Normal file
99
lab/serving/api.py
Normal file
99
lab/serving/api.py
Normal file
@@ -0,0 +1,99 @@
|
||||
"""API REST de prediction de consommation electrique (FastAPI).
|
||||
|
||||
Le service orchestre la chaine de prediction :
|
||||
1. recevoir la requete (client_id, date) ;
|
||||
2. recuperer les features du client (feature store simule) ;
|
||||
3. charger le modele promu (Model Registry, au demarrage) ;
|
||||
4. calculer la prediction et repondre en JSON.
|
||||
|
||||
Lancement : uvicorn lab.serving.api:app --host 0.0.0.0 --port 8000
|
||||
Swagger UI : /docs
|
||||
"""
|
||||
|
||||
import logging
|
||||
from contextlib import asynccontextmanager
|
||||
|
||||
from fastapi import FastAPI, HTTPException
|
||||
|
||||
from . import features, registry
|
||||
from .schemas import (
|
||||
BatchPredictionRequest,
|
||||
BatchPredictionResponse,
|
||||
PredictionRequest,
|
||||
PredictionResponse,
|
||||
)
|
||||
|
||||
logging.basicConfig(level=logging.INFO)
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
# Modele charge une fois au demarrage (couteux) et reutilise a chaque requete.
|
||||
_state: dict[str, registry.LoadedModel] = {}
|
||||
|
||||
|
||||
@asynccontextmanager
|
||||
async def lifespan(app: FastAPI):
|
||||
_state["model"] = registry.load_champion()
|
||||
yield
|
||||
_state.clear()
|
||||
|
||||
|
||||
app = FastAPI(
|
||||
title="Electricity Consumption Prediction API",
|
||||
description="Expose le modele promu (MLflow Model Registry) via une API REST.",
|
||||
version="1.0.0",
|
||||
lifespan=lifespan,
|
||||
)
|
||||
|
||||
|
||||
def _get_model() -> registry.LoadedModel:
|
||||
model = _state.get("model")
|
||||
if model is None: # modele indisponible au demarrage
|
||||
raise HTTPException(status_code=503, detail="Modele non charge.")
|
||||
return model
|
||||
|
||||
|
||||
def _predict(client_id: str, model: registry.LoadedModel) -> PredictionResponse:
|
||||
feats = features.get_features(client_id) # peut lever UnknownClientError
|
||||
prediction = model.predict_one(feats)
|
||||
return PredictionResponse(
|
||||
client_id=client_id,
|
||||
prediction_kwh=prediction,
|
||||
model_name=model.name,
|
||||
model_version=model.version,
|
||||
)
|
||||
|
||||
|
||||
@app.get("/health", summary="Verification de l'etat du service")
|
||||
def health() -> dict[str, str]:
|
||||
"""Endpoint de sante : renvoie 200 si le service repond."""
|
||||
return {"status": "ok"}
|
||||
|
||||
|
||||
@app.post("/predict", response_model=PredictionResponse, summary="Prediction unitaire")
|
||||
def predict(request: PredictionRequest) -> PredictionResponse:
|
||||
model = _get_model()
|
||||
try:
|
||||
return _predict(request.client_id, model)
|
||||
except features.UnknownClientError:
|
||||
# Partie 3 : client inconnu -> erreur cliente, pas un plantage du service.
|
||||
raise HTTPException(
|
||||
status_code=404,
|
||||
detail=f"Client inconnu : aucune feature pour '{request.client_id}'.",
|
||||
)
|
||||
|
||||
|
||||
@app.post(
|
||||
"/predict/batch",
|
||||
response_model=BatchPredictionResponse,
|
||||
summary="Prediction en batch (plusieurs clients)",
|
||||
)
|
||||
def predict_batch(request: BatchPredictionRequest) -> BatchPredictionResponse:
|
||||
model = _get_model()
|
||||
predictions: list[PredictionResponse] = []
|
||||
unknown: list[str] = []
|
||||
for client_id in request.client_ids:
|
||||
try:
|
||||
predictions.append(_predict(client_id, model))
|
||||
except features.UnknownClientError:
|
||||
unknown.append(client_id)
|
||||
return BatchPredictionResponse(predictions=predictions, unknown_client_ids=unknown)
|
||||
66
lab/serving/features.py
Normal file
66
lab/serving/features.py
Normal file
@@ -0,0 +1,66 @@
|
||||
"""Recuperation des features de prediction.
|
||||
|
||||
Dans un systeme reel, ces valeurs proviendraient d'un feature store / d'une base
|
||||
alimentee par le pipeline de calcul de features (lags, moyennes glissantes) sur
|
||||
l'historique de consommation. On separe volontairement cette etape du calcul de la
|
||||
prediction : la source des features peut changer sans toucher au modele.
|
||||
|
||||
Ici on SIMULE cette recuperation par un simple dictionnaire Python, comme demande
|
||||
par l'enonce. Les valeurs sont des observations reelles (derniere ligne connue de
|
||||
quelques clients dans data/test.parquet).
|
||||
"""
|
||||
|
||||
from .. import constants
|
||||
|
||||
|
||||
class UnknownClientError(KeyError):
|
||||
"""Aucune feature disponible pour ce client (identifiant inconnu)."""
|
||||
|
||||
|
||||
# feature store simule : client_id -> {feature: valeur}
|
||||
FEATURE_STORE: dict[str, dict[str, float]] = {
|
||||
"MT_124": {
|
||||
"lag_1d": 107.656,
|
||||
"lag_7d": 25.120,
|
||||
"lag_30d": 70.574,
|
||||
"lag_365d": 25.120,
|
||||
"rolling_mean_7d": 65.870,
|
||||
"rolling_mean_30d": 71.310,
|
||||
},
|
||||
"MT_156": {
|
||||
"lag_1d": 13.149,
|
||||
"lag_7d": 13.929,
|
||||
"lag_30d": 21.577,
|
||||
"lag_365d": 8.935,
|
||||
"rolling_mean_7d": 16.648,
|
||||
"rolling_mean_30d": 19.720,
|
||||
},
|
||||
"MT_158": {
|
||||
"lag_1d": 30.739,
|
||||
"lag_7d": 16.608,
|
||||
"lag_30d": 34.094,
|
||||
"lag_365d": 6.574,
|
||||
"rolling_mean_7d": 19.067,
|
||||
"rolling_mean_30d": 21.486,
|
||||
},
|
||||
"MT_159": {
|
||||
"lag_1d": 23.305,
|
||||
"lag_7d": 24.741,
|
||||
"lag_30d": 21.386,
|
||||
"lag_365d": 5.333,
|
||||
"rolling_mean_7d": 11.707,
|
||||
"rolling_mean_30d": 13.619,
|
||||
},
|
||||
}
|
||||
|
||||
|
||||
def get_features(client_id: str) -> dict[str, float]:
|
||||
"""Renvoyer les features du client, ordonnees comme a l'entrainement du modele.
|
||||
|
||||
Leve UnknownClientError si le client est inconnu.
|
||||
"""
|
||||
if client_id not in FEATURE_STORE:
|
||||
raise UnknownClientError(client_id)
|
||||
raw = FEATURE_STORE[client_id]
|
||||
# On respecte l'ordre des colonnes attendu par le modele (SERVING_FEATURES).
|
||||
return {feature: raw[feature] for feature in constants.SERVING_FEATURES}
|
||||
43
lab/serving/registry.py
Normal file
43
lab/serving/registry.py
Normal file
@@ -0,0 +1,43 @@
|
||||
"""Chargement du modele promu depuis le MLflow Model Registry (par alias)."""
|
||||
|
||||
import logging
|
||||
from dataclasses import dataclass
|
||||
|
||||
import mlflow
|
||||
import pandas as pd
|
||||
|
||||
from .. import constants
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
|
||||
@dataclass
|
||||
class LoadedModel:
|
||||
"""Modele charge + metadonnees de version, garde en memoire par l'API."""
|
||||
|
||||
model: mlflow.pyfunc.PyFuncModel
|
||||
name: str
|
||||
version: str
|
||||
|
||||
def predict_one(self, features: dict[str, float]) -> float:
|
||||
"""Prediction pour un jeu de features (1 ligne)."""
|
||||
frame = pd.DataFrame([features], columns=constants.SERVING_FEATURES)
|
||||
return float(self.model.predict(frame)[0])
|
||||
|
||||
|
||||
def load_champion() -> LoadedModel:
|
||||
"""Charger la version pointee par l'alias `MODEL_ALIAS` du modele enregistre.
|
||||
|
||||
On identifie le modele par `models:/<nom>@<alias>` : le meme code sert n'importe
|
||||
quelle version promue, sans redeploiement, en changeant seulement l'alias cote MLflow.
|
||||
"""
|
||||
name = constants.REGISTERED_MODEL_NAME
|
||||
alias = constants.MODEL_ALIAS
|
||||
uri = f"models:/{name}@{alias}"
|
||||
logger.info(f"Chargement du modele {uri}")
|
||||
|
||||
client = mlflow.MlflowClient()
|
||||
version = client.get_model_version_by_alias(name=name, alias=alias)
|
||||
model = mlflow.pyfunc.load_model(uri)
|
||||
logger.info(f"Modele charge : {name} v{version.version}")
|
||||
return LoadedModel(model=model, name=name, version=str(version.version))
|
||||
44
lab/serving/schemas.py
Normal file
44
lab/serving/schemas.py
Normal file
@@ -0,0 +1,44 @@
|
||||
"""Schemas Pydantic du service de prediction (contrat d'entree/sortie de l'API)."""
|
||||
|
||||
import datetime
|
||||
|
||||
from pydantic import BaseModel, Field
|
||||
|
||||
|
||||
class PredictionRequest(BaseModel):
|
||||
"""Ce que le client fournit : QUI et QUAND, pas les features (calculees cote serveur)."""
|
||||
|
||||
client_id: str = Field(
|
||||
...,
|
||||
description="Identifiant du client (ex. 'MT_124').",
|
||||
examples=["MT_124"],
|
||||
)
|
||||
date: datetime.date | None = Field(
|
||||
default=None,
|
||||
description="Date de la prediction (parametre de requete). Optionnelle.",
|
||||
examples=["2015-01-01"],
|
||||
)
|
||||
|
||||
|
||||
class BatchPredictionRequest(BaseModel):
|
||||
"""Prediction pour plusieurs clients en une seule requete (inference batch)."""
|
||||
|
||||
client_ids: list[str] = Field(
|
||||
...,
|
||||
description="Liste d'identifiants clients.",
|
||||
examples=[["MT_124", "MT_156", "MT_158"]],
|
||||
)
|
||||
date: datetime.date | None = None
|
||||
|
||||
|
||||
class PredictionResponse(BaseModel):
|
||||
client_id: str
|
||||
prediction_kwh: float = Field(description="Consommation predite (kWh).")
|
||||
model_name: str
|
||||
model_version: str
|
||||
|
||||
|
||||
class BatchPredictionResponse(BaseModel):
|
||||
predictions: list[PredictionResponse]
|
||||
# Clients ignores (features introuvables) : on ne fait pas echouer tout le lot.
|
||||
unknown_client_ids: list[str] = Field(default_factory=list)
|
||||
0
lab/split/__init__.py
Normal file
0
lab/split/__init__.py
Normal file
35
lab/split/cli.py
Normal file
35
lab/split/cli.py
Normal file
@@ -0,0 +1,35 @@
|
||||
import logging
|
||||
|
||||
import pandas as pd
|
||||
|
||||
from .. import constants
|
||||
|
||||
logging.basicConfig(level=logging.INFO)
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
|
||||
def main():
|
||||
# Entrees : dataset source deja prepare (hors git)
|
||||
feature_file_path = constants.SOURCE_DIR / constants.FEATURE_FILENAME
|
||||
target_file_path = constants.SOURCE_DIR / constants.TARGET_FILENAME
|
||||
|
||||
logger.info(f"Read dataset from {feature_file_path} and {target_file_path}")
|
||||
df_features = pd.read_parquet(feature_file_path)
|
||||
df_target = pd.read_parquet(target_file_path)
|
||||
|
||||
df = df_features.join(df_target)
|
||||
timestamps = df.index.get_level_values("timestamp")
|
||||
|
||||
# Sorties : splits versionnes par DVC dans le depot
|
||||
constants.DATASET_DIR.mkdir(parents=True, exist_ok=True)
|
||||
logger.info(f"Split strategy: {constants.CHOSEN_SPLIT_STRATEGY.value}")
|
||||
for dataset_name, (start_date, end_date) in constants.DATASET_SPLIT_DATES[constants.CHOSEN_SPLIT_STRATEGY].items():
|
||||
file_path = constants.DATASET_DIR / f"{dataset_name}.parquet"
|
||||
mask = (timestamps.date >= start_date) & (timestamps.date <= end_date)
|
||||
df_split = df[mask]
|
||||
logger.info(f"Split dataset into {dataset_name} ({start_date} -> {end_date}) with shape: {df_split.shape}")
|
||||
df_split.to_parquet(file_path)
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
main()
|
||||
70
pipeline.py
Normal file
70
pipeline.py
Normal file
@@ -0,0 +1,70 @@
|
||||
"""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()
|
||||
14
requirements.txt
Normal file
14
requirements.txt
Normal file
@@ -0,0 +1,14 @@
|
||||
# Dependances runtime du pipeline (Partie 4 : "installer les dependances").
|
||||
# Versions figees sur l'environnement de reference de la VM (venv /opt/venvs/mlops,
|
||||
# Python 3.14) pour garantir la reproductibilite et la coherence avec le serveur
|
||||
# MLflow (3.13.0) et l'API de service (chargement du modele enregistre).
|
||||
mlflow==3.13.0
|
||||
scikit-learn==1.9.0
|
||||
scipy==1.17.1
|
||||
numpy==2.4.6
|
||||
pandas==2.3.3
|
||||
pyarrow==24.0.0
|
||||
typer==0.26.7
|
||||
boto3==1.43.0
|
||||
botocore==1.43.0
|
||||
s3fs==2026.6.0
|
||||
10
serve.sh
Executable file
10
serve.sh
Executable file
@@ -0,0 +1,10 @@
|
||||
#!/usr/bin/env bash
|
||||
# Lance l'API de prediction (TP03). A executer depuis la racine du depot (/home/user/tp).
|
||||
# Charge .env (MLflow + creds S3 pour recuperer le modele promu depuis le Registry).
|
||||
set -euo pipefail
|
||||
cd "$(dirname "$0")"
|
||||
set -a
|
||||
# shellcheck disable=SC1091
|
||||
source .env
|
||||
set +a
|
||||
exec /opt/venvs/mlops/bin/uvicorn lab.serving.api:app --host 0.0.0.0 --port 8000
|
||||
646
tp05_evidently_drift_2011_vs_2014.html
Normal file
646
tp05_evidently_drift_2011_vs_2014.html
Normal file
File diff suppressed because one or more lines are too long
597
tp_module5_monitoring_derive.ipynb
Normal file
597
tp_module5_monitoring_derive.ipynb
Normal file
File diff suppressed because one or more lines are too long
Reference in New Issue
Block a user