Robusta DAG:ar med felhantering och återförsök

Bygg datapipelines med Airflow

Volker Janz

Senior Developer Advocate at Astronomer

Varför uppgifter misslyckas

 

Problem med uppgiftsanslutningar

 

  • API-timeouts och hastighetsbegränsningar
  • Kortvariga nätverksavbrott och DNS-fördröjningar
  • Flera uppgifter som tävlar om samma resurs
Bygg datapipelines med Airflow

Återförsök

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

 

  • Uppgiften försöker upp till 3 gånger vid fel
  • Väntar 2 minuter mellan varje försök
Bygg datapipelines med Airflow

Exponentiell backoff

 

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

Exponentiell backoff

 

  • Fördröjningen fördubblas efter varje försök
  • Ger externa tjänster tid att återhämta sig
Bygg datapipelines med 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, logg-URL
  • Skicka vidare till Slack, PagerDuty eller annat aviseringsverktyg
Bygg datapipelines med Airflow

Callbacks på DAG- och uppgiftsnivå

DAG-nivå (övergripande)

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

Uppgiftsnivå (specifik)

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

 

  • Uppgiftsnivå åsidosätter DAG-nivå
  • Kombinera båda för lageruppbyggd avisering
Bygg datapipelines med Airflow

max_consecutive_failed_dag_runs

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

Misslyckade DAG-körningar följt av automatisk paus

  • Minskar aviseringströtthet vid återkommande fel
  • DAG:en återupptas när den manuellt aktiveras igen
Bygg datapipelines med Airflow

Allt tillsammans

@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():
        ...
  • Återförsök + backoff hanterar tillfälliga fel
  • Callbacks aviserar rätt personer
  • Automatisk paus dämpar bruset
Bygg datapipelines med Airflow

Nu kör vi en övning!

Bygg datapipelines med Airflow

Preparing Video For Download...