Triggere și toleranță la erori

Introducere în Apache Airflow în Python

Mike Metzger

Data Engineer

Construirea de pipeline-uri robuste

  • Triggere - controlează când pot rula task-urile
    • Definesc task-uri care se execută în condiții specifice
    • Trigger implicit - succes total

Ilustrație a triggerelor care controlează când rulează task-urile

  • Toleranța la erori - permite recuperarea în cadrul unui task
    • Reîncercarea unui task individual

Ilustrație a toleranței la erori care reîncearcă un task eșuat

Introducere în Apache Airflow în Python

Reguli de trigger

  • Verifică starea task-urilor anterioare înainte de a porni un task
  • Permite definirea dacă task-urile ar trebui să continue
  • Implicit, toate task-urile anterioare trebuie să se finalizeze cu succes
  • Regulile de trigger permit modificarea acestui comportament

Ilustrație a regulilor de trigger care verifică stările task-urilor anterioare

Introducere în Apache Airflow în Python

Utilizări ale regulilor de trigger

  • Task-uri de notificare
  • Task-uri de curățare
  • Execuție condiționată

Ilustrație a cazurilor de utilizare a regulilor de trigger, cum ar fi notificările și curățarea

Introducere în Apache Airflow în Python

Reguli de trigger principale

  • all_success - Toate task-urile anterioare s-au finalizat cu succes, fără omisiuni
  • all_failed - Toate task-urile anterioare au eșuat
  • all_done - Toate task-urile anterioare s-au finalizat, indiferent de rezultat
  • one_failed - Cel puțin un task a eșuat
  • one_success - Cel puțin un task a reușit
  • none_failed - Toate task-urile anterioare s-au finalizat cu succes sau au fost omise
Introducere în Apache Airflow în Python

Implementarea regulilor de trigger

  • from airflow.utils.trigger_rule import TriggerRule
  • Regulile de trigger se aplică decoratorului @task
  • trigger_rule=TriggerRule.<TriggerRule enum>
@task(trigger_rule=TriggerRule.ALL_SUCCESS)
def run_if_everything_succeeds:
  print('All previous tasks succeeded!')
Introducere în Apache Airflow în Python

Atribute de toleranță la erori

  • retries - parametru de task care definește de câte ori Airflow reîncearcă un task eșuat înainte de a-l marca ca eșuat
  • retry_delay - parametru de task care acceptă un timedelta definind intervalul de așteptare între reîncercări
@task(
  trigger_rule=TriggerRule.ONE_SUCCESS,
  retries=3,
  retry_delay=timedelta(minutes=5)
)
Introducere în Apache Airflow în Python

TriggerDagRunOperator

  • Permite unui DAG să declanșeze un alt DAG
  • Furnizează cod și DAG-uri
  • from airflow.providers.standard.operators.trigger_dagrun import TriggerDagRunOperator
  • trigger_dag_id - Trebuie să corespundă cu dag_id
  • wait_for_completion - Blochează task-urile următoare până la finalizare
  • poke_interval - Frecvența de verificare
  • conf - Date transmise DAG-ului copil
Introducere în Apache Airflow în Python

Exemplu TriggerDagRun

trigger_child = TriggerDagRunOperator(
        task_id="trigger_child_pipeline",
        trigger_dag_id="child_pipeline_dag",   
        wait_for_completion=True,              
        poke_interval=30,                      
        conf={                                 
            "source": "s3://my-bucket/raw/",
        },
    )

task1 >> trigger_child    # Run first task, then kick off child Dag
trigger_child >> cleanup  # Run a cleanup task after child Dag completes
Introducere în Apache Airflow în Python

Să exersăm!

Introducere în Apache Airflow în Python

Preparing Video For Download...