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 real2from datetime import datetime, timedelta3from airflow import DAG4from airflow.operators.python import PythonOperator5from airflow.operators.bash import BashOperator67# Configuración común para todas las tareas del DAG8default_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}1819with 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 AM25 catchup=False,26 max_active_runs=1, # Solo una ejecución a la vez27 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-semana2# ┌───────────── 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# * * * * *910# Ejemplos:11# "0 6 * * *" → Cada día a las 6:0012# "30 8 * * 1-5" → Lunes a viernes a las 8:3013# "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:0015# "*/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 operators2from airflow.operators.python import PythonOperator, BranchPythonOperator3from airflow.operators.bash import BashOperator4from airflow.operators.empty import EmptyOperator56def 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 datos9 registros = 500010 print(f"Extraídos {registros} registros")11 # Pasar datos a la siguiente tarea via XCom12 context['ti'].xcom_push(key='num_registros', value=registros)13 return registros1415def 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 tarea20 return 'procesamiento_ligero'2122def procesar_ligero():23 print("Procesando con Pandas (dataset pequeño)")2425def procesar_pesado():26 print("Procesando con Spark (dataset grande)")2728# Dentro del DAG:29extraer = PythonOperator(task_id='extraer', python_callable=extraer_datos)3031decidir = BranchPythonOperator(32 task_id='decidir_ruta',33 python_callable=decidir_procesamiento,34)3536ligero = 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')3940cargar = BashOperator(task_id='cargar', bash_command='echo "Cargando al warehouse..."')4142# Dependencias43extraer >> 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 datos2def 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 datos89# Tarea que CONSUME datos10def 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 defecto1516 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.py2from datetime import datetime, timedelta3from airflow import DAG4from airflow.operators.python import PythonOperator5from airflow.operators.bash import BashOperator67default_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}1415def extraer_ventas(**context):16 """Simula extracción de ventas del día anterior."""17 fecha = context['ds'] # Airflow inyecta la fecha de ejecución18 print(f"Extrayendo ventas para fecha: {fecha}")19 # En producción: leer de API, S3, base de datos20 path_output = f"/tmp/ventas_raw_{fecha}.csv"21 # Simular escritura22 print(f"Datos guardados en {path_output}")23 return path_output # Se guarda como XCom automáticamente2425def 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álculos31 path_output = f"/tmp/ventas_gold_{fecha}.parquet"32 print(f"Datos transformados guardados en {path_output}")33 return path_output3435def 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 INSERT42 print(f"Datos del {fecha} cargados exitosamente")4344def 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")5051with 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 6AM57 catchup=False,58 max_active_runs=1,59 tags=['produccion', 'ventas', 'etl'],60) as dag:6162 extraer = PythonOperator(63 task_id='extraer',64 python_callable=extraer_ventas,65 )6667 transformar = PythonOperator(68 task_id='transformar',69 python_callable=transformar_ventas,70 )7172 cargar = PythonOperator(73 task_id='cargar',74 python_callable=cargar_warehouse,75 )7677 validar = PythonOperator(78 task_id='validar',79 python_callable=validar_calidad,80 )8182 notificar = BashOperator(83 task_id='notificar_exito',84 bash_command='echo "Pipeline de ventas completado para {{ ds }}"',85 )8687 # Flujo: extraer → transformar → cargar → validar → notificar88 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
- 01.Logical Date (ds): la fecha que el DAG está procesando. NO es cuándo se ejecutó, es QUÉ fecha procesa. Fundamental para idempotencia.
- 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.
- 03.Pools: limitan cuántas tareas pueden ejecutarse en paralelo (ej: max 3 conexiones a la base de datos).
- 04.SLAs: alertan si una tarea no se completó en el tiempo esperado.
- 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 logica2docker exec airflow-standalone airflow dags test etl_clientes 2024-03-1534# 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
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.
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".
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
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...