НачатьНачать бесплатно

Реализация правила запуска

После создания рабочего процесса вы понимаете, что DAG выиграет от получения уведомлений при сбое хотя бы одной задачи. Вы решаете добавить задачу, которая реализует проверку one failed и оповещает вас о сбое любой задачи в DAG.

Все остальные задачи уже определены, а объекты task и dag импортированы заранее.

Это упражнение является частью курса

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

Посмотреть курс

Инструкции к упражнению

  • Импортируйте необходимую библиотеку для работы с правилами запуска.
  • Добавьте соответствующий атрибут правила запуска к задаче notify_on_failure.
  • Настройте атрибут так, чтобы задача запускалась при сбое одной или нескольких вышестоящих задач.
  • Установите notify_on_failure как нижестоящую зависимость двух задач преобразования данных.

Интерактивное практическое упражнение

Попробуйте выполнить это упражнение, дополнив этот пример кода.

# Import TriggerRule
from airflow.utils.____ import ____

@dag(schedule="@daily", start_date=datetime(2026, 5, 1))
def etl_pipeline():

    # Trigger notify_on_failure when any upstream task fails
    @task(____=TriggerRule.____)
    def notify_on_failure(**context) -> None:
        dag_id = context["dag"].dag_id
        run_id = context["run_id"]
        print(f"ALERT: A task failed in DAG '{dag_id}', run '{run_id}'. Sending notification...")

    # Set notify_on_failure downstream of both transform tasks
    [transform_users(), transform_orders()] ____ notify_on_failure()
    
etl_pipeline()
Редактировать и запускать код