DAGs robustos con gestión de fallos y reintentos

Creación de canalizaciones de datos con Airflow

Volker Janz

Senior Developer Advocate at Astronomer

Por qué fallan las tareas

 

Problemas de conexión de tareas

 

  • Timeouts de API y límites de velocidad
  • Cortes breves de red y retrasos en la resolución de dominios
  • Varias tareas compiten por el mismo recurso
Creación de canalizaciones de datos con Airflow

Reintentos

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

 

  • Reintenta la tarea hasta 3 veces si falla
  • Espera 2 minutos entre intentos
Creación de canalizaciones de datos con Airflow

Retroceso exponencial

 

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

Retroceso exponencial

 

  • El retraso se duplica tras cada reintento
  • Da tiempo a que servicios externos se recuperen
Creación de canalizaciones de datos con Airflow

on_failure_callback

def alert_on_failure(context):
    dag_id = context["dag"].dag_id
    task_id = context["ti"].task_id
    print(f"ALERTA: {task_id} en {dag_id} ha fallado!")

@task(on_failure_callback=alert_on_failure)
def fetch_weather():
    ...

 

  • Diccionario context: dag, ti, exception, log URL
  • Enruta a Slack, PagerDuty o cualquier herramienta de alertas
Creación de canalizaciones de datos con Airflow

Callbacks a nivel de Dag vs tarea

Nivel Dag (general)

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

Nivel tarea (específico)

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

 

  • El nivel tarea sobrescribe al nivel Dag
  • Usa ambos para alertas en capas
Creación de canalizaciones de datos con Airflow

max_consecutive_failed_dag_runs

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

Ejecuciones de Dag fallidas seguidas de pausa automática

  • Evita la fatiga de alertas por fallos persistentes
  • El Dag se reanuda al quitar la pausa manualmente
Creación de canalizaciones de datos con Airflow

Todo junto

@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():
        ...
  • Reintentos + backoff manejan fallos transitorios
  • Los callbacks avisan a las personas adecuadas
  • La pausa automática reduce el ruido
Creación de canalizaciones de datos con Airflow

¡Vamos a practicar!

Creación de canalizaciones de datos con Airflow

Preparing Video For Download...