From d86224a0f71798f4375a7fd4890aef2b63afa22a Mon Sep 17 00:00:00 2001 From: Meryemel-gham Date: Tue, 22 Sep 2026 11:55:13 +0200 Subject: [PATCH 1/2] feat(airflow): orchestre l'import de la Mock API --- .github/workflows/airflow.yml | 7 ++-- README.md | 2 +- docker-compose.yml | 15 ++++++--- docs/architecture/00-vue-ensemble.md | 7 ++-- docs/architecture/10-infra.md | 8 ++++- etl/README.md | 10 ++++-- etl/airflow/dags/mock_api_import.py | 50 ++++++++++++++++++++++++++++ etl/airflow/tests/test_dags.py | 45 ++++++++++++++++++++++++- 8 files changed, 127 insertions(+), 17 deletions(-) create mode 100644 etl/airflow/dags/mock_api_import.py diff --git a/.github/workflows/airflow.yml b/.github/workflows/airflow.yml index 23b6af0..0d652d9 100644 --- a/.github/workflows/airflow.yml +++ b/.github/workflows/airflow.yml @@ -93,11 +93,12 @@ jobs: # `--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 + # Les commandes des DAGs `alertes`, historique et API Mock sont couvertes. + - name: Vérifie que les quatre 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 + && env -u VIRTUAL_ENV uv run --no-sync python -m app.etl.historical_import --help + && env -u VIRTUAL_ENV uv run --no-sync python -m app.etl.mock_api_import --help" \ No newline at end of file diff --git a/README.md b/README.md index daa08f4..2d04a38 100644 --- a/README.md +++ b/README.md @@ -50,7 +50,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, import historique) +│ ├── dags/ DAGs d'orchestration (pipeline ML, alertes, imports historique et API Mock) │ ├── 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 ffa7238..5960dfa 100644 --- a/docker-compose.yml +++ b/docker-compose.yml @@ -34,12 +34,19 @@ x-airflow-common: &airflow-common # memes identifiants que le backend en attendant. ML_DATABASE_URL: postgresql+psycopg://${POSTGRES_USER}:${POSTGRES_PASSWORD}@db:5432/${POSTGRES_DB} MLFLOW_TRACKING_URI: sqlite:////opt/ml/state/mlflow.db - # Le DAG `alertes` lance le backend en sous-processus : il lit `DATABASE_URL`, en - # dialecte asyncpg, là où le pipeline ML lit `ML_DATABASE_URL`. + # Les DAGs backend lisent `DATABASE_URL` en dialecte asyncpg, là où le pipeline ML + # utilise `ML_DATABASE_URL`. DATABASE_URL: postgresql+asyncpg://${POSTGRES_USER}:${POSTGRES_PASSWORD}@db:5432/${POSTGRES_DB} - # Clé distincte de celle de l'API : la détection ne signe ni ne vérifie aucun jeton, et - # Airflow permet d'exécuter du code depuis son interface (cf. ADR 0008). + + # Clé distincte de celle de l'API : les traitements lancés par Airflow ne signent ni ne + # vérifient aucun jeton. Airflow permet d'exécuter du code depuis son interface (ADR 0008). APP_SECRET_KEY: ${AIRFLOW_APP_SECRET_KEY:-} + + # Configuration utilisée par `app.etl.mock_api_import` dans le scheduler Airflow. + APP_MOCK_API_BASE_URL: ${APP_MOCK_API_BASE_URL:-https://api-mock.charlieandre.fr} + APP_MOCK_API_USERNAME: ${APP_MOCK_API_USERNAME:-} + APP_MOCK_API_PASSWORD: ${APP_MOCK_API_PASSWORD:-} + APP_MOCK_API_TIMEOUT_SECONDS: ${APP_MOCK_API_TIMEOUT_SECONDS:-10} volumes: - ./etl/airflow/dags:/opt/airflow/dags - ./etl/airflow/plugins:/opt/airflow/plugins diff --git a/docs/architecture/00-vue-ensemble.md b/docs/architecture/00-vue-ensemble.md index 0aec09c..ea30741 100644 --- a/docs/architecture/00-vue-ensemble.md +++ b/docs/architecture/00-vue-ensemble.md @@ -70,11 +70,10 @@ 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 : quatre DAGs tournent, deux pour +Le lien `airflow --> db` est maintenant en trait plein : cinq 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), 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. +génération des recommandations (issue #116), `historical_import` pour le dataset historique +(issue #119) et `mock_api_import` pour l'ingestion horaire de l'API Mock (issue #15). Le lien `prom -.-> api` de même : l'API expose bien `/metrics` au format Prometheus, mais aucun collecteur ne vient le lire. diff --git a/docs/architecture/10-infra.md b/docs/architecture/10-infra.md index 204067b..f4485c8 100644 --- a/docs/architecture/10-infra.md +++ b/docs/architecture/10-infra.md @@ -54,7 +54,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 de l'api-server : c'est le scheduler qui a besoin du volume `airflow_ml_state` (modèle, magasin MLflow). -### Airflow (issues #115, #116 et #119) +### Airflow (issues #15, #115, #116 et #119) Quatre services (Airflow 3.3), `docker compose profiles` non utilisés (démarrage explicite via `make airflow-up`, pas dans `make dev`) : @@ -87,12 +87,18 @@ l'[ADR 0008](../adr/0008-airflow-execute-le-code-du-backend.md). | `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` | +| `mock_api_import` | `0 * * * *` | `app.etl.mock_api_import`, dans `/opt/backend/.venv` ; importe l'intervalle horaire Airflow précédent depuis l'API Mock | 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. +Le DAG `mock_api_import` exécute le pipeline API Mock toutes les heures. Il transmet +`data_interval_start` et `data_interval_end` au script backend et charge les mesures dans les +mêmes tables `site` et `reading` que le pipeline historique. `max_active_runs=1` empêche deux +intervalles de s'exécuter simultanément. + **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 diff --git a/etl/README.md b/etl/README.md index d789d6b..f952283 100644 --- a/etl/README.md +++ b/etl/README.md @@ -670,9 +670,13 @@ recommandations (`alertes`, issue #116), ainsi que l'import historique 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. +`/opt/data/raw`. -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 `mock_api_import` exécute `app.etl.mock_api_import` toutes les heures. Chaque exécution +traite l'intervalle Airflow précédent. Les deux pipelines normalisent leurs données vers les +tables communes `site` et `reading`, tout en conservant leur source (`csv` ou `api_history`). + +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`, `alertes.py`, `historical_import.py` et +`mock_api_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/mock_api_import.py b/etl/airflow/dags/mock_api_import.py new file mode 100644 index 0000000..07973e2 --- /dev/null +++ b/etl/airflow/dags/mock_api_import.py @@ -0,0 +1,50 @@ +"""DAG d'import périodique des données de l'API Mock EnerVision (issue #15). + +Orchestre le pipeline existant `app.etl.mock_api_import` sans dupliquer sa logique ETL. +Chaque exécution traite l'intervalle horaire Airflow précédent. + +Le pipeline backend reste responsable de la validation, de la normalisation, du suivi de la +qualité, de l'idempotence et du chargement dans PostgreSQL/TimescaleDB. +""" + +from __future__ import annotations + +from datetime import datetime, timedelta + +from airflow.providers.standard.operators.bash import BashOperator +from airflow.sdk import DAG + +# Le backend possède son propre environnement uv dans l'image Airflow (ADR 0008). +COMMANDE_BACKEND = "cd /opt/backend && env -u VIRTUAL_ENV uv run --no-sync python -m" + +# L'API et le pipeline backend plafonnent une réponse à 1 000 lectures par site. +LIMITE_LECTURES = 60 + +# Deux reprises donnent trois tentatives au total. Même dans le pire cas, l'exécution reste +# inférieure au pas horaire du DAG. +NOMBRE_REPRISES = 2 +DELAI_ENTRE_REPRISES = timedelta(minutes=2) +PLAFOND_PAR_TENTATIVE = timedelta(minutes=10) + +with DAG( + dag_id="mock_api_import", + description="Importe chaque heure les données de l'API Mock dans site et reading.", + schedule="@hourly", + start_date=datetime(2026, 1, 1), + catchup=False, + # Deux exécutions simultanées pourraient demander et traiter le même intervalle. + max_active_runs=1, + tags=["etl", "mock-api"], +) as dag: + BashOperator( + task_id="import_mock_api", + bash_command=( + f"{COMMANDE_BACKEND} app.etl.mock_api_import " + '--start-time "{{ data_interval_start.isoformat() }}" ' + '--end-time "{{ data_interval_end.isoformat() }}" ' + f"--limit {LIMITE_LECTURES}" + ), + retries=NOMBRE_REPRISES, + retry_delay=DELAI_ENTRE_REPRISES, + execution_timeout=PLAFOND_PAR_TENTATIVE, + ) diff --git a/etl/airflow/tests/test_dags.py b/etl/airflow/tests/test_dags.py index 60978b5..282b7d5 100644 --- a/etl/airflow/tests/test_dags.py +++ b/etl/airflow/tests/test_dags.py @@ -10,13 +10,20 @@ from airflow.sdk import BaseOperator DAGS_FOLDER = Path(__file__).resolve().parent.parent / "dags" -DAG_IDS = ["ml_train", "ml_score", "alertes", "historical_import"] +DAG_IDS = [ + "ml_train", + "ml_score", + "alertes", + "historical_import", + "mock_api_import", +] TACHES = [ ("ml_train", "train"), ("ml_score", "score"), ("alertes", "detection"), ("alertes", "recommandations"), ("historical_import", "import_historical"), + ("mock_api_import", "import_mock_api"), ] @@ -52,6 +59,10 @@ def test_historical_import_has_no_schedule(dagbag: DagBag) -> None: assert dagbag.dags["historical_import"].schedule is None +def test_mock_api_import_runs_every_hour(dagbag: DagBag) -> None: + assert dagbag.dags["mock_api_import"].timetable.expression == "0 * * * *" + + 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 @@ -84,6 +95,20 @@ def test_historical_import_uses_the_expected_source_files(dagbag: DagBag) -> Non assert "--metadata /opt/data/raw/dataset_metadata.json" in commande +def test_mock_api_import_calls_the_existing_backend_module(dagbag: DagBag) -> None: + commande = dagbag.dags["mock_api_import"].get_task("import_mock_api").bash_command + + assert "app.etl.mock_api_import" in commande + + +def test_mock_api_import_uses_the_airflow_data_interval(dagbag: DagBag) -> None: + commande = dagbag.dags["mock_api_import"].get_task("import_mock_api").bash_command + + assert '--start-time "{{ data_interval_start.isoformat() }}"' in commande + assert '--end-time "{{ data_interval_end.isoformat() }}"' in commande + assert "--limit 60" 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). @@ -95,6 +120,12 @@ def test_historical_import_runs_in_the_backend_environment(dagbag: DagBag) -> No assert "/opt/backend" in commande +def test_mock_api_import_runs_in_the_backend_environment(dagbag: DagBag) -> None: + commande = dagbag.dags["mock_api_import"].get_task("import_mock_api").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. @@ -136,6 +167,14 @@ def duree_au_pire(tache: BaseOperator) -> timedelta: return (tache.retries + 1) * tache.execution_timeout + tache.retries * tache.retry_delay +def test_mock_api_import_worst_case_stays_below_its_hourly_step( + dagbag: DagBag, +) -> None: + tache = dagbag.dags["mock_api_import"].get_task("import_mock_api") + + assert duree_au_pire(tache) < timedelta(hours=1) + + 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. @@ -159,6 +198,10 @@ def test_historical_import_retries_after_a_transient_failure(dagbag: DagBag) -> assert dagbag.dags["historical_import"].get_task("import_historical").retries >= 1 +def test_mock_api_import_retries_after_a_transient_failure(dagbag: DagBag) -> None: + assert dagbag.dags["mock_api_import"].get_task("import_mock_api").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 From e118c008bf653b521dec9561cf001c9655ecf410 Mon Sep 17 00:00:00 2001 From: Meryemel-gham Date: Tue, 22 Sep 2026 16:24:05 +0200 Subject: [PATCH 2/2] fix(airflow): fiabilise l'import horaire de la Mock API --- README.md | 2 +- apps/backend/app/etl/mock_api_import.py | 14 +++++--- .../backend/tests/etl/test_mock_api_import.py | 35 +++++++++++++++++++ docker-compose.yml | 13 ++++--- docs/architecture/00-vue-ensemble.md | 12 ++++--- docs/architecture/10-infra.md | 16 ++++++--- etl/README.md | 16 +++++---- etl/airflow/dags/mock_api_import.py | 23 ++++++++---- etl/airflow/tests/test_dags.py | 22 ++++++++---- 9 files changed, 113 insertions(+), 40 deletions(-) diff --git a/README.md b/README.md index 2d04a38..236d8f6 100644 --- a/README.md +++ b/README.md @@ -23,7 +23,7 @@ Ce que la documentation apporte à chacun : [docs/architecture/00-vue-ensemble.m | Backend | FastAPI, Python 3.14 | `apps/backend` | En place | | Frontend | Angular 22, Node 26 | `apps/frontend` | En place | | Base | PostgreSQL 17 + TimescaleDB | `db` | En place | -| ETL | Apache Airflow | `etl/airflow` | Quatre DAGs | +| ETL | Apache Airflow | `etl/airflow` | Cinq DAGs | | Infra | Terraform (k3s single-node) | `infra/terraform` | Initialise | | Reverse proxy | Nginx, TLS | `infra/proxy` | En place | | CI/CD | GitHub Actions | `.github/workflows` | En place | diff --git a/apps/backend/app/etl/mock_api_import.py b/apps/backend/app/etl/mock_api_import.py index 0d5d6be..65fcea5 100644 --- a/apps/backend/app/etl/mock_api_import.py +++ b/apps/backend/app/etl/mock_api_import.py @@ -45,15 +45,19 @@ CAPACITY_BOUNDS = (0.0, 100_000.0) def create_mock_api_client() -> httpx.AsyncClient: settings = get_settings() - if settings.mock_api_username is None or settings.mock_api_password is None: + username = settings.mock_api_username + password = ( + settings.mock_api_password.get_secret_value() + if settings.mock_api_password is not None + else None + ) + + if not username or not username.strip() or not password or not password.strip(): raise ValueError("Les identifiants de l'API Mock ne sont pas configurés.") return httpx.AsyncClient( base_url=settings.mock_api_base_url.rstrip("/"), - auth=( - settings.mock_api_username, - settings.mock_api_password.get_secret_value(), - ), + auth=(username, password), timeout=settings.mock_api_timeout_seconds, ) diff --git a/apps/backend/tests/etl/test_mock_api_import.py b/apps/backend/tests/etl/test_mock_api_import.py index cdcff55..fa654bc 100644 --- a/apps/backend/tests/etl/test_mock_api_import.py +++ b/apps/backend/tests/etl/test_mock_api_import.py @@ -283,6 +283,41 @@ def test_create_mock_api_client_requires_credentials( mock_api_import.create_mock_api_client() +@pytest.mark.parametrize( + ("username", "password_value"), + [ + ("", "test-password"), + ("test-user", ""), + (" ", "test-password"), + ("test-user", " "), + ], +) +def test_create_mock_api_client_rejects_empty_credentials( + monkeypatch: pytest.MonkeyPatch, + username: str, + password_value: str, +) -> None: + password = MagicMock() + password.get_secret_value.return_value = password_value + + settings = SimpleNamespace( + mock_api_username=username, + mock_api_password=password, + ) + + monkeypatch.setattr( + mock_api_import, + "get_settings", + lambda: settings, + ) + + with pytest.raises( + ValueError, + match="Les identifiants de l'API Mock ne sont pas configurés", + ): + mock_api_import.create_mock_api_client() + + async def test_create_mock_api_client_uses_configuration( monkeypatch: pytest.MonkeyPatch, ) -> None: diff --git a/docker-compose.yml b/docker-compose.yml index 5960dfa..b6d9f5d 100644 --- a/docker-compose.yml +++ b/docker-compose.yml @@ -42,11 +42,6 @@ x-airflow-common: &airflow-common # vérifient aucun jeton. Airflow permet d'exécuter du code depuis son interface (ADR 0008). APP_SECRET_KEY: ${AIRFLOW_APP_SECRET_KEY:-} - # Configuration utilisée par `app.etl.mock_api_import` dans le scheduler Airflow. - APP_MOCK_API_BASE_URL: ${APP_MOCK_API_BASE_URL:-https://api-mock.charlieandre.fr} - APP_MOCK_API_USERNAME: ${APP_MOCK_API_USERNAME:-} - APP_MOCK_API_PASSWORD: ${APP_MOCK_API_PASSWORD:-} - APP_MOCK_API_TIMEOUT_SECONDS: ${APP_MOCK_API_TIMEOUT_SECONDS:-10} volumes: - ./etl/airflow/dags:/opt/airflow/dags - ./etl/airflow/plugins:/opt/airflow/plugins @@ -174,6 +169,14 @@ services: airflow-scheduler: <<: *airflow-common command: scheduler + environment: + <<: *airflow-common-env + # LocalExecutor exécute les tâches dans le scheduler : lui seul a besoin des + # identifiants de l'API Mock. + APP_MOCK_API_BASE_URL: ${APP_MOCK_API_BASE_URL:-https://api-mock.charlieandre.fr} + APP_MOCK_API_USERNAME: ${APP_MOCK_API_USERNAME:-} + APP_MOCK_API_PASSWORD: ${APP_MOCK_API_PASSWORD:-} + APP_MOCK_API_TIMEOUT_SECONDS: ${APP_MOCK_API_TIMEOUT_SECONDS:-10} depends_on: db: condition: service_healthy diff --git a/docs/architecture/00-vue-ensemble.md b/docs/architecture/00-vue-ensemble.md index ea30741..ad49c47 100644 --- a/docs/architecture/00-vue-ensemble.md +++ b/docs/architecture/00-vue-ensemble.md @@ -74,6 +74,7 @@ Le lien `airflow --> db` est maintenant en trait plein : cinq DAGs tournent, deu 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), `historical_import` pour le dataset historique (issue #119) et `mock_api_import` pour l'ingestion horaire de l'API Mock (issue #15). +La réconciliation globale des données provenant des deux sources reste à 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,15 +89,16 @@ 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)). Provisionnement de la VM par Terraform, qui installe Docker, prépare les deux environnements et enregistre le runner, jamais appliqué ([ADR 0010](../adr/0010-terraform-provisionne-github-actions-deploie.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. 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 | +| ETL | Apache Airflow | `etl/airflow` | `En cours` | Webserver et scheduler avec LocalExecutor via Docker Compose, sur une base PostgreSQL dédiée. Cinq DAGs sont présents : `ml_train`, `ml_score`, `alertes`, `historical_import` et `mock_api_import`. L'import historique reste manuel et l'import API Mock est exécuté chaque heure. La réconciliation globale des deux sources reste à compléter dans l'issue #15. | | CI/CD | GitHub Actions | `.github/workflows` | `En cours` | 7 workflows, 19 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, formatage et validation du Terraform. Déploiement continu vers la VM ENI écrit par `deploy.yml`, `dev` en recette et `main` en production après approbation ([ADR 0009](../adr/0009-deux-environnements-compose-sur-la-vm-eni.md)), mais jamais exécuté : la machine n'est pas provisionnée et le runner n'y est pas enregistré. Détail dans [50-cicd.md](50-cicd.md) | ## Flux bout en bout -Statut : `En cours`. **Le chemin de lecture tourne** : base, API et frontend. **Le chemin -d'ingestion dessiné ci-dessous n'existe pas** : les trois DAGs livrés (`ml_train`, `ml_score`, -issue #115 ; `alertes`, issue #116) orchestrent le pipeline ML et la détection d'alertes, pas -l'ingestion, qui reste lancée à la main par les scripts d'import (issues #15 et #16). +Statut : `En cours`. **Le chemin de lecture tourne** entre la base, l'API et le frontend. +**Le chemin d'ingestion est maintenant orchestré par Airflow** : `historical_import` charge +manuellement le dataset CSV/JSON et `mock_api_import` collecte périodiquement les mesures de +l'API Mock. La réconciliation globale des données provenant des deux sources reste à compléter +dans l'issue #15. ```mermaid sequenceDiagram diff --git a/docs/architecture/10-infra.md b/docs/architecture/10-infra.md index f4485c8..ff7fa0f 100644 --- a/docs/architecture/10-infra.md +++ b/docs/architecture/10-infra.md @@ -87,17 +87,23 @@ l'[ADR 0008](../adr/0008-airflow-execute-le-code-du-backend.md). | `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` | -| `mock_api_import` | `0 * * * *` | `app.etl.mock_api_import`, dans `/opt/backend/.venv` ; importe l'intervalle horaire Airflow précédent depuis l'API Mock | +| `mock_api_import` | `45 * * * *` | `app.etl.mock_api_import`, dans `/opt/backend/.venv` ; importe l'heure précédant son déclenchement depuis l'API Mock | 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. -Le DAG `mock_api_import` exécute le pipeline API Mock toutes les heures. Il transmet -`data_interval_start` et `data_interval_end` au script backend et charge les mesures dans les -mêmes tables `site` et `reading` que le pipeline historique. `max_active_runs=1` empêche deux -intervalles de s'exécuter simultanément. +Le DAG `mock_api_import` exécute le pipeline API Mock toutes les heures, à la minute `:45`. +Un `CronTriggerTimetable` explicite lui attribue un intervalle d'une heure, y compris lors d'un +déclenchement manuel. Il transmet cet intervalle au script backend et charge les mesures dans +les tables communes `site` et `reading`. Le décalage à `:45` laisse quinze minutes avant le +scoring exécuté à l'heure pile, puis quinze minutes supplémentaires avant les alertes à `:15`. +`max_active_runs=1` empêche deux exécutions du DAG de se chevaucher. + +Le DAG conserve `catchup=False` pour éviter un rattrapage massif depuis sa date de démarrage. +Une interruption du scheduler peut donc créer un intervalle manquant, qui devra être rejoué +explicitement par une opération de backfill. **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 f952283..94f8ba2 100644 --- a/etl/README.md +++ b/etl/README.md @@ -663,14 +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), 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 tourne désormais réellement (`etl/airflow/`, `make airflow-up`) et orchestre cinq DAGs : +le pipeline ML (`ml_train` et `ml_score`, issue #115), la détection d'alertes et la génération +des recommandations (`alertes`, issue #116), l'import historique (`historical_import`, +issue #119) et l'import périodique de l'API Mock (`mock_api_import`, issue #15). -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`. +Le DAG `mock_api_import` s'exécute chaque heure, à la minute `:45`. Il appelle +`app.etl.mock_api_import` avec un intervalle explicite d'une heure et une limite de 1 000 lectures +par site. Les deux pipelines normalisent leurs données vers les tables communes `site` et +`reading`, tout en conservant leur source (`csv` ou `api_history`). La réconciliation globale +des deux sources reste à compléter dans l'issue #15. Le DAG `mock_api_import` exécute `app.etl.mock_api_import` toutes les heures. Chaque exécution traite l'intervalle Airflow précédent. Les deux pipelines normalisent leurs données vers les diff --git a/etl/airflow/dags/mock_api_import.py b/etl/airflow/dags/mock_api_import.py index 07973e2..f0043d0 100644 --- a/etl/airflow/dags/mock_api_import.py +++ b/etl/airflow/dags/mock_api_import.py @@ -1,7 +1,7 @@ """DAG d'import périodique des données de l'API Mock EnerVision (issue #15). Orchestre le pipeline existant `app.etl.mock_api_import` sans dupliquer sa logique ETL. -Chaque exécution traite l'intervalle horaire Airflow précédent. +Chaque exécution traite l'heure précédant son déclenchement. Le pipeline backend reste responsable de la validation, de la normalisation, du suivi de la qualité, de l'idempotence et du chargement dans PostgreSQL/TimescaleDB. @@ -13,12 +13,14 @@ from datetime import datetime, timedelta from airflow.providers.standard.operators.bash import BashOperator from airflow.sdk import DAG +from airflow.timetables.trigger import CronTriggerTimetable # Le backend possède son propre environnement uv dans l'image Airflow (ADR 0008). COMMANDE_BACKEND = "cd /opt/backend && env -u VIRTUAL_ENV uv run --no-sync python -m" -# L'API et le pipeline backend plafonnent une réponse à 1 000 lectures par site. -LIMITE_LECTURES = 60 +# Le pipeline backend et l'API acceptent au maximum 1 000 lectures par site. +# Cette marge évite de perdre silencieusement une lecture si une heure en contient plus de 60. +LIMITE_LECTURES = 1000 # Deux reprises donnent trois tentatives au total. Même dans le pire cas, l'exécution reste # inférieure au pas horaire du DAG. @@ -26,10 +28,19 @@ NOMBRE_REPRISES = 2 DELAI_ENTRE_REPRISES = timedelta(minutes=2) PLAFOND_PAR_TENTATIVE = timedelta(minutes=10) +# L'intervalle est déclaré explicitement pour ne pas dépendre de la valeur du paramètre Airflow +# `create_cron_data_intervals`. Le déclenchement à :45 laisse quinze minutes avant `ml_score`, +# exécuté à l'heure pile, puis avant `alertes`, exécuté à :15. +PLANIFICATION = CronTriggerTimetable( + "45 * * * *", + timezone="UTC", + interval=timedelta(hours=1), +) + with DAG( dag_id="mock_api_import", description="Importe chaque heure les données de l'API Mock dans site et reading.", - schedule="@hourly", + schedule=PLANIFICATION, start_date=datetime(2026, 1, 1), catchup=False, # Deux exécutions simultanées pourraient demander et traiter le même intervalle. @@ -40,8 +51,8 @@ with DAG( task_id="import_mock_api", bash_command=( f"{COMMANDE_BACKEND} app.etl.mock_api_import " - '--start-time "{{ data_interval_start.isoformat() }}" ' - '--end-time "{{ data_interval_end.isoformat() }}" ' + "--start-time \"{{ data_interval_start.strftime('%Y-%m-%dT%H:%M:%S') }}\" " + "--end-time \"{{ data_interval_end.strftime('%Y-%m-%dT%H:%M:%S') }}\" " f"--limit {LIMITE_LECTURES}" ), retries=NOMBRE_REPRISES, diff --git a/etl/airflow/tests/test_dags.py b/etl/airflow/tests/test_dags.py index 282b7d5..572b057 100644 --- a/etl/airflow/tests/test_dags.py +++ b/etl/airflow/tests/test_dags.py @@ -1,12 +1,13 @@ """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 datetime import datetime, timedelta from pathlib import Path import pytest from airflow.dag_processing.dagbag import DagBag from airflow.sdk import BaseOperator +from airflow.timetables.trigger import CronTriggerTimetable DAGS_FOLDER = Path(__file__).resolve().parent.parent / "dags" @@ -59,8 +60,17 @@ def test_historical_import_has_no_schedule(dagbag: DagBag) -> None: assert dagbag.dags["historical_import"].schedule is None -def test_mock_api_import_runs_every_hour(dagbag: DagBag) -> None: - assert dagbag.dags["mock_api_import"].timetable.expression == "0 * * * *" +def test_mock_api_import_uses_an_explicit_hourly_interval(dagbag: DagBag) -> None: + timetable = dagbag.dags["mock_api_import"].timetable + + assert isinstance(timetable, CronTriggerTimetable) + assert timetable.serialize()["expression"] == "45 * * * *" + + manual_interval = timetable.infer_manual_data_interval( + run_after=datetime.fromisoformat("2026-09-22T12:30:00+00:00"), + ) + + assert manual_interval.end - manual_interval.start == timedelta(hours=1) def test_ml_train_task_calls_the_training_module(dagbag: DagBag) -> None: @@ -104,9 +114,9 @@ def test_mock_api_import_calls_the_existing_backend_module(dagbag: DagBag) -> No def test_mock_api_import_uses_the_airflow_data_interval(dagbag: DagBag) -> None: commande = dagbag.dags["mock_api_import"].get_task("import_mock_api").bash_command - assert '--start-time "{{ data_interval_start.isoformat() }}"' in commande - assert '--end-time "{{ data_interval_end.isoformat() }}"' in commande - assert "--limit 60" in commande + assert "--start-time \"{{ data_interval_start.strftime('%Y-%m-%dT%H:%M:%S') }}\"" in commande + assert "--end-time \"{{ data_interval_end.strftime('%Y-%m-%dT%H:%M:%S') }}\"" in commande + assert "--limit 1000" in commande @pytest.mark.parametrize("task_id", ["detection", "recommandations"])