← ⚙️ MapReduce y Spark
intermedio🧪 Lab práctico

3.4 · Lab 3: Tu primer job de Spark en local

⏱ 60 minMódulo 3: Procesamiento Distribuido

🎯 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 -version funciona pero PySpark falla al arrancar, define JAVA_HOME apuntando a la carpeta del JDK (no a bin). También verás un warning ruidoso sobre winutils.exe y HADOOP_HOME: en modo local es inofensivo, ignóralo.

¿No quieres instalar nada? Abre un cuaderno en Google Colab y pon !pip install pyspark en 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"]
El shuffle es la frontera entre stages: todo lo anterior se ejecuta sin mover datos entre nodos.

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:

  1. Jobs: verás un job por cada acción que ejecutaste (count, takeOrdered, show). Fíjate en la duración de cada uno.
  2. 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.
  3. 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.

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?

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?

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é?

🧹 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

RecursoTipoPor qué leerlo
RDD Programming GuideDocumentación oficialExplica el catálogo completo de transformaciones y acciones sobre RDDs, y la diferencia entre operaciones narrow y wide
Spark SQL, DataFrames and Datasets GuideDocumentación oficialLa API que usarás el 95% del tiempo: lectura de fuentes, funciones de columna y cómo actúa Catalyst
Web UIDocumentación oficialDetalla cada pestaña del panel de localhost:4040 y qué métrica mirar para diagnosticar un job lento