Robustní DAGy s ošetřením chyb a opakováním

Tvorba datových pipeline s Airflow

Volker Janz

Senior Developer Advocate at Astronomer

Proč úlohy selhávají

 

Problémy s připojením úloh

 

  • Timeouty a limity API
  • Krátké výpadky sítě a zpoždění DNS
  • Více úloh soutěžících o stejný zdroj
Tvorba datových pipeline s Airflow

Opakování (retries)

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

 

  • Úloha se při selhání opakuje až
  • Mezi pokusy čeká 2 minuty
Tvorba datových pipeline s Airflow

Exponential backoff

 

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

Exponential backoff

 

  • Prodleva se po každém pokusu zdvojnásobí
  • Dává externím službám čas na zotavení
Tvorba datových pipeline s 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():
    ...

 

  • Slovník context: dag, ti, exception, URL logu
  • Přesměruj na Slack, PagerDuty nebo jiný alertovací nástroj
Tvorba datových pipeline s Airflow

Callbacky na úrovni DAGu vs. úlohy

Na úrovni DAGu (zachytí vše)

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

Na úrovni úlohy (konkrétní)

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

 

  • Úroveň úlohy přebíjí úroveň DAGu
  • Kombinuj obě pro vícevrstvé alertování
Tvorba datových pipeline s Airflow

max_consecutive_failed_dag_runs

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

Neúspěšné spuštění DAGu následované automatickým pozastavením

  • Zabraňuje zahlcení alerty při opakovaných selháních
  • DAG se obnoví po ručním odpauzování
Tvorba datových pipeline s Airflow

Všechno dohromady

@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():
        ...
  • Retries + backoff zvládají přechodná selhání
  • Callbacky upozorní správné lidi
  • Auto-pauza zastaví zbytečný hluk
Tvorba datových pipeline s Airflow

Pojďme cvičit!

Tvorba datových pipeline s Airflow

Preparing Video For Download...