🌀 Capítulo 11 · Nivel Intermedio-Avanzado
Apache Airflow
Apache Airflow es el orquestador de workflows de datos más usado del mundo. Permite programar, monitorear y gestionar pipelines complejos usando Python. Aprender Airflow bien es crucial para cualquier Data Engineer.
Conceptos Core de Airflow
- DAG (Directed Acyclic Graph): El pipeline. Un grafo donde los nodos son Tasks y las aristas son dependencias. "Acyclic" = no hay ciclos (no hay loops infinitos).
- Task: La unidad de trabajo. Puede ser un Operator, Sensor o TaskGroup.
- Operator: Define lo que hace una Task. BashOperator, PythonOperator, SparkSubmitOperator, etc.
- Sensor: Operator especial que espera una condición (archivo en S3, dato en DB, hora específica).
- DagRun: Una instancia de ejecución de un DAG para un momento específico (logical_date).
- TaskInstance: Una instancia de ejecución de una Task en un DagRun específico. Tiene estados: queued, running, success, failed, skipped.
- Scheduler: El componente que decide cuándo ejecutar los DAGs y las Tasks.
- Executor: Cómo se ejecutan las Tasks. LocalExecutor (mismo proceso), CeleryExecutor (workers distribuidos), KubernetesExecutor (pod por task).
Arquitectura de Airflow
🏗️ Componentes de Airflow
- Web Server: UI para monitorear DAGs, TaskInstances, logs. Flask + Gunicorn.
- Scheduler: Lee DAGs del filesystem, evalúa cuándo correr, crea DagRuns y TaskInstances, los envía al Executor.
- Executor: Gestiona dónde y cómo corren las Tasks. En producción: CeleryExecutor o KubernetesExecutor.
- Workers: Procesos que ejecutan las Tasks reales. En CeleryExecutor, son workers de Celery.
- Metadata DB: PostgreSQL o MySQL. Almacena el estado de todos los DAGs, Tasks, variables y conexiones.
- DAGs Folder: Directorio donde el Scheduler escanea los archivos Python con las definiciones de DAGs.
DAGs con Taskflow API
"""
Airflow 2.x TaskFlow API: DAGs más limpios y pythónicos
La mejor manera de escribir DAGs en 2026
"""
from airflow.decorators import dag, task
from airflow.operators.empty import EmptyOperator
from datetime import datetime, timedelta
from typing import Any
import logging
logger = logging.getLogger(__name__)
# ── DAG COMPLETO: ETL Pipeline con TaskFlow API ────────────────────────────
@dag(
dag_id='orders_etl_pipeline',
description='Ingesta y transforma órdenes diariamente',
schedule='0 3 * * *', # Cron: 3am todos los días
# schedule=timedelta(hours=6), # Alternativa: cada 6 horas
start_date=datetime(2026, 1, 1),
catchup=False, # No ejecutar runs pasados al activar el DAG
max_active_runs=1, # Solo 1 run a la vez (evita conflictos)
max_active_tasks=8, # Max tasks paralelas por DAG
default_args={
'owner': 'data-team',
'depends_on_past': False,
'email_on_failure': True,
'email': ['data-alerts@company.com'],
'retries': 2,
'retry_delay': timedelta(minutes=5),
'retry_exponential_backoff': True,
'execution_timeout': timedelta(hours=2),
},
tags=['etl', 'orders', 'daily'],
doc_md="""
## Orders ETL Pipeline
Ingesta órdenes desde la API de operaciones, limpia y carga en el Data Warehouse.
### SLA: 5am UTC
### Owner: Data Engineering Team
### Runbook: https://wiki.company.com/runbooks/orders-etl
"""
)
def orders_etl_pipeline():
@task(task_id='extract_orders_api')
def extract_orders(logical_date=None, **context) -> list[dict]:
"""Extrae órdenes de la API para el día de la ejecución."""
from api_client import OrdersAPIClient
date_str = logical_date.strftime('%Y-%m-%d')
logger.info(f"Extracting orders for date: {date_str}")
client = OrdersAPIClient(api_key="{{ var.value.orders_api_key }}")
orders = client.get_orders_for_date(date_str)
logger.info(f"Extracted {len(orders)} orders")
return orders # Se pasa automáticamente como XCom
@task(task_id='validate_orders')
def validate_orders(orders: list[dict]) -> list[dict]:
"""Valida la calidad de los datos."""
valid_orders = []
invalid_count = 0
for order in orders:
# Reglas de validación
if not order.get('order_id'):
invalid_count += 1
continue
if order.get('amount', -1) < 0:
invalid_count += 1
continue
valid_orders.append(order)
logger.info(f"Validation: {len(valid_orders)} valid, {invalid_count} invalid")
# Alertar si más del 5% son inválidos
if len(orders) > 0 and invalid_count / len(orders) > 0.05:
raise ValueError(f"Too many invalid orders: {invalid_count}/{len(orders)}")
return valid_orders
@task(task_id='transform_orders')
def transform_orders(orders: list[dict]) -> list[dict]:
"""Aplica transformaciones de negocio."""
import pandas as pd
df = pd.DataFrame(orders)
df['amount'] = df['amount'].round(2)
df['order_date'] = pd.to_datetime(df['order_date']).dt.date
df['revenue_category'] = pd.cut(
df['amount'],
bins=[0, 100, 500, float('inf')],
labels=['small', 'medium', 'large']
)
return df.to_dict(orient='records')
@task(task_id='load_to_warehouse')
def load_to_warehouse(orders: list[dict], **context) -> dict:
"""Carga las órdenes en el Data Warehouse."""
from warehouse import DataWarehouse
dw = DataWarehouse()
result = dw.upsert_orders(orders)
logger.info(f"Loaded {result['inserted']} new, {result['updated']} updated orders")
return result
@task(task_id='trigger_dbt_run', trigger_rule='all_success')
def trigger_dbt_run(**context) -> str:
"""Dispara los modelos dbt después de la carga."""
import subprocess
result = subprocess.run(
["dbt", "run", "--models", "marts.finance", "--profiles-dir", "/opt/airflow/dbt"],
capture_output=True, text=True
)
if result.returncode != 0:
raise Exception(f"dbt failed: {result.stderr}")
return result.stdout
@task(task_id='send_success_notification', trigger_rule='all_success')
def send_notification(load_result: dict, **context) -> None:
"""Notifica el éxito del pipeline."""
from slack_sdk import WebClient
client = WebClient(token="{{ var.value.slack_token }}")
client.chat_postMessage(
channel="#data-alerts",
text=f"✅ Orders ETL completed: {load_result['inserted']} new orders loaded"
)
# ── DEFINIR EL FLUJO ──────────────────────────────────────────────────────
raw_orders = extract_orders()
valid_orders = validate_orders(raw_orders)
transformed_orders = transform_orders(valid_orders)
load_result = load_to_warehouse(transformed_orders)
dbt_run = trigger_dbt_run()
notification = send_notification(load_result)
# load_result → ambos en paralelo
load_result >> [dbt_run, notification]
# Instanciar el DAG
pipeline = orders_etl_pipeline()
Operators y Sensors
"""
Los Operators más usados en Data Engineering
"""
from airflow.operators.python import PythonOperator, BranchPythonOperator
from airflow.operators.bash import BashOperator
from airflow.operators.empty import EmptyOperator
from airflow.providers.amazon.aws.operators.glue import GlueJobOperator
from airflow.providers.google.cloud.operators.bigquery import BigQueryInsertJobOperator
from airflow.providers.databricks.operators.databricks import DatabricksRunNowOperator
from airflow.providers.amazon.aws.sensors.s3 import S3KeySensor
from airflow.providers.http.sensors.http import HttpSensor
from airflow.sensors.time_sensor import TimeSensor
from airflow.sensors.external_task import ExternalTaskSensor
with DAG('operators_demo', schedule='@daily', start_date=datetime(2026,1,1), catchup=False) as dag:
# ── PYTHON OPERATOR ────────────────────────────────────────────────────
extract_task = PythonOperator(
task_id='extract_data',
python_callable=extract_orders_function,
op_kwargs={'date': '{{ ds }}', 'limit': 1000}, # Jinja templates
)
# ── BASH OPERATOR ─────────────────────────────────────────────────────
run_dbt = BashOperator(
task_id='run_dbt_models',
bash_command="""
cd /opt/dbt/project &&
dbt run --models marts.finance --vars '{"run_date": "{{ ds }}"}'
""",
env={'DBT_PROFILES_DIR': '/opt/dbt/profiles'},
)
# ── BRANCH: Decidir qué rama ejecutar ─────────────────────────────────
def choose_branch(**context):
logical_date = context['logical_date']
if logical_date.weekday() == 0: # Lunes
return 'full_refresh_task'
return 'incremental_task'
branching = BranchPythonOperator(
task_id='decide_refresh_strategy',
python_callable=choose_branch,
)
full_refresh = EmptyOperator(task_id='full_refresh_task')
incremental = EmptyOperator(task_id='incremental_task')
join = EmptyOperator(task_id='join_branches', trigger_rule='none_failed_min_one_success')
branching >> [full_refresh, incremental] >> join
# ── S3 SENSOR: Esperar que un archivo esté disponible ─────────────────
wait_for_file = S3KeySensor(
task_id='wait_for_upstream_file',
bucket_name='my-data-lake',
bucket_key='incoming/orders/{{ ds }}/data.parquet',
aws_conn_id='aws_default',
timeout=60 * 60 * 4, # Timeout: 4 horas
poke_interval=60, # Verificar cada 60 segundos
mode='reschedule', # Libera el worker mientras espera (más eficiente)
)
# ── EXTERNAL TASK SENSOR: Esperar otro DAG ────────────────────────────
wait_for_upstream_dag = ExternalTaskSensor(
task_id='wait_for_customers_dag',
external_dag_id='customers_etl_pipeline',
external_task_id='load_to_warehouse', # Task específica
allowed_states=['success'],
execution_date_fn=lambda dt: dt, # Same logical_date
timeout=3600,
poke_interval=60,
mode='reschedule',
)
# ── DATABRICKS OPERATOR ────────────────────────────────────────────────
spark_job = DatabricksRunNowOperator(
task_id='run_spark_job',
databricks_conn_id='databricks_default',
job_id='{{ var.value.spark_orders_job_id }}',
notebook_params={'run_date': '{{ ds }}', 'mode': 'incremental'},
)
# ── BIGQUERY OPERATOR ──────────────────────────────────────────────────
bq_query = BigQueryInsertJobOperator(
task_id='run_bq_aggregation',
configuration={
"query": {
"query": """
SELECT customer_id, SUM(amount) as total
FROM `project.dataset.orders`
WHERE DATE(created_at) = '{{ ds }}'
GROUP BY customer_id
""",
"useLegacySql": False,
"destinationTable": {
"projectId": "my-project",
"datasetId": "analytics",
"tableId": "customer_daily_revenue",
},
"writeDisposition": "WRITE_TRUNCATE"
}
},
gcp_conn_id='google_cloud_default',
)
# Dependencias entre tasks
wait_for_file >> wait_for_upstream_dag >> extract_task >> run_dbt >> spark_job >> bq_query
XComs y TaskGroups
"""
XComs: comunicación entre Tasks
TaskGroups: organizar Tasks en grupos visuales en la UI
"""
from airflow.utils.task_group import TaskGroup
from airflow.models import XCom
# ── XCOMS: Pasar datos entre Tasks ─────────────────────────────────────────
# Método 1: TaskFlow API (automático, recomendado)
@task
def get_count() -> int:
return 42 # Se guarda automáticamente en XCom
@task
def process_data(count: int) -> None:
print(f"Processing {count} records") # Recibe el XCom automáticamente
count = get_count()
process_data(count)
# Método 2: xcom_push / xcom_pull (API clásica)
def push_data(**context):
context['ti'].xcom_push(key='file_path', value='s3://bucket/file.parquet')
def pull_data(**context):
file_path = context['ti'].xcom_pull(task_ids='push_task', key='file_path')
print(f"File: {file_path}")
# ⚠️ CUIDADO con XComs grandes:
# XComs se almacenan en la metadata DB (PostgreSQL) → solo para datos pequeños!
# Para DataFrames grandes: guardar en S3/GCS y pasar solo la URL
# ── TASKGROUPS: Organizar visualmente el DAG ───────────────────────────────
with DAG('etl_with_task_groups', schedule='@daily', start_date=datetime(2026,1,1), catchup=False) as dag:
start = EmptyOperator(task_id='start')
end = EmptyOperator(task_id='end', trigger_rule='all_success')
with TaskGroup('extraction', tooltip='Extract from all sources') as extract_group:
extract_orders = PythonOperator(task_id='extract_orders', python_callable=get_orders)
extract_customers = PythonOperator(task_id='extract_customers', python_callable=get_customers)
extract_products = PythonOperator(task_id='extract_products', python_callable=get_products)
# Corren en paralelo dentro del grupo
with TaskGroup('transformation', tooltip='Clean and transform data') as transform_group:
clean_orders = PythonOperator(task_id='clean_orders', python_callable=clean_orders_fn)
clean_customers = PythonOperator(task_id='clean_customers', python_callable=clean_customers_fn)
join_data = PythonOperator(task_id='join_all', python_callable=join_fn)
[clean_orders, clean_customers] >> join_data
with TaskGroup('loading', tooltip='Load to Data Warehouse') as load_group:
load_silver = PythonOperator(task_id='load_silver', python_callable=load_silver_fn)
load_gold = BashOperator(task_id='run_dbt', bash_command='dbt run --models marts')
load_silver >> load_gold
start >> extract_group >> transform_group >> load_group >> end
Mejores Prácticas de Airflow
✅ Reglas de Oro para DAGs en Producción
- Idempotencia: Re-ejecutar una Task con los mismos parámetros debe dar el mismo resultado. Usa
OVERWRITEen lugar deAPPENDcuando sea posible. - Tasks pequeñas y atómicas: Cada Task debe hacer UNA cosa. Facilita retries y debugging.
- No poner lógica en el nivel del DAG: Las importaciones pesadas y la lógica de negocio deben estar dentro de las funciones de las Tasks, no en el nivel del módulo (el Scheduler lee el DAG frecuentemente).
- Usar Connections y Variables: Nunca hardcodear credenciales o configuraciones en el DAG. Usar la metadata DB de Airflow para esto.
- catchup=False: Para DAGs nuevos, a menos que necesites realmente backfill.
- max_active_runs: Limitar a 1-3 para evitar condiciones de carrera en los datos.
- SLA Monitoring: Configurar SLAs con email/Slack alerts para DAGs críticos.
# ── VARIABLES Y CONNECTIONS (buenas prácticas) ────────────────────────────
from airflow.models import Variable
from airflow.hooks.base import BaseHook
# Variables: configuración (accesible desde la UI)
api_key = Variable.get('orders_api_key', default_var='dev-key')
batch_size = int(Variable.get('etl_batch_size', default_var=1000))
# Con deserialización JSON automática
config = Variable.get('pipeline_config', deserialize_json=True)
# En la UI: {"batch_size": 5000, "timeout": 3600, "retry": true}
# Connections: credenciales de sistemas externos
from airflow.providers.postgres.hooks.postgres import PostgresHook
# En la UI configuras: host, port, user, password, schema
# En el código:
pg_hook = PostgresHook(postgres_conn_id='warehouse_production')
df = pg_hook.get_pandas_df("SELECT * FROM orders WHERE date = %(date)s", parameters={'date': '2026-01-01'})
# ── TEMPLATES Y MACROS ────────────────────────────────────────────────────
# Airflow tiene Jinja templating para la mayoría de argumentos de Operators
# Variables disponibles:
# {{ ds }} → YYYY-MM-DD de la ejecución lógica
# {{ ds_nodash }} → YYYYMMDD (sin guiones)
# {{ ts }} → ISO 8601 timestamp
# {{ macros.ds_add(ds, 7) }} → ds + 7 días
# {{ var.value.my_var }} → Variable de Airflow
# {{ conn.my_conn.host }} → Atributo de una Connection
# {{ logical_date }} → datetime object
bash_with_template = BashOperator(
task_id='copy_daily_file',
bash_command='aws s3 cp s3://source/{{ ds }}/data.csv s3://dest/{{ ds }}/data.csv',
)
Alternativas Modernas a Airflow
| Herramienta | Enfoque | Pros | Contras | Cuándo Usar |
|---|---|---|---|---|
| Apache Airflow | DAG-based, Python | Maduro, enorme ecosistema, UI potente | Complejo de operar, curva de aprendizaje | El estándar. Casi siempre. |
| Prefect | Python-first, flows | API más limpia, más pythónico, cloud managed | Menor ecosistema de providers | Startups, equipo Python-heavy |
| Dagster | Asset-based, software-defined assets | Lineage nativo, testing, observability | Curva de aprendizaje, diferente paradigma | Data Platform teams maduros |
| Mage | Pipelines interactivos | UI con notebooks, fácil de empezar | Menos maduro para producción | Prototipado rápido, equipos pequeños |
| dbt + Airflow | SQL transformations + scheduling | El stack más común en 2026 | Dos herramientas que aprender | Analytics Engineering stack |