fix(etl): borne le DAG alertes sur son pire cas et couvre sa seconde commande en CI
Airflow / Lint et intégrité des DAGs (push) Successful in 56s
Airflow / Construction de l'image (push) Successful in 3m20s

`execution_timeout` plafonne une tentative, pas la tâche. Avec deux reprises, quinze
minutes par tentative autorisaient quarante-neuf minutes par tâche et quatre-vingt-dix-huit
pour l'enchaînement, quand le commentaire annonçait une somme tenant sous le pas horaire.
Le plafond passe à cinq minutes, ce qui borne le pire cas à trente-huit minutes, et le test
d'intégrité calcule désormais ce pire cas plutôt que la somme des plafonds : reprises et
délais d'attente compris, c'est la durée qu'un `max_active_runs=1` fait payer à l'exécution
suivante.

La CI vérifie aussi `app.cli generate-recommendations --help` sans réseau. C'est la seconde
commande du DAG, et son import tire FastAPI, les repositories et les services, donc une part
de l'environnement `/opt/backend` que la détection seule ne touche pas.

`10-infra.md` nomme enfin ce que le décalage de quinze minutes ne garantit pas : le plafond
de `ml_score` valant trente minutes, un scoring qui déborde prive la règle `anomaly` de la
prédiction de l'heure, qu'elle ne retrouvera au passage suivant que si sa fenêtre la couvre
encore.
This commit is contained in:
Johan LEROY
2026-09-21 13:28:23 +02:00
parent 306c5a52e5
commit f8d08c8686
4 changed files with 34 additions and 17 deletions
+6 -3
View File
@@ -91,8 +91,11 @@ jobs:
bash -c "cd /opt/ml && env -u VIRTUAL_ENV uv run --no-sync python -m enervision_ml.train --help" bash -c "cd /opt/ml && env -u VIRTUAL_ENV uv run --no-sync python -m enervision_ml.train --help"
# `--help` sort par argparse avant `get_settings()` : ni base ni secret requis, et # `--help` sort par argparse avant `get_settings()` : ni base ni secret requis, et
# l'import du module prouve que l'environnement /opt/backend est complet. # l'import du module prouve que l'environnement /opt/backend est complet. Les deux
- name: Vérifie que la détection d'alertes s'importe sans réseau # commandes du DAG `alertes` sont couvertes, `app.cli` tirant tout FastAPI derrière lui.
- name: Vérifie que les deux commandes du DAG alertes s'importent sans réseau
run: > run: >
docker run --rm --network none enervision-airflow:ci docker run --rm --network none enervision-airflow:ci
bash -c "cd /opt/backend && env -u VIRTUAL_ENV uv run --no-sync python -m app.detection.internal_alerts --help" bash -c "cd /opt/backend
&& env -u VIRTUAL_ENV uv run --no-sync python -m app.detection.internal_alerts --help
&& env -u VIRTUAL_ENV uv run --no-sync python -m app.cli generate-recommendations --help"
+11 -4
View File
@@ -80,10 +80,17 @@ l'[ADR 0008](../adr/0008-airflow-execute-le-code-du-backend.md).
`prediction` du même instant, que `ml_score` écrit à l'heure pile. Le décalage laisse le scoring `prediction` du même instant, que `ml_score` écrit à l'heure pile. Le décalage laisse le scoring
finir. Aucune dépendance n'est déclarée entre les deux DAGs pour autant, ni `ExternalTaskSensor` ni finir. Aucune dépendance n'est déclarée entre les deux DAGs pour autant, ni `ExternalTaskSensor` ni
tâche greffée : quatre règles de détection sur cinq ne touchent pas au modèle, et un modèle jamais tâche greffée : quatre règles de détection sur cinq ne touchent pas au modèle, et un modèle jamais
entraîné ne doit pas priver le parc de ses alertes. Ses deux tâches s'enchaînent en revanche entraîné ne doit pas priver le parc de ses alertes. Le décalage est donc une convention et non une
(`recommendation.alert_id` est une clé étrangère `NOT NULL`), et toutes deux sont rejouables sans garantie : le plafond de `ml_score` est de 30 minutes, et un scoring qui déborde de `:15` prive
risque : l'idempotence est portée par la base, `uq_alert_source_reference` et `anomaly` de la `prediction` de l'heure, qu'elle ne retrouvera au passage suivant que si sa fenêtre
`uq_recommendation_alert_rule`. la couvre encore. Les quatre autres règles ne s'en aperçoivent pas.
Ses deux tâches s'enchaînent en revanche (`recommendation.alert_id` est une clé étrangère `NOT
NULL`), et toutes deux sont rejouables sans risque : l'idempotence est portée par la base,
`uq_alert_source_reference` et `uq_recommendation_alert_rule`. Chacune a 2 tentatives, 2 minutes
d'attente entre elles et un plafond de 5 minutes **par tentative** : au pire, reprises comprises,
l'enchaînement occupe 38 minutes, ce qui le garde sous le pas horaire qu'un `max_active_runs=1`
rend contraignant.
`airflow-init` s'appuie sur l'entrypoint de l'image (`_AIRFLOW_DB_MIGRATE`, `airflow-init` s'appuie sur l'entrypoint de l'image (`_AIRFLOW_DB_MIGRATE`,
`_AIRFLOW_WWW_USER_*`) plutôt que sur un script maison : l'entrypoint porte le code de sortie, une `_AIRFLOW_WWW_USER_*`) plutôt que sur un script maison : l'entrypoint porte le code de sortie, une
+3 -3
View File
@@ -26,9 +26,9 @@ COMMANDE_BACKEND = "cd /opt/backend && env -u VIRTUAL_ENV uv run --no-sync pytho
# `uq_alert_source_reference` et `uq_recommendation_alert_rule`) : reprendre ne duplique rien. # `uq_alert_source_reference` et `uq_recommendation_alert_rule`) : reprendre ne duplique rien.
TENTATIVES = 2 TENTATIVES = 2
DELAI_ENTRE_TENTATIVES = timedelta(minutes=2) DELAI_ENTRE_TENTATIVES = timedelta(minutes=2)
# La somme des deux plafonds reste sous le pas horaire : une exécution pendue ne doit pas # `execution_timeout` vaut par tentative : c'est le pire cas des deux tâches enchaînées, reprises
# empiéter sur la suivante. # et délais compris, qui doit tenir sous le pas horaire. Les tests d'intégrité en font le calcul.
PLAFOND_PAR_TACHE = timedelta(minutes=15) PLAFOND_PAR_TACHE = timedelta(minutes=5)
with DAG( with DAG(
dag_id="alertes", dag_id="alertes",
+14 -7
View File
@@ -5,6 +5,7 @@ from datetime import timedelta
from pathlib import Path from pathlib import Path
import pytest import pytest
from airflow.models.baseoperator import BaseOperator
from airflow.models.dagbag import DagBag from airflow.models.dagbag import DagBag
DAGS_FOLDER = Path(__file__).resolve().parent.parent / "dags" DAGS_FOLDER = Path(__file__).resolve().parent.parent / "dags"
@@ -106,14 +107,20 @@ def test_ml_score_execution_timeout_stays_below_its_hourly_step(dagbag: DagBag)
assert timeout < timedelta(hours=1) assert timeout < timedelta(hours=1)
def test_alertes_execution_timeouts_stay_below_its_hourly_step(dagbag: DagBag) -> None: def duree_au_pire(tache: BaseOperator) -> timedelta:
# Les deux taches s'enchainent : c'est leur somme qui doit tenir dans le pas horaire. # `execution_timeout` plafonne une tentative, pas la tache : deux reprises occupent trois
plafonds = [ # plafonds et deux delais d'attente.
dagbag.dags["alertes"].get_task(task_id).execution_timeout assert tache.execution_timeout is not None
for task_id in ("detection", "recommandations") 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")
] ]
assert all(plafond is not None for plafond in plafonds) assert sum((duree_au_pire(tache) for tache in taches), timedelta()) < timedelta(hours=1)
assert sum(plafonds, timedelta()) < timedelta(hours=1)
def test_ml_score_retries_after_a_transient_failure(dagbag: DagBag) -> None: def test_ml_score_retries_after_a_transient_failure(dagbag: DagBag) -> None: