feat(etl,ml): orchestre l'entrainement et le scoring LightGBM via deux DAGs Airflow
This commit is contained in:
+2
-2
@@ -344,6 +344,6 @@ uv run ruff check app\etl tests\etl
|
||||
|
||||
L'import historique constitue la première brique du pipeline Data EnerVision.
|
||||
|
||||
La prochaine étape consiste à orchestrer les traitements ETL avec Apache Airflow, puis à préparer les données nécessaires à l'entraînement du modèle de Machine Learning.
|
||||
Airflow tourne désormais réellement (`etl/airflow/`, `make airflow-up`), mais orchestre pour l'instant le pipeline ML (`ml_train`/`ml_score`, issue #115), pas encore ce pipeline ETL : orchestrer `historical_import.py` (normalisation et chargement micro-batch, issues #15/#16) reste à faire.
|
||||
|
||||
Airflow sera utilisé comme orchestrateur des traitements existants et ne remplacera pas la logique métier déjà implémentée dans le pipeline ETL.
|
||||
Le principe reste le même que documenté à l'origine : Airflow orchestre les traitements existants sans remplacer leur logique métier, cf. `etl/airflow/dags/ml_train.py`/`ml_score.py` pour un exemple concret de ce patron (des `BashOperator` qui invoquent le script tel quel).
|
||||
@@ -0,0 +1,40 @@
|
||||
# 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.
|
||||
FROM apache/airflow:2.10.4-python3.12
|
||||
|
||||
# LightGBM est compile contre libgomp (OpenMP), absent de l'image de base (minimale, sans
|
||||
# toolchain de compilation). Sans lui : `OSError: libgomp.so.1: cannot open shared object file`
|
||||
# au premier `import lightgbm`, seulement au moment ou une tache tourne reellement.
|
||||
USER root
|
||||
RUN apt-get update \
|
||||
&& apt-get install --no-install-recommends -y libgomp1 \
|
||||
&& apt-get clean \
|
||||
&& rm -rf /var/lib/apt/lists/*
|
||||
# Pre-cree, appartenant a `airflow` : docker-compose y monte un volume nomme partage entre
|
||||
# `ml_train` et `ml_score` (le modele ecrit par l'un, lu par l'autre). Un volume nomme herite des
|
||||
# permissions du repertoire qu'il recouvre a son premier montage ; sans ce chown prealable, il
|
||||
# serait cree root:root et illisible par le conteneur, qui tourne en `airflow` (uid 50000).
|
||||
RUN mkdir -p /opt/ml/state && chown -R airflow:root /opt/ml
|
||||
USER airflow
|
||||
|
||||
# L'image de base embarque deja un `uv`, mais trop ancien (0.4.29) pour le format de verrou de
|
||||
# `ml/uv.lock`. On le remplace par la version deja pinnee ailleurs dans le depot
|
||||
# (apps/backend/Dockerfile).
|
||||
COPY --from=ghcr.io/astral-sh/uv:0.11.26 /uv /home/airflow/.local/bin/uv
|
||||
|
||||
ENV UV_COMPILE_BYTECODE=1 \
|
||||
UV_LINK_MODE=copy \
|
||||
UV_PROJECT_ENVIRONMENT=/opt/ml/.venv
|
||||
|
||||
WORKDIR /opt/ml
|
||||
|
||||
COPY --chown=airflow:root ml/pyproject.toml ml/uv.lock ./
|
||||
RUN uv sync --locked --no-install-project --no-dev
|
||||
|
||||
COPY --chown=airflow:root ml/enervision_ml ./enervision_ml
|
||||
RUN uv sync --locked --no-dev
|
||||
|
||||
WORKDIR /opt/airflow
|
||||
@@ -0,0 +1,33 @@
|
||||
"""DAG de scoring horaire du modele LightGBM (issue #115).
|
||||
|
||||
Planifie toutes les heures, au rythme documente par `enervision_ml.score` (score le prochain pas
|
||||
horaire par site). Reutilise le modele ecrit par `ml_train` (DAG separe, declenche a la main) :
|
||||
ce DAG ne reentraine jamais rien. Si aucun modele n'a encore ete entraine, la tache echoue
|
||||
(`FileNotFoundError`) plutot que de rester silencieuse.
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
from datetime import datetime
|
||||
|
||||
from airflow.models.dag import DAG
|
||||
from airflow.operators.bash import BashOperator
|
||||
|
||||
MODEL_PATH = "/opt/ml/state/models/lightgbm-consumption.txt"
|
||||
|
||||
with DAG(
|
||||
dag_id="ml_score",
|
||||
description="Score le prochain pas horaire par site (enervision_ml.score).",
|
||||
schedule="@hourly",
|
||||
start_date=datetime(2026, 1, 1),
|
||||
catchup=False,
|
||||
tags=["ml"],
|
||||
) as dag:
|
||||
# `--frozen --no-dev` : 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 "
|
||||
f"--model {MODEL_PATH}"
|
||||
),
|
||||
)
|
||||
@@ -0,0 +1,36 @@
|
||||
"""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.
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
from datetime import datetime
|
||||
|
||||
from airflow.models.dag import DAG
|
||||
from airflow.operators.bash import BashOperator
|
||||
|
||||
MODEL_PATH = "/opt/ml/state/models/lightgbm-consumption.txt"
|
||||
MLFLOW_TRACKING_URI = "sqlite:////opt/ml/state/mlflow.db"
|
||||
|
||||
with DAG(
|
||||
dag_id="ml_train",
|
||||
description="Entraine le modele LightGBM de prevision de consommation (enervision_ml.train).",
|
||||
schedule=None,
|
||||
start_date=datetime(2026, 1, 1),
|
||||
catchup=False,
|
||||
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.
|
||||
BashOperator(
|
||||
task_id="train",
|
||||
bash_command=(
|
||||
"cd /opt/ml && uv run --frozen --no-dev python -m enervision_ml.train "
|
||||
f"--model-output {MODEL_PATH} --mlflow-tracking-uri {MLFLOW_TRACKING_URI}"
|
||||
),
|
||||
)
|
||||
@@ -0,0 +1,32 @@
|
||||
[project]
|
||||
name = "enervision-airflow"
|
||||
version = "0.1.0"
|
||||
description = "DAGs d'orchestration EnerVision (Airflow)"
|
||||
requires-python = ">=3.12,<3.13"
|
||||
dependencies = [
|
||||
"apache-airflow==2.10.4",
|
||||
]
|
||||
|
||||
[dependency-groups]
|
||||
dev = [
|
||||
"ruff>=0.16.7",
|
||||
"pytest>=9.1.1",
|
||||
]
|
||||
|
||||
[tool.uv]
|
||||
package = false
|
||||
|
||||
[tool.ruff]
|
||||
line-length = 100
|
||||
target-version = "py312"
|
||||
src = ["dags", "tests"]
|
||||
|
||||
[tool.ruff.lint]
|
||||
select = ["E", "W", "F", "I", "N", "UP", "B", "SIM", "RUF"]
|
||||
|
||||
[tool.ruff.format]
|
||||
quote-style = "double"
|
||||
|
||||
[tool.pytest.ini_options]
|
||||
testpaths = ["tests"]
|
||||
addopts = "-q"
|
||||
@@ -0,0 +1,16 @@
|
||||
"""Isole Airflow d'un `~/airflow` reel : `AIRFLOW_HOME` doit etre pose avant le premier `import
|
||||
airflow`, donc ici plutot que dans une fixture (les fixtures s'executent trop tard, apres que les
|
||||
modules de test aient deja importe `airflow`)."""
|
||||
|
||||
import os
|
||||
from pathlib import Path
|
||||
|
||||
_AIRFLOW_HOME = Path(__file__).resolve().parent / ".airflow_home"
|
||||
_AIRFLOW_HOME.mkdir(exist_ok=True)
|
||||
|
||||
os.environ.setdefault("AIRFLOW_HOME", str(_AIRFLOW_HOME))
|
||||
os.environ.setdefault("AIRFLOW__CORE__LOAD_EXAMPLES", "False")
|
||||
os.environ.setdefault("AIRFLOW__CORE__UNIT_TEST_MODE", "True")
|
||||
os.environ.setdefault(
|
||||
"AIRFLOW__DATABASE__SQL_ALCHEMY_CONN", f"sqlite:///{_AIRFLOW_HOME / 'airflow.db'}"
|
||||
)
|
||||
@@ -0,0 +1,52 @@
|
||||
"""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 pathlib import Path
|
||||
|
||||
import pytest
|
||||
from airflow.models.dagbag import DagBag
|
||||
|
||||
DAGS_FOLDER = Path(__file__).resolve().parent.parent / "dags"
|
||||
|
||||
|
||||
@pytest.fixture(scope="module")
|
||||
def dagbag() -> DagBag:
|
||||
return DagBag(dag_folder=str(DAGS_FOLDER), include_examples=False)
|
||||
|
||||
|
||||
def test_dags_folder_has_no_import_error(dagbag: DagBag) -> None:
|
||||
assert dagbag.import_errors == {}
|
||||
|
||||
|
||||
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)
|
||||
assert dagbag.dags["ml_train"].timetable.summary == "None"
|
||||
|
||||
|
||||
def test_ml_score_runs_every_hour() -> 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 * * * *"
|
||||
|
||||
|
||||
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
|
||||
|
||||
|
||||
def test_ml_score_task_calls_the_scoring_module(dagbag: DagBag) -> None:
|
||||
tache = dagbag.dags["ml_score"].get_task("score")
|
||||
assert "enervision_ml.score" in tache.bash_command
|
||||
|
||||
|
||||
def test_ml_score_reuses_the_model_path_written_by_ml_train(dagbag: DagBag) -> None:
|
||||
entrainement = dagbag.dags["ml_train"].get_task("train").bash_command
|
||||
scoring = dagbag.dags["ml_score"].get_task("score").bash_command
|
||||
chemin_modele = "/opt/ml/state/models/lightgbm-consumption.txt"
|
||||
|
||||
assert chemin_modele in entrainement
|
||||
assert chemin_modele in scoring
|
||||
Generated
+1970
File diff suppressed because it is too large
Load Diff
Reference in New Issue
Block a user