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:
| Componente | Qué hace | Ejemplos |
|---|---|---|
| Transformer | Toma un DataFrame y devuelve otro DataFrame transformado | VectorAssembler, StringIndexer (ya ajustado), un modelo entrenado |
| Estimator | Se ajusta (fit) a los datos y produce un Transformer | LogisticRegression, 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"]
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:
| Evaluador | Métrica | Uso típico |
|---|---|---|
BinaryClassificationEvaluator | AUC-ROC, AUC-PR | Clasificación binaria desbalanceada |
MulticlassClassificationEvaluator | Accuracy, F1, precisión, recall | Clasificación multiclase |
RegressionEvaluator | RMSE, 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.
🔍 El orden de las etapas del pipeline importa: primero VectorAssembler (necesita las columnas sueltas), luego el clasificador (que consume la columna 'features'). El evaluador necesita labelCol='churn' por defecto si el CSV usa ese nombre... revisa: labelCol por defecto es 'label', así que pásalo explícito.
from pyspark.sql import SparkSession
from pyspark.ml import Pipeline
from pyspark.ml.feature import VectorAssembler
from pyspark.ml.classification import LogisticRegression
from pyspark.ml.evaluation import MulticlassClassificationEvaluator
spark = SparkSession.builder.appName("churn").getOrCreate()
# 1. Cargar datos
datos = spark.read.csv("churn.csv", header=True, inferSchema=True)
# 2. Ensamblar features
assembler = VectorAssembler(
inputCols=["edad", "antiguedad_meses", "uso_mensual"],
outputCol="features"
)
# 3. Definir el clasificador (la etiqueta es 'churn')
lr = LogisticRegression(labelCol="churn", featuresCol="features", maxIter=20)
# 4. Pipeline con las etapas en orden
pipeline = Pipeline(stages=[assembler, lr])
# 5. Split 80/20 con semilla fija
train, test = datos.randomSplit([0.8, 0.2], seed=42)
# 6. Entrenar y predecir
modelo = pipeline.fit(train)
predicciones = modelo.transform(test)
# 7. Evaluar F1
evaluator = MulticlassClassificationEvaluator(
labelCol="churn", predictionCol="prediction", metricName="f1"
)
f1 = evaluator.evaluate(predicciones)
print(f"F1-score en test: {f1:.4f}")Comprueba lo aprendido
Comprueba que lo pillaste
En MLlib, ¿qué diferencia hay entre un Estimator y un Transformer?
Un Estimator (p. ej. LogisticRegression) aprende de los datos con fit() y devuelve un modelo, que es un Transformer. Un Transformer (p. ej. VectorAssembler o el modelo entrenado) toma un DataFrame y devuelve otro con transform().
Comprueba que lo pillaste
¿Por qué es necesario VectorAssembler antes de entrenar un clasificador en MLlib?
Los algoritmos de MLlib no leen columnas sueltas: esperan una columna vectorial llamada 'features'. VectorAssembler concatena las columnas de entrada en ese vector. La conversión de texto a índices la hace StringIndexer.
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?
Un clasificador que siempre dice 'no fraude' ya logra 99.5% de accuracy sin detectar nada. Con desbalance extremo hay que mirar precisión, recall y F1 sobre la clase minoritaria (y curvas PR/AUC).
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 contransform). - VectorAssembler une las columnas numéricas en el vector
featuresque exigen los algoritmos. - El mismo pipeline sirve para entrenamiento y producción:
fitsobre train,transformsobre 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
| Recurso | Tipo | Por qué leerlo |
|---|---|---|
| Apache Spark (Wikipedia en español) | Docs | Panorama general en 10 minutos; lee solo la sección “MLlib” para situar la librería dentro del ecosistema Spark. |
| MLlib: Main Guide — ML Pipelines | Docs | La referencia oficial de Transformer, Estimator y Pipeline; con leer “Pipeline components” y el ejemplo de código basta. (en inglés) |
| Learning Spark, 2nd Edition | Libro | Descarga 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 tuning | Docs | Para cuando llegues a CrossValidator: lee la sección “Cross-Validation” y el ParamGridBuilder del ejemplo. (en inglés) |