Robuste Dags mit Fehlerbehandlung und Retries

Data-Pipelines mit Airflow aufbauen

Volker Janz

Senior Developer Advocate at Astronomer

Warum Tasks fehlschlagen

 

Aufgaben: Verbindungsprobleme

 

  • API-Timeouts und Rate Limits
  • Kurze Netzunterbrechungen und Verzögerungen bei der Namensauflösung
  • Mehrere Tasks konkurrieren um dieselbe Ressource
Data-Pipelines mit Airflow aufbauen

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

 

  • Task versucht es bei Fehlern bis zu 3-mal erneut
  • Wartet 2 Minuten zwischen den Versuchen
Data-Pipelines mit Airflow aufbauen

Exponentielles Backoff

 

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

Exponentielles Backoff

 

  • Verzögerung verdoppelt sich nach jedem Retry
  • Gibt externen Diensten Zeit zur Erholung
Data-Pipelines mit Airflow aufbauen

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
  • Leite an Slack, PagerDuty oder jedes Alerting-Tool weiter
Data-Pipelines mit Airflow aufbauen

Callbacks auf Dag- vs. Task-Level

Dag-Level (Catch-all)

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

Task-Level (spezifisch)

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

 

  • Task-Level überschreibt Dag-Level
  • Nutze beides für gestaffelte Alerts
Data-Pipelines mit Airflow aufbauen

max_consecutive_failed_dag_runs

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

Fehlgeschlagene Dag-Runs mit Auto-Pause danach

  • Verhindert Alert-Müdigkeit bei dauerhaften Fehlern
  • Dag läuft weiter, wenn manuell wieder aktiviert
Data-Pipelines mit Airflow aufbauen

Alles zusammenführen

@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 handhaben vorübergehende Fehler
  • Callbacks alarmieren die richtigen Personen
  • Auto-Pause reduziert den Lärm
Data-Pipelines mit Airflow aufbauen

Lass uns üben!

Data-Pipelines mit Airflow aufbauen

Preparing Video For Download...