Wyzwalacze i tolerancja na błędy

Wprowadzenie do Apache Airflow w Pythonie

Mike Metzger

Data Engineer

Tworzenie niezawodnych potoków

  • Wyzwalacze – kontrolują, kiedy zadania mogą być uruchamiane
    • Definiują zadania wykonywane pod określonymi warunkami
    • Domyślny wyzwalacz – wszystkie zakończone sukcesem

Ilustracja wyzwalaczy kontrolujących uruchamianie zadań

  • Tolerancja na błędy – umożliwia odzyskiwanie w obrębie zadania
    • Ponowienie pojedynczego zadania

Ilustracja tolerancji na błędy – ponowienie nieudanego zadania

Wprowadzenie do Apache Airflow w Pythonie

Reguły wyzwalacza

  • Sprawdza stan poprzednich zadań przed uruchomieniem kolejnego
  • Pozwala określić, czy zadania powinny być kontynuowane
  • Domyślnie wszystkie poprzednie zadania muszą zakończyć się sukcesem
  • Reguły wyzwalacza umożliwiają zmianę tego zachowania

Ilustracja reguł wyzwalacza sprawdzających stany zadań nadrzędnych

Wprowadzenie do Apache Airflow w Pythonie

Zastosowania reguł wyzwalacza

  • Zadania powiadomień
  • Zadania czyszczenia
  • Warunkowe wykonanie

Ilustracja zastosowań reguł wyzwalacza: powiadomienia i czyszczenie

Wprowadzenie do Apache Airflow w Pythonie

Kluczowe reguły wyzwalacza

  • all_success – każde poprzednie zadanie zakończone sukcesem, bez pominięć
  • all_failed – każde poprzednie zadanie zakończone niepowodzeniem
  • all_done – każde poprzednie zadanie ukończone, niezależnie od wyniku
  • one_failed – co najmniej jedno zadanie zakończone niepowodzeniem
  • one_success – co najmniej jedno zadanie zakończone sukcesem
  • none_failed – wszystkie poprzednie zadania zakończone sukcesem lub pominięte
Wprowadzenie do Apache Airflow w Pythonie

Implementacja reguły wyzwalacza

  • from airflow.utils.trigger_rule import TriggerRule
  • Reguły wyzwalacza stosuje się do dekoratora @task
  • trigger_rule=TriggerRule.<TriggerRule enum>
@task(trigger_rule=TriggerRule.ALL_SUCCESS)
def run_if_everything_succeeds:
  print('All previous tasks succeeded!')
Wprowadzenie do Apache Airflow w Pythonie

Atrybuty tolerancji na błędy

  • retries - parametr zadania określający, ile razy Airflow ponowi nieudane zadanie przed oznaczeniem go jako nieudane
  • retry_delay - parametr zadania przyjmujący timedelta definiujący czas oczekiwania między próbami
@task(
  trigger_rule=TriggerRule.ONE_SUCCESS,
  retries=3,
  retry_delay=timedelta(minutes=5)
)
Wprowadzenie do Apache Airflow w Pythonie

TriggerDagRunOperator

  • Umożliwia uruchomienie jednego DAG-a przez inny
  • Dostarcza kod i DAG-i
  • from airflow.providers.standard.operators.trigger_dagrun import TriggerDagRunOperator
  • trigger_dag_id – musi odpowiadać dag_id
  • wait_for_completion – blokuje kolejne zadania do zakończenia
  • poke_interval – jak często sprawdzać stan
  • conf – dane przekazywane do podrzędnego DAG-a
Wprowadzenie do Apache Airflow w Pythonie

Przykład 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
Wprowadzenie do Apache Airflow w Pythonie

Czas na ćwiczenia!

Wprowadzenie do Apache Airflow w Pythonie

Preparing Video For Download...