docs: documente l'ingestion depuis l'API Mock
This commit is contained in:
+332
-69
@@ -6,10 +6,14 @@ système qui en découle.
|
||||
|
||||
## Ce que couvre ce document
|
||||
|
||||
**Dix tables applicatives existent** : quatre pour l'authentification, six pour les données
|
||||
d'énergie, dont l'hypertable `reading`. Les sections marquées `Fait` relèvent le code. Celles
|
||||
marquées `Cible` décrivent ce qui n'est pas écrit, au premier rang desquelles la chaîne
|
||||
d'ingestion, les agrégats continus, la compression et la rétention.
|
||||
**Douze tables applicatives existent** : six pour l'authentification et six pour les données
|
||||
d'énergie, dont l'hypertable `reading`.
|
||||
|
||||
Les sections marquées `Fait` relèvent du code déjà implémenté. Les sections marquées `Cible`
|
||||
décrivent les éléments prévus mais pas encore réalisés.
|
||||
|
||||
L'ingestion des deux sources de données du MVP est maintenant implémentée. L'orchestration
|
||||
Airflow, les agrégats continus, la compression et la rétention restent des cibles.
|
||||
|
||||
## Trois emplacements, trois rôles
|
||||
|
||||
@@ -35,8 +39,16 @@ Statut : `Fait`.
|
||||
- `db/init/100-extensions.sql` crée l'extension `timescaledb`.
|
||||
- `db/init/110-test-database.sql` crée `enervision_test`, dont le nom est attendu en dur par
|
||||
`apps/backend/tests/conftest.py`.
|
||||
- Cinq révisions Alembic. La première, `5353c0e4f094`, **ne crée aucune table** : elle
|
||||
établit `alembic_version` et refuse de s'appliquer si l'extension manque :
|
||||
- Six révisions Alembic sont actuellement appliquées.
|
||||
- La première, `5353c0e4f094`, **ne crée aucune table** : elle établit `alembic_version`
|
||||
et refuse de s'appliquer si l'extension TimescaleDB manque.
|
||||
- Les révisions suivantes créent les tables liées à l'authentification :
|
||||
`app_user`, `login_attempt`, `audit_log` et `refresh_token`.
|
||||
- La révision `e6d2026091501` crée les six tables Data et déclare l'hypertable `reading`.
|
||||
- La révision `c0adab96238c` ajoute les tables `password_reset_attempt`
|
||||
et `password_reset_token`.
|
||||
|
||||
La garde de la première migration est :
|
||||
|
||||
```sql
|
||||
IF NOT EXISTS (SELECT 1 FROM pg_extension WHERE extname = 'timescaledb') THEN
|
||||
@@ -47,37 +59,58 @@ END IF;
|
||||
Cette garde forme paire avec le 503 de `/api/v1/health/ready`. Un bootstrap sauté ne se voit pas
|
||||
au démarrage de l'API : ces deux gardes le rendent visible tôt, des deux côtés.
|
||||
|
||||
Les trois suivantes créent les tables de l'authentification, décrites plus bas : `app_user`,
|
||||
puis `login_attempt` et `audit_log`, puis `refresh_token`. La cinquième, `e6d2026091501`, crée
|
||||
les six tables de données décrites en fin de document et déclare l'hypertable `reading`.
|
||||
|
||||
## Cycle de vie d'une mesure
|
||||
|
||||
Statut : `Cible`, sauf l'hypertable `reading` qui existe. Ni l'ingestion, ni les agrégats
|
||||
continus, ni la compression, ni la rétention ne sont écrits.
|
||||
Statut : `Partiellement fait`.
|
||||
|
||||
Les mécanismes d'ingestion sont maintenant implémentés pour les deux sources de données du MVP :
|
||||
|
||||
- le dataset historique CSV/JSON avec `historical_import.py` ;
|
||||
- l'API Mock avec `mock_api_import.py`.
|
||||
|
||||
Les traitements sont actuellement exécutables directement depuis le backend.
|
||||
|
||||
L'orchestration avec Apache Airflow reste une cible, tout comme les agrégats continus,
|
||||
la compression et les politiques de rétention.
|
||||
|
||||
```mermaid
|
||||
flowchart LR
|
||||
src["Source de mesures"] -.-> ing["Ingestion Airflow"]
|
||||
ing -.-> hy[("Hypertable reading")]
|
||||
csv["CSV + JSON"] --> hist["historical_import.py"]
|
||||
mock["API Mock"] --> api["mock_api_import.py"]
|
||||
|
||||
hist --> hy[("Hypertable reading")]
|
||||
api --> hy
|
||||
|
||||
airflow["Airflow"] -.-> hist
|
||||
airflow -.-> api
|
||||
|
||||
hy -.-> agg[("Agrégat continu")]
|
||||
hy -.-> comp["Compression"]
|
||||
hy -.-> ret["Rétention"]
|
||||
agg -.-> api["API FastAPI"]
|
||||
|
||||
agg -.-> backend["API FastAPI"]
|
||||
agg -.-> graf["Grafana"]
|
||||
```
|
||||
|
||||
Les lectures de l'API et de Grafana visent l'agrégat continu, pas la table brute : c'est tout
|
||||
l'intérêt de TimescaleDB, et cela doit rester vrai quand les volumes augmenteront.
|
||||
Les flèches pleines représentent les traitements actuellement implémentés.
|
||||
|
||||
Les flèches pointillées représentent les éléments encore prévus comme cibles.
|
||||
|
||||
Les lectures futures de l'API et de Grafana visent l'agrégat continu plutôt que la table brute
|
||||
lorsque cette partie TimescaleDB sera mise en place.
|
||||
|
||||
## Tables d'authentification
|
||||
|
||||
Statut : `Fait`. Elles ne sont pas des séries temporelles et n'ont donc rien à voir avec les
|
||||
hypertables ; elles vivent dans `apps/backend/alembic/`, qui porte le schéma exposé par l'API.
|
||||
Statut : `Fait`.
|
||||
|
||||
Elles ne sont pas des séries temporelles et n'ont donc rien à voir avec les hypertables ;
|
||||
elles vivent dans `apps/backend/alembic/`, qui porte le schéma exposé par l'API.
|
||||
|
||||
```mermaid
|
||||
erDiagram
|
||||
APP_USER ||--o{ REFRESH_TOKEN : ouvre
|
||||
APP_USER ||--o{ PASSWORD_RESET_TOKEN : recoit
|
||||
|
||||
APP_USER {
|
||||
uuid id PK
|
||||
string email UK
|
||||
@@ -88,6 +121,7 @@ erDiagram
|
||||
bool must_change_password
|
||||
timestamptz credentials_changed_at
|
||||
}
|
||||
|
||||
REFRESH_TOKEN {
|
||||
uuid id PK
|
||||
uuid family_id
|
||||
@@ -99,6 +133,7 @@ erDiagram
|
||||
text revoked_reason
|
||||
uuid replaced_by
|
||||
}
|
||||
|
||||
LOGIN_ATTEMPT {
|
||||
bigint id PK
|
||||
timestamptz occurred_at
|
||||
@@ -106,6 +141,7 @@ erDiagram
|
||||
inet client_ip
|
||||
text outcome
|
||||
}
|
||||
|
||||
AUDIT_LOG {
|
||||
bigint id PK
|
||||
timestamptz occurred_at
|
||||
@@ -114,31 +150,47 @@ erDiagram
|
||||
text action
|
||||
jsonb detail
|
||||
}
|
||||
|
||||
PASSWORD_RESET_ATTEMPT {
|
||||
bigint id PK
|
||||
timestamptz occurred_at
|
||||
string email_tried
|
||||
inet client_ip
|
||||
}
|
||||
|
||||
PASSWORD_RESET_TOKEN {
|
||||
uuid id PK
|
||||
uuid user_id FK
|
||||
bytea token_hash UK
|
||||
timestamptz issued_at
|
||||
timestamptz expires_at
|
||||
timestamptz consumed_at
|
||||
inet client_ip
|
||||
text user_agent
|
||||
}
|
||||
```
|
||||
|
||||
Quatre choix de modélisation portent une intention et se défendent seuls :
|
||||
Plusieurs choix de modélisation portent une intention précise :
|
||||
|
||||
- **`app_user` et non `user`** : `user` est un mot réservé PostgreSQL, raccourci de
|
||||
`CURRENT_USER`. Le nom rappelle en prime qu'il s'agit d'un compte applicatif, par opposition
|
||||
au rôle PostgreSQL qui portera le cantonnement de l'ETL.
|
||||
- **`credentials_changed_at`, une seule colonne**, couvre le changement de mot de passe, le
|
||||
changement de rôle et la désactivation. Un compteur de version ne dirait rien à un humain qui
|
||||
lit un audit.
|
||||
- **`refresh_token.expires_at` est absolu et hérité** du prédécesseur à chaque rotation. S'il
|
||||
glissait, la promesse de sept jours serait fictive et une session active ne finirait jamais.
|
||||
- **`audit_log.actor_id` n'a aucune clé étrangère**, et `actor_email` comme `actor_role` sont
|
||||
dénormalisés. Une contrainte `ON DELETE SET NULL` déclencherait un `UPDATE` que le déclencheur
|
||||
d'ajout seul refuserait. Voir l'[ADR 0004](../adr/0004-journal-d-audit-en-ajout-seul.md).
|
||||
`CURRENT_USER`. Le nom rappelle aussi qu'il s'agit d'un compte applicatif.
|
||||
- **`credentials_changed_at`, une seule colonne**, couvre notamment le changement de mot de passe,
|
||||
le changement de rôle et la désactivation.
|
||||
- **`refresh_token.expires_at` est absolu et hérité** du prédécesseur à chaque rotation.
|
||||
- **`audit_log.actor_id` n'a aucune clé étrangère** afin de conserver les informations d'audit
|
||||
même si l'entité d'origine évolue.
|
||||
- `password_reset_token` ne stocke que l'empreinte du jeton et jamais sa valeur directement.
|
||||
- `password_reset_attempt` est séparée de `audit_log`, car son volume peut être piloté
|
||||
par des demandes externes répétées.
|
||||
|
||||
`audit_log` porte deux déclencheurs qui refusent `UPDATE`, `DELETE` et `TRUNCATE`. Elle n'est
|
||||
donc **pas** une hypertable : une politique de rétention émettrait des `DELETE` qu'ils
|
||||
refuseraient. `login_attempt`, à l'inverse, est faite pour se purger, puisque son volume est
|
||||
piloté par l'attaquant.
|
||||
`audit_log` porte des déclencheurs qui refusent `UPDATE`, `DELETE` et `TRUNCATE`.
|
||||
Elle n'est donc **pas** une hypertable.
|
||||
|
||||
## Gabarit de révision créant une hypertable
|
||||
|
||||
Conforme à la règle de l'ADR 0001 : table et hypertable dans la même révision. La révision
|
||||
`e6d2026091501` en est l'exemple réel, réduit ici à l'essentiel.
|
||||
Conforme à la règle de l'ADR 0001 : table et hypertable dans la même révision.
|
||||
|
||||
La révision `e6d2026091501` en est l'exemple réel, réduit ici à l'essentiel.
|
||||
|
||||
```python
|
||||
def upgrade() -> None:
|
||||
@@ -182,24 +234,28 @@ colonne de temps : les index déclarés dans la révision le couvrent déjà.
|
||||
|
||||
## Questions ouvertes
|
||||
|
||||
Elles relèvent du jalon J2, « valider le périmètre retenu ». Le schéma est livré : ce qui suit
|
||||
porte sur son exploitation, plus sur sa forme.
|
||||
Elles portent maintenant principalement sur l'exploitation du schéma :
|
||||
|
||||
- **Quelle granularité** à l'ingestion : la seconde, la minute, le quart d'heure.
|
||||
- **Quels agrégats continus**, et sur quelles fenêtres.
|
||||
- **Quelle profondeur de rétention** en données brutes, et à partir de quand on compresse.
|
||||
- **Multi-tenant ou non** : un site appartient-il à un client, et faut-il cloisonner les lectures.
|
||||
- **Quelle granularité** conserver à long terme à l'ingestion : seconde, minute ou quart d'heure.
|
||||
- **Quels agrégats continus** créer et sur quelles fenêtres.
|
||||
- **Quelle profondeur de rétention** conserver en données brutes et à partir de quand compresser.
|
||||
- **Multi-tenant ou non** : un site appartient-il à un client et faut-il cloisonner les lectures.
|
||||
|
||||
## Modélisation détaillée des données
|
||||
|
||||
Cette modélisation prend en compte les fichiers CSV historiques,
|
||||
leurs métadonnées JSON et les données de l’API Mock.
|
||||
Elle comprend six tables, depuis le stockage des mesures
|
||||
jusqu’aux recommandations proposées à l’utilisateur.
|
||||
Cette modélisation prend en compte :
|
||||
|
||||
- les fichiers CSV historiques ;
|
||||
- leurs métadonnées JSON ;
|
||||
- les données de l'API Mock.
|
||||
|
||||
Elle comprend six tables Data, depuis le stockage des mesures jusqu'aux recommandations proposées
|
||||
à l'utilisateur.
|
||||
|
||||
### Schéma de données
|
||||
|
||||
Le diagramme ci-dessous présente les tables et leurs relations.
|
||||
|
||||
La révision `e6d2026091501` les crée.
|
||||
|
||||

|
||||
@@ -208,21 +264,20 @@ La révision `e6d2026091501` les crée.
|
||||
|
||||
### Description des tables
|
||||
|
||||
Chaque table remplit un rôle précis dans le traitement et l’exploitation
|
||||
des données.
|
||||
Chaque table remplit un rôle précis dans le traitement et l'exploitation des données.
|
||||
|
||||
| Table | Rôle | Origine des informations |
|
||||
|---|---|---|
|
||||
| `dataset` | Identifier les jeux historiques, retrouver leurs fichiers et conserver leurs métadonnées | Archive CSV/JSON et informations ajoutées lors de l’import |
|
||||
| `dataset` | Identifier les jeux historiques, retrouver leurs fichiers et conserver leurs métadonnées | Archive CSV/JSON et informations ajoutées lors de l'import |
|
||||
| `site` | Regrouper les informations des sites : identifiant, nom, type et caractéristiques disponibles | CSV et API Mock `/api/v1/sites` |
|
||||
| `reading` | Stocker les mesures, leur provenance, leur qualité et les éventuelles valeurs imputées | CSV et API Mock `/current` et `/readings` |
|
||||
| `prediction` | Conserver les prévisions, leur période cible et la référence du modèle utilisé | Traitements ML d’EnerVision |
|
||||
| `prediction` | Conserver les prévisions, leur période cible et la référence du modèle utilisé | Traitements ML d'EnerVision |
|
||||
| `alert` | Enregistrer les alertes, leur type, leur gravité et leur message | API Mock `/alerts` et détections EnerVision |
|
||||
| `recommendation` | Proposer des actions et expliquer la règle qui les motive | Règles métier d’EnerVision |
|
||||
| `recommendation` | Proposer des actions et expliquer la règle qui les motive | Règles métier d'EnerVision |
|
||||
|
||||
Les anomalies historiques décrites dans les JSON sont conservées
|
||||
dans `dataset.metadata`. Elles servent à l’analyse des données
|
||||
et ne sont pas considérées comme des alertes actuelles.
|
||||
Les anomalies historiques décrites dans les JSON sont conservées dans `dataset.metadata`.
|
||||
|
||||
Elles servent à l'analyse des données et ne sont pas considérées comme des alertes actuelles.
|
||||
|
||||
### Relations entre les tables
|
||||
|
||||
@@ -232,15 +287,21 @@ et ne sont pas considérées comme des alertes actuelles.
|
||||
- Une alerte peut être associée à une prévision du même site.
|
||||
- Une alerte peut donner lieu à plusieurs recommandations.
|
||||
|
||||
## Ingestion des données historiques
|
||||
# Ingestion des données historiques
|
||||
|
||||
Le MVP EnerVision initialise les données énergétiques à partir du dataset fourni dans le cadre du projet.
|
||||
Statut : `Fait`.
|
||||
|
||||
Le dataset de référence contient 122 647 mesures issues de 7 sites et couvre la période du 1er janvier 2023 au 31 décembre 2024.
|
||||
Le MVP EnerVision initialise les données énergétiques à partir du dataset fourni dans le cadre
|
||||
du projet.
|
||||
|
||||
Les fichiers sources CSV et JSON sont nécessaires uniquement pour l'initialisation des données. Ils ne sont pas versionnés dans Git et sont placés localement dans `data/raw/`.
|
||||
Le dataset de référence contient 122 647 mesures issues de 7 sites et couvre la période
|
||||
du 1er janvier 2023 au 31 décembre 2024.
|
||||
|
||||
### Architecture du flux
|
||||
Les fichiers sources CSV et JSON sont nécessaires uniquement pour l'initialisation des données.
|
||||
|
||||
Ils ne sont pas versionnés dans Git et sont placés localement dans `data/raw/`.
|
||||
|
||||
## Architecture du flux historique
|
||||
|
||||
```text
|
||||
Dataset CSV + métadonnées JSON
|
||||
@@ -271,17 +332,26 @@ Dataset CSV + métadonnées JSON
|
||||
|
||||
Le pipeline est développé en Python.
|
||||
|
||||
Pandas est utilisé pour l'extraction, la validation et la préparation des données. SQLAlchemy Async assure le chargement transactionnel dans PostgreSQL/TimescaleDB.
|
||||
Pandas est utilisé pour l'extraction, la validation et la préparation des données.
|
||||
|
||||
SQLAlchemy Async assure le chargement transactionnel dans PostgreSQL/TimescaleDB.
|
||||
|
||||
Une empreinte SHA-256 permet d'identifier le dataset utilisé et d'assurer sa traçabilité.
|
||||
|
||||
Les valeurs manquantes sont conservées pendant l'ingestion afin de préserver les données sources. Aucune imputation n'est réalisée à cette étape.
|
||||
Les valeurs manquantes sont conservées pendant l'ingestion afin de préserver les données sources.
|
||||
|
||||
Aucune imputation n'est réalisée à cette étape.
|
||||
|
||||
Le chargement des mesures est effectué par batches de 1 000 lignes.
|
||||
|
||||
Les données provenant du dataset CSV sont identifiées par `source = "csv"` et associées à leur `dataset_id`.
|
||||
Les données provenant du dataset CSV sont identifiées par :
|
||||
|
||||
### Résultats validés
|
||||
```text
|
||||
source = "csv"
|
||||
dataset_id = identifiant du dataset
|
||||
```
|
||||
|
||||
## Résultats validés pour l'historique
|
||||
|
||||
Le chargement de référence a permis d'obtenir :
|
||||
|
||||
@@ -290,14 +360,207 @@ Le chargement de référence a permis d'obtenir :
|
||||
- 122 647 mesures ;
|
||||
- 0 doublon détecté dans le dataset source.
|
||||
|
||||
L'idempotence a également été vérifiée par une deuxième exécution du pipeline : aucune nouvelle mesure n'a été créée et le nombre de `reading` est resté à 122 647.
|
||||
L'idempotence a également été vérifiée par une deuxième exécution du pipeline :
|
||||
aucune nouvelle mesure n'a été créée et le nombre de `reading` est resté à 122 647.
|
||||
|
||||
La procédure détaillée d'installation, d'exécution, de validation et de contrôle du pipeline est disponible dans `etl/README.md`.
|
||||
La procédure détaillée d'installation, d'exécution, de validation et de contrôle du pipeline
|
||||
est disponible dans `etl/README.md`.
|
||||
|
||||
### Évolution prévue
|
||||
# Ingestion depuis l'API Mock
|
||||
|
||||
L'étape suivante consiste à orchestrer les traitements Data avec Apache Airflow.
|
||||
Statut : `Fait`.
|
||||
|
||||
L'orchestration réutilisera la logique ETL existante afin de séparer la logique de traitement de la planification, du suivi des exécutions et de la gestion des erreurs.
|
||||
La deuxième source du pipeline Data est l'API Mock EnerVision.
|
||||
|
||||
Le pipeline servira ensuite de base à la préparation des données nécessaires au modèle de Machine Learning.
|
||||
Le traitement est implémenté dans :
|
||||
|
||||
```text
|
||||
apps/backend/app/etl/mock_api_import.py
|
||||
```
|
||||
|
||||
## Endpoints utilisés
|
||||
|
||||
Le pipeline récupère les informations des sites depuis :
|
||||
|
||||
```text
|
||||
GET /api/v1/sites
|
||||
```
|
||||
|
||||
puis les mesures historiques simulées depuis :
|
||||
|
||||
```text
|
||||
GET /api/v1/readings
|
||||
```
|
||||
|
||||
Pour `/api/v1/readings`, les informations suivantes sont envoyées :
|
||||
|
||||
```text
|
||||
site_id
|
||||
start_time
|
||||
end_time
|
||||
limit
|
||||
```
|
||||
|
||||
Les paramètres de ligne de commande disponibles pour l'import sont :
|
||||
|
||||
```text
|
||||
--start-time
|
||||
--end-time
|
||||
--limit
|
||||
--dry-run
|
||||
```
|
||||
|
||||
## Flux d'ingestion API Mock
|
||||
|
||||
```text
|
||||
API Mock
|
||||
|
|
||||
+-----+------+
|
||||
| |
|
||||
v v
|
||||
/sites /readings
|
||||
| |
|
||||
+-----+------+
|
||||
|
|
||||
v
|
||||
mock_api_import.py
|
||||
|
|
||||
v
|
||||
Transformation
|
||||
+ qualité data
|
||||
|
|
||||
v
|
||||
PostgreSQL / TimescaleDB
|
||||
| |
|
||||
v v
|
||||
site reading
|
||||
```
|
||||
|
||||
Les informations des sites sont insérées ou mises à jour dans `site`.
|
||||
|
||||
Les mesures sont enregistrées dans l'hypertable `reading` avec :
|
||||
|
||||
```text
|
||||
source = "api_history"
|
||||
dataset_id = NULL
|
||||
```
|
||||
|
||||
Les données provenant de l'API Mock ne sont donc pas associées à un enregistrement de la table
|
||||
`dataset`.
|
||||
|
||||
La réponse source reçue depuis l'API est conservée dans :
|
||||
|
||||
```text
|
||||
raw_data
|
||||
```
|
||||
|
||||
## Qualité des données de l'API Mock
|
||||
|
||||
Les valeurs `NULL` ne sont pas remplacées pendant l'ingestion.
|
||||
|
||||
Les informations suivantes fournies par l'API sont conservées :
|
||||
|
||||
```text
|
||||
data_quality
|
||||
null_reasons
|
||||
```
|
||||
|
||||
Cette conservation permet de distinguer une valeur manquante d'une valeur réelle égale à zéro
|
||||
et de garder les informations liées aux éventuelles défaillances de capteurs.
|
||||
|
||||
Aucune imputation n'est réalisée pendant cette phase :
|
||||
|
||||
```text
|
||||
imputed_values = NULL
|
||||
imputation_method = NULL
|
||||
```
|
||||
|
||||
## Validation de l'import API Mock
|
||||
|
||||
Un scénario de validation a été exécuté pour les 7 sites sur la période :
|
||||
|
||||
```text
|
||||
15/06/2024 12:00 UTC
|
||||
à
|
||||
15/06/2024 13:00 UTC
|
||||
```
|
||||
|
||||
avec :
|
||||
|
||||
```text
|
||||
limit = 60
|
||||
```
|
||||
|
||||
Résultat :
|
||||
|
||||
```text
|
||||
7 sites
|
||||
60 lectures par site
|
||||
420 lectures récupérées
|
||||
```
|
||||
|
||||
Les données ont été chargées dans PostgreSQL/TimescaleDB puis contrôlées directement en base.
|
||||
|
||||
Les contrôles ont confirmé :
|
||||
|
||||
- `source = "api_history"` ;
|
||||
- `dataset_id = NULL` ;
|
||||
- la conservation des valeurs `NULL` ;
|
||||
- la conservation de `data_quality` ;
|
||||
- la conservation de `null_reasons` ;
|
||||
- la conservation de `raw_data`.
|
||||
|
||||
L'idempotence a été vérifiée en rejouant le même import.
|
||||
|
||||
Une mesure déjà présente n'est pas ajoutée une seconde fois.
|
||||
|
||||
Les tests automatisés couvrent également :
|
||||
|
||||
- la récupération des sites ;
|
||||
- les paramètres envoyés à `/api/v1/readings` ;
|
||||
- les réponses HTTP en erreur ;
|
||||
- le format de la réponse ;
|
||||
- la transformation des mesures ;
|
||||
- les valeurs manquantes ;
|
||||
- la qualité des données ;
|
||||
- la conservation des données sources ;
|
||||
- l'idempotence en base.
|
||||
|
||||
# Évolution prévue
|
||||
|
||||
La prochaine étape consiste à orchestrer les deux mécanismes d'ingestion avec Apache Airflow.
|
||||
|
||||
```text
|
||||
CSV / JSON ----------------+
|
||||
|
|
||||
v
|
||||
+------------------+
|
||||
| Airflow |
|
||||
+------------------+
|
||||
|
|
||||
+----------------+----------------+
|
||||
| |
|
||||
v v
|
||||
historical_import.py mock_api_import.py
|
||||
| |
|
||||
+----------------+----------------+
|
||||
|
|
||||
v
|
||||
PostgreSQL / TimescaleDB
|
||||
```
|
||||
|
||||
Airflow servira à :
|
||||
|
||||
- planifier les traitements ;
|
||||
- définir leur ordre d'exécution ;
|
||||
- suivre leur état ;
|
||||
- gérer et remonter les erreurs ;
|
||||
- faciliter les exécutions récurrentes.
|
||||
|
||||
Airflow ne remplacera pas la logique ETL déjà implémentée.
|
||||
|
||||
Les scripts Python resteront responsables de l'extraction, de la validation, de la transformation
|
||||
et du chargement des données.
|
||||
|
||||
Le pipeline servira ensuite de base à la préparation des données nécessaires au modèle
|
||||
de Machine Learning.
|
||||
+307
-22
@@ -2,9 +2,10 @@
|
||||
|
||||
## Objectif
|
||||
|
||||
Le pipeline ETL EnerVision permet d'intégrer les données énergétiques historiques dans PostgreSQL/TimescaleDB.
|
||||
Le pipeline ETL EnerVision permet d'intégrer les données énergétiques dans PostgreSQL/TimescaleDB à partir de deux sources :
|
||||
|
||||
Cette première étape du pipeline Data permet de charger le dataset fourni dans le cadre du projet, contenant les mesures énergétiques de 7 sites sur la période du 1er janvier 2023 au 31 décembre 2024.
|
||||
- le dataset historique CSV/JSON fourni dans le cadre du projet ;
|
||||
- l'API Mock EnerVision.
|
||||
|
||||
Le pipeline assure :
|
||||
|
||||
@@ -12,12 +13,15 @@ Le pipeline assure :
|
||||
- la validation de leur structure et de leur cohérence ;
|
||||
- la normalisation des données nécessaires au stockage ;
|
||||
- le suivi de la qualité des données ;
|
||||
- la traçabilité du dataset importé ;
|
||||
- la traçabilité des données importées ;
|
||||
- le chargement des données dans PostgreSQL/TimescaleDB ;
|
||||
- la conservation des valeurs manquantes et des informations de qualité ;
|
||||
- l'idempotence du chargement afin d'éviter la création de doublons.
|
||||
|
||||
## Données sources
|
||||
|
||||
### Dataset historique
|
||||
|
||||
Le dataset est fourni par le formateur dans le cadre du projet EnerVision.
|
||||
|
||||
Il contient les deux fichiers suivants :
|
||||
@@ -29,7 +33,7 @@ dataset_metadata.json
|
||||
|
||||
Ces fichiers sont nécessaires une seule fois pour initialiser les données historiques de l'environnement.
|
||||
|
||||
Ils ne sont pas versionnés dans Git. Chaque membre de l'équipe récupère manuellement une fois les fichiers fournis par le formateur et les place dans :
|
||||
Ils ne sont pas versionnés dans Git. Chaque membre de l'équipe récupère manuellement les fichiers fournis par le formateur et les place dans :
|
||||
|
||||
```text
|
||||
data/raw/
|
||||
@@ -47,14 +51,26 @@ data/
|
||||
|
||||
Le fichier `.gitkeep` est versionné afin de conserver le répertoire `data/raw/` dans Git. Les fichiers CSV et JSON sont ignorés par Git.
|
||||
|
||||
### API Mock
|
||||
|
||||
La deuxième source est l'API Mock EnerVision.
|
||||
|
||||
Elle permet de récupérer :
|
||||
|
||||
- les informations des sites avec `GET /api/v1/sites` ;
|
||||
- les mesures simulées avec `GET /api/v1/readings`.
|
||||
|
||||
L'API Mock est utilisée pour compléter les données historiques avec des mesures simulées récupérées sur une période donnée.
|
||||
|
||||
## Technologies utilisées
|
||||
|
||||
| Technologie | Utilisation |
|
||||
|---|---|
|
||||
| Python | Développement du pipeline ETL |
|
||||
| Pandas | Lecture, validation et transformation des données |
|
||||
| JSON | Lecture des métadonnées du dataset |
|
||||
| hashlib / SHA-256 | Identification, intégrité et traçabilité du dataset |
|
||||
| Pandas | Lecture, validation et transformation du dataset historique |
|
||||
| JSON | Lecture des métadonnées et conservation des données sources |
|
||||
| HTTPX | Appels HTTP asynchrones vers l'API Mock |
|
||||
| hashlib / SHA-256 | Identification, intégrité et traçabilité du dataset historique |
|
||||
| SQLAlchemy Async | Connexion et chargement asynchrone en base |
|
||||
| PostgreSQL | Stockage relationnel |
|
||||
| TimescaleDB | Stockage des séries temporelles énergétiques |
|
||||
@@ -62,11 +78,14 @@ Le fichier `.gitkeep` est versionné afin de conserver le répertoire `data/raw/
|
||||
| Alembic | Gestion des migrations du schéma |
|
||||
| uv | Gestion et exécution de l'environnement Python |
|
||||
| Ruff | Contrôle de la qualité du code |
|
||||
| mypy | Vérification du typage |
|
||||
| Pytest | Tests automatisés |
|
||||
|
||||
## Fonctionnement du pipeline
|
||||
# Import du dataset historique
|
||||
|
||||
Le script principal d'import se trouve dans :
|
||||
## Fonctionnement du pipeline historique
|
||||
|
||||
Le script d'import se trouve dans :
|
||||
|
||||
```text
|
||||
apps/backend/app/etl/historical_import.py
|
||||
@@ -207,7 +226,7 @@ Valeurs manquantes identifiées :
|
||||
| `humidity_percent` | 3 423 |
|
||||
| `solar_irradiance_wm2` | 3 964 |
|
||||
|
||||
## Exécution en dry-run
|
||||
## Exécution historique en dry-run
|
||||
|
||||
Depuis le dossier :
|
||||
|
||||
@@ -227,7 +246,7 @@ uv run python -m app.etl.historical_import `
|
||||
|
||||
Aucune donnée n'est écrite dans la base pendant cette exécution.
|
||||
|
||||
## Chargement réel
|
||||
## Chargement historique réel
|
||||
|
||||
Depuis `apps/backend/` :
|
||||
|
||||
@@ -249,7 +268,7 @@ Chargement : 2000/122647
|
||||
Chargement : 122647/122647
|
||||
```
|
||||
|
||||
## Résultats obtenus
|
||||
## Résultats obtenus pour le dataset historique
|
||||
|
||||
Après le chargement initial, les contrôles en base ont confirmé :
|
||||
|
||||
@@ -266,7 +285,7 @@ Le premier import a créé :
|
||||
nouvelles lectures : 122647
|
||||
```
|
||||
|
||||
## Idempotence
|
||||
## Idempotence du dataset historique
|
||||
|
||||
Le pipeline a été exécuté une deuxième fois avec exactement le même dataset afin de vérifier son idempotence.
|
||||
|
||||
@@ -280,7 +299,7 @@ nouvelles lectures : 0
|
||||
|
||||
Une nouvelle exécution du même import ne crée donc pas de mesures supplémentaires pour le dataset testé.
|
||||
|
||||
## Vérifications SQL
|
||||
## Vérifications SQL du dataset historique
|
||||
|
||||
Depuis la racine du projet, vérifier le nombre d'enregistrements avec :
|
||||
|
||||
@@ -302,21 +321,213 @@ Vérifier la source des mesures avec :
|
||||
docker compose exec db psql -U enervision -d enervision -c "SELECT source, COUNT(*) FROM reading GROUP BY source ORDER BY source;"
|
||||
```
|
||||
|
||||
Résultat attendu :
|
||||
Résultat attendu pour le dataset historique :
|
||||
|
||||
```text
|
||||
csv | 122647
|
||||
```
|
||||
|
||||
## Tests et qualité
|
||||
# Import depuis l'API Mock
|
||||
|
||||
Les tests automatisés du pipeline sont situés dans :
|
||||
## Fonctionnement
|
||||
|
||||
Le script d'import de l'API Mock se trouve dans :
|
||||
|
||||
```text
|
||||
apps/backend/app/etl/mock_api_import.py
|
||||
```
|
||||
|
||||
Le flux est le suivant :
|
||||
|
||||
```text
|
||||
API Mock
|
||||
|
|
||||
+-----+------+
|
||||
| |
|
||||
v v
|
||||
/sites /readings
|
||||
| |
|
||||
+-----+------+
|
||||
|
|
||||
v
|
||||
mock_api_import.py
|
||||
|
|
||||
v
|
||||
Transformation
|
||||
+ qualité data
|
||||
|
|
||||
v
|
||||
PostgreSQL / TimescaleDB
|
||||
| |
|
||||
v v
|
||||
site reading
|
||||
```
|
||||
|
||||
Le pipeline commence par récupérer les sites avec :
|
||||
|
||||
```text
|
||||
GET /api/v1/sites
|
||||
```
|
||||
|
||||
Il récupère ensuite les mesures de chaque site avec :
|
||||
|
||||
```text
|
||||
GET /api/v1/readings
|
||||
```
|
||||
|
||||
Les paramètres envoyés à `/api/v1/readings` sont :
|
||||
|
||||
```text
|
||||
site_id
|
||||
start_time
|
||||
end_time
|
||||
limit
|
||||
```
|
||||
|
||||
Le paramètre `limit` doit être compris entre 1 et 1000.
|
||||
|
||||
## Configuration de l'API Mock
|
||||
|
||||
La connexion à l'API Mock est configurée avec les variables d'environnement suivantes :
|
||||
|
||||
```text
|
||||
APP_MOCK_API_BASE_URL
|
||||
APP_MOCK_API_USERNAME
|
||||
APP_MOCK_API_PASSWORD
|
||||
APP_MOCK_API_TIMEOUT_SECONDS
|
||||
```
|
||||
|
||||
Les identifiants réels ne sont pas versionnés dans Git.
|
||||
|
||||
Les fichiers `.env.example` indiquent uniquement les variables nécessaires à l'exécution.
|
||||
|
||||
## Transformation des mesures API
|
||||
|
||||
Les mesures provenant de l'API Mock sont enregistrées dans `reading` avec :
|
||||
|
||||
```text
|
||||
source = "api_history"
|
||||
dataset_id = NULL
|
||||
```
|
||||
|
||||
Les mesures provenant de l'API ne sont donc pas rattachées à un dataset historique.
|
||||
|
||||
Le timestamp reçu depuis l'API est converti en `datetime` avec timezone avant le chargement.
|
||||
|
||||
La réponse source est conservée dans :
|
||||
|
||||
```text
|
||||
raw_data
|
||||
```
|
||||
|
||||
afin de préserver la donnée reçue et faciliter la traçabilité.
|
||||
|
||||
## Qualité des données API
|
||||
|
||||
Les valeurs `NULL` fournies par l'API sont conservées telles quelles.
|
||||
|
||||
Une valeur manquante n'est pas transformée en zéro et la mesure n'est pas supprimée.
|
||||
|
||||
Le pipeline conserve également :
|
||||
|
||||
```text
|
||||
data_quality
|
||||
null_reasons
|
||||
```
|
||||
|
||||
Les niveaux de qualité possibles sont :
|
||||
|
||||
```text
|
||||
good
|
||||
partial
|
||||
degraded
|
||||
critical
|
||||
```
|
||||
|
||||
Aucune imputation n'est réalisée pendant l'ingestion :
|
||||
|
||||
```text
|
||||
imputed_values = NULL
|
||||
imputation_method = NULL
|
||||
```
|
||||
|
||||
Cette stratégie permet de distinguer une véritable valeur nulle ou manquante d'une consommation égale à zéro et de conserver les informations liées aux défaillances de capteurs.
|
||||
|
||||
## Dry-run de l'API Mock
|
||||
|
||||
Le mode `--dry-run` permet de tester la connexion, la récupération des sites et la récupération des mesures sans écrire dans PostgreSQL.
|
||||
|
||||
Depuis `apps/backend/` :
|
||||
|
||||
```powershell
|
||||
uv run python -m app.etl.mock_api_import `
|
||||
--start-time "2024-06-15T12:00:00" `
|
||||
--end-time "2024-06-15T13:00:00" `
|
||||
--limit 60 `
|
||||
--dry-run
|
||||
```
|
||||
|
||||
## Chargement réel depuis l'API Mock
|
||||
|
||||
Depuis `apps/backend/` :
|
||||
|
||||
```powershell
|
||||
uv run python -m app.etl.mock_api_import `
|
||||
--start-time "2024-06-15T12:00:00" `
|
||||
--end-time "2024-06-15T13:00:00" `
|
||||
--limit 60
|
||||
```
|
||||
|
||||
## Résultat validé pour l'API Mock
|
||||
|
||||
Le scénario de validation utilisé couvre la période :
|
||||
|
||||
```text
|
||||
15/06/2024 12:00 UTC
|
||||
à
|
||||
15/06/2024 13:00 UTC
|
||||
```
|
||||
|
||||
avec une limite de 60 lectures par site.
|
||||
|
||||
Résultat obtenu :
|
||||
|
||||
```text
|
||||
sites récupérés : 7
|
||||
lectures par site : 60
|
||||
lectures récupérées : 420
|
||||
source : api_history
|
||||
dataset_id : NULL
|
||||
```
|
||||
|
||||
Les contrôles effectués directement dans PostgreSQL/TimescaleDB ont confirmé :
|
||||
|
||||
- l'enregistrement des mesures dans `reading` ;
|
||||
- la présence des 7 sites ;
|
||||
- `source = "api_history"` ;
|
||||
- `dataset_id = NULL` ;
|
||||
- la conservation des valeurs `NULL` ;
|
||||
- la conservation de `data_quality` ;
|
||||
- la conservation de `null_reasons` ;
|
||||
- la conservation de la donnée source dans `raw_data`.
|
||||
|
||||
## Idempotence de l'import API Mock
|
||||
|
||||
Le même import a été exécuté plusieurs fois afin de vérifier qu'une mesure déjà présente n'est pas créée une seconde fois.
|
||||
|
||||
L'idempotence repose sur la contrainte d'unicité de la table `reading` et sur la gestion des conflits lors de l'insertion.
|
||||
|
||||
Un test d'intégration automatisé vérifie également ce comportement.
|
||||
|
||||
# Tests et qualité
|
||||
|
||||
Les tests automatisés des pipelines ETL sont situés dans :
|
||||
|
||||
```text
|
||||
apps/backend/tests/etl/
|
||||
```
|
||||
|
||||
Ils couvrent notamment :
|
||||
Les tests de l'import historique couvrent notamment :
|
||||
|
||||
- la validation du dataset ;
|
||||
- les colonnes obligatoires ;
|
||||
@@ -328,22 +539,96 @@ Ils couvrent notamment :
|
||||
- la construction des mesures destinées à la BDD ;
|
||||
- le respect des contraintes du modèle de données.
|
||||
|
||||
Les tests de l'import API Mock couvrent notamment :
|
||||
|
||||
- la récupération des sites ;
|
||||
- l'appel à `/api/v1/readings` ;
|
||||
- les paramètres `site_id`, `start_time`, `end_time` et `limit` ;
|
||||
- la gestion des erreurs HTTP ;
|
||||
- la validation du format de la réponse ;
|
||||
- la transformation des mesures ;
|
||||
- la conservation des valeurs `NULL` ;
|
||||
- la conservation de `data_quality` et `null_reasons` ;
|
||||
- `source = "api_history"` ;
|
||||
- `dataset_id = NULL` ;
|
||||
- la conservation de `raw_data` ;
|
||||
- l'idempotence du chargement.
|
||||
|
||||
Exécuter les tests ETL :
|
||||
|
||||
```powershell
|
||||
uv run pytest tests\etl -v
|
||||
```
|
||||
|
||||
Exécuter les tests unitaires de l'import API Mock :
|
||||
|
||||
```powershell
|
||||
uv run pytest tests\etl\test_mock_api_import.py -v
|
||||
```
|
||||
|
||||
Exécuter le test d'intégration de l'import API Mock :
|
||||
|
||||
```powershell
|
||||
uv run pytest tests\etl\test_mock_api_import.py -m integration -v
|
||||
```
|
||||
|
||||
Contrôler la qualité du code :
|
||||
|
||||
```powershell
|
||||
uv run ruff check app\etl tests\etl
|
||||
```
|
||||
|
||||
## Suite du pipeline Data
|
||||
Contrôler le typage :
|
||||
|
||||
L'import historique constitue la première brique du pipeline Data EnerVision.
|
||||
```powershell
|
||||
uv run mypy app
|
||||
```
|
||||
|
||||
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.
|
||||
Exécuter la suite complète avec le seuil de couverture :
|
||||
|
||||
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.
|
||||
```powershell
|
||||
uv run pytest --cov-fail-under=85
|
||||
```
|
||||
|
||||
Lors de la validation de l'import API Mock :
|
||||
|
||||
```text
|
||||
8 tests unitaires passés
|
||||
1 test d'intégration passé
|
||||
```
|
||||
|
||||
La suite backend complète a également été validée avec une couverture supérieure au seuil de 85 %.
|
||||
|
||||
# Suite du pipeline Data
|
||||
|
||||
Deux sources de données sont maintenant prises en charge :
|
||||
|
||||
```text
|
||||
Dataset CSV/JSON
|
||||
|
|
||||
v
|
||||
historical_import.py
|
||||
|
|
||||
+-----------------+
|
||||
|
|
||||
v
|
||||
PostgreSQL / TimescaleDB
|
||||
^
|
||||
|
|
||||
+-----------------+
|
||||
|
|
||||
mock_api_import.py
|
||||
^
|
||||
|
|
||||
API Mock
|
||||
```
|
||||
|
||||
La logique d'extraction, de transformation et de chargement est donc disponible pour les deux sources de données du MVP.
|
||||
|
||||
La prochaine étape consiste à orchestrer ces traitements avec Apache Airflow.
|
||||
|
||||
Airflow permettra de planifier les traitements, gérer leur ordre d'exécution, suivre leur état et remonter les erreurs.
|
||||
|
||||
Airflow ne remplacera pas la logique ETL Python existante. Les scripts actuels resteront responsables de l'extraction, de la validation, de la transformation et du chargement.
|
||||
|
||||
Le pipeline Data servira ensuite à préparer les données nécessaires au modèle de Machine Learning.
|
||||
Reference in New Issue
Block a user