Implementacja reguły wyzwalania
Po utworzeniu przepływu pracy zauważasz, że DAG zyska na tym, jeśli będzie wysyłać powiadomienia w przypadku niepowodzenia co najmniej jednego zadania. Postanawiasz zaimplementować zadanie korzystające z kontroli one failed, które powiadomi cię o niepowodzeniu któregokolwiek zadania w DAG-u.
Wszystkie pozostałe zadania zostały już zdefiniowane, a obiekty task i dag są zaimportowane.
To ćwiczenie jest częścią kursu
Wprowadzenie do Apache Airflow w Pythonie
Instrukcje do ćwiczenia
- Zaimportuj odpowiednią bibliotekę do obsługi reguł wyzwalania.
- Dodaj odpowiedni atrybut reguły wyzwalania do zadania
notify_on_failure. - Ustaw atrybut tak, aby zadanie uruchamiało się, gdy co najmniej jedno zadanie upstream zakończy się niepowodzeniem.
- Ustaw
notify_on_failurejako zależność downstream obu zadań transformacji.
Interaktywne ćwiczenie praktyczne
Spróbuj tego ćwiczenia, uzupełniając ten przykładowy kod.
# 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()