DAG robustes : gestion des échecs et reprises

Créer des pipelines de données avec Airflow

Volker Janz

Senior Developer Advocate at Astronomer

Pourquoi les tâches échouent

 

Problèmes de connexion des tâches

 

  • Expirations d'API et limites de débit
  • Brèves coupures réseau et retards de résolution de noms de domaine
  • Plusieurs tâches en concurrence pour la même ressource
Créer des pipelines de données avec Airflow

Réessais

@task(
    retries=3,
    retry_delay=timedelta(minutes=2),
)
def fetch_weather():
    response = requests.get("https://api.weather.com/forecast")
    response.raise_for_status()
    return response.json()

 

  • Réessaie la tâche jusqu'à 3 fois en cas d'échec
  • Attend 2 minutes entre chaque tentative
Créer des pipelines de données avec Airflow

Recul exponentiel

 

@task(
    retries=3,
    retry_delay=timedelta(minutes=2),
    retry_exponential_backoff=2.0,
)
def fetch_weather():
    ...

Recul exponentiel

 

  • Le délai double après chaque réessai
  • Laisse le temps aux services externes de se rétablir
Créer des pipelines de données avec Airflow

on_failure_callback

def alert_on_failure(context):
    dag_id = context["dag"].dag_id
    task_id = context["ti"].task_id
    print(f"ALERTE : {task_id} dans {dag_id} a échoué !")

@task(on_failure_callback=alert_on_failure)
def fetch_weather():
    ...

 

  • Dictionnaire context : dag, ti, exception, URL du journal
  • Acheminer vers Slack, PagerDuty ou tout outil d'alerte
Créer des pipelines de données avec Airflow

Rappels : niveau DAG vs tâche

Niveau DAG (fourre-tout)

@dag(
    on_failure_callback=alert_team,
)
def my_pipeline():
    ...

Niveau tâche (spécifique)

@task(
    on_failure_callback=page_oncall,
)
def critical_step():
    ...

 

  • Le niveau tâche a priorité sur le niveau DAG
  • Utilisez les deux pour des alertes en couches
Créer des pipelines de données avec Airflow

max_consecutive_failed_dag_runs

@dag(max_consecutive_failed_dag_runs=3)
def monitoring_pipeline():
    ...

Exécutions de DAG échouées suivies d'une pause auto

  • Évite la fatigue d'alerte lors d'échecs répétés
  • Le DAG reprend après réactivation manuelle
Créer des pipelines de données avec Airflow

Tout rassembler

@dag(
    max_consecutive_failed_dag_runs=3,
    on_failure_callback=alert_team,
)
def weather_pipeline():

    @task(
        retries=3,
        retry_delay=timedelta(minutes=2),
        retry_exponential_backoff=True,
        on_failure_callback=page_oncall,
    )
    def fetch_weather():
        ...
  • Réessais + recul gèrent les pannes passagères
  • Les rappels préviennent les bonnes personnes
  • La pause auto réduit le bruit
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...