Реализация правила запуска
После создания рабочего процесса вы понимаете, что 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()