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
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_failureca 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()