From 308769b3255122e7a04925f4c41e77e79ca6975e Mon Sep 17 00:00:00 2001 From: Meryemel-gham Date: Mon, 21 Sep 2026 14:40:21 +0200 Subject: [PATCH] feat(etl): orchestre l'import historique avec Airflow --- docker-compose.yml | 1 + etl/airflow/dags/historical_import.py | 43 +++++++++++++++++++++++++++ 2 files changed, 44 insertions(+) create mode 100644 etl/airflow/dags/historical_import.py diff --git a/docker-compose.yml b/docker-compose.yml index da19f43..6f2c1c4 100644 --- a/docker-compose.yml +++ b/docker-compose.yml @@ -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 diff --git a/etl/airflow/dags/historical_import.py b/etl/airflow/dags/historical_import.py new file mode 100644 index 0000000..5fb119e --- /dev/null +++ b/etl/airflow/dags/historical_import.py @@ -0,0 +1,43 @@ +"""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" + +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), + )