Monitorowanie, alerty i wywołania zwrotne

Wprowadzenie do Apache Airflow w Pythonie

Mike Metzger

Data Engineer

Cykl życia DAG-a

 

  • Każdy przebieg DAG-a i zadanie przechodzi przez sekwencję stanów
    • W kolejce
    • Uruchomione
    • Stan końcowy: sukces, błąd, pominięcie
  • Airflow śledzi stany i przejścia

Diagram cyklu życia przebiegu DAG-a: od kolejki przez uruchomienie do stanu końcowego

Wprowadzenie do Apache Airflow w Pythonie

Wywołania zwrotne

  • Funkcje wywoływane automatycznie przez Airflow przy przebiegu DAG-a lub zadaniu
  • Wywołania zwrotne dla określonych przejść stanów
    • on_failure_callback – przebieg DAG-a / zadanie kończy się błędem
    • on_success_callback – przebieg DAG-a / zadanie kończy się sukcesem
  • Wywołania zwrotne specyficzne dla zadań
    • on_retry_callback
    • on_skipped_callback
    • on_execute_callback

Ilustracja wywołań zwrotnych Airflow wyzwalanych przy przejściach stanów

Wprowadzenie do Apache Airflow w Pythonie

Kontekst wywołania zwrotnego

  • Airflow automatycznie przekazuje słownik context
  • context zawiera informacje o przebiegu DAG-a / zadaniu
    • context["dag"].dag_id – nazwa DAG-a
    • context["task_instance"].task_id – nazwa zadania
    • context["logical_date"] – data przebiegu DAG-a
  • Funkcja wywołania zwrotnego musi przyjmować 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.")
Wprowadzenie do Apache Airflow w Pythonie

Przykład wywołania zwrotnego

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.
Wprowadzenie do Apache Airflow w Pythonie

Notyfikatory

  • Funkcje powiązane z wywołaniami zwrotnymi
  • Wysyłają alerty do systemów zewnętrznych
  • Dostępne notyfikatory:
    • SmtpNotifier – wysyła alerty e-mail
    • SlackNotifier – publikuje wiadomości na kanale Slack
    • Notyfikatory dedykowane – PagerDuty, OpsGenie

Ilustracja notyfikatorów Airflow wysyłających alerty do systemów zewnętrznych

Wprowadzenie do Apache Airflow w Pythonie

SmtpNotifier

  • Wysyła e-mail przy wywołaniu zwrotnym
  • W bibliotece airflow.providers.smtp.notifications.smtp
  • Wymaga atrybutów from_email i to
  • Obsługuje subject, html_content i inne
    @dag(dag_id=`sales_etl_dag`,
       on_failure_callback=SmtpNotifier(
           to='[email protected]',
           from_email='[email protected]',
           subject='Dag sales_etl_dag has failed'
       )
    )
    
Wprowadzenie do Apache Airflow w Pythonie

Dziennik audytu

  • Użyj dziennika audytu do monitorowania Airflow
  • Zawiera sekwencję wszystkich zdarzeń w instancji Airflow

Strona dziennika audytu Airflow z listą zdarzeń systemowych według typu

Wprowadzenie do Apache Airflow w Pythonie

Czas na ćwiczenia!

Wprowadzenie do Apache Airflow w Pythonie

Preparing Video For Download...