Déclencheurs et tolérance aux pannes

Introduction à Apache Airflow en Python

Mike Metzger

Data Engineer

Créer des pipelines robustes

  • Déclencheurs – contrôlent quand les tâches peuvent s'exécuter
    • Définissent des tâches exécutées selon des conditions précises
    • Déclencheur par défaut – tout réussi

Illustration des déclencheurs contrôlant quand les tâches peuvent s'exécuter

  • Tolérance aux pannes – permet la reprise au sein d'une tâche
    • Relancer une tâche individuellement

Illustration de la tolérance aux pannes relançant une tâche en échec

Introduction à Apache Airflow en Python

Règles de déclenchement

  • Vérifie l'état des tâches en amont avant qu'une tâche commence
  • Vous permet de décider si les tâches doivent continuer
  • Par défaut, toutes les tâches précédentes doivent réussir avant de poursuivre
  • Les règles de déclenchement permettent de changer cela

Illustration des règles de déclenchement vérifiant l'état des tâches en amont

Introduction à Apache Airflow en Python

Usages des règles de déclenchement

  • Tâches d'avis
  • Tâches de nettoyage
  • Exécution conditionnelle

Illustration des cas d'usage des règles de déclenchement, comme les tâches d'avis et de nettoyage

Introduction à Apache Airflow en Python

Règles de déclenchement clés

  • all_success – Toutes les tâches précédentes ont réussi, aucune omission
  • all_failed – Toutes les tâches précédentes ont échoué
  • all_done – Toutes les tâches précédentes sont terminées, peu importe le résultat
  • one_failed – Au moins une tâche a échoué
  • one_success – Au moins une tâche a réussi
  • none_failed – Toutes les tâches précédentes ont réussi ou ont été omises
Introduction à Apache Airflow en Python

Mise en œuvre des règles de déclenchement

  • from airflow.utils.trigger_rule import TriggerRule
  • Les règles de déclenchement s'appliquent au décorateur @task
  • trigger_rule=TriggerRule.<TriggerRule enum>
@task(trigger_rule=TriggerRule.ALL_SUCCESS)
def run_if_everything_succeeds:
  print('All previous tasks succeeded!')
Introduction à Apache Airflow en Python

Attributs de tolérance aux pannes

  • retries – paramètre de tâche qui définit combien de fois Airflow doit relancer une tâche en échec avant de la marquer échouée
  • retry_delay – paramètre de tâche acceptant un timedelta définissant l'attente entre les relances
@task(
  trigger_rule=TriggerRule.ONE_SUCCESS,
  retries=3,
  retry_delay=timedelta(minutes=5)
)
Introduction à Apache Airflow en Python

TriggerDagRunOperator

  • Permet à un Dag d'en déclencher un autre
  • Fournit du code et des Dag
  • from airflow.providers.standard.operators.trigger_dagrun import TriggerDagRunOperator
  • trigger_dag_id – Doit correspondre à dag_id
  • wait_for_completion – Bloque les tâches suivantes jusqu'à la fin
  • poke_interval – Fréquence de vérification
  • conf – Données à transmettre au Dag enfant
Introduction à Apache Airflow en Python

Exemple TriggerDagRun

trigger_child = TriggerDagRunOperator(
        task_id="trigger_child_pipeline",
        trigger_dag_id="child_pipeline_dag",   
        wait_for_completion=True,              
        poke_interval=30,                      
        conf={                                 
            "source": "s3://my-bucket/raw/",
        },
    )

task1 >> trigger_child    # Exécuter la première tâche, puis déclencher le Dag enfant
trigger_child >> cleanup  # Exécuter une tâche de nettoyage après la fin du Dag enfant
Introduction à Apache Airflow en Python

Passons à la pratique !

Introduction à Apache Airflow en Python

Preparing Video For Download...