diff --git a/.github/workflows/airflow.yml b/.github/workflows/airflow.yml index d5a722a..8ae13d1 100644 --- a/.github/workflows/airflow.yml +++ b/.github/workflows/airflow.yml @@ -91,11 +91,12 @@ jobs: 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 + # 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.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/README.md b/README.md index 56093e9..e05ac18 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/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/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..4e5aaa4 100644 --- a/docs/architecture/10-infra.md +++ b/docs/architecture/10-infra.md @@ -51,7 +51,7 @@ 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`) : @@ -75,6 +75,12 @@ l'[ADR 0008](../adr/0008-airflow-execute-le-code-du-backend.md). | `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 diff --git a/etl/README.md b/etl/README.md index 698ffca..d789d6b 100644 --- a/etl/README.md +++ b/etl/README.md @@ -663,8 +663,16 @@ mock_api_import.py La logique d'extraction, de transformation et de chargement est donc 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 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). -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). +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. + +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 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 new file mode 100644 index 0000000..a668578 --- /dev/null +++ b/etl/airflow/dags/historical_import.py @@ -0,0 +1,45 @@ +"""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" +SOURCE_TIMEZONE = "UTC" +BATCH_SIZE = 1000 + +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), + ) diff --git a/etl/airflow/tests/test_dags.py b/etl/airflow/tests/test_dags.py index 555849d..96bf030 100644 --- a/etl/airflow/tests/test_dags.py +++ b/etl/airflow/tests/test_dags.py @@ -10,12 +10,13 @@ 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"), ] @@ -47,6 +48,10 @@ def test_alertes_runs_after_the_hourly_scoring(dagbag: DagBag) -> None: 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,12 +72,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. @@ -133,6 +155,10 @@ def test_alertes_retries_after_a_transient_failure(dagbag: DagBag, task_id: str) 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