refactor(etl): migre les DAGs et leurs tests vers le SDK Airflow 3
DAG et BaseOperator viennent d'airflow.sdk, BashOperator du provider standard (airflow.operators.bash n'est plus qu'un alias déprécié). DagBag s'importe depuis airflow.dag_processing et ne prend plus include_examples, les exemples étant déjà coupés par la configuration de conftest.py. Les timetables n'exposent plus summary : la planification se lit sur dag.schedule (None) et timetable.expression (cron normalisé).
This commit is contained in:
@@ -15,8 +15,8 @@ from __future__ import annotations
|
|||||||
|
|
||||||
from datetime import datetime, timedelta
|
from datetime import datetime, timedelta
|
||||||
|
|
||||||
from airflow.models.dag import DAG
|
from airflow.providers.standard.operators.bash import BashOperator
|
||||||
from airflow.operators.bash import BashOperator
|
from airflow.sdk import DAG
|
||||||
|
|
||||||
# Le backend a son propre environnement uv dans l'image (ADR 0008). `--no-sync` et
|
# Le backend a son propre environnement uv dans l'image (ADR 0008). `--no-sync` et
|
||||||
# `env -u VIRTUAL_ENV` : cf. `ml_train.py`, même raisonnement.
|
# `env -u VIRTUAL_ENV` : cf. `ml_train.py`, même raisonnement.
|
||||||
|
|||||||
@@ -10,8 +10,8 @@ from __future__ import annotations
|
|||||||
|
|
||||||
from datetime import datetime, timedelta
|
from datetime import datetime, timedelta
|
||||||
|
|
||||||
from airflow.models.dag import DAG
|
from airflow.providers.standard.operators.bash import BashOperator
|
||||||
from airflow.operators.bash import BashOperator
|
from airflow.sdk import DAG
|
||||||
|
|
||||||
MODEL_PATH = "/opt/ml/state/models/lightgbm-consumption.txt"
|
MODEL_PATH = "/opt/ml/state/models/lightgbm-consumption.txt"
|
||||||
|
|
||||||
|
|||||||
@@ -11,8 +11,8 @@ from __future__ import annotations
|
|||||||
|
|
||||||
from datetime import datetime, timedelta
|
from datetime import datetime, timedelta
|
||||||
|
|
||||||
from airflow.models.dag import DAG
|
from airflow.providers.standard.operators.bash import BashOperator
|
||||||
from airflow.operators.bash import BashOperator
|
from airflow.sdk import DAG
|
||||||
|
|
||||||
MODEL_PATH = "/opt/ml/state/models/lightgbm-consumption.txt"
|
MODEL_PATH = "/opt/ml/state/models/lightgbm-consumption.txt"
|
||||||
MLFLOW_TRACKING_URI = "sqlite:////opt/ml/state/mlflow.db"
|
MLFLOW_TRACKING_URI = "sqlite:////opt/ml/state/mlflow.db"
|
||||||
|
|||||||
@@ -5,8 +5,8 @@ from datetime import timedelta
|
|||||||
from pathlib import Path
|
from pathlib import Path
|
||||||
|
|
||||||
import pytest
|
import pytest
|
||||||
from airflow.models.baseoperator import BaseOperator
|
from airflow.dag_processing.dagbag import DagBag
|
||||||
from airflow.models.dagbag import DagBag
|
from airflow.sdk import BaseOperator
|
||||||
|
|
||||||
DAGS_FOLDER = Path(__file__).resolve().parent.parent / "dags"
|
DAGS_FOLDER = Path(__file__).resolve().parent.parent / "dags"
|
||||||
|
|
||||||
@@ -21,7 +21,7 @@ TACHES = [
|
|||||||
|
|
||||||
@pytest.fixture(scope="module")
|
@pytest.fixture(scope="module")
|
||||||
def dagbag() -> DagBag:
|
def dagbag() -> DagBag:
|
||||||
return DagBag(dag_folder=str(DAGS_FOLDER), include_examples=False)
|
return DagBag(dag_folder=str(DAGS_FOLDER))
|
||||||
|
|
||||||
|
|
||||||
def test_dags_folder_has_no_import_error(dagbag: DagBag) -> None:
|
def test_dags_folder_has_no_import_error(dagbag: DagBag) -> None:
|
||||||
@@ -33,18 +33,18 @@ def test_every_expected_dag_is_discovered(dagbag: DagBag) -> None:
|
|||||||
|
|
||||||
|
|
||||||
def test_ml_train_has_no_schedule(dagbag: DagBag) -> None:
|
def test_ml_train_has_no_schedule(dagbag: DagBag) -> None:
|
||||||
assert dagbag.dags["ml_train"].timetable.summary == "None"
|
assert dagbag.dags["ml_train"].schedule is None
|
||||||
|
|
||||||
|
|
||||||
def test_ml_score_runs_every_hour(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.
|
# `@hourly` est un alias Airflow pour ce cron, c'est sous cette forme que la timetable le rend.
|
||||||
assert dagbag.dags["ml_score"].timetable.summary == "0 * * * *"
|
assert dagbag.dags["ml_score"].timetable.expression == "0 * * * *"
|
||||||
|
|
||||||
|
|
||||||
def test_alertes_runs_after_the_hourly_scoring(dagbag: DagBag) -> None:
|
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`
|
# 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.
|
# du meme instant, que `ml_score` ecrit a l'heure pile.
|
||||||
assert dagbag.dags["alertes"].timetable.summary == "15 * * * *"
|
assert dagbag.dags["alertes"].timetable.expression == "15 * * * *"
|
||||||
|
|
||||||
|
|
||||||
def test_ml_train_task_calls_the_training_module(dagbag: DagBag) -> None:
|
def test_ml_train_task_calls_the_training_module(dagbag: DagBag) -> None:
|
||||||
|
|||||||
Reference in New Issue
Block a user