From d86224a0f71798f4375a7fd4890aef2b63afa22a Mon Sep 17 00:00:00 2001 From: Meryemel-gham Date: Tue, 22 Sep 2026 11:55:13 +0200 Subject: [PATCH] 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