fix(etl): fiabilise airflow-init, borne les DAGs ML et ajoute la CI Airflow
This commit is contained in:
@@ -0,0 +1 @@
|
||||
3.12
|
||||
@@ -1,8 +1,7 @@
|
||||
# Image Airflow EnerVision : ajoute le projet ml/ dans son propre environnement Python 3.14,
|
||||
# distinct du Python 3.12 qui fait tourner Airflow lui-meme, pour que les DAGs puissent lancer
|
||||
# `uv run python -m enervision_ml.train`/`.score` en sous-processus (cf. docs/architecture/
|
||||
# 20-backend.md, section Détection d'alertes internes pour le meme raisonnement applique a
|
||||
# app/detection). Airflow ne devient jamais un consommateur direct de LightGBM/MLflow.
|
||||
# distinct du Python 3.12 qui fait tourner Airflow lui-meme (apache-airflow 2.10 ne supporte pas
|
||||
# 3.14), pour que les DAGs puissent lancer `uv run python -m enervision_ml.train`/`.score` en
|
||||
# sous-processus. Airflow ne devient jamais un consommateur direct de LightGBM/MLflow.
|
||||
FROM apache/airflow:2.10.4-python3.12
|
||||
|
||||
# LightGBM est compile contre libgomp (OpenMP), absent de l'image de base (minimale, sans
|
||||
|
||||
@@ -8,7 +8,7 @@ ce DAG ne reentraine jamais rien. Si aucun modele n'a encore ete entraine, la ta
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
from datetime import datetime
|
||||
from datetime import datetime, timedelta
|
||||
|
||||
from airflow.models.dag import DAG
|
||||
from airflow.operators.bash import BashOperator
|
||||
@@ -21,13 +21,21 @@ with DAG(
|
||||
schedule="@hourly",
|
||||
start_date=datetime(2026, 1, 1),
|
||||
catchup=False,
|
||||
# Deux scorings qui se chevauchent inseraient en meme temps dans `prediction` (pas de contrainte
|
||||
# d'unicite sur `(site_id, target_at)`, chaque run garde sa ligne).
|
||||
max_active_runs=1,
|
||||
tags=["ml"],
|
||||
) as dag:
|
||||
# `--frozen --no-dev` : cf. `ml_train.py`, meme raisonnement.
|
||||
# `--no-sync`, `env -u VIRTUAL_ENV` : cf. `ml_train.py`, meme raisonnement.
|
||||
BashOperator(
|
||||
task_id="score",
|
||||
bash_command=(
|
||||
"cd /opt/ml && uv run --frozen --no-dev python -m enervision_ml.score "
|
||||
"cd /opt/ml && env -u VIRTUAL_ENV uv run --no-sync python -m enervision_ml.score "
|
||||
f"--model {MODEL_PATH}"
|
||||
),
|
||||
# Un incident transitoire sur Postgres ne doit pas faire perdre le creneau horaire.
|
||||
retries=2,
|
||||
retry_delay=timedelta(minutes=2),
|
||||
# Bien en dessous du pas horaire : un scoring pendu ne doit pas empieter sur le suivant.
|
||||
execution_timeout=timedelta(minutes=30),
|
||||
)
|
||||
|
||||
@@ -1,14 +1,15 @@
|
||||
"""DAG d'entrainement du modele LightGBM (issue #115).
|
||||
|
||||
Pas de planification : reentrainer est couteux et sa cadence n'est pas une decision prise
|
||||
(cf. `docs/architecture/20-backend.md`). Declenchement manuel depuis l'UI ou la CLI Airflow en
|
||||
attendant. `ml_score` (DAG separe, planifie toutes les heures) reutilise le modele que ce DAG
|
||||
ecrit, il ne reentraine jamais rien lui-meme.
|
||||
Pas de planification : reentrainer est couteux et sa cadence n'est pas une decision prise, en
|
||||
particulier tant que `train.py` ecrase le modele sans comparer ses metriques a l'ancien (cf.
|
||||
`docs/architecture/10-infra.md`, section Airflow). Declenchement manuel depuis l'UI ou la CLI
|
||||
Airflow en attendant. `ml_score` (DAG separe, planifie toutes les heures) reutilise le modele que
|
||||
ce DAG ecrit, il ne reentraine jamais rien lui-meme.
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
from datetime import datetime
|
||||
from datetime import datetime, timedelta
|
||||
|
||||
from airflow.models.dag import DAG
|
||||
from airflow.operators.bash import BashOperator
|
||||
@@ -22,15 +23,21 @@ with DAG(
|
||||
schedule=None,
|
||||
start_date=datetime(2026, 1, 1),
|
||||
catchup=False,
|
||||
# Deux entrainements simultanes ecriraient le meme fichier modele.
|
||||
max_active_runs=1,
|
||||
tags=["ml"],
|
||||
) as dag:
|
||||
# `--frozen --no-dev` : l'environnement `/opt/ml/.venv` est fige a la construction de l'image
|
||||
# (groupe `dev` exclu). Sans `--no-dev` ici, `uv run` resynchronise ruff/mypy a chaque
|
||||
# execution : un acces reseau evitable, sur le chemin d'execution d'une tache planifiee.
|
||||
# `--no-sync` : l'environnement `/opt/ml/.venv` est fige a la construction de l'image, `uv run`
|
||||
# ne le resynchronise pas (sinon `enervision-ml` est reconstruit a chaque tache).
|
||||
# `env -u VIRTUAL_ENV` : l'image de base positionne celui d'Airflow, que `uv` signale a chaque
|
||||
# execution sans qu'il change quoi que ce soit.
|
||||
BashOperator(
|
||||
task_id="train",
|
||||
bash_command=(
|
||||
"cd /opt/ml && uv run --frozen --no-dev python -m enervision_ml.train "
|
||||
"cd /opt/ml && env -u VIRTUAL_ENV uv run --no-sync python -m enervision_ml.train "
|
||||
f"--model-output {MODEL_PATH} --mlflow-tracking-uri {MLFLOW_TRACKING_URI}"
|
||||
),
|
||||
# Un entrainement complet dure quelques minutes ; une connexion pendue ne doit pas
|
||||
# immobiliser un slot du scheduler indefiniment.
|
||||
execution_timeout=timedelta(hours=1),
|
||||
)
|
||||
|
||||
@@ -1,6 +1,7 @@
|
||||
"""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
|
||||
|
||||
import pytest
|
||||
@@ -22,14 +23,12 @@ def test_every_expected_dag_is_discovered(dagbag: DagBag) -> None:
|
||||
assert set(dagbag.dag_ids) == {"ml_train", "ml_score"}
|
||||
|
||||
|
||||
def test_ml_train_has_no_schedule() -> None:
|
||||
dagbag = DagBag(dag_folder=str(DAGS_FOLDER), include_examples=False)
|
||||
def test_ml_train_has_no_schedule(dagbag: DagBag) -> None:
|
||||
assert dagbag.dags["ml_train"].timetable.summary == "None"
|
||||
|
||||
|
||||
def test_ml_score_runs_every_hour() -> 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.
|
||||
dagbag = DagBag(dag_folder=str(DAGS_FOLDER), include_examples=False)
|
||||
assert dagbag.dags["ml_score"].timetable.summary == "0 * * * *"
|
||||
|
||||
|
||||
@@ -50,3 +49,34 @@ def test_ml_score_reuses_the_model_path_written_by_ml_train(dagbag: DagBag) -> N
|
||||
|
||||
assert chemin_modele in entrainement
|
||||
assert chemin_modele in scoring
|
||||
|
||||
|
||||
@pytest.mark.parametrize("dag_id", ["ml_train", "ml_score"])
|
||||
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`.
|
||||
assert dagbag.dags[dag_id].max_active_runs == 1
|
||||
|
||||
|
||||
@pytest.mark.parametrize(("dag_id", "task_id"), [("ml_train", "train"), ("ml_score", "score")])
|
||||
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
|
||||
|
||||
|
||||
def test_ml_score_execution_timeout_stays_below_its_hourly_step(dagbag: DagBag) -> None:
|
||||
timeout = dagbag.dags["ml_score"].get_task("score").execution_timeout
|
||||
assert timeout is not None
|
||||
assert timeout < timedelta(hours=1)
|
||||
|
||||
|
||||
def test_ml_score_retries_after_a_transient_failure(dagbag: DagBag) -> None:
|
||||
assert dagbag.dags["ml_score"].get_task("score").retries >= 1
|
||||
|
||||
|
||||
@pytest.mark.parametrize(("dag_id", "task_id"), [("ml_train", "train"), ("ml_score", "score")])
|
||||
def test_tasks_never_resync_the_baked_environment(
|
||||
dagbag: DagBag, dag_id: str, task_id: str
|
||||
) -> None:
|
||||
# Sans `--no-sync`, `uv run` reconstruit `enervision-ml` a chaque execution.
|
||||
assert "--no-sync" in dagbag.dags[dag_id].get_task(task_id).bash_command
|
||||
|
||||
Reference in New Issue
Block a user