DAG-uri robuste cu gestionarea erorilor și reîncercări

Construirea pipeline-urilor de date cu Airflow

Volker Janz

Senior Developer Advocate at Astronomer

De ce eșuează taskurile

 

Probleme de conexiune între taskuri

 

  • Timeout-uri API și limite de rată
  • Întreruperi scurte de rețea și întârzieri DNS
  • Mai multe taskuri care concurează pentru aceeași resursă
Construirea pipeline-urilor de date cu Airflow

Reîncercări

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

 

  • Taskul reîncearcă de până la 3 ori la eșec
  • Așteaptă 2 minute între fiecare încercare
Construirea pipeline-urilor de date cu Airflow

Backoff exponențial

 

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

Backoff exponențial

 

  • Intervalul se dublează după fiecare reîncercare
  • Oferă serviciilor externe timp să se recupereze
Construirea pipeline-urilor de date cu 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 log
  • Trimite alerte în Slack, PagerDuty sau alt sistem
Construirea pipeline-urilor de date cu Airflow

Callback-uri la nivel de DAG vs. task

La nivel de DAG (general)

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

La nivel de task (specific)

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

 

  • Nivelul de task suprascrie nivelul de DAG
  • Folosește ambele pentru alerte pe niveluri
Construirea pipeline-urilor de date cu Airflow

max_consecutive_failed_dag_runs

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

Rulări eșuate ale DAG-ului urmate de pauză automată

  • Reduce oboseala de alerte cauzată de eșecuri persistente
  • DAG-ul se reia când este reactivat manual
Construirea pipeline-urilor de date cu Airflow

Totul împreună

@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():
        ...
  • Reîncercările + backoff gestionează eșecurile tranzitorii
  • Callback-urile alertează persoanele potrivite
  • Pauza automată reduce zgomotul
Construirea pipeline-urilor de date cu Airflow

Hai să exersăm!

Construirea pipeline-urilor de date cu Airflow

Preparing Video For Download...