Compare commits

..
Author SHA1 Message Date
Meryemel-gham fe100653d0 fix(etl): traite les retours de revue du DAG historique
Airflow / Lint et intégrité des DAGs (push) Successful in 57s
Airflow / Construction de l'image (push) Successful in 2m52s
SonarQube / build-back (push) Successful in 1m17s
SonarQube / test-ml (push) Failing after 1m50s
SonarQube / build-front (push) Successful in 9m44s
SonarQube / test-back (push) Failing after 55s
SonarQube / test-front (push) Failing after 5m3s
SonarQube / SonarQube (push) Skipped
2026-09-21 16:05:55 +02:00
Meryemel-gham cb4df846cb docs(etl): documente l'orchestration de l'import historique 2026-09-21 14:41:38 +02:00
Meryemel-gham ae4082c584 ci(etl): valide l'import historique dans l'image Airflow 2026-09-21 14:41:26 +02:00
Meryemel-gham f18d4f9ef9 test(etl): couvre le DAG d'import historique 2026-09-21 14:41:08 +02:00
Meryemel-gham 308769b325 feat(etl): orchestre l'import historique avec Airflow 2026-09-21 14:40:21 +02:00
18 changed files with 168 additions and 226 deletions
+5 -4
View File
@@ -91,11 +91,12 @@ jobs:
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
# l'import du module prouve que l'environnement /opt/backend est complet. Les deux
# 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
# 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.
- name: Vérifie que les trois commandes backend s'importent sans réseau
run: >
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
&& 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"
+2 -2
View File
@@ -21,7 +21,7 @@ Ce que la documentation apporte à chacun : [docs/architecture/00-vue-ensemble.m
| Backend | FastAPI, Python 3.14 | `apps/backend` | Initialise |
| Frontend | Angular 22, Node 24 LTS | `apps/frontend` | Tableau de bord |
| Base | PostgreSQL 17 + TimescaleDB | `db` | Initialise |
| ETL | Apache Airflow | `etl/airflow` | Trois DAGs |
| ETL | Apache Airflow | `etl/airflow` | Quatre DAGs |
| Infra | Terraform (k3s single-node) | `infra/terraform` | Initialise |
| Reverse proxy | Nginx, TLS | `infra/proxy` | En place |
| CI/CD | GitHub Actions | `.github/workflows` | Backend en place |
@@ -48,7 +48,7 @@ L'etat detaille de chaque brique et les vues d'architecture sont dans
│ ├── migrations/ Migrations SQL versionnees
│ └── seeds/ Jeux de donnees de reference
├── etl/airflow/
│ ├── dags/ DAGs d'orchestration (pipeline ML, alertes)
│ ├── dags/ DAGs d'orchestration (pipeline ML, alertes, import historique)
│ ├── plugins/ Operateurs et hooks maison
│ ├── include/ Requetes SQL et ressources des DAGs
│ └── tests/ Tests d'integrite des DAGs
+1 -4
View File
@@ -1,4 +1,4 @@
# Conventions de tests unitaires : Frontend
# Conventions de tests unitaires — Frontend
## Outil
Vitest (intégré nativement à Angular CLI, pas d'installation à faire).
@@ -83,6 +83,3 @@ describe('MonComposant', () => {
## Lancer les tests
- Développement (mode watch) : `npm test`
- Rapport de couverture (CI) : `npm run test:ci -- --coverage`, puis ouvrir `coverage/index.html`
- Un fichier ou un dossier seulement :
`npx ng test --watch=false --coverage=false --include=src/app/core/services/alerts.service.spec.ts`
(répéter `--include` pour plusieurs cibles ; un dossier joue tous ses specs)
@@ -2,63 +2,53 @@ import { Alert } from '../../shared/models/alert.model';
export const ALERTS_FIXTURE: Alert[] = [
{
alert_id: 5,
alert_id: 'ALR-SITE002-1718458320',
timestamp: '2026-09-15T11:12:00',
site_id: 'SITE002',
timestamp: '2026-09-15T11:12:00Z',
type: 'threshold',
severity: 'critical',
message: 'Puissance appelée 812.5 kW au-dessus de la capacité du site (720.0 kW)',
type: 'outage',
message: 'Risque de surcharge sur Usine Lyon Vénissieux',
value: 812.5,
threshold: 720.0,
metric: 'consumption_kw',
prediction_id: null,
},
{
alert_id: 4,
alert_id: 'ALR-SITE003-1718458321',
timestamp: '2026-09-15T11:05:00',
site_id: 'SITE003',
timestamp: '2026-09-15T11:05:00Z',
type: 'outage',
severity: 'critical',
message: 'Aucune lecture depuis 5:00:00 (dernière lecture : 2026-09-15T06:05:00+00:00)',
value: null,
threshold: null,
metric: null,
prediction_id: null,
},
{
alert_id: 3,
site_id: 'SITE005',
timestamp: '2026-09-15T10:47:00Z',
type: 'spike',
severity: 'high',
message: 'Variation brutale entre deux lectures consécutives (260.0 kW -> 410.0 kW)',
value: 410.0,
threshold: 260.0,
metric: 'consumption_kw',
prediction_id: null,
},
{
alert_id: 2,
site_id: 'SITE006',
timestamp: '2026-09-15T10:30:00Z',
type: 'sensor',
severity: 'medium',
message: 'Qualité de mesure degraded (capteur hors ligne, valeur nulle)',
value: null,
threshold: null,
metric: null,
prediction_id: null,
message: 'Perte réseau totale sur Data Center Marseille',
value: 0,
threshold: 0,
},
{
alert_id: 1,
alert_id: 'ALR-SITE005-1718458322',
timestamp: '2026-09-15T10:47:00',
site_id: 'SITE005',
severity: 'high',
type: 'threshold',
message: 'Usine Toulouse approche de son seuil de capacité',
value: 410.0,
threshold: 480.0,
},
{
alert_id: 'ALR-SITE006-1718458323',
timestamp: '2026-09-15T10:30:00',
site_id: 'SITE006',
severity: 'medium',
type: 'sensor',
message: 'Capteur de température défaillant sur Bureau Lille',
value: 0,
threshold: 0,
},
{
alert_id: 'ALR-SITE004-1718458324',
timestamp: '2026-09-15T09:58:00',
site_id: 'SITE004',
timestamp: '2026-09-15T09:58:00Z',
type: 'anomaly',
severity: 'low',
message: 'Écart de 13% entre la consommation mesurée (62.0 kWh) et la prévision (55.0 kWh)',
type: 'anomaly',
message: 'Comportement de consommation inhabituel sur Bureau Bordeaux',
value: 62.0,
threshold: 55.0,
metric: 'consumption_kwh',
prediction_id: 42,
},
];
@@ -3,20 +3,6 @@ import { provideHttpClient } from '@angular/common/http';
import { provideHttpClientTesting, HttpTestingController } from '@angular/common/http/testing';
import { AlertsService } from './alerts.service';
import { environment } from '../../../environments/environment';
import { Alert } from '../../shared/models/alert.model';
const ALERT_API: Alert = {
alert_id: 1,
site_id: 'site-1',
timestamp: '2026-09-16T00:00:00Z',
type: 'threshold',
severity: 'high',
message: 'Dépassement du seuil configuré',
value: 812.5,
threshold: 720.0,
metric: 'consumption_kw',
prediction_id: null,
};
describe('AlertsService', () => {
let service: AlertsService;
@@ -32,35 +18,26 @@ describe('AlertsService', () => {
afterEach(() => httpMock.verify());
it("appelle le bon endpoint sans paramètre et retourne un tableau d'alertes", () => {
let result: Alert[] = [];
it("appelle le bon endpoint et retourne un tableau d'alertes", () => {
let result: unknown;
service.getAlerts().subscribe((r) => (result = r));
const req = httpMock.expectOne(
(r) => r.url === `${environment.apiUrl}/alerts` && r.method === 'GET',
);
expect(req.request.params.keys()).toEqual([]);
req.flush([ALERT_API]);
const req = httpMock.expectOne(`${environment.apiUrl}/alerts`);
expect(req.request.method).toBe('GET');
expect(result.length).toBe(1);
expect(result[0].alert_id).toBe(1);
expect(result[0].prediction_id).toBeNull();
});
req.flush([
{
alert_id: 'ALR-TEST-1',
timestamp: '2026-09-15T12:00:00',
site_id: 'SITE001',
severity: 'high',
type: 'threshold',
message: 'Test',
value: 100,
threshold: 90,
},
]);
it('transmet les filtres site_id et severity en paramètres de requête', () => {
service.getAlerts({ site_id: 'SITE001', severity: 'high' }).subscribe();
const req = httpMock.expectOne((r) => r.url === `${environment.apiUrl}/alerts`);
expect(req.request.params.get('site_id')).toBe('SITE001');
expect(req.request.params.get('severity')).toBe('high');
req.flush([]);
});
it('ne pose pas de paramètre pour un filtre omis', () => {
service.getAlerts({ site_id: 'SITE001' }).subscribe();
const req = httpMock.expectOne((r) => r.url === `${environment.apiUrl}/alerts`);
expect(req.request.params.has('severity')).toBe(false);
req.flush([]);
expect((result as unknown[]).length).toBe(1);
});
});
@@ -1,25 +1,13 @@
import { Service, inject } from '@angular/core';
import { HttpClient, HttpParams } from '@angular/common/http';
import { HttpClient } from '@angular/common/http';
import { environment } from '../../../environments/environment';
import { Alert, AlertSeverity } from '../../shared/models/alert.model';
export interface AlertFilters {
site_id?: string;
severity?: AlertSeverity;
}
import { Alert } from '../../shared/models/alert.model';
@Service()
export class AlertsService {
private http = inject(HttpClient);
getAlerts(filters: AlertFilters = {}) {
let params = new HttpParams();
if (filters.site_id) {
params = params.set('site_id', filters.site_id);
}
if (filters.severity) {
params = params.set('severity', filters.severity);
}
return this.http.get<Alert[]>(`${environment.apiUrl}/alerts`, { params });
getAlerts() {
return this.http.get<Alert[]>(`${environment.apiUrl}/alerts`);
}
}
@@ -1,37 +0,0 @@
import {
LIBELLE_PAR_SEVERITE,
LIBELLE_PAR_TYPE,
SEVERITES,
TON_PAR_SEVERITE,
TYPES_ALERTE,
UNITE_PAR_METRIQUE,
} from './alert-presentation';
describe('alert-presentation', () => {
it('distingue le ton des sévérités high et critical', () => {
expect(TON_PAR_SEVERITE.high).toBe('danger');
expect(TON_PAR_SEVERITE.critical).toBe('critical');
expect(TON_PAR_SEVERITE.high).not.toBe(TON_PAR_SEVERITE.critical);
});
it("n'affiche pas une alerte faible avec le ton de succès", () => {
expect(TON_PAR_SEVERITE.low).toBe('neutral');
expect(TON_PAR_SEVERITE.medium).toBe('warning');
});
it('donne un libellé français à chaque sévérité et à chaque type', () => {
for (const severite of SEVERITES) {
expect(LIBELLE_PAR_SEVERITE[severite]).toBeTruthy();
}
for (const type of TYPES_ALERTE) {
expect(LIBELLE_PAR_TYPE[type]).toBeTruthy();
}
expect(SEVERITES.length).toBe(4);
expect(TYPES_ALERTE.length).toBe(5);
});
it('associe une unité à chaque métrique du contrat', () => {
expect(UNITE_PAR_METRIQUE.consumption_kw).toBe('kW');
expect(UNITE_PAR_METRIQUE.consumption_kwh).toBe('kWh');
});
});
@@ -1,41 +0,0 @@
import { BadgeTone } from '../components/ui/badge/badge';
import { AlertMetric, AlertSeverity, AlertType } from './alert.model';
// Pourquoi : `low` en neutre plutôt qu'en vert, une alerte faible reste une alerte ; le vert se
// lisait comme « tout va bien » à côté des rouges.
export const TON_PAR_SEVERITE: Record<AlertSeverity, BadgeTone> = {
low: 'neutral',
medium: 'warning',
high: 'danger',
critical: 'critical',
};
export const LIBELLE_PAR_SEVERITE: Record<AlertSeverity, string> = {
low: 'Faible',
medium: 'Moyenne',
high: 'Élevée',
critical: 'Critique',
};
export const LIBELLE_PAR_TYPE: Record<AlertType, string> = {
spike: 'Pic de consommation',
threshold: 'Seuil dépassé',
anomaly: 'Anomalie',
outage: 'Coupure',
sensor: 'Capteur',
};
export const UNITE_PAR_METRIQUE: Record<AlertMetric, string> = {
consumption_kw: 'kW',
consumption_kwh: 'kWh',
};
export const SEVERITES: readonly AlertSeverity[] = ['low', 'medium', 'high', 'critical'];
export const TYPES_ALERTE: readonly AlertType[] = [
'spike',
'threshold',
'anomaly',
'outage',
'sensor',
];
@@ -1,16 +1,13 @@
export type AlertSeverity = 'low' | 'medium' | 'high' | 'critical';
export type AlertType = 'spike' | 'threshold' | 'anomaly' | 'outage' | 'sensor';
export type AlertMetric = 'consumption_kw' | 'consumption_kwh';
export interface Alert {
alert_id: number;
site_id: string;
alert_id: string;
timestamp: string;
type: AlertType;
site_id: string;
severity: AlertSeverity;
type: AlertType;
message: string;
value: number | null;
threshold: number | null;
metric: AlertMetric | null;
prediction_id: number | null;
value: number;
threshold: number;
}
-15
View File
@@ -29,18 +29,3 @@
color: var(--color-disabled);
margin-top: 0.25rem;
}
// Piège : le chevron est un SVG en data URI, où aucun token CSS n'est lisible ; sa couleur
// reprend en dur la valeur de --color-text-muted.
.form-select {
@extend .form-input;
padding-right: 2.25rem;
color: var(--color-text);
background-color: var(--color-surface);
background-image: url("data:image/svg+xml,%3Csvg xmlns='http://www.w3.org/2000/svg' viewBox='0 0 20 20' fill='none' stroke='%236b7280' stroke-width='1.5' stroke-linecap='round' stroke-linejoin='round'%3E%3Cpath d='M6 8l4 4 4-4'/%3E%3C/svg%3E");
background-repeat: no-repeat;
background-position: right 0.6rem center;
background-size: 1rem;
appearance: none;
cursor: pointer;
}
-1
View File
@@ -3,7 +3,6 @@
{
"compileOnSave": false,
"compilerOptions": {
"strict": true,
"noImplicitOverride": true,
"noPropertyAccessFromIndexSignature": true,
"noImplicitReturns": true,
+1
View File
@@ -36,6 +36,7 @@ x-airflow-common: &airflow-common
volumes:
- ./etl/airflow/dags:/opt/airflow/dags
- ./etl/airflow/plugins:/opt/airflow/plugins
- ./data/raw:/opt/data/raw:ro
- airflow_logs:/opt/airflow/logs
- airflow_ml_state:/opt/ml/state
restart: unless-stopped
+5 -4
View File
@@ -70,10 +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
[30-frontend.md](30-frontend.md).
Le lien `airflow --> db` est maintenant en trait plein : trois DAGs tournent, deux pour
Le lien `airflow --> db` est maintenant en trait plein : quatre 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
génération des recommandations (issue #116), cf. plus bas et [20-backend.md](20-backend.md). Le
reste du périmètre Airflow envisagé (ingestion, issues #15/#16) reste en pointillé, non construit.
génération des recommandations (issue #116), et `historical_import` pour l'ingestion du dataset
historique (issue #119). L'orchestration de l'import API Mock et la réconciliation globale des
deux sources restent à compléter dans l'issue #15.
Le lien `prom -.-> api` de même : l'API expose bien `/metrics` au format Prometheus, mais aucun
collecteur ne vient le lire.
@@ -88,7 +89,7 @@ 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 |
| 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)). 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 |
| ETL | Apache Airflow | `etl/airflow` | `En cours` | Webserver + scheduler (LocalExecutor) tournent via docker-compose, base de métadonnées Postgres dédiée. Trois DAGs en sous-processus `uv run` : `ml_train` manuel et `ml_score` `@hourly` pour le pipeline ML (issue #115), `alertes` à `15 * * * *` pour la détection et les recommandations (issue #116, [ADR 0008](../adr/0008-airflow-execute-le-code-du-backend.md)). L'ingestion (issues #15/#16) n'a pas encore de DAG |
| 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 |
| CI/CD | GitHub Actions | `.github/workflows` | `En cours` | 5 workflows, 16 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. Détail dans [50-cicd.md](50-cicd.md). **Aucun job de déploiement** (#21) |
## Flux bout en bout
+7 -1
View File
@@ -51,7 +51,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 du webserver :
c'est le scheduler qui a besoin du volume `airflow_ml_state` (modèle, magasin MLflow).
### Airflow (issues #115 et #116)
### Airflow (issues #115, #116 et #119)
Trois services, `docker compose profiles` non utilisés (démarrage explicite via `make
airflow-up`, pas dans `make dev`) :
@@ -75,6 +75,12 @@ l'[ADR 0008](../adr/0008-airflow-execute-le-code-du-backend.md).
| `ml_train` | manuelle | `enervision_ml.train`, 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` |
| `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` |
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
`./data/raw:/opt/data/raw:ro` permet au scheduler de lire les fichiers CSV/JSON sans pouvoir les
modifier.
**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
@@ -24,13 +24,12 @@ seule fois dans `src/styles.scss`. Disponibles partout sans import supplémentai
| `--shadow-card` | Ombre portée des cartes |
| `--space-1` à `--space-5` | Échelle d'espacement (0.35rem à 2.5rem) |
Les classes de formulaire partagées (`.form-label`, `.form-input`, `.form-select`, `.form-hint`)
sont dans `apps/frontend/src/styles/_forms.scss`, importées globalement de la même façon. Elles
s'appliquent directement à des `<label>`/`<input>`/`<select>` natifs, liés par `formControlName` ou
par un simple `(change)` : pas de composant `ControlValueAccessor` dédié, le gain n'en vaut pas la
complexité pour des formulaires aussi simples que ceux de ce projet. `.form-select` habille un
`<select>` natif avec la bordure et le focus de `.form-input`, plus un chevron. Les erreurs de
formulaire, elles, s'affichent via `<ev-alert severity="danger">`, pas une classe dédiée.
Les classes de formulaire partagées (`.form-label`, `.form-input`, `.form-hint`) sont dans
`apps/frontend/src/styles/_forms.scss`, importées globalement de la même façon. Elles
s'appliquent directement à des `<label>`/`<input>` natifs liés par `formControlName` : pas de
composant `ControlValueAccessor` dédié, le gain n'en vaut pas la complexité pour des formulaires
aussi simples que ceux de ce projet. Les erreurs de formulaire, elles, s'affichent via
`<ev-alert severity="danger">`, pas une classe dédiée.
La classe `.auth-page` (`apps/frontend/src/styles/_auth-page.scss`, importée globalement) porte
le fond dégradé et le centrage commun aux pages d'authentification (`login`, `change-password`,
+10 -2
View File
@@ -663,8 +663,16 @@ 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.
Airflow tourne désormais réellement (`etl/airflow/`, `make airflow-up`) et orchestre le pipeline ML (`ml_train`/`ml_score`, issue #115) ainsi que la détection d'alertes et la génération des recommandations (`alertes`, issue #116). Il n'orchestre pas encore ces deux imports : `historical_import.py` et `mock_api_import.py` (normalisation et chargement micro-batch, issues #15/#16) restent à faire.
Airflow tourne désormais réellement (`etl/airflow/`, `make airflow-up`) et orchestre le pipeline
ML (`ml_train`/`ml_score`, issue #115), la détection d'alertes et la génération des
recommandations (`alertes`, issue #116), ainsi que l'import historique
(`historical_import`, issue #119).
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` montrent le patron retenu (des `BashOperator` qui invoquent le script tel quel, dans l'environnement `uv` que l'image embarque pour lui).
Le DAG `historical_import` est déclenché manuellement. Il exécute
`app.etl.historical_import` avec les fichiers montés en lecture seule depuis `data/raw` vers
`/opt/data/raw`. L'orchestration de l'import API Mock et la réconciliation globale des deux
sources restent couvertes par 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 pipeline Data servira ensuite à préparer les données nécessaires au modèle de Machine Learning.
+45
View File
@@ -0,0 +1,45 @@
"""DAG d'import du dataset historique EnerVision (issue #119).
Orchestre le pipeline existant `app.etl.historical_import` sans dupliquer sa logique ETL.
Le dataset historique sert à initialiser l'environnement : le DAG reste donc manuel.
Le backend est exécuté dans l'environnement `/opt/backend` embarqué dans l'image Airflow,
sur le même patron que le DAG `alertes` (ADR 0008).
"""
from __future__ import annotations
from datetime import datetime, timedelta
from airflow.models.dag import DAG
from airflow.operators.bash import BashOperator
COMMANDE_BACKEND = "cd /opt/backend && env -u VIRTUAL_ENV uv run --no-sync python -m"
CSV_PATH = "/opt/data/raw/all_sites_combined.csv"
METADATA_PATH = "/opt/data/raw/dataset_metadata.json"
SOURCE_TIMEZONE = "UTC"
BATCH_SIZE = 1000
with DAG(
dag_id="historical_import",
description="Importe le dataset historique CSV/JSON dans dataset, site et reading.",
schedule=None,
start_date=datetime(2026, 1, 1),
catchup=False,
max_active_runs=1,
tags=["etl", "historical"],
) as dag:
BashOperator(
task_id="import_historical",
bash_command=(
f"{COMMANDE_BACKEND} app.etl.historical_import "
f"--csv {CSV_PATH} "
f"--metadata {METADATA_PATH} "
"--source-timezone UTC "
"--batch-size 1000"
),
retries=1,
retry_delay=timedelta(minutes=2),
execution_timeout=timedelta(minutes=30),
)
+27 -1
View File
@@ -10,12 +10,13 @@ from airflow.models.dagbag import DagBag
DAGS_FOLDER = Path(__file__).resolve().parent.parent / "dags"
DAG_IDS = ["ml_train", "ml_score", "alertes"]
DAG_IDS = ["ml_train", "ml_score", "alertes", "historical_import"]
TACHES = [
("ml_train", "train"),
("ml_score", "score"),
("alertes", "detection"),
("alertes", "recommandations"),
("historical_import", "import_historical"),
]
@@ -47,6 +48,10 @@ def test_alertes_runs_after_the_hourly_scoring(dagbag: DagBag) -> None:
assert dagbag.dags["alertes"].timetable.summary == "15 * * * *"
def test_historical_import_has_no_schedule(dagbag: DagBag) -> None:
assert dagbag.dags["historical_import"].timetable.summary == "None"
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
@@ -67,12 +72,29 @@ def test_alertes_recommendation_task_calls_the_backend_cli(dagbag: DagBag) -> No
assert "app.cli generate-recommendations" in tache.bash_command
def test_historical_import_calls_the_existing_backend_module(dagbag: DagBag) -> None:
tache = dagbag.dags["historical_import"].get_task("import_historical")
assert "app.etl.historical_import" in tache.bash_command
def test_historical_import_uses_the_expected_source_files(dagbag: DagBag) -> None:
commande = dagbag.dags["historical_import"].get_task("import_historical").bash_command
assert "--csv /opt/data/raw/all_sites_combined.csv" in commande
assert "--metadata /opt/data/raw/dataset_metadata.json" in commande
@pytest.mark.parametrize("task_id", ["detection", "recommandations"])
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).
assert "/opt/backend" in dagbag.dags["alertes"].get_task(task_id).bash_command
def test_historical_import_runs_in_the_backend_environment(dagbag: DagBag) -> None:
commande = dagbag.dags["historical_import"].get_task("import_historical").bash_command
assert "/opt/backend" in commande
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
# tant que la detection n'a pas ecrit.
@@ -133,6 +155,10 @@ def test_alertes_retries_after_a_transient_failure(dagbag: DagBag, task_id: str)
assert dagbag.dags["alertes"].get_task(task_id).retries >= 1
def test_historical_import_retries_after_a_transient_failure(dagbag: DagBag) -> None:
assert dagbag.dags["historical_import"].get_task("import_historical").retries >= 1
@pytest.mark.parametrize(("dag_id", "task_id"), TACHES)
def test_tasks_never_resync_the_baked_environment(
dagbag: DagBag, dag_id: str, task_id: str