Files
ENI-projet-piscine/etl/airflow/dags/mock_api_import.py
T
Johan LEROY cbbfaf4910 fix(etl): importer une seule mesure par heure depuis l'API Mock
L'API Mock ne renvoie pas les mesures d'une période : elle génère `limit` points
répartis sur l'intervalle demandé (1 000 par heure avec --limit 1000, un toutes
les 3,6 s). Le DAG aurait écrit 7 000 lignes par heure et par environnement,
alors que le dataset historique a une mesure horaire et que les features ML
décalent par ligne : `shift(168)`, le retard d'une semaine, serait devenu un
retard de dix minutes, sans erreur visible au scoring ni au réentraînement.

Avec --limit 1, l'API renvoie la mesure de :00 de chaque heure, au pas du CSV.
Constaté sur la recette le 23/09 avant la réactivation des DAGs.
2026-09-23 15:31:51 +02:00

62 lines
2.5 KiB
Python

"""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"
# Contrainte : l'API Mock génère `limit` points répartis sur l'intervalle. Un seul donne la mesure
# de :00, au pas horaire du CSV que les features ML supposent (`shift(168)` compte des lignes).
LIMITE_LECTURES = 1
# 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,
)