From 308769b3255122e7a04925f4c41e77e79ca6975e Mon Sep 17 00:00:00 2001 From: Meryemel-gham Date: Mon, 21 Sep 2026 14:40:21 +0200 Subject: [PATCH 1/5] feat(etl): orchestre l'import historique avec Airflow --- docker-compose.yml | 1 + etl/airflow/dags/historical_import.py | 43 +++++++++++++++++++++++++++ 2 files changed, 44 insertions(+) create mode 100644 etl/airflow/dags/historical_import.py diff --git a/docker-compose.yml b/docker-compose.yml index da19f43..6f2c1c4 100644 --- a/docker-compose.yml +++ b/docker-compose.yml @@ -36,6 +36,7 @@ x-airflow-common: &airflow-common volumes: - ./etl/airflow/dags:/opt/airflow/dags - ./etl/airflow/plugins:/opt/airflow/plugins + - ./data/raw:/opt/data/raw:ro - airflow_logs:/opt/airflow/logs - airflow_ml_state:/opt/ml/state restart: unless-stopped diff --git a/etl/airflow/dags/historical_import.py b/etl/airflow/dags/historical_import.py new file mode 100644 index 0000000..5fb119e --- /dev/null +++ b/etl/airflow/dags/historical_import.py @@ -0,0 +1,43 @@ +"""DAG d'import du dataset historique EnerVision (issue #119). + +Orchestre le pipeline existant `app.etl.historical_import` sans dupliquer sa logique ETL. +Le dataset historique sert à initialiser l'environnement : le DAG reste donc manuel. + +Le backend est exécuté dans l'environnement `/opt/backend` embarqué dans l'image Airflow, +sur le même patron que le DAG `alertes` (ADR 0008). +""" + +from __future__ import annotations + +from datetime import datetime, timedelta + +from airflow.models.dag import DAG +from airflow.operators.bash import BashOperator + +COMMANDE_BACKEND = "cd /opt/backend && env -u VIRTUAL_ENV uv run --no-sync python -m" + +CSV_PATH = "/opt/data/raw/all_sites_combined.csv" +METADATA_PATH = "/opt/data/raw/dataset_metadata.json" + +with DAG( + dag_id="historical_import", + description="Importe le dataset historique CSV/JSON dans dataset, site et reading.", + schedule=None, + start_date=datetime(2026, 1, 1), + catchup=False, + max_active_runs=1, + tags=["etl", "historical"], +) as dag: + BashOperator( + task_id="import_historical", + bash_command=( + f"{COMMANDE_BACKEND} app.etl.historical_import " + f"--csv {CSV_PATH} " + f"--metadata {METADATA_PATH} " + "--source-timezone UTC " + "--batch-size 1000" + ), + retries=1, + retry_delay=timedelta(minutes=2), + execution_timeout=timedelta(minutes=30), + ) From f18d4f9ef9bcdb5e988e4e09cf099d5780e90d09 Mon Sep 17 00:00:00 2001 From: Meryemel-gham Date: Mon, 21 Sep 2026 14:41:08 +0200 Subject: [PATCH 2/5] test(etl): couvre le DAG d'import historique --- etl/airflow/tests/test_dags.py | 51 ++++++++++++++++++++++------------ 1 file changed, 33 insertions(+), 18 deletions(-) diff --git a/etl/airflow/tests/test_dags.py b/etl/airflow/tests/test_dags.py index 555849d..1f066f9 100644 --- a/etl/airflow/tests/test_dags.py +++ b/etl/airflow/tests/test_dags.py @@ -1,5 +1,8 @@ -"""Tests d'integrite des DAGs : s'importent sans erreur, structure attendue. Pas d'execution -reelle des taches (ca reclamerait le conteneur avec `uv`/`enervision_ml`), juste la definition.""" +"""Tests d'integrite des DAGs : s'importent sans erreur et ont la structure attendue. + +Les tests ne lancent pas reellement les traitements : ils valident uniquement la definition +des DAGs et leurs commandes. +""" from datetime import timedelta from pathlib import Path @@ -10,12 +13,14 @@ from airflow.models.dagbag import DagBag DAGS_FOLDER = Path(__file__).resolve().parent.parent / "dags" -DAG_IDS = ["ml_train", "ml_score", "alertes"] +DAG_IDS = ["ml_train", "ml_score", "alertes", "historical_import"] + TACHES = [ ("ml_train", "train"), ("ml_score", "score"), ("alertes", "detection"), ("alertes", "recommandations"), + ("historical_import", "import_historical"), ] @@ -37,16 +42,17 @@ def test_ml_train_has_no_schedule(dagbag: DagBag) -> None: def test_ml_score_runs_every_hour(dagbag: DagBag) -> None: - # `@hourly` est un alias Airflow pour ce cron, c'est sous cette forme que `.summary` le rend. assert dagbag.dags["ml_score"].timetable.summary == "0 * * * *" def test_alertes_runs_after_the_hourly_scoring(dagbag: DagBag) -> None: - # Le decalage n'est pas cosmetique : la regle `anomaly` compare une lecture a la `prediction` - # du meme instant, que `ml_score` ecrit a l'heure pile. assert dagbag.dags["alertes"].timetable.summary == "15 * * * *" +def test_historical_import_has_no_schedule(dagbag: DagBag) -> None: + assert dagbag.dags["historical_import"].timetable.summary == "None" + + def test_ml_train_task_calls_the_training_module(dagbag: DagBag) -> None: tache = dagbag.dags["ml_train"].get_task("train") assert "enervision_ml.train" in tache.bash_command @@ -67,15 +73,29 @@ def test_alertes_recommendation_task_calls_the_backend_cli(dagbag: DagBag) -> No assert "app.cli generate-recommendations" in tache.bash_command +def test_historical_import_calls_the_existing_backend_module(dagbag: DagBag) -> None: + tache = dagbag.dags["historical_import"].get_task("import_historical") + assert "app.etl.historical_import" in tache.bash_command + + +def test_historical_import_uses_the_expected_source_files(dagbag: DagBag) -> None: + commande = dagbag.dags["historical_import"].get_task("import_historical").bash_command + + assert "--csv /opt/data/raw/all_sites_combined.csv" in commande + assert "--metadata /opt/data/raw/dataset_metadata.json" in commande + + @pytest.mark.parametrize("task_id", ["detection", "recommandations"]) def test_alertes_tasks_run_in_the_backend_environment(dagbag: DagBag, task_id: str) -> None: - # Le backend a son propre venv dans l'image, distinct de celui de ml/ (ADR 0008). assert "/opt/backend" in dagbag.dags["alertes"].get_task(task_id).bash_command +def test_historical_import_runs_in_the_backend_environment(dagbag: DagBag) -> None: + commande = dagbag.dags["historical_import"].get_task("import_historical").bash_command + assert "/opt/backend" in commande + + def test_alertes_generates_recommendations_after_detecting(dagbag: DagBag) -> None: - # `recommendation.alert_id` est une cle etrangere `NOT NULL` : la generation n'a rien a lire - # tant que la detection n'a pas ecrit. assert dagbag.dags["alertes"].get_task("detection").downstream_task_ids == {"recommandations"} @@ -90,14 +110,11 @@ def test_ml_score_reuses_the_model_path_written_by_ml_train(dagbag: DagBag) -> N @pytest.mark.parametrize("dag_id", DAG_IDS) def test_no_two_runs_of_a_dag_overlap(dagbag: DagBag, dag_id: str) -> None: - # Deux entrainements ecriraient le meme fichier modele, deux scorings inseriraient en meme - # temps dans `prediction`, deux detections analyseraient la meme fenetre. assert dagbag.dags[dag_id].max_active_runs == 1 @pytest.mark.parametrize(("dag_id", "task_id"), TACHES) def test_every_task_has_an_execution_timeout(dagbag: DagBag, dag_id: str, task_id: str) -> None: - # Sans plafond, une connexion pendue immobilise un slot du scheduler indefiniment. assert dagbag.dags[dag_id].get_task(task_id).execution_timeout is not None @@ -108,15 +125,11 @@ def test_ml_score_execution_timeout_stays_below_its_hourly_step(dagbag: DagBag) def duree_au_pire(tache: BaseOperator) -> timedelta: - # `execution_timeout` plafonne une tentative, pas la tache : deux reprises occupent trois - # plafonds et deux delais d'attente. assert tache.execution_timeout is not None return (tache.retries + 1) * tache.execution_timeout + tache.retries * tache.retry_delay def test_alertes_worst_case_stays_below_its_hourly_step(dagbag: DagBag) -> None: - # Les deux taches s'enchainent : c'est leur somme, reprises comprises, qui doit tenir dans le - # pas horaire, sinon `max_active_runs=1` fait attendre l'execution suivante. taches = [ dagbag.dags["alertes"].get_task(task_id) for task_id in ("detection", "recommandations") ] @@ -129,13 +142,15 @@ def test_ml_score_retries_after_a_transient_failure(dagbag: DagBag) -> None: @pytest.mark.parametrize("task_id", ["detection", "recommandations"]) def test_alertes_retries_after_a_transient_failure(dagbag: DagBag, task_id: str) -> None: - # Les deux commandes sont idempotentes en base, une reprise ne duplique rien. assert dagbag.dags["alertes"].get_task(task_id).retries >= 1 +def test_historical_import_retries_after_a_transient_failure(dagbag: DagBag) -> None: + assert dagbag.dags["historical_import"].get_task("import_historical").retries >= 1 + + @pytest.mark.parametrize(("dag_id", "task_id"), TACHES) def test_tasks_never_resync_the_baked_environment( dagbag: DagBag, dag_id: str, task_id: str ) -> None: - # Sans `--no-sync`, `uv run` reconstruit le projet a chaque execution. assert "--no-sync" in dagbag.dags[dag_id].get_task(task_id).bash_command From ae4082c584d7cdc9b03dc287e0a4d54b594b1399 Mon Sep 17 00:00:00 2001 From: Meryemel-gham Date: Mon, 21 Sep 2026 14:41:26 +0200 Subject: [PATCH 3/5] ci(etl): valide l'import historique dans l'image Airflow --- .github/workflows/airflow.yml | 10 ++++------ 1 file changed, 4 insertions(+), 6 deletions(-) diff --git a/.github/workflows/airflow.yml b/.github/workflows/airflow.yml index d5a722a..d98b780 100644 --- a/.github/workflows/airflow.yml +++ b/.github/workflows/airflow.yml @@ -90,12 +90,10 @@ jobs: docker run --rm --network none enervision-airflow:ci bash -c "cd /opt/ml && env -u VIRTUAL_ENV uv run --no-sync python -m enervision_ml.train --help" - # `--help` sort par argparse avant `get_settings()` : ni base ni secret requis, et - # l'import du module prouve que l'environnement /opt/backend est complet. Les deux - # commandes du DAG `alertes` sont couvertes, `app.cli` tirant tout FastAPI derrière lui. - - name: Vérifie que les deux commandes du DAG alertes s'importent sans réseau + # `--help` valide que le module historique et ses dépendances sont présents dans + # l'environnement backend embarqué, sans nécessiter PostgreSQL ni accès réseau. + - name: Vérifie que l'import historique s'importe sans réseau run: > docker run --rm --network none enervision-airflow:ci bash -c "cd /opt/backend - && env -u VIRTUAL_ENV uv run --no-sync python -m app.detection.internal_alerts --help - && env -u VIRTUAL_ENV uv run --no-sync python -m app.cli generate-recommendations --help" + && env -u VIRTUAL_ENV uv run --no-sync python -m app.etl.historical_import --help" \ No newline at end of file From cb4df846cb0a036de4d763aa62905ed10f60d862 Mon Sep 17 00:00:00 2001 From: Meryemel-gham Date: Mon, 21 Sep 2026 14:41:38 +0200 Subject: [PATCH 4/5] docs(etl): documente l'orchestration de l'import historique --- README.md | 4 +- docs/architecture/00-vue-ensemble.md | 9 +- docs/architecture/10-infra.md | 376 ++++++++++++++++++++++----- etl/README.md | 192 +++++++++++++- 4 files changed, 506 insertions(+), 75 deletions(-) diff --git a/README.md b/README.md index cc139dc..ad1c07b 100644 --- a/README.md +++ b/README.md @@ -21,7 +21,7 @@ Ce que la documentation apporte à chacun : [docs/architecture/00-vue-ensemble.m | Backend | FastAPI, Python 3.14 | `apps/backend` | Initialise | | Frontend | Angular 22, Node 24 LTS | `apps/frontend` | Tableau de bord | | Base | PostgreSQL 17 + TimescaleDB | `db` | Initialise | -| ETL | Apache Airflow | `etl/airflow` | Trois DAGs | +| ETL | Apache Airflow | `etl/airflow` | Quatre DAGs | | Infra | Terraform (k3s single-node) | `infra/terraform` | Initialise | | Reverse proxy | Nginx, TLS | `infra/proxy` | En place | | CI/CD | GitHub Actions | `.github/workflows` | Backend en place | @@ -48,7 +48,7 @@ L'etat detaille de chaque brique et les vues d'architecture sont dans │ ├── migrations/ Migrations SQL versionnees │ └── seeds/ Jeux de donnees de reference ├── etl/airflow/ -│ ├── dags/ DAGs d'orchestration (pipeline ML, alertes) +│ ├── dags/ DAGs d'orchestration (pipeline ML, alertes, import historique) │ ├── plugins/ Operateurs et hooks maison │ ├── include/ Requetes SQL et ressources des DAGs │ └── tests/ Tests d'integrite des DAGs diff --git a/docs/architecture/00-vue-ensemble.md b/docs/architecture/00-vue-ensemble.md index 7bdd1da..2c9a70e 100644 --- a/docs/architecture/00-vue-ensemble.md +++ b/docs/architecture/00-vue-ensemble.md @@ -70,10 +70,11 @@ Le lien `front -.-> api` reste en pointillé : le frontend appelle bien une API, intercepteur répond à sa place tant que les endpoints n'existent pas. Voir [30-frontend.md](30-frontend.md). -Le lien `airflow --> db` est maintenant en trait plein : trois DAGs tournent, deux pour +Le lien `airflow --> db` est maintenant en trait plein : quatre DAGs tournent, deux pour l'entraînement et le scoring du modèle ML (issue #115), un pour la détection d'alertes et la -génération des recommandations (issue #116), cf. plus bas et [20-backend.md](20-backend.md). Le -reste du périmètre Airflow envisagé (ingestion, issues #15/#16) reste en pointillé, non construit. +génération des recommandations (issue #116), et `historical_import` pour l'ingestion du dataset +historique (issue #119). L'orchestration de l'import API Mock et la réconciliation globale des +deux sources restent à compléter dans l'issue #15. Le lien `prom -.-> api` de même : l'API expose bien `/metrics` au format Prometheus, mais aucun collecteur ne vient le lire. @@ -88,7 +89,7 @@ collecteur ne vient le lire. | ML | LightGBM, MLflow | `ml` | `En cours` | Pipeline d'entraînement et de scoring (`enervision_ml.train`/`.score`, features par lags/moyennes glissantes partagées entre les deux, baseline de persistance saisonnière, suivi MLflow local), exposé en lecture via `GET /predictions`, orchestré par Airflow (`ml_train`/`ml_score`). Voir [ADR 0005](../adr/0005-modele-prediction-lightgbm.md) et [ML-START.md](../ML-START.md). Surveillance de dérive (EC06, #44/#45) pas encore construite | | Infra | Docker Compose, Nginx, Terraform, k3s single-node | `infra`, `docker-compose.prod.yml` | `En cours` | Reverse proxy et overlay de déploiement écrits et validés, jamais lancés sur le serveur ([ADR 0007](../adr/0007-terminaison-tls-et-reverse-proxy-nginx.md)). Module d'installation k3s jamais appliqué, aucune ressource Kubernetes déclarée | | Monitoring | Prometheus, Grafana, Alertmanager | `monitoring` | `Cible` | Rien, hors le `/metrics` exposé par l'API | -| ETL | Apache Airflow | `etl/airflow` | `En cours` | Webserver + scheduler (LocalExecutor) tournent via docker-compose, base de métadonnées Postgres dédiée. Trois DAGs en sous-processus `uv run` : `ml_train` manuel et `ml_score` `@hourly` pour le pipeline ML (issue #115), `alertes` à `15 * * * *` pour la détection et les recommandations (issue #116, [ADR 0008](../adr/0008-airflow-execute-le-code-du-backend.md)). L'ingestion (issues #15/#16) n'a pas encore de DAG | +| ETL | Apache Airflow | `etl/airflow` | `En cours` | Webserver + scheduler (LocalExecutor) tournent via docker-compose, base de métadonnées Postgres dédiée. Quatre DAGs en sous-processus `uv run` : `ml_train`, `ml_score`, `alertes` et `historical_import`. Le DAG historique orchestre `app.etl.historical_import` et charge `dataset`, `site` et `reading`. L'orchestration API Mock reste à compléter dans #15 | | CI/CD | GitHub Actions | `.github/workflows` | `En cours` | 5 workflows, 16 jobs : lint, typage, tests avec seuil de couverture bloquant, tests d'intégration sur TimescaleDB réel, audit de dépendances, SAST Bandit, quality gate SonarCloud, intégrité des DAGs Airflow. Détail dans [50-cicd.md](50-cicd.md). **Aucun job de déploiement** (#21) | ## Flux bout en bout diff --git a/docs/architecture/10-infra.md b/docs/architecture/10-infra.md index c6121e7..a8f3624 100644 --- a/docs/architecture/10-infra.md +++ b/docs/architecture/10-infra.md @@ -51,86 +51,334 @@ Trois pièges sont documentés en tête du `docker-compose.yml`, ils ne se devin - `LocalExecutor` exécute les tâches comme sous-processus du **scheduler**, jamais du webserver : c'est le scheduler qui a besoin du volume `airflow_ml_state` (modèle, magasin MLflow). -### Airflow (issues #115 et #116) +### Airflow (issues #115, #116 et #119) -Trois services, `docker compose profiles` non utilisés (démarrage explicite via `make -airflow-up`, pas dans `make dev`) : +Trois services Airflow sont définis dans `docker-compose.yml`. Les `docker compose profiles` ne +sont pas utilisés : le démarrage reste explicite via `make airflow-up` et Airflow ne fait pas +partie de la boucle `make dev`. | Service | Rôle | Points notables | |---|---|---| -| `airflow-init` | Migre la base de métadonnées, crée le compte admin | Conteneur jetable (`restart: "no"`), ne redémarre jamais. `webserver`/`scheduler` attendent qu'il se termine avec succès | -| `airflow-webserver` | UI, port `8080` | `LocalExecutor` : n'exécute aucune tâche lui-même | -| `airflow-scheduler` | Planifie et **exécute** les tâches (`LocalExecutor`) | Les DAGs y tournent en sous-processus (`uv run --no-sync python -m ...`), c'est lui qui a besoin du volume `airflow_ml_state` | +| `airflow-init` | Migre la base de métadonnées et crée le compte admin | Conteneur jetable (`restart: "no"`). `webserver` et `scheduler` attendent qu'il se termine avec succès | +| `airflow-webserver` | Interface Airflow sur le port `8080` | Avec `LocalExecutor`, il n'exécute aucune tâche lui-même | +| `airflow-scheduler` | Planifie et exécute les tâches | Les DAGs tournent en sous-processus avec `LocalExecutor` | -Construits depuis `etl/airflow/Dockerfile`, contexte `.` (racine du repo, pas `etl/airflow/`) : -l'image doit pouvoir `COPY` les sources de `ml/` **et** de `apps/backend/` pour se synchroniser -deux environnements Python **3.14** (`/opt/ml/.venv` et `/opt/backend/.venv`, `uv sync --locked` à -la construction), distincts du Python 3.12 qui fait tourner Airflow lui-même. Les DAGs shellent -vers ces venvs plutôt que d'importer LightGBM, MLflow ou SQLAlchemy dans le process Airflow. -Le choix et ses contreparties sont dans -l'[ADR 0008](../adr/0008-airflow-execute-le-code-du-backend.md). +`LocalExecutor` exécute les tâches dans le scheduler et non dans le webserver. -| DAG | Planification | Ce qu'il lance, et où | +Les services sont construits depuis : + +```text +etl/airflow/Dockerfile +``` + +avec la racine du dépôt comme contexte Docker. + +Airflow 2.10.4 fonctionne avec Python 3.12, tandis que le pipeline ML et le backend utilisent des +dépendances Python 3.14. + +L'image Airflow embarque donc deux environnements distincts : + +```text +/opt/ml/.venv +/opt/backend/.venv +``` + +Le premier contient le pipeline Machine Learning. + +Le second contient le backend EnerVision utilisé par les DAGs `alertes` et +`historical_import`. + +Cette séparation évite d'installer directement LightGBM, MLflow ou les dépendances SQLAlchemy du +backend dans l'environnement Python utilisé par Airflow. + +Le choix est décrit dans +[l'ADR 0008](../adr/0008-airflow-execute-le-code-du-backend.md). + +Les DAGs actuellement présents sont : + +| DAG | Planification | Ce qu'il exécute | |---|---|---| -| `ml_train` | manuelle | `enervision_ml.train`, dans `/opt/ml/.venv` | -| `ml_score` | `0 * * * *` | `enervision_ml.score`, dans `/opt/ml/.venv` | -| `alertes` | `15 * * * *` | `app.detection.internal_alerts` puis `app.cli generate-recommendations`, dans `/opt/backend/.venv` | +| `ml_train` | manuelle | `enervision_ml.train` dans `/opt/ml/.venv` | +| `ml_score` | `0 * * * *` | `enervision_ml.score` dans `/opt/ml/.venv` | +| `alertes` | `15 * * * *` | `app.detection.internal_alerts` puis `app.cli generate-recommendations` dans `/opt/backend/.venv` | +| `historical_import` | manuelle | `app.etl.historical_import` dans `/opt/backend/.venv` | -**Pourquoi `alertes` tourne à la quinzième minute.** Sa règle `anomaly` compare une lecture à la -`prediction` du même instant, que `ml_score` écrit à l'heure pile. Le décalage laisse le scoring -finir. Aucune dépendance n'est déclarée entre les deux DAGs pour autant, ni `ExternalTaskSensor` ni -tâche greffée : quatre règles de détection sur cinq ne touchent pas au modèle, et un modèle jamais -entraîné ne doit pas priver le parc de ses alertes. Le décalage est donc une convention et non une -garantie : le plafond de `ml_score` est de 30 minutes, et un scoring qui déborde de `:15` prive -`anomaly` de la `prediction` de l'heure, qu'elle ne retrouvera au passage suivant que si sa fenêtre -la couvre encore. Les quatre autres règles ne s'en aperçoivent pas. +#### Import historique -Ses deux tâches s'enchaînent en revanche (`recommendation.alert_id` est une clé étrangère `NOT -NULL`), et toutes deux sont rejouables sans risque : l'idempotence est portée par la base, -`uq_alert_source_reference` et `uq_recommendation_alert_rule`. Chacune a 2 tentatives, 2 minutes -d'attente entre elles et un plafond de 5 minutes **par tentative** : au pire, reprises comprises, -l'enchaînement occupe 38 minutes, ce qui le garde sous le pas horaire qu'un `max_active_runs=1` -rend contraignant. +Le DAG : -`airflow-init` s'appuie sur l'entrypoint de l'image (`_AIRFLOW_DB_MIGRATE`, -`_AIRFLOW_WWW_USER_*`) plutôt que sur un script maison : l'entrypoint porte le code de sortie, une -migration ratée (typiquement la base `airflow` absente, cf. ci-dessous) fait échouer le service et -`webserver`/`scheduler` ne démarrent pas sur une base non migrée. Le mot de passe du compte admin -passe par l'environnement, jamais par `argv` (ni `ps`, ni `docker compose config`). +```text +historical_import +``` -Les variables `AIRFLOW_*` ne sont volontairement pas en `${VAR:?}` : Compose interpole le fichier -entier avant de filtrer les services, une variable requise manquante casserait `make db-up`, -`make dev`... pour tout poste dont le `.env` est antérieur. Elles valent `${VAR:-}` et c'est -`airflow-init` qui refuse de démarrer (clé Fernet, clé Flask, mot de passe ou -`AIRFLOW_APP_SECRET_KEY` vides). +est défini dans : -Le conteneur reçoit deux variables du backend en plus de `ML_DATABASE_URL` : `DATABASE_URL`, en -dialecte asyncpg, et `APP_SECRET_KEY`, alimentée par `AIRFLOW_APP_SECRET_KEY`. Cette dernière est -**délibérément différente** de celle de l'API. La configuration du backend refuse de se construire -sans clé, mais la détection ne signe ni ne vérifie aucun jeton : un Airflow compromis, qui permet -déjà d'exécuter du code depuis son interface, ne doit pas livrer par-dessus la clé de signature -des JWT. +```text +etl/airflow/dags/historical_import.py +``` -**Pourquoi `ml_train` est manuel.** Réentraîner est coûteux et sa cadence n'est pas une décision -prise. Surtout, `train.py` écrase le modèle sans comparer ses métriques à celles de l'ancien : un -cron déploierait silencieusement un modèle dégradé. Tant que ce garde-fou n'existe pas, le -déclenchement reste humain. `ml_score`, lui, est planifié à l'heure, avec `max_active_runs=1` -(pas deux scorings simultanés dans `prediction`), 2 tentatives et un plafond de 30 minutes. +Il ne contient aucune logique d'import propre. -CI : `.github/workflows/airflow.yml` (Python 3.12 via `etl/airflow/.python-version`) lance lint et -tests d'intégrité des DAGs, et construit l'image (elle `COPY` `ml/` et `apps/backend/`, une -modification de l'un ou de l'autre peut donc la casser, d'où leurs chemins dans les déclencheurs) -avant de vérifier que les deux environnements s'y importent sans réseau. +Il utilise un `BashOperator` pour exécuter le module backend existant : -Piège à connaître : sur un volume `pgdata` déjà peuplé (poste de dev existant plutôt que premier -`make db-up`), `db/init/120-airflow-database.sql` ne se rejoue pas (PostgreSQL n'exécute -`docker-entrypoint-initdb.d/` que sur un volume vide). Créer la base `airflow` à la main une fois : -`docker compose exec db psql -U $POSTGRES_USER -d $POSTGRES_DB -c "CREATE DATABASE airflow;"`. +```text +app.etl.historical_import +``` -`libgomp1` est installé explicitement dans l'image (`apt-get`, en root) : l'image Airflow de base -est minimale et n'embarque pas la runtime OpenMP dont LightGBM a besoin, sans quoi l'erreur -(`OSError: libgomp.so.1`) n'apparaît qu'à la première tâche réellement exécutée, pas à la -construction de l'image. +dans l'environnement : + +```text +/opt/backend/.venv +``` + +Le flux est donc : + +```text +Airflow scheduler + | + v +historical_import + | + v +BashOperator + | + v +app.etl.historical_import + | + v +PostgreSQL / TimescaleDB +``` + +Les fichiers historiques locaux sont montés en lecture seule dans les services Airflow : + +```text +./data/raw:/opt/data/raw:ro +``` + +Les chemins utilisés depuis le conteneur sont : + +```text +/opt/data/raw/all_sites_combined.csv +/opt/data/raw/dataset_metadata.json +``` + +Le montage en lecture seule empêche les traitements Airflow de modifier les fichiers sources. + +Le dataset historique sert uniquement à initialiser les données de l'environnement. + +Le DAG utilise donc : + +```text +schedule = None +catchup = False +max_active_runs = 1 +``` + +Il est déclenché manuellement. + +`max_active_runs = 1` empêche deux imports du même dataset de s'exécuter simultanément. + +La tâche possède également : + +```text +retries = 1 +retry_delay = 2 minutes +execution_timeout = 30 minutes +``` + +Le traitement historique étant idempotent, une reprise après une erreur transitoire ne doit pas +créer de doublons. + +Deux exécutions manuelles successives ont été validées avec succès. + +Après les deux exécutions, PostgreSQL/TimescaleDB contenait toujours : + +```text +122647 +``` + +lectures avec : + +```text +source = "csv" +``` + +La deuxième exécution n'a donc pas dupliqué les lectures historiques. + +#### DAG alertes + +Le DAG `alertes` est planifié à la quinzième minute de chaque heure. + +La règle `anomaly` compare une lecture à la `prediction` du même instant, que `ml_score` écrit à +l'heure pile. Le décalage laisse donc du temps au scoring pour terminer. + +Aucune dépendance Airflow explicite n'est cependant déclarée entre `ml_score` et `alertes`. +Quatre règles de détection sur cinq ne dépendent pas du modèle, et l'absence d'un modèle entraîné +ne doit pas empêcher les autres alertes d'être produites. + +Les deux tâches du DAG `alertes` s'enchaînent : + +```text +detection + | + v +recommandations +``` + +`recommendation.alert_id` étant une clé étrangère `NOT NULL`, la génération des recommandations +est exécutée après la détection. + +Les traitements sont idempotents en base grâce aux contraintes : + +```text +uq_alert_source_reference +uq_recommendation_alert_rule +``` + +Chaque tâche possède deux tentatives, deux minutes d'attente entre les tentatives et un plafond +de cinq minutes par tentative. + +#### DAGs ML + +`ml_train` reste manuel. + +Réentraîner le modèle est coûteux et `train.py` remplace actuellement le modèle existant sans +comparer automatiquement les métriques du nouveau modèle avec celles du précédent. + +Tant que ce mécanisme de sélection n'existe pas, le réentraînement reste déclenché humainement. + +`ml_score` est planifié toutes les heures et réutilise le modèle produit par `ml_train`. + +Il utilise : + +```text +max_active_runs = 1 +retries = 2 +execution_timeout = 30 minutes +``` + +Deux scorings ne peuvent donc pas s'exécuter simultanément sur les mêmes données. + +#### Configuration Airflow + +`airflow-init` s'appuie sur l'entrypoint de l'image Airflow avec : + +```text +_AIRFLOW_DB_MIGRATE +_AIRFLOW_WWW_USER_* +``` + +Une migration de la base Airflow qui échoue fait échouer `airflow-init`. + +Le `webserver` et le `scheduler` dépendent du succès de ce service et ne démarrent donc pas sur +une base de métadonnées non initialisée. + +Les variables Airflow sont fournies depuis le fichier `.env`. + +Les secrets ne sont pas passés dans les arguments des processus. + +Les variables `AIRFLOW_*` ne sont volontairement pas déclarées avec `${VAR:?}` dans le bloc +commun de Docker Compose : Compose interpole le fichier complet même lorsqu'un seul service est +démarré. + +La validation des secrets nécessaires est réalisée par `airflow-init`. + +Le scheduler reçoit également les variables nécessaires aux traitements backend : + +```text +DATABASE_URL +APP_SECRET_KEY +``` + +`APP_SECRET_KEY` est alimentée par : + +```text +AIRFLOW_APP_SECRET_KEY +``` + +Cette clé est distincte de celle utilisée par l'API EnerVision. + +#### Base de métadonnées Airflow + +Airflow utilise une base PostgreSQL dédiée : + +```text +airflow +``` + +Elle est créée lors de l'initialisation de PostgreSQL par : + +```text +db/init/120-airflow-database.sql +``` + +Sur un volume `pgdata` déjà existant, les scripts de `docker-entrypoint-initdb.d` ne sont pas +rejoués automatiquement. + +Dans ce cas, la base peut être créée manuellement une fois : + +```powershell +docker compose exec db psql -U enervision -d enervision -c "CREATE DATABASE airflow;" +``` + +#### CI Airflow + +Le workflow : + +```text +.github/workflows/airflow.yml +``` + +utilise Python 3.12 via : + +```text +etl/airflow/.python-version +``` + +Il vérifie : + +```text +formatage Ruff +analyse statique Ruff +tests d'intégrité des DAGs +construction de l'image Airflow +``` + +La construction de l'image embarque : + +```text +ml/ +apps/backend/ +``` + +Une modification de ces composants peut donc casser l'image Airflow. + +La CI vérifie également sans accès réseau que les commandes utilisées par les DAGs sont +importables depuis leurs environnements respectifs. + +Pour l'import historique, elle exécute notamment : + +```text +python -m app.etl.historical_import --help +``` + +depuis `/opt/backend`. + +Cette vérification permet de détecter une dépendance backend manquante ou un environnement Docker +incomplet sans avoir besoin de démarrer PostgreSQL. + +#### Dépendance système LightGBM + +`libgomp1` est installé explicitement dans l'image Airflow. + +LightGBM dépend de cette bibliothèque OpenMP. + +Sans elle, l'image Docker pourrait être construite correctement mais l'import de LightGBM +échouerait au moment de l'exécution avec une erreur liée à : + +```text +libgomp.so.1 +``` ## Machine cible, exécution Docker diff --git a/etl/README.md b/etl/README.md index 698ffca..15d2915 100644 --- a/etl/README.md +++ b/etl/README.md @@ -639,7 +639,7 @@ La suite backend complète a également été validée avec une couverture supé ## Suite du pipeline Data -Deux sources de données sont maintenant prises en charge : +Deux sources de données sont prises en charge par la logique ETL du backend : ```text Dataset CSV/JSON @@ -661,10 +661,192 @@ mock_api_import.py API Mock ``` -La logique d'extraction, de transformation et de chargement est donc disponible pour les deux sources de données du MVP. +La logique d'extraction, de validation, de transformation et de chargement est disponible pour +les deux sources de données du MVP. -Airflow tourne désormais réellement (`etl/airflow/`, `make airflow-up`) et orchestre le pipeline ML (`ml_train`/`ml_score`, issue #115) ainsi que la détection d'alertes et la génération des recommandations (`alertes`, issue #116). Il n'orchestre pas encore ces deux imports : `historical_import.py` et `mock_api_import.py` (normalisation et chargement micro-batch, issues #15/#16) restent à faire. +Airflow tourne réellement dans `etl/airflow/` et orchestre désormais quatre DAGs : -Airflow permet de planifier les traitements, gérer leur ordre d'exécution, suivre leur état et remonter les erreurs. Il ne remplace pas la logique ETL Python existante : les scripts actuels restent responsables de l'extraction, de la validation, de la transformation et du chargement. `etl/airflow/dags/ml_train.py`, `ml_score.py` et `alertes.py` montrent le patron retenu (des `BashOperator` qui invoquent le script tel quel, dans l'environnement `uv` que l'image embarque pour lui). +```text +ml_train +ml_score +alertes +historical_import +``` -Le pipeline Data servira ensuite à préparer les données nécessaires au modèle de Machine Learning. +Les DAGs `ml_train` et `ml_score` orchestrent le pipeline Machine Learning (issue #115). + +Le DAG `alertes` orchestre la détection des alertes et la génération des recommandations +(issue #116). + +Le DAG `historical_import` orchestre l'import du dataset historique CSV/JSON (issue #119). + +### Orchestration de l'import historique + +Le DAG historique est défini dans : + +```text +etl/airflow/dags/historical_import.py +``` + +Il ne réimplémente aucune logique ETL. Il déclenche directement le module existant : + +```text +app.etl.historical_import +``` + +Le flux d'exécution est le suivant : + +```text +data/raw/ +├── all_sites_combined.csv +└── dataset_metadata.json + | + v +Airflow + | + v +DAG historical_import + | + v +BashOperator + | + v +app.etl.historical_import + | + v +PostgreSQL / TimescaleDB + | + +--> dataset + +--> site + +--> reading +``` + +Les fichiers historiques locaux sont montés dans les conteneurs Airflow en lecture seule : + +```text +./data/raw:/opt/data/raw:ro +``` + +Le DAG utilise les chemins suivants : + +```text +/opt/data/raw/all_sites_combined.csv +/opt/data/raw/dataset_metadata.json +``` + +Le montage en lecture seule évite qu'un traitement Airflow puisse modifier les fichiers sources. + +Le backend est déjà embarqué dans l'image Airflow dans son propre environnement Python : + +```text +/opt/backend/.venv +``` + +Le DAG utilise un `BashOperator` avec le même principe que le DAG `alertes` : + +```text +cd /opt/backend +env -u VIRTUAL_ENV +uv run --no-sync python -m app.etl.historical_import +``` + +Airflow reste ainsi responsable de l'orchestration tandis que le backend reste responsable de +l'extraction, de la validation, de la transformation et du chargement. + +### Planification + +Le dataset historique sert à initialiser l'environnement et n'est pas une source périodique. + +Le DAG est donc configuré avec : + +```text +schedule = None +catchup = False +max_active_runs = 1 +``` + +Le déclenchement est manuel depuis l'interface Airflow ou avec la CLI. + +`max_active_runs = 1` empêche deux imports historiques de s'exécuter simultanément. + +La tâche `import_historical` définit également : + +```text +retries = 1 +retry_delay = 2 minutes +execution_timeout = 30 minutes +``` + +Le retry permet de reprendre le traitement après une erreur transitoire, notamment une +indisponibilité temporaire de PostgreSQL. + +Le pipeline historique étant idempotent, une nouvelle exécution ne doit pas dupliquer les mesures +déjà présentes. + +### Déclenchement et suivi + +Le DAG peut être déclenché avec : + +```powershell +docker compose exec airflow-scheduler airflow dags trigger historical_import +``` + +Les exécutions peuvent être consultées avec : + +```powershell +docker compose exec airflow-scheduler airflow dags list-runs -d historical_import +``` + +### Validation de l'orchestration + +L'orchestration a été validée localement avec Docker Compose et le `LocalExecutor` Airflow. + +Le scheduler détecte les quatre DAGs : + +```text +alertes +historical_import +ml_score +ml_train +``` + +Deux exécutions manuelles successives du DAG `historical_import` ont été réalisées. + +Les deux exécutions se sont terminées avec : + +```text +state = success +``` + +Les exécutions ont été traitées séquentiellement. + +Après les deux exécutions, PostgreSQL/TimescaleDB contenait toujours : + +```text +122647 +``` + +lectures historiques avec : + +```text +source = "csv" +``` + +Le nombre de lectures n'a donc pas doublé après la seconde exécution. + +Cette validation confirme que l'orchestration Airflow réutilise correctement +`app.etl.historical_import` et conserve l'idempotence du pipeline historique. + +### Suite + +L'import de l'API Mock est déjà disponible côté backend avec : + +```text +app.etl.mock_api_import +``` + +Son orchestration Airflow ainsi que la réconciliation globale entre les sources historique et +API Mock restent à compléter dans l'issue #15. + +Airflow ne remplace pas les pipelines Python existants : il orchestre leur exécution, leur +planification, les reprises sur erreur et leur suivi. \ No newline at end of file From fe100653d05a3b342ce0f9947bf4897b183ad6b8 Mon Sep 17 00:00:00 2001 From: Meryemel-gham Date: Mon, 21 Sep 2026 16:04:58 +0200 Subject: [PATCH 5/5] fix(etl): traite les retours de revue du DAG historique --- .github/workflows/airflow.yml | 9 +- docs/architecture/10-infra.md | 398 +++++--------------------- etl/README.md | 198 +------------ etl/airflow/dags/historical_import.py | 2 + etl/airflow/tests/test_dags.py | 23 +- 5 files changed, 115 insertions(+), 515 deletions(-) diff --git a/.github/workflows/airflow.yml b/.github/workflows/airflow.yml index d98b780..8ae13d1 100644 --- a/.github/workflows/airflow.yml +++ b/.github/workflows/airflow.yml @@ -90,10 +90,13 @@ jobs: docker run --rm --network none enervision-airflow:ci bash -c "cd /opt/ml && env -u VIRTUAL_ENV uv run --no-sync python -m enervision_ml.train --help" - # `--help` valide que le module historique et ses dépendances sont présents dans - # l'environnement backend embarqué, sans nécessiter PostgreSQL ni accès réseau. - - name: Vérifie que l'import historique s'importe sans réseau + # `--help` sort par argparse avant `get_settings()` : ni base ni secret requis, et + # l'import des modules prouve que l'environnement /opt/backend est complet. + # Les deux commandes du DAG `alertes` et la commande du DAG historique sont couvertes. + - name: Vérifie que les trois commandes backend s'importent sans réseau run: > docker run --rm --network none enervision-airflow:ci bash -c "cd /opt/backend + && env -u VIRTUAL_ENV uv run --no-sync python -m app.detection.internal_alerts --help + && env -u VIRTUAL_ENV uv run --no-sync python -m app.cli generate-recommendations --help && env -u VIRTUAL_ENV uv run --no-sync python -m app.etl.historical_import --help" \ No newline at end of file diff --git a/docs/architecture/10-infra.md b/docs/architecture/10-infra.md index a8f3624..4e5aaa4 100644 --- a/docs/architecture/10-infra.md +++ b/docs/architecture/10-infra.md @@ -53,332 +53,90 @@ Trois pièges sont documentés en tête du `docker-compose.yml`, ils ne se devin ### Airflow (issues #115, #116 et #119) -Trois services Airflow sont définis dans `docker-compose.yml`. Les `docker compose profiles` ne -sont pas utilisés : le démarrage reste explicite via `make airflow-up` et Airflow ne fait pas -partie de la boucle `make dev`. +Trois services, `docker compose profiles` non utilisés (démarrage explicite via `make +airflow-up`, pas dans `make dev`) : | Service | Rôle | Points notables | |---|---|---| -| `airflow-init` | Migre la base de métadonnées et crée le compte admin | Conteneur jetable (`restart: "no"`). `webserver` et `scheduler` attendent qu'il se termine avec succès | -| `airflow-webserver` | Interface Airflow sur le port `8080` | Avec `LocalExecutor`, il n'exécute aucune tâche lui-même | -| `airflow-scheduler` | Planifie et exécute les tâches | Les DAGs tournent en sous-processus avec `LocalExecutor` | +| `airflow-init` | Migre la base de métadonnées, crée le compte admin | Conteneur jetable (`restart: "no"`), ne redémarre jamais. `webserver`/`scheduler` attendent qu'il se termine avec succès | +| `airflow-webserver` | UI, port `8080` | `LocalExecutor` : n'exécute aucune tâche lui-même | +| `airflow-scheduler` | Planifie et **exécute** les tâches (`LocalExecutor`) | Les DAGs y tournent en sous-processus (`uv run --no-sync python -m ...`), c'est lui qui a besoin du volume `airflow_ml_state` | -`LocalExecutor` exécute les tâches dans le scheduler et non dans le webserver. +Construits depuis `etl/airflow/Dockerfile`, contexte `.` (racine du repo, pas `etl/airflow/`) : +l'image doit pouvoir `COPY` les sources de `ml/` **et** de `apps/backend/` pour se synchroniser +deux environnements Python **3.14** (`/opt/ml/.venv` et `/opt/backend/.venv`, `uv sync --locked` à +la construction), distincts du Python 3.12 qui fait tourner Airflow lui-même. Les DAGs shellent +vers ces venvs plutôt que d'importer LightGBM, MLflow ou SQLAlchemy dans le process Airflow. +Le choix et ses contreparties sont dans +l'[ADR 0008](../adr/0008-airflow-execute-le-code-du-backend.md). -Les services sont construits depuis : - -```text -etl/airflow/Dockerfile -``` - -avec la racine du dépôt comme contexte Docker. - -Airflow 2.10.4 fonctionne avec Python 3.12, tandis que le pipeline ML et le backend utilisent des -dépendances Python 3.14. - -L'image Airflow embarque donc deux environnements distincts : - -```text -/opt/ml/.venv -/opt/backend/.venv -``` - -Le premier contient le pipeline Machine Learning. - -Le second contient le backend EnerVision utilisé par les DAGs `alertes` et -`historical_import`. - -Cette séparation évite d'installer directement LightGBM, MLflow ou les dépendances SQLAlchemy du -backend dans l'environnement Python utilisé par Airflow. - -Le choix est décrit dans -[l'ADR 0008](../adr/0008-airflow-execute-le-code-du-backend.md). - -Les DAGs actuellement présents sont : - -| DAG | Planification | Ce qu'il exécute | +| DAG | Planification | Ce qu'il lance, et où | |---|---|---| -| `ml_train` | manuelle | `enervision_ml.train` dans `/opt/ml/.venv` | -| `ml_score` | `0 * * * *` | `enervision_ml.score` dans `/opt/ml/.venv` | -| `alertes` | `15 * * * *` | `app.detection.internal_alerts` puis `app.cli generate-recommendations` dans `/opt/backend/.venv` | -| `historical_import` | manuelle | `app.etl.historical_import` dans `/opt/backend/.venv` | - -#### Import historique - -Le DAG : - -```text -historical_import -``` - -est défini dans : - -```text -etl/airflow/dags/historical_import.py -``` - -Il ne contient aucune logique d'import propre. - -Il utilise un `BashOperator` pour exécuter le module backend existant : - -```text -app.etl.historical_import -``` - -dans l'environnement : - -```text -/opt/backend/.venv -``` - -Le flux est donc : - -```text -Airflow scheduler - | - v -historical_import - | - v -BashOperator - | - v -app.etl.historical_import - | - v -PostgreSQL / TimescaleDB -``` - -Les fichiers historiques locaux sont montés en lecture seule dans les services Airflow : - -```text -./data/raw:/opt/data/raw:ro -``` - -Les chemins utilisés depuis le conteneur sont : - -```text -/opt/data/raw/all_sites_combined.csv -/opt/data/raw/dataset_metadata.json -``` - -Le montage en lecture seule empêche les traitements Airflow de modifier les fichiers sources. - -Le dataset historique sert uniquement à initialiser les données de l'environnement. - -Le DAG utilise donc : - -```text -schedule = None -catchup = False -max_active_runs = 1 -``` - -Il est déclenché manuellement. - -`max_active_runs = 1` empêche deux imports du même dataset de s'exécuter simultanément. - -La tâche possède également : - -```text -retries = 1 -retry_delay = 2 minutes -execution_timeout = 30 minutes -``` - -Le traitement historique étant idempotent, une reprise après une erreur transitoire ne doit pas -créer de doublons. - -Deux exécutions manuelles successives ont été validées avec succès. - -Après les deux exécutions, PostgreSQL/TimescaleDB contenait toujours : - -```text -122647 -``` - -lectures avec : - -```text -source = "csv" -``` - -La deuxième exécution n'a donc pas dupliqué les lectures historiques. - -#### DAG alertes - -Le DAG `alertes` est planifié à la quinzième minute de chaque heure. - -La règle `anomaly` compare une lecture à la `prediction` du même instant, que `ml_score` écrit à -l'heure pile. Le décalage laisse donc du temps au scoring pour terminer. - -Aucune dépendance Airflow explicite n'est cependant déclarée entre `ml_score` et `alertes`. -Quatre règles de détection sur cinq ne dépendent pas du modèle, et l'absence d'un modèle entraîné -ne doit pas empêcher les autres alertes d'être produites. - -Les deux tâches du DAG `alertes` s'enchaînent : - -```text -detection - | - v -recommandations -``` - -`recommendation.alert_id` étant une clé étrangère `NOT NULL`, la génération des recommandations -est exécutée après la détection. - -Les traitements sont idempotents en base grâce aux contraintes : - -```text -uq_alert_source_reference -uq_recommendation_alert_rule -``` - -Chaque tâche possède deux tentatives, deux minutes d'attente entre les tentatives et un plafond -de cinq minutes par tentative. - -#### DAGs ML - -`ml_train` reste manuel. - -Réentraîner le modèle est coûteux et `train.py` remplace actuellement le modèle existant sans -comparer automatiquement les métriques du nouveau modèle avec celles du précédent. - -Tant que ce mécanisme de sélection n'existe pas, le réentraînement reste déclenché humainement. - -`ml_score` est planifié toutes les heures et réutilise le modèle produit par `ml_train`. - -Il utilise : - -```text -max_active_runs = 1 -retries = 2 -execution_timeout = 30 minutes -``` - -Deux scorings ne peuvent donc pas s'exécuter simultanément sur les mêmes données. - -#### Configuration Airflow - -`airflow-init` s'appuie sur l'entrypoint de l'image Airflow avec : - -```text -_AIRFLOW_DB_MIGRATE -_AIRFLOW_WWW_USER_* -``` - -Une migration de la base Airflow qui échoue fait échouer `airflow-init`. - -Le `webserver` et le `scheduler` dépendent du succès de ce service et ne démarrent donc pas sur -une base de métadonnées non initialisée. - -Les variables Airflow sont fournies depuis le fichier `.env`. - -Les secrets ne sont pas passés dans les arguments des processus. - -Les variables `AIRFLOW_*` ne sont volontairement pas déclarées avec `${VAR:?}` dans le bloc -commun de Docker Compose : Compose interpole le fichier complet même lorsqu'un seul service est -démarré. - -La validation des secrets nécessaires est réalisée par `airflow-init`. - -Le scheduler reçoit également les variables nécessaires aux traitements backend : - -```text -DATABASE_URL -APP_SECRET_KEY -``` - -`APP_SECRET_KEY` est alimentée par : - -```text -AIRFLOW_APP_SECRET_KEY -``` - -Cette clé est distincte de celle utilisée par l'API EnerVision. - -#### Base de métadonnées Airflow - -Airflow utilise une base PostgreSQL dédiée : - -```text -airflow -``` - -Elle est créée lors de l'initialisation de PostgreSQL par : - -```text -db/init/120-airflow-database.sql -``` - -Sur un volume `pgdata` déjà existant, les scripts de `docker-entrypoint-initdb.d` ne sont pas -rejoués automatiquement. - -Dans ce cas, la base peut être créée manuellement une fois : - -```powershell -docker compose exec db psql -U enervision -d enervision -c "CREATE DATABASE airflow;" -``` - -#### CI Airflow - -Le workflow : - -```text -.github/workflows/airflow.yml -``` - -utilise Python 3.12 via : - -```text -etl/airflow/.python-version -``` - -Il vérifie : - -```text -formatage Ruff -analyse statique Ruff -tests d'intégrité des DAGs -construction de l'image Airflow -``` - -La construction de l'image embarque : - -```text -ml/ -apps/backend/ -``` - -Une modification de ces composants peut donc casser l'image Airflow. - -La CI vérifie également sans accès réseau que les commandes utilisées par les DAGs sont -importables depuis leurs environnements respectifs. - -Pour l'import historique, elle exécute notamment : - -```text -python -m app.etl.historical_import --help -``` - -depuis `/opt/backend`. - -Cette vérification permet de détecter une dépendance backend manquante ou un environnement Docker -incomplet sans avoir besoin de démarrer PostgreSQL. - -#### Dépendance système LightGBM - -`libgomp1` est installé explicitement dans l'image Airflow. - -LightGBM dépend de cette bibliothèque OpenMP. - -Sans elle, l'image Docker pourrait être construite correctement mais l'import de LightGBM -échouerait au moment de l'exécution avec une erreur liée à : - -```text -libgomp.so.1 -``` +| `ml_train` | manuelle | `enervision_ml.train`, dans `/opt/ml/.venv` | +| `ml_score` | `0 * * * *` | `enervision_ml.score`, dans `/opt/ml/.venv` | +| `alertes` | `15 * * * *` | `app.detection.internal_alerts` puis `app.cli generate-recommendations`, dans `/opt/backend/.venv` | +| `historical_import` | manuelle | `app.etl.historical_import`, dans `/opt/backend/.venv` ; les fichiers de `data/raw` sont montés en lecture seule dans `/opt/data/raw` | + +Le DAG `historical_import` réutilise le pipeline historique existant sans dupliquer sa logique. +Il reste manuel, car le dataset sert à initialiser l'environnement. Le montage +`./data/raw:/opt/data/raw:ro` permet au scheduler de lire les fichiers CSV/JSON sans pouvoir les +modifier. + +**Pourquoi `alertes` tourne à la quinzième minute.** Sa règle `anomaly` compare une lecture à la +`prediction` du même instant, que `ml_score` écrit à l'heure pile. Le décalage laisse le scoring +finir. Aucune dépendance n'est déclarée entre les deux DAGs pour autant, ni `ExternalTaskSensor` ni +tâche greffée : quatre règles de détection sur cinq ne touchent pas au modèle, et un modèle jamais +entraîné ne doit pas priver le parc de ses alertes. Le décalage est donc une convention et non une +garantie : le plafond de `ml_score` est de 30 minutes, et un scoring qui déborde de `:15` prive +`anomaly` de la `prediction` de l'heure, qu'elle ne retrouvera au passage suivant que si sa fenêtre +la couvre encore. Les quatre autres règles ne s'en aperçoivent pas. + +Ses deux tâches s'enchaînent en revanche (`recommendation.alert_id` est une clé étrangère `NOT +NULL`), et toutes deux sont rejouables sans risque : l'idempotence est portée par la base, +`uq_alert_source_reference` et `uq_recommendation_alert_rule`. Chacune a 2 tentatives, 2 minutes +d'attente entre elles et un plafond de 5 minutes **par tentative** : au pire, reprises comprises, +l'enchaînement occupe 38 minutes, ce qui le garde sous le pas horaire qu'un `max_active_runs=1` +rend contraignant. + +`airflow-init` s'appuie sur l'entrypoint de l'image (`_AIRFLOW_DB_MIGRATE`, +`_AIRFLOW_WWW_USER_*`) plutôt que sur un script maison : l'entrypoint porte le code de sortie, une +migration ratée (typiquement la base `airflow` absente, cf. ci-dessous) fait échouer le service et +`webserver`/`scheduler` ne démarrent pas sur une base non migrée. Le mot de passe du compte admin +passe par l'environnement, jamais par `argv` (ni `ps`, ni `docker compose config`). + +Les variables `AIRFLOW_*` ne sont volontairement pas en `${VAR:?}` : Compose interpole le fichier +entier avant de filtrer les services, une variable requise manquante casserait `make db-up`, +`make dev`... pour tout poste dont le `.env` est antérieur. Elles valent `${VAR:-}` et c'est +`airflow-init` qui refuse de démarrer (clé Fernet, clé Flask, mot de passe ou +`AIRFLOW_APP_SECRET_KEY` vides). + +Le conteneur reçoit deux variables du backend en plus de `ML_DATABASE_URL` : `DATABASE_URL`, en +dialecte asyncpg, et `APP_SECRET_KEY`, alimentée par `AIRFLOW_APP_SECRET_KEY`. Cette dernière est +**délibérément différente** de celle de l'API. La configuration du backend refuse de se construire +sans clé, mais la détection ne signe ni ne vérifie aucun jeton : un Airflow compromis, qui permet +déjà d'exécuter du code depuis son interface, ne doit pas livrer par-dessus la clé de signature +des JWT. + +**Pourquoi `ml_train` est manuel.** Réentraîner est coûteux et sa cadence n'est pas une décision +prise. Surtout, `train.py` écrase le modèle sans comparer ses métriques à celles de l'ancien : un +cron déploierait silencieusement un modèle dégradé. Tant que ce garde-fou n'existe pas, le +déclenchement reste humain. `ml_score`, lui, est planifié à l'heure, avec `max_active_runs=1` +(pas deux scorings simultanés dans `prediction`), 2 tentatives et un plafond de 30 minutes. + +CI : `.github/workflows/airflow.yml` (Python 3.12 via `etl/airflow/.python-version`) lance lint et +tests d'intégrité des DAGs, et construit l'image (elle `COPY` `ml/` et `apps/backend/`, une +modification de l'un ou de l'autre peut donc la casser, d'où leurs chemins dans les déclencheurs) +avant de vérifier que les deux environnements s'y importent sans réseau. + +Piège à connaître : sur un volume `pgdata` déjà peuplé (poste de dev existant plutôt que premier +`make db-up`), `db/init/120-airflow-database.sql` ne se rejoue pas (PostgreSQL n'exécute +`docker-entrypoint-initdb.d/` que sur un volume vide). Créer la base `airflow` à la main une fois : +`docker compose exec db psql -U $POSTGRES_USER -d $POSTGRES_DB -c "CREATE DATABASE airflow;"`. + +`libgomp1` est installé explicitement dans l'image (`apt-get`, en root) : l'image Airflow de base +est minimale et n'embarque pas la runtime OpenMP dont LightGBM a besoin, sans quoi l'erreur +(`OSError: libgomp.so.1`) n'apparaît qu'à la première tâche réellement exécutée, pas à la +construction de l'image. ## Machine cible, exécution Docker diff --git a/etl/README.md b/etl/README.md index 15d2915..d789d6b 100644 --- a/etl/README.md +++ b/etl/README.md @@ -639,7 +639,7 @@ La suite backend complète a également été validée avec une couverture supé ## Suite du pipeline Data -Deux sources de données sont prises en charge par la logique ETL du backend : +Deux sources de données sont maintenant prises en charge : ```text Dataset CSV/JSON @@ -661,192 +661,18 @@ mock_api_import.py API Mock ``` -La logique d'extraction, de validation, de transformation et de chargement est disponible pour -les deux sources de données du MVP. +La logique d'extraction, de transformation et de chargement est donc disponible pour les deux sources de données du MVP. -Airflow tourne réellement dans `etl/airflow/` et orchestre désormais quatre DAGs : +Airflow tourne désormais réellement (`etl/airflow/`, `make airflow-up`) et orchestre le pipeline +ML (`ml_train`/`ml_score`, issue #115), la détection d'alertes et la génération des +recommandations (`alertes`, issue #116), ainsi que l'import historique +(`historical_import`, issue #119). -```text -ml_train -ml_score -alertes -historical_import -``` +Le DAG `historical_import` est déclenché manuellement. Il exécute +`app.etl.historical_import` avec les fichiers montés en lecture seule depuis `data/raw` vers +`/opt/data/raw`. L'orchestration de l'import API Mock et la réconciliation globale des deux +sources restent couvertes par l'issue #15. -Les DAGs `ml_train` et `ml_score` orchestrent le pipeline Machine Learning (issue #115). +Airflow permet de planifier les traitements, gérer leur ordre d'exécution, suivre leur état et remonter les erreurs. Il ne remplace pas la logique ETL Python existante : les scripts actuels restent responsables de l'extraction, de la validation, de la transformation et du chargement. `etl/airflow/dags/ml_train.py`, `ml_score.py` et `alertes.py` et `historical_import.py` montrent le patron retenu (des `BashOperator` qui invoquent le script tel quel, dans l'environnement `uv` que l'image embarque pour lui). -Le DAG `alertes` orchestre la détection des alertes et la génération des recommandations -(issue #116). - -Le DAG `historical_import` orchestre l'import du dataset historique CSV/JSON (issue #119). - -### Orchestration de l'import historique - -Le DAG historique est défini dans : - -```text -etl/airflow/dags/historical_import.py -``` - -Il ne réimplémente aucune logique ETL. Il déclenche directement le module existant : - -```text -app.etl.historical_import -``` - -Le flux d'exécution est le suivant : - -```text -data/raw/ -├── all_sites_combined.csv -└── dataset_metadata.json - | - v -Airflow - | - v -DAG historical_import - | - v -BashOperator - | - v -app.etl.historical_import - | - v -PostgreSQL / TimescaleDB - | - +--> dataset - +--> site - +--> reading -``` - -Les fichiers historiques locaux sont montés dans les conteneurs Airflow en lecture seule : - -```text -./data/raw:/opt/data/raw:ro -``` - -Le DAG utilise les chemins suivants : - -```text -/opt/data/raw/all_sites_combined.csv -/opt/data/raw/dataset_metadata.json -``` - -Le montage en lecture seule évite qu'un traitement Airflow puisse modifier les fichiers sources. - -Le backend est déjà embarqué dans l'image Airflow dans son propre environnement Python : - -```text -/opt/backend/.venv -``` - -Le DAG utilise un `BashOperator` avec le même principe que le DAG `alertes` : - -```text -cd /opt/backend -env -u VIRTUAL_ENV -uv run --no-sync python -m app.etl.historical_import -``` - -Airflow reste ainsi responsable de l'orchestration tandis que le backend reste responsable de -l'extraction, de la validation, de la transformation et du chargement. - -### Planification - -Le dataset historique sert à initialiser l'environnement et n'est pas une source périodique. - -Le DAG est donc configuré avec : - -```text -schedule = None -catchup = False -max_active_runs = 1 -``` - -Le déclenchement est manuel depuis l'interface Airflow ou avec la CLI. - -`max_active_runs = 1` empêche deux imports historiques de s'exécuter simultanément. - -La tâche `import_historical` définit également : - -```text -retries = 1 -retry_delay = 2 minutes -execution_timeout = 30 minutes -``` - -Le retry permet de reprendre le traitement après une erreur transitoire, notamment une -indisponibilité temporaire de PostgreSQL. - -Le pipeline historique étant idempotent, une nouvelle exécution ne doit pas dupliquer les mesures -déjà présentes. - -### Déclenchement et suivi - -Le DAG peut être déclenché avec : - -```powershell -docker compose exec airflow-scheduler airflow dags trigger historical_import -``` - -Les exécutions peuvent être consultées avec : - -```powershell -docker compose exec airflow-scheduler airflow dags list-runs -d historical_import -``` - -### Validation de l'orchestration - -L'orchestration a été validée localement avec Docker Compose et le `LocalExecutor` Airflow. - -Le scheduler détecte les quatre DAGs : - -```text -alertes -historical_import -ml_score -ml_train -``` - -Deux exécutions manuelles successives du DAG `historical_import` ont été réalisées. - -Les deux exécutions se sont terminées avec : - -```text -state = success -``` - -Les exécutions ont été traitées séquentiellement. - -Après les deux exécutions, PostgreSQL/TimescaleDB contenait toujours : - -```text -122647 -``` - -lectures historiques avec : - -```text -source = "csv" -``` - -Le nombre de lectures n'a donc pas doublé après la seconde exécution. - -Cette validation confirme que l'orchestration Airflow réutilise correctement -`app.etl.historical_import` et conserve l'idempotence du pipeline historique. - -### Suite - -L'import de l'API Mock est déjà disponible côté backend avec : - -```text -app.etl.mock_api_import -``` - -Son orchestration Airflow ainsi que la réconciliation globale entre les sources historique et -API Mock restent à compléter dans l'issue #15. - -Airflow ne remplace pas les pipelines Python existants : il orchestre leur exécution, leur -planification, les reprises sur erreur et leur suivi. \ No newline at end of file +Le pipeline Data servira ensuite à préparer les données nécessaires au modèle de Machine Learning. diff --git a/etl/airflow/dags/historical_import.py b/etl/airflow/dags/historical_import.py index 5fb119e..a668578 100644 --- a/etl/airflow/dags/historical_import.py +++ b/etl/airflow/dags/historical_import.py @@ -18,6 +18,8 @@ COMMANDE_BACKEND = "cd /opt/backend && env -u VIRTUAL_ENV uv run --no-sync pytho CSV_PATH = "/opt/data/raw/all_sites_combined.csv" METADATA_PATH = "/opt/data/raw/dataset_metadata.json" +SOURCE_TIMEZONE = "UTC" +BATCH_SIZE = 1000 with DAG( dag_id="historical_import", diff --git a/etl/airflow/tests/test_dags.py b/etl/airflow/tests/test_dags.py index 1f066f9..96bf030 100644 --- a/etl/airflow/tests/test_dags.py +++ b/etl/airflow/tests/test_dags.py @@ -1,8 +1,5 @@ -"""Tests d'integrite des DAGs : s'importent sans erreur et ont la structure attendue. - -Les tests ne lancent pas reellement les traitements : ils valident uniquement la definition -des DAGs et leurs commandes. -""" +"""Tests d'integrite des DAGs : s'importent sans erreur, structure attendue. Pas d'execution +reelle des taches (ca reclamerait le conteneur avec `uv`/`enervision_ml`), juste la definition.""" from datetime import timedelta from pathlib import Path @@ -14,7 +11,6 @@ from airflow.models.dagbag import DagBag DAGS_FOLDER = Path(__file__).resolve().parent.parent / "dags" DAG_IDS = ["ml_train", "ml_score", "alertes", "historical_import"] - TACHES = [ ("ml_train", "train"), ("ml_score", "score"), @@ -42,10 +38,13 @@ def test_ml_train_has_no_schedule(dagbag: DagBag) -> None: def test_ml_score_runs_every_hour(dagbag: DagBag) -> None: + # `@hourly` est un alias Airflow pour ce cron, c'est sous cette forme que `.summary` le rend. assert dagbag.dags["ml_score"].timetable.summary == "0 * * * *" def test_alertes_runs_after_the_hourly_scoring(dagbag: DagBag) -> None: + # Le decalage n'est pas cosmetique : la regle `anomaly` compare une lecture a la `prediction` + # du meme instant, que `ml_score` ecrit a l'heure pile. assert dagbag.dags["alertes"].timetable.summary == "15 * * * *" @@ -87,6 +86,7 @@ def test_historical_import_uses_the_expected_source_files(dagbag: DagBag) -> Non @pytest.mark.parametrize("task_id", ["detection", "recommandations"]) def test_alertes_tasks_run_in_the_backend_environment(dagbag: DagBag, task_id: str) -> None: + # Le backend a son propre venv dans l'image, distinct de celui de ml/ (ADR 0008). assert "/opt/backend" in dagbag.dags["alertes"].get_task(task_id).bash_command @@ -96,6 +96,8 @@ def test_historical_import_runs_in_the_backend_environment(dagbag: DagBag) -> No def test_alertes_generates_recommendations_after_detecting(dagbag: DagBag) -> None: + # `recommendation.alert_id` est une cle etrangere `NOT NULL` : la generation n'a rien a lire + # tant que la detection n'a pas ecrit. assert dagbag.dags["alertes"].get_task("detection").downstream_task_ids == {"recommandations"} @@ -110,11 +112,14 @@ def test_ml_score_reuses_the_model_path_written_by_ml_train(dagbag: DagBag) -> N @pytest.mark.parametrize("dag_id", DAG_IDS) def test_no_two_runs_of_a_dag_overlap(dagbag: DagBag, dag_id: str) -> None: + # Deux entrainements ecriraient le meme fichier modele, deux scorings inseriraient en meme + # temps dans `prediction`, deux detections analyseraient la meme fenetre. assert dagbag.dags[dag_id].max_active_runs == 1 @pytest.mark.parametrize(("dag_id", "task_id"), TACHES) def test_every_task_has_an_execution_timeout(dagbag: DagBag, dag_id: str, task_id: str) -> None: + # Sans plafond, une connexion pendue immobilise un slot du scheduler indefiniment. assert dagbag.dags[dag_id].get_task(task_id).execution_timeout is not None @@ -125,11 +130,15 @@ def test_ml_score_execution_timeout_stays_below_its_hourly_step(dagbag: DagBag) def duree_au_pire(tache: BaseOperator) -> timedelta: + # `execution_timeout` plafonne une tentative, pas la tache : deux reprises occupent trois + # plafonds et deux delais d'attente. assert tache.execution_timeout is not None return (tache.retries + 1) * tache.execution_timeout + tache.retries * tache.retry_delay def test_alertes_worst_case_stays_below_its_hourly_step(dagbag: DagBag) -> None: + # Les deux taches s'enchainent : c'est leur somme, reprises comprises, qui doit tenir dans le + # pas horaire, sinon `max_active_runs=1` fait attendre l'execution suivante. taches = [ dagbag.dags["alertes"].get_task(task_id) for task_id in ("detection", "recommandations") ] @@ -142,6 +151,7 @@ def test_ml_score_retries_after_a_transient_failure(dagbag: DagBag) -> None: @pytest.mark.parametrize("task_id", ["detection", "recommandations"]) def test_alertes_retries_after_a_transient_failure(dagbag: DagBag, task_id: str) -> None: + # Les deux commandes sont idempotentes en base, une reprise ne duplique rien. assert dagbag.dags["alertes"].get_task(task_id).retries >= 1 @@ -153,4 +163,5 @@ def test_historical_import_retries_after_a_transient_failure(dagbag: DagBag) -> def test_tasks_never_resync_the_baked_environment( dagbag: DagBag, dag_id: str, task_id: str ) -> None: + # Sans `--no-sync`, `uv run` reconstruit le projet a chaque execution. assert "--no-sync" in dagbag.dags[dag_id].get_task(task_id).bash_command