Zacznij terazZacznij za darmo

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

Zobacz kurs

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_failure jako 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()
Edytuj i uruchom kod