DAG robusti con gestione errori e retry

Creare data pipeline con Airflow

Volker Janz

Senior Developer Advocate at Astronomer

Perché i task falliscono

 

Problemi di connessione tra task

 

  • Timeout API e limiti di rate
  • Brevi interruzioni di rete e ritardi DNS
  • Più task che competono per la stessa risorsa
Creare data pipeline con Airflow

Retry

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

 

  • Il task riprova fino a 3 volte in caso di errore
  • Attende 2 minuti tra un tentativo e l'altro
Creare data pipeline con Airflow

Backoff esponenziale

 

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

Backoff esponenziale

 

  • Il ritardo raddoppia dopo ogni retry
  • Dà tempo ai servizi esterni di riprendersi
Creare data pipeline con 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():
    ...

 

  • Dict context: dag, ti, exception, URL dei log
  • Invia avvisi a Slack, PagerDuty o altri tool
Creare data pipeline con Airflow

Callback a livello Dag vs task

A livello Dag (catch-all)

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

A livello task (specifico)

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

 

  • Il livello task sovrascrive quello Dag
  • Usali entrambi per avvisi a livelli
Creare data pipeline con Airflow

max_consecutive_failed_dag_runs

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

Esecuzioni Dag fallite seguite da pausa automatica

  • Evita la alert fatigue da errori persistenti
  • Il Dag riparte quando lo riattivi manualmente
Creare data pipeline con Airflow

Mettiamo tutto insieme

@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():
        ...
  • Retry + backoff gestiscono errori transitori
  • Le callback avvisano le persone giuste
  • La pausa automatica riduce il rumore
Creare data pipeline con Airflow

Esercitiamoci!

Creare data pipeline con Airflow

Preparing Video For Download...