From cbbfaf4910f52896498d980b772b3327091071ab Mon Sep 17 00:00:00 2001 From: Johan LEROY Date: Wed, 23 Sep 2026 15:31:51 +0200 Subject: [PATCH 1/2] fix(etl): importer une seule mesure par heure depuis l'API Mock MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit L'API Mock ne renvoie pas les mesures d'une période : elle génère `limit` points répartis sur l'intervalle demandé (1 000 par heure avec --limit 1000, un toutes les 3,6 s). Le DAG aurait écrit 7 000 lignes par heure et par environnement, alors que le dataset historique a une mesure horaire et que les features ML décalent par ligne : `shift(168)`, le retard d'une semaine, serait devenu un retard de dix minutes, sans erreur visible au scoring ni au réentraînement. Avec --limit 1, l'API renvoie la mesure de :00 de chaque heure, au pas du CSV. Constaté sur la recette le 23/09 avant la réactivation des DAGs. --- etl/README.md | 6 ++++-- etl/airflow/dags/mock_api_import.py | 6 +++--- etl/airflow/tests/test_dags.py | 2 +- 3 files changed, 8 insertions(+), 6 deletions(-) diff --git a/etl/README.md b/etl/README.md index 94f8ba2..cabf674 100644 --- a/etl/README.md +++ b/etl/README.md @@ -669,8 +669,10 @@ des recommandations (`alertes`, issue #116), l'import historique (`historical_im issue #119) et l'import périodique de l'API Mock (`mock_api_import`, issue #15). 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 +`app.etl.mock_api_import` avec un intervalle explicite d'une heure et une limite d'une lecture +par site. L'API Mock génère autant de points que la limite demandée, répartis sur l'intervalle : +un seul donne la mesure de :00, au pas horaire du dataset historique, que les features ML +supposent en décalant par ligne. 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. diff --git a/etl/airflow/dags/mock_api_import.py b/etl/airflow/dags/mock_api_import.py index f0043d0..ca50254 100644 --- a/etl/airflow/dags/mock_api_import.py +++ b/etl/airflow/dags/mock_api_import.py @@ -18,9 +18,9 @@ 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" -# 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 +# Contrainte : l'API Mock génère `limit` points répartis sur l'intervalle. Un seul donne la mesure +# de :00, au pas horaire du CSV que les features ML supposent (`shift(168)` compte des lignes). +LIMITE_LECTURES = 1 # Deux reprises donnent trois tentatives au total. Même dans le pire cas, l'exécution reste # inférieure au pas horaire du DAG. diff --git a/etl/airflow/tests/test_dags.py b/etl/airflow/tests/test_dags.py index d03dd97..eb45364 100644 --- a/etl/airflow/tests/test_dags.py +++ b/etl/airflow/tests/test_dags.py @@ -118,7 +118,7 @@ def test_mock_api_import_uses_the_airflow_data_interval(dagbag: DagBag) -> None: 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 + assert commande.endswith("--limit 1") @pytest.mark.parametrize("task_id", ["detection", "recommandations"]) From 6c09beeb3ca11b7463a19cb14502dac2446c948b Mon Sep 17 00:00:00 2001 From: Johan LEROY Date: Wed, 23 Sep 2026 15:39:04 +0200 Subject: [PATCH 2/2] =?UTF-8?q?fix(etl):=20demander=20la=20mesure=20de=20l?= =?UTF-8?q?'heure=20pile=20=C3=A0=20l'API=20Mock?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Correctif de cbbfaf4, dont le message affirmait à tort une mesure à :00. Avec --limit 1, l'API Mock renvoie le point de début d'intervalle : le DAG, déclenché à :45 sur [:45 - 1 h, :45], aurait écrit ses mesures à :45, au pas horaire mais hors de la grille du dataset historique. L'intervalle part désormais de l'heure pile du déclenchement : un run à 13:45 demande [13:00, 13:45] et importe la mesure de 13:00, avant ml_score à 14:00. Vérifié contre l'API Mock en recette et par le rendu du gabarit Jinja. --- docs/architecture/10-infra.md | 2 +- etl/README.md | 10 +++++----- etl/airflow/dags/mock_api_import.py | 8 ++++---- etl/airflow/tests/test_dags.py | 4 ++-- 4 files changed, 12 insertions(+), 12 deletions(-) diff --git a/docs/architecture/10-infra.md b/docs/architecture/10-infra.md index 8cea7bd..9b47cc7 100644 --- a/docs/architecture/10-infra.md +++ b/docs/architecture/10-infra.md @@ -89,7 +89,7 @@ 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` | `45 * * * *` | `app.etl.mock_api_import`, dans `/opt/backend/.venv` ; importe l'heure précédant son déclenchement depuis l'API Mock | +| `mock_api_import` | `45 * * * *` | `app.etl.mock_api_import`, dans `/opt/backend/.venv` ; importe depuis l'API Mock la mesure de l'heure pile précédant son déclenchement | | `derive` | `30 5 * * *` | `app.monitoring.drift`, dans `/opt/backend/.venv` ; quotidien parce que sa fenêtre couvre 168 h, et sans reprise parce qu'une dérive n'est pas une panne passagère | Le DAG `historical_import` réutilise le pipeline historique existant sans dupliquer sa logique. diff --git a/etl/README.md b/etl/README.md index cabf674..00bf60b 100644 --- a/etl/README.md +++ b/etl/README.md @@ -669,15 +669,15 @@ des recommandations (`alertes`, issue #116), l'import historique (`historical_im issue #119) et l'import périodique de l'API Mock (`mock_api_import`, issue #15). 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 d'une lecture -par site. L'API Mock génère autant de points que la limite demandée, répartis sur l'intervalle : -un seul donne la mesure de :00, au pas horaire du dataset historique, que les features ML -supposent en décalant par ligne. Les deux pipelines normalisent leurs données vers les tables communes `site` et +`app.etl.mock_api_import` sur l'intervalle qui va de l'heure pile à son déclenchement, avec une +limite d'une lecture par site. L'API Mock génère autant de points que la limite demandée, +répartis sur l'intervalle et le premier à son début : un seul donne la mesure de :00, au pas +horaire du dataset historique, que les features ML supposent en décalant par ligne. 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 +importe la mesure de l'heure pile qui précède son déclenchement. 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 diff --git a/etl/airflow/dags/mock_api_import.py b/etl/airflow/dags/mock_api_import.py index ca50254..f62aa66 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'heure précédant son déclenchement. +Chaque exécution importe la mesure de l'heure pile qui précède 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. @@ -18,8 +18,8 @@ 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" -# Contrainte : l'API Mock génère `limit` points répartis sur l'intervalle. Un seul donne la mesure -# de :00, au pas horaire du CSV que les features ML supposent (`shift(168)` compte des lignes). +# Contrainte : l'API Mock génère `limit` points répartis sur l'intervalle, le premier à son début. +# Un seul, depuis l'heure pile, donne la mesure de :00 au pas du CSV que suppose `shift(168)`. LIMITE_LECTURES = 1 # Deux reprises donnent trois tentatives au total. Même dans le pire cas, l'exécution reste @@ -51,7 +51,7 @@ with DAG( task_id="import_mock_api", bash_command=( f"{COMMANDE_BACKEND} app.etl.mock_api_import " - "--start-time \"{{ data_interval_start.strftime('%Y-%m-%dT%H:%M:%S') }}\" " + "--start-time \"{{ data_interval_end.strftime('%Y-%m-%dT%H:00:00') }}\" " "--end-time \"{{ data_interval_end.strftime('%Y-%m-%dT%H:%M:%S') }}\" " f"--limit {LIMITE_LECTURES}" ), diff --git a/etl/airflow/tests/test_dags.py b/etl/airflow/tests/test_dags.py index eb45364..473bda3 100644 --- a/etl/airflow/tests/test_dags.py +++ b/etl/airflow/tests/test_dags.py @@ -113,10 +113,10 @@ def test_mock_api_import_calls_the_existing_backend_module(dagbag: DagBag) -> No assert "app.etl.mock_api_import" in commande -def test_mock_api_import_uses_the_airflow_data_interval(dagbag: DagBag) -> None: +def test_mock_api_import_asks_for_the_on_the_hour_reading(dagbag: DagBag) -> None: commande = dagbag.dags["mock_api_import"].get_task("import_mock_api").bash_command - assert "--start-time \"{{ data_interval_start.strftime('%Y-%m-%dT%H:%M:%S') }}\"" in commande + assert "--start-time \"{{ data_interval_end.strftime('%Y-%m-%dT%H:00:00') }}\"" in commande assert "--end-time \"{{ data_interval_end.strftime('%Y-%m-%dT%H:%M:%S') }}\"" in commande assert commande.endswith("--limit 1")