Saltar al contenido

lección 7

DAGs, operators y scheduling: Airflow en la práctica

Domina los conceptos centrales de Airflow: estructura de DAGs, tipos de operators, scheduling con cron y XCom para pasar datos.

60 min

### Anatomía de un DAG de producción

En la lección anterior creamos DAGs mínimos para verificar la instalación. Ahora vamos a construir un DAG de verdad: con múltiples tareas, scheduling automático, paso de datos entre tareas, reintentos y alertas. El tipo de DAG que verás en una empresa real.

Un DAG de producción tiene tres partes bien diferenciadas: la configuración (default_args, schedule, timeouts), la definición de tareas (operators), y las dependencias (el grafo). Veamos cada una en detalle.

1# dags/ventas_diarias.py — Un DAG de producción real
2from datetime import datetime, timedelta
3from airflow import DAG
4from airflow.operators.python import PythonOperator
5from airflow.operators.bash import BashOperator
6
7# Configuración común para todas las tareas del DAG
8default_args = {
9 'owner': 'data-team',
10 'depends_on_past': False,
11 'email': ['data-alerts@empresa.com'],
12 'email_on_failure': True,
13 'email_on_retry': False,
14 'retries': 3,
15 'retry_delay': timedelta(minutes=5),
16 'execution_timeout': timedelta(hours=1),
17}
18
19with DAG(
20 dag_id='ventas_diarias_etl',
21 default_args=default_args,
22 description='Pipeline diario de ventas: ingesta → transformación → warehouse',
23 start_date=datetime(2024, 1, 1),
24 schedule='0 6 * * *', # Todos los días a las 6:00 AM
25 catchup=False,
26 max_active_runs=1, # Solo una ejecución a la vez
27 tags=['produccion', 'ventas'],
28) as dag:
29 pass # Tareas definidas abajo

Configuración de un DAG de producción: retries, timeouts, alertas, schedule

Desglosemos los campos más importantes:

  • default_args: configuración heredada por todas las tareas. Evita repetir retries, timeouts, etc. en cada tarea.
  • schedule: expresión cron. "0 6 * * *" = minuto 0, hora 6, todos los días, todos los meses, todos los días de la semana. Es decir: 6AM cada día.
  • catchup=False: si el DAG estuvo pausado 3 días, NO ejecuta los 3 días perdidos al despausarse. En producción casi siempre quieres False.
  • max_active_runs=1: impide que se lancen dos ejecuciones simultáneas del mismo DAG. Evita conflictos de recursos.
  • retries=3, retry_delay=5min: si una tarea falla, reintenta 3 veces con 5 minutos entre intentos.

### Expresiones cron: el lenguaje del scheduling

Cron lleva desde 1979 programando tareas en Unix (Version 7 Unix, Bell Labs; la version de Vixie que usamos hoy es de 1987). Su sintaxis es críptica pero universal — la encontrarás en Airflow, GitHub Actions, Kubernetes, y cualquier scheduler que se precie. Merece la pena aprenderla:

1# Formato cron: minuto hora día-mes mes día-semana
2# ┌───────────── minuto (0-59)
3# │ ┌───────────── hora (0-23)
4# │ │ ┌───────────── día del mes (1-31)
5# │ │ │ ┌───────────── mes (1-12)
6# │ │ │ │ ┌───────────── día de la semana (0-6, 0=domingo)
7# │ │ │ │ │
8# * * * * *
9
10# Ejemplos:
11# "0 6 * * *" → Cada día a las 6:00
12# "30 8 * * 1-5" → Lunes a viernes a las 8:30
13# "0 */2 * * *" → Cada 2 horas (00:00, 02:00, 04:00...)
14# "0 9 1 * *" → El día 1 de cada mes a las 9:00
15# "*/15 * * * *" → Cada 15 minutos

Cheatsheet de cron. Usa crontab.guru para verificar tus expresiones.

### Operators: los bloques de construcción

Un Operator es una plantilla de tarea. Define QUÉ tipo de trabajo se ejecuta. Airflow viene con docenas de operators built-in, y puedes instalar más con providers (plugins para servicios externos como AWS, GCP, Slack, etc.). Los más usados:

  • PythonOperator: ejecuta una función Python. El más flexible y el que más usarás.
  • BashOperator: ejecuta un comando de shell. Útil para scripts existentes o comandos del sistema.
  • EmptyOperator (DummyOperator): no hace nada. Sirve como punto de unión en DAGs complejos.
  • BranchPythonOperator: como PythonOperator pero devuelve el task_id de la siguiente tarea a ejecutar. Es el Choice de Airflow.
  • Providers: S3Operator, PostgresOperator, SlackOperator, etc. Vienen en paquetes pip separados.
1# Ejemplo completo con varios tipos de operators
2from airflow.operators.python import PythonOperator, BranchPythonOperator
3from airflow.operators.bash import BashOperator
4from airflow.operators.empty import EmptyOperator
5
6def extraer_datos(**context):
7 """Extrae datos y devuelve cuántos registros hay."""
8 # En producción aquí leerías de una API, S3, o base de datos
9 registros = 5000
10 print(f"Extraídos {registros} registros")
11 # Pasar datos a la siguiente tarea via XCom
12 context['ti'].xcom_push(key='num_registros', value=registros)
13 return registros
14
15def decidir_procesamiento(**context):
16 """Decide qué rama tomar según el volumen."""
17 registros = context['ti'].xcom_pull(task_ids='extraer', key='num_registros')
18 if registros > 10000:
19 return 'procesamiento_pesado' # task_id de la siguiente tarea
20 return 'procesamiento_ligero'
21
22def procesar_ligero():
23 print("Procesando con Pandas (dataset pequeño)")
24
25def procesar_pesado():
26 print("Procesando con Spark (dataset grande)")
27
28# Dentro del DAG:
29extraer = PythonOperator(task_id='extraer', python_callable=extraer_datos)
30
31decidir = BranchPythonOperator(
32 task_id='decidir_ruta',
33 python_callable=decidir_procesamiento,
34)
35
36ligero = PythonOperator(task_id='procesamiento_ligero', python_callable=procesar_ligero)
37pesado = PythonOperator(task_id='procesamiento_pesado', python_callable=procesar_pesado)
38unir = EmptyOperator(task_id='unir', trigger_rule='none_failed_min_one_success')
39
40cargar = BashOperator(task_id='cargar', bash_command='echo "Cargando al warehouse..."')
41
42# Dependencias
43extraer >> decidir >> [ligero, pesado] >> unir >> cargar

Patrón de branching en Airflow: BranchPythonOperator + EmptyOperator como punto de unión

### XCom: pasar datos entre tareas

En Step Functions, los datos fluyen automáticamente del output de un estado al input del siguiente. En Airflow no es tan automático — necesitas XCom (Cross-Communication). XCom es un sistema de mensajería ligero entre tareas: una tarea "pushea" un valor y otra lo "pullea".

1# Tarea que PRODUCE datos
2def extraer(**context):
3 datos = {"registros": 5000, "fuente": "api_ventas"}
4 # Forma explícita:
5 context['ti'].xcom_push(key='resultado_extraccion', value=datos)
6 # O forma implícita (el return se guarda como XCom con key='return_value'):
7 return datos
8
9# Tarea que CONSUME datos
10def transformar(**context):
11 # Pullear XCom de la tarea 'extraer'
12 datos = context['ti'].xcom_pull(task_ids='extraer', key='resultado_extraccion')
13 # O si usaste return:
14 datos = context['ti'].xcom_pull(task_ids='extraer') # key='return_value' por defecto
15
16 print(f"Transformando {datos['registros']} registros de {datos['fuente']}")

XCom: el mecanismo para pasar datos entre tareas en Airflow

XCom NO es para pasar datasets grandes. Se almacena en la base de datos de metadatos de Airflow (SQLite/Postgres). Límite práctico: unos pocos KB. Para datasets usa archivos intermedios en S3/disco. XCom es para metadatos: paths de archivos, conteos, flags de estado.

### Un DAG completo de producción: ETL de ventas diarias

Juntemos todo en un DAG realista. Este es el tipo de DAG que mantendrías en una empresa: scheduling diario, paso de metadatos entre tareas, reintentos y notificación de fallos.

1# dags/etl_ventas_produccion.py
2from datetime import datetime, timedelta
3from airflow import DAG
4from airflow.operators.python import PythonOperator
5from airflow.operators.bash import BashOperator
6
7default_args = {
8 'owner': 'data-team',
9 'retries': 3,
10 'retry_delay': timedelta(minutes=5),
11 'email_on_failure': True,
12 'email': ['data-team@empresa.com'],
13}
14
15def extraer_ventas(**context):
16 """Simula extracción de ventas del día anterior."""
17 fecha = context['ds'] # Airflow inyecta la fecha de ejecución
18 print(f"Extrayendo ventas para fecha: {fecha}")
19 # En producción: leer de API, S3, base de datos
20 path_output = f"/tmp/ventas_raw_{fecha}.csv"
21 # Simular escritura
22 print(f"Datos guardados en {path_output}")
23 return path_output # Se guarda como XCom automáticamente
24
25def transformar_ventas(**context):
26 """Limpia y transforma los datos extraídos."""
27 path_raw = context['ti'].xcom_pull(task_ids='extraer')
28 fecha = context['ds']
29 print(f"Transformando datos de {path_raw}")
30 # En producción: Pandas, limpieza, cálculos
31 path_output = f"/tmp/ventas_gold_{fecha}.parquet"
32 print(f"Datos transformados guardados en {path_output}")
33 return path_output
34
35def cargar_warehouse(**context):
36 """Carga datos transformados al warehouse."""
37 path_gold = context['ti'].xcom_pull(task_ids='transformar')
38 fecha = context['ds']
39 print(f"Cargando {path_gold} al warehouse")
40 # En producción: COPY INTO, INSERT, etc.
41 # IDEMPOTENTE: DELETE WHERE fecha = X, luego INSERT
42 print(f"Datos del {fecha} cargados exitosamente")
43
44def validar_calidad(**context):
45 """Ejecuta checks de calidad post-carga."""
46 fecha = context['ds']
47 print(f"Validando datos del {fecha} en el warehouse")
48 # En producción: contar registros, verificar NULLs, etc.
49 print("Calidad OK: 0 anomalías detectadas")
50
51with DAG(
52 dag_id='etl_ventas_produccion',
53 default_args=default_args,
54 description='Pipeline diario de ventas: API → Gold → Warehouse',
55 start_date=datetime(2024, 1, 1),
56 schedule='0 6 * * *', # Cada día a las 6AM
57 catchup=False,
58 max_active_runs=1,
59 tags=['produccion', 'ventas', 'etl'],
60) as dag:
61
62 extraer = PythonOperator(
63 task_id='extraer',
64 python_callable=extraer_ventas,
65 )
66
67 transformar = PythonOperator(
68 task_id='transformar',
69 python_callable=transformar_ventas,
70 )
71
72 cargar = PythonOperator(
73 task_id='cargar',
74 python_callable=cargar_warehouse,
75 )
76
77 validar = PythonOperator(
78 task_id='validar',
79 python_callable=validar_calidad,
80 )
81
82 notificar = BashOperator(
83 task_id='notificar_exito',
84 bash_command='echo "Pipeline de ventas completado para {{ ds }}"',
85 )
86
87 # Flujo: extraer → transformar → cargar → validar → notificar
88 extraer >> transformar >> cargar >> validar >> notificar

DAG de producción completo. Nota cómo {{ ds }} inyecta la fecha — es templating Jinja2.

Cuidado con context["ds"]: NO es el dia en que el DAG se ejecuta. Es la fecha que el DAG esta PROCESANDO (logical_date). Una ejecucion diaria cubre el intervalo [dia D, dia D+1) y se lanza cuando el intervalo se cierra — es decir, al dia siguiente. Asi que el DAG que corre a las 6AM del 15 de enero tiene ds = "2024-01-14". Es justo lo que quieres para la idempotencia (siempre procesa la misma fecha), pero si necesitas saber "que dia es HOY", usa context["data_interval_end"], que si cae el dia de la ejecucion.

### Conceptos clave que debes dominar

  1. 01.Logical Date (ds): la fecha que el DAG está procesando. NO es cuándo se ejecutó, es QUÉ fecha procesa. Fundamental para idempotencia.
  2. 02.Trigger Rule: por defecto una tarea solo se ejecuta si TODAS sus dependencias tuvieron éxito. Puedes cambiarlo: all_success, all_failed, one_success, none_failed, etc.
  3. 03.Pools: limitan cuántas tareas pueden ejecutarse en paralelo (ej: max 3 conexiones a la base de datos).
  4. 04.SLAs: alertan si una tarea no se completó en el tiempo esperado.
  5. 05.Task Groups: agrupan tareas visualmente en la UI sin cambiar la lógica.

### Como probar tus DAGs

El comando que mas vas a usar es airflow dags test: ejecuta el DAG entero en local, con XCom de verdad, sin scheduler y sin despausar nada. Es como pulsar un boton de "ejecutar todo ahora".

1# Ejecutar un DAG entero con una fecha logica
2docker exec airflow-standalone airflow dags test etl_clientes 2024-03-15
3
4# Ejecutar una sola tarea (ojo: el XCom NO se comparte entre tasks test)
5docker exec airflow-standalone airflow tasks test etl_clientes extraer_clientes 2024-03-15

airflow dags test ejecuta el DAG de punta a punta. tasks test ejecuta una tarea suelta

Zonas horarias: Airflow interpreta el schedule en UTC por defecto. En Espana eso son las 7:00 en invierno y las 8:00 en verano. Si quieres pensar en hora local, dale al DAG su zona con pendulum: start_date=pendulum.datetime(2024, 1, 1, tz="Europe/Madrid"). Con eso, "las 6:00" son las 6:00 de Madrid todo el anio.

Una nota sobre email_on_failure: lo dejamos en los DAGs porque es lo que veras en produccion, pero en este montaje local (standalone con SQLite) no hay servidor de correo configurado. El aviso no va a llegar. Cuando quieras que funcione de verdad, configura el bloque [smtp] de Airflow o usa un on_failure_callback que escriba en Slack.

## ejercicios

[01]

Construir un DAG ETL con schedule y XCom

Crea un DAG "etl_clientes" que se ejecute cada día a las 7:30AM. Tiene 3 tareas: extraer_clientes (devuelve path vía XCom), transformar (lee el path de XCom) y cargar. Usa default_args con 2 retries.

Cargando editor...
[02]

DAG con BranchPythonOperator

Crea un DAG que extraiga datos y luego decida: si es lunes (weekday=0) ejecuta "reporte_semanal", cualquier otro día ejecuta "reporte_diario". Ambas ramas convergen en "enviar_email".

Cargando editor...
[03]

Traducir requisitos de negocio a cron

Traduce estos requisitos a expresiones cron y comprueba que disparan cuando crees. El codigo las verifica mostrando los proximos 4 disparos.

💡 Resultado esperado

1. 0 * * * *          -> Mon 11/03 01:00 | Mon 11/03 02:00 | Mon 11/03 03:00 | Mon 11/03 04:00
2. 0 8 * * 1-5        -> Mon 11/03 08:00 | Tue 12/03 08:00 | Wed 13/03 08:00 | Thu 14/03 08:00
3. 0 0 1 * *          -> Mon 01/04 00:00 | Wed 01/05 00:00 | Sat 01/06 00:00 | Mon 01/07 00:00
4. */15 9-17 * * *    -> Mon 11/03 09:00 | Mon 11/03 09:15 | Mon 11/03 09:30 | Mon 11/03 09:45
Cargando editor...

Regístrate para guardar tu progreso.

## comentarios

Reporta erratas, ayuda a otros o comparte tu opinión. Sé constructivo.

Inicia sesión para comentar y responder.

cargando comentarios...