Apache Spark: Procesamiento de Big Data desde Cero

Featured image for article 60

Apache Spark es un motor distribuido para procesar datos cuando una sola máquina deja de ser suficiente o cuando un flujo necesita repartir trabajo entre varios nodos. Puede ejecutar transformaciones por lotes, consultas SQL, procesamiento continuo y tareas de machine learning dentro de un mismo ecosistema.

La documentación vigente de Apache presenta DataFrames y Spark SQL como las APIs modernas para datos estructurados, mientras que los RDD siguen disponibles como una abstracción de bajo nivel. Esta guía explica cómo pensar en Spark, no solo cómo copiar un comando.

Cuándo tiene sentido usar Spark

Spark resulta apropiado cuando el volumen, la velocidad o el coste de procesamiento justifican distribuir el trabajo. Ejemplos frecuentes son transformar archivos históricos de gran tamaño, combinar eventos de muchas fuentes, preparar tablas analíticas o calcular agregaciones que no caben cómodamente en memoria.

No es la respuesta automática para cualquier dataset. Para un CSV de unos pocos megabytes, pandas, DuckDB o una consulta SQL suelen ser más simples. La distribución añade planificación, transferencia de datos, serialización, registros y operación de clústeres.

Arquitectura: driver, ejecutores y tareas

Una aplicación Spark tiene un driver que construye el plan de ejecución y coordina el trabajo. Los executors ejecutan tareas sobre particiones de datos y conservan resultados intermedios cuando corresponde. El gestor del clúster asigna los recursos; Spark puede desplegarse en modo standalone, YARN o Kubernetes.

Una transformación como filter o select es perezosa: Spark registra la operación, pero no procesa los datos hasta que una acción como count, write o collect exige un resultado. Esto permite optimizar el plan completo antes de ejecutarlo.

DataFrames, Spark SQL y RDD

API Uso recomendado Consideración
DataFrame Transformaciones estructuradas en Python, Scala o Java Permite al optimizador entender columnas y tipos
Spark SQL Equipos que trabajan con consultas y modelos tabulares Se integra con catálogos, vistas y funciones SQL
RDD Lógica de bajo nivel que no encaja en la API estructurada Ofrece menos oportunidades de optimización automática

Primer ejemplo con PySpark

from pyspark.sql import SparkSession
from pyspark.sql import functions as F

spark = SparkSession.builder.appName("ventas-mensuales").getOrCreate()

ventas = (
    spark.read
    .option("header", True)
    .option("inferSchema", True)
    .csv("datos/ventas/*.csv")
)

resumen = (
    ventas
    .filter(F.col("estado") == "confirmada")
    .withColumn("mes", F.date_trunc("month", F.col("fecha")))
    .groupBy("mes", "pais")
    .agg(
        F.sum("importe").alias("ventas"),
        F.countDistinct("cliente_id").alias("clientes")
    )
)

resumen.write.mode("overwrite").parquet("salidas/ventas_mensuales")

El ejemplo lee varios archivos, filtra registros, deriva el mes, agrega dos métricas y escribe Parquet. En producción conviene declarar el esquema en lugar de inferirlo: evita lecturas extra y detecta cambios de tipo antes.

Particiones y operaciones costosas

Una partición es la unidad que Spark asigna a una tarea. Muy pocas particiones dejan recursos sin utilizar; demasiadas crean tareas pequeñas con sobrecarga. El tamaño adecuado depende del formato, el clúster y la operación.

Los shuffles redistribuyen datos entre ejecutores. Aparecen en agregaciones, ordenaciones y muchos joins. Son costosos porque implican red, disco y memoria. Antes de aumentar máquinas, revisa el plan con explain(), filtra pronto, selecciona solo columnas necesarias y evita claves extremadamente desbalanceadas.

Joins sin sorpresas

Si una tabla es pequeña, un broadcast join puede enviarla a cada ejecutor y evitar redistribuir la tabla grande. No fuerces esta estrategia sin medir el tamaño real: una tabla que no cabe en la memoria de los ejecutores puede provocar fallos.

clientes = spark.read.parquet("datos/clientes")
paises = spark.read.parquet("datos/catalogo_paises")

resultado = clientes.join(F.broadcast(paises), "pais_id", "left")

Procesamiento continuo con Structured Streaming

Structured Streaming expresa un flujo como una tabla que recibe filas. Puedes aplicar operaciones de DataFrame y escribir resultados por microbatches o mediante los modos compatibles con la fuente y el destino. Para datos con retraso, define marcas de agua y una política explícita de duplicados; de lo contrario, el estado puede crecer sin límite.

Errores habituales al empezar

  • Usar collect() sobre un resultado grande y saturar la memoria del driver.
  • Crear funciones Python fila por fila cuando existe una función nativa de Spark.
  • Guardar miles de archivos diminutos que luego ralentizan las lecturas.
  • Persistir cada DataFrame aunque se utilice una sola vez.
  • Ignorar el sesgo de claves en joins y agregaciones.
  • Confundir más ejecutores con menor coste total sin medir tiempo y consumo.

Ruta de aprendizaje recomendada

  1. Practica DataFrames localmente con un conjunto pequeño.
  2. Aprende a leer planes lógicos y físicos con explain().
  3. Compara Parquet con CSV y observa el efecto del esquema y la selección de columnas.
  4. Prueba joins, agregaciones y reparticionado con datos sesgados.
  5. Solo después pasa a despliegue, monitoreo y ajuste del clúster.

Consulta la documentación oficial de Apache Spark para verificar versiones, plataformas y APIs compatibles. Si primero necesitas diseñar el recorrido completo de los datos, revisa también la guía de data pipelines.

Preguntas frecuentes

¿Spark reemplaza una base de datos?

No. Spark procesa datos; normalmente lee y escribe en sistemas de almacenamiento, catálogos, lagos o bases de datos.

¿PySpark es más lento que Scala?

Las operaciones nativas de DataFrame se ejecutan en el motor de Spark. La penalización aparece sobre todo cuando se mueve lógica fila a fila al proceso Python.

¿Necesito un clúster para aprender?

No. El modo local permite practicar APIs y planes. Un clúster se vuelve necesario para estudiar distribución, fallos y rendimiento real.

Subir