Niezawodne DAG-i: obsługa błędów i ponawianie zadań

Budowanie potoków danych z Airflow

Volker Janz

Senior Developer Advocate at Astronomer

Dlaczego zadania kończą się błędem?

 

Problemy z połączeniem zadań

 

  • Przekroczenia limitu czasu i limity zapytań API
  • Krótkie przerwy w sieci i opóźnienia rozwiązywania nazw domen
  • Wiele zadań rywalizujących o ten sam zasób
Budowanie potoków danych z Airflow

Ponawianie zadań

@task(
    retries=3,
    retry_delay=timedelta(minutes=2),
)
def fetch_weather():
    response = requests.get("https://api.weather.com/forecast")
    response.raise_for_status()
    return response.json()

 

  • Zadanie ponawia próbę maksymalnie 3 razy po błędzie
  • Między kolejnymi próbami odczekuje 2 minuty
Budowanie potoków danych z Airflow

Wykładnicze wydłużanie opóźnień

 

@task(
    retries=3,
    retry_delay=timedelta(minutes=2),
    retry_exponential_backoff=2.0,
)
def fetch_weather():
    ...

Wykładnicze wydłużanie opóźnień

 

  • Opóźnienie podwaja się po każdej próbie
  • Daje zewnętrznym usługom czas na odzyskanie sprawności
Budowanie potoków danych z Airflow

on_failure_callback

def alert_on_failure(context):
    dag_id = context["dag"].dag_id
    task_id = context["ti"].task_id
    print(f"ALERT: {task_id} in {dag_id} failed!")

@task(on_failure_callback=alert_on_failure)
def fetch_weather():
    ...

 

  • Słownik context: dag, ti, exception, URL logu
  • Powiadomienia możesz wysyłać do Slacka, PagerDuty lub innego narzędzia
Budowanie potoków danych z Airflow

Callbacki na poziomie DAG-a i zadania

Na poziomie DAG-a (ogólny)

@dag(
    on_failure_callback=alert_team,
)
def my_pipeline():
    ...

Na poziomie zadania (szczegółowy)

@task(
    on_failure_callback=page_oncall,
)
def critical_step():
    ...

 

  • Callback na poziomie zadania nadpisuje callback DAG-a
  • Używaj obu, aby tworzyć wielowarstwowe alerty
Budowanie potoków danych z Airflow

max_consecutive_failed_dag_runs

@dag(max_consecutive_failed_dag_runs=3)
def monitoring_pipeline():
    ...

Nieudane uruchomienia DAG-a zakończone automatycznym wstrzymaniem

  • Zapobiega zalewowi alertów przy powtarzających się błędach
  • DAG wznawia działanie po ręcznym odblokowaniu
Budowanie potoków danych z Airflow

Wszystko razem

@dag(
    max_consecutive_failed_dag_runs=3,
    on_failure_callback=alert_team,
)
def weather_pipeline():

    @task(
        retries=3,
        retry_delay=timedelta(minutes=2),
        retry_exponential_backoff=True,
        on_failure_callback=page_oncall,
    )
    def fetch_weather():
        ...
  • Ponawianie i wydłużanie opóźnień radzą sobie z przejściowymi błędami
  • Callbacki powiadamiają odpowiednie osoby
  • Automatyczne wstrzymanie eliminuje zbędny szum alertów
Budowanie potoków danych z Airflow

Czas na praktykę!

Budowanie potoków danych z Airflow

Preparing Video For Download...