Tester du code Airflow

Créer des pipelines de données avec Airflow

Volker Janz

Senior Developer Advocate at Astronomer

Trois niveaux de tests

Pyramide de tests Dag

  • Intégrité : le Dag se charge-t-il sans erreur ?
  • Unitaire : la logique d'affaires produit-elle les bons résultats ?
  • Intégration : le Dag complet s'exécute-t-il de bout en bout ?
Créer des pipelines de données avec Airflow

Tests d'intégrité avec DagBag

from airflow.models import DagBag

dag_bag = DagBag(include_examples=False)
def test_no_import_errors(): assert len(dag_bag.import_errors) == 0
def test_dag_loaded(): assert "daily_etl" in dag_bag.dags
Créer des pipelines de données avec Airflow

Pourquoi les Dags échouent à l'import

ModuleNotFoundError

  • ModuleNotFoundError : fournisseur manquant ou mauvais chemin d'import
  • NameError : variable ou fonction renommée, ou coquille
  • ImportError : importations circulaires entre fichiers de Dag
  • Les tests d'intégrité aident à éviter ces problèmes
Créer des pipelines de données avec Airflow

Tests unitaires des fonctions de tâches

Dans le fichier du Dag :

def clean_record(record):
    return {
        "name": record["name"].strip(),
        "email": record["email"].lower(),
    }

@task
def transform(records):
    return [clean_record(r)
            for r in records]
  • Extraire la logique d'affaires

Dans le fichier de test :

from dags.data_cleaning import (
    clean_record,
)

def test_strips_whitespace():
    result = clean_record(
      {"name": "  Alice  ",
       "email": "[email protected]"}
    )
    assert result["name"] == "Alice"
  • Cibler les tests unitaires sur la logique d'affaires
Créer des pipelines de données avec Airflow

Tests d'intégration avec dag.test()

import pytest
from airflow.models import DagBag
from pendulum import datetime

dag_bag = DagBag(include_examples=False)

def test_etl_pipeline(): dag = dag_bag.get_dag("etl_output") assert dag is not None dag.test(logical_date=datetime(2026, 1, 15)) output = Path("/tmp/etl_results.json") assert output.exists() results = json.loads(output.read_text()) assert len(results) == 2
  • Exécuter le Dag réel avec une entrée contrôlée et valider la sortie
Créer des pipelines de données avec Airflow

Tests en CI

Pipeline CI

  • Intégrité + unitaire : à chaque commit (rapide)
  • Intégration : lors des demandes de tirage ou la nuit (plus lent)
  • Aucun Dag ne passe en production sans réussir les trois
Créer des pipelines de données avec Airflow

Passons à la pratique !

Créer des pipelines de données avec Airflow

Preparing Video For Download...