Mostrando entradas con la etiqueta Machine Learning. Mostrar todas las entradas
Mostrando entradas con la etiqueta Machine Learning. Mostrar todas las entradas

martes, 20 de diciembre de 2016

XGBOOST & Hadoop/YARN (II): Ejemplo

En la anterior entrada vimos cómo instalar la librería XGBOOST sobre CentOS con soporte HDFS. Bien, en ésta trataremos de ver su ejecución a través de algún ejemplo, eso sí, sin entrar a valorar el resultado o si se puede mejorar el modelo, variables, etc, simplemente se trata de demostrar la funcionalidad de poder ejecutar la librería XGBOOST en modo distribuido.

DMLC-YARN.JAR
Lo primero que deberemos hacer es crear el paquete dmlc-yarn.jar necesario para la ejecución de XGBOOST en modo distribuido, YARN.

     cd /usr/local/xgboost/dmlc-core/tracker/yarn/
     ./build.sh // NO preocuparse por los warnings.

     ls -al dmlc-yarn.jar // Debemos comprobar que se crea el fichero dmlc-yarn.jar
     rw-r--r-- 1 root root 21292 dic 16 13:39 dmlc-yarn.jar



Ejemplo

Ejemplo de cómo ejecutar XGBOOST en modo YARN con los conjuntos de datos almacenados previamente en HDFS. Además, el modelo resultante también será salvado en el HDFS.

Nota: El ejemplo mostrado a continuación fue ejecutado bajo el usuario root, totalmente desaconsejable su uso, pero así que cada uno lo adapte a su entorno ;)

1) Creo la estructura de directorios necesarios para el ejemplo en el HDFS:
     hadoop fs -mkdir /user/root/xgboost
     hadoop fs -mkdir /user/root/xgboost/data
     hadoop fs -mkdir /user/root/xgboost/data/train
     hadoop fs -mkdir /user/root/xgboost/data/test
     hadoop fs -mkdir /user/root/xgboost/model
     hadoop fs -chmod 777 /user/root/xgboost/model

2) Cargo los siguientes datasets de ejemplo que vienen, por defecto, con XGBOOST.
     hadoop fs -put /usr/local/xgboost/demo/data/agaricus.txt.train /user/root/xgboost/data/train/
     hadoop fs -put /usr/local/xgboost/demo/data/agaricus.txt.test /user/root/xgboost/data/test/

3) Creo la configuración necesaria para llevar a cabo el ejemplo.
     export LD_LIBRARY_PATH=/usr/hdp/2.4.2.0-258/usr/lib/
     vim /usr/local/xgboost/demo/distributed-training/hdfs.conf
          # General Parameters, see comment for each definition
          # choose the booster, can be gbtree or gblinear
          booster = gbtree
          # choose logistic regression loss function for binary classification
          objective = binary:logistic

          # Tree Booster Parameters
          # step size shrinkage
          eta = 1.0
          # minimum loss reduction required to make a further partition
          gamma = 1.0
          # minimum sum of instance weight(hessian) needed in a child
          min_child_weight = 1
          # maximum depth of a tree
          max_depth = 3

          # Task Parameters
          # the number of round to do boosting
          num_round = 2
          # 0 means do not save any model except the final round model
          save_period = 0
          # The path of training data
          data = "hdfs://namenode01.domain.com/user/root/xgboost/data/train"
          # The path of validation data, used to monitor training process
          eval[test] = "hdfs://namenode01.domain.com/root/xgboost/data/test"
          model_dir = "hdfs://namenode01.domain.com/user/root/xgboost/model"
          # evaluate on training data as well each round
          eval_train = 1

¡¡¡Importante!!! Incluir en el path del HDFS el FQDN de nuestro NameNode principal o bien el del grupo si hemos configurado la alta disponibilidad.

4) Ahora ya sí, se puede proceder a ejecutar el ejemplo:
/usr/local/xgboost/dmlc-core/tracker/dmlc-submit --cluster yarn --num-workers 2 /usr/local/xgboost/xgboost /usr/local/xgboost/demo/distributed-training/hdfs.conf
2016-12-16 13:43:19,708 INFO start listen on 192.168.4.245:9091
16/12/16 13:43:21 WARN util.NativeCodeLoader: Unable to load native-hadoop library for your platform... using builtin-java classes where applicable
16/12/16 13:43:22 WARN shortcircuit.DomainSocketFactory: The short-circuit local reads feature cannot be used because libhadoop cannot be loaded.
16/12/16 13:43:23 INFO impl.TimelineClientImpl: Timeline service address: http://namenode01.domain.com:8188/ws/v1/timeline/
16/12/16 13:43:23 INFO client.RMProxy: Connecting to ResourceManager at namenode01.domain.com/192.168.4.245:8050
16/12/16 13:43:24 INFO dmlc.Client: jobname=DMLC[nworker=2]:xgboost,username=root
16/12/16 13:43:24 INFO dmlc.Client: Submitting application application_1481716095353_0075
16/12/16 13:43:24 INFO impl.YarnClientImpl: Submitted application application_1481716095353_0075
2016-12-16 13:43:29,410 INFO @tracker All of 2 nodes getting started
2016-12-16 13:43:33,723 INFO [13:43:33] [0] test-error:0.016139 train-error:0.014433
2016-12-16 13:43:33,904 INFO [13:43:33] [1] test-error:0.000000 train-error:0.001228
2016-12-16 13:43:34,253 INFO @tracker All nodes finishes job
2016-12-16 13:43:34,253 INFO @tracker 4.84310793877 secs between node start and job finish
Application application_1481716095353_0075 finished with state FINISHED at 1481888615119

5) Comprobar la existencia del modelo recién creado. El nombre del fichero, por defecto, sigue el patrón <num_rounds>.model siendo <num_rounds> el valor establecido para dicha variable en el fichero de configuracion.
     hadoop fs -ls /user/root/xgboost/model
     Found 1 items
     -rw-r--r-- 2 yarn root 1501 2016-12-16 13:43 /user/root/xgboost/model/0002.model

El modelo puede ser cargado posteriormente en R, Python o Julia.

lunes, 19 de diciembre de 2016

XGBOOST & Hadoop/YARN (I)

Este primer tutorial trata de explicar los pasos necesarios para desplegar la librería XGBOOST sobre CentOS con soporte HDFS, y más concretamente sobre un clúster Hadoop / YARN, pues pese a existir la "Installation Guide" en su página principal sobre cómo hacerlo, ésta 'sólo' cubre los sistemas operativos Ubuntu/Debian, Windows y OSX.

Entorno

Antes de pasar a realizar cualquier tipo de acción, me gustaría detallar cual es el entorno sobre el que voy a trabajar y desplegar la librería XGBOOST,

  • CentOS 6.7
  • Hortonworks Data Platform (HDP) v.2.4.2 => Hadoop v2.7.1

Requisitos

G++ >= 4.6

Decir que CentOS 6.7 dispone de un compilador GCC y G++ bastante 'desactualizados' en sus repositorios oficiales, por lo que recurrí al siguiente procedimiento para la instalación de una versión 'más reciente' y superior incluso a la requerida:

     yum install wget git -y // En caso de no disponer previamente de estos paquetes
     cd /etc/yum.repos.d
     wget https://people.centos.org/tru/devtools-2/devtools-2.repo

Instalamos a continuación los siguientes paquetes:
     devtoolset-2-gcc.x86_64
     devtoolset-2-gcc-c++.x86_64
     devtoolset-2-gcc-plugin-devel.x86_64
     devtoolset-2-binutils.x86_64
     devtoolset-2-binutils-devel.x86_64

Una vez instalados, comprobar la versión desplegada:
     /opt/rh/devtoolset-2/root/usr/bin/gcc -v
     …
     gcc version 4.8.2 20140120 (Red Hat 4.8.2-15) (GCC)

Python v2.7
Además, en CentOS 6.7, por defecto, el Python incluido en los repositorios suele ser el de la versión 'obsoleta' 2.6.6, por lo que también deberemos actualizarla para poder trabajar posteriormente con la librería XGBOOST. Las versiones aceptadas por XGBOOST son Python 2.7 o superior, o Python 3.4 o superior.

En mi caso vamos a recurrir a la versión 2.7.6 de Python:
     yum -y install zlib-devel bzip2-devel openssl-devel ncurses-devel sqlite-devel
     cd /opt
     wget --no-check-certificate https://www.python.org/ftp/python/2.7.6/Python-2.7.6.tar.xz
     tar xf Python-2.7.6.tar.xz
     cd Python-2.7.6
     ./configure --prefix=/usr/local
     make && make altinstall

¡¡¡Importante!!! Usar altinstall en lugar de install porque sino acabaremos con dos versiones diferentes de Python instaladas en nuestro sistema y ambas nombradas Python.

Tras este tipo de instalación convivirán en nuestro sistema ambas versiones:
     python -V
     Python 2.6.6

     python2.7 -V
     Python 2.7.6

Pip 2.7
Para facilitar las futuras instalaciones de paquetes de Python, es aconsejable también actualizar la versión de Pip.
     wget https://bitbucket.org/pypa/setuptools/downloads/ez_setup.py
     /usr/local/bin/python2.7 ez_setup.py
     /usr/local/bin/easy_install-2.7 pip

Comprobamos su correcto funcionamiento gracias a la instalación del paquete argparse que será necesario posteriormente para ejecutar y obtener la ayuda del comando xgboost.
     pip2.7 install argparse

Java
Debido a la instalación de la suite HDP de Hortonworks ya dispondremos de una versión de Java instalada en nuestro clúster, pero no viene mal repasar y asegurarse de su correcta instalación y configuración de las variables de entorno, $JAVA_HOME y $PATH.

[XGBOOST]
XGBOOST con soporte HDFS

Una vez cumplidos con los requerimientos, podemos pasar a construir la librería compartida de XGBOOST sobre CentOS con soporte HDFS.

     cd /usr/local
     git clone --recursive https://github.com/dmlc/xgboost
     cd xgboost

     cp make/config.mk ./config.mk

     vim config.mk // Descomentar las siguientes líneas y modificar su contenido de acuerdo a:

          export CC=/opt/rh/devtoolset-2/root/usr/bin/gcc

          export CPP=/opt/rh/devtoolset-2/root/usr/bin/cpp

          export CXX=/opt/rh/devtoolset-2/root/usr/bin/c++
          USE_HDFS = 1

     cd dmlc-core/
     cp make/config.mk config.mk
     vim config.mk // Descomentar las siguientes líneas y modificar su contenido de acuerdo a:
          export CC=/opt/rh/devtoolset-2/root/usr/bin/gcc
          export CPP=/opt/rh/devtoolset-2/root/usr/bin/cpp
          export CXX=/opt/rh/devtoolset-2/root/usr/bin/c++
          USE_HDFS = 1

     cd /usr/local/xgboost
     make clean_all
     make -j4


Para comprobar que se ha construido la librería satisfactoriamente, deberemos asegurarnos que se ha generado el archivo libxgboost.so, en nuestro caso:

     ls /usr/local/xgboost/lib/libxgboost.so

     -rwxr-xr-x 1 root root 2374684 dic 15 16:48 lib/libxgboost.so

Una vez finalizada su construcción satisfactoriamente, deberemos editar el fichero dmlc-submit para que use la versión 2.7 de Python en lugar de la 2.6. Para ello bastará con modificar la primera línea de dicho fichero:
     vim /usr/local/xgboost/dmlc-core/tracker/dmlc-submit
          #!/usr/bin/env python2.7

Problema hdfs.h / libdmlc.a
En un primer intento a la hora de construir la librería, me surgió el error que podéis observar más abajo. Realmente el fallo no se debe a la falta del archivo o libreria libdmlc.a. Observando con detenimiento un poco más arriba del comentado error, detecté otro fallo "hdfs.h: No existe el fichero o el directorio" el cual desencadena el error final.

Problema:
     In file included from src/io.cc:16:0:
     src/io/hdfs_filesys.h:10:18: fatal error: hdfs.h: No existe el fichero o el directorio
     #include <hdfs.h>
     …
     c++: error: dmlc-core/libdmlc.a: No existe el fichero o el directorio
     make: *** [lib/libxgboost.so] Error 1
     make: *** Se espera a que terminen otras tareas....

Solución: Me bastó con buscar la librería hdfs.h en el sistema y modificar la definición de la siguiente línea dentro del fichero de configuración dmlc.mk:
     vim /usr/local/xgboost/dmlc-core/make/dmlc.mk
          HDFS_INC_PATH=/usr/hdp/2.4.2.0-258/usr/include

Por útlimo, volvemos a ejecutar los comandos necesarios para construir la librería compartida de XGBOOST.
     make clean_all
     make -j4


[+Info] Para ver un ejemplo de su ejecución en modo YARN: XGBOOST & Hadoop/YARN (II): Ejemplo

viernes, 11 de marzo de 2016

[Spark MLlib] Filtrado Colaborativo (II): Recomendación por Usuario

Continuamos con este breve tutorial o introducción al filtrado colaborativo con Spark. Bien, la semana pasada en la entrada [Spark MLlib] Filtrado Colaborativo (I): Sistema de Recomendación construimos nuestro sistema de recomendación o más bien el generador de modelos. Esta semana tocará hacer uso de él y recomendar productos, en nuestro ejemplo películas, a nuestros usuarios a partir de sus votaciones.

Continuaremos también trabajando con la estructura de directorios y archivos ya definido, lo que vendría a ser nuestro proyecto dentro de nuestra herramienta IDE favorita. Antes de continuar, me gustaría matizar que este tutorial NO está orientado a la programación en scala ni a sus best practices a la hora de desarrollar. Todos los ejemplos son mejorables desde la primera línea. Por ejemplo, algo tan sencillo como recibir los parámetros con los que juego como argumentos de entrada, definir un único método principal, pues voy declarando funciones main por cada funcionalidad o algoritmo que trato de explicar, etc, pero lo dicho, no, no me voy a entretener ni por un instante en escribir el código correctamente, sólo pretendo dar unas pinceladas sobre cómo funcionan ciertos algoritmos o cómo alcanzar una determinada funcionalidad.

1) En esta ocasión creamos el archivo getRecommendationsByUserId.scala dentro de la misma carpeta donde creamos nuestro generador de modelos, es decir, en el directorio Carpeta_Proyecto/src/main/scala/mllib/collaborativeFiltering/

package mllib.collaborativeFiltering

import org.apache.log4j.{Level, Logger}
import org.apache.spark.mllib.recommendation.{ALS, Rating, MatrixFactorizationModel}
import org.apache.spark.{SparkContext, SparkConf}

object getRecommendationsByUserId {
  def main (args: Array[String]): Unit = {

    Logger.getLogger("org.apache.spark").setLevel(Level.WARN)
    Logger.getLogger("org.eclipse.jetty.server").setLevel(Level.OFF)

    val sparkConfig = new SparkConf().setMaster("local[2]").setAppName("getRecommendationsByUserId")
    val sc = new SparkContext(sparkConfig)

    val mlRatings = sc.textFile("data/movielens/ratings.csv")
    val ratings = mlRatings.map { line =>
      // userId,movieId,rating,timestamp
      val fields = line.split(",")
      Rating(fields(0).toInt, fields(1).toInt, fields(2).toDouble)
    }.cache()

    val mlMovies = sc.textFile("data/movielens/movies.csv")
    val movies = mlMovies.map { line =>
      // movieId,title,genres
      val fields = line.split(",")
      (fields(0).toInt, fields(1))
    }.collect.toMap

    val model = MatrixFactorizationModel.load(sc, "target/tmp/myCollaborativeFilter")    // Recuperando el modelo

    val userId = 25560     // Seleccionamos el usuario id. Interesante pasarlo al script como argumento ;)
    val moviesForUser = ratings.keyBy(_.user).lookup(userId)
    println("El usuario ha votado a " + moviesForUser.size + " peliculas:")
    moviesForUser.sortBy(-_.rating).take(10).map(rating => (movies(rating.product), rating.rating)).foreach(println)    // Observamos que tipo de películas le gustan al usuario

    val topKRecs = model.recommendProducts(userId,10)    // El método para obtener las recomendaciones ya viene implementado de serie en la librería. Obtenemos el top10 de productos, en nuestro caso películas, recomendados. También podríamos pasar el número de recomendaciones, topN  ;) ;)
    topKRecs.map(rating => (movies(rating.product), rating.rating)).foreach(println)

  }
}

2) Ejecutar la aplicación
# sbt compile
...
[success] Total time: 1 s, completed 11-mar-2016 10:14:13

# sbt run
...
Multiple main classes detected, select one to run:

 [1] mllib.DecisionTrees.classificationFlightsOnTime
 [2] mllib.collaborativeFiltering.getRecommendationsByUserId
 [3] mllib.collaborativeFiltering.recommendation

Enter number: 2
...
El usuario ha votado a 50 peliculas:
Star Wars: Episode IV - A New Hope (1977),5.0)
(Star Wars: Episode V - The Empire Strikes Back (1980),5.0)
("Princess Bride,5.0)
(Star Wars: Episode VI - Return of the Jedi (1983),5.0)
(Life Is Beautiful (La Vita è bella) (1997),5.0)
(X-Men (2000),5.0)
(Unbreakable (2000),5.0)
("Lord of the Rings: The Fellowship of the Ring,5.0)
(Spider-Man (2002),5.0)

("Lord of the Rings: The Two Towers,5.0)
...
("Roman Spring of Mrs. Stone,11.756657663464846)
(Babar The Movie (1989),11.171622844283377)
(Goodbye to Language 3D (2014),9.747998979226367)
(Garfield's Halloween Adventure (1985),9.736430143083723)
(Jimmy Carter Man from Plains (2007),9.72594699681999)
(Dangerous (1935),9.256450056766873)
(Bad Ronald (1974),9.089006045667578)
(End of the Line (2007),9.057430462126002)
(In the City of Sylvia (En la ciudad de Sylvia) (2007),8.963930587070152)
(Nine Dead Gay Guys (2003),8.895781428393946)
...
[success] Total time: 317 s, completed 11-mar-2016 10:41:22

Muy sencillo, ¿verdad?

Entradas relacionadas:

    [Spark MLlib] Filtrado Colaborativo (III): Recomendación por Ítem (pendiente)

viernes, 4 de marzo de 2016

[Spark MLlib] Filtrado Colaborativo (I): Sistema de Recomendación

Continuando mi aventura dentro del mundo del Machine Learning y tras [Spark MLlib] Árboles de Decisión: Ejemplo de Clasificación hoy voy a realizar y tratar de explicar el típico ejemplo de un sistema de recomendación gracias al uso de filtrado colaborativo (collaborative filtering).

Para ello me apoyaré en el 'desconocidísimo' dataset MovieLens y de nuevo en la librería MLlib de Spark que incluye para este tipo de soluciones el algoritmo ALS. También podríamos implementar dicho sistema de recomendación con Mahout, poderosísima herramienta de Machine Learning dentro del ecosistema de Big Data, que personalmente y para este problema en concreto presenta un mayor abanico de posibilidades de solución o al menos de algoritmos para atacarlo.
Collaborative Filtering
Filtrado Colaborativo
http://www.gravityrd.com/technology
Bueno, ¡al turrón!

1) Partiendo de la estructura de directorios del ejemplo anterior, ver [Spark MLlib] Árboles de Decisión: Ejemplo de Clasificación, crearemos la carpeta Carpeta_Proyecto/src/main/scala/mllib/collaborativeFiltering.

2) Nos descargaremos cualquiera de los conjunto de datos de MovieLens, en mi caso elegí "ml-latest.zip (size: 144 MB)", y extraeremos su contenido en el directorio Carpeta_Proyecto/data/movielens.

2.1) Antes de continuar deberemos realizar una pequeña edición sobre los ficheros ratings.csv y movies.csv. La primera línea de ambos ficheros es la definición de sus columnas por lo que deberemos eliminarla.

3) Fichero build.sbt. No sufre ninguna modificación respecto al del anterior ejemplo.

4) Copiar-pegar el siguiente código dentro del fichero recommendation.scala ubicado en el directorio Carpeta_Proyecto/src/main/scala/mllib/collaborativeFiltering/

package mllib.collaborativeFiltering

import org.apache.log4j.{Level, Logger}
import org.apache.spark.rdd.RDD
import org.apache.spark.{SparkConf, SparkContext}
import org.apache.spark.mllib.recommendation.{ALS, Rating, MatrixFactorizationModel}

object recommendation {
  def main (args: Array[String]): Unit = {

    Logger.getLogger("org.apache.spark").setLevel(Level.WARN)
    Logger.getLogger("org.eclipse.jetty.server").setLevel(Level.OFF)

    val sparkConfig = new SparkConf().setMaster("local[2]").setAppName("recommendation")
    val sc = new SparkContext(sparkConfig)    // Inicializar la aplicación de Spark mediante su objeto necesario

    // Leer el fichero de votaciones a partir del cual se crearán por cada línea un objeto de tipo Ratings cuyos argumentos deben ser: identificador de usuario, id. de producto (en nuestro caso moviesid), y valoración. Todos los argumentos deben ser de tipo numérico. En caso de no serlo alguno de ellos, deberemos en primer lugar proceder a su codificación o mapeo. Ej. variables carrierMap, originMap o destMap de la entrada anterior.
    val mlRatings = sc.textFile("data/movielens/ratings.csv")
    val ratings = mlRatings.map { line =>
      // userId,movieId,rating,timestamp
      val fields = line.split(",")
      Rating(fields(0).toInt, fields(1).toInt, fields(2).toDouble)
    }

    // Dividir el dataset inicial de valoraciones en tres conjuntos: entrenamiento, trest y validación
    val splits = ratings.randomSplit(Array(0.6, 0.2, 0.2))
    val (trainingData, testData, validationData) = (splits(0), splits(1), splits(2))

    // En el ejemplo que se propone en la web de Spark al respecto se implementa un único modelo, es decir, se determina el valor de las variables necesarias y se lleva a cabo su ejecución contra el conjunto de datos completo. Mientras que aquí he complicado un poco el código y lo que se busca es jugar con diferentes valores de las variables y con el conjunto de datos de entrenamiento. Se ejecuta de manera secuencial cada una de las combinaciones posibles quedándose al final con el modelo que mejor resultado o validación Root Mean Squared Error (RMSE) obtenga. Inconveniente, el tiempo de ejecución.

    val ranks = List(10, 15)    // Valor razonable entre 10-200. A mayor valor, generalmente mejor resultado, pero a costa de un mayor consumo de memoria
    val lambdas = List(0.01, 1.0)    // Su valor se debe establecer en función de las aproximaciones o resultados que vayamos obteniendo
    val numIters = List(10, 20)    // Valor habitual: 10; Con pocas iteraciones se consigue converger hacia una solución razonable
    var bestModel: Option[MatrixFactorizationModel] = None
    var bestValidationRmse = Double.MaxValue
    var bestRank = 0
    var bestLambda = -1.0
    var bestNumIter = -1
    for (rank <- ranks; lambda <- lambdas; numIter <- numIters) {
      val model = ALS.train(trainingData, rank, numIter, lambda)    // Ejecución del algoritmo ALS y construcción del modelo respecto del conjunto de datos de entrenamiento
      val validationRmse = computeRmse(model, validationData, false)    // Comparación o validación Root Mean Squared Error del modelo recién creado contra el conjunto de datos de validación
      println("RMSE (validation) = " + validationRmse + " for the model trained with rank = "
        + rank + ", lambda = " + lambda + ", and numIter = " + numIter + ".")
      if (validationRmse < bestValidationRmse) {Compute
        bestModel = Some(model)
        bestValidationRmse = validationRmse
        bestRank = rank
        bestLambda = lambda
        bestNumIter = numIter
      }
    }

    val testRmse = computeRmse(bestModel.get, testData, false)    // Ejecutar el modelo ganador contra el conjunto de datos de test

    println("The best model was trained with rank = " + bestRank + " and lambda = " + bestLambda
      + ", and numIter = " + bestNumIter + ", and its RMSE on the test set is " + testRmse + ".")

    bestModel.get.save(sc, "target/tmp/myCollaborativeFilter")    // Salvar el modelo

    sc.stop()    // Detener la aplicación de Spark
  }

  // Función para el cálculo del error cuadrático medio ó RMSE

  // En primer lugar se ejecuta el modelo en cuestión contra el conjunto de datos traspasado. A continuación se comparan los valores predecidos con las valoraciones reales hechas por los usuarios, obteniéndose así el error cuadrático medio final.
  def computeRmse(model: MatrixFactorizationModel, data: RDD[Rating], implicitPrefs: Boolean)
    : Double = {

    def mapPredictedRating(r: Double): Double = {
      if (implicitPrefs) math.max(math.min(r, 1.0), 0.0) else r
    }

    val predictions: RDD[Rating] = model.predict(data.map(x => (x.user, x.product)))
    val predictionsAndRatings = predictions.map{ x =>
      ((x.user, x.product), mapPredictedRating(x.rating))
    }.join(data.map(x => ((x.user, x.product), x.rating))).values
    math.sqrt(predictionsAndRatings.map(x => (x._1 - x._2) * (x._1 - x._2)).mean())
  }
}

5) Ejecutar la aplicación
# sbt compile
...
[success] Total time: 5 s, completed 04-mar-2016 09:59:18

# sbt run
...
Multiple main classes detected, select one to run:   // En el caso de haber continuado a partir del proyecto de la anterior entrada, [Spark MLlib] Árboles de Decisión: Ejemplo de Clasificación, nos aparecerá este mensaje.

 [1] mllib.DecisionTrees.classificationFlightsOnTime
 [2] mllib.collaborativeFiltering.recommendation
Enter number: 2    // Seleccionaremos la opción que corresponda a nuestro sistema de recomendación
...
RMSE (validation) = 0.8652560855462604 for the model trained with rank = 10, lambda = 0.01, and numIter = 10.
RMSE (validation) = 0.8629134178168192 for the model trained with rank = 10, lambda = 0.01, and numIter = 20.
RMSE (validation) = 1.3281799325759402 for the model trained with rank = 10, lambda = 1.0, and numIter = 10.
RMSE (validation) = 1.3281784808552248 for the model trained with rank = 10, lambda = 1.0, and numIter = 20.
RMSE (validation) = 0.8842091343625045 for the model trained with rank = 15, lambda = 0.01, and numIter = 10.
RMSE (validation) = 0.8788401499772798 for the model trained with rank = 15, lambda = 0.01, and numIter = 20.
RMSE (validation) = 1.3281800451766226 for the model trained with rank = 15, lambda = 1.0, and numIter = 10.
RMSE (validation) = 1.3281800690626133 for the model trained with rank = 15, lambda = 1.0, and numIter = 20.

The best model was trained with rank = 10 and lambda = 0.01, and numIter = 20, and its RMSE on the test set is 0.8628875199400583.
...
16/03/04 10:18:13 INFO FileOutputCommitter: Saved output of task 'attempt_201603041018_1009_m_000000_0' to file:/Carpeta_Proyecto/target/tmp/myCollaborativeFilter/data/user/_temporary/0/task_201603041018_1009_m_000000
16/03/04 10:18:13 INFO FileOutputCommitter: Saved output of task 'attempt_201603041018_1009_m_000001_0' to file:/Carpeta_Proyecto/target/tmp/myCollaborativeFilter/data/user/_temporary/0/task_201603041018_1009_m_000001
...

Con esto habremos conseguido definir nuestro modelo, salvarlo y así poder recuperarlo, como veremos en futuros capítulos, para por ejemplo hacer recomendaciones online.

Siguientes capítulos:

    [Spark MLlib] Filtrado Colaborativo (II): Recomendación por Usuario

    [Spark MLlib] Filtrado Colaborativo (III): Recomendación por Ítem (pendiente)

viernes, 26 de febrero de 2016

[Spark MLlib] Árboles de Decisión: Ejemplo de Clasificación

Al poco de adentrarnos en esto del mundo del Big Data y de los primeros conceptos con los que nos familiarizaremos serán sus famosas 3 V's: variedad, velocidad y volumen; aunque en seguida una cuarta saldrá a la luz, valor. Esta última es a la que personalmente suelo relacionar con Machine Learning. Es verdad que también se puede sacar valor del Big Data, por ejemplo, como mera capa de almacenamiento (volumen), pero las empresas hoy en día buscan descubrir, nutrirse y enriquecerse gracias al valor del dato. Mientras tanto, las tres primeras las suelo relacionar más con el mundo de los sistemas y de su arquitectura de los cuales provengo.

Así pues, el área del Machine Learning me queda lejos, bastante lejos, pero como es de costumbre en la informática toca familiarizarse con todo, formación continua e inacabable, como si no fuera suficiente con el ecosistema y las nuevas herramientas que emergen (Flink, NiFi, Arrow...). Pues bien, ahora y como culturilla me encuentro tratando de adquirir unos mínimos conocimientos al respecto. Y que mejor forma de hacerlo que dejar constancia de los ejercicios que vaya realizando, pues me gusta aprender a través de ejemplos, tratando, o creyendo más bien, haberlos entendido gracias a su explicación poco a poco, por lo que si algún data scientist o experto en el área lee alguna "burrada" en esta o futuras entradas que por favor me corrija.

Como podéis intuir por el título de la entrada he empezado con los famosos árboles de decisión. Se trata de un modelo predictivo cuyo punto fuerte recae sobre la sencillez en su representación gráfica así como por el formato de sus reglas en casi un lenguaje natural.

Bueno, voy a dejarme de tanta literatura y centrarme en el ejemplo, pues si queréis más información o detalle basta con que introduzcáis esos conceptos en vuestro buscador favorito y tendréis miles de entradas al respecto.

El ejemplo mostrado a continuación se encuentra implementado en lenguaje Scala (v2.11.7), sbt (v0.13.9) y hace uso además de la librería MLlib que incluye Spark (v1.6). Doy por hecho la correcta instalación de todas estas herramientas en vuestro entorno, pero si os puede ayudar en algo, preguntad en los comentarios. Por mi parte y en lo que respecta a esta entrada trataré de ser lo más genérico posible, abstrayéndome del sistema operativo o de la herramienta de desarrollo utilizados.

La finalidad del ejemplo consiste en construir un modelo predictivo que nos ayude a anticiparnos a conocer si un vuelo será retrasado o no, entendiendo por ello un retraso en la llegada superior a los 40 minutos. El ejemplo está basado en el de la entrada: MapR Blog - Apache Spark Machine Learning Tutorial.

1) Empezaremos creando una nueva carpeta para nuestro proyecto, Carpeta_Proyecto. Dentro de esta crearemos los siguientes subdirectorios: data y src/main/scala/mllib/DecisionTrees.

2) De la siguiente URL Bureau of Transportation Statistics-Airline On-Time Performance generaremos y descargaremos nuestro conjunto de datos.

En la parte superior podréis filtrar por geografía (estado), año y mes. También podéis seleccionar más o menos columnas que como veréis luego en el ejemplo, nos quedaremos con las que consideremos interesantes. Para este ejemplo se han de seleccionar las columnas: Dayofmonth, DayOfWeek, Carrier, Tailnum, FlightNum, OriginAirportID, OriginState, DestAirportID, DestState, CRSDepTime, DepTime, DepDelayMinutes, CRSArrTime, ArrTime, ArrDelay, CRSElapsedTime y Distance.

Ubicaremos el dataset recien generado en Carpeta_Proyecto/data. En mi caso como me descargué los datos filtrados por el año 2015 y mes de diciembre lo nombre ontime-201512.csv

3) En el directorio raíz del proyecto crearemos el archivo build.sbt:
import sbt._
import Keys._

name := "Carpeta_Proyecto"   // Reemplazar por el título del proyecto que le hayáis asignado

version := "1.0"

scalaVersion := "2.11.7"   // Configurar de acuerdo a vuestra versión de scala

libraryDependencies ++= Seq(
  "org.apache.spark" % "spark-core_2.11" % "1.6.0",   // Configurar de acuerdo a vuestra versión de scala y spark
  "org.apache.spark" % "spark-mllib_2.11" % "1.6.0"   // Configurar de acuerdo a vuestra versión de scala y spark
)

4) Por último, cread el fichero classificationFlightsOnTime.scala dentro de Carpeta_Proyecto/src/main/scala/mllib/DecisionTrees/ y copiar-pegar el código de la aplicación.

package mllib.DecisionTrees

import org.apache.spark.mllib.linalg.Vectors
import org.apache.spark.mllib.regression.LabeledPoint
import org.apache.spark.mllib.tree.DecisionTree
import org.apache.spark.{SparkConf, SparkContext}

object classificationFlightsOnTime {
  def main (args: Array[String]): Unit = {

    val sparkConfig = new SparkConf()
      .setMaster("local[2]")   // Configuración del clúster de Spark
      .setAppName("ClasificacionFlightsOnTime")   // Nombre otorgado a la app de Spark durante su ejecución
    val sc = new SparkContext(sparkConfig)  // Inicializar el necesario objeto de Spark Context

    val textRDD = sc.textFile("data/ontime-201512.csv")   // Asegurarse en poner la ruta y nombre correctos del fichero generado en el punto 2)
    val flightsRDD = textRDD.map(parseFlight).cache() // Gracias a la función parseFlight y clase Flight definiremos cual es el esquema (variables y tipos) de nuestro fichero CSV

    // A continuación crearemos una serie de asociaciones entre el valor de ciertas variables o columnas de tipo cadena con un identificador único, pudiendo reconvertir así más adelante el valor de dichas columnas de tipo string a entero.
    // La primera de estas asociaciones es para la columna Carrier
    var carrierMap: Map[String, Int] = Map()
    var index: Int = 0
    flightsRDD.map(flight => flight.carrier).distinct.collect.foreach(x => { carrierMap += (x -> index); index += 1 })

    // Idem para la columna OriginState del CSV que en la clase Flight se corresponde con el atributo origin
    var originMap: Map[String, Int] = Map()
    var index1: Int = 0
    flightsRDD.map(flight => flight.origin).distinct.collect.foreach(x => { originMap += (x -> index1); index1 += 1 })

    // Idem para DestState respecto del CSV o dest en la clase Flight
    var destMap: Map[String, Int] = Map()
    var index2: Int = 0
    flightsRDD.map(flight => flight.dest).distinct.collect.foreach(x => { destMap += (x -> index2); index2 += 1 })

    //  El siguiente paso será construir el array o vector con aquellas características que sean de nuestro interés. Es aquí por lo tanto donde filtraremos o quitaremos aquellas columnas que hayamos podido seleccionar de más a la hora de descargar nuestro dataset. También se trata del punto en donde convertiremos las columnas/abritutos Carrier, Origin y Dest de tipo cadena a entero, pues el vector finalmente que habremos de pasar al algoritmo tiene que ser de tipo numérico.
    val mlprep = flightsRDD.map(flight => {
      val monthday = flight.dofM.toInt - 1 // category
      val weekday = flight.dofW.toInt - 1 // category
      val crsdeptime1 = flight.crsdeptime.toInt
      val crsarrtime1 = flight.crsarrtime.toInt
      val carrier1 = carrierMap(flight.carrier) // category
      val crselapsedtime1 = flight.crselapsedtime.toDouble
      val origin1 = originMap(flight.origin) // category
      val dest1 = destMap(flight.dest) // category
      val delayed = if (flight.depdelaymins.toDouble > 40) 1.0 else 0.0    // Recordar nuestra definición de vuelo retrasado: si el retraso de éste era mayor a 40 minutos. Bien, aquí es donde transformaremos también ese retraso de minutos a verdadero o falso de acuerdo a dicha condición.
      Array(delayed.toDouble, monthday.toDouble, weekday.toDouble, crsdeptime1.toDouble, crsarrtime1.toDouble,
        carrier1.toDouble, crselapsedtime1.toDouble, origin1.toDouble, dest1.toDouble)
    })

    // Ahora nos encargaremos de asociar cada uno de los vectores anteriormente creados a lo que se denomina como LabeledPoint, que podríamos decir que se trata de la etiqueta que clasifica al vector en cuestión. En nuestro caso el LabeledPoint será 1.0 o 0.0 de acuerdo a si el vuelo fue retrasado o no, respectivamente. Gracias a esta acción construimos el conjunto de datos final.
    val mldata = mlprep.map(x => LabeledPoint(x(0), Vectors.dense(x(1), x(2), x(3), x(4), x(5), x(6), x(7), x(8))))

    // Partimos nuestro reciente creado conjunto de datos en los conjuntos de entrenamiento y test al 70% y 30% respectivamente.
    val splits = mldata.randomSplit(Array(0.7, 0.3))
    val (trainingData, testData) = (splits(0), splits(1))

    // Configuramos los parámetros que necesita el árbol de decisión
    var categoricalFeaturesInfo = Map[Int, Int]()    // Identificará que características del vector (segundo parámetro de la variable mldata) son categóricas y que rango de valores pueden tomar
    categoricalFeaturesInfo += (0 -> 31)    // El 0 se refiere a la posición 0 del vector que se corresponde con la columna monthday y 31 al rango de valores diferentes que puede tomar ésta 
    categoricalFeaturesInfo += (1 -> 7)    // En este caso, 1 igual columna weekday y lógicamente puede tomar 7 valores diferentes
    categoricalFeaturesInfo += (4 -> carrierMap.size)    // 4 igual a columna carrier1 y gracias a haber convertido sus valores de tipo string a entero, sabemos de antemano cuántos valores diferentes puede tomar
    categoricalFeaturesInfo += (6 -> originMap.size)    // Columna origin1
    categoricalFeaturesInfo += (7 -> destMap.size)    // Columna dest1

    // Resto de parámetros del árbol de decisión
    val numClasses = 2    // Vuelo retrasado o no
    val impurity = "gini"    // Tipo de medición de impurezas respecto a la homogeneidad de las etiquetas en el nodo
    val maxDepth = 4    // Profundidad del árbol
    val maxBins = 7000    // Número máximo de "bins"

    // Entrenamiento del modelo
    val model = DecisionTree.trainClassifier(trainingData, numClasses, categoricalFeaturesInfo, impurity, maxDepth, maxBins)

    // Visualización del árbol de decisión resultante
    println("Learned classification tree model:\n" + model.toDebugString)

    // Evaluación del modelo contra el conjunto de datos de test
    val labelAndPreds = testData.map { point =>
      val prediction = model.predict(point.features)
      (point.label, prediction)
    }

    // Cálculo del número de predicciones incorrectas
    val wrongPrediction =(labelAndPreds.filter{
      case (label, prediction) => ( label !=prediction)
    })
    println("Wrong predicions: " + wrongPrediction.count())

    // Cálculo del ratio de predicciones incorrectas
    val ratioWrong=wrongPrediction.count().toDouble/testData.count()
    println("Ratio wrong predictions: " + ratioWrong)
  }

    // Función para parsear las líneas del fichero en la clase Flights
  def parseFlight(str: String): Flight = {
    val line = str.split(",")
    Flight(line(0), line(1), line(2), line(3), line(4).toInt, line(5), line(6), line(7), line(8), line(9).toDouble,
      line(10).toDouble, line(11).toDouble, line(12).toDouble, line(13).toDouble, line(14).toDouble, line(15).toDouble,
      line(16).toInt)
  }
}

    // Definición del esquema de la clase Flights
case class Flight(dofM: String, dofW: String, carrier: String, tailnum: String, flnum: Int, org_id: String,
                  origin: String, dest_id: String, dest: String, crsdeptime: Double, deptime: Double,
                  depdelaymins: Double, crsarrtime: Double, arrtime: Double, arrdelay: Double, crselapsedtime: Double,
                  dist: Int)

5) Llega la hora de probar nuestra aplicación de Spark. Ejecutar desde la línea de comandos
# cd /ruta/a/Carpeta_Proyecto
# sbt compile
...
[success] Total time: 4 s, completed 26-feb-2016 10:27:14

# sbt run
...
Learned classification tree model:
DecisionTreeModel classifier of depth 4 with 31 nodes
  If (feature 0 in {0.0,1.0,2.0,3.0,4.0,5.0,6.0,7.0,8.0,9.0,10.0,11.0,12.0,13.0,14.0,15.0,16.0,17.0,18.0,19.0,20.0})
   If (feature 2 <= 1035.0)
    If (feature 4 in {0.0,1.0,2.0,3.0})
     If (feature 0 in {0.0,1.0,2.0,3.0,4.0,5.0,6.0,7.0,8.0,9.0,10.0,11.0})
      Predict: 0.0
     Else (feature 0 not in {0.0,1.0,2.0,3.0,4.0,5.0,6.0,7.0,8.0,9.0,10.0,11.0})
      Predict: 0.0
    Else (feature 4 not in {0.0,1.0,2.0,3.0})
     If (feature 1 in {0.0})
      Predict: 0.0
     Else (feature 1 not in {0.0})
      Predict: 0.0
   Else (feature 2 > 1035.0)
    If (feature 0 in {0.0,1.0,2.0,3.0,4.0,5.0,6.0,7.0,8.0,9.0,10.0,11.0})
     If (feature 0 in {0.0,1.0,2.0})
      Predict: 0.0
     Else (feature 0 not in {0.0,1.0,2.0})
      Predict: 0.0
    Else (feature 0 not in {0.0,1.0,2.0,3.0,4.0,5.0,6.0,7.0,8.0,9.0,10.0,11.0})
     If (feature 4 in {0.0,1.0,2.0,3.0})
      Predict: 0.0
     Else (feature 4 not in {0.0,1.0,2.0,3.0})
      Predict: 0.0
  Else (feature 0 not in {0.0,1.0,2.0,3.0,4.0,5.0,6.0,7.0,8.0,9.0,10.0,11.0,12.0,13.0,14.0,15.0,16.0,17.0,18.0,19.0,20.0})
   If (feature 2 <= 1120.0)
    If (feature 2 <= 903.0)
     If (feature 0 in {21.0,22.0,23.0,24.0,25.0})
      Predict: 0.0
     Else (feature 0 not in {21.0,22.0,23.0,24.0,25.0})
      Predict: 0.0
    Else (feature 2 > 903.0)
     If (feature 0 in {21.0,22.0,23.0,24.0,25.0,26.0,27.0,28.0,29.0})
      Predict: 0.0
     Else (feature 0 not in {21.0,22.0,23.0,24.0,25.0,26.0,27.0,28.0,29.0})
      Predict: 0.0
   Else (feature 2 > 1120.0)
    If (feature 0 in {21.0,22.0,23.0,24.0,25.0,26.0,27.0,28.0,29.0})
     If (feature 0 in {21.0,22.0,23.0,24.0,25.0})
      Predict: 0.0
     Else (feature 0 not in {21.0,22.0,23.0,24.0,25.0})
      Predict: 0.0
    Else (feature 0 not in {21.0,22.0,23.0,24.0,25.0,26.0,27.0,28.0,29.0})
     If (feature 7 in {1.0,2.0,3.0,4.0,5.0,6.0,7.0,8.0,9.0,10.0,11.0,12.0,13.0,14.0,15.0,16.0,17.0,18.0,19.0,20.0,21.0,22.0,23.0,24.0,25.0,26.0,27.0,28.0,29.0,30.0,31.0,32.0,33.0,34.0,35.0,36.0,37.0,38.0,39.0,40.0,41.0,42.0,43.0,44.0,45.0,46.0})
      Predict: 0.0
     Else (feature 7 not in {1.0,2.0,3.0,4.0,5.0,6.0,7.0,8.0,9.0,10.0,11.0,12.0,13.0,14.0,15.0,16.0,17.0,18.0,19.0,20.0,21.0,22.0,23.0,24.0,25.0,26.0,27.0,28.0,29.0,30.0,31.0,32.0,33.0,34.0,35.0,36.0,37.0,38.0,39.0,40.0,41.0,42.0,43.0,44.0,45.0,46.0})
      Predict: 0.0
...
Wrong predicions: 14141
...
Ratio wrong predictions: 0.09830378866875217
...

5.1) Ahora por ejemplo establecemos la profundidad del árbol a 3:
...
Wrong predicions: 14207
...
Ratio wrong predictions: 0.09855022197558269

5.2) Profundidad del árbol a 7:
...
Wrong predicions: 14150
...
Ratio wrong predictions: 0.09854995751556601

5.3) Y a 9:
...
Wrong predicions: 14193
...
Ratio wrong predictions: 0.09897627581974644

Bueno, pues ya tendríamos construido nuestro modelo, aunque con la tasa o ratio de acierto reflejado... me deja muchas dudas de su utilidad, todo sería afinarlo más, ¿mayor o menor profundidad? De acuerdo a los valores de ratio observados, podríamos establecer un valor bajo, pues son prácticamente idénticos y ante un mayor nivel de profundidad lo único que se consigue es complicarnos la comprensión del árbol de decisión, ¿otras variables? Es posible. He ahí el buen hacer o no de los data scientist y de sus modelos. ¡Ah! Y de la CALIDAD de los datos, ese "pequeño e insignificante" detalle ;)

Por cierto, si el modelo resultante fuera considerado correcto y se quisiera hacer uso de él en el futuro, podríamos salvarlo a disco, en este ejemplo model.save(sc, "target/tmp/myDecisionTreeClassificationFlightsOnTime"), y para recuperarlo, val model = DecisionTreeModel.load(sc, "target/tmp/myDecisionTreeClassificationFlightsOnTime"), evitando así el tener que volver a entrenar el modelo.

Espero que os haya servido o al menos para haceros una idea de cómo trabajar con Spark MLlib y con árboles de decisión.