Merge remote-tracking branch 'origin/dev' into feat/reconciliation-dag
This commit is contained in:
+12
-14
@@ -668,20 +668,18 @@ le pipeline ML (`ml_train` et `ml_score`, issue #115), la détection d'alertes e
|
||||
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 `mock_api_import` s'exécute chaque heure, à la minute `:45`, sur un intervalle explicite
|
||||
d'une heure. L'API Mock génère autant de points que la limite demandée, répartis sur
|
||||
l'intervalle : `app.etl.mock_api_import.limit_for_window()` dérive donc `limit` de la fenêtre
|
||||
reçue (une lecture par site pour cette fenêtre d'1h) plutôt que de dépendre d'une valeur fixée à
|
||||
la main côté DAG, et refuse une fenêtre qui ne couvre pas un nombre entier d'heures. La fenêtre
|
||||
`[:45, :45)` place cette lecture à :45, pas à :00 (l'API place son premier point au début de la
|
||||
fenêtre demandée), un décalage constant sans effet sur les lags positionnels ML ni sur les
|
||||
jointures en aval. 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 entre les
|
||||
deux sources (issue #15) est close : voir `docs/architecture/40-data.md`.
|
||||
|
||||
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`).
|
||||
Le DAG `mock_api_import` s'exécute chaque heure, à la minute `:45`, sur une fenêtre qui part de
|
||||
l'heure pile précédant son déclenchement jusqu'à l'instant du déclenchement lui-même (pas
|
||||
l'intervalle Airflow `data_interval_start`/`end` tel quel). L'API Mock génère autant de points que
|
||||
la limite demandée, répartis sur la fenêtre et le premier à son début :
|
||||
`app.etl.mock_api_import.limit_for_window()` dérive donc `limit` de la fenêtre reçue (une seule
|
||||
lecture ici, ancrée sur l'heure pile) plutôt que de dépendre d'une valeur fixée à la main côté
|
||||
DAG, et refuse une fenêtre qui ne démarre pas pile sur l'heure. Une fenêtre calée sur l'intervalle
|
||||
Airflow tel quel (`[:45, :45)`) placerait cette lecture à :45, hors de la grille horaire du reste
|
||||
du schéma (vérifié empiriquement contre l'API Mock) ; partir de l'heure pile évite ce décalage.
|
||||
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 entre les deux sources
|
||||
(issue #15) est close : voir `docs/architecture/40-data.md`.
|
||||
|
||||
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).
|
||||
|
||||
@@ -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.
|
||||
@@ -26,11 +26,11 @@ 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. Conséquence vérifiée empiriquement
|
||||
# sur l'API Mock (cf. `limit_for_window()` dans `app.etl.mock_api_import`) : la fenêtre importée
|
||||
# est `[:45, :45)`, donc chaque lecture atterrit à :45, pas à :00, un décalage constant sur
|
||||
# toute la série, sans effet sur les lags positionnels ni sur les jointures en aval (`ml_score`
|
||||
# vise la dernière lecture + 1h, `derive` joint à l'égalité), cf. 40-data.md.
|
||||
# exécuté à l'heure pile, puis avant `alertes`, exécuté à :15. La fenêtre demandée à l'API Mock
|
||||
# (voir `bash_command` ci-dessous) ne suit pas cet intervalle Airflow tel quel : elle part de
|
||||
# l'heure pile qui précède le déclenchement, pas de `data_interval_start`, pour que l'unique
|
||||
# lecture demandée (`app.etl.mock_api_import.limit_for_window()`) atterrisse à :00 et non à :45
|
||||
# (vérifié empiriquement sur l'API Mock), au pas horaire du reste du schéma, cf. 40-data.md.
|
||||
PLANIFICATION = CronTriggerTimetable(
|
||||
"45 * * * *",
|
||||
timezone="UTC",
|
||||
@@ -51,12 +51,16 @@ 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` part de l'heure pile qui précède le déclenchement, pas de
|
||||
# `data_interval_start` : sur `[:45, :45)`, l'API aurait placé son unique lecture
|
||||
# à :45, hors de la grille horaire du reste du schéma (vérifié empiriquement).
|
||||
"--start-time \"{{ data_interval_end.strftime('%Y-%m-%dT%H:00:00') }}\" "
|
||||
"--end-time \"{{ data_interval_end.strftime('%Y-%m-%dT%H:%M:%S') }}\""
|
||||
# Pas de --limit : app.etl.mock_api_import.limit_for_window() le dérive de la
|
||||
# fenêtre (une lecture/heure), et refuse une fenêtre qui ne couvre pas un nombre
|
||||
# entier d'heures ou dépasse le plafond de l'API. Porter la règle dans le code,
|
||||
# pas dans ce DAG, évite qu'un appel manuel oublie de la respecter.
|
||||
# fenêtre (ici plus courte qu'une heure, donc une seule lecture, ancrée sur
|
||||
# --start-time) et refuse une fenêtre qui ne démarre pas pile sur l'heure. Porter
|
||||
# la règle dans le code, pas dans ce DAG, évite qu'un appel manuel oublie de la
|
||||
# respecter.
|
||||
),
|
||||
retries=NOMBRE_REPRISES,
|
||||
retry_delay=DELAI_ENTRE_REPRISES,
|
||||
|
||||
@@ -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
|
||||
# Pas de --limit : app.etl.mock_api_import.limit_for_window() le dérive de la fenêtre.
|
||||
assert "--limit" not in commande
|
||||
|
||||
Reference in New Issue
Block a user