🎯 Objetivo
Resolver el mismo WordCount dos veces —con RDDs y con DataFrames— sobre un texto real, y leer en la Spark UI los stages, el shuffle y el DAG que ya viste en teoría.
✅ Antes de empezar
Spark corre sobre la JVM, así que necesitas Java aunque escribas Python. Comprueba tu entorno antes de instalar nada:
java -version # necesitas Java 17 o 21
python --version # necesitas Python 3.10 o superior
Si Java responde, instala PySpark:
pip install pyspark
Windows: si
java -versionfunciona pero PySpark falla al arrancar, defineJAVA_HOMEapuntando a la carpeta del JDK (no abin). También verás un warning ruidoso sobrewinutils.exeyHADOOP_HOME: en modo local es inofensivo, ignóralo.
¿No quieres instalar nada? Abre un cuaderno en Google Colab y pon
!pip install pysparken la primera celda; todo el lab funciona igual, salvo la Spark UI del Paso 5.
Paso 1: Crear la SparkSession
La SparkSession es tu puerta de entrada a Spark: crea el driver, reserva los executors y abre la Spark UI. Con master("local[*]") le dices que use tu máquina como si fuera un clúster de un solo nodo, con tantos hilos como cores tengas.
from pyspark.sql import SparkSession
spark = (
SparkSession.builder
.master("local[*]")
.appName("lab3")
.getOrCreate()
)
sc = spark.sparkContext
sc.setLogLevel("WARN") # menos ruido en consola
print("Version de Spark:", spark.version)
print("Paralelismo por defecto (cores):", sc.defaultParallelism)
Paso 2: Descargar el texto y cargarlo como RDD
Necesitamos un archivo lo bastante grande para que el paralelismo se note. Usamos El Quijote desde Project Gutenberg (dominio público, ~2 MB). sc.textFile crea un RDD donde cada elemento es una línea del archivo.
curl -L -o quijote.txt https://www.gutenberg.org/files/2000/2000-0.txt
Si esa URL falla, entra a https://www.gutenberg.org/ebooks/2000 y descarga la versión Plain Text UTF-8. También sirve cualquier .txt grande que ya tengas.
lineas = sc.textFile("quijote.txt")
print("Numero de lineas:", lineas.count())
print("Particiones del RDD:", lineas.getNumPartitions())
for linea in lineas.take(3):
print(repr(linea))
Paso 3: WordCount con RDDs
Este es el “hola mundo” de los macrodatos y el mismo algoritmo de la lección de MapReduce. La fase Map es todo lo que ocurre por línea sin mirar a las demás (flatMap, filter, map); la fase Reduce es reduceByKey, que necesita juntar todas las apariciones de una palabra en el mismo nodo.
import re
conteo = (
lineas
.flatMap(lambda l: re.findall(r"[a-záéíóúñü]+", l.lower())) # MAP: linea -> palabras
.filter(lambda p: len(p) > 3) # MAP: descartar palabras cortas
.map(lambda p: (p, 1)) # MAP: emitir par (clave, 1)
.reduceByKey(lambda a, b: a + b) # REDUCE: sumar por clave
)
top10 = conteo.takeOrdered(10, key=lambda par: -par[1]) # accion
for palabra, veces in top10:
print(palabra, veces)
Paso 4: El mismo WordCount con DataFrames
Ahora el mismo resultado, pero declarando qué quieres en vez de cómo calcularlo. explode convierte cada arreglo de palabras en varias filas, igual que hacía flatMap, y groupBy().count() reemplaza al reduceByKey.
from pyspark.sql.functions import col, explode, split, lower, length, desc
df = spark.read.text("quijote.txt") # una columna llamada "value"
palabras = (
df
.select(explode(split(lower(col("value")), r"[^a-záéíóúñü]+")).alias("palabra"))
.filter(length(col("palabra")) > 3)
)
top10_df = palabras.groupBy("palabra").count().orderBy(desc("count"))
top10_df.show(10, truncate=False)
graph LR A["textFile<br/>lineas"] --> B["flatMap<br/>palabras"] B --> C["map<br/>palabra, 1"] C -.->|shuffle| D["reduceByKey<br/>suma por palabra"] D --> E["takeOrdered<br/>top 10"]
Paso 5: Leer la Spark UI
Mientras la sesión siga viva, Spark expone un panel web con todo lo que ejecutó. Es la herramienta número uno para diagnosticar jobs lentos, y aquí vas a ver el diagrama de arriba dibujado con tus datos reales.
# Manten la sesion abierta mientras exploras la UI
input("Abre http://localhost:4040 y pulsa Enter para continuar...")
Abre http://localhost:4040 y recorre estas pestañas en orden:
- Jobs: verás un job por cada acción que ejecutaste (
count,takeOrdered,show). Fíjate en la duración de cada uno. - Stages: entra al job del
takeOrdered. Tiene dos stages, no uno. La columna Shuffle Write del primero y Shuffle Read del segundo te dicen cuántos bytes cruzaron la red. - DAG Visualization: dentro del job, despliega este panel. La línea que separa los dos bloques es exactamente el
reduceByKey.
Paso 6: Comparar los planes con explain
Un DataFrame no ejecuta lo que escribiste, sino lo que Catalyst decidió que era equivalente y más barato. explain(True) te muestra los cuatro planes: analizado, optimizado y físico.
top10_df.explain(True)
# El RDD no tiene plan optimizable, solo su linaje de dependencias
print(conteo.toDebugString().decode())
🧪 Reto final
🧪 Ejercicio
Top 10 de palabras largas sin stopwords
Sobre el mismo quijote.txt y usando DataFrames, obtén el top 10 de palabras de MÁS de 6 letras, excluyendo una lista de stopwords en español. Muestra el resultado ordenado de mayor a menor frecuencia y termina cerrando la sesión.
🔍 Reutiliza split + explode del Paso 4. Para excluir stopwords usa ~col('palabra').isin(lista) dentro de un filter, y para el largo usa length(col('palabra')) > 6.
from pyspark.sql import SparkSession
from pyspark.sql.functions import col, explode, split, lower, length, desc
spark = (
SparkSession.builder
.master("local[*]")
.appName("reto-lab3")
.getOrCreate()
)
STOPWORDS = [
"aquellos", "aquellas", "nosotros", "vosotros", "aquella",
"tambien", "porque", "cuando", "aunque", "mientras",
"entonces", "despues", "siempre", "todavia", "ninguno",
]
df = spark.read.text("quijote.txt")
resultado = (
df
.select(explode(split(lower(col("value")), r"[^a-záéíóúñü]+")).alias("palabra"))
.filter(length(col("palabra")) > 6)
.filter(~col("palabra").isin(STOPWORDS))
.groupBy("palabra")
.count()
.orderBy(desc("count"))
)
resultado.show(10, truncate=False)
# Version imprimible sin f-strings
for fila in resultado.take(10):
print(fila["palabra"], "->", fila["count"])
spark.stop()Comprueba lo aprendido
Comprueba que lo pillaste
En el Paso 3 encadenaste flatMap, filter, map y reduceByKey, y la celda terminó al instante. ¿En qué momento se leyó realmente el archivo?
flatMap, filter, map y reduceByKey son transformaciones lazy: solo añaden pasos al DAG. Recién con takeOrdered, que es una acción, Spark planifica, divide en stages y lee el archivo.
Comprueba que lo pillaste
En la pestaña Stages viste que el job del WordCount tiene dos stages en lugar de uno. ¿Qué operación marca esa frontera?
Un stage termina donde empieza un shuffle. flatMap, filter y map son operaciones estrechas (narrow): cada partición se procesa sola. reduceByKey es ancha (wide): reagrupa por clave a través de la red y obliga a abrir un stage nuevo.
Comprueba que lo pillaste
El Paso 6 mostró un plan optimizado para el DataFrame pero solo un linaje de RDDs para la versión con lambdas. ¿Por qué?
Catalyst optimiza porque entiende qué hace cada expresión: puede podar columnas, empujar filtros o reordenar operaciones. Ante una lambda arbitraria de Python no puede razonar sobre su contenido, así que la ejecuta tal cual.
🧹 Limpieza
Cierra la sesión al terminar: libera los executors, la memoria reservada y el puerto 4040. Si no lo haces, la siguiente sesión arrancará en el puerto 4041 y te confundirás de UI.
spark.stop()
📚 Lecturas y fuentes
| Recurso | Tipo | Por qué leerlo |
|---|---|---|
| RDD Programming Guide | Documentación oficial | Explica el catálogo completo de transformaciones y acciones sobre RDDs, y la diferencia entre operaciones narrow y wide |
| Spark SQL, DataFrames and Datasets Guide | Documentación oficial | La API que usarás el 95% del tiempo: lectura de fuentes, funciones de columna y cómo actúa Catalyst |
| Web UI | Documentación oficial | Detalla cada pestaña del panel de localhost:4040 y qué métrica mirar para diagnosticar un job lento |