Merge pull request #146 from ineszang/feat/dag-mock-api-import
feat(airflow): orchestre l'import de la Mock API
This commit is contained in:
@@ -93,11 +93,12 @@ jobs:
|
|||||||
|
|
||||||
# `--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 des modules prouve que l'environnement /opt/backend est complet.
|
# l'import des modules prouve que l'environnement /opt/backend est complet.
|
||||||
# Les deux commandes du DAG `alertes` et la commande du DAG historique sont couvertes.
|
# Les commandes des DAGs `alertes`, historique et API Mock sont couvertes.
|
||||||
- name: Vérifie que les trois commandes backend s'importent sans réseau
|
- name: Vérifie que les quatre commandes backend 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
|
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.detection.internal_alerts --help
|
||||||
&& env -u VIRTUAL_ENV uv run --no-sync python -m app.cli generate-recommendations --help
|
&& env -u VIRTUAL_ENV uv run --no-sync python -m app.cli generate-recommendations --help
|
||||||
&& env -u VIRTUAL_ENV uv run --no-sync python -m app.etl.historical_import --help"
|
&& env -u VIRTUAL_ENV uv run --no-sync python -m app.etl.historical_import --help
|
||||||
|
&& env -u VIRTUAL_ENV uv run --no-sync python -m app.etl.mock_api_import --help"
|
||||||
@@ -23,7 +23,7 @@ Ce que la documentation apporte à chacun : [docs/architecture/00-vue-ensemble.m
|
|||||||
| Backend | FastAPI, Python 3.14 | `apps/backend` | En place |
|
| Backend | FastAPI, Python 3.14 | `apps/backend` | En place |
|
||||||
| Frontend | Angular 22, Node 26 | `apps/frontend` | En place |
|
| Frontend | Angular 22, Node 26 | `apps/frontend` | En place |
|
||||||
| Base | PostgreSQL 17 + TimescaleDB | `db` | En place |
|
| Base | PostgreSQL 17 + TimescaleDB | `db` | En place |
|
||||||
| ETL | Apache Airflow | `etl/airflow` | Quatre DAGs |
|
| ETL | Apache Airflow | `etl/airflow` | Cinq DAGs |
|
||||||
| Infra | Terraform (k3s single-node) | `infra/terraform` | Initialise |
|
| Infra | Terraform (k3s single-node) | `infra/terraform` | Initialise |
|
||||||
| Reverse proxy | Nginx, TLS | `infra/proxy` | En place |
|
| Reverse proxy | Nginx, TLS | `infra/proxy` | En place |
|
||||||
| CI/CD | GitHub Actions | `.github/workflows` | En place |
|
| CI/CD | GitHub Actions | `.github/workflows` | En place |
|
||||||
@@ -50,7 +50,7 @@ L'etat detaille de chaque brique et les vues d'architecture sont dans
|
|||||||
│ ├── migrations/ Migrations SQL versionnees
|
│ ├── migrations/ Migrations SQL versionnees
|
||||||
│ └── seeds/ Jeux de donnees de reference
|
│ └── seeds/ Jeux de donnees de reference
|
||||||
├── etl/airflow/
|
├── etl/airflow/
|
||||||
│ ├── dags/ DAGs d'orchestration (pipeline ML, alertes, import historique)
|
│ ├── dags/ DAGs d'orchestration (pipeline ML, alertes, imports historique et API Mock)
|
||||||
│ ├── plugins/ Operateurs et hooks maison
|
│ ├── plugins/ Operateurs et hooks maison
|
||||||
│ ├── include/ Requetes SQL et ressources des DAGs
|
│ ├── include/ Requetes SQL et ressources des DAGs
|
||||||
│ └── tests/ Tests d'integrite des DAGs
|
│ └── tests/ Tests d'integrite des DAGs
|
||||||
|
|||||||
@@ -45,15 +45,19 @@ CAPACITY_BOUNDS = (0.0, 100_000.0)
|
|||||||
def create_mock_api_client() -> httpx.AsyncClient:
|
def create_mock_api_client() -> httpx.AsyncClient:
|
||||||
settings = get_settings()
|
settings = get_settings()
|
||||||
|
|
||||||
if settings.mock_api_username is None or settings.mock_api_password is None:
|
username = settings.mock_api_username
|
||||||
|
password = (
|
||||||
|
settings.mock_api_password.get_secret_value()
|
||||||
|
if settings.mock_api_password is not None
|
||||||
|
else None
|
||||||
|
)
|
||||||
|
|
||||||
|
if not username or not username.strip() or not password or not password.strip():
|
||||||
raise ValueError("Les identifiants de l'API Mock ne sont pas configurés.")
|
raise ValueError("Les identifiants de l'API Mock ne sont pas configurés.")
|
||||||
|
|
||||||
return httpx.AsyncClient(
|
return httpx.AsyncClient(
|
||||||
base_url=settings.mock_api_base_url.rstrip("/"),
|
base_url=settings.mock_api_base_url.rstrip("/"),
|
||||||
auth=(
|
auth=(username, password),
|
||||||
settings.mock_api_username,
|
|
||||||
settings.mock_api_password.get_secret_value(),
|
|
||||||
),
|
|
||||||
timeout=settings.mock_api_timeout_seconds,
|
timeout=settings.mock_api_timeout_seconds,
|
||||||
)
|
)
|
||||||
|
|
||||||
|
|||||||
@@ -283,6 +283,41 @@ def test_create_mock_api_client_requires_credentials(
|
|||||||
mock_api_import.create_mock_api_client()
|
mock_api_import.create_mock_api_client()
|
||||||
|
|
||||||
|
|
||||||
|
@pytest.mark.parametrize(
|
||||||
|
("username", "password_value"),
|
||||||
|
[
|
||||||
|
("", "test-password"),
|
||||||
|
("test-user", ""),
|
||||||
|
(" ", "test-password"),
|
||||||
|
("test-user", " "),
|
||||||
|
],
|
||||||
|
)
|
||||||
|
def test_create_mock_api_client_rejects_empty_credentials(
|
||||||
|
monkeypatch: pytest.MonkeyPatch,
|
||||||
|
username: str,
|
||||||
|
password_value: str,
|
||||||
|
) -> None:
|
||||||
|
password = MagicMock()
|
||||||
|
password.get_secret_value.return_value = password_value
|
||||||
|
|
||||||
|
settings = SimpleNamespace(
|
||||||
|
mock_api_username=username,
|
||||||
|
mock_api_password=password,
|
||||||
|
)
|
||||||
|
|
||||||
|
monkeypatch.setattr(
|
||||||
|
mock_api_import,
|
||||||
|
"get_settings",
|
||||||
|
lambda: settings,
|
||||||
|
)
|
||||||
|
|
||||||
|
with pytest.raises(
|
||||||
|
ValueError,
|
||||||
|
match="Les identifiants de l'API Mock ne sont pas configurés",
|
||||||
|
):
|
||||||
|
mock_api_import.create_mock_api_client()
|
||||||
|
|
||||||
|
|
||||||
async def test_create_mock_api_client_uses_configuration(
|
async def test_create_mock_api_client_uses_configuration(
|
||||||
monkeypatch: pytest.MonkeyPatch,
|
monkeypatch: pytest.MonkeyPatch,
|
||||||
) -> None:
|
) -> None:
|
||||||
|
|||||||
+14
-4
@@ -34,12 +34,14 @@ x-airflow-common: &airflow-common
|
|||||||
# memes identifiants que le backend en attendant.
|
# memes identifiants que le backend en attendant.
|
||||||
ML_DATABASE_URL: postgresql+psycopg://${POSTGRES_USER}:${POSTGRES_PASSWORD}@db:5432/${POSTGRES_DB}
|
ML_DATABASE_URL: postgresql+psycopg://${POSTGRES_USER}:${POSTGRES_PASSWORD}@db:5432/${POSTGRES_DB}
|
||||||
MLFLOW_TRACKING_URI: sqlite:////opt/ml/state/mlflow.db
|
MLFLOW_TRACKING_URI: sqlite:////opt/ml/state/mlflow.db
|
||||||
# Le DAG `alertes` lance le backend en sous-processus : il lit `DATABASE_URL`, en
|
# Les DAGs backend lisent `DATABASE_URL` en dialecte asyncpg, là où le pipeline ML
|
||||||
# dialecte asyncpg, là où le pipeline ML lit `ML_DATABASE_URL`.
|
# utilise `ML_DATABASE_URL`.
|
||||||
DATABASE_URL: postgresql+asyncpg://${POSTGRES_USER}:${POSTGRES_PASSWORD}@db:5432/${POSTGRES_DB}
|
DATABASE_URL: postgresql+asyncpg://${POSTGRES_USER}:${POSTGRES_PASSWORD}@db:5432/${POSTGRES_DB}
|
||||||
# Clé distincte de celle de l'API : la détection ne signe ni ne vérifie aucun jeton, et
|
|
||||||
# Airflow permet d'exécuter du code depuis son interface (cf. ADR 0008).
|
# Clé distincte de celle de l'API : les traitements lancés par Airflow ne signent ni ne
|
||||||
|
# vérifient aucun jeton. Airflow permet d'exécuter du code depuis son interface (ADR 0008).
|
||||||
APP_SECRET_KEY: ${AIRFLOW_APP_SECRET_KEY:-}
|
APP_SECRET_KEY: ${AIRFLOW_APP_SECRET_KEY:-}
|
||||||
|
|
||||||
volumes:
|
volumes:
|
||||||
- ./etl/airflow/dags:/opt/airflow/dags
|
- ./etl/airflow/dags:/opt/airflow/dags
|
||||||
- ./etl/airflow/plugins:/opt/airflow/plugins
|
- ./etl/airflow/plugins:/opt/airflow/plugins
|
||||||
@@ -167,6 +169,14 @@ services:
|
|||||||
airflow-scheduler:
|
airflow-scheduler:
|
||||||
<<: *airflow-common
|
<<: *airflow-common
|
||||||
command: scheduler
|
command: scheduler
|
||||||
|
environment:
|
||||||
|
<<: *airflow-common-env
|
||||||
|
# LocalExecutor exécute les tâches dans le scheduler : lui seul a besoin des
|
||||||
|
# identifiants de l'API Mock.
|
||||||
|
APP_MOCK_API_BASE_URL: ${APP_MOCK_API_BASE_URL:-https://api-mock.charlieandre.fr}
|
||||||
|
APP_MOCK_API_USERNAME: ${APP_MOCK_API_USERNAME:-}
|
||||||
|
APP_MOCK_API_PASSWORD: ${APP_MOCK_API_PASSWORD:-}
|
||||||
|
APP_MOCK_API_TIMEOUT_SECONDS: ${APP_MOCK_API_TIMEOUT_SECONDS:-10}
|
||||||
depends_on:
|
depends_on:
|
||||||
db:
|
db:
|
||||||
condition: service_healthy
|
condition: service_healthy
|
||||||
|
|||||||
@@ -70,11 +70,11 @@ Le lien `front -.-> api` reste en pointillé : le frontend appelle bien une API,
|
|||||||
intercepteur répond à sa place tant que les endpoints n'existent pas. Voir
|
intercepteur répond à sa place tant que les endpoints n'existent pas. Voir
|
||||||
[30-frontend.md](30-frontend.md).
|
[30-frontend.md](30-frontend.md).
|
||||||
|
|
||||||
Le lien `airflow --> db` est maintenant en trait plein : quatre DAGs tournent, deux pour
|
Le lien `airflow --> db` est maintenant en trait plein : cinq DAGs tournent, deux pour
|
||||||
l'entraînement et le scoring du modèle ML (issue #115), un pour la détection d'alertes et la
|
l'entraînement et le scoring du modèle ML (issue #115), un pour la détection d'alertes et la
|
||||||
génération des recommandations (issue #116), et `historical_import` pour l'ingestion du dataset
|
génération des recommandations (issue #116), `historical_import` pour le dataset historique
|
||||||
historique (issue #119). L'orchestration de l'import API Mock et la réconciliation globale des
|
(issue #119) et `mock_api_import` pour l'ingestion horaire de l'API Mock (issue #15).
|
||||||
deux sources restent à compléter dans l'issue #15.
|
La réconciliation globale des données provenant des deux sources reste à compléter dans l'issue #15.
|
||||||
|
|
||||||
Le lien `prom -.-> api` de même : l'API expose bien `/metrics` au format Prometheus, mais aucun
|
Le lien `prom -.-> api` de même : l'API expose bien `/metrics` au format Prometheus, mais aucun
|
||||||
collecteur ne vient le lire.
|
collecteur ne vient le lire.
|
||||||
@@ -89,15 +89,16 @@ collecteur ne vient le lire.
|
|||||||
| ML | LightGBM, MLflow | `ml` | `En cours` | Pipeline d'entraînement et de scoring (`enervision_ml.train`/`.score`, features par lags/moyennes glissantes partagées entre les deux, baseline de persistance saisonnière, suivi MLflow local), exposé en lecture via `GET /predictions`, orchestré par Airflow (`ml_train`/`ml_score`). Voir [ADR 0005](../adr/0005-modele-prediction-lightgbm.md) et [ML-START.md](../ML-START.md). Surveillance de dérive (EC06, #44/#45) pas encore construite |
|
| ML | LightGBM, MLflow | `ml` | `En cours` | Pipeline d'entraînement et de scoring (`enervision_ml.train`/`.score`, features par lags/moyennes glissantes partagées entre les deux, baseline de persistance saisonnière, suivi MLflow local), exposé en lecture via `GET /predictions`, orchestré par Airflow (`ml_train`/`ml_score`). Voir [ADR 0005](../adr/0005-modele-prediction-lightgbm.md) et [ML-START.md](../ML-START.md). Surveillance de dérive (EC06, #44/#45) pas encore construite |
|
||||||
| Infra | Docker Compose, Nginx, Terraform, k3s single-node | `infra`, `docker-compose.prod.yml` | `En cours` | Reverse proxy et overlay de déploiement écrits et validés, jamais lancés sur le serveur ([ADR 0007](../adr/0007-terminaison-tls-et-reverse-proxy-nginx.md)). Provisionnement de la VM par Terraform, qui installe Docker, prépare les deux environnements et enregistre le runner, jamais appliqué ([ADR 0010](../adr/0010-terraform-provisionne-github-actions-deploie.md)). Module d'installation k3s jamais appliqué, aucune ressource Kubernetes déclarée |
|
| Infra | Docker Compose, Nginx, Terraform, k3s single-node | `infra`, `docker-compose.prod.yml` | `En cours` | Reverse proxy et overlay de déploiement écrits et validés, jamais lancés sur le serveur ([ADR 0007](../adr/0007-terminaison-tls-et-reverse-proxy-nginx.md)). Provisionnement de la VM par Terraform, qui installe Docker, prépare les deux environnements et enregistre le runner, jamais appliqué ([ADR 0010](../adr/0010-terraform-provisionne-github-actions-deploie.md)). Module d'installation k3s jamais appliqué, aucune ressource Kubernetes déclarée |
|
||||||
| Monitoring | Prometheus, Grafana, Alertmanager | `monitoring` | `Cible` | Rien, hors le `/metrics` exposé par l'API |
|
| Monitoring | Prometheus, Grafana, Alertmanager | `monitoring` | `Cible` | Rien, hors le `/metrics` exposé par l'API |
|
||||||
| ETL | Apache Airflow | `etl/airflow` | `En cours` | Webserver + scheduler (LocalExecutor) tournent via docker-compose, base de métadonnées Postgres dédiée. Quatre DAGs en sous-processus `uv run` : `ml_train`, `ml_score`, `alertes` et `historical_import`. Le DAG historique orchestre `app.etl.historical_import` et charge `dataset`, `site` et `reading`. L'orchestration API Mock reste à compléter dans #15 |
|
| ETL | Apache Airflow | `etl/airflow` | `En cours` | Webserver et scheduler avec LocalExecutor via Docker Compose, sur une base PostgreSQL dédiée. Cinq DAGs sont présents : `ml_train`, `ml_score`, `alertes`, `historical_import` et `mock_api_import`. L'import historique reste manuel et l'import API Mock est exécuté chaque heure. La réconciliation globale des deux sources reste à compléter dans l'issue #15. |
|
||||||
| CI/CD | GitHub Actions | `.github/workflows` | `En cours` | 7 workflows, 19 jobs : lint, typage, tests avec seuil de couverture bloquant, tests d'intégration sur TimescaleDB réel, audit de dépendances, SAST Bandit, quality gate SonarCloud, intégrité des DAGs Airflow, formatage et validation du Terraform. Déploiement continu vers la VM ENI écrit par `deploy.yml`, `dev` en recette et `main` en production après approbation ([ADR 0009](../adr/0009-deux-environnements-compose-sur-la-vm-eni.md)), mais jamais exécuté : la machine n'est pas provisionnée et le runner n'y est pas enregistré. Détail dans [50-cicd.md](50-cicd.md) |
|
| CI/CD | GitHub Actions | `.github/workflows` | `En cours` | 7 workflows, 19 jobs : lint, typage, tests avec seuil de couverture bloquant, tests d'intégration sur TimescaleDB réel, audit de dépendances, SAST Bandit, quality gate SonarCloud, intégrité des DAGs Airflow, formatage et validation du Terraform. Déploiement continu vers la VM ENI écrit par `deploy.yml`, `dev` en recette et `main` en production après approbation ([ADR 0009](../adr/0009-deux-environnements-compose-sur-la-vm-eni.md)), mais jamais exécuté : la machine n'est pas provisionnée et le runner n'y est pas enregistré. Détail dans [50-cicd.md](50-cicd.md) |
|
||||||
|
|
||||||
## Flux bout en bout
|
## Flux bout en bout
|
||||||
|
|
||||||
Statut : `En cours`. **Le chemin de lecture tourne** : base, API et frontend. **Le chemin
|
Statut : `En cours`. **Le chemin de lecture tourne** entre la base, l'API et le frontend.
|
||||||
d'ingestion dessiné ci-dessous n'existe pas** : les trois DAGs livrés (`ml_train`, `ml_score`,
|
**Le chemin d'ingestion est maintenant orchestré par Airflow** : `historical_import` charge
|
||||||
issue #115 ; `alertes`, issue #116) orchestrent le pipeline ML et la détection d'alertes, pas
|
manuellement le dataset CSV/JSON et `mock_api_import` collecte périodiquement les mesures de
|
||||||
l'ingestion, qui reste lancée à la main par les scripts d'import (issues #15 et #16).
|
l'API Mock. La réconciliation globale des données provenant des deux sources reste à compléter
|
||||||
|
dans l'issue #15.
|
||||||
|
|
||||||
```mermaid
|
```mermaid
|
||||||
sequenceDiagram
|
sequenceDiagram
|
||||||
|
|||||||
@@ -54,7 +54,7 @@ Trois pièges sont documentés en tête du `docker-compose.yml`, ils ne se devin
|
|||||||
- `LocalExecutor` exécute les tâches comme sous-processus du **scheduler**, jamais de l'api-server :
|
- `LocalExecutor` exécute les tâches comme sous-processus du **scheduler**, jamais de l'api-server :
|
||||||
c'est le scheduler qui a besoin du volume `airflow_ml_state` (modèle, magasin MLflow).
|
c'est le scheduler qui a besoin du volume `airflow_ml_state` (modèle, magasin MLflow).
|
||||||
|
|
||||||
### Airflow (issues #115, #116 et #119)
|
### Airflow (issues #15, #115, #116 et #119)
|
||||||
|
|
||||||
Quatre services (Airflow 3.3), `docker compose profiles` non utilisés (démarrage explicite via `make
|
Quatre services (Airflow 3.3), `docker compose profiles` non utilisés (démarrage explicite via `make
|
||||||
airflow-up`, pas dans `make dev`) :
|
airflow-up`, pas dans `make dev`) :
|
||||||
@@ -87,12 +87,24 @@ l'[ADR 0008](../adr/0008-airflow-execute-le-code-du-backend.md).
|
|||||||
| `ml_score` | `0 * * * *` | `enervision_ml.score`, dans `/opt/ml/.venv` |
|
| `ml_score` | `0 * * * *` | `enervision_ml.score`, dans `/opt/ml/.venv` |
|
||||||
| `alertes` | `15 * * * *` | `app.detection.internal_alerts` puis `app.cli generate-recommendations`, dans `/opt/backend/.venv` |
|
| `alertes` | `15 * * * *` | `app.detection.internal_alerts` puis `app.cli generate-recommendations`, dans `/opt/backend/.venv` |
|
||||||
| `historical_import` | manuelle | `app.etl.historical_import`, dans `/opt/backend/.venv` ; les fichiers de `data/raw` sont montés en lecture seule dans `/opt/data/raw` |
|
| `historical_import` | manuelle | `app.etl.historical_import`, dans `/opt/backend/.venv` ; les fichiers de `data/raw` sont montés en lecture seule dans `/opt/data/raw` |
|
||||||
|
| `mock_api_import` | `45 * * * *` | `app.etl.mock_api_import`, dans `/opt/backend/.venv` ; importe l'heure précédant son déclenchement depuis l'API Mock |
|
||||||
|
|
||||||
Le DAG `historical_import` réutilise le pipeline historique existant sans dupliquer sa logique.
|
Le DAG `historical_import` réutilise le pipeline historique existant sans dupliquer sa logique.
|
||||||
Il reste manuel, car le dataset sert à initialiser l'environnement. Le montage
|
Il reste manuel, car le dataset sert à initialiser l'environnement. Le montage
|
||||||
`./data/raw:/opt/data/raw:ro` permet au scheduler de lire les fichiers CSV/JSON sans pouvoir les
|
`./data/raw:/opt/data/raw:ro` permet au scheduler de lire les fichiers CSV/JSON sans pouvoir les
|
||||||
modifier.
|
modifier.
|
||||||
|
|
||||||
|
Le DAG `mock_api_import` exécute le pipeline API Mock toutes les heures, à la minute `:45`.
|
||||||
|
Un `CronTriggerTimetable` explicite lui attribue un intervalle d'une heure, y compris lors d'un
|
||||||
|
déclenchement manuel. Il transmet cet intervalle au script backend et charge les mesures dans
|
||||||
|
les tables communes `site` et `reading`. Le décalage à `:45` laisse quinze minutes avant le
|
||||||
|
scoring exécuté à l'heure pile, puis quinze minutes supplémentaires avant les alertes à `:15`.
|
||||||
|
`max_active_runs=1` empêche deux exécutions du DAG de se chevaucher.
|
||||||
|
|
||||||
|
Le DAG conserve `catchup=False` pour éviter un rattrapage massif depuis sa date de démarrage.
|
||||||
|
Une interruption du scheduler peut donc créer un intervalle manquant, qui devra être rejoué
|
||||||
|
explicitement par une opération de backfill.
|
||||||
|
|
||||||
**Pourquoi `alertes` tourne à la quinzième minute.** Sa règle `anomaly` compare une lecture à la
|
**Pourquoi `alertes` tourne à la quinzième minute.** Sa règle `anomaly` compare une lecture à la
|
||||||
`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
|
||||||
|
|||||||
+15
-9
@@ -663,16 +663,22 @@ mock_api_import.py
|
|||||||
|
|
||||||
La logique d'extraction, de transformation et de chargement est donc 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 désormais réellement (`etl/airflow/`, `make airflow-up`) et orchestre le pipeline
|
Airflow tourne désormais réellement (`etl/airflow/`, `make airflow-up`) et orchestre cinq DAGs :
|
||||||
ML (`ml_train`/`ml_score`, issue #115), la détection d'alertes et la génération des
|
le pipeline ML (`ml_train` et `ml_score`, issue #115), la détection d'alertes et la génération
|
||||||
recommandations (`alertes`, issue #116), ainsi que l'import historique
|
des recommandations (`alertes`, issue #116), l'import historique (`historical_import`,
|
||||||
(`historical_import`, issue #119).
|
issue #119) et l'import périodique de l'API Mock (`mock_api_import`, issue #15).
|
||||||
|
|
||||||
Le DAG `historical_import` est déclenché manuellement. Il exécute
|
Le DAG `mock_api_import` s'exécute chaque heure, à la minute `:45`. Il appelle
|
||||||
`app.etl.historical_import` avec les fichiers montés en lecture seule depuis `data/raw` vers
|
`app.etl.mock_api_import` avec un intervalle explicite d'une heure et une limite de 1 000 lectures
|
||||||
`/opt/data/raw`. L'orchestration de l'import API Mock et la réconciliation globale des deux
|
par site. Les deux pipelines normalisent leurs données vers les tables communes `site` et
|
||||||
sources restent couvertes par l'issue #15.
|
`reading`, tout en conservant leur source (`csv` ou `api_history`). La réconciliation globale
|
||||||
|
des deux sources reste à compléter dans l'issue #15.
|
||||||
|
|
||||||
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 `mock_api_import` exécute `app.etl.mock_api_import` toutes les heures. Chaque exécution
|
||||||
|
traite l'intervalle Airflow précédent. Les deux pipelines normalisent leurs données vers les
|
||||||
|
tables communes `site` et `reading`, tout en conservant leur source (`csv` ou `api_history`).
|
||||||
|
|
||||||
|
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`, `alertes.py`, `historical_import.py` et
|
||||||
|
`mock_api_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 pipeline Data servira ensuite à préparer les données nécessaires au modèle de Machine Learning.
|
Le pipeline Data servira ensuite à préparer les données nécessaires au modèle de Machine Learning.
|
||||||
|
|||||||
@@ -0,0 +1,61 @@
|
|||||||
|
"""DAG d'import périodique des données de l'API Mock EnerVision (issue #15).
|
||||||
|
|
||||||
|
Orchestre le pipeline existant `app.etl.mock_api_import` sans dupliquer sa logique ETL.
|
||||||
|
Chaque exécution traite l'heure précédant son déclenchement.
|
||||||
|
|
||||||
|
Le pipeline backend reste responsable de la validation, de la normalisation, du suivi de la
|
||||||
|
qualité, de l'idempotence et du chargement dans PostgreSQL/TimescaleDB.
|
||||||
|
"""
|
||||||
|
|
||||||
|
from __future__ import annotations
|
||||||
|
|
||||||
|
from datetime import datetime, timedelta
|
||||||
|
|
||||||
|
from airflow.providers.standard.operators.bash import BashOperator
|
||||||
|
from airflow.sdk import DAG
|
||||||
|
from airflow.timetables.trigger import CronTriggerTimetable
|
||||||
|
|
||||||
|
# Le backend possède son propre environnement uv dans l'image Airflow (ADR 0008).
|
||||||
|
COMMANDE_BACKEND = "cd /opt/backend && env -u VIRTUAL_ENV uv run --no-sync python -m"
|
||||||
|
|
||||||
|
# Le pipeline backend et l'API acceptent au maximum 1 000 lectures par site.
|
||||||
|
# Cette marge évite de perdre silencieusement une lecture si une heure en contient plus de 60.
|
||||||
|
LIMITE_LECTURES = 1000
|
||||||
|
|
||||||
|
# Deux reprises donnent trois tentatives au total. Même dans le pire cas, l'exécution reste
|
||||||
|
# inférieure au pas horaire du DAG.
|
||||||
|
NOMBRE_REPRISES = 2
|
||||||
|
DELAI_ENTRE_REPRISES = timedelta(minutes=2)
|
||||||
|
PLAFOND_PAR_TENTATIVE = timedelta(minutes=10)
|
||||||
|
|
||||||
|
# L'intervalle est déclaré explicitement pour ne pas dépendre de la valeur du paramètre Airflow
|
||||||
|
# `create_cron_data_intervals`. Le déclenchement à :45 laisse quinze minutes avant `ml_score`,
|
||||||
|
# exécuté à l'heure pile, puis avant `alertes`, exécuté à :15.
|
||||||
|
PLANIFICATION = CronTriggerTimetable(
|
||||||
|
"45 * * * *",
|
||||||
|
timezone="UTC",
|
||||||
|
interval=timedelta(hours=1),
|
||||||
|
)
|
||||||
|
|
||||||
|
with DAG(
|
||||||
|
dag_id="mock_api_import",
|
||||||
|
description="Importe chaque heure les données de l'API Mock dans site et reading.",
|
||||||
|
schedule=PLANIFICATION,
|
||||||
|
start_date=datetime(2026, 1, 1),
|
||||||
|
catchup=False,
|
||||||
|
# Deux exécutions simultanées pourraient demander et traiter le même intervalle.
|
||||||
|
max_active_runs=1,
|
||||||
|
tags=["etl", "mock-api"],
|
||||||
|
) as dag:
|
||||||
|
BashOperator(
|
||||||
|
task_id="import_mock_api",
|
||||||
|
bash_command=(
|
||||||
|
f"{COMMANDE_BACKEND} app.etl.mock_api_import "
|
||||||
|
"--start-time \"{{ data_interval_start.strftime('%Y-%m-%dT%H:%M:%S') }}\" "
|
||||||
|
"--end-time \"{{ data_interval_end.strftime('%Y-%m-%dT%H:%M:%S') }}\" "
|
||||||
|
f"--limit {LIMITE_LECTURES}"
|
||||||
|
),
|
||||||
|
retries=NOMBRE_REPRISES,
|
||||||
|
retry_delay=DELAI_ENTRE_REPRISES,
|
||||||
|
execution_timeout=PLAFOND_PAR_TENTATIVE,
|
||||||
|
)
|
||||||
@@ -1,22 +1,30 @@
|
|||||||
"""Tests d'integrite des DAGs : s'importent sans erreur, structure attendue. Pas d'execution
|
"""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."""
|
reelle des taches (ca reclamerait le conteneur avec `uv`/`enervision_ml`), juste la definition."""
|
||||||
|
|
||||||
from datetime import timedelta
|
from datetime import datetime, timedelta
|
||||||
from pathlib import Path
|
from pathlib import Path
|
||||||
|
|
||||||
import pytest
|
import pytest
|
||||||
from airflow.dag_processing.dagbag import DagBag
|
from airflow.dag_processing.dagbag import DagBag
|
||||||
from airflow.sdk import BaseOperator
|
from airflow.sdk import BaseOperator
|
||||||
|
from airflow.timetables.trigger import CronTriggerTimetable
|
||||||
|
|
||||||
DAGS_FOLDER = Path(__file__).resolve().parent.parent / "dags"
|
DAGS_FOLDER = Path(__file__).resolve().parent.parent / "dags"
|
||||||
|
|
||||||
DAG_IDS = ["ml_train", "ml_score", "alertes", "historical_import"]
|
DAG_IDS = [
|
||||||
|
"ml_train",
|
||||||
|
"ml_score",
|
||||||
|
"alertes",
|
||||||
|
"historical_import",
|
||||||
|
"mock_api_import",
|
||||||
|
]
|
||||||
TACHES = [
|
TACHES = [
|
||||||
("ml_train", "train"),
|
("ml_train", "train"),
|
||||||
("ml_score", "score"),
|
("ml_score", "score"),
|
||||||
("alertes", "detection"),
|
("alertes", "detection"),
|
||||||
("alertes", "recommandations"),
|
("alertes", "recommandations"),
|
||||||
("historical_import", "import_historical"),
|
("historical_import", "import_historical"),
|
||||||
|
("mock_api_import", "import_mock_api"),
|
||||||
]
|
]
|
||||||
|
|
||||||
|
|
||||||
@@ -52,6 +60,19 @@ def test_historical_import_has_no_schedule(dagbag: DagBag) -> None:
|
|||||||
assert dagbag.dags["historical_import"].schedule is None
|
assert dagbag.dags["historical_import"].schedule is None
|
||||||
|
|
||||||
|
|
||||||
|
def test_mock_api_import_uses_an_explicit_hourly_interval(dagbag: DagBag) -> None:
|
||||||
|
timetable = dagbag.dags["mock_api_import"].timetable
|
||||||
|
|
||||||
|
assert isinstance(timetable, CronTriggerTimetable)
|
||||||
|
assert timetable.serialize()["expression"] == "45 * * * *"
|
||||||
|
|
||||||
|
manual_interval = timetable.infer_manual_data_interval(
|
||||||
|
run_after=datetime.fromisoformat("2026-09-22T12:30:00+00:00"),
|
||||||
|
)
|
||||||
|
|
||||||
|
assert manual_interval.end - manual_interval.start == timedelta(hours=1)
|
||||||
|
|
||||||
|
|
||||||
def test_ml_train_task_calls_the_training_module(dagbag: DagBag) -> None:
|
def test_ml_train_task_calls_the_training_module(dagbag: DagBag) -> None:
|
||||||
tache = dagbag.dags["ml_train"].get_task("train")
|
tache = dagbag.dags["ml_train"].get_task("train")
|
||||||
assert "enervision_ml.train" in tache.bash_command
|
assert "enervision_ml.train" in tache.bash_command
|
||||||
@@ -84,6 +105,20 @@ def test_historical_import_uses_the_expected_source_files(dagbag: DagBag) -> Non
|
|||||||
assert "--metadata /opt/data/raw/dataset_metadata.json" in commande
|
assert "--metadata /opt/data/raw/dataset_metadata.json" in commande
|
||||||
|
|
||||||
|
|
||||||
|
def test_mock_api_import_calls_the_existing_backend_module(dagbag: DagBag) -> None:
|
||||||
|
commande = dagbag.dags["mock_api_import"].get_task("import_mock_api").bash_command
|
||||||
|
|
||||||
|
assert "app.etl.mock_api_import" in commande
|
||||||
|
|
||||||
|
|
||||||
|
def test_mock_api_import_uses_the_airflow_data_interval(dagbag: DagBag) -> None:
|
||||||
|
commande = dagbag.dags["mock_api_import"].get_task("import_mock_api").bash_command
|
||||||
|
|
||||||
|
assert "--start-time \"{{ data_interval_start.strftime('%Y-%m-%dT%H:%M:%S') }}\"" in commande
|
||||||
|
assert "--end-time \"{{ data_interval_end.strftime('%Y-%m-%dT%H:%M:%S') }}\"" in commande
|
||||||
|
assert "--limit 1000" in commande
|
||||||
|
|
||||||
|
|
||||||
@pytest.mark.parametrize("task_id", ["detection", "recommandations"])
|
@pytest.mark.parametrize("task_id", ["detection", "recommandations"])
|
||||||
def test_alertes_tasks_run_in_the_backend_environment(dagbag: DagBag, task_id: str) -> None:
|
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).
|
# Le backend a son propre venv dans l'image, distinct de celui de ml/ (ADR 0008).
|
||||||
@@ -95,6 +130,12 @@ def test_historical_import_runs_in_the_backend_environment(dagbag: DagBag) -> No
|
|||||||
assert "/opt/backend" in commande
|
assert "/opt/backend" in commande
|
||||||
|
|
||||||
|
|
||||||
|
def test_mock_api_import_runs_in_the_backend_environment(dagbag: DagBag) -> None:
|
||||||
|
commande = dagbag.dags["mock_api_import"].get_task("import_mock_api").bash_command
|
||||||
|
|
||||||
|
assert "/opt/backend" in commande
|
||||||
|
|
||||||
|
|
||||||
def test_alertes_generates_recommendations_after_detecting(dagbag: DagBag) -> None:
|
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
|
# `recommendation.alert_id` est une cle etrangere `NOT NULL` : la generation n'a rien a lire
|
||||||
# tant que la detection n'a pas ecrit.
|
# tant que la detection n'a pas ecrit.
|
||||||
@@ -136,6 +177,14 @@ def duree_au_pire(tache: BaseOperator) -> timedelta:
|
|||||||
return (tache.retries + 1) * tache.execution_timeout + tache.retries * tache.retry_delay
|
return (tache.retries + 1) * tache.execution_timeout + tache.retries * tache.retry_delay
|
||||||
|
|
||||||
|
|
||||||
|
def test_mock_api_import_worst_case_stays_below_its_hourly_step(
|
||||||
|
dagbag: DagBag,
|
||||||
|
) -> None:
|
||||||
|
tache = dagbag.dags["mock_api_import"].get_task("import_mock_api")
|
||||||
|
|
||||||
|
assert duree_au_pire(tache) < timedelta(hours=1)
|
||||||
|
|
||||||
|
|
||||||
def test_alertes_worst_case_stays_below_its_hourly_step(dagbag: DagBag) -> None:
|
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
|
# 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.
|
# pas horaire, sinon `max_active_runs=1` fait attendre l'execution suivante.
|
||||||
@@ -159,6 +208,10 @@ def test_historical_import_retries_after_a_transient_failure(dagbag: DagBag) ->
|
|||||||
assert dagbag.dags["historical_import"].get_task("import_historical").retries >= 1
|
assert dagbag.dags["historical_import"].get_task("import_historical").retries >= 1
|
||||||
|
|
||||||
|
|
||||||
|
def test_mock_api_import_retries_after_a_transient_failure(dagbag: DagBag) -> None:
|
||||||
|
assert dagbag.dags["mock_api_import"].get_task("import_mock_api").retries >= 1
|
||||||
|
|
||||||
|
|
||||||
@pytest.mark.parametrize(("dag_id", "task_id"), TACHES)
|
@pytest.mark.parametrize(("dag_id", "task_id"), TACHES)
|
||||||
def test_tasks_never_resync_the_baked_environment(
|
def test_tasks_never_resync_the_baked_environment(
|
||||||
dagbag: DagBag, dag_id: str, task_id: str
|
dagbag: DagBag, dag_id: str, task_id: str
|
||||||
|
|||||||
Reference in New Issue
Block a user