← 🤖 ML a Escala
avanzado

5.1 · Machine Learning distribuido con Spark MLlib

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

El problema del machine learning a escala

Entrenar un modelo de machine learning con unos miles de registros es trivial: scikit-learn lo hace en memoria en segundos. Pero cuando el dataset tiene cientos de millones de filas (clicks de usuarios, transacciones, lecturas de sensores), ya no cabe en la memoria de una sola máquina y entrenar de forma secuencial puede tardar días.

Spark MLlib es la librería de machine learning distribuido de Apache Spark. Su idea central es ejecutar el algoritmo de entrenamiento en paralelo sobre el clúster: cada nodo procesa una partición de los datos y los gradientes (o estadísticos) se agregan entre nodos en cada iteración.

MLlib moderno trabaja sobre DataFrames (el paquete pyspark.ml), no sobre los antiguos RDD (pyspark.mllib, en desuso). Todo lo que veremos aquí usa la API de DataFrames.

La abstracción clave: Transformer y Estimator

MLlib organiza todo el flujo de trabajo en dos tipos de componentes:

ComponenteQué haceEjemplos
TransformerToma un DataFrame y devuelve otro DataFrame transformadoVectorAssembler, StringIndexer (ya ajustado), un modelo entrenado
EstimatorSe ajusta (fit) a los datos y produce un TransformerLogisticRegression, RandomForestClassifier, StringIndexer (sin ajustar)

Un Pipeline es simplemente una cadena ordenada de estos componentes. Al llamar a pipeline.fit(train), Spark ejecuta cada etapa en orden: los Estimators se ajustan y los Transformers transforman, pasando el DataFrame resultante a la siguiente etapa. El resultado (pipelineModel) es a su vez un Transformer que se aplica con .transform(test).

graph LR
RAW["DataFrame crudo"] --> SI["StringIndexer<br/>(Estimator)"]
SI --> VA["VectorAssembler<br/>(Transformer)"]
VA --> SCALER["StandardScaler<br/>(Estimator)"]
SCALER --> LR["LogisticRegression<br/>(Estimator)"]
LR --> MODEL["PipelineModel<br/>(Transformer)"]
MODEL --> PRED["Predicciones<br/>sobre datos nuevos"]
Un pipeline de MLlib: los datos fluyen por etapas de preparación y terminan en un modelo que genera predicciones.

La ventaja práctica es enorme: el mismo pipeline se aplica al entrenamiento y a la producción, evitando el error clásico de que los datos de entrenamiento se preprocesan de una forma y los de predicción de otra.

Feature engineering con VectorAssembler

Los algoritmos de MLlib no aceptan columnas sueltas: esperan una única columna llamada features que contiene un vector denso o disperso. VectorAssembler combina varias columnas numéricas en ese vector:

from pyspark.ml.feature import VectorAssembler

assembler = VectorAssembler(
    inputCols=["edad", "ingresos", "num_compras"],
    outputCol="features"
)

Para variables categóricas se usa StringIndexer (etiqueta → índice numérico) y, si se quiere evitar un orden artificial, OneHotEncoder. Para escalar magnitudes muy distintas, StandardScaler normaliza cada componente del vector.

Train/test split y entrenamiento

La división entrenamiento/prueba es una línea, pero en distribuido conviene fijar la semilla para reproducibilidad:

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

Sobre train ajustamos el pipeline completo. Los clasificadores más usados en MLlib son:

  • Regresión logística (LogisticRegression): lineal, rápida, buen punto de partida. La optimización se distribuye con gradiente descendente en mini-lotes a través del clúster.
  • Random forest (RandomForestClassifier): conjunto de árboles de decisión; cada árbol se entrena en paralelo sobre subconjuntos de datos y atributos. Robusto y poco sensible a hiperparámetros, pero más costoso.

Evaluación de modelos

MLlib incluye evaluadores en pyspark.ml.evaluation:

EvaluadorMétricaUso típico
BinaryClassificationEvaluatorAUC-ROC, AUC-PRClasificación binaria desbalanceada
MulticlassClassificationEvaluatorAccuracy, F1, precisión, recallClasificación multiclase
RegressionEvaluatorRMSE, MAE, R²Regresión
from pyspark.ml.evaluation import MulticlassClassificationEvaluator

evaluator = MulticlassClassificationEvaluator(metricName="f1")
f1 = evaluator.evaluate(predicciones)

No te fíes solo del accuracy cuando las clases están desbalanceadas. Un modelo que siempre predice la clase mayoritaria puede tener 95% de accuracy y no servir para nada. Mira F1, precisión y recall por clase.

Para buscar hiperparámetros a escala, MLlib ofrece CrossValidator con ParamGridBuilder: entrena y evalúa cada combinación de hiperparámetros con validación cruzada, repartiendo todo ese trabajo por el clúster.

Ejercicio: pipeline de clasificación completo

🧪 Ejercicio

Clasificación de churn con MLlib

Completa el pipeline para predecir si un cliente abandonará (churn = 1) usando PySpark MLlib. Debes: (1) ensamblar las columnas numéricas en un vector de features con VectorAssembler, (2) crear un Pipeline con el ensamblador y un LogisticRegression, (3) hacer el split 80/20, (4) entrenar con fit sobre train, (5) predecir sobre test y (6) evaluar el F1-score. Asume que existe una SparkSession llamada spark y que los datos están en 'churn.csv' con las columnas edad, antiguedad_meses, uso_mensual, churn.

Comprueba lo aprendido

Comprueba que lo pillaste

En MLlib, ¿qué diferencia hay entre un Estimator y un Transformer?

Comprueba que lo pillaste

¿Por qué es necesario VectorAssembler antes de entrenar un clasificador en MLlib?

Comprueba que lo pillaste

Un modelo de detección de fraude acierta el 99% de los casos, pero el fraude solo representa el 0.5% de las transacciones. ¿Cuál es la evaluación más adecuada?

Resumen

  • MLlib distribuye el entrenamiento de modelos por el clúster usando la API de DataFrames (pyspark.ml).
  • Todo el flujo se organiza en Pipelines de Estimators (que aprenden con fit) y Transformers (que aplican cambios con transform).
  • VectorAssembler une las columnas numéricas en el vector features que exigen los algoritmos.
  • El mismo pipeline sirve para entrenamiento y producción: fit sobre train, transform sobre test.
  • La evaluación se hace con evaluadores (F1, AUC, RMSE) y la búsqueda de hiperparámetros con CrossValidator, también distribuida.
  • Con clases desbalanceadas, el accuracy solo es una métrica engañosa.

📚 Lecturas y fuentes

RecursoTipoPor qué leerlo
Apache Spark (Wikipedia en español)DocsPanorama general en 10 minutos; lee solo la sección “MLlib” para situar la librería dentro del ecosistema Spark.
MLlib: Main Guide — ML PipelinesDocsLa referencia oficial de Transformer, Estimator y Pipeline; con leer “Pipeline components” y el ejemplo de código basta. (en inglés)
Learning Spark, 2nd EditionLibroDescarga gratuita; capítulos 10 y 11 cubren MLlib y la puesta en producción de modelos. Salta el resto. (en inglés)
ML Tuning: model selection and hyperparameter tuningDocsPara cuando llegues a CrossValidator: lee la sección “Cross-Validation” y el ParamGridBuilder del ejemplo. (en inglés)