fix(etl): demander la mesure de l'heure pile à l'API Mock
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.
This commit is contained in:
@@ -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` |
|
| `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` |
|
| `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` |
|
| `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 |
|
| `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.
|
Le DAG `historical_import` réutilise le pipeline historique existant sans dupliquer sa logique.
|
||||||
|
|||||||
+5
-5
@@ -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).
|
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 d'une lecture
|
`app.etl.mock_api_import` sur l'intervalle qui va de l'heure pile à son déclenchement, avec une
|
||||||
par site. L'API Mock génère autant de points que la limite demandée, répartis sur l'intervalle :
|
limite d'une lecture par site. L'API Mock génère autant de points que la limite demandée,
|
||||||
un seul donne la mesure de :00, au pas horaire du dataset historique, que les features ML
|
répartis sur l'intervalle et le premier à son début : un seul donne la mesure de :00, au pas
|
||||||
supposent en décalant par ligne. Les deux pipelines normalisent leurs données vers les tables communes `site` et
|
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.
|
||||||
|
|
||||||
Le DAG `mock_api_import` exécute `app.etl.mock_api_import` toutes les heures. Chaque exécution
|
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`).
|
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
|
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
|
||||||
|
|||||||
@@ -1,7 +1,7 @@
|
|||||||
"""DAG d'import périodique des données de l'API Mock EnerVision (issue #15).
|
"""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.
|
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
|
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.
|
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).
|
# 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"
|
||||||
|
|
||||||
# Contrainte : l'API Mock génère `limit` points répartis sur l'intervalle. Un seul donne la mesure
|
# Contrainte : l'API Mock génère `limit` points répartis sur l'intervalle, le premier à son début.
|
||||||
# de :00, au pas horaire du CSV que les features ML supposent (`shift(168)` compte des lignes).
|
# Un seul, depuis l'heure pile, donne la mesure de :00 au pas du CSV que suppose `shift(168)`.
|
||||||
LIMITE_LECTURES = 1
|
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
|
||||||
@@ -51,7 +51,7 @@ with DAG(
|
|||||||
task_id="import_mock_api",
|
task_id="import_mock_api",
|
||||||
bash_command=(
|
bash_command=(
|
||||||
f"{COMMANDE_BACKEND} app.etl.mock_api_import "
|
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') }}\" "
|
"--end-time \"{{ data_interval_end.strftime('%Y-%m-%dT%H:%M:%S') }}\" "
|
||||||
f"--limit {LIMITE_LECTURES}"
|
f"--limit {LIMITE_LECTURES}"
|
||||||
),
|
),
|
||||||
|
|||||||
@@ -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
|
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
|
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 "--end-time \"{{ data_interval_end.strftime('%Y-%m-%dT%H:%M:%S') }}\"" in commande
|
||||||
assert commande.endswith("--limit 1")
|
assert commande.endswith("--limit 1")
|
||||||
|
|
||||||
|
|||||||
Reference in New Issue
Block a user