🌀 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.

⏱️ Lectura: ~60 min
🎯 Nivel: Intermedio → Avanzado
📌 Airflow: 2.9+ (Taskflow API)

Conceptos Core de Airflow

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 OVERWRITE en lugar de APPEND cuando 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

HerramientaEnfoqueProsContrasCuá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