← ⚙️ MapReduce y Spark
intermedio

3.3 · Apache Spark: procesamiento en memoria

⏱ 25 minMódulo 3: Procesamiento Distribuido

¿Qué es Apache Spark?

Apache Spark es un motor de procesamiento distribuido creado en UC Berkeley (AMPLab, 2009) como respuesta a las limitaciones de MapReduce. Su idea central: mantener los datos intermedios en memoria RAM en lugar de escribirlos a disco entre etapas. En cargas iterativas puede ser de 10× a 100× más rápido que Hadoop MapReduce.

Spark corre sobre YARN, Mesos, Kubernetes o en modo standalone, y lee datos de HDFS, S3, bases de datos y más.

RDD vs DataFrame vs Dataset

Spark ofrece tres abstracciones de datos, de menor a mayor nivel:

AbstracciónQué esVentajasDesventajas
RDDColección distribuida de objetos, particionada y tolerante a fallosControl total, bajo nivelSin optimizador; serialización costosa
DataFrameTabla distribuida con esquema (filas y columnas con nombre)Optimizador Catalyst, API concisa, multi-lenguajeMenos tipado en Python
DatasetDataFrame con tipado fuerte (solo Scala/Java)Seguridad de tipos en compilaciónNo existe en Python

El RDD (Resilient Distributed Dataset) es la base interna: todo en Spark se traduce a operaciones sobre RDDs. Pero en la práctica se recomienda usar DataFrames, porque el optimizador Catalyst reescribe y acelera las consultas automáticamente.

En PySpark solo trabajas con RDDs y DataFrames: los Datasets tipados requieren un lenguaje compilado (Scala o Java).

Transformaciones (lazy) vs acciones

Las operaciones de Spark se dividen en dos categorías:

  • Transformaciones (filter, map, groupBy, join): describen qué calcular, pero no ejecutan nada. Son lazy (perezosas): solo construyen un plan.
  • Acciones (collect, count, show, write): disparan la ejecución real del plan en el clúster.

Esta evaluación perezosa permite a Spark optimizar el plan completo antes de ejecutar: por ejemplo, empujar filtros cerca de la fuente de datos o elegir el mejor orden de joins.

El DAG de ejecución

Cuando se invoca una acción, Spark construye un DAG (grafo acíclico dirigido) de todas las transformaciones, lo divide en stages (etapas) separadas por operaciones de shuffle, y cada stage en tasks que corren en paralelo sobre las particiones del dataset.

graph TD
subgraph Plan["DAG del job"]
  A["textFile: leer log"] --> B["filter: solo ERROR"]
  B --> C["map: extraer codigo"]
  C --> D["groupBy: contar por codigo"]
  D --> E["collect: accion"]
end
subgraph Stage1["Stage 1 (sin shuffle)"]
  T1["task 1: particion 1<br/>filter + map"]
  T2["task 2: particion 2<br/>filter + map"]
  T3["task 3: particion 3<br/>filter + map"]
end
subgraph Stage2["Stage 2 (tras shuffle)"]
  T4["task 4: reduce por clave"]
  T5["task 5: reduce por clave"]
end
A -.-> T1
A -.-> T2
A -.-> T3
T1 --> SH["shuffle por red"]
T2 --> SH
T3 --> SH
SH --> T4
SH --> T5
Un DAG de Spark: transformaciones encadenadas se agrupan en stages separados por el shuffle; cada stage se ejecuta como tasks paralelas sobre particiones.

Las particiones son las unidades de paralelismo: un DataFrame de 200 particiones puede procesarse con hasta 200 tareas simultáneas en distintos nodos.

¿Por qué Spark es más rápido que MapReduce?

  • Memoria en lugar de disco: los resultados intermedios se cachean en RAM; MapReduce los escribe a HDFS tras cada etapa.
  • DAG en lugar de jobs encadenados: Spark ejecuta un pipeline completo en un solo job, sin materializar cada paso intermedio.
  • Optimizador Catalyst: analiza el plan lógico y lo reescribe (filtros tempranos, poda de columnas).
  • API rica: SQL, streaming (Structured Streaming), ML (MLlib) y grafos (GraphX) sobre el mismo motor.

PySpark básico

PySpark es la API de Python para Spark. Un programa típico crea una SparkSession, carga datos en un DataFrame y aplica transformaciones:

from pyspark.sql import SparkSession

spark = SparkSession.builder.appName("demo").getOrCreate()

df = spark.read.json("/datos/eventos.json")   # lazy: aún no lee nada
errores = df.filter(df.nivel == "ERROR")       # lazy
conteo = errores.groupBy("servicio").count()   # lazy
conteo.show()                                   # acción: ejecuta el DAG

Ejercicio: pipeline con PySpark

🧪 Ejercicio

Análisis de ventas con PySpark

Dado un DataFrame `ventas` con columnas `region`, `producto` y `monto`, escribe un pipeline PySpark que: (1) filtre las ventas con monto mayor a 100, (2) cree una columna `monto_con_igv` (monto * 1.18), y (3) agrupe por región calculando el total vendido y el número de operaciones, ordenado de mayor a menor total. Termina mostrando el resultado.

Comprueba lo aprendido

Comprueba que lo pillaste

Ejecutas `df.filter(col('edad') > 30)` en PySpark y el comando termina al instante sin tocar los datos. ¿Por qué?

Comprueba que lo pillaste

¿Cuál es la principal razón por la que Spark supera en velocidad a Hadoop MapReduce en cargas iterativas?

Comprueba que lo pillaste

¿Qué ventaja ofrece un DataFrame frente a un RDD en Spark?

Resumen

  • Spark procesa datos distribuidos manteniendo resultados intermedios en memoria, siendo mucho más rápido que MapReduce.
  • Abstracciones: RDD (bajo nivel), DataFrame (tabla con esquema, recomendada) y Dataset (tipado, solo Scala/Java).
  • Las transformaciones son lazy y las acciones disparan la ejecución; esto permite optimizar el plan completo.
  • El motor construye un DAG dividido en stages por los shuffles, y ejecuta tasks paralelas sobre las particiones.
  • PySpark lleva Spark a Python con una API fluida: filter, withColumn, groupBy().agg(), show().

📚 Lecturas y fuentes

RecursoTipoPor qué leerlo
Apache Spark (Wikipedia en español)ArtículoPanorama en español: introducción y “Componentes” (Spark SQL, MLlib, GraphX) para ubicar el motor en el ecosistema.
Spark SQL and DataFrames (docs oficiales)DocsLa API que usarás en la práctica. Lee “Getting Started” y “Datasets and DataFrames”; salta las guías de Hive y JDBC. (en inglés)
RDD Programming Guide (docs oficiales)DocsPara entender qué hay debajo del DataFrame. Basta con “RDD Operations” (transformaciones vs. acciones) y “RDD Persistence” (cache). (en inglés)
Resilient Distributed Datasets (Zaharia et al., NSDI 2012)PaperEl paper que inventó el RDD. Secciones 2 (abstracción) y 3 (API y linaje para tolerancia a fallos); la 6 es evaluación y puedes saltarla. (en inglés)