← 🤖 ML a Escala
avanzado🧪 Lab práctico

5.4 · Lab 5: Construye un recomendador con Spark MLlib

⏱ 65 minMódulo 5: Analítica y ML a Escala

🎯 Objetivo

Entrenar un recomendador de películas con ALS de Spark MLlib sobre datos reales de MovieLens, medir su error con RMSE y mirar con tus propios ojos si las recomendaciones tienen sentido.

✅ Antes de empezar

Necesitas tres cosas. Ninguna tarda más de cinco minutos.

1. Java 17 o superior. Spark corre sobre la JVM. Comprueba:

java -version

Si no ves nada o ves una versión anterior a 17, instala un JDK (Temurin, Zulu o el de tu gestor de paquetes).

2. PySpark.

pip install pyspark

Si ya hiciste el Lab 3 del módulo 3, PySpark ya está instalado y no tienes que tocar nada.

3. El dataset. Lo descargamos desde Python en el Paso 1, así que de momento no hagas nada.

¿Problemas con Java en tu máquina? Abre un cuaderno en Google Colab, ejecuta !pip install pyspark en la primera celda y sigue el lab tal cual: Colab ya trae una JVM compatible. Todo el código de esta lección funciona igual allí.

Trabajaremos con MovieLens ml-latest-small: unas 100.000 valoraciones, 610 usuarios y 9.742 películas en apenas 1 MB. Es pequeño a propósito, para que el ciclo entrenar-evaluar-corregir sea rápido. Cuando quieras sentir la escala de verdad existe ml-25m, con 25 millones de valoraciones: el mismo código, otro mundo en tiempos de ejecución.

Paso 1: Descarga y descomprime el dataset

Lo hacemos con urllib y zipfile de la librería estándar en vez de wget o curl, para que funcione igual en Windows, Linux y Mac. El zip pesa alrededor de 1 MB, así que la descarga es instantánea.

# descargar_datos.py
import urllib.request
import zipfile
import os

URL = "https://files.grouplens.org/datasets/movielens/ml-latest-small.zip"
ZIP = "ml-latest-small.zip"

if not os.path.exists(ZIP):
    print("Descargando MovieLens ml-latest-small...")
    urllib.request.urlretrieve(URL, ZIP)

with zipfile.ZipFile(ZIP) as z:
    z.extractall(".")

print("Archivos extraidos:")
for nombre in sorted(os.listdir("ml-latest-small")):
    ruta = os.path.join("ml-latest-small", nombre)
    print(f"  {nombre:16s} {os.path.getsize(ruta) / 1024:8.1f} KB")

Ejecútalo:

python descargar_datos.py

Paso 2: Crea la SparkSession y carga los datos

Toda aplicación Spark arranca con una SparkSession. Le pedimos que infiera el esquema (inferSchema) para que rating llegue como número y no como texto: ALS exige columnas numéricas y falla si le pasas strings.

# recomendador.py
from pyspark.sql import SparkSession

spark = (SparkSession.builder
         .appName("RecomendadorMovieLens")
         .master("local[*]")
         .getOrCreate())

spark.sparkContext.setLogLevel("WARN")

ratings = (spark.read
           .option("header", "true")
           .option("inferSchema", "true")
           .csv("ml-latest-small/ratings.csv"))

movies = (spark.read
          .option("header", "true")
          .option("inferSchema", "true")
          .csv("ml-latest-small/movies.csv"))

ratings.printSchema()
movies.printSchema()

print("Valoraciones:", ratings.count())
print("Usuarios:    ", ratings.select("userId").distinct().count())
print("Peliculas:   ", movies.count())

ratings.show(5)

Paso 3: Explora y mide la esparsidad

Antes de modelar, mira los datos. Aquí está el corazón conceptual del filtrado colaborativo: la matriz usuario-película está casi vacía, y el trabajo del modelo es rellenar los huecos. Calcúlalo explícitamente en vez de creerlo de memoria.

from pyspark.sql import functions as F

# Distribucion de las valoraciones
print("--- Distribucion de ratings ---")
ratings.groupBy("rating").count().orderBy("rating").show()

# Peliculas con mas valoraciones (minimo 50 para que la media signifique algo)
print("--- Mas valoradas ---")
(ratings.groupBy("movieId")
        .agg(F.count("*").alias("n"), F.round(F.avg("rating"), 2).alias("media"))
        .filter(F.col("n") >= 50)
        .join(movies, "movieId")
        .orderBy(F.desc("n"))
        .select("title", "n", "media")
        .show(10, truncate=False))

# Esparsidad de la matriz usuario-pelicula
n_ratings = ratings.count()
n_users = ratings.select("userId").distinct().count()
n_movies = ratings.select("movieId").distinct().count()
celdas = n_users * n_movies
densidad = n_ratings / celdas

print("Celdas posibles:", celdas)
print("Celdas rellenas:", n_ratings)
print("Densidad: {:.2f}%".format(densidad * 100))
print("Esparsidad: {:.2f}%".format((1 - densidad) * 100))
graph LR
A["ratings.csv<br/>100k valoraciones"] --> B["Split<br/>80 / 20"]
B --> C["Entrenar ALS<br/>factorizacion matricial"]
B --> D["Test"]
C --> E["Predecir"]
D --> E
E --> F["RMSE"]
C --> G["Top 10 por usuario"]
G --> H["Join con movies.csv<br/>titulos legibles"]
El pipeline de ML a escala: los mismos pasos que en scikit-learn, pero cada uno distribuido.

Paso 4: Divide en entrenamiento y prueba

Un modelo evaluado con los mismos datos con los que aprendió siempre parece brillante. Reservamos un 20% que el modelo no verá nunca durante el entrenamiento. Fijamos la semilla para que la partición sea reproducible entre ejecuciones.

train, test = ratings.randomSplit([0.8, 0.2], seed=42)

print("Entrenamiento:", train.count())
print("Prueba:       ", test.count())

train.cache()  # lo vamos a recorrer varias veces al entrenar

Paso 5: Entrena el modelo ALS

ALS alterna mínimos cuadrados hasta encontrar los vectores latentes de usuarios y películas. Tres parámetros mandan: rank (cuántos factores latentes), maxIter (cuántas alternancias) y regParam (la regularización que evita el sobreajuste).

Pero el parámetro que arruina más labs es otro: coldStartStrategy. En el conjunto de prueba aparecerán usuarios o películas que el modelo no vio al entrenar, y sin vector latente su predicción es NaN. Como cualquier operación con NaN da NaN, el RMSE final sale NaN y parece que todo está roto. Con coldStartStrategy="drop" esas filas se descartan de las predicciones y la métrica vuelve a ser un número.

from pyspark.ml.recommendation import ALS

als = ALS(
    userCol="userId",
    itemCol="movieId",
    ratingCol="rating",
    rank=10,              # factores latentes
    maxIter=10,           # iteraciones de alternancia
    regParam=0.1,         # regularizacion
    coldStartStrategy="drop",   # sin esto el RMSE sale NaN
    nonnegative=True,
    seed=42,
)

modelo = als.fit(train)

print("Factores de usuario:", modelo.userFactors.count())
print("Factores de item:   ", modelo.itemFactors.count())
modelo.userFactors.show(3, truncate=80)

Paso 6: Evalúa con RMSE

RegressionEvaluator compara la columna prediction con el rating real y devuelve la raíz del error cuadrático medio. Es el número que te dice, en promedio, cuánto se equivoca el modelo en la misma escala que los ratings.

from pyspark.ml.evaluation import RegressionEvaluator

predicciones = modelo.transform(test)
predicciones.select("userId", "movieId", "rating", "prediction").show(10)

evaluador = RegressionEvaluator(
    metricName="rmse",
    labelCol="rating",
    predictionCol="prediction",
)

rmse = evaluador.evaluate(predicciones)
print("RMSE = {:.4f}".format(rmse))

# Baseline: predecir siempre la media global
media = train.agg(F.avg("rating")).first()[0]
base = test.withColumn("prediction", F.lit(media))
print("RMSE del baseline (media global) = {:.4f}".format(evaluador.evaluate(base)))

Si tu RMSE sale NaN, vuelve al Paso 5: te falta coldStartStrategy="drop".

Paso 7: Genera recomendaciones con títulos de verdad

Un RMSE es una abstracción. La prueba real es mirar la lista. recommendForAllUsers(10) devuelve el top 10 de cada usuario, pero como IDs numéricos, que no dicen nada. Hacemos explode del array y un join con movies.csv para leer títulos, y luego comparamos las recomendaciones de un usuario con lo que ese usuario ya puntuó alto.

USUARIO = 42

recomendaciones = modelo.recommendForAllUsers(10)
print("Usuarios con recomendaciones:", recomendaciones.count())

top = (recomendaciones
       .filter(F.col("userId") == USUARIO)
       .select(F.explode("recommendations").alias("r"))
       .select(F.col("r.movieId").alias("movieId"),
               F.col("r.rating").alias("prediccion"))
       .join(movies, "movieId")
       .orderBy(F.desc("prediccion")))

print("--- Recomendado para el usuario", USUARIO, "---")
top.select("title", "genres", F.round("prediccion", 2).alias("prediccion")).show(10, truncate=False)

print("--- Lo que el usuario", USUARIO, "ya puntuo alto ---")
(ratings.filter((F.col("userId") == USUARIO) & (F.col("rating") >= 4.5))
        .join(movies, "movieId")
        .select("title", "genres", "rating")
        .orderBy(F.desc("rating"))
        .show(10, truncate=False))

🧪 Reto final

🧪 Ejercicio

Búsqueda en rejilla para bajar el RMSE

El RMSE de 0,9 no es intocable. Haz una búsqueda en rejilla pequeña sobre rank (prueba 5, 10 y 20) y regParam (prueba 0.05, 0.1 y 0.2): 9 combinaciones. Entrena un ALS por combinación sobre train, evalúa cada uno con RMSE sobre test, imprime una tabla ordenada de peor a mejor y quédate con el modelo ganador. Usa siempre coldStartStrategy igual a drop y la misma semilla para que la comparación sea justa. Bonus: reentrena el ganador sobre el 100% de los datos antes de ponerlo en producción.

Comprueba lo aprendido

Comprueba que lo pillaste

Entrenas un ALS, evalúas con RegressionEvaluator y el RMSE sale NaN. ¿Cuál es la causa más probable?

Comprueba que lo pillaste

En este lab la matriz usuario-película tiene una densidad de aproximadamente el 1,7%. ¿Por qué esa esparsidad no impide que ALS funcione?

🧹 Limpieza

Cierra la sesión para liberar los ejecutores y el puerto de la interfaz web de Spark.

spark.stop()
print("Sesion cerrada")

🎓 Has terminado el curso

Enhorabuena: has cerrado el arco completo. Empezaste con los fundamentos (qué hace grande a un dato, OLTP frente a OLAP), bajaste al almacenamiento distribuido (HDFS, NoSQL, el teorema CAP), subiste al procesamiento (MapReduce, YARN, Spark), pasaste por el streaming (Kafka, tiempo real, Lambda y Kappa) y has acabado entrenando machine learning a escala sobre datos reales.

Tres siguientes pasos concretos:

  1. Monta un clúster de verdad en la nube (EMR, Dataproc o Databricks Community) y ejecuta este mismo script sin local[*], para ver el trabajo repartido entre nodos.
  2. Repite el pipeline con ml-25m: 25 millones de valoraciones. El código no cambia; los tiempos, las particiones y tu paciencia sí. Ahí es donde se entiende para qué sirve Spark.
  3. Certifícate: Databricks Certified Associate Developer for Apache Spark, o Google Cloud Professional Data Engineer. Ambas validan justo lo que has practicado aquí.

📚 Lecturas y fuentes

RecursoTipoPor qué leerlo
Collaborative Filtering — Spark MLlibDocsLa referencia oficial de ALS; lee “Cold-start strategy” y “Explicit vs. implicit feedback” para entender los parámetros que tocaste en el Paso 5. (en inglés)
MovieLens Datasets — GroupLensDatasetLa página del dataset, con el README que describe cada CSV y el enlace a ml-25m para repetir el lab a escala real. (en inglés)
Matrix Factorization Techniques for Recommender Systems (Koren, Bell y Volinsky, 2009)PaperExplica qué está haciendo ALS por dentro y por qué ganó el Netflix Prize; con “A Basic Matrix Factorization Model” te sobra. (en inglés)
ALS — API de PySparkDocsLista completa de parámetros de ALS en Python; útil para el reto final y para saber qué más hay además de rank y regParam. (en inglés)