fix(etl): traite les retours de revue du DAG historique
Airflow / Lint et intégrité des DAGs (push) Successful in 57s
Airflow / Construction de l'image (push) Successful in 2m52s
SonarQube / build-back (push) Successful in 1m17s
SonarQube / test-ml (push) Failing after 1m50s
SonarQube / build-front (push) Successful in 9m44s
SonarQube / test-back (push) Failing after 55s
SonarQube / test-front (push) Failing after 5m3s
SonarQube / SonarQube (push) Skipped
Airflow / Lint et intégrité des DAGs (push) Successful in 57s
Airflow / Construction de l'image (push) Successful in 2m52s
SonarQube / build-back (push) Successful in 1m17s
SonarQube / test-ml (push) Failing after 1m50s
SonarQube / build-front (push) Successful in 9m44s
SonarQube / test-back (push) Failing after 55s
SonarQube / test-front (push) Failing after 5m3s
SonarQube / SonarQube (push) Skipped
This commit is contained in:
+12
-186
@@ -639,7 +639,7 @@ La suite backend complète a également été validée avec une couverture supé
|
||||
|
||||
## Suite du pipeline Data
|
||||
|
||||
Deux sources de données sont prises en charge par la logique ETL du backend :
|
||||
Deux sources de données sont maintenant prises en charge :
|
||||
|
||||
```text
|
||||
Dataset CSV/JSON
|
||||
@@ -661,192 +661,18 @@ mock_api_import.py
|
||||
API Mock
|
||||
```
|
||||
|
||||
La logique d'extraction, de validation, de transformation et de chargement est disponible pour
|
||||
les deux sources de données du MVP.
|
||||
La logique d'extraction, de transformation et de chargement est donc disponible pour les deux sources de données du MVP.
|
||||
|
||||
Airflow tourne réellement dans `etl/airflow/` et orchestre désormais quatre DAGs :
|
||||
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).
|
||||
|
||||
```text
|
||||
ml_train
|
||||
ml_score
|
||||
alertes
|
||||
historical_import
|
||||
```
|
||||
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.
|
||||
|
||||
Les DAGs `ml_train` et `ml_score` orchestrent le pipeline Machine Learning (issue #115).
|
||||
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 `alertes` orchestre la détection des alertes et la génération des recommandations
|
||||
(issue #116).
|
||||
|
||||
Le DAG `historical_import` orchestre l'import du dataset historique CSV/JSON (issue #119).
|
||||
|
||||
### Orchestration de l'import historique
|
||||
|
||||
Le DAG historique est défini dans :
|
||||
|
||||
```text
|
||||
etl/airflow/dags/historical_import.py
|
||||
```
|
||||
|
||||
Il ne réimplémente aucune logique ETL. Il déclenche directement le module existant :
|
||||
|
||||
```text
|
||||
app.etl.historical_import
|
||||
```
|
||||
|
||||
Le flux d'exécution est le suivant :
|
||||
|
||||
```text
|
||||
data/raw/
|
||||
├── all_sites_combined.csv
|
||||
└── dataset_metadata.json
|
||||
|
|
||||
v
|
||||
Airflow
|
||||
|
|
||||
v
|
||||
DAG historical_import
|
||||
|
|
||||
v
|
||||
BashOperator
|
||||
|
|
||||
v
|
||||
app.etl.historical_import
|
||||
|
|
||||
v
|
||||
PostgreSQL / TimescaleDB
|
||||
|
|
||||
+--> dataset
|
||||
+--> site
|
||||
+--> reading
|
||||
```
|
||||
|
||||
Les fichiers historiques locaux sont montés dans les conteneurs Airflow en lecture seule :
|
||||
|
||||
```text
|
||||
./data/raw:/opt/data/raw:ro
|
||||
```
|
||||
|
||||
Le DAG utilise les chemins suivants :
|
||||
|
||||
```text
|
||||
/opt/data/raw/all_sites_combined.csv
|
||||
/opt/data/raw/dataset_metadata.json
|
||||
```
|
||||
|
||||
Le montage en lecture seule évite qu'un traitement Airflow puisse modifier les fichiers sources.
|
||||
|
||||
Le backend est déjà embarqué dans l'image Airflow dans son propre environnement Python :
|
||||
|
||||
```text
|
||||
/opt/backend/.venv
|
||||
```
|
||||
|
||||
Le DAG utilise un `BashOperator` avec le même principe que le DAG `alertes` :
|
||||
|
||||
```text
|
||||
cd /opt/backend
|
||||
env -u VIRTUAL_ENV
|
||||
uv run --no-sync python -m app.etl.historical_import
|
||||
```
|
||||
|
||||
Airflow reste ainsi responsable de l'orchestration tandis que le backend reste responsable de
|
||||
l'extraction, de la validation, de la transformation et du chargement.
|
||||
|
||||
### Planification
|
||||
|
||||
Le dataset historique sert à initialiser l'environnement et n'est pas une source périodique.
|
||||
|
||||
Le DAG est donc configuré avec :
|
||||
|
||||
```text
|
||||
schedule = None
|
||||
catchup = False
|
||||
max_active_runs = 1
|
||||
```
|
||||
|
||||
Le déclenchement est manuel depuis l'interface Airflow ou avec la CLI.
|
||||
|
||||
`max_active_runs = 1` empêche deux imports historiques de s'exécuter simultanément.
|
||||
|
||||
La tâche `import_historical` définit également :
|
||||
|
||||
```text
|
||||
retries = 1
|
||||
retry_delay = 2 minutes
|
||||
execution_timeout = 30 minutes
|
||||
```
|
||||
|
||||
Le retry permet de reprendre le traitement après une erreur transitoire, notamment une
|
||||
indisponibilité temporaire de PostgreSQL.
|
||||
|
||||
Le pipeline historique étant idempotent, une nouvelle exécution ne doit pas dupliquer les mesures
|
||||
déjà présentes.
|
||||
|
||||
### Déclenchement et suivi
|
||||
|
||||
Le DAG peut être déclenché avec :
|
||||
|
||||
```powershell
|
||||
docker compose exec airflow-scheduler airflow dags trigger historical_import
|
||||
```
|
||||
|
||||
Les exécutions peuvent être consultées avec :
|
||||
|
||||
```powershell
|
||||
docker compose exec airflow-scheduler airflow dags list-runs -d historical_import
|
||||
```
|
||||
|
||||
### Validation de l'orchestration
|
||||
|
||||
L'orchestration a été validée localement avec Docker Compose et le `LocalExecutor` Airflow.
|
||||
|
||||
Le scheduler détecte les quatre DAGs :
|
||||
|
||||
```text
|
||||
alertes
|
||||
historical_import
|
||||
ml_score
|
||||
ml_train
|
||||
```
|
||||
|
||||
Deux exécutions manuelles successives du DAG `historical_import` ont été réalisées.
|
||||
|
||||
Les deux exécutions se sont terminées avec :
|
||||
|
||||
```text
|
||||
state = success
|
||||
```
|
||||
|
||||
Les exécutions ont été traitées séquentiellement.
|
||||
|
||||
Après les deux exécutions, PostgreSQL/TimescaleDB contenait toujours :
|
||||
|
||||
```text
|
||||
122647
|
||||
```
|
||||
|
||||
lectures historiques avec :
|
||||
|
||||
```text
|
||||
source = "csv"
|
||||
```
|
||||
|
||||
Le nombre de lectures n'a donc pas doublé après la seconde exécution.
|
||||
|
||||
Cette validation confirme que l'orchestration Airflow réutilise correctement
|
||||
`app.etl.historical_import` et conserve l'idempotence du pipeline historique.
|
||||
|
||||
### Suite
|
||||
|
||||
L'import de l'API Mock est déjà disponible côté backend avec :
|
||||
|
||||
```text
|
||||
app.etl.mock_api_import
|
||||
```
|
||||
|
||||
Son orchestration Airflow ainsi que la réconciliation globale entre les sources historique et
|
||||
API Mock restent à compléter dans l'issue #15.
|
||||
|
||||
Airflow ne remplace pas les pipelines Python existants : il orchestre leur exécution, leur
|
||||
planification, les reprises sur erreur et leur suivi.
|
||||
Le pipeline Data servira ensuite à préparer les données nécessaires au modèle de Machine Learning.
|
||||
|
||||
@@ -18,6 +18,8 @@ COMMANDE_BACKEND = "cd /opt/backend && env -u VIRTUAL_ENV uv run --no-sync pytho
|
||||
|
||||
CSV_PATH = "/opt/data/raw/all_sites_combined.csv"
|
||||
METADATA_PATH = "/opt/data/raw/dataset_metadata.json"
|
||||
SOURCE_TIMEZONE = "UTC"
|
||||
BATCH_SIZE = 1000
|
||||
|
||||
with DAG(
|
||||
dag_id="historical_import",
|
||||
|
||||
@@ -1,8 +1,5 @@
|
||||
"""Tests d'integrite des DAGs : s'importent sans erreur et ont la structure attendue.
|
||||
|
||||
Les tests ne lancent pas reellement les traitements : ils valident uniquement la definition
|
||||
des DAGs et leurs commandes.
|
||||
"""
|
||||
"""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 pathlib import Path
|
||||
@@ -14,7 +11,6 @@ from airflow.models.dagbag import DagBag
|
||||
DAGS_FOLDER = Path(__file__).resolve().parent.parent / "dags"
|
||||
|
||||
DAG_IDS = ["ml_train", "ml_score", "alertes", "historical_import"]
|
||||
|
||||
TACHES = [
|
||||
("ml_train", "train"),
|
||||
("ml_score", "score"),
|
||||
@@ -42,10 +38,13 @@ def test_ml_train_has_no_schedule(dagbag: DagBag) -> None:
|
||||
|
||||
|
||||
def test_ml_score_runs_every_hour(dagbag: DagBag) -> None:
|
||||
# `@hourly` est un alias Airflow pour ce cron, c'est sous cette forme que `.summary` le rend.
|
||||
assert dagbag.dags["ml_score"].timetable.summary == "0 * * * *"
|
||||
|
||||
|
||||
def test_alertes_runs_after_the_hourly_scoring(dagbag: DagBag) -> None:
|
||||
# Le decalage n'est pas cosmetique : la regle `anomaly` compare une lecture a la `prediction`
|
||||
# du meme instant, que `ml_score` ecrit a l'heure pile.
|
||||
assert dagbag.dags["alertes"].timetable.summary == "15 * * * *"
|
||||
|
||||
|
||||
@@ -87,6 +86,7 @@ def test_historical_import_uses_the_expected_source_files(dagbag: DagBag) -> Non
|
||||
|
||||
@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).
|
||||
assert "/opt/backend" in dagbag.dags["alertes"].get_task(task_id).bash_command
|
||||
|
||||
|
||||
@@ -96,6 +96,8 @@ def test_historical_import_runs_in_the_backend_environment(dagbag: DagBag) -> No
|
||||
|
||||
|
||||
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.
|
||||
assert dagbag.dags["alertes"].get_task("detection").downstream_task_ids == {"recommandations"}
|
||||
|
||||
|
||||
@@ -110,11 +112,14 @@ def test_ml_score_reuses_the_model_path_written_by_ml_train(dagbag: DagBag) -> N
|
||||
|
||||
@pytest.mark.parametrize("dag_id", DAG_IDS)
|
||||
def test_no_two_runs_of_a_dag_overlap(dagbag: DagBag, dag_id: str) -> None:
|
||||
# Deux entrainements ecriraient le meme fichier modele, deux scorings inseriraient en meme
|
||||
# temps dans `prediction`, deux detections analyseraient la meme fenetre.
|
||||
assert dagbag.dags[dag_id].max_active_runs == 1
|
||||
|
||||
|
||||
@pytest.mark.parametrize(("dag_id", "task_id"), TACHES)
|
||||
def test_every_task_has_an_execution_timeout(dagbag: DagBag, dag_id: str, task_id: str) -> None:
|
||||
# Sans plafond, une connexion pendue immobilise un slot du scheduler indefiniment.
|
||||
assert dagbag.dags[dag_id].get_task(task_id).execution_timeout is not None
|
||||
|
||||
|
||||
@@ -125,11 +130,15 @@ def test_ml_score_execution_timeout_stays_below_its_hourly_step(dagbag: DagBag)
|
||||
|
||||
|
||||
def duree_au_pire(tache: BaseOperator) -> timedelta:
|
||||
# `execution_timeout` plafonne une tentative, pas la tache : deux reprises occupent trois
|
||||
# plafonds et deux delais d'attente.
|
||||
assert tache.execution_timeout is not None
|
||||
return (tache.retries + 1) * tache.execution_timeout + tache.retries * tache.retry_delay
|
||||
|
||||
|
||||
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.
|
||||
taches = [
|
||||
dagbag.dags["alertes"].get_task(task_id) for task_id in ("detection", "recommandations")
|
||||
]
|
||||
@@ -142,6 +151,7 @@ def test_ml_score_retries_after_a_transient_failure(dagbag: DagBag) -> None:
|
||||
|
||||
@pytest.mark.parametrize("task_id", ["detection", "recommandations"])
|
||||
def test_alertes_retries_after_a_transient_failure(dagbag: DagBag, task_id: str) -> None:
|
||||
# Les deux commandes sont idempotentes en base, une reprise ne duplique rien.
|
||||
assert dagbag.dags["alertes"].get_task(task_id).retries >= 1
|
||||
|
||||
|
||||
@@ -153,4 +163,5 @@ def test_historical_import_retries_after_a_transient_failure(dagbag: DagBag) ->
|
||||
def test_tasks_never_resync_the_baked_environment(
|
||||
dagbag: DagBag, dag_id: str, task_id: str
|
||||
) -> None:
|
||||
# Sans `--no-sync`, `uv run` reconstruit le projet a chaque execution.
|
||||
assert "--no-sync" in dagbag.dags[dag_id].get_task(task_id).bash_command
|
||||
|
||||
Reference in New Issue
Block a user