ÎncepețiÎncepe gratuit

Implementarea unei reguli de declanșare

După ce ai creat un workflow, îți dai seama că DAG-ul ar beneficia de notificări atunci când cel puțin un task eșuează. Decizi să implementezi un task care aplică verificarea one failed în DAG-ul tău, pentru a primi o alertă dacă vreun task eșuează.

Toate celelalte taskuri au fost deja definite, iar obiectele task și dag sunt deja importate pentru tine.

Acest exercițiu face parte din cursul

Introducere în Apache Airflow în Python

Vezi cursul

Instrucțiuni pentru exercițiu

  • Importă biblioteca corespunzătoare pentru a folosi regulile de declanșare.
  • Adaugă atributul corespunzător regulii de declanșare la taskul notify_on_failure.
  • Setează atributul astfel încât taskul să se declanșeze când unul sau mai multe taskuri din amonte eșuează.
  • Setează notify_on_failure ca dependență din aval față de cele două taskuri de transformare.

Exercițiu interactiv practic

Încearcă acest exercițiu completând acest cod de exemplu.

# 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()
Editează și rulează codul