¿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ón | Qué es | Ventajas | Desventajas |
|---|---|---|---|
| RDD | Colección distribuida de objetos, particionada y tolerante a fallos | Control total, bajo nivel | Sin optimizador; serialización costosa |
| DataFrame | Tabla distribuida con esquema (filas y columnas con nombre) | Optimizador Catalyst, API concisa, multi-lenguaje | Menos tipado en Python |
| Dataset | DataFrame con tipado fuerte (solo Scala/Java) | Seguridad de tipos en compilación | No 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
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.
🔍 Usa df.filter() o df.where() para filtrar, df.withColumn() para añadir columnas, y funciones de pyspark.sql.functions como sum() y count() dentro de agg().
from pyspark.sql import SparkSession
from pyspark.sql.functions import col, sum as _sum, count
spark = SparkSession.builder.appName("ventas").getOrCreate()
datos = [
("Norte", "laptop", 1200.0),
("Sur", "mouse", 25.0),
("Norte", "teclado", 150.0),
("Sur", "monitor", 800.0),
("Lima", "laptop", 1500.0),
("Lima", "mouse", 30.0),
]
ventas = spark.createDataFrame(datos, ["region", "producto", "monto"])
resultado = (
ventas
.filter(col("monto") > 100) # transformacion (lazy)
.withColumn("monto_con_igv", col("monto") * 1.18) # transformacion (lazy)
.groupBy("region") # transformacion (lazy)
.agg(
_sum("monto_con_igv").alias("total_ventas"),
count("*").alias("num_operaciones"),
)
.orderBy(col("total_ventas").desc())
)
resultado.show() # accion: aqui se ejecuta todo el DAG
spark.stop()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é?
Las transformaciones en Spark son perezosas: filter, map y groupBy solo añaden pasos al plan lógico. Nada se ejecuta hasta llamar a una acción como show(), count() o collect(), momento en que se construye y ejecuta el DAG.
Comprueba que lo pillaste
¿Cuál es la principal razón por la que Spark supera en velocidad a Hadoop MapReduce en cargas iterativas?
MapReduce persiste los resultados intermedios en HDFS después de cada job. Spark los mantiene en memoria y ejecuta el pipeline completo como un DAG, eliminando el costo dominante de I/O a disco en algoritmos iterativos.
Comprueba que lo pillaste
¿Qué ventaja ofrece un DataFrame frente a un RDD en Spark?
Al conocer el esquema (columnas y tipos), el optimizador Catalyst puede aplicar técnicas como empujar filtros a la fuente o podar columnas innecesarias. Con RDDs, Spark solo ve objetos opacos y no puede optimizar.
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
| Recurso | Tipo | Por qué leerlo |
|---|---|---|
| Apache Spark (Wikipedia en español) | Artículo | Panorama en español: introducción y “Componentes” (Spark SQL, MLlib, GraphX) para ubicar el motor en el ecosistema. |
| Spark SQL and DataFrames (docs oficiales) | Docs | La 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) | Docs | Para 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) | Paper | El 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) |