Le DAG Versioning dans Airflow 3.0
Airflow 3.0 introduit le DAG Versioning : chaque exécution d'un DAG est associée à la version du code telle qu'elle était au moment de la planification. Fini les surprises où un DAG en cours d'exécution est affecté par un déploiement de code en production.
Comment ça fonctionne
Airflow capture automatiquement un snapshot du DAG au moment du scheduling. Si vous déployez une nouvelle version du DAG pendant qu'une exécution est en cours, les tâches restantes continuent avec l'ancienne version. Les nouvelles exécutions utilisent la version mise à jour.
from airflow.sdk import DAG, task
from datetime import datetime
# Version 1 du DAG (déployée lundi)
with DAG(
dag_id='etl_pipeline',
schedule='@daily',
start_date=datetime(2025, 4, 1),
) as dag:
@task
def extraire():
"""Extraction depuis la source A."""
return {'source': 'A', 'lignes': 1000}
@task
def transformer(data: dict):
"""Transformation v1 : nettoyage simple."""
data['lignes_valides'] = int(data['lignes'] * 0.95)
return data
@task
def charger(data: dict):
print(f"Chargement de {data['lignes_valides']} lignes")
charger(transformer(extraire()))
# Si on déploie une v2 mardi qui modifie transformer(),
# l'exécution de lundi (en cours) garde la v1.
# L'exécution de mardi utilisera la v2.
Consulter les versions
L'interface web d'Airflow 3.0 affiche la version associée à chaque DAG Run. On peut aussi la consulter via l'API REST ou la CLI.
# Lister les versions d'un DAG
airflow dags list-versions --dag-id etl_pipeline
# Voir quelle version est associée à un DAG Run
airflow dags list-runs --dag-id etl_pipeline
# Affiche : run_id | dag_version | state | start_date
