Триггеры и отказоустойчивость

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

Mike Metzger

Data Engineer

Построение надёжных пайплайнов

  • Триггеры — управляют запуском задач
    • Определяют условия выполнения задач
    • Триггер по умолчанию — все завершились успешно

Иллюстрация управления запуском задач с помощью триггеров

  • Отказоустойчивость — позволяет восстановить отдельную задачу
    • Повторный запуск задачи

Иллюстрация повторного запуска упавшей задачи при отказоустойчивости

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

Правила триггеров

  • Проверяет состояние предыдущих задач перед запуском
  • Позволяет задать условие продолжения выполнения
  • По умолчанию все предыдущие задачи должны завершиться успешно
  • Правила триггеров позволяют изменить это поведение

Иллюстрация проверки состояний вышестоящих задач с помощью правил триггеров

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

Применение правил триггеров

  • Задачи уведомления
  • Задачи очистки
  • Условное выполнение

Иллюстрация сценариев применения правил триггеров: уведомления и очистка

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

Основные правила триггеров

  • all_success — все предыдущие задачи выполнены успешно, без пропусков
  • all_failed — все предыдущие задачи завершились с ошибкой
  • all_done — все предыдущие задачи завершены, независимо от результата
  • one_failed — хотя бы одна задача завершилась с ошибкой
  • one_success — хотя бы одна задача выполнена успешно
  • none_failed — все предыдущие задачи выполнены успешно или пропущены
Введение в Apache Airflow на Python

Реализация правил триггеров

  • from airflow.utils.trigger_rule import TriggerRule
  • Правила триггеров применяются к декоратору @task
  • trigger_rule=TriggerRule.<TriggerRule enum>
@task(trigger_rule=TriggerRule.ALL_SUCCESS)
def run_if_everything_succeeds:
  print('All previous tasks succeeded!')
Введение в Apache Airflow на Python

Атрибуты отказоустойчивости

  • retries — параметр задачи, определяющий количество повторных попыток Airflow перед тем, как пометить задачу как неудавшуюся
  • retry_delay — параметр задачи, принимающий timedelta и задающий время ожидания между попытками
@task(
  trigger_rule=TriggerRule.ONE_SUCCESS,
  retries=3,
  retry_delay=timedelta(minutes=5)
)
Введение в Apache Airflow на Python

TriggerDagRunOperator

  • Позволяет одному DAG запускать другой DAG
  • Предоставляет код и DAG-и
  • from airflow.providers.standard.operators.trigger_dagrun import TriggerDagRunOperator
  • trigger_dag_id — должен совпадать с dag_id
  • wait_for_completion — блокирует дальнейшие задачи до завершения
  • poke_interval — частота проверки
  • conf — данные, передаваемые дочернему DAG
Введение в Apache Airflow на Python

Пример 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    # Run first task, then kick off child Dag
trigger_child >> cleanup  # Run a cleanup task after child Dag completes
Введение в Apache Airflow на Python

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

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

Preparing Video For Download...