Utilizando MLLib en Scala
Los siguientes fragmentos de código pueden ser ejecutados en spark-shell.
Clasificación Binaria
El siguiente fragmento de código ilustra cómo cargar un conjunto de datos de ejemplo, ejecutar un algoritmo de entrenamiento en estos datos utilizando un método estático en el objeto del algoritmo, y realizar predicciones con el modelo resultante para calcular el error de entrenamiento.
import org.apache.spark.SparkContext
import org.apache.spark.mllib.classification.SVMWithSGD
import org.apache.spark.mllib.regression.LabeledPoint
// Cargar y parsear el archivo de datos
val datos = sc.textFile("mllib/datos_ejemplo_svm.txt")
val datosParseados = datos.map { linea =>
val partes = linea.split(' ')
LabeledPoint(partes(0).toDouble, partes.tail.map(x => x.toDouble).toArray)
}
// Ejecutar algoritmo de entrenamiento para construir el modelo
val iteraciones = 20
val modelo = SVMWithSGD.train(datosParseados, iteraciones)
// Evaluar el modelo en ejemplos de entrenamiento y calcular el error de entrenamiento
val etiquetasPredicciones = datosParseados.map { punto =>
val prediccion = modelo.predict(punto.features)
(punto.label, prediccion)
}
val errorEntrenamiento = etiquetasPredicciones.filter(r => r._1 != r._2).count.toDouble / datosParseados.count
println("Error de Entrenamiento = " + errorEntrenamiento)
El método SVMWithSGD.train() por defecto realiza regularización L2 con el parámetro de regularización establecido en 1.0. Si queremos configurar este algoritmo, podemos personalizar SVMWithSGD aún más creando un nuevo objeto directamente y llamando a métodos establecedores. Todos los demás algoritmos de MLlib admiten esta personalización. Por ejemplo, el siguiente código produce una variante con regularización L1 de SVMs con parámetro de regularización establecido en 0.1, y ejecuta el algoritmo de entrenamiento durante 200 iteraciones.
import org.apache.spark.mllib.optimization.L1Updater
val svmAlgoritmo = new SVMWithSGD()
svmAlgoritmo.optimizer.setNumIterations(200)
.setRegParam(0.1)
.setUpdater(new L1Updater)
val modeloL1 = svmAlgoritmo.run(datosParseados)
Regresión Lineal
El siguiente ejemplo demuestra cómo cargar datos de entrenamiento, parsearlos como un RDD de LabeledPoint. El ejemplo luego utiliza LinearRegressionWithSGD para construir un modelo lineal simple para predecir valores de etiqueta. Calculamos el Error Cuadrático Medio al final para evaluar la bondad de ajuste.
import org.apache.spark.mllib.regression.LinearRegressionWithSGD
import org.apache.spark.mllib.regression.LabeledPoint
// Cargar y parsear los datos
val datos = sc.textFile("mllib/datos/ridge/lpsa.data")
val datosParseados = datos.map { linea =>
val partes = linea.split(',')
LabeledPoint(partes(0).toDouble, partes(1).split(' ').map(x => x.toDouble).toArray)
}
// Construir el modelo
val iteraciones = 20
val modelo = LinearRegressionWithSGD.train(datosParseados, iteraciones)
// Evaluar el modelo en ejemplos de entrenamiento y calcular el error de entrenamiento
val valoresPredicciones = datosParseados.map { punto =>
val prediccion = modelo.predict(punto.features)
(punto.label, prediccion)
}
val ECM = valoresPredicciones.map{ case(v, p) => math.pow((v - p), 2)}.reduce(_ + _)/valoresPredicciones.count
println("Error Cuadrático Medio de entrenamiento = " + ECM)
De manera similar puedes utilizar RidgeRegressionWithSGD y LassoWithSGD y comparar los Errores Cuadráticos Medios de entrenamiento.
Clustering
En el siguiente ejemplo, después de cargar y parsear los datos, utilizamos el objeto KMeans para agrupar los datos en dos clústeres. El número de clústeres deseados se pasa al algoritmo. Luego calculamos la Suma de Errores Cuadráticos Dentro del Conjunto (WSSSE). Puede reducir esta medida de error aumentando k. De hecho, el óptimo k suele ser donde hay un "codo" en la gráfica de WSSSE.
import org.apache.spark.mllib.clustering.KMeans
// Cargar y parsear los datos
val datos = sc.textFile("datos_kmeans.txt")
val datosParseados = datos.map( _.split(' ').map(_.toDouble))
// Agrupar los datos en dos clases usando KMeans
val iteraciones = 20
val numClusters = 2
val clusters = KMeans.train(datosParseados, numClusters, iteraciones)
// Evaluar el clustering calculando la Suma de Errores Cuadráticos Dentro del Conjunto
val WSSSE = clusters.computeCost(datosParseados)
println("Suma de Errores Cuadráticos Dentro del Conjunto = " + WSSSE)
Filtrado Colaborativo
En el siguiente ejemplo cargamos datos de calificación. Cada fila consiste en un usuario, un producto y una calificación. Utilizamos el método ALS.train() por defecto que asume que las calificaciones son explícitas. Evaluamos el modelo de recomendación midiendo el Error Cuadrático Medio de la predicción de calificaciones.
import org.apache.spark.mllib.recommendation.ALS
import org.apache.spark.mllib.recommendation.Rating
// Cargar y parsear los datos
val datos = sc.textFile("mllib/datos/als/test.data")
val calificaciones = datos.map(_.split(',') match {
case Array(usuario, item, tasa) => Rating(usuario.toInt, item.toInt, tasa.toDouble)
})
// Construir el modelo de recomendación usando ALS
val iteraciones = 20
val modelo = ALS.train(calificaciones, 1, 20, 0.01)
// Evaluar el modelo en datos de calificación
val usuariosProductos = calificaciones.map{ case Rating(usuario, producto, tasa) => (usuario, producto)}
val predicciones = modelo.predict(usuariosProductos).map{
case Rating(usuario, producto, tasa) => ((usuario, producto), tasa)
}
val tasasPredicciones = calificaciones.map{
case Rating(usuario, producto, tasa) => ((usuario, producto), tasa)
}.join(predicciones)
val ECM = tasasPredicciones.map{
case ((usuario, producto), (t1, t2)) => math.pow((t1- t2), 2)
}.reduce(_ + _)/tasasPredicciones.count
println("Error Cuadrático Medio = " + ECM)
Si la matriz de calificación se deriva de otra fuente de información (es decir, se infiere de otras señales), puede usar el método trainImplicit para obtener mejores resultados.
val modelo = ALS.trainImplicit(calificaciones, 1, 20, 0.01)
Utilizando MLLib en Java
Todos los métodos de MLlib usan tipos amigables con Java, por lo que puedes importar y llamar desde Java de la misma manera que en Scala. La única advertencia es que los métodos toman objetos RDD de Scala, mientras que la API de Java de Spark utiliza una clase JavaRDD separada. Puedes convertir un Java RDD a uno de Scala llamando a .rdd() en tu objeto JavaRDD.
Utilizando MLLib en Python
Los siguientes ejemplos pueden ser probados en el shell de PySpark.
Clasificación Binaria
El siguiente ejemplo muestra cómo cargar un conjunto de datos de ejemplo, construir un modelo de Regresión Logística, y realizar predicciones con el modelo resultante para calcular el error de entrenamiento.
from pyspark.mllib.classification import LogisticRegressionWithSGD
from numpy import array
# Cargar y parsear los datos
datos = sc.textFile("mllib/datos/sample_svm_data.txt")
datosParseados = datos.map(lambda linea: array([float(x) for x in linea.split(' ')]))
modelo = LogisticRegressionWithSGD.train(datosParseados)
# Construir el modelo
etiquetasPredicciones = datosParseados.map(lambda punto: (int(punto.item(0)),
modelo.predict(punto.take(range(1, punto.size)))))
# Evaluar el modelo en datos de entrenamiento
errorEntrenamiento = etiquetasPredicciones.filter(lambda (v, p): v != p).count() / float(datosParseados.count())
print("Error de Entrenamiento = " + str(errorEntrenamiento))
Regresión Lineal
El siguiente ejemplo demuestra cómo cargar datos de entrenamiento, parsearlos como un RDD de LabeledPoint. El ejemplo luego utiliza LinearRegressionWithSGD para construir un modelo lineal simple para predecir valores de etiqueta. Calculamos el Error Cuadrático Medio al final para evaluar la bondad de ajuste.
from pyspark.mllib.regression import LinearRegressionWithSGD
from numpy import array
# Cargar y parsear los datos
datos = sc.textFile("mllib/datos/ridge-data/lpsa.data")
datosParseados = datos.map(lambda linea: array([float(x) for x in linea.replace(',', ' ').split(' ')]))
# Construir el modelo
modelo = LinearRegressionWithSGD.train(datosParseados)
# Evaluar el modelo en datos de entrenamiento
valoresPredicciones = datosParseados.map(lambda punto: (punto.item(0),
modelo.predict(punto.take(range(1, punto.size)))))
ECM = valoresPredicciones.map(lambda (v, p): (v - p)**2).reduce(lambda x, y: x + y)/valoresPredicciones.count()
print("Error Cuadrático Medio = " + str(ECM))
Clustering
En el siguiente ejemplo, después de cargar y parsear los datos, utilizamos el objeto KMeans para agrupar los datos en dos clústeres. El número de clústeres deseados se pasa al algoritmo. Luego calculamos la Suma de Errores Cuadráticos Dentro del Conjunto (WSSSE). Puede reducir esta medida de error aumentando k. De hecho, el óptimo k suele ser donde hay un "codo" en la gráfica de WSSSE.
from pyspark.mllib.clustering import KMeans
from numpy import array
from math import sqrt
# Cargar y parsear los datos
datos = sc.textFile("kmeans_data.txt")
datosParseados = datos.map(lambda linea: array([float(x) for x in linea.split(' ')]))
# Construir el modelo (agrupar los datos)
clusters = KMeans.train(datosParseados, 2, maxIterations=10,
runs=30, initialization_mode="random")
# Evaluar el clustering calculando la Suma de Errores Cuadráticos Dentro del Conjunto
def error(punto):
centro = clusters.centers[clusters.predict(punto)]
return sqrt(sum([x**2 for x in (punto - centro)]))
WSSSE = datosParseados.map(lambda punto: error(punto)).reduce(lambda x, y: x + y)
print("Suma de Errores Cuadráticos Dentro del Conjunto = " + str(WSSSE))
De manera similar puedes utilizar RidgeRegressionWithSGD y LassoWithSGD y comparar los Errores Cuadráticos Medios de entrenamiento.
Filtrado Colaborativo
En el siguiente ejemplo cargamos datos de calificación. Cada fila consiste en un usuario, un producto y una calificación. Utilizamos el método ALS.train() por defecto que asume que las calificaciones son explícitas. Evaluamos la recomendación midiendo el Error Cuadrático Medio de la predicción de calificaciones.
from pyspark.mllib.recommendation import ALS
from numpy import array
# Cargar y parsear los datos
datos = sc.textFile("mllib/datos/als/test.data")
calificaciones = datos.map(lambda linea: array([float(x) for x in linea.split(',')]))
# Construir el modelo de recomendación usando Mínimos Cuadrados Alternados
modelo = ALS.train(calificaciones, 1, 20)
# Evaluar el modelo en datos de entrenamiento
datosPrueba = calificaciones.map(lambda p: (int(p[0]), int(p[1])))
predicciones = modelo.predictAll(datosPrueba).map(lambda r: ((r[0], r[1]), r[2]))
tasasPredicciones = calificaciones.map(lambda r: ((r[0], r[1]), r[2])).join(predicciones)
ECM = tasasPredicciones.map(lambda r: (r[1][0] - r[1][1])**2).reduce(lambda x, y: x + y)/tasasPredicciones.count()
print("Error Cuadrático Medio = " + str(ECM))
Si la matriz de calificación se deriva de otra fuente de información (es decir, se infiere de otras señales), puede usar el método trainImplicit para obtener mejores resultados.
# Construir el modelo de recomendación usando Mínimos Cuadrados Alternados basado en calificaciones implícitas
modelo = ALS.trainImplicit(calificaciones, 1, 20)