fix(etl): importer une seule mesure par heure depuis l'API Mock

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.
This commit is contained in:
Johan LEROY
2026-09-23 15:31:51 +02:00
parent c2f360c591
commit cbbfaf4910
3 changed files with 8 additions and 6 deletions
+4 -2
View File
@@ -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). 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 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 `app.etl.mock_api_import` avec un intervalle explicite d'une heure et une limite d'une lecture
par site. Les deux pipelines normalisent leurs données vers les tables communes `site` et 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 `reading`, tout en conservant leur source (`csv` ou `api_history`). La réconciliation globale
des deux sources reste à compléter dans l'issue #15. des deux sources reste à compléter dans l'issue #15.
+3 -3
View File
@@ -18,9 +18,9 @@ from airflow.timetables.trigger import CronTriggerTimetable
# Le backend possède son propre environnement uv dans l'image Airflow (ADR 0008). # 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" 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. # Contrainte : l'API Mock génère `limit` points répartis sur l'intervalle. Un seul donne la mesure
# Cette marge évite de perdre silencieusement une lecture si une heure en contient plus de 60. # de :00, au pas horaire du CSV que les features ML supposent (`shift(168)` compte des lignes).
LIMITE_LECTURES = 1000 LIMITE_LECTURES = 1
# Deux reprises donnent trois tentatives au total. Même dans le pire cas, l'exécution reste # Deux reprises donnent trois tentatives au total. Même dans le pire cas, l'exécution reste
# inférieure au pas horaire du DAG. # inférieure au pas horaire du DAG.
+1 -1
View File
@@ -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 "--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 "--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"]) @pytest.mark.parametrize("task_id", ["detection", "recommandations"])