Surveillance, alertes et fonctions de rappel

Introduction à Apache Airflow en Python

Mike Metzger

Data Engineer

Cycle de vie d'un Dag

 

  • Chaque exécution de Dag et chaque tâche passe par une suite d'états
    • En file d'attente
    • En cours d'exécution
    • État final : Réussite, Échec, Ignoré
  • Airflow suit les états et leurs transitions

Schéma du cycle de vie d'une exécution de Dag : en file d'attente, en cours, puis état final

Introduction à Apache Airflow en Python

Rappels (callbacks)

  • Fonctions sur une exécution de Dag ou une tâche qu'Airflow appelle automatiquement
  • Rappels pour des transitions d'état précises
    • on_failure_callback - Échec de l'exécution / tâche
    • on_success_callback - Réussite de l'exécution / tâche
  • Rappels propres aux tâches
    • on_retry_callback
    • on_skipped_callback
    • on_execute_callback

Illustration des rappels Airflow déclenchés lors de transitions d'état

Introduction à Apache Airflow en Python

Contexte du rappel

  • Airflow transmet automatiquement un dictionnaire context
  • context contient des infos sur l'exécution / la tâche
    • context["dag"].dag_id - Nom du Dag
    • context["task_instance"].task_id - Nom de la tâche
    • context["logical_date"] - Date de l'exécution du Dag
  • La fonction de rappel doit accepter context
def alert_on_failure(context):

dag_id = context["dag"].dag_id task_id = context["task_instance"].task_id print(f"Task {task_id} in DAG {dag_id} has failed.")
Introduction à Apache Airflow en Python

Exemple de rappel

def alert_on_failure(context):
    dag_id = context["dag"].dag_id
    task_id = context["task_instance"].task_id
    print(f"Task {task_id} in Dag {dag_id} has failed.")

@dag(on_failure_callback=alert_on_failure) def sales_etl_dag(): @task() def data_import_task(): raise ValueError("Simulated failure")
# Task data_import_task in Dag sales_etl_dag has failed.
Introduction à Apache Airflow en Python

Notifiers

  • Fonctions pouvant être liées aux rappels
  • Envoient des alertes à des systèmes externes
  • Plusieurs « notifiers » disponibles :
    • SmtpNotifier - Envoie des alertes par courriel
    • SlackNotifier - Publie des messages dans un canal Slack
    • Spécifiques à des applis : PagerDuty, OpsGenie

Illustration des « notifiers » Airflow envoyant des alertes vers des systèmes externes

Introduction à Apache Airflow en Python

SmtpNotifier

  • Envoie un courriel lors d'un rappel
  • Dans la bibliothèque airflow.providers.smtp.notifications.smtp
  • Exige les attributs from_email et to
  • Peut inclure subject, html_content, etc.
    @dag(dag_id=`sales_etl_dag`,
       on_failure_callback=SmtpNotifier(
           to='[email protected]',
           from_email='[email protected]',
           subject='Dag sales_etl_dag has failed'
       )
    )
    
Introduction à Apache Airflow en Python

Journal d'audit

  • Utilisez le journal d'audit pour surveiller Airflow
  • Contient la suite de tous les événements de l'instance Airflow

Page du journal d'audit d'Airflow listant les événements système par type

Introduction à Apache Airflow en Python

Passons à la pratique !

Introduction à Apache Airflow en Python

Preparing Video For Download...