Мониторинг, оповещения и обратные вызовы

Введение в Apache Airflow на Python

Mike Metzger

Data Engineer

Жизненный цикл DAG

 

  • Каждый запуск DAG и задача проходят через последовательность состояний
    • В очереди
    • Выполняется
    • Финальное состояние: успешно, ошибка, пропущено
  • Airflow отслеживает состояния и переходы

Диаграмма жизненного цикла запуска DAG: от очереди к выполнению и финальному состоянию

Введение в Apache Airflow на Python

Обратные вызовы

  • Функции для запуска DAG или задачи, которые Airflow вызывает автоматически
  • Обратные вызовы для конкретных переходов состояний
    • on_failure_callback — запуск DAG / задача завершились ошибкой
    • on_success_callback — запуск DAG / задача выполнены успешно
  • Обратные вызовы для отдельных задач
    • on_retry_callback
    • on_skipped_callback
    • on_execute_callback

Иллюстрация обратных вызовов Airflow, срабатывающих при переходах состояний

Введение в Apache Airflow на Python

Контекст обратного вызова

  • Airflow автоматически передаёт словарь context
  • context содержит информацию о запуске DAG / задаче
    • context["dag"].dag_id — имя DAG
    • context["task_instance"].task_id — имя задачи
    • context["logical_date"] — дата запуска DAG
  • Функция обратного вызова должна принимать 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.")
Введение в Apache Airflow на Python

Пример обратного вызова

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.
Введение в Apache Airflow на Python

Уведомители

  • Функции, которые можно привязать к обратным вызовам
  • Отправляют оповещения во внешние системы
  • Доступно несколько уведомителей:
    • SmtpNotifier — отправляет оповещения по электронной почте
    • SlackNotifier — публикует сообщения в канал Slack
    • Уведомители для приложений — PagerDuty, OpsGenie

Иллюстрация уведомителей Airflow, отправляющих оповещения во внешние системы

Введение в Apache Airflow на Python

SmtpNotifier

  • Отправляет письмо при срабатывании обратного вызова
  • Находится в библиотеке airflow.providers.smtp.notifications.smtp
  • Требует атрибуты from_email и to
  • Можно указать subject, html_content и др.
    @dag(dag_id=`sales_etl_dag`,
       on_failure_callback=SmtpNotifier(
           to='[email protected]',
           from_email='[email protected]',
           subject='Dag sales_etl_dag has failed'
       )
    )
    
Введение в Apache Airflow на Python

Журнал аудита

  • Используйте журнал аудита для мониторинга Airflow
  • Содержит полную последовательность событий в экземпляре Airflow

Страница журнала аудита Airflow со списком системных событий по типам

Введение в Apache Airflow на Python

Давайте потренируемся!

Введение в Apache Airflow на Python

Preparing Video For Download...