Monitorizare, alerte și callback-uri

Introducere în Apache Airflow în Python

Mike Metzger

Data Engineer

Ciclul de viață al DAG-ului

 

  • Fiecare rulare de DAG și sarcină parcurge o secvență de stări
    • În așteptare
    • În execuție
    • Stare finală: Succes, Eșec, Omis
  • Airflow urmărește stările și tranzițiile

Diagramă a ciclului de viață al unui DAG, de la așteptare la execuție și stare finală

Introducere în Apache Airflow în Python

Callback-uri

  • Funcții apelate automat de Airflow la rularea DAG-ului sau a sarcinii
  • Callback-uri pentru tranziții de stare specifice
    • on_failure_callback - Rulare DAG / sarcină eșuată
    • on_success_callback - Rulare DAG / sarcină reușită
  • Callback-uri specifice sarcinilor
    • on_retry_callback
    • on_skipped_callback
    • on_execute_callback

Ilustrație a callback-urilor Airflow declanșate la tranziții de stare

Introducere în Apache Airflow în Python

Contextul callback-ului

  • Airflow transmite automat un dicționar context
  • context conține informații despre rularea DAG-ului / sarcinii
    • context["dag"].dag_id - Numele DAG-ului
    • context["task_instance"].task_id - Numele sarcinii
    • context["logical_date"] - Data rulării DAG-ului
  • Funcția callback trebuie să accepte 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.")
Introducere în Apache Airflow în Python

Exemplu de callback

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.
Introducere în Apache Airflow în Python

Notificatori

  • Funcții asociate callback-urilor
  • Trimit alerte către sisteme externe
  • Mai mulți notificatori disponibili:
    • SmtpNotifier - Trimite alerte prin email
    • SlackNotifier - Postează mesaje pe un canal Slack
    • Notificatori specifici - PagerDuty, OpsGenie

Ilustrație a notificatorilor Airflow care trimit alerte către sisteme externe

Introducere în Apache Airflow în Python

SmtpNotifier

  • Trimite un email la callback
  • În biblioteca airflow.providers.smtp.notifications.smtp
  • Necesită atributele from_email și to
  • Poate include 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'
       )
    )
    
Introducere în Apache Airflow în Python

Audit Log

  • Utilizați Audit Log pentru a monitoriza Airflow
  • Conține secvența tuturor evenimentelor din instanța Airflow

Pagina de audit log din Airflow cu lista evenimentelor de sistem pe tip

Introducere în Apache Airflow în Python

Să exersăm!

Introducere în Apache Airflow în Python

Preparing Video For Download...