Muchos pipelines de datos empiezan como un puñado de cron jobs y terminan siendo un problema de fiabilidad: tareas que fallan en silencio, sin reintentos automáticos, sin alertas claras, hasta que alguien descubre a las tres de la mañana que un análisis crítico no se ejecutó. Apache Airflow existe justamente para resolver ese caos de orquestación.
El problema de fondo con cron
0 2 * * * python /scripts/extract_data.py
0 3 * * * python /scripts/transform_data.py
0 4 * * * python /scripts/load_data.py
Este patrón tiene varios problemas que se notan tarde: si extract_data.py falla, transform_data.py se ejecuta igual, con datos viejos o incorrectos. No hay forma sencilla de saber si algo corrió sin revisar logs a mano. Si una tarea tarda más de lo esperado, se solapa con la siguiente. Y la lógica de reintento hay que escribirla desde cero cada vez.
Airflow resuelve todo esto de raíz: dependencias explícitas entre tareas, reintentos configurables, alertas, y una interfaz visual para ver el estado de cada ejecución.
Los conceptos básicos: DAG, Task, XCom, Sensor
Un DAG (grafo acíclico dirigido) es tu workflow: una serie de tareas con dependencias explícitas entre ellas.
from airflow import DAG
from airflow.operators.python import PythonOperator
from datetime import datetime
def extract():
print("Extrayendo datos...")
return {"rows": 1000}
def transform(ti):
extracted_data = ti.xcom_pull(task_ids='extract')
print(f"Transformando {extracted_data['rows']} filas...")
def load(ti):
transformed = ti.xcom_pull(task_ids='transform')
print("Cargando...")
with DAG(
'etl_pipeline',
start_date=datetime(2026, 1, 1),
schedule_interval='0 2 * * *',
) as dag:
task_extract = PythonOperator(task_id='extract', python_callable=extract)
task_transform = PythonOperator(task_id='transform', python_callable=transform)
task_load = PythonOperator(task_id='load', python_callable=load)
task_extract >> task_transform >> task_load
Una Task es la unidad de trabajo individual (PythonOperator, BashOperator, SqlOperator...). XCom es el mecanismo por el que una tarea pasa datos a la siguiente. Un Sensor espera a que se cumpla una condición externa —que llegue un archivo a S3, por ejemplo— antes de continuar:
from airflow.sensors.s3 import S3KeySensor
wait_for_file = S3KeySensor(
task_id='wait_for_file',
bucket_name='my-bucket',
poke_interval=60,
)
task_extract >> wait_for_file >> task_transform
Reintentos y alertas
PythonOperator(
task_id='api_call',
python_callable=fetch_api,
retries=3,
retry_delay=timedelta(minutes=5),
on_failure_callback=alert_slack,
)
Esa función alert_slack recibe el contexto de la ejecución fallida y puede notificar directamente al canal correspondiente:
def alert_slack(context):
task = context['task']
dag = context['dag']
message = f"La tarea {task} falló en el DAG {dag} el {context['execution_date']}"
client.chat_postMessage(channel='#alerts', text=message)
Validar los datos, no solo que el pipeline corrió
Que un pipeline se ejecute sin errores no significa que los datos sean correctos. Great Expectations permite añadir validaciones explícitas como parte del DAG:
from great_expectations_provider.operators.great_expectations import GreatExpectationsOperator
validate = GreatExpectationsOperator(
task_id='validate_data',
datasource_name='my_data',
expectation_suite_name='user_data_expectations',
data_context_root_dir='/path/to/great_expectations',
)
task_extract >> task_transform >> validate >> task_load
Esto comprueba automáticamente cosas como si las columnas esperadas existen, si los tipos son correctos, si hay valores fuera de rango, duplicados o nulos donde no deberían. Si la validación falla, la tarea falla y se dispara el reintento, con auditoría completa de qué pasó.
Un caso típico: sincronización de inventario
Es habitual ver este patrón en equipos de retail: un proceso de sincronización de inventario con millones de referencias, corriendo como cron jobs manuales, que se desincroniza una o dos horas cada mañana y falla por timeout de forma más o menos regular, con horas de investigación cada vez que pasa. Al migrar ese mismo pipeline a Airflow —con reintentos automáticos, validación de conteos antes de cargar, y alertas claras cuando algo falla— el tiempo dedicado a arreglar problemas manualmente baja drásticamente, porque la mayoría de fallos se resuelven solos con el reintento o se detectan antes de propagarse.
Cómo desplegarlo
Self-hosted con Docker: control total, pero el overhead operacional (uptime, escalado, backups) recae en tu equipo.
FROM apache/airflow:latest-python3.9
RUN pip install great-expectations
Astronomer: Airflow gestionado, con auto-escalado, backups y monitoring incluidos. Tiene sentido cuando el coste de gestionarlo tú mismo supera el de pagar por el servicio gestionado.
Herramientas
Docker para desarrollo local, PostgreSQL como base de datos de metadatos de Airflow, Great Expectations para validación de datos, y la API de Slack para notificaciones.
Para tu primer DAG
- Instala Airflow con docker-compose
- Crea un DAG simple con tres PythonOperators
- Define las dependencias con el operador
>> - Pruébalo localmente con
airflow dags test - Añade lógica de reintento (
retries=3) - Añade monitoring (
on_failure_callback) - Despliega a producción, self-hosted o gestionado
Una vez montado esto, el caos de la automatización se convierte en orden: las tareas corren cuando tienen que correr, las alertas son claras, y queda auditoría completa de cada ejecución.
Fuentes:






Comentarios (0)
Deja un comentario
No hay comentarios aún. ¡Sé el primero en comentar!