Skip to content
Volver al blog
Python
Airflow
Data Engineering
ETL

Construyendo Pipelines Escalables con Python

Patrones de diseño y herramientas probadas en producción para crear pipelines de datos que escalen sin romper el banco ni el equipo.

15 de mayo de 20258 min de lectura

Cuando empecé a diseñar pipelines de datos en producción, mi mayor error fue pensar que escalar era solo cuestión de agregar más recursos. Spoiler: no lo es.

Arquitectura de un pipeline escalable con Python y Airflow

El problema real del escalado

Un pipeline que procesa 1 GB funciona diferente que uno que procesa 1 TB. No es solo velocidad — es cómo manejas fallos, reintentos, dependencias y observabilidad cuando el volumen crece diez veces de la noche a la mañana.

# Anti-patrón: leer todo en memoria
def procesar_datos():
    df = pd.read_csv("archivo_enorme.csv")  # 💀 OOM cuando el archivo es grande
    return transformar(df)

# Patrón correcto: procesar en chunks
def procesar_datos_chunked(chunk_size: int = 10_000):
    for chunk in pd.read_csv("archivo_enorme.csv", chunksize=chunk_size):
        yield transformar(chunk)

Principios que aplico en producción

1. Idempotencia ante todo

Cada tarea de tu pipeline debe poder ejecutarse N veces con el mismo resultado. Esto no es opcional — es lo que permite que los reintentos automáticos de Airflow no corrompan tus datos.

def cargar_a_bigquery(dataset: str, tabla: str, fecha: str):
    """Usa WRITE_TRUNCATE para partición específica — idempotente por diseño."""
    job_config = bigquery.LoadJobConfig(
        write_disposition="WRITE_TRUNCATE",
        time_partitioning=bigquery.TimePartitioning(field="fecha"),
    )
    client.load_table_from_dataframe(df, f"{dataset}.{tabla}${fecha}", job_config=job_config)

2. Fail fast, retry smart

Los errores transitorios (timeouts de red, rate limits de API) se recuperan solos con reintentos. Los errores de datos requieren intervención humana. Distinguir entre ambos es crítico.

from airflow.decorators import task
from tenacity import retry, stop_after_attempt, wait_exponential

@retry(stop=stop_after_attempt(3), wait=wait_exponential(multiplier=1, min=4, max=10))
def llamar_api_externa(endpoint: str) -> dict:
    response = requests.get(endpoint, timeout=30)
    response.raise_for_status()
    return response.json()

3. Observabilidad como ciudadana de primera clase

Un pipeline sin métricas es un pipeline que eventualmente fallará en silencio. Instrumenta desde el inicio:

  • Tiempo de ejecución por etapa
  • Número de registros procesados vs esperados
  • Tasa de error por fuente de datos

Herramientas que uso en 2025

Logo de Apache Airflow

| Herramienta | Para qué | |-------------|----------| | Airflow | Orquestación y scheduling | | dbt | Transformaciones SQL versionadas | | Polars | Alternativa a Pandas para grandes volúmenes | | BigQuery | Data warehouse serverless en GCP | | Prefect | Orquestación alternativa con mejor DX |

Conclusión

La escalabilidad en pipelines de datos no es un problema de infraestructura — es un problema de diseño. Invierte tiempo en hacer cada etapa idempotente, en distinguir tipos de error y en instrumentar desde el día uno. El resto se puede agregar después.

¿Qué herramientas usas en tus pipelines? Me encantaría leer tu perspectiva.