scieee AI-readable full text Open interactive document viewer

Predicción de series temporales con k-NN sobre Spark

Águila Sánchez, Sebastián; Montero Cobo de Guzmán, Sergio

Abstract

Este trabajo ha tenido como objetivo la implementación de un algoritmo de predicción k-NN para series temporales, haciendo uso del paradigma de procesamiento distribuido Spark, con el fin conseguir una solución Big Data. Conseguir esta solución es muy interesante, puesto que es un algoritmo de aprendizaje basado en instancias no genera un modelo para realizar las predicciones, sino que se apoya de un histórico. Un histórico lo suficientemente grande podría desbordar la memoria de cualquier computador. Con el uso del algoritmo distribuido, se puede procesar sin limitaciones este histórico por muy grande que sea, permitiendo tener en cuenta más datos para componer la solución, obteniendo mejores predicciones que las que se obtendrían usando una solución mono-máquina con un histórico limitado.

Full text

Universidad Complutense de Madrid Facultad de Informática Predicción de series temporales con k-NN sobre Spark Alumnos Sebastián Águila Sánchez Sergio Montero Cobo de Guzmán Director del proyecto Javier Arroyo Gallardo Codirector del proyecto Albert Meco Alías Trabajo de Fin de Grado Ingeniería de Software Septiembre de 2018 2 Índice general Palabras clave 5 Resumen 7 1. Introducción 9 1.1. Importancia de las soluciones Big Data en algoritmos de aprendizaje automático . 9 1.2. Aprendizaje basado en instancias . . . . . . . . . . . . . . . . . . . . . . . . . . 10 1.3. Motivación ..................................... 10 1.4. Objetivos ...................................... 12 2. Contexto tecnológico 15 2.1. ApacheHadoop................................... 15 2.1.1. HDFS.................................... 16 2.1.2. YARN.................................... 17 2.1.3. MapReduce................................. 17 2.2. ApacheSpark.................................... 17 2.2.1. Arquitectura interna . . . . . . . . . . . . . . . . . . . . . . . . . . . . 18 2.2.2. Flujo de una aplicación Spark . . . . . . . . . . . . . . . . . . . . . . . 19 2.2.3. Tipos de datos y objetos . . . . . . . . . . . . . . . . . . . . . . . . . . 24 2.2.4. Ventajas de Spark respecto a MapReduce . . . . . . . . . . . . . . . . . 25 2.3. Presentación de la tecnología empleada . . . . . . . . . . . . . . . . . . . . . . 27 2.4. Recursos utilizados en el desarrollo . . . . . . . . . . . . . . . . . . . . . . . . 28 2.4.1. Ordenadorportátil ............................. 28 2.4.2. Máquinas.................................. 28 2.4.3. ClústerHadoop............................... 28 2.4.4. AmazonEMR ............................... 29 3. Algoritmo k-Nearest Neighbours 31 3.1. Descripción del Algoritmo . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 32 3.1.1. Parámetros de predicción . . . . . . . . . . . . . . . . . . . . . . . . . . 32 3.1.2. Medida de similitud . . . . . . . . . . . . . . . . . . . . . . . . . . . . 33 3.1.3. Generación de predicciones . . . . . . . . . . . . . . . . . . . . . . . . 33 3.2. Entrenamiento y validación . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 34 3 4ÍNDICE GENERAL 3.2.1. Medidadelerror .............................. 35 4. Implementación y optimización del algoritmo 37 4.1. Implementación paralela del algoritmo . . . . . . . . . . . . . . . . . . . . . . . 37 4.1.1. Preparación de los datos . . . . . . . . . . . . . . . . . . . . . . . . . . 38 4.1.2. Búsqueda de los kvecinos ......................... 41 4.1.3. Función de reducción . . . . . . . . . . . . . . . . . . . . . . . . . . . . 42 4.2. Funcionesprincipales................................ 43 4.2.1. KNN_Next................................. 43 4.2.2. KNN-Past.................................. 45 4.2.3. KNN_Optim ................................ 46 5. Experimentos y comparación con la implementación mono-máquina 51 5.1. Series temporales consideradas . . . . . . . . . . . . . . . . . . . . . . . . . . . 51 5.1.1. Bitcoin Historical Data . . . . . . . . . . . . . . . . . . . . . . . . . . . 51 5.1.2. SunspotDaily................................ 51 5.1.3. SunspotMonthly.............................. 52 5.2. Alcance de la solución mono-máquina en R vs solución Spark . . . . . . . . . . 52 5.3. Escenarios...................................... 52 5.3.1. Seriepequeña................................ 52 5.3.2. Seriemedia................................. 54 5.3.3. SerieGrande ................................ 58 5.4. Conclusiones .................................... 60 6. Trabajo Futuro 61 Bibliografía 63 Agradecimientos 65 Palabras clave Palabras clave en Español Regresión Predicción k-NN Series temporales Spark RDD Keywords in English Regression Prediction k-NN Time Series Spark RDD 5 6Palabras clave Resumen Resumen en Español Este trabajo ha tenido como objetivo la implementación de un algoritmo de predicción k-NN para series temporales, haciendo uso del paradigma de procesamiento distribuido Spark, con el fin conseguir una solución Big Data. Conseguir esta solución es muy interesante, puesto que es un algoritmo de aprendizaje basado en instancias no genera un modelo para realizar las predicciones, sino que se apoya de un histórico. Un histórico lo suficientemente grande podría desbordar la memoria de cualquier computador. Con el uso del algoritmo distribuido, se puede procesar sin limitaciones este histórico por muy grande que sea, permitiendo tener en cuenta más datos para componer la solución, obteniendo mejores predicciones que las que se obtendrían usando una solución mono-máquina con un histórico limitado. Summary in English The objective of this work was to implement a k-NN forecasting algorithm for time series, using the distributed processing paradigm Spark, in order to achieve a Big Data solution. This solution is very useful, since it’s an instance based learning algorithm and it does not generate a model to make predictions, instead it uses historical data. With enough historical data, it is possible to overflow the memory of any computer.. With the use of the distributed algorithm, this historical can be processed without limitations, no matter how big it is, allowing to have more data to compose the solution, obtaining better predictions than what would be obtained using a mono-machine solution with a limited historical. 7 8Resumen Capítulo 1 Introducción La búsqueda de la capacidad para explotar correctamente los datos con el fin de sacarles un valor presenta un grave problema, la gran cantidad y diversidad de datos que tenemos al alcance hacen imposible su procesamiento haciendo uso de la computación tradicional mono-máquina, debido a que superan la capacidad de computación y memoria de un cualquier computador, por muy potente que sea. 1.1. Importancia de las soluciones Big Data en algoritmos de aprendizaje automático Dado un escenario en el que sea de vital importancia procesar la mayor cantidad de datos posibles, se manifiesta la necesidad de tener un sistema que emplee una computación distribuida con una escalabilidad horizontal sencilla y económica, que permita realizar el procesamiento de la totalidad de la información accesible con el fin de intentar sacarle valor. Con el fin de intentar sacar valor a los datos, es habitual que se analicen con algoritmos de aprendizaje automático, estos permiten extraer comportamientos o tendencias que realizando un buen uso de esta información consiguen añadir valor adicional a los datos. En un escenario Big Data en el que hay infinidad de volumen de datos que necesitan ser procesados y adaptados con velocidad, es imposible usar estos algoritmos debido a que el escenario sobrepasa la capacidad de computación y de memoria de la máquina que realiza el procesamiento. Este es el motivo por el cual los algoritmos de aprendizaje automático están siendo integrados con soluciones de computación distribuida en entornos Big Data, este tipo de soluciones permiten mejorar las características de los algoritmos de aprendizaje automático, haciéndolos más robustos y potentes, permitiendo procesar más instancias en las predicciones usando una cantidad ilimitada de ejemplos de entrenamiento. Al usar una computación distribuida el procesamiento se paraleliza en varias máquinas consiguiendo generan las predicciones en el menor tiempo posible. Cualquier tipo de algoritmo de aprendizaje automático se puede beneficiar de la capacidad de procesamiento y la paralelización que ofrece esta tecnología, pero solo tiene sentido aplicar este tipo de computación si hay una enorme cantidad de datos. 9 16 CAPÍTULO 2. CONTEXTO TECNOLÓGICO Figura 2.1: Componentes del ecosistema de Hadoop con su jerarquía de dependencias. 2.1.1. HDFS Es un componente que engloba la capa de almacenamiento distribuido, perteneciente a Hadoop Common, que es la base de Hadoop y se encarga de proveer un sistema de ficheros distribuido, escalable, seguro y tolerante a fallos [8]. Hay que tener en cuenta que no es un sistema de ficheros con propósito general, el conjunto de operaciones permitidas está limitado, no se pueden modificar los datos y esto es debido a que el sistema está pensado para hacer lecturas. Entre muchas de sus virtudes caben destacar: Mejora el rendimiento de los motores de procesamiento distribuido del ecosistema de Hadoop como MapReduce oSpark, esto es debido a que el procesamiento se lleva a la máquina donde se encuentra físicamente el dato. Este fenómeno se le conoce como data locality, consiste en elegir nodos que tiene el dato físicamente dentro del clúster para ser los que realicen el procesamiento, con ello se evitan latencias y transferencias de red. Tiene jerarquía de directorios como el sistema tradicional de Linux, con su gestión de permisos y usuarios por directorio. Permite el almacenamiento de enormes cantidades de datos, con un almacenamiento fácilmente escalable y barato, basta con añadir más máquinas o ampliar la memoria de las ya existentes. Es un sistema de ficheros fiable gracias a la replicación de los datos que evitan su pérdida y mejora el rendimiento de las lecturas. 2.2. APACHE SPARK 17 2.1.2. YARN Es un gestor de recursos del clúster para motores de procesamiento distribuidos, puede ser usado además de por MapReduce por diferentes motores de procesamiento como Spark u otros componentes del ecosistema de Hadoop [8]. Estos componentes a través de una API solicitan los recursos que necesitan del clúster para llevar a cabo sus trabajos. Además monitoriza el uso de los recursos y los trabajos que están siendo ejecutados dentro del clúster. 2.1.3. MapReduce Aunque finalmente no se ha elegido como motor de procesamiento distribuido para la implementación del algoritmo, es interesante darle una mención. Los motivos de su descarte para la implementación del algoritmo se exponen en las secciones 2.2.4 y 2.3. MapReduce surge como solución al procesamiento de grandes volúmenes de datos que no podían ser procesados con un modelo mono-máquina. Es un framework para el procesamiento de algoritmos paralelizables y distribuidos. Pertenece a Hadoop Commons. Usa el paradigma de programación MapReduce. Este paradigma está inspirado por la estrategia de programación Divide y Vencerás, que consiste en dividir el problema original en n subproblemas más sencillos para resolverlos de forma individual y luego combinar todas las soluciones [8]. MapReduce consta de tres fases: 1. Map: se paralelizan y se distribuyen los datos en formato clave-valor en cada máquina y se les aplica una función. 2. Shuffle & Short: se aplica a la salida de la fase Map un algoritmo de shuffle and sort que se encarga de agrupar los datos por clave para que sirvan como entrada a la fase Reduce. 3. Reduce: Es el último paso, en la fase Reduce, por clave se combinan las salidas de los Map aplicando una función determinada. Un trabajo MapReduce solo requiere implementar la parte Map, la parte Reduce no es obligatoria, debido a que en muchos contextos no es necesario hacer una reducción de los datos. Por ejemplo, si queremos implementar un algoritmo que haga el volcado de una base de datos al sistema de ficheros HDFS, solo tendríamos que implementar la parte Map que se encargue de sacar los datos y los mande al sistema de ficheros. MapReduce permite cualquier entrada de datos, sea de componentes de Hadoop como HDFS o componentes externos como bases de datos. 2.2. Apache Spark Spark es un motor de procesamiento distribuido con un conjunto de librerías para el procesamiento paralelo en un clúster, es open-source, escalable y tolerante a fallos [7]. Surge debido a la necesidad de implementar un motor de procesamiento distribuido que fuera eficiente para tareas en las que se hagan varias iteraciones sobre un conjunto de datos, estas tareas 18 CAPÍTULO 2. CONTEXTO TECNOLÓGICO usando MapReduce son totalmente ineficientes puesto que para cada iteración hay que lanzar otra tarea MapReduce, véase figura 2.5. Tiene su propio motor de procesamiento distribuido, no usa MapReduce como hacen otros elementos del ecosistema de Hadoop como Hive o Pig. Este motor tiene la capacidad de mantener en memoria gran cantidad de datos. Algunos de los motivos para elegir Spark son: Trabajar con Spark resulta atractivo puesto que se adapta a un abanico bastante amplio de usuarios como programadores o analistas. Aparte de Java soporta más variedad de lenguajes que MapReduce, como Python, Scala o R. Este último lenguaje presenta ciertos problemas debido a que no soporta la totalidad de la funcionalidad de Spark. Permite trabajar con los datos al estilo de tablas de SQL gracias a las librerías de SparkSQL. Permite el procesamiento de datos en near-real-time usando las librerías de SparkStreaming, este tipo de procesamiento consiste en : dado un dato en Streaming realizar su procesamiento en un instante de tiempo inferior a un minuto. Incluye librerías de machine learning como MLlib o bibliotecas de procesamiento gráfico como GraphX. Soporta la mayoría de sistemas de almacenamiento distribuidos como HDFS, Amazon S3, Azure Storage. Cuenta con una gran comunidad de desarrolladores que mantienen y mejoran el proyecto con nuevas librerías y mejoras. 2.2.1. Arquitectura interna Las aplicaciones de Spark están formadas por un proceso Driver (es el que gestiona todo el trabajo de Spark) y un conjunto de ejecutores (son un conjunto de procesos que realizan el procesamiento distribuido y son gestionados por el Driver). El Driver ejecuta la función principal del programa dentro de un nodo y se encarga de hacer un análisis de los trabajos, distribuir y programar la ejecución del trabajo Spark en los ejecutores. Spark usa un gestor de recursos (YARN, Mesos, Standalone ...) para aprovisionar los ejecutores necesarios para realizar el procesamiento de los datos. Apache Spark soporta varios lenguajes de programación pero su forma de trabajar con ellos difiere en función de si puede ser ejecutado en una JVM (Java Virtual Machine) o no. Java y Scala pueden ser ejecutados directamente pero el resto de los lenguajes soportados (Python y R) se traducen a un código que puede ser interpretado por las JVM de los ejecutores. Modos de despliegue Las aplicaciones Spark se distribuir a lo largo del clúster apoyándose en gestores de recursos. algunos de los gestores de recursos son: 2.2. APACHE SPARK 19 Standalone mode: Es un gestor de recursos creado específicamente para Apache Spark, es una buena opción si en el clúster solo se ejecutan tareas de Spark debido a que es sencillo y rápido, pero si en el clúster tenemos una administración y una gestión concreta de los trabajos con prioridades o límites de recursos y además se ejecutan otros procesos distribuidos que no sean de Spark, no sería una buena elección. Apache Hadoop YARN: En apartados anteriores se habló de este gestor de recursos, las ventajas que presenta es que está integrado por defecto en un clúster Hadoop y su complejidad permite que coexistan Spark con otros motores de procesamiento. Modos de ejecución Las aplicaciones Spark se pueden ejecutar de 3 formas: Modo Local: Es el modo de ejecución en la propia máquina local, el paralelismo se realiza mediante hilos, es el modo indicado para hacer pruebas o desarrollos de una forma rápida. Modo Cliente: En este modo la aplicación despliega en el propio cliente el Driver y es el cliente el que mantiene el proceso, si se cae la conexión del cliente con el clúster la aplicación se cierra. Esta es la desventaja de este modo pero este el único modo distribuido que soporta Spark Shell, esta herramienta es muy útil para realzar el desarrollo y las pruebas ya que permite hacer sentencias interactivas de una forma muy rápida. Modo Clúster: En este modo el cliente despliega un ejecutable dependiente del lenguaje de programación. El gestor de recursos lanza en un nodo trabajador el proceso Driver de Spark y es el gestor el que tiene que mantener el proceso. En este modo de ejecución no necesitamos que esté activo el cliente durante todo el proceso. 2.2.2. Flujo de una aplicación Spark El código que lanzamos en Spark lanza una aplicación de Spark, la aplicación esta por compuesta por uno o varios Jobs que generalmente se ejecutan uno detrás de otro. Estos Jobs están compuestos por Stages, que a su vez se dividen en Tasks. Spark automáticamente analiza Stages y Tasks con el fin de realizar el procesamiento de una forma óptima. Para entender el flujo de una aplicación Spark es importante explicar los siguientes elementos de Spark: Jobs Un trabajo en Spark es un conjunto de Stages, habrá tantos Stages como shuffles (reparticiones) se hagan con los datos. SparkSession Es el punto de entrada de la aplicación desde Spark 2.0, para lanzar un proceso Spark es necesario obtener la instancia desde el Driver de este objeto. Desde este objeto se hacen las lecturas 20 CAPÍTULO 2. CONTEXTO TECNOLÓGICO del conjunto de datos de entrada en el formato que se prefiera RDD, DataFrame ... Particiones Cuando hacemos una lectura de los datos con en el Driver con el objeto SparkSession para fines de paralelización distribuida a lo largo del clúster se realizan particiones sobre el conjunto de los datos de entrada. Las particiones se distribuyen en los distintos ejecutores que se lanzan con la aplicación. Cada partición se procesa en un Task del ejecutor, un Task no es más que un procesamiento paralelo que ocurre en un ejecutor. Es conveniente tener al menos las mismas particiones que ejecutores. Si los datos son muy grandes es conveniente añadir más particiones debido a que con la configuración por defecto el shuffle no puede superar más de 2 Gb. Las particiones son objeto de optimización en los trabajos de Spark. Transformaciones En Spark los objetos son inmutables, es decir que no se pueden modificar, en vez de esto se realizan una lista de modificaciones para posteriormente crear otro objeto derivado con los cambios, a eso se le llama transformación. La generación de estos nuevos objetos llevan un modo de ejecución lazy, se siguen añadiendo modificaciones a la lista de cambios (plan de ejecución) sin generar el nuevo objeto hasta que no se necesita forzosamente. En jerga de Spark, no se genera hasta que no ocurre una Action, mientras que no ocurre todas las transformaciones se realizan en el mismo Stage. Las transformaciones van generando un plan de ejecución que no es ejecutado hasta que ocurre una Action. Las transformaciones dan lugar a dos tipos de dependencias: Narrow dependency : esta dependencia surge cuando las transformaciones se realizan sobre las particiones 1 a 1 y no necesitan repartición de los datos, se van añadiendo a él plan de ejecución. Wide dependency : Exigen de una repartición de los datos, desemboca a realizar una Action que despliega todo el plan de ejecución dando como resultado un nuevo Stage. Spark con su plan de ejecución optimiza todas las transformaciones, pero no tiene en cuenta los datos que se encuentran entre Stages. En caso de reutilización es conveniente hacer un cacheo. Actions Es el desencadenante que fuerza a el lanzamiento del plan de ejecución con las transformaciones que están previstas sobre los datos para un Stage. Se produce cuando hay una Wide dependency en el plan de ejecución que fuerza a generar los datos. Operaciones como count, collects o escribir el resultado de las transformaciones generan una Action. Las Actions siempre devuelven resultados. 2.2. APACHE SPARK 21 Figura 2.2: Captura de una parte del flujo del algoritmo. Como podemos observar en la figura 2.2 cada vez que se realiza un partitionBy (es una Action) este da lugar a un nuevo Stage, esto se debe a que la función partitionBy necesita realizar una repartición de los datos, esta función genera una Wide dependency. Esta figura es una visualización de un DAG generado en un trabajo de Spark, representa el flujo previsto por el planificador de tareas, los cuadros azules representan las funciones que se ejecutan y las líneas representan las transformaciones de los datos. Usar esta información que proporciona la visualización del DAG es muy útil para optimizar el código ya que nos dan información sobre funciones que se realizan en un mismo Stage o Stages que dependen unos de otros, si somos capaces de interpretar correctamente la visualización se pueden evitar cuellos de botella innecesarios y se puede optimizar el código haciendo un uso adecuado de las funciones, todas las funciones que se realizan en un mismo Stage son optimizadas automáticamente pero las que están entre Stages no se optimizan, es tarea del desarrollador visualizar y optimizar los procesos que ocurren entre Stages. Stages Son un conjunto de tareas que se puede o transformaciones que se puede llevar a cabo sin realizar una repartición de los datos, esto ocurre mientras que no se tenga un Wide dependency que fuerce a realizar un shuffle de los datos. Acumulando estas tareas se optimizan los cálculos evitando reprocesamientos y reparticiones innecesarias. 22 CAPÍTULO 2. CONTEXTO TECNOLÓGICO Figura 2.3: Secuencia con la generación de un Stage. En la figura 2.3 podemos observar que se generan tantas Tasks como particiones tenga el RDD, estas serán N. Todas estas Tasks (1 ... N) se van procesando según van quedando huecos libres en los ejecutores en un mismo Stage. Como podemos observar en la figura cada ejecutor es capaz de procesar tres tareas de forma paralela, teniendo dos ejecutores como los de la imagen, podemos paralelizar hasta 6 tareas, cada tarea procesa una partición de los datos. Al terminar procesarse estas Task que se visualizan en la imagen (Tasks 1-6), irían entrando las siguientes Tasks 7-N, sin seguir ningún orden pero manteniendo el procesamiento máximo de seis Tasks de forma paralela. El factor de paralelismo (cantidad de tareas paralelas) de cada ejecutor depende del número vCores que tenga asignados, en este ejemplo como máximo cada ejecutor puede procesar tres Task por lo tanto podemos deducir que cada ejecutor cuenta con 3 v-Cores. En el ejemplo de la figura 2.4 se puede observar mejor esta relación de v-Cores y Tasks. Task Las tareas son operaciones entre bloques de datos, estas transformaciones ocurren en un mismo ejecutor que procesa una partición de los datos, cuantas más particiones se tengan más paralelismo habrá y mejores tiempos obtendremos. 2.2. APACHE SPARK 23 Figura 2.4: Distribución de las Tasks que confluyen de forma concurrente en un clúster La figura 2.4 representa dos Jobs que están haciendo uso del clúster de forma concurrente. El clúster esta formado por tres máquinas, dos con 4v-Cores y una con 6v-Cores. Como podemos deducir de la imagen, los Jobs usan configuraciones diferentes, para Job1 se están usando ejecutores de 4v-Cores en cambio para Job2 de 2v-Cores. Se puede observar que la cantidad máxima de Tasks que se pueden ejecutar en un Job de forma paralela es su número de ejecutores multiplicado por el número de v-Cores por ejecutor. Dentro de los ejecutores se lanza una JVM con tantos procesos paralelos como vCores tenga la configuración de ejecutor. Es importante cuadrar bien la configuración de memoria y vCores por ejecutor para aprovechar mejor los recursos del clúster. Podemos observar en la imagen que si el Job2 hubiese usado ejecutores de 2vCores en vez de 4vCores, hubiera cambiado la distribución de ejecutores teniendo en la máquinas 2 y 3 dos ejecutores de 2v cores en vez de uno de 4v-cores, permitiendo desplegar un ejecutor adicional en la máquina 1 (solo tenia 2v-Cores disponibles,ahora si hay recursos para la nueva configuración de ejecutor), permitiendo procesar hasta 10 Tasks de forma paralela en vez de 8, aprovechando mejor los recursos del clúster. Por otro lado tenemos el Job2 que no ha tenido configurado correctamente el número de particiones de los datos, esto se observa en el ejecutor 3-1, solo procesa una tarea a pesar de que su configuración le hubiera permitido procesar dos. Esto no permite aprovechar correctamente los recursos del clúster incrementando los tiempos de procesamiento, es necesario ajustar correctamente las particiones en función de la configuración de los ejecutores. Usar una configuración con muchos vCores fuerza a tener menos ejecutores que tienen gran cantidad de memoria. Una configuración con muchos vCores nos ofrece las ventajas del paralelismo en una JVM como el uso de memoria compartida pero si contamos con demasiada memoria por ejecutor, el Garbage Collector ralentiza demasiado el proceso al gestionar una gran cantidad 24 CAPÍTULO 2. CONTEXTO TECNOLÓGICO de memoria. Usar una configuración con muchos ejecutores fuerza a inutilizar gran parte de la memoria empleándolas en las JVM que tiene que gestionar cada ejecutor en vez de usarlas para el procesamiento, con esta configuración es más sencillo hacer uso de todos los recursos del clúster pero no consigue un aprovechamiento óptimo. Hay que buscar un equilibrio entre ambas para conseguir una buena gestión de los recursos y un buen aprovechamiento, evitando que gestione demasiada memoria por JVM e intentando tener las menos JVM posibles para evitar replicar variables de Broadcast en cada JVM y gastar memoria en la gestión del proceso. En nuestro caso particular utilizamos variables de Broadcast bastante grandes y en ellas guardamos toda la parte de las instancias temporales que integran la parte de entrenamiento (véase figura 3.1), esta variable de Broadcast se distribuye a cada ejecutor. La configuración más adecuada a utilizar en este algoritmo es tener menos ejecutores con más vCores usando un número de particiones que garantice el aprovechamiento del clúster, una configuración correcta sería : noparticiones =noejecutores ∗vCoresejecutor. Menos ejecutores implica menos réplicas de la variable de Broadcast en el clúster (una réplica por ejecutor). Este es un claro ejemplo que refuerza la idea que se presentó en la introducción, los algoritmos distribuidos dependen en gran medida de su configuración. 2.2.3. Tipos de datos y objetos RDD Es la abstracción básica que tiene Spark de un conjunto de datos distribuido en varias máquinas, es tolerante a fallos y se puede cachear para optimizar las aplicaciones. DataFrame Un DataFrame es un tipo de datos que representan la estructura de una tabla que tiene filas y columnas, pertenecen a la librería SparkSQL. Tiene su esquema con el tipo de dato de cada una de las columnas. Simplifica las operaciones representando los RDD como si fueran una tabla. Accumulators Los acumuladores son variables compartidas que solo tienen definida una operación de add, como contadores. Broadcast Variables Una variable de Broadcast es una variable que se serializa y que se envía a cada ejecutor para que las tareas puedan usarla en el caso de necesitarlas, son similares al Distributed Cache de MapReduce. Generalmente se almacenan en memoria, pero en el caso de que no hubiera memoria se escribirían en disco. 2.2. APACHE SPARK 25 2.2.4. Ventajas de Spark respecto a MapReduce A la hora de elegir un framework de procesamiento distribuido para el algoritmo k-NN hemos tenido en cuenta las siguientes ventajas que nos ofrece Apache Spark respecto a MapReduce: Para trabajos iterativos sobre un conjunto de datos en los que se reutiliza información es bastante más rápido que MapReduce. Gracias al planificador de tareas DAG, se puede trazar y monitorizar de una forma precisa los trabajos que realiza por aplicación. Esto es muy útil a la hora de optimizar el código. Es más sencillo que MapReduce para el programador, además de Java soporta varios lenguajes de programación como Python, Scala o R. Permite trabajar de forma sencilla con datos en memoria, los cacheos de datos que se reutilizan mejoran los tiempos de procesamiento del algoritmo. Desarrollar con Spark es muy rápido ya que cuenta con un modo interactivo que a través de una shell se pueden desarrollar y testear funciones de forma rápida y sencilla. Otras características que ofrece Apache Spark y que no han sido usadas en el proyecto son: Tener librerías que facilitan la programación como SparkSQL que permite hacer consultas al estilo SQL. Contar con librerías como MLlib para tipos de datos y algoritmos usados en machine learning. Es importante mencionar que en un intento de reutilizar partes del algoritmo ya implementadas en R [1], pudimos comprobar que la API de Spark en R esta limitada a el uso de SparkSQL y a trabajar con funciones de tabla, no permite trabajar con los datos a bajo nivel, esto nos hizo descartar la opción de usar R y reutilizar parte de la solución implementada por nuestros compañeros. Para nuestro caso en particular, que realizamos procesamientos iterativos sobre un conjunto de datos (véase la función de la figura 4.9) Spark ha sido la mejor opción pero debemos tener en cuenta que en otros contextos en los que se den procesamientos lineales simples de grandes cantidades de datos lo mejor es usar MapReduce. 32 CAPÍTULO 3. ALGORITMO K-NEAREST NEIGHBOURS 3.1. Descripción del Algoritmo Como hemos comentado en anterioridad, este algoritmo k-NN predice sobre una serie temporal [3]. Para realizar las predicciones se realizan los siguientes pasos [2]: 1. Se preparan los datos, agrupando dinstantes temporales del pasado para poder calcular la distancia entre instantes temporales. Cada instante temporal se representa por sus d últimos valores. De esa forma una serie temporal se transforma en una matriz de instantes temporales descritos cada uno por d valores (o variables). 2. Calculamos las distancias de cada todos los instantes temporales que queremos predecir contra todos sus instantes del pasado usando medidas de similitud, estas distancias serán almacenadas en lo que conoceremos como matriz de distancias. 3. Elegimos en la matriz de distancias los kinstantes del pasado con menor distancia (k vecinos) a los instantes temporales que queremos predecir. 4. Una vez obtenemos los vecinos, usamos el valor de predicción de los kvecinos más cercanos y calculamos una media ponderada que dará como resultado la predicción al instante temporal que estamos evaluando. Una vez completados estos pasos hay que entender cómo funcionan exactamente los cálculos para completar el algoritmo. Este usa exactamente dos que son: Distancias o medidas de similitud: Se utilizan para encontrar los conjuntos más parecidos (vecinos) al instante temporal que se está evaluando para poder predecir sobre los mismos revisando todo el pasado al instante evaluado. Generación de la predicción: Ya obtenidos los vecinos más cercanos, es decir, los kinstantes en el pasado de la serie más parecidos al instante actual que estamos evaluando, utilizamos el siguiente valor de cada uno de los kinstantes y realizamos una media ponderada por la cercanía para obtener una predicción calculada con la predicción de los kvecinos. A continuación, hablaremos de estos cálculos que se han implementado con más en profundidad. 3.1.1. Parámetros de predicción Cuando aplicamos el algoritmo k-NN han de fijarse los parámetros kyda los cuales haremos referencia durante todo el capítulo el significado de cada parámetro es el siguiente: Parámetro k: Este parámetro fija el número de vecinos más cercanos al instante temporal que estemos evaluando, estos kvecinos servirán para realizar la predicción de cada instante temporal en función de la media de los valores de predicción cada kvecino. Parámetro d: Sirve para preparar los datos para el algoritmo k-NN, fijar cuantos el número de retardos, es decir, valores anteriores y consecutivos, que usamos para describir cada instante temporal. 3.1. DESCRIPCIÓN DEL ALGORITMO 33 3.1.2. Medida de similitud Para conocer los elementos más cercanos al los conjuntos que estamos evaluando usamos una medida de similitud con la cual generamos una matriz de distancias para poder encontrar en ella los vecinos más cercanos a uno dado, en nuestra implementación quedan recogidas tres tipos de distancia que son: Distancia Euclídea: Es la medida de distancia más común a utilizar, se deduce del teorema de Pitágoras, calcula la distancia más corta entre dos puntos X y Y de un espacio n-dimensional, su fórmula matemática es: d(X, Y ) = v u u t n X i=1 (xi−yi)2(3.1) Distancia Manhattan: Es una distancia relacionada con la distancia euclídea pero en vez de hacer la distancia más corta en diagonal, realiza la medida por manzanas, es decir, calcula la distancia entre dos puntos X y Y de un espacio n-dimensional en función de cuadrículas equivalentes, su fórmula matemática es: d(X, Y ) = n X i=1 |xi−yi|(3.2) Distancia Canberra: Es un medida derivada de la distancia de Manhattan, calcula la distancia entre dos puntos X y Y de un espacio n-dimensional en función de cuadrículas equivalentes aplicando una ponderación, su fórmula matemática es: d(X, Y ) = n X i=1 |xi−yi| |xi|+|yi|(3.3) Estos cálculos de distancias métricas quedan recogidos en el código fuente como distance.py lo que lo hace extensible para futuras métricas que se quieran incluir. 3.1.3. Generación de predicciones Cuando obtenemos los valores futuros observados de cada vecino, es necesario realizar una media ponderada de los mismos para poder obtener una predicción temporal sobre el dato de entrada e instante temporal que estamos evaluando, para ello usamos diferentes funciones de ponderación las cuales describiremos a continuación: Ponderación equitativa: es el método más simple y usado de todos, consiste en darle el mismo peso a todos los vecinos de una forma equitativa de la distancia que tengan hacia el conjunto a evaluar, la fórmula de la media aplicando esta ponderación es la siguiente: X=1 k k X i=1 xi=x1+x2+· · · +xk k(3.4) Donde kes el número de vecinos, Xel resultado de la media y xicorresponde a la predicción de cada vecino i. 34 CAPÍTULO 3. ALGORITMO K-NEAREST NEIGHBOURS Ponderación por proximidad: Este método consiste en darle un peso a la predicción de cada vecino según la distancia a la que se encuentre del conjunto evaluado. La fórmula con este método de ponderación es la siguiente: X=Pk i=1 xi∗1 yi Pk i=1 1 yi =x1∗1 y1+x2∗1 y2+· · · +xk∗1 yk 1 y1+1 y2+· · · +1 yk (3.5) Siendo kel número de vecinos, Xel resultado de la media , xila predicción de cada vecino ieyila distancia al conjunto evaluado i. Ponderación Lineal: Esta ponderación aplica un peso a cada elemento en función del orden distancia que tiene, es decir, aplica un mayor peso a los elementos más próximos y menor a los más lejanos sin tener en cuenta la distancia exacta a la que se encuentra, solo el orden de distancia, su fórmula aplicada es la siguiente: X=Pk i=1 xi∗(k−i+ 1) Pk i=1(k−i+ 1) =x1∗k+x2∗k−1 + · · · +xk k+k−1 + · · · + 1 (3.6) Donde kes el número de vecinos, Xel resultado de la media y xicorresponde al elemento de predicción de cada vecino i. Los cálculos de las medias ponderadas quedan recogidos en el código fuente como weights.py lo que lo hace extensible para futuras ponderaciones que se quieran incluir. 3.2. Entrenamiento y validación El algoritmo k-NN hace uso del aprendizaje basado en instancias que consiste en entrenar la generación de valores de las nuevas instancias basándose en las ya observadas con anterioridad dividiendo la serie temporal en varias partes. Normalmente un conjunto de datos se divide en dos partes, entrenamiento y validación, uno para cada propósito. En series temporales los datos están ordenados temporalmente y eso sucede también con estos conjuntos: el de entrenamiento consiste en un primer tramo de la serie y el de validación en el tramo siguiente. Sin embargo, el k-NN necesita pasado para buscar instancias y por eso se divide en tres siendo el primer tramo de inicialización. A continuación hablaremos de cada segmento detenidamente: Inicialización: Al ser un algoritmo basado en instancias, necesitamos un tramo desde el principio la serie temporal hasta cierto punto que se crea conveniente, estos datos serán un pasado donde podremos buscar los kvecinos para predecir instantes temporales, suele ser la fase que más instancias tiene sobre la serie temporal. Entrenamiento: Consiste en realizar predicciones desde el siguiente instante temporal del segmento de inicialización hasta el principio del segmento de validación de la serie temporal. Trata de encontrar los valores de ky de dque minimizan el error, para ello se prueba un rango de kdesde 1akmax y para ddesde 1admax, al tener que probar todas las combinaciones de kydse genera un gran coste en tiempo. Con esto conseguiremos buscar los 3.2. ENTRENAMIENTO Y VALIDACIÓN 35 parámetros para el algoritmo más óptimos a la hora de realizar las predicciones, los cuales serán los que cometan el menor error posible en las predicciones. Logramos calcular el error aplicando diferentes medidas del error, hablaremos de ellas en el siguiente apartado. Validación: Una vez obtenemos unos parámetros kydóptimos para la serie en la que estamos prediciendo, basados en lo aprendido en la fase de entrenamiento, se realiza las predicciones para todos instantes temporales desde el siguiente instante temporal a la fase de entrenamiento hasta el final de la serie. Tras conseguir predecir estos instantes temporales, se comparan con los resultados reales de cada instante, consiguiendo un índice de error para cada uno de estos instantes temporales. Figura 3.1: División de las diferentes segmentos de una serie temporal En la figura 3.1 podemos observar como se dividen las fases del algoritmo a lo largo de la serie temporal, para realizar estas divisiones de la serie, dispondremos del parámetro init para fijar la instancia temporal donde empezaremos a predecir con el algoritmo k-NN según en que fase nos encontremos. 3.2.1. Medida del error A la hora de validar las predicciones del algoritmo, es necesario elegir los valores adecuados para las diferentes combinaciones de kyden la fase de entrenamiento. Para poder hacer esto es necesario comprobar los errores entre la predicción obtenida y el valor real en cada instancia de la fase calculando a través de medidas del error un valor de error entre el valor real y la predicción. Sea para nel número de elementos, yiel valor real al instante temporal evaluado ieˆyila predicción a ese instante i, se definen las siguientes medidas de error: 36 CAPÍTULO 3. ALGORITMO K-NEAREST NEIGHBOURS Error medio (ME): Consiste en hacer en hacer una media usando la diferencia entre el valor real y la predicción, es poco aconsejable de utilizar ya que los valores negativos y positivos pueden cancelar unos a otros, su fórmula matemática es: ME =Pn i=1 yi−ˆyi n(3.7) Error medio absoluto (MAE): Está basado en el error medio, pero para evitar estas anulaciones en el sumatorio trata los datos en valor absoluto dando un resultado de error más fiable: MAE =Pn i=1 |yi−ˆyi| n(3.8) Raíz cuadrada del error cuadrático medio (RMSE): Es un método para el análisis del error en regresión, penaliza más que el MAE los errores grandes. Es muy utilizado en el análisis de temperatura y clima, su fórmula matemática es la siguiente: RMSE =sPn i=1(yi−ˆyi)2 n(3.9) Error medio porcentual (MPE): Se trata de una medida del error en porcentaje de precisión,no debe utilizarse si hay valores reales iguales a 0, como en el error medio, pueden darse el casos de que haya valores que se anulen entre sí: MPE =100 % n n X i=1 yi−ˆyi yi (3.10) Error medio porcentual absoluto (MAPE): Funciona igual que el error medio porcentual, solo que, da un porcentaje más certero al ser los datos en valor absoluto evitando así, la anulación de errores entre sí: MAPE =100 % n n X i=1 |yi−ˆyi yi |(3.11) Los cálculos de error métrico quedan recogidos en el código fuente como error.py para el sumatorio de valores y en KNN_Optim.py para su división, la implementación fue diseñada así, ya que spark tiene que realizar primero el sumatorio de diferencias para acto seguido aplicar la operación con el número de datos ncorrespondiente al error usado. Capítulo 4 Implementación y optimización del algoritmo La implementación de nuestro k-NN en Apache Spark queda guardado en el repositorio con la URL: https://github.com/saguila/KNN-clústerComputed/tree/master/ TFG protegido bajo una licencia Apache 2.0, en el mismo se encuentra la descripción y funciones comentadas y explicadas para su funcionamiento incluyendo versión de PySpark utilizada. A continuación hablaremos de su implementación del algoritmo k-NN tal y como fue descrito en la sección 3, las funciones principales requieren siempre como parámetros la serie temporal (rdd), los parámetros dyky el número de elementos de la serie (n) que es calculado usando la función de Spark count(), el resto de parámetros tienen definición por defecto luego no son obligatorios para que la función se ejecute [6]. Las funciones que implementan el algoritmo son las siguientes: KNN_Next: Función que predice el algoritmo k-NN para el siguiente instante de la serie temporal. KNN_Optim: Función que implementa el entrenamiento del algoritmo de para optimizar los parámetros kyd, usando en un rango determinado, diferentes combinaciones de kyd. KNN_Past: Predice valores correspondientes a diferentes fases temporal de la serie. Para poder hacer lectura de las series temporales volcaremos los datos sobre una estructura de datos llamada RDD (Resilient Distributed Dataset) que sirve para poder procesar la serie temporal de forma distribuida utilizando Spark, se implementan tres pasos del algoritmo usando esta estructura de manera distribuida, estos son : La preparación de los datos en función del parámetro d El cálculo de la matriz de distancias La generación de la predicción 4.1. Implementación paralela del algoritmo En este apartado, hablaremos acerca de las funciones distribuidas utilizadas por todas las funciones principales, explicando el contenido y funcionamiento de cada una de ellas. 37 38 CAPÍTULO 4. IMPLEMENTACIÓN Y OPTIMIZACIÓN DEL ALGORITMO 4.1.1. Preparación de los datos Para preparar los datos para comenzar el algoritmo, hay que crear un RDD de la serie temporal tal y como podemos observar en el cuadro 4.1 que hace referencia una serie temporal con el instante temporal ty su valor correspondiente y. Instante temporal t1t2t3t4· · · tn−2tn−1tn V alor y1y2y3y4· · · yn−2yn−1yn Cuadro 4.1: Mostramos una entrada de un RDD con una serie temporal (t) con su valor (y) A continuación, realizamos un proceso agrupación de los datos de la serie tal y como observamos en el cuadro 4.2, se agrupan los valores de dinstancias temporales, fijando para cada fila un identificador para cada agrupación siendo identificador 1 el primer instante de la serie temporal procesada y n-d la última instancia de la serie procesada. El valor de predicción de cada agrupación hará correspondencia al siguiente valor de la serie temporal (siendo este yd+id). Identificador V alores Predicción 1y1y2· · · yd−2yd−1ydyd+1 2y2y3· · · yd−1ydyd+1 yd+2 3y3y4· · · ydyd+1 yd+2 yd+3 4y4y5· · · yd+1 yd+2 yd+3 yd+4 . . .. . .. . ..... . .. . .. . .. . . n−d−3yn−d−3yn−d−2· · · yn−5yn−4yn−3yn−2 n−d−2yn−d−2yn−d−1· · · yn−4yn−3yn−2yn−1 n−d−1yn−d−1yn−d· · · yn−3yn−2yn−1yn n−d yn−dyn−d+1 · · · yn−2yn−1ynNone Cuadro 4.2: RDD después de prepararse la serie en función de d Diseñar esta función fue algo complicado ya que Spark trabaja de forma distribuida cada instancia de la serie, la solución inicial unificar varias copias de la serie temporal igual al valor del parámetro dy agruparlas entre sí para conseguir los datos tratados. Esta solución ocupaba una gran cantidad de espacio en memoria acabando por exigir un clúster más potente de lo necesario para preparar los datos de la serie, además exigía un gran coste en tiempo al tener que redistribuir dcada RDD. Buscando alternativas, se optó por realizar una solución proporcionando a cada valor varios identificadores correspondientes a las agrupaciones que pertenecía dicho valor, y acto seguido, agrupar los valores por cada identificador dado quedando como se describió en el cuadro 4.2, no exigía necesidad de hacer copias de la serie temporal lo cual permitió reducir el coste en espacio y tiempo considerablemente. Esta función de preparación de datos se encuentra definida en las función dRdd descrita en el fichero fuente dRdd.py, recibe por entrada, el RDD que contiene la serie temporal (rdd), el parámetro del algoritmo dy el número de instantes que contiene la serie temporal (n), devolviendo como salida la serie temporal agrupada. 4.1. IMPLEMENTACIÓN PARALELA DEL ALGORITMO 39 Figura 4.1: Diagrama de la función de agrupación de instantes temporales RDD. En los diagramas 4.1, 4.5, 4.6 y 4.9 dentro de la función principal se representan en color: Gris : Para una función determinada. Rojo : Para un RDD, colección de objetos, que está distribuida en varias máquinas. Morado : Para variables de Broadcast. Amarillo : Para variables de ámbito local. 40 CAPÍTULO 4. IMPLEMENTACIÓN Y OPTIMIZACIÓN DEL ALGORITMO Figura 4.2: Explicación en detalle del algoritmo de agrupación de instancias temporales dRdd. La figura 4.2 representa una ejecución del algoritmo con d=3. A, B, C, D y E son las instancias temporales que se van a agrupar. En el glosario podemos observar los tipos de datos que tenemos en cada uno de los pasos hasta obtener la agrupación temporal. En la figura podemos observar los siguientes pasos: 1. En este paso se realiza un zipWithIndex para poner un índice a las instancias temporales para tener una referencia que se usará posteriormente. 2. Se realiza un Map en el que se le pasa la función dRddMapTransformFunction, se encarga de coger los índices generados en el paso anterior y mediante una lógica crear arrays con las instancias temporales duplicadas necesarias para formar la agrupación. 3. Se separan las instancias duplicadas en vectores pasando a ser tuplas índice-instancia usando una función de flatMap. 4. Se agrupan las instancias temporales por indice usando GroupByKey, esta agrupación nos devuelve los índices y un iterador a las instancias agrupadas, esta estructura no nos interesa, necesitamos tener los datos de una forma directa. 4.1. IMPLEMENTACIÓN PARALELA DEL ALGORITMO 41 5. Se usa una función de MapValues para trabajar solo con los valores (los iteradores a las instancias) ya que la clave (el índice o referencia temporal) la queremos seguir manteniendo, sobre este MapValues se aplica la funciondRddConsolidateIterator, se encarga de transformar los iteradores a instancias en una lista ordenada con las instancias agrupadas, estas representan la serie temporal agrupada final. 4.1.2. Búsqueda de los kvecinos A la hora de calcular una matriz de distancias, se genera un proceso de coste elevado en tiempo, ya que exige comparar dos series temporales una que contiene los instantes temporales a predecir y otra con la serie temporal completa para poder obtener la distancia de cada instante temporal a predecir con su pasado obteniendo los kvecinos más cercanos de cada instancia a predecir. En un principio se implementó almacenando las dos series en dos RDDs, al calcularse las distancia entre dos series distribuidas en varios ejecutores exigía mucho coste en tiempo, ya que ambos ambas series temporales debían redistribuirse entre los ejecutores varias veces. Este problema llevo a investigar alternativas mejores dando como resultado al uso de las variables de broadcast [4] explicadas en el sección 2.3.3. Figura 4.3: Trabajo de spark distribuyendo las variables de broadcast para generar la matriz de distancia distribuida Aplicando este método conseguimos mayor velocidad en el algoritmo gracias a que, al tener todos los instantes a predecir en todos los ejecutores, evitamos tener que hacer varias redistribuciones de la serie. 48 CAPÍTULO 4. IMPLEMENTACIÓN Y OPTIMIZACIÓN DEL ALGORITMO 2.5 Abrimos el bucle anidado hasta k(línea 8) y volvemos a empezar el proceso hasta realizar estos pasos dveces, una vez realizados devolveremos la variable optimData que almacena el menor error encontrado en las combinaciones, la ky la dóptimas en una tupla de tres datos. . Figura 4.7: Trabajo del bucle externo Ahora vamos a explicar el funcionamiento del bucle anidado hasta k, este bucle se utiliza para realizar las combinaciones con la dactual con todos los rangos de kposibles buscando el menor error posible, así evitamos que rehacer más veces de las necesarias la matriz de distancia, que es el proceso más costoso del algoritmo. Podemos observar en la figura 4.8 el esquema de funcionamiento del bucle interno, a continuación explicamos paso a paso el funcionamiento del bucle: 1. Abrimos un bucle anidado hasta k(línea 8). 4.2. FUNCIONES PRINCIPALES 49 2. Realizamos la predicción de los jvecinos más cercanos y calculamos el error usando un tipo de medida de error. Se almacena en una tupla de tres valores para comprobar si es el menor error que llamaremos check con los valores: error a comprobar, valor de kactual (j) y valor dactual (i). 3. Comprobamos si el error en la tupla check tiene menor error que la tupla óptima elegida actualmente, de ser así, almacenamos check en optimData. 4. Repetimos este proceso hasta probar todas los valores fijados en el rango de k. Figura 4.8: Trabajo del bucle interno Como podremos observar en la explicación esta función la consideraremos crítica ya que, según aumentan las combinaciones a analizar de kyd, siendo dla que aumentará drásticamente el coste en tiempo del algoritmo. 50 CAPÍTULO 4. IMPLEMENTACIÓN Y OPTIMIZACIÓN DEL ALGORITMO Figura 4.9: Diagrama de flujo KNN_Optim Capítulo 5 Experimentos y comparación con la implementación mono-máquina A la hora de realizar predicciones temporales podemos encontrarnos diferentes tamaños de series temporales lo cual conlleva a diferentes conclusiones de que método usar, distribuido o mono-máquina. En este capítulo hablaremos y analizaremos la eficiencia de cada método según el tamaño de datos de la serie temporal. 5.1. Series temporales consideradas En este apartado presentaremos las series temporales que hemos usado para probar el funcionamiento del algoritmo. 5.1.1. Bitcoin Historical Data Contiene un histórico por minuto del precio de la cripto moneda Bitcoin en Dólares desde enero de 2012 hasta julio de 2018 con un total de 3.255.876 registros. https://www.kaggle.com/mczielinski/bitcoin-historical-data/home. La predicción de los valores futuros de Bitcoin con la solución implementada presenta un problema grave, debido a que es una serie con tendencia estrictamente creciente a la que no se la puede aplicar directamente el algoritmo, puesto que no se observan patrones ni buenos vecinos para usar en las predicciones. Con el fin de hacerla más útil se ha hecho una transformación, filtrando los registros pertenecientes al periodo inicial en el que la criptomoneda era estable quedando un total de 2.188.132 empezando desde julio del 2012 hasta julio de 2018. Además, se ha transformado la serie y se ha convertido en la serie de los rendimientos logarítmicos del valor del bitcoin. Con esto hemos conseguido transformar un problema de predicción del próximo valor del Bitcoin a un problema de predicción de periodos de rendimientos haciendo útil la aplicación del algoritmo. 5.1.2. Sunspot Daily Contiene un histórico con las observaciones diarias del número de manchas solares desde el 1 de enero de 1818 hasta el 30 Junio de 2018 con un total de 73.230 registros. 51 52CAPÍTULO 5. EXPERIMENTOS Y COMPARACIÓN CON LA IMPLEMENTACIÓN MONO-MÁQUINA http://www.sidc.be/silso/datafiles. Esta serie al contrario que la anterior, directamente se pueden observar patrones, no ha necesitado ninguna transformación de los datos para hacer predicciones útiles con el algoritmo, lo único que se ha realizado fue una adaptación a la entrada esperada por el programa, eliminando campos que no pertenecen a la serie temporal. 5.1.3. Sunspot Monthly Igual que el anterior contiene un histórico con las observaciones diarias del número de manchas solares pero con una periodicidad mensual desde el enero de 1 de enero de 1749 hasta Junio de 2018 con un total de 3.233 registros. http://www.sidc.be/silso/datafiles. Como paso con el conjunto de datos con periodicidad diaria, se ha realizado fue una adaptación a la entrada esperada por el programa, eliminando campos que no pertenecen a la serie temporal. 5.2. Alcance de la solución mono-máquina en R vs solución Spark La solución implementada por los compañeros en Rpara una máquina [1] esta limitada por la cantidad de memoria que ocupa la matriz de distancias que se genera para resolver el algoritmo. Esta matriz con una cantidad de registros determinada puede superar la capacidad de una máquina, debido a que cualquier computador de 64 bits como mucho puede referenciar 8 Tb en memoria ram. Con una máquina con características similares a la que se ha usado en el desarrollo (véase sección 2.4.1), se pueden procesar unos 44000 instantes temporales como máximo que ocupan 7,3 Gb de memoria ram. El clúster donde se realizarán las pruebas para la solución de Spark con 4 máquinas serán tal y como se describe en la sección 2.4.2. 5.3. Escenarios Para poder analizar la eficiencia entre entre el algoritmo distribuido contra mono-máquina usaremos diferentes históricos según su tamaño y analizaremos las funciones una a una para realizar una comparación de coste en tiempo entre ambos, todos los históricos nombrados y funciones a probar en este punto se encuentran almacenados en el repositorio de GitHub de la sección 5.1. 5.3.1. Serie pequeña En estas pruebas utilizaremos principalmente el histórico Sunspot.monthly y un parámetro init de 400 datos, luego realizaremos la predicción sobre 2.833 . A continuación mostraremos cada una de las funciones y su coste en tiempo del algoritmo. 5.3. ESCENARIOS 53 KNN_Next Spark d k100101102103 1002,11 1,95 1,76 1,93 1011,6 1,71 1,73 1,53 1022,15 2,17 1,9 1,83 1034,69 4,99 4,38 5,5 Cuadro 5.1: Coste en tiempo de la aplicación con Spark en segundos para KNN_Next R d k100101102103 1000,03 0,03 0,03 0,03 1010,05 0,05 0,05 0,05 1020,31 0,32 0,32 0,33 1037,1 8,17 8,31 8,57 Cuadro 5.2: Coste en tiempo de la aplicación con Ren segundos para KNN_Next Observando los resultados se puede ver que claramente que la programación mono-máquina es más rápida cuando hay parámetros menores, según va creciendo (d) va aumentando el coste progresivamente llegando a un punto en el que es mejor hacer programación distribuida. El valor de (k) no ha tenido apenas influencia luego podemos deducir que el número de vecinos prácticamente, no afecta al algoritmo. KNN_Past Spark d k100101102103 1008,89 9,37 13,38 17,34 10127,45 27,87 27,89 27,91 102133,7 141,29 143,91 147,24 103950,02 851,11 919,04 829,14 Cuadro 5.3: Coste en tiempo de la aplicación con Spark en segundos para KNN_Past 54CAPÍTULO 5. EXPERIMENTOS Y COMPARACIÓN CON LA IMPLEMENTACIÓN MONO-MÁQUINA R d k100101102103 1000,25 0,31 0,25 0,32 1010,27 0,29 0,32 0,30 1020,57 0,54 0,61 0,55 1037,14 7,16 7,16 7,25 Cuadro 5.4: Coste en tiempo de la aplicación con Ren segundos para KNN_Past Tras realizar las pruebas en ambas máquinas, observamos que spark empeora dramáticamente según van aumentando el retardo (d), es bastante peor en tiempo en comparación con R. Pero hay que tener en cuenta que las soluciones distribuidas son mucho más costosas y contamos con un clúster de prueba no productivo, por lo tanto este retardo se puede evitar con un clúster más potente ya que podríamos emplear más ejecutores para el algoritmo. KNN_Optim Spark d k100101102 1009,24 8,7 27,24 101167,19 343,8 337,28 1028135,74 8635,23 9038,16 Cuadro 5.5: Coste en tiempo de la aplicación con Spark en segundos para KNN_Optim R d k100101102 1000,30 0,45 3,79 1012,39 5,69 40,49 10259,83 136,10 448,6 Cuadro 5.6: Coste en tiempo de la aplicación con Ren segundos para KNN_Optim Al igual que pasaba con el KNN_Past los tiempos son peores, tenemos el mismo problema el clúster de pruebas es demasiado sencillo, la complejidad y el tamaño de los datos merma los tiempos. 5.3.2. Serie media Para este apartado se utilizo una serie temporal de pruebas sunspot.daily descrito en el punto 5.1.2. 5.3. ESCENARIOS 55 Probamos a ejecutarlo varias veces con diferentes configuraciones, usando Amazon Elastic MapReduce con un clúster con 5 máquinas m5.xlarge workers con 4 cores y 16 gigas por máquina que hacen un total de 80 gigas y 20 cores. tras ejecutar todas las funciones con un valor de 3 para los parámetros d,ky un init de 48820, el algoritmo tardo 4.369 segundos (1 hora y 12 minutos aproximadamente) en ejecutarse usando las tres funciones distribuyendo los tiempos así: KNN_Optim 52 min KNN_Past 21 min KNN_Next 10 seg Cuadro 5.7: Coste en tiempo de la aplicación con Spark para cada función Se ha realizado una sola ejecución de todas las funciones debido a los costes del alquiler del clustér los cuales encarecen su precio por horas de uso. Tras finalizar las pruebas con Spark, probamos a ejecutar en mono-máquina con la misma serie temporal dando como resultado un fallo de ejecución debido a que la matriz de distancia excedía el límite de memoria máxima (7.3 GB), se amplio el límite y se volvió a probar dando lugar a un "Screen Freeze"del ordenador utilizado, luego no se consiguió probar. Al no poder realizar la prueba sobre mono-máquina, no se pudo deducir una mejora o empeoramiento en coste en tiempo. Como no se pudo demostrar la velocidad se realizó un recorte al histórico para que fuera aceptable usando una mono máquina demostrando que, al haber usado un clúster con mayor capacidad de memoria que la mono-máquina se consiguió una predicción más certera como podemos ver en los resultados: 56CAPÍTULO 5. EXPERIMENTOS Y COMPARACIÓN CON LA IMPLEMENTACIÓN MONO-MÁQUINA Figura 5.1: Diferencia entre predicción y realidad para un fragmento de 2000 instantes temporales. Figura 5.2: Validación del error absoluto entre Spark yRcon 9 combinaciones. 5.3. ESCENARIOS 57 Para estas pruebas hemos probado el histórico completo en Spark y los últimos 12000 instantes temporales del histórico para Rprediciendo 2000 instantes temporales tanto para entrenamiento como para validación, utilizando el error métrico MAE y un valor 3 para los parámetros kyd dando un total de 9 combinaciones para la búsqueda de los parámetros óptimos. MAE Train Test R10,575 8,139 Spark 7,387 6,558 Cuadro 5.8: Diferencia de error con MAE entre RySpark En la figura 5.1 podemos observar la diferencia entre los datos predichos por cada método, se ve una similitud no muy clara entre la realidad y las predicciones lo cual ha requerido comprobar con una validación del error, en las predicciones que podemos observar en la figura 5.2 y el cuadro 5.8 que demuestran que, al tener un histórico con mayor cantidad de instantes temporales, Spark ha realizado unas predicciones más certeras que las realizadas en R. Figura 5.3: Validación del error absoluto entre Spark yRcon 100 combinaciones Se repitieron las mismas pruebas alterando los parámetros kyda 10 dando un total de 100 combinaciones para buscar los parámetros óptimos. Tal y como observamos en la figura 5.3 con los que pudieron comprobar la importancia de aumentar las combinaciones en el entrenamiento, observamos que Spark sigue siendo más precisa al disponer de un histórico mayor, pero ambas 64 Bibliografía Agradecimiento A los tutores del proyecto Javier Arroyo Gallardo y a Albert Meco Alías, por darnos la oportunidad de participar en este proyecto, por la guía y el apoyo recibido. ACoreNetworks S.L por facilitarnos las maquinas usadas para el desarrollo y las pruebas del TFG. ADaniel Francisco Bastarrica Lacalle y a Javier Berdecio Trigueros, por su trabajo que nos ha servido de apoyo y comparativa a nuestro proyecto. 65