Реалізація правила спрацювання (Trigger Rule)
Після створення робочого процесу ви розумієте, що вашому Dag стане в пригоді інформування, якщо принаймні одне з завдань завершиться помилкою. Ви вирішуєте додати завдання, яке реалізує перевірку one failed у вашому Dag, щоб сповіщати вас, якщо будь-яке завдання в Dag падає.
Усі інші завдання вже визначено, а об'єкти task і dag імпортовано для вас.
Ця вправа є частиною курсу
Вступ до Apache Airflow на Python
Інструкції до вправи
- Імпортуйте відповідну бібліотеку для використання правил спрацювання (trigger rules).
- Додайте потрібний атрибут правила спрацювання до завдання
notify_on_failure. - Налаштуйте атрибут так, щоб завдання спрацьовувало, коли один або більше вхідних (upstream) завдань завершуються помилкою.
- Задайте
notify_on_failureяк залежність униз за потоком (downstream) для двох завдань перетворення.
Інтерактивна практична вправа
Спробуйте виконати цю вправу, доповнивши цей зразок коду.
# 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()