Kom igångKom igång gratis

Implementera en triggerregel

När du granskar ditt arbetsflöde inser du att DAG:en skulle bli bättre med notifieringar när minst en uppgift misslyckas. Du väljer att lägga till en uppgift som implementerar kontrollen one failed i din DAG, så att du får en avisering om någon uppgift misslyckas.

Alla övriga uppgifter är redan definierade och objekten task och dag är importerade åt dig.

Den här övningen är en del av kursen

Introduktion till Apache Airflow i Python

Visa kurs

Övningsinstruktioner

  • Importera lämpligt bibliotek för att använda triggerregler.
  • Lägg till rätt triggerregelattribut i uppgiften notify_on_failure.
  • Ställ in attributet så att uppgiften triggas när en eller flera uppströmsuppgifter misslyckas.
  • Sätt notify_on_failure som ett nedströmsberoende till de två transformeringsuppgifterna.

Interaktiv övning med praktiskt arbete

Testa den här övningen genom att slutföra den här exempelkoden.

# 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()
Redigera och kör kod