Fusionne dev dans test/integration-api-db-ml
Quatre conflits, tous additifs, nés du DAG `mock_api_import` (#146) arrivé sur `dev` pendant que cette branche ajoutait `derive` : la liste des DAGs du README, celle de la vue d'ensemble et du tableau d'infrastructure, et `DAG_IDS`/`TACHES` dans les tests d'intégrité. Les six DAGs sont conservés de part et d'autre. Collision que git ne voyait pas : `dev` a reçu un ADR 0011 et un 0012 (procédure de déploiement, état de la VM ENI) pendant que cette branche en ajoutait un autre sous le même numéro. L'ADR de la surveillance de dérive devient 0013, avec ses onze références, et la table de `docs/README.md` reprend les trois.
This commit is contained in:
+15
-9
@@ -663,16 +663,22 @@ mock_api_import.py
|
||||
|
||||
La logique d'extraction, de transformation et de chargement est donc disponible pour les deux sources de données du MVP.
|
||||
|
||||
Airflow tourne désormais réellement (`etl/airflow/`, `make airflow-up`) et orchestre le pipeline
|
||||
ML (`ml_train`/`ml_score`, issue #115), la détection d'alertes et la génération des
|
||||
recommandations (`alertes`, issue #116), ainsi que l'import historique
|
||||
(`historical_import`, issue #119).
|
||||
Airflow tourne désormais réellement (`etl/airflow/`, `make airflow-up`) et orchestre cinq DAGs :
|
||||
le pipeline ML (`ml_train` et `ml_score`, issue #115), la détection d'alertes et la génération
|
||||
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 `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.
|
||||
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
|
||||
`reading`, tout en conservant leur source (`csv` ou `api_history`). La réconciliation globale
|
||||
des deux sources reste à compléter dans l'issue #15.
|
||||
|
||||
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.
|
||||
|
||||
@@ -0,0 +1,61 @@
|
||||
"""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.
|
||||
|
||||
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
|
||||
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
|
||||
|
||||
# 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)
|
||||
|
||||
# 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.
|
||||
PLANIFICATION = CronTriggerTimetable(
|
||||
"45 * * * *",
|
||||
timezone="UTC",
|
||||
interval=timedelta(hours=1),
|
||||
)
|
||||
|
||||
with DAG(
|
||||
dag_id="mock_api_import",
|
||||
description="Importe chaque heure les données de l'API Mock dans site et reading.",
|
||||
schedule=PLANIFICATION,
|
||||
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.strftime('%Y-%m-%dT%H:%M:%S') }}\" "
|
||||
"--end-time \"{{ data_interval_end.strftime('%Y-%m-%dT%H:%M:%S') }}\" "
|
||||
f"--limit {LIMITE_LECTURES}"
|
||||
),
|
||||
retries=NOMBRE_REPRISES,
|
||||
retry_delay=DELAI_ENTRE_REPRISES,
|
||||
execution_timeout=PLAFOND_PAR_TENTATIVE,
|
||||
)
|
||||
@@ -1,22 +1,31 @@
|
||||
"""Tests d'integrite des DAGs : s'importent sans erreur, structure attendue. Pas d'execution
|
||||
reelle des taches (ca reclamerait le conteneur avec `uv`/`enervision_ml`), juste la definition."""
|
||||
|
||||
from datetime import timedelta
|
||||
from datetime import datetime, timedelta
|
||||
from pathlib import Path
|
||||
|
||||
import pytest
|
||||
from airflow.dag_processing.dagbag import DagBag
|
||||
from airflow.sdk import BaseOperator
|
||||
from airflow.timetables.trigger import CronTriggerTimetable
|
||||
|
||||
DAGS_FOLDER = Path(__file__).resolve().parent.parent / "dags"
|
||||
|
||||
DAG_IDS = ["ml_train", "ml_score", "alertes", "historical_import", "derive"]
|
||||
DAG_IDS = [
|
||||
"ml_train",
|
||||
"ml_score",
|
||||
"alertes",
|
||||
"historical_import",
|
||||
"mock_api_import",
|
||||
"derive",
|
||||
]
|
||||
TACHES = [
|
||||
("ml_train", "train"),
|
||||
("ml_score", "score"),
|
||||
("alertes", "detection"),
|
||||
("alertes", "recommandations"),
|
||||
("historical_import", "import_historical"),
|
||||
("mock_api_import", "import_mock_api"),
|
||||
("derive", "derive"),
|
||||
]
|
||||
|
||||
@@ -53,6 +62,19 @@ def test_historical_import_has_no_schedule(dagbag: DagBag) -> None:
|
||||
assert dagbag.dags["historical_import"].schedule is None
|
||||
|
||||
|
||||
def test_mock_api_import_uses_an_explicit_hourly_interval(dagbag: DagBag) -> None:
|
||||
timetable = dagbag.dags["mock_api_import"].timetable
|
||||
|
||||
assert isinstance(timetable, CronTriggerTimetable)
|
||||
assert timetable.serialize()["expression"] == "45 * * * *"
|
||||
|
||||
manual_interval = timetable.infer_manual_data_interval(
|
||||
run_after=datetime.fromisoformat("2026-09-22T12:30:00+00:00"),
|
||||
)
|
||||
|
||||
assert manual_interval.end - manual_interval.start == timedelta(hours=1)
|
||||
|
||||
|
||||
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
|
||||
@@ -85,6 +107,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.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
|
||||
|
||||
|
||||
@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).
|
||||
@@ -96,6 +132,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.
|
||||
@@ -137,6 +179,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.
|
||||
@@ -160,6 +210,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
|
||||
|
||||
|
||||
def test_derive_runs_once_a_day(dagbag: DagBag) -> None:
|
||||
assert dagbag.dags["derive"].timetable.expression == "30 5 * * *"
|
||||
|
||||
|
||||
Reference in New Issue
Block a user