Triggers och feltolerans

Introduktion till Apache Airflow i Python

Mike Metzger

Data Engineer

Bygg robusta pipelines

  • Triggers – styr när uppgifter kan köras
    • Definierar uppgifter som körs under specifika villkor
    • Standardtrigger – all success

Illustration av triggers som styr när uppgifter kan köras

  • Feltolerans – möjliggör återhämtning inom en uppgift
    • Försök om enskild uppgift igen

Illustration av feltolerans som försöker om en misslyckad uppgift

Introduktion till Apache Airflow i Python

Triggerregler

  • Kontrollerar statusen för tidigare uppgifter innan en uppgift startar
  • Låter dig styra om uppgifter ska fortsätta
  • Som standard måste alla tidigare uppgifter slutföras utan fel
  • Triggerregler låter dig ändra det

Illustration av triggerregler som kontrollerar uppströmsuppgifters status

Introduktion till Apache Airflow i Python

Användningsfall för triggerregler

  • Aviseringsuppgifter
  • Rensningsuppgifter
  • Villkorsstyrd körning

Illustration av användningsfall för triggerregler, som aviseringar och rensningsuppgifter

Introduktion till Apache Airflow i Python

Viktiga triggerregler

  • all_success – Alla tidigare uppgifter slutfördes utan fel, inga överhoppade
  • all_failed – Alla tidigare uppgifter misslyckades
  • all_done – Alla tidigare uppgifter slutfördes, oavsett utfall
  • one_failed – Minst en uppgift misslyckades
  • one_success – Minst en uppgift lyckades
  • none_failed – Alla tidigare uppgifter lyckades eller hoppades över
Introduktion till Apache Airflow i Python

Implementering av triggerregler

  • from airflow.utils.trigger_rule import TriggerRule
  • Triggerregler tillämpas på @task-dekoratorn
  • trigger_rule=TriggerRule.<TriggerRule enum>
@task(trigger_rule=TriggerRule.ALL_SUCCESS)
def run_if_everything_succeeds:
  print('All previous tasks succeeded!')
Introduktion till Apache Airflow i Python

Attribut för feltolerans

  • retries – uppgiftsparameter som anger hur många gånger Airflow ska försöka köra en misslyckad uppgift innan den markeras som misslyckad
  • retry_delay – uppgiftsparameter som tar emot ett timedelta och anger väntetiden mellan försöken
@task(
  trigger_rule=TriggerRule.ONE_SUCCESS,
  retries=3,
  retry_delay=timedelta(minutes=5)
)
Introduktion till Apache Airflow i Python

TriggerDagRunOperator

  • Låter en DAG starta en annan DAG
  • Tillhandahåller kod och DAG:ar
  • from airflow.providers.standard.operators.trigger_dagrun import TriggerDagRunOperator
  • trigger_dag_id – Måste matcha dag_id
  • wait_for_completion – Blockerar fler uppgifter tills körningen är klar
  • poke_interval – Hur ofta status kontrolleras
  • conf – Data som skickas till barn-DAG
Introduktion till Apache Airflow i Python

Exempel på 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
Introduktion till Apache Airflow i Python

Nu kör vi en övning!

Introduktion till Apache Airflow i Python

Preparing Video For Download...