Implementace pravidla spuštění (Trigger Rule)
Po vytvoření workflow zjistíš, že by DAG těžil z notifikací v případě, kdy selže alespoň jeden úkol. Rozhodneš se proto implementovat úkol s kontrolou one failed, který tě upozorní, pokud v DAGu selže jakýkoli úkol.
Všechny ostatní úkoly jsou již definované a objekty task i dag jsou za tebe naimportované.
Toto cvičení je součástí kurzu
Úvod do Apache Airflow v Pythonu
Pokyny k cvičení
- Naimportuj příslušnou knihovnu pro práci s pravidly spuštění (trigger rules).
- Přidej odpovídající atribut trigger rule k úkolu
notify_on_failure. - Nastav atribut tak, aby se úkol spustil, když selže jeden nebo více upstream úkolů.
- Nastav
notify_on_failurejako downstream závislost obou transformačních úkolů.
Interaktivní cvičení na vyzkoušení si v praxi
Vyzkoušejte si toto cvičení dokončením tohoto ukázkového kódu.
# 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()