Robuuste Dags met foutafhandeling en retries

Data-pijplijnen bouwen met Airflow

Volker Janz

Senior Developer Advocate at Astronomer

Waarom taken falen

 

Problemen met taakverbinding

 

  • API-time-outs en rate limits
  • Korte netwerkstoringen en vertraging bij domeinnaamresolutie
  • Meerdere taken die om dezelfde resource strijden
Data-pijplijnen bouwen met Airflow

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

 

  • Probeert taak tot 3 keer opnieuw bij falen
  • Wacht 2 minuten tussen pogingen
Data-pijplijnen bouwen met Airflow

Exponentiële backoff

 

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

Exponentiële backoff

 

  • Wachtijd verdubbelt na elke retry
  • Geeft externe services tijd om te herstellen
Data-pijplijnen bouwen met 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():
    ...

 

  • context-dict: dag, ti, exception, log-URL
  • Stuur door naar Slack, PagerDuty of elk alertingtool
Data-pijplijnen bouwen met Airflow

Callbacks op dag- vs taakniveau

Dag-niveau (vangt alles af)

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

Taak-niveau (specifiek)

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

 

  • Taak-niveau overschrijft dag-niveau
  • Gebruik beide voor gelaagde alerts
Data-pijplijnen bouwen met Airflow

max_consecutive_failed_dag_runs

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

Mislukte Dag-runs gevolgd door auto-pauze

  • Voorkomt alert-moeheid door aanhoudende fouten
  • Dag gaat verder na handmatig unpausen
Data-pijplijnen bouwen met Airflow

Alles combineren

@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 vangen tijdelijke fouten op
  • Callbacks waarschuwen de juiste mensen
  • Auto-pauze dempt de ruis
Data-pijplijnen bouwen met Airflow

Laten we oefenen!

Data-pijplijnen bouwen met Airflow

Preparing Video For Download...