scieee AI-readable full text Open interactive document viewer

Optimización de la Entrada Salida mediante librerías y lenguajes paralelos

Larrosa-Jiménez, Rafael

Abstract

Uno de los grandes retos de la HPC (High Performance Computing) consiste en optimizar el subsistema de Entrada/Salida, (E/S), o I/O (Input/Output). Ken Batcher resume este hecho en la siguiente frase: "Un supercomputador es un dispositivo que convierte los problemas limitados por la potencia de cálculo en problemas limitados por la E/S" ("A Supercomputer is a device for turning compute-bound problems into I/O-bound problems") . En otras palabras, el cuello de botella ya no reside tanto en el procesamiento de los datos como en la disponibilidad de los mismos. Además, este problema se exacerbará con la llegada del Exascale y la popularización de las aplicaciones Big Data. En este contexto, esta tesis contribuye a mejorar el rendimiento y la facilidad de uso del subsistema de E/S de los sistemas de supercomputación. Principalmente se proponen dos contribuciones al respecto: i) una interfaz de E/S desarrollada para el lenguaje Chapel que mejora la productividad del programador a la hora de codificar las operaciones de E/S; y ii) una implementación optimizada del almacenamiento de datos de secuencias genéticas. Con más detalle, la primera contribución estudia y analiza distintas optimizaciones de la E/S en Chapel, al tiempo que provee a los usuarios de una interfaz simple para el acceso paralelo y distribuido a los datos contenidos en ficheros. Por tanto, contribuimos tanto a aumentar la productividad de los desarrolladores, como a que la implementación sea lo más óptima posible. La segunda contribución también se enmarca dentro de los problemas de E/S, pero en este caso se centra en mejorar el almacenamiento de los datos de secuencias genéticas, incluyendo su compresión, y en permitir un uso eficiente de esos datos por parte de las aplicaciones existentes, permitiendo una recuperación eficiente tanto de forma secuencial como aleatoria. Adicionalmente, proponemos una implementación paralela basada en Chapel.

Full text

Departamento de Arquitectura de Computadores TESIS DOCTORAL Optimizaci´ on de la Entrada Salida mediante librer´ ıas y lenguajes paralelos Rafael Larrosa Jim´ enez Noviembre de 2015 Dirigida por: Mar´ ıa ´ Angeles Gonz´ alez Navarro, Rafael Asenjo Plaza AUTOR: Rafael Larrosa Jiménez http://orcid.org/0000-0002-4166-3041 EDITA: Publicaciones y Divulgación Científica. Universidad de Málaga Esta obra está bajo una licencia de Creative Commons Reconocimiento-NoComercial-SinObraDerivada 4.0 Internacional: Cualquier parte de esta obra se puede reproducir sin autorización pero con el reconocimiento y atribución de los autores. No se puede hacer uso comercial de la obra y no se puede alterar, transformar o hacer obras derivadas. http://creativecommons.org/licenses/by-nc-nd/4.0/legalcode Esta Tesis Doctoral está depositada en el Repositorio Institucional de la Universidad de Málaga (RIUMA): riuma.uma.es Dra. D˜ na. Ma´ Angeles Gonz´ alez Navarro. Profesora Titular del Departamento de Arquitectura de Computadores de la Universidad de M´ alaga. Dr. D. Rafael Asenjo Plaza. Profesor Titular del Departamento de Arquitectura de Computadores de la Universidad de M´ alaga. CERTIFICAN: Que la memoria titulada “Optimizaci´ on de la Entrada Salida mediante librer´ ıas y lenguajes paralelos”, ha sido realizada por D. Rafael Larrosa Jim´ enez bajo nuestra direcci´ on en el Departamento de Arquitectura de Computadores de la Universidad de M´ alaga y constituye la Tesis que presenta para optar al grado de Doctor en Ingenier´ ıa Inform´ atica. M´ alaga, Noviembre de 2015 Dra. D˜ na. Ma´ Angeles Gonz´ alez Navarro. Codirectora de la tesis. Dr. D. Rafael Asenjo Plaza. Codirector de la tesis. . A mis hijos Rafa y Javi, y a Macarena, para quienes esta tesis ha durado desde siempre. . Agradecimientos En primer lugar, me gustar´ ıa dar las gracias a mis directores: Mar´ ıa ´ Angeles Gonz´ alez Navarro y Rafael Asenjo Plaza, sin sus ´ animos y dedicaci´ on esta tesis no habr´ ıa sido posible. A continuaci´ on al Departamento de Arquitectura de Computadores y las personas que lo componen, Carmen, quien consigue que para todos sea sencillo lo complicado, ayudando siempre, todos los profesores y personal del departamento, Paco, Juanjo y MaCarmen, siempre haciendo que los problemas tecnol´ ogicos desaparezcan. Aunque much´ ısimos merecen estar aqu´ ı, voy a resumirlo en algunos, Guillermo, siempre dispuesto a concretar hasta los m´ as m´ ınimos detalles de todo, desde la fotograf´ ıa a la m´ usica, Nicol´ as, compa˜ nero de fatigas, Mario, Gerardo, Pablo, Eladio, Sergio, Mar´ ıa Antonia, Felipe, Siham, Juli´ an, Francisco, Oswaldo, Ujaldon, Andr´ es, Manuel S´ anchez, Juan, Emilio, Oscar, Sonia, Javier, Chema, Inmaculada, Eligius y especialmente a Julio Villalba, por toda su ayuda durante el Proyecto Fin de Carrera. Agradecer a Dar´ ıo todas las decisiones discutidas y consensuadas para conseguir que todo funcione de forma ´ optima, espero que sigamos optimizando sistemas por mucho tiempo. A Gonzalo todas las discusiones sem´ anticas y sint´ acticas, como sobre cu´ ando usar n´ umeros romanos y cu´ ando no, y a Roc´ ıo todo el conocimiento y buen hacer bioinform´ atico, que incluso aparece reflejado en esta tesis. Tambi´ en a Pepi, Diego, Roc´ ıo Romero, Ana, Shanti, Rosario, Hicham, Isabel, Pedro y Marina por todas las animadas conversaciones. A la mesa redonda de relatividad y mec´ anica cu´ antica, Rafael Miranda, Francisco Villatoro, Jos´ e Galindo y Juan Ignacio Ramos, por todas las discusiones y ense˜ nanzas durante tanto tiempo. A mis amigos de siempre, Caro, Mancebo, Luengo, Caballero, Alberto, Jos´ e Luis, y tantos otros que mencionar, por los buenos ratos. A todos mis profesores, desde Infantil hasta la Universidad, por haber sabido siempre inculcar las ganas de aprender. Quiero agradecer tambi´ en al equipo de Chapel su buen trabajo, especialmente a Brad II Agradecimientos Chamberlain por dirigirlo y estar siempre al tanto, y a Vassily Litvinov por su ayuda y discusiones en todo aquello relacionado con las distribuciones. Tambi´ en a Alberto Sanz, por su ayuda y por hacer m´ as sencilla la programaci´ on. Adem´ as quisiera agradecer el acceso a los recursos de computaci´ on usados en esta tesis, a la empresa Cray por el acceso al ordenador Crow, al SCBI (SuperComputaci´ on y BioInform´ atica) de la Universidad de M´ alaga por el acceso a Picasso, y al Oak Ridge Leadership Facility del Oak Ridge National Laboratory, que est´ a mantenido por la Oficina de la Ciencia del Departamento de Energ´ ıa de los EEUU bajo el contrato n´ umero DE-AC05-00OR22725, el acceso a Jaguar, Titan y EOS. Agradecer tambi´ en a David Knaak su ayuda para conocer mejor la arquitectura de los Cray, y las optimizaciones en sus implementaciones de la E/S, contestando las preguntas que se le hicieron, que de otra manera no podr´ ıamos haber obtenido. A Maribel y Pepe, por toda su ayuda y cari˜ no. A Silvia y Jos´ e Luis, por no desesperar ante mis r´ eplicas recursivas, y a mis sobrinos, Lu´ ıs y Julio, que me han mostrado otra perspectiva de eso que llamamos ”inform´ atica”. Gracias a toda mi familia, y especialmente a mi padre y mi madre, Rafael y Antonia, que me han apoyado y soportado siempre, y que a´ un ahora casi no se creen que por fin esta tesis llegue a su fin. A mi hermano Antonio, por toda su ayuda durante tanto tiempo y todos los buenos ratos que hemos pasado y que continuaremos pasando, y a Eva, por aportar felicidad. Y muy especialmente a Rafa y a Javi, que no se han enterado muy bien de a qu´ e dedicaba tanto tiempo, vuestra presencia me llena de felicidad. Gracias por sorprenderme cada d´ ıa, os quiero. Y gracias de todo coraz´ on a Macarena, desde que te conoc´ ı aquel 8 de enero has llenado de felicidad, cari˜ no y amor mi vida. Te quiero. Resumen Uno de los grandes retos de la HPC (High Performance Computing) consiste en optimizar el subsistema de Entrada/Salida, (E/S), o I/O (Input/Output). Ken Batcher resume este hecho en la siguiente frase: “Un supercomputador es un dispositivo que convierte los problemas limitados por la potencia de c´ alculo en problemas limitados por la E/S”1. En otras palabras, el cuello de botella ya no reside tanto en el procesamiento de los datos como en la disponibilidad de los mismos. Adem´ as, este problema se exacerbar´ a con la llegada del Exascale y la popularizaci´ on de las aplicaciones Big Data. En este contexto, esta tesis contribuye a mejorar el rendimiento y la facilidad de uso del subsistema de E/S de los sistemas de supercomputaci´ on. Principalmente se proponen dos contribuciones al respecto: i) una interfaz de E/S desarrollada para el lenguaje Chapel que mejora la productividad del programador a la hora de codificar las operaciones de E/S; y ii) una implementaci´ on optimizada del almacenamiento de datos de secuencias gen´ eticas. Con m´ as detalle, la primera contribuci´ on estudia y analiza distintas optimizaciones de la E/S en Chapel, al tiempo que provee a los usuarios de una interfaz simple para el acceso paralelo y distribuido a los datos contenidos en ficheros. Por tanto, contribuimos tanto a aumentar la productividad de los desarrolladores, como a que la implementaci´ on sea lo m´ as ´ optima posible. La segunda contribuci´ on tambi´ en se enmarca dentro de los problemas de E/S, pero en este caso se centra en mejorar el almacenamiento de los datos de secuencias gen´ eticas, incluyendo su compresi´ on, y en permitir un uso eficiente de esos datos por parte de las aplicaciones existentes, permitiendo una recuperaci´ on eficiente tanto de forma secuencial como aleatoria. Adicionalmente, proponemos una implementaci´ on paralela basada en Chapel. 1“A supercomputer is a device for turning compute-bound problems into I/O-bound problems.” ´ Indice de figuras 1.1. Configuraci´ on b´ asica de GPFS usando una SAN. . . . . . . . . . . . . 10 1.2. Configuraci´ on GPFS usando servidores de E/S. . . . . . . . . . . . . . 11 1.3. Un ejemplo de sistema Lustre, con dos OSSs redundantes. . . . . . . . 15 1.4. Escritura secuencial en Lustre al escribir en un stripe........... 18 1.5. C´ odigo OpenMP para calcular pi en paralelo [13] . . . . . . . . . . . . 24 1.6. C´ odigo MPI para calcular pi en paralelo [13]. . . . . . . . . . . . . . . 26 1.7. C´ alculo paralelo de pi en el lenguaje de alta productividad Chapel. . . . 27 1.8. C´ odigo MPI IO para escribir en paralelo un array distribuido. . . . . . . 28 1.9. Arquitectura E/S de Crow. . . . . . . . . . . . . . . . . . . . . . . . . 30 1.10. Enlaces del chip SeaStar2+, incluyendo su conexi´ on a un nodo. . . . . . 31 1.11. Arquitectura hardware del sistema de almacenamiento Spider I [110]. . 32 1.12. Arquitectura de Spider II [112]. . . . . . . . . . . . . . . . . . . . . . . 34 1.13. Arquitectura SION original [111]. . . . . . . . . . . . . . . . . . . . . 36 1.14. Arquitectura Spider II y SION [112]. . . . . . . . . . . . . . . . . . . . 37 1.15. Esquema de un blade perteneciente a un Cray XC30. . . . . . . . . . . 38 1.16. Fotograf´ ıa de un blade de un Cray XC30. . . . . . . . . . . . . . . . . 38 1.17. Ejemplo de canales de comunicaci´ on entre dos blades de un XC30. . . . 39 1.18. Esquema global de un Cray XC30. . . . . . . . . . . . . . . . . . . . . 40 1.19. Ordenador Picasso de la UMA. . . . . . . . . . . . . . . . . . . . . . . 41 2.1. Declaraci´ on en Chapel de una matriz distribuida en bloques. . . . . . . 44 XI XII ´ INDICE DE FIGURAS 2.2. Ejemplo de distribuci´ on de bloques sobre cuatro locales. ........ 46 2.3. Distribuci´ on en bloques unidimensional de una matriz en tres locales. . 47 2.4. Distribuci´ on en bloques bidimensional de una matriz en cuatro locales. . 47 2.5. Ejemplo de asignaci´ on parcial entre arrays en Chapel. . . . . . . . . . . 48 2.6. Funci´ on de GASNet para realizar una transferencia bulk-strided. . . . . 52 2.7. Ejemplo de copia entre dos arrays tridimensionales. . . . . . . . . . . . 54 2.8. Otro ejemplo de copia entre dos arrays tridimensionales. . . . . . . . . 55 2.9. C´ odigo para realizar transferencias agregadas. . . . . . . . . . . . . . . 56 2.10. Asinganci´ on entre arrays A y B cuando s´ on de tipo DR oBD....... 57 2.11. C´ odigo Chapel para ilustrar las funciones de correspondencia. . . . . . 59 2.12. Asignaci´ on entre arrays con 4 locales. . . . . . . . . . . . . . . . . . . 59 2.13. Tiempo en segundos del algoritmo PARACR en Titan. . . . . . . . . . 61 2.14. Dise˜ no original del interfaz para realizar E/S en Chapel. . . . . . . . . . 65 2.15. Nuevo interfaz para ficheros distribuidos en Chapel. . . . . . . . . . . . 65 2.16. Otro posible interfaz para ficheros distribuidos en Chapel. . . . . . . . . 66 2.17. Ejemplo de compartici´ on de dos stripes en Lustre al escribir una matriz. 69 2.18. Arquitectura global de E/S con agregadores. . . . . . . . . . . . . . . . 71 2.19. Arquitectura global de E/S incluyendo m´ as paralelismo. . . . . . . . . . 72 2.20. Pseudoc´ odigo Chapel de E/S paralela. . . . . . . . . . . . . . . . . . . 74 2.21. Paralelismo 1: Escritura en paralelo en los OSTs . . . . . . . . . . . . . 75 2.22. Paralelismo 2: Paralelismo en la agregaci´ on................ 77 2.23. Paralelismo 3: Paralelismo entre E/S y agregaci´ on. ........... 78 2.24. Paralelismo 4. Paralelismo en la escritura en un OST . . . . . . . . . . 79 2.25. Ancho de banda para las mejores combinaciones (0,x,x,x). . . . . . . . 88 2.26. Ancho de banda para las mejores combinaciones (1,x,x,x). . . . . . . . 89 2.27. Ancho de banda para las combinaciones (1,x,x,1) con 8 OSTs. . . . . . 91 2.28. Comparaci´ on entre los algoritmos (1,1,x,0) y MPI IO. . . . . . . . . . . 92 3.1. Coste de secuenciar un mill´ on de bases (una megabase) [139]. . . . . . 101 ´ INDICE DE FIGURAS XIII 3.2. FormatoFQbin............................... 105 3.3. Comparaci´ on entre distintos algoritmos de compresi´ on.......... 111 3.4. Aceleraci´ on obtenida con FQbin al leer. . . . . . . . . . . . . . . . . . 112 3.5. Declaraci´ on en Chapel de las variables usadas en la librer´ ıa FQbin. . . . 117 3.6. Tiempos y aceleraci´ on de la implementaci´ on Chapel de FQbin para un locale.................................... 118 3.7. Variables usadas en la librer´ ıa FQbin distribuida. . . . . . . . . . . . . 119 3.8. Cambios a efectuar en los bucles para que sean distribuidos. . . . . . . 119 3.9. Tiempos de la implementaci´ on distribuida en Chapel de FQbin. . . . . . 120 A.1. Superordenador Jaguar. . . . . . . . . . . . . . . . . . . . . . . . . . . 132 A.2. Superordenador Titan. . . . . . . . . . . . . . . . . . . . . . . . . . . 132 A.3. Arquitectura del chip Gemini. . . . . . . . . . . . . . . . . . . . . . . 133 A.4. Comparaci´ on entre las redes Seastar2 y Gemini. . . . . . . . . . . . . . 134 A.5. Distribuci´ on de los nodos LNET en Titan [43] . . . . . . . . . . . . . . 134 ´ Indice de Tablas 1.1. Proyecci´ on sobre los principales atributos de un sistema Exascale. . . . 2 2.1. Resumen de las distintas fuentes de paralelismo que explotaremos de formacombinada.............................. 71 3.1. Tiempo de acceso para la primera y ´ ultima secuencia del fichero . . . . 112 3.2. Comparativa de tiempos de lectura de FQbin vs FASTA. . . . . . . . . 113 XV Prefacio There ain’t no such thing as a free lunch Hist´ oricamente hemos estado acostumbrados a que el software vaya cada vez m´ as r´ apido, aprovechando la mayor velocidad de las CPUs que a˜ no tras a˜ no han incorporado mejoras tecnol´ ogicas y arquitecturales. Los programas simplemente se ejecutaban en menos tiempo sin que hiciera falta ning´ un esfuerzo para conseguir ese mejor rendimiento. Pero la “comida gratis” se acab´ o. En el a˜ no 2005 se public´ o el art´ ıculo “The Free Lunch Is Over: A fundamental turn toward concurrency in software”2[116], donde se argumenta que esa etapa termin´ o, dando paso a otra, en la que en vez de aumentar la velocidad de las CPUs los fabricantes se concentran en la concurrencia, es decir, en incrementar el n´ umero de cores disponibles. Esto a su vez implica re-implementar las aplicaciones as´ ı como desarrollar nuevo software del sistema (compiladores y runtime) para explotar los nuevos niveles de concurrencia en las modernas arquitecturas. M´ as a´ un. Con el advenimiento de la era de los sistemas Exascale al final de esta d´ ecada, aparecen nuevos desaf´ ıos que han de ser resueltos para poder escalar desde los actuales sistemas Petascale. Entre otras cosas, el aumento del rendimiento se espera que se consiga aumentando dr´ asticamente el n´ umero de cores por nodo. Esto dar´ a lugar a sistemas que expongan del orden de 1 ×109threads concurrentes. Para poder explotar este elevad´ ısimo nivel de concurrencia se hace necesario el desarrollo de nuevos modelos de programaci´ on. Y no s´ olo eso, adem´ as habr´ a que resolver el hecho de que la memoria y el ancho de banda por core se ver´ an dr´ asticamente reducidos. Puesto que el coste del movimiento de datos, tanto desde el punto de vista del consumo de energ´ ıa como del rendimiento, no mejorar´ a al ritmo de los flops, se exigir´ a que los algoritmos traten de minimizar el movimiento de datos. Ello requerir´ a del soporte de nuevas funcionalidades 2“La comida gratis se ha acabado: un giro fundamental hacia la concurrencia en el software” XVII XVIII Prefacio de E/S a todos los niveles: desde el SoC a la memoria, desde la memoria al nodo que sirve la E/S y desde el nodo de E/S al disco. Por todas estas razones, los nuevos modelos de programaci´ on que se dise˜ nen no s´ olo han de ser capaces de gestionar eficientemente billones de elementos de computaci´ on concurrentes, sino que adem´ as deben de permitir un uso eficiente de la E/S. El problema de la E/S viene de intentar conjugar tres factores, la velocidad, la capacidad y la fiabilidad. Los dispositivos de almacenamiento (cintas, discos duros rotacionales, discos SSD, y pr´ oximamente NVRAMs) tienen unos par´ ametros dados por la tecnolog´ ıa disponible en cada momento. Para lograr mayores capacidades y velocidades se suelen agrupar varios dispositivos de almacenamiento, formando dispositivos de bloques, lo que proporciona un acceso transparente y unificado a multitud de dispositivos individuales, sumando las capacidades, y hasta cierto punto las velocidades, de todos ellos. Pero este agrupamiento tiene l´ ımites tanto f´ ısicos (los racks se llenan), como el´ ectricos (el consumo se vuelve prohibitivo), como a nivel del n´ umero de dispositivos que se pueden conectar en un bus sin saturarlo. Para resolver los l´ ımites f´ ısicos del hardware en velocidad, capacidad y fiabilidad se recurre a soluciones software que aumentan la complejidad de los sistemas. Esto resulta en un incremento significativo en la dificultad de programar productivamente la E/S, lo que a su vez provoca que muchos usuarios terminen optando por el uso de las operaciones POSIX est´ andar. Desde hace ya d´ ecadas, los sistemas de ficheros se comparten desde m´ ultiples nodos de computaci´ on, llegando en algunos casos a situaciones en que el mismo sistema de ficheros sea accedido simult´ aneamente desde decenas de miles de nodos. Dichas arquitecturas de almacenamiento requieren de redes de comunicaciones que han ido creciendo en complejidad al tiempo que los sistemas han ido escalando. Estos sistemas de ficheros paralelos constituyen una capa intermedia que ocultan las operaciones de acceso concurrente a los dispositivos de E/S, pero en muchos casos requieren que el programador, si busca rendimiento, sea responsable de afinar algunos par´ ametros no triviales de la arquitectura del sistema de ficheros o incluso que se haga cargo de aspectos como asegurar la consistencia de accesos concurrentes a ficheros compartidos por varios nodos. En este trabajo no abogamos por optimizar manualmente los programas para explotar todo el potencial de los modernos sistemas paralelos de E/S, sino por ofrecer lenguajes, herramientas y librer´ ıas que simplifiquen la obtenci´ on del m´ aximo rendimiento sin aumentar la complejidad de la programaci´ on. Con todo esto, esta tesis est´ a guiada por dos objetivos principales. Por un lado presentar una interfaz de operaciones de E/S de alta productividad y alto rendimiento, con el objetivo de facilitar la programaci´ on y asegurar el rendimiento a la hora de implementar operaciones de E/S paralelas. Por otro lado, y en el contexto del almacenamiento de secuencias gen´ eticas, proponemos una librer´ ıa de E/S, actualmente en explotaci´ on, que implementa un nuevo formato de fichero de almacenamiento de secuencias gen´ omicas. Prefacio XIX En los dos casos nos apoyaremos en el lenguaje paralelo Chapel, inicialmente creado en Cray Inc., y que est´ a ganando aceptaci´ on a la hora de codificar en poco tiempo aplicaciones paralelas. En esta tesis en particular, contribuimos a la incorporaci´ on de operaciones paralelas de E/S en Chapel. Estudiamos y evaluamos diversas estrategias para implementar la E/S paralela, analizando en cada caso las distintas fuentes de paralelismo en el sistema, e incluso c´ omo interact´ uan entre s´ ı, presentando resultados obtenidos en un supercomputador de Cray. Como caso de uso, validamos la idoneidad de Chapel para la gesti´ on paralela de la E/S en problemas de almacenamiento de secuenciamiento gen´ omico, para los que hemos desarrollado un nuevo formato m´ as eficiente, lo que constituye la segunda contribuci´ on de la tesis. La novedad de este formato es que incluye la compresi´ on de datos y permite una recuperaci´ on eficiente de las secuencias tanto de forma secuencial como aleatoria. La organizaci´ on de esta memoria es la siguiente. En el cap´ ıtulo 1 presentamos el estado del arte en los sistemas de E/S, revisamos las distintas soluciones existentes desde el punto de vista hardware, resumimos las caracter´ ısticas y funcionalidades de los actuales sistemas de ficheros paralelos, presentamos los lenguajes de alta productividad, entre los que se encuentra Chapel as´ ı como las limitaciones de los mismos a la hora de realizar operaciones de E/S paralelas. En el cap´ ıtulo 2 proponemos una interfaz paralela de E/S en Chapel, damos detalles de su implementaci´ on y la validamos en una plataforma de supercomputaci´ on. En el cap´ ıtulo 3 proponemos un nuevo formato de almacenamiento de datos obtenidos por secuenciaci´ on gen´ etica m´ as eficiente que los publicados hasta la fecha, y analizamos su implementaci´ on paralela en Chapel. Concluimos en el cap´ ıtulo 4 sintetizando las principales aportaciones de esta tesis y discutiendo posibles l´ ıneas de trabajo futuro. 6 Cap´ ıtulo 1. Estudio del estado del arte en los sistemas de E/S Area Network). La conexi´ on directa (DAS), como su nombre indica, est´ a caracterizada por conectar, a trav´ es de un bus, el almacenamiento al host, sin usar ning´ un elemento de red (hub, router,switch, etc.). Pero si tenemos el sistema de almacenamiento conectado a un s´ olo nodo s´ olo ese nodo podr´ a acceder al almacenamiento. Para solucionar ese problema aparecen las NAS [55], que permiten a otros nodos el acceso a los datos usando protocolos c´ omo NFS [104] o samba (CIFS) [95] Para ello hay que conectar el almacenamiento, mediante una DAS o una SAN, a un servidor, y este proveer´ a a los clientes de una interfaz de acceso a los datos a trav´ es de la red, permitiendo a los clientes un acceso sencillo a los ficheros y a los datos que contienen. En resumen, lo que provee una NAS es un sistema de ficheros. Con el tiempo se crearon en las organizaciones multitud de islas de servidores con almacenamiento local, por lo que surgi´ o la necesidad de consolidar el almacenamiento. Para conseguir esto se desarrollaron las SAN [117], permitiendo la conexi´ on entre un sistema de almacenamiento y varios servidores, ya que una SAN provee dispositivos de bloques, y no sistemas de ficheros. Se suele usar el protocolo SCSI sobre determinados medios f´ ısicos (InfiniBand (IB), ethernet, FC,...). El hecho de que un sistema de almacenamiento proporcione dispositivos de bloque a los distintos servidores supone varias ventajas: i) se optimizan los recursos hardware; ii) se unifica el almacenamiento secundario para todos los servidores; iii) se reduce el espacio sobrante; iv) reduce el coste; y v) facilita la administraci´ on. Posteriormente, si hay necesidad de compartir datos con clientes, se puede montar una NAS, creando un sistema h´ ıbrido. Pero si muchos clientes necesitan acceder a grandes cantidades de datos el servidor NAS se convierte en el cuello de botella. La soluci´ on obvia es conectar los clientes directamente a la SAN, pero eso tiene el problema de que la SAN s´ olo provee acceso a nivel de bloque, no existe ninguna capa que provea una interfaz de sistema de ficheros en una SAN. Para tener un acceso coherente a esos sistemas de bloque aparecieron los sistemas de fichero de disco compartido (shared-disk file systems), que se encargan de proporcionar un acceso coherente a los datos, usando servidores de metadatos, y un acceso directo desde los clientes a los bloques que contienen los datos de los ficheros. De esa forma se evita que los datos tengan que pasar por un servidor antes de llegar a su destino, yendo directamente de los dispositivos de E/S al nodo que los solicit´ o. Eso significa que los clientes necesitan un acceso directo, a nivel de bloque, a la SAN. Esto se puede conseguir uni´ endolos a trav´ es de alguna red de interconexi´ on, generalmente usando fibre channel. Esta soluci´ on no es factible en el ´ ambito de HPC (High Performance Computing) donde podemos tener miles de nodos necesitando acceder a los datos y sin posibilidad de acceso directo a la SAN, ya que a pesar de tener sus propias redes de interconexi´ on de alta velocidad, estas no suelen ser compatibles con las redes que usan las SAN para dar acceso a los bloques. La soluci´ on actual al problema de acceso a los datos en entornos 1.1. Introducci´ on a la aceleraci´ on de la E/S 7 HPC se basa en los sistemas de ficheros distribuidos y paralelos (distributed parallel file systems). Veremos varios ejemplos en la secci´ on 1.2, e implementaciones concretas a partir de la secci´ on 1.4.2. En esta tesis se ha usado el sistema de ficheros Lustre [12] por ser el m´ as usado en sistemas HPC [109], como explicamos en la secci´ on 1.2.4. En todas las soluciones mencionadas el acceso a los ficheros se realiza a trav´ es de una interfaz POSIX [127], lo que permite que las aplicaciones sean independientes del hardware, del sistema de ficheros y del sistema operativo usados, consiguiendo una mayor productividad del programador al permitir el uso de los programas en una mayor variedad de entornos. 1.1.4. El sistema de almacenamiento terciario Aunque en este trabajo no se ha usado ning´ un sistema de almacenamiento terciario, por completitud se realiza aqu´ ı una introducci´ on a su uso y funcionamiento, ya que se usa en muchos sistemas HPC para el almacenamiento masivo de datos. En el almacenamiento terciario podemos mencionar los medios (cintas, discos bluray,...), los dispositivos de acceso a los medios (unidades de cinta, unidades blu-ray,...), y finalmente suele existir un mecanismo de intercambio autom´ atico. As´ ı podemos encontrar que m´ ultiples medios est´ an depositados en casillas dentro de un armario, junto a uno o m´ as dispositivos de acceso, y uno o m´ as brazos rob´ oticos se encargan del movimiento de los medios dentro del armario, entre las distintas posiciones que puede tomar. Estos armarios suelen tener una o m´ as casillas que ejercen de interfaz con el mundo exterior, de forma que se pueden extraer e incorporar cintas sin necesidad de interrumpir su operaci´ on. Encontramos que la caracter´ ıstica que define el almacenamiento terciario es que el acceso a los medios se realiza usando un brazo rob´ otico, que se encargar´ a de mover el medio de almacenamiento de datos, cinta o disco, a una unidad que se encargar´ a de acceder a los datos que contiene, lo que aumenta la latencia de acceso. Una vez insertado el disco o la cinta en la unidad de acceso, hay que moverlo hasta que el cabezal de acceso a los datos est´ a sobre la posici´ on en la que est´ an almacenados los datos, accederlos, ponerlo en un estado estable (en el caso de las cintas rebobinarlas) sacar el medio de la unidad y almacenarlo en un slot libre de la librer´ ıa rob´ otica. Todos esos pasos influyen en el rendimiento (ancho de banda y la latencia) y hacen complicado el realizar una previsi´ on de los tiempos de acceso. Aunque no vaya a desaparecer el almacenamiento terciario, lo que si est´ a cambiando es el prop´ osito con el que se emplea. Conforme los discos mejoran y se hacen m´ as baratos es m´ as l´ ogico emplearlos para hacer backups, pero las cintas a´ un tienen sus ventajas, como el poder mover con facilidad el backup fuera de las instalaciones y el menor coste, 8 Cap´ ıtulo 1. Estudio del estado del arte en los sistemas de E/S dependiendo del volumen de los datos que se necesite almacenar, y de los requisitos de recuperaci´ on de esos datos. En el futuro se estima que las cintas se usen como segunda o tercera l´ ınea de actuaci´ on en los backups. Los discos ´ opticos tampoco se pueden desechar como opci´ on de almacenamiento, ya que si el volumen de datos a acceder no es muy grande son mucho m´ as r´ apidos que las cintas, al permitir accesos aleatorios. As´ ı podemos encontrar que Facebook est´ a usando discos bluray para almacenar im´ agenes que se acceden poco, creando lo que llaman un cold storage (almacenamiento fr´ ıo) [130], ya que el 97 % del contenido recibe s´ olo un 29 % de las peticiones. 1.2. Sistemas de ficheros paralelos Los sistemas de ficheros paralelos convencionales suelen usar varios servidores con arrays de discos locales para servir los ficheros, pudiendo un fichero almacenarse de una de estas tres formas: i) en un s´ olo servidor, ii) replicado en varios de forma ´ ıntegra, o iii) repartido entre varios, de forma que cada uno de los servidores contiene un trozo. Se podr´ ıa pensar que si cada fichero se almacena en un s´ olo servidor, el sistema est´ a desaprovechando recursos, ya que no estar´ ıa consiguiendo aprovechar todo el ancho de banda existente. Pero hay que tener en cuenta que un sistema HPC en explotaci´ on puede estar orientado a permitir la ejecuci´ on concurrente de muchas aplicaciones, por lo que se producen multitud de operaciones de E/S simult´ aneamente a distintos ficheros. En estos escenarios se suele suponer un reparto homog´ eneo de peticiones entre los distintos servidores de E/S , por lo que si los ficheros son peque˜ nos, no distribuirlos es buena una opci´ on. Otra opci´ on, en caso de ficheros de gran tama˜ no, pasa por distribuirlos uniformemente entre los servidores de E/S. Sin embargo, si tenemos muchos servidores ser´ ıa ineficiente dividir los ficheros entre todos, ya que se podr´ ıa provocar una sobrecarga de peticiones hacia todos los servidores, por lo que en la pr´ actica el n´ umero de servidores que se utiliza para distribuir un fichero suele ser relativamente bajo. En general, los sistemas de ficheros paralelos se aprovechan de que existen multitud de aplicaciones en ejecuci´ on simult´ aneamente para, adem´ as de intentar optimizar el acceso a cada uno de los ficheros, optimizar el throughput total del sistema. Vamos a ver en las siguientes secciones la arquitectura y caracter´ ısticas principales de los sistemas de ficheros paralelos que constituyen el estado del arte en los grandes sistemas HPC en al actualidad: GPFS (General Parallel File System), Ceph, EOS y Lustre, este ´ ultimo el m´ as popular en la mayor´ ıa de los grandes sistemas HPC. Todos ellos son distribuidos y permiten el acceso a los datos desde distintos servidores en paralelo. Entre la lista de sistemas de ficheros paralelos a analizar no hemos incluido algunos sistemas de ficheros muy populares actualmente, como el GFS (Google File System) o el HDFS (Hadoop Distributed File System) debido a que son sistemas de ficheros hechos a medi- 1.2. Sistemas de ficheros paralelos 9 da para su uso en entornos concretos. En particular, est´ an dise˜ nados para ser usados en aplicaciones que se basan en el paradigma map-reduce, y por tanto no son de prop´ osito general. Uno de los objetivos de estos sistemas de ficheros es acercar los c´ alculos a los datos en vez de los datos a los c´ alculos, lo que puede ser una estrategia prometedora en el ´ ambito de aplicaciones data-intensive. 1.2.1. GPFS GPFS [107] es un sistema de ficheros paralelo y distribuido para clusters desarrollado por IBM. Intenta emular el comportamiento de un sistema de ficheros POSIX, haciendo creer a cada nodo del cluster que tiene un sistema de ficheros local. El nombre ha cambiado recientemente de GPFS a Spectrum Scale. Un sistema de ficheros GPFS est´ a formado por un conjunto de arrays de discos que contienen los datos y los metadatos de los sistemas de ficheros. Cada sistema de ficheros se puede acceder desde todos los nodos del cluster usando el interfaz est´ andar POSIX. GPFS se encarga de mantener la coherencia y consistencia de m´ ultiples accesos desde distintos nodos, usando bloqueos a nivel de byte, administraci´ on distribuida de bloqueos y cuadernos de bit´ acora (journaling). Gracias a todo ello no hace falta modificar las aplicaciones POSIX para poder ejecutarlas en un sistema que use GPFS. Adem´ as del interfaz POSIX, GPFS a˜ nade un conjunto de interfaces para aumentar la funcionalidad que ofrece a las aplicaciones, con el objetivo de que desde la aplicaci´ on se puedan indicar pistas sobre el patr´ on de acceso a los ficheros. Esa informaci´ on se usar´ a para optimizar el acceso a los discos, las precargas o el uso de las caches [100]. Adem´ as GPFS tambi´ en es capaz de detectar algunos patrones de acceso en tiempo de ejecuci´ on. Por ejemplo, puede detectar un patr´ on de acceso secuencial, secuencial inverso y aleatorio. Cuando se crea un sistema de ficheros GPFS se le asignan conjuntos de dispositivos de bloques, llamados NSD (Network Shared Disks) en la terminolog´ ıa GPFS. Cada NSD puede ser accedido por uno o m´ as nodos del cluster GPFS, de forma que se puede tener redundancia y en caso de que caiga un servidor otro se encargue de ese NSD. El problema de la consistencia debido el acceso simultaneo desde distintos nodos a los datos y los metadatos se soluciona usando un sistema de tokens administrados de forma distribuida que act´ uan como cerrojos (locks), y se encargan de coordinar los accesos a los NSD. La responsabilidad de la administraci´ on de los tokens se asigna din´ amicamente a uno o m´ as nodos del cluster GPFS, esto permite una mayor escalabilidad cuando se est´ an usando muchos ficheros en momentos de mucha carga de trabajo. Existen tres formas b´ asicas de configurar un cluster GPFS [45]: i) usando discos compartidos entre todos los nodos; ii) usando servidores de E/S; iii) mezclando las dos anteriores. 10 Cap´ ıtulo 1. Estudio del estado del arte en los sistemas de E/S Figura 1.1: Configuraci´ on b´ asica de GPFS usando discos compartidos a trav´ es de una SAN. La configuraci´ on m´ as b´ asica para tener un cluster GPFS es usando discos compartidos, donde el almacenamiento est´ a conectado a trav´ es de una SAN a todas las m´ aquinas del cluster, como se puede ver en la figura 1.1. Esto significa que los dispositivos de bloques se pueden acceder directamente desde todos los nodos usando un protocolo como SCSI o similar, a trav´ es de una red como fibre channel, infiniband, u otras, y usando elementos de interconexi´ on como switches. Esta configuraci´ on, usando discos compartidos, se suele usar para dar servicio de NFS a trav´ es de cNFS (clustered NFS), de forma que los nodos GPFS son los que sirven en paralelo el protocolo NFS. El problema de esta configuraci´ on es que conforme los requerimientos de capacidad y procesamiento crecen, las tecnolog´ ıas de conexi´ on pueden no ser las apropiadas para un cluster con muchos nodos, por lo que en sistemas muy grandes se suele usar una de las otras dos configuraciones. La configuraci´ on m´ as usual en sistemas HPC consiste en usar servidores GPFS de E/S que comparten los discos a trav´ es de una SAN, y que sirven los sistemas de ficheros al resto de nodos como se ve en la figura 1.2. El acceso al sistema de ficheros desde los nodos se realiza usando un protocolo llamado NSD (Network Shared Disk), que provee una interfaz a nivel de bloques sobre las redes que haya disponibles, por ejemplo usando TCP/IP con ethernet o verbs con infiniband. En cualquier caso, para las aplicaciones que usan el interfaz GPFS no hay ninguna diferencia, salvo la de rendimiento, entre acceder los datos directamente por la SAN, o usando servidores de E/S (llamados NSD servers en la terminolog´ ıa de GPFS). La decisi´ on de cu´ antos nodos configurar como servidores de E/S est´ a basada en los requerimientos de rendimiento, la arquitectura de la red y las prestaciones del sistema de almacenamiento. 1.2. Sistemas de ficheros paralelos 11 Figura 1.2: Configuraci´ on GPFS usando servidores de E/S, que sirven los datos a los clientes, y a su vez acceden a los bloques contenidos en discos compartidos a trav´ es de una SAN. Otra configuraci´ on posible es combinar las dos anteriores, algunos nodos se usar´ an como servidores de E/S y otros acceder´ an directamente a los datos. Esto es sencillo de configurar, ya que a los nodos se les puede indicar m´ ultiples caminos para el acceso a los sistemas de ficheros, usando por defecto siempre el camino m´ as r´ apido. Si un nodo tiene un HBA (Host Bus Adapter) para acceder a la SAN, lo usar´ a, y si no, o si se corta la conexi´ on por el HBA, entonces acceder´ a a trav´ es de los servidores de E/S, usando otra de las redes disponibles, como ethernet. Usando esta configuraci´ on de acceso mixto, en el que algunos nodos acceden de forma directa y otros no, se pueden configurar de forma directa los grandes consumidores o productores de datos, por ejemplo los nodos de backups, y el resto configurarlos para que accedan a trav´ es de los servidores de NSD. La decisi´ on de usar una conexi´ on directa a la SAN, usar la soluci´ on con servidores de E/S, o la soluci´ on mixta, es una decisi´ on tanto de rendimiento como econ´ omica. El conectar todos los nodos del sistema mediante una SAN exige poner una red de interco- 12 Cap´ ıtulo 1. Estudio del estado del arte en los sistemas de E/S nexi´ on con una tecnolog´ ıa apropiada para obtener un buen rendimiento, pero dependiendo de para que se vaya a usar, eso no ser´ a rentable, ya que el cuello de botella no ser´ a la red de interconexi´ on. Si adem´ as se consiguen paralelizar los accesos desde distintos nodos la necesidad de tener una red de muy alta velocidad entre los nodos ser´ a menor a´ un. As´ ı, cuando hay muchos clientes, ser´ a m´ as rentable el usar servidores de E/S para el acceso al sistema de ficheros GPFS. 1.2.2. Ceph Ceph [136] es un sistema de ficheros paralelo que sigue a POSIX de forma relajada. Hay dos caracter´ ısticas en las que se desv´ ıa del est´ andar: en la medida del espacio ocupado por ficheros distribuidos y en la atomicidad de las escrituras desde distintos nodos. El aspecto diferenciador con otros sistemas de ficheros paralelos, desde el punto de vista de la arquitectura, es que en Ceph se separan los datos y metadatos, sustituyendo el interfaz tradicional del sistema de bloques con uno en el que los clientes acceden a los objetos por su nombre, dejando que los dispositivos realicen la decisi´ on de c´ omo distribuir los bloques de datos a bajo nivel. As´ ı, los clientes comunican las operaciones de metadatos (como open, rename,..) al MDS (MetaData Server) mientras las operaciones con datos (escrituras y lecturas) se realizan directamente en los OSD (Object Storage Device). El punto fuerte de Ceph es la eliminaci´ on de las tablas de asignaci´ on de ficheros (file allocation table, fat), reemplaz´ andolas con una funci´ on CRUSH (Controlled Replication Under Scalable Hashing) que se encarga de distribuir los datos entre los OSD de una forma semialeatoria. Para realizar esa distribuci´ on la funci´ on CRUSH [137] se comporta de forma parecida a un hash, pero un hash normal no ser´ ıa efectivo debido a la naturaleza din´ amica del almacenamiento, en el que se van a˜ nadiendo y desapareciendo recursos conforme las necesidades cambian, o se van rompiendo. La idea central de Ceph es distribuir los datos para realizar un uso eficiente tanto del almacenamiento como de los anchos de banda. Esto se consigue haciendo que los datos que se van incorporando al sistema se vayan distribuyendo de forma aleatoria entre todos los nodos, cuando se incorporan nuevos dispositivos o nodos, se migran subconjuntos aleatorios a los nuevos dispositivos y cuando desaparecen dispositivos se redistribuyen de forma uniforme los datos que conten´ ıan. El objetivo de todo esto es evitar los desbalanceos y las asimetr´ ıas de cargas (que los datos m´ as accedidos no est´ en bien repartidos entre todos los dispositivos) entre los dispositivos existentes. Para conseguir estos objetivos, Ceph primero mapea los objetos en grupos de posicionamiento (Placement Groups oPGs) usando una funci´ on de hash simple junto a una m´ ascara de bits que controla el n´ umero de PGs. Los PGs se asignan a los OSD usando la funci´ on CRUSH, que mapea cada PGs a una lista ordenada de OSD en los que se almacenar´ an las r´ eplicas, usando una funci´ on de distribuci´ on semialeatoria. Para localizar un objeto, la funci´ on CRUSH s´ olo requie- 1.2. Sistemas de ficheros paralelos 13 re el PGs y el mapa de clusters OSD cluster map, que es una descripci´ on compacta y jer´ arquica de los dispositivos que forman el cluster de almacenamiento. Esto tiene dos ventajas principales, una es que est´ a completamente distribuido, cualquier elemento del sistema puede calcular de forma independiente la localizaci´ on de cualquier objeto, y la otra ventaja es que el mapa es muy raro que se actualice, con lo que se consigue eliminar cualquier intercambio de metadatos relacionados con esta distribuci´ on. Un problema es mantener la funci´ on CRUSH aunque cambien los mapas. La soluci´ on es etiquetar las distintas versiones de los mapas de un sistema, e indicar cu´ al es el mapa a usar en cada caso. Una de las caracter´ ısticas m´ as importantes es la replicaci´ on de los datos, ya que se intenta abaratar los costes no usando RAIDs de discos, u otros complicados y caros sistemas de redundancia. La soluci´ on pasa por escribir cada dato en m´ as de un sitio, y cuando el primario falla se intenta acceder a los replicados. El OSD primario solicita realizar copias en otros OSDs, y ´ estas no se dan por realizadas hasta que todas las r´ eplicas han sido procesadas [136]. Desde el punto de vista de la existencia de problemas en un OSD, se tienen dos variables en Ceph [135] que determinan la vitalidad del OSD, la accesibilidad y la asignaci´ on de datos. Un OSD que no est´ e accesible se marca como down, y sus responsabilidades pasan al siguiente OSD del grupo de posicionamiento (PG) al que pertenece. Si no se recupera en poco tiempo se marca como out, se queda fuera de la distribuci´ on de datos, y otro OSD se une a cada uno de los grupos de posicionamiento PG a los que pertenec´ ıa el OSD ca´ ıdo, y replica su contenido. Los clientes que ten´ ıan operaciones pendientes con ese OSD simplemente las reenv´ ıan al nuevo OSD primario. 1.2.3. EOS El CERN (Conseil Europ´ een pour la Recherche Nucl´ eaire) es uno de los mayores generadores de datos del planeta, por lo que tiene unas necesidades de almacenamiento especiales, que les ha llevado a buscar soluciones a medida [39]. Experimentos como Atlas oAlice generan Petabytes de datos que tienen que ser accedidos y procesados por miles de f´ ısicos de todo el planeta. Los datos est´ an almacenados en dos infraestructuras, CASTOR y EOS [39]. CASTOR se encarga del almacenamiento en cinta, y EOS del almacenamiento en discos, siendo ambos considerados sistemas de administraci´ on de contenidos (Content Management System, CMS). EOS [2] es un sistema de ficheros paralelo desarrollado y mantenido principalmente por el CERN desde 2010. Sigue una sem´ antica POSIX simplificada, con lo que no es POSIX, s´ olo soporta lo que les interesa, que es almacenar y procesar los datos que generan los instrumentos del CERN [1], que suelen generar ficheros de gran tama˜ no. Para 14 Cap´ ıtulo 1. Estudio del estado del arte en los sistemas de E/S el acceso a esos datos se usan diferentes interfaces, como XrootD [34], https, etc. Desde 2013 XrootD se usa como la forma principal de acceso a los datos de Atlas [53]. La arquitectura de EOS est´ a basada en tres elementos [2]: i) Un administrador (MGM) que se ejecuta en un nodo, pudiendo tener otro de fail-over ; ii) una cola de mensajes Message Queue, MQ, que coordina los mensajes as´ ıncronos del sistema, y que se suele ejecutar en el MGM; iii) un componente de almacenamiento de ficheros File Storage Component, FSC, que almacena los datos y los transfiere desde y hacia los clientes, y se ejecuta en los servidores de ficheros. As´ ı, un cliente le pide al MGM una operaci´ on sobre un fichero, y ´ este le redirige al FST adecuado que se encargar´ a de la operaci´ on. 1.2.4. Lustre La primera versi´ on de Lustre [12] se present´ o en diciembre del a˜ no 2003, y a partir de ese momento ha ido evolucionando, apoyado por empresas, centros de investigaci´ on y universidades. Lustre ha evolucionado desde entonces, con una mejora continua del rendimiento, de la funcionalidad, y de la fiabilidad. Desde el principio, Lustre es un sistema de ficheros paralelo y de c´ odigo fuente abierto (open source). Lustre es usado en aproximadamente la mitad de los 100 ordenadores m´ as r´ apidos de forma continua desde 2005, y por la mayor´ ıa de los sistemas HPC del planeta [109]. Lustre es escalable, y un ´ unico sistema Lustre puede dar servicio a varios clusters independientes con decenas de miles de clientes y decenas de Petabytes de almacenamiento distribuido en miles de servidores. Ese es el caso, por ejemplo, de uno de los sistemas de E/S que hemos usado en este trabajo, Spider, y que describimos en la secci´ on 1.4.2. Es por todo lo anterior por lo que se ha decidido usar Lustre como el sistema de ficheros con el que realizaremos la evaluaci´ on experimental de nuestra librer´ ıa de E/S paralela para Chapel. Los componentes de un sistema Lustre son: i) un servidor de administraci´ on (Management Server, MGS), que contiene la informaci´ on de la configuraci´ on del sistema; ii) un servidor de metadatos (MetaData Server, MDS) que administra los nombres, directorios y metadatos en general del sistema de ficheros; iii) servidores de almacenamiento de objetos (Object Storage Servers, OSS) que proveen los servicios de E/S de los ficheros. Los metadatos se almacenan en un (MetaData Target, MDT), mientras los datos se almacenan en (Object Storage Targets, OSTs). Cada OSS puede servir varios OSTs. Tanto los MDTs como los OSTs son interfaces a dispositivos de bloques, generalmente LUNs (Logical Unit Number) SCSI, que los MDS y OSS acceder´ an mediante una SAN, tal y como se explic´ o en la secci´ on 1.1.3. En la figura 1.3 se puede ver un sistema Lustre que combina varias de las soluciones 1.2. Sistemas de ficheros paralelos 15 Clients OSS Servers and OSTs MDS Servers MDT OST OST OST OST Figura 1.3: Un ejemplo de sistema Lustre, con un MDT, cuatro OSTs, dos MDS y tres OSSs. Existen canales redundantes entre los dos ´ ultimos OSTs y OSSs para proporcionar tolerancia a fallos hardware o software. que se suelen usar. Por un lado el MGS se suele instalar en los nodos MDS, ya que consume muy pocos recursos, tanto de espacio como de ancho de banda. Un sistema Lustre s´ olo puede tener un MDT y un MDS, por lo que se pueden usar varios nodos para tener un respaldo en caso de que falle el que se est´ a usando. Es por eso por lo que aparecen dos nodos MDS que comparten un MDT en la figura 1.3. Es habitual que un OSS administre varios OSTs. Aunque un OST s´ olo puede estar administrado por un OSS, es normal que tengan conexiones a varios OSSs por si uno falla que otro se pueda hacer cargo de ´ el. Las redes de interconexi´ on pueden ser ethernet, infiniband u otras. Podemos encontrar dos fuentes de paralelismo en el sistema de ficheros Lustre. Por un lado tenemos el uso de m´ ultiples OSTs en paralelo gracias a accesos concurrentes de distintos procesos. Por otro lado, el acceso a un ´ unico fichero se realiza en paralelo ya que normalmente est´ a particionado en trozos (llamados stripes), repartidos entre m´ ultiples OSTs. Esto provoca efectos anti-intuitivos en el acceso a los datos, ya que, por ejemplo, lo normal es que sea m´ as r´ apido usar m´ as OSTs que usar menos. Sin embargo, en [29] se puede leer que conforme se usan m´ as OSTs el sistema es mucho m´ as lento. Por ejemplo, con 104 agregadores, usando 16 OSTs obtienen un ancho de banda de 1200 MB/s, mientras usando 2 OSTs es de unos 3000 MB/s. Es decir, puede ser m´ as ´ optimo 22 Cap´ ıtulo 1. Estudio del estado del arte en los sistemas de E/S de Lustre de forma exhaustiva, como por ejemplo en [143][109][98][97][112][111][99]. Sin embargo, ninguno de los trabajos anteriores estudia los problemas de mejora de la productividad que afrontamos en esta tesis y que ser´ an abordados en los pr´ oximos cap´ ıtulos. 1.3. Lenguajes de alta productividad Como hemos comentado anteriormente, el paso a sistemas Exascale requerir´ a del desarrollo de nuevos modelos de programaci´ on capaces de gestionar billones de elementos de procesamiento concurrentes as´ ı como de manejar eficientemente la E/S. El desaf´ ıo mayor aparecer´ a en las grandes aplicaciones de tipo data-intensive (Big-Data). Una opci´ on para programar productivamente los sistemas de Exascale pasa por crear nuevos lenguajes de programaci´ on, y de hecho es uno de los objetivos tanto del proyecto HPCS como del programa FETHPC-H2020. Sin embargo, es una actividad altamente arriesgada, ya que se puede invertir mucho esfuerzo en su creaci´ on sin que creen un impacto en la comunidad cient´ ıfica. Podemos encontrar ejemplos de proyectos anteriores de gran envergadura, en los que se ha intentado crear nuevos lenguajes de programaci´ on que mejoraran los ya existentes. Por ejemplo, entre las d´ ecadas de 1970 y 1980 se intent´ o substituir al COBOL y al Fortran, los lenguajes m´ as populares en aquellos momentos, por un nuevo lenguaje, el Ada, pero a pesar de todos los recursos que se invirtieron en el proyecto, no lleg´ o a ser lo popular que se esperaba, limit´ andose hoy en d´ ıa su uso al ´ ambito del software empotrado. Otro proyecto similar al HPCS, el de 5th-generation, se cre´ o en Jap´ on en la d´ ecada de los 80, donde se dieron un plazo de diez a˜ nos para encontrar un nuevo modelo de programaci´ on m´ as productivo. Este esfuerzo dio como resultado el CPL (Concurrent Logic Programming), un dialecto de Prolog orientado al paralelismo y a obtener un alto rendimiento, pero que no ha transcendido m´ as all´ a del ´ ambito educativo. El motivo de esos fracasos se puede encontrar en que la comunidad de programadores ha valorado m´ as la portabilidad, el rendimiento y la evoluci´ on incremental de los lenguajes ya en uso, antes que la elegancia, la expresividad o incluso su facilidad de uso de uno nuevo, al menos en la programaci´ on de sistemas HPC. El dise˜ nar una aplicaci´ on que se ejecuta en un sistema de memoria distribuida presenta el gran inconveniente de que, puesto que el espacio de direcciones est´ a repartido entre los nodos del sistema, cada nodo tiene que gestionar el acceso a posiciones de memoria remotas (es decir, alojadas en otro nodo) a trav´ es de rutinas de comunicaciones de bajo nivel. Inicialmente la programaci´ on en estos sistemas se basaba en el uso de librer´ ıas como PVM (Parallel Virtual Machine) [54] o MPI (Message Passing Interface) [52], que exponen al programador un API para la gesti´ on de comunicaciones remotas 1.3. Lenguajes de alta productividad 23 punto a punto, comunicaciones globales colectivas, operaciones de sincronizaci´ on, etc. Pero los algoritmos implementados con estas librer´ ıas son complicados de programar, depurar y mantener. Para simplificar la gesti´ on de los accesos remotos, se propuso un nuevo modelo de programaci´ on paralelo denominado PGAS (Partitioned global Address Space) [140], que ofrece al programador la visi´ on de un espacio de direcciones global y l´ ogicamente compartido entre los procesos, dejando en manos del compilador y runtime la gesti´ on de los accesos remotos, puesto que los datos est´ an f´ ısicamente distribuidos entre los nodos. En este modelo se inspiran lenguajes como UPC (Unified Parallel C), [24], X10 [18], Fortress [40], o Chapel [15]. Este paradigma simplifica enormemente la tarea de programaci´ on desde el punto de vista del usuario. El modelo de programaci´ on PGAS intenta conjugar el rendimiento que aporta el acceso a datos alojados localmente en cada nodo, con la simplicidad y programabilidad de un modelo de memoria compartida, puesto que ofrece un espacio de direccionamiento global, que es directamente accesible por cualquier proceso. 1.3.1. Chapel como ejemplo de lenguaje PGAS Chapel [15] es un lenguaje paralelo que forma parte del programa HPCS (High Productivity Computing Systems) del DARPA, para mejorar la productividad del usuario en sistemas masivamente paralelos. Chapel soporta el modelo de programaci´ on PGAS (Partitioned Global Address Space). Para ello provee una visi´ on global de las estructuras de datos, as´ ı como una visi´ on global del control. Una ventaja de Chapel respecto al resto de lenguajes HPCS, es que es un lenguaje multiresoluci´ on. Con esto nos referimos al hecho de que ofrece un buen n´ umero de funcionalidades de alto y bajo nivel que permiten trabajar con distintos niveles de abstracci´ on y control, en un entorno unificado. Por ejemplo las funcionalidades de bajo nivel permiten mayor control del hardware y que el programador pueda acceder a detalles de implementaci´ on del runtime. Por el contrario, las funcionalidades de alto nivel abstraen el hardware y permiten definir operaciones de manera global, sin preocuparse por introducir comunicaciones para los accesos remotos, sincronizaciones para los accesos compartidos, etc. De hecho, las funcionalidades de alto nivel en Chapel est´ an basadas en las de bajo nivel, para asegurar que todas son interoperables. Un ejemplo de multiresoluci´ on es el soporte que ofrece Chapel para expresar paralelismo de tarea o paralelismo de datos. En general, el paralelismo de tarea ofrece funcionalidades de bajo nivel como crear tareas, despacharlas, sincronizarlas, etc. Por el contrario, el paralelismo de datos ofrece funcionalidades de alto nivel como la definici´ on de distribuciones de datos y la gesti´ on autom´ atica del reparto del trabajo y de las comu- 24 Cap´ ıtulo 1. Estudio del estado del arte en los sistemas de E/S nicaciones. Por ejemplo, una funcionalidad de alto nivel que ayuda a los programadores a razonar sobre localidad es el concepto de locale, que en Chapel es un tipo predefinido. Un locale es una representaci´ on abstracta, para una arquitectura concreta, de los datos que est´ an f´ ısicamente compartidos dentro de un nodo (son locales a ´ el). El acceso a un dato por parte de una tarea es local si la tarea y el dato est´ an mapeados en el mismo locale, y remote en cualquier otro caso. En una arquitectura paralela convencional un locale har´ a referencia a un nodo del sistema donde el espacio de direcciones es compartido por todos los cores/threads del nodo. El programador puede controlar, usando cl´ ausulas espec´ ıficas de bajo nivel, el que una tarea o thread se ejecute en un determinado locale. En los cap´ ıtulos 2 y 3 se usar´ a frecuentemente el t´ ermino locale en referencia a los nodos de un sistema HPC. Para m´ as detalles consultar [15]. #include <omp.h> #include <iostream> #include <cmath> using namespace std; int main() { int i, n, chunk; float a[1000], b[1000],result; // Some initializations, done sequentially n = 1000; chunk = 100; result = 0.0; for (i=0; i < n; i++) { a[i] = 1.0/(i+1); b[i] = 1.0/(i+1); } #pragma omp parallel for default(shared) private(i) schedule(static,chunk ) reduction(+:result) for (i=0; i < n; i++) result = result + (a[i] *b[i]); std::cout << "Final result=" << result << "so pi=" << sqrt(6*result) << std::endl; } Figura 1.5: C´ odigo OpenMP para calcular pi en paralelo [13] 1.3. Lenguajes de alta productividad 25 En la figura 1.5 podemos ver un c´ odigo paralelo para calcular el n´ umero pi realizado en OpenMP. En la figura 1.6 tenemos el mismo algoritmo programado en MPI, y finalmente en la figura 1.7 podemos ver el mismo algoritmo implementado en Chapel. En este ´ ultimo caso vemos que el algoritmo s´ olo ocupa una l´ ınea, la 2, que, al igual que en los lenguajes anteriores, se encarga de ejecutar N iteraciones de forma paralela, reparti´ endolas autom´ aticamente entre todos los cores que haya disponibles. Despu´ es se reducen los resultados sumando todos los valores generados en los distintos threads, y finalmente se multiplica la reducci´ on por 6 y se hace la ra´ ız cuadrada, obteniendo tras esa operaci´ on una aproximaci´ on de pi. En el ejemplo de Chapel, las variables que van almacenando los resultados intermedios, y que ser´ an usados en la reducci´ on, se definen de forma impl´ ıcita, es decir, no aparecen. Adem´ as el c´ odigo Chapel, permite cambiar el valor de N sin necesidad de cambiar el c´ odigo. Para ello, en la l´ ınea 1 de la figura 1.7 se define N como de tipo config, lo que permite que se pueda pasar ese valor como un argumento en la l´ ınea de comandos al invocar el ejecutable. 1.3.2. E/S en lenguajes PGAS En general los lenguajes de programaci´ on PGAS presentan una interfaz directa para hacer E/S, con un acceso directo a la interfaz POSIX, pudiendo abrir (open), buscar (seek), leer (read), escribir (write) y cerrar (close) los ficheros, pero es muy complicado el usar esas funciones de forma distribuida y eficiente, sin crear contenci´ on o condiciones de carrera en los accesos, por lo que en general estos lenguajes no aprovechan las posibilidades ofrecidas por los sistemas de ficheros paralelos. Por ejemplo, tanto X10 [18] como Fortress [40] tienen s´ olo implementado la interfaz POSIX [127] para realizar E/S, y no tienen a´ un ni en proyecto la implementaci´ on de un nuevo interfaz. En el lenguaje UPC s´ ı crearon un documento [37] cuya ´ ultima versi´ on es del a˜ no 2006, donde se definen una serie de propuestas de adiciones y extensiones al lenguaje UPC, con la idea de que una vez que est´ en implementadas y sean estables se incorporen a las especificaciones del lenguaje. Estas funciones se a˜ nadieron a la versi´ on 1.2 de las especificaciones del lenguaje UPC [22], siendo implementadas en algunos compiladores [128], pero no aparecen en la ´ ultima versi´ on, la 1.3 [23] y no se garantiza su soporte a largo plazo. Entre las caracter´ ısticas b´ asicas de esas propuestas est´ an que las nuevas funciones ser´ an colectivas, lo que significa que ser´ an llamadas a la vez desde todos los threads del programa en ejecuci´ on. Adem´ as se apuesta por un modelo de consistencia d´ ebil, por lo que se ha de cerrar el fichero o realizar una operaci´ on de sincronizaci´ on expl´ ıcita en todos los threads para garantizar que cualquier thread pueda acceder correctamente a los datos almacenados por otros. Una restricci´ on de estas propuestas es que no se podr´ an usar a la vez las funciones POSIX y las definidas por UPC-IO. 26 Cap´ ıtulo 1. Estudio del estado del arte en los sistemas de E/S #include <cstdio> #include <cmath> #include <cstdlib> #include <mpi.h> #include <unistd.h> #include <string.h> using namespace std; int main(int argc, char **argv) { int i,n=1000, nthreads; double p_sum,sum; char *CPU_name; int loop_min, loop_max, tid; MPI_Init(&argc,&argv); // get the thread ID MPI_Comm_rank(MPI_COMM_WORLD, &tid); // get the number of threads MPI_Comm_size(MPI_COMM_WORLD, &nthreads); n=1000; sum = 0.0; p_sum = 0.0; CPU_name = (char *)calloc(80,sizeof(char)); gethostname(CPU_name,80); printf("thread %d running on machine = %s\n",tid,CPU_name); loop_min = 1 + (int)((long)(tid + 0) *(long)(n)/(long)nthreads); loop_max = (int)((long)(tid + 1) *(long)(n)/(long)nthreads); printf("thread %d loop_min= %i loop_max= %i\n",tid,loop_min, loop_max); for(i=loop_min;i<loop_max;i++) p_sum += 1.0/(i*i); printf("thread %d partial sum= %f\n", tid,p_sum); MPI_Reduce(&p_sum,&sum,1,MPI_DOUBLE,MPI_SUM,0,MPI_COMM_WORLD); if (tid == 0) printf("sum = %f so pi = %f\n",sum,sqrt(6*sum)); MPI_Finalize(); } Figura 1.6: C´ odigo MPI para calcular pi en paralelo [13]. 1.3. Lenguajes de alta productividad 27 1config const N = 1000; 2var pi = sqrt(6.0 *(+ reduce [ i in 1..N ] 1.0/(i*i) )); 3write ("pi =",pi); Figura 1.7: C´ alculo paralelo de pi en el lenguaje de alta productividad Chapel. En general, la mayor´ ıa de los lenguajes PGAS cuando quieren gestionar eficientemente las operaciones de E/S recurren a una librer´ ıa que act´ ue de intermediaria entre la aplicaci´ on y el sistema de ficheros. La librer´ ıa m´ as usada en estos casos es MPI IO. 1.3.2.1. MPI IO MPI IO es la librer´ ıa m´ as usada para realizar accesos concurrentes sobre sistemas de ficheros paralelos. Su desarrollo comenz´ o en 1994 en IBM, la NASA la adopt´ o en 1996 y en 1997 se incorpor´ o al est´ andar MPI-2. En general se ha realizado poca investigaci´ on en MPI-IO debido a que es complejo obtener un buen rendimiento [36][29][143]. Adem´ as el interfaz en s´ ı no se puede cambiar sin romper la compatibilidad, por lo que es complicado plantear mejoras del interfaz. Podemos ver en la figura 1.8 un ejemplo de c´ odigo MPI-IO para escribir de forma paralela en un fichero. Antes, hay que inicializar la librer´ ıa MPI, incluyendo llamadas a MPI Init,MPI Comm rank,MPI Comm size, y varias a MPI Bcast para difundir los datos que son comunes a todos los procesos MPI. En la l´ ınea 3 vemos una llamada aMPI Dims create, que crear´ a un sistema cartesiano de procesadores para describir la disposici´ on virtual de los nodos que implementar´ an la E/S. Tras eso, en la l´ ınea 6, se llama a la funci´ on MPI Type create darray para crear un nuevo tipo de fichero que es devuelto en file type, y que contiene la descripci´ on de c´ omo se proyectar´ an los datos de las distintas tareas que se van a encargar de la E/S. Tras crear el nuevo tipo de fichero, en la l´ ınea 7 se comunica a MPI para que pueda ser usado. Tras eso abrimos el fichero en la l´ ınea 10, con lo que creamos una nueva vista del fichero, para que cada proceso sepa c´ omo mapear sus datos. Esta funci´ on debe ser llamada desde todos los procesos, para que cada uno obtenga su propio mapeo de los datos. Tras definir la vista, en la l´ ınea 12 se llama a la funci´ on que realmente realiza el acceso al sistema de ficheros en paralelo desde todos los nodos implicados, MPI - File write all (o a MPI File read all en caso de que sea una lectura) para que se realice la escritura (o lectura). Tras los accesos todos los procesos pueden cerrar el fichero, como se realiza en la l´ ınea 13. La librer´ ıa MPI IO tiene optimizaciones particulares para varios sistemas de ficheros. 28 Cap´ ıtulo 1. Estudio del estado del arte en los sistemas de E/S ... 1 MPI_Comm_size(MPI_COMM_WORLD, \&pool_size); 2 ndims= number_of_dimensions_of_the_array; 3 MPI_Dims_create(pool_size, ndims, array_of_psizes); ... 4 MPI_Comm_rank(MPI_COMM_WORLD, \&my_rank); 5 order = MPI_ORDER_C; 6 MPI_Type_create_darray(pool_size, my_rank, ndims, array_of_gsizes, array_of_distribs, array_of_dargs, array_of_psizes, order, MPI_INT, \&file_type); 7 MPI_Type_commit(\&file_type); ... 8 MPI_Type_extent(file_type, \&file_type_extent); 9 MPI_Type_size(file_type, \&file_type_size); ... 10 file_open_error = MPI_File_open(MPI_COMM_WORLD, file_name, MPI_MODE_CREATE | MPI_MODE_WRONLY, MPI_INFO_NULL, \&fh); ... 11 MPI_File_set_view(fh, 0, MPI_INT, file_type,"native", MPI_INFO_NULL) ; ... 12 file_write_error=MPI_File_write_all(fh, write_buffer, write_buffer_size, MPI_INT, \&status); ... 13 MPI_File_close(\&fh); ... Figura 1.8: C´ odigo MPI IO para escribir en paralelo un array distribuido. En la pr´ actica, dependiendo del sistema de ficheros subyacente y de la distribuci´ on de los datos, las capas m´ as bajas de MPI IO pueden decidir que o bien todos los nodos realicen operaciones en paralelo de E/S o bien que s´ olo un subconjunto de los nodos realicen las operaciones de E/S. Esos nodos son los denominados agregadores [101]. Tambi´ en en algunas implementaciones, se da la posibilidad de elegir el algoritmo para implementar la E/S, realizando la elecci´ on a trav´ es de una variable de entorno. Por ejemplo, en la implementaci´ on de MPI IO en sistemas de Cray se usa la variable MPICH MPIIO CB ALIGN para elegir entre distintas estrategias de E/S. 1.4. Sistemas usados durante el desarrollo de esta tesis Los sistemas con los que hemos trabajado en esta tesis tienen en com´ un que sus sistemas de almacenamiento se basan en sistemas de ficheros distribuidos. Inicialmente usamos Crow, un sistema de desarrollo y pruebas interno de la empresa Cray. Posterior- 1.4. Sistemas usados durante el desarrollo de esta tesis 29 mente usamos Jaguar, el que en ese momento era el n´ umero 1 de la lista top500 [88], instalado y operado por el OLCF (Oak Ridge Leadership Computing Facility) del ORLN (Oak Ridge National Laboratory) y que ten´ ıa un sistema de ficheros con nombre propio, Spider [110]. En la descripci´ on que realizamos a continuaci´ on comenzamos por Crow, el ordenador de pruebas de Cray. A continuaci´ on describimos el entorno del OLCF en Jaguar, donde primero detallamos el sistema de E/S (Spider), despu´ es la red de interconexi´ on (SION) y las plataformas de c´ alculo (EOS) del OLCF de las que presentamos resultados en este trabajo. Finalmente describimos en este cap´ ıtulo la generaci´ on actual de Picasso, el superordenador de la Universidad de M´ alaga, ubicado en el SCBI, en el que tambi´ en se han obtenido resultados. Otros equipos que tambi´ en hemos usado en este trabajo, pero para los que no presentamos resultados, se describen en el ap´ endice A.1. En particular, todos esos equipos se basan tanto en el sistema de E/S Spider como en la red de interconexi´ on SION. En ese ap´ endice se incluye una descripci´ on del proyecto Summit, que estar´ a en marcha en el OLCF en el a˜ no 2018 y cuyo avance m´ as significativo ser´ a en su arquitectura de E/S. 1.4.1. Crow Al comienzo de este trabajo ten´ ıamos la necesidad, a la hora de probar los algoritmos, de realizar las pruebas en un ordenador que dispusiera de un sistema de ficheros distribuido, adem´ as de soportar Chapel, el lenguaje en el que nos hemos basado. Desde la empresa Cray nos dejaron acceso a uno de los ordenadores que usan internamente para realizar pruebas y desarrollar c´ odigo: Crow. Es un Cray XT5, sistema representativo de clusters de tama˜ no peque˜ no y medio en HPC, optimizado para cargas intensivas de c´ alculo. Crow est´ a compuesto por 20 nodos con 2 CPUs AMD Shanghai de 4 cores, a 2.8 GHz y 32 Gigabytes de RAM. Cada nodo tiene 8 cores, sumando en total 160 cores. En el apartado de software encontramos que usan el CLE (Cray Linux Environment), que es un SLES 11 (SuSE Linux Enterprise Server) con modificaciones realizadas por Cray para que se adapte perfectamente a su hardware, un sistema de colas PBS (Portable Batch System) de la empresa Altair, y como sistema de ficheros el Lustre-Cray/1.8.2. El sistema de ficheros en ese caso estaba formado por un MDT (Metadata Target) y tres OSTs (Object Storage Target). Tras contactar con David Knaak, el autor de [71], ´ este nos inform´ o que s´ olo hab´ ıa dos servidores para Lustre, por lo que el MDT y un OST compart´ ıan un servidor y los otros dos OSTs compart´ ıan el otro servidor, lo que nos ayud´ o a entender los resultados que hab´ ıamos obtenido en las pruebas preliminares, en las que observamos que resultaba ventajoso usar un ´ unico agregador por OST, y que cada agregador construyera bloques 30 Cap´ ıtulo 1. Estudio del estado del arte en los sistemas de E/S Locales MDT OST 0 OST 1 OST 2 OSS OSS and MDS CRAY XT5 Figura 1.9: Arquitectura E/S de Crow. de tama˜ no del stripe de Lustre para a partir de ellos realizar las transferencias de E/S con el OST asignado. La arquitectura de E/S de Crow la podemos ver en la figura 1.9. Los Cray modelo XT5 usan una topolog´ ıa de torus 3D para su red de interconexi´ on entre los nodos, implementada con chips Seastar II+, que incluyen un chip PowerPC 440, un motor DMA y varios enlaces hypertransport que interconectan tanto los buses hypertransport internos del nodo como los chips Seastar II+ de los nodos cercanos, ofreciendo una velocidad de 9.6 Gigabytes por segundo en cada enlace. La topolog´ ıa de torus 3D tiene varias ventajas: i) no hace falta tener switches en la red de interconexi´ on, con lo que el consumo de energ´ ıa es m´ as bajo; ii) con un torus se usan exactamente los mismos cables para todas las conexiones, del mismo tipo y la misma longitud, lo que hace que sean sim´ etricas y por tanto m´ as simples de implementar y mantener; iii) en un torus se tienen caminos alternativos si un nodo falla. Se puede ver un esquema de sus interconexiones en la figura 1.10. Adem´ as incluye correcci´ on de errores para interceptar y corregir los errores que se produzcan en las comunicaciones entre nodos, y una implementaci´ on de la librer´ ıa portals en firmware para optimizar los pasos de mensajes entre los nodos, por ejemplo al usar MPI. 1.4.2. Spider I y II Originalmente cada sistema del OLCF dispon´ ıa de un sistema de almacenamiento local, lo que originaba varios problemas, entre ellos la necesidad de copiar los ficheros 1.4. Sistemas usados durante el desarrollo de esta tesis 31 Figura 1.10: Enlaces del chip SeaStar2+, incluyendo su conexi´ on a un nodo. para ser usados en distintos sistemas, y de crear, construir y mantener distintos sistemas de ficheros de alto rendimiento, con el coste que ello implica [110]. Es por ello que se decidi´ o dividir los sistemas de ficheros en tres partes, una para los datos m´ as importantes de los usuarios, que se accede por NFS y no est´ a visible desde los nodos de c´ alculo; otra para realizar los c´ alculos, llamada Spider, que sea de alta velocidad y baja latencia, y que se considera un almacenamiento temporal (scratch); y otra permanente de alta capacidad, con una alta latencia de acceso, un HSM (Hierarchical Storage Management) con almacenamiento terciario. As´ ı, aunque exista un sistema de ficheros convencional (montado por NFS) para los datos m´ as importantes de los usuarios, ´ este no est´ a accesible desde los nodos de c´ alculo, sino s´ olo desde los servidores de login, nodos de transferencia de ficheros, y en general desde todos los ordenadores menos los de c´ alculo. Para los nodos de c´ alculo se usa el sistema de E/S Spider, que provee un sistema ´ unico para las necesidades de almacenamiento de todos los sistemas existentes en el OLCF [110], evitando las copias de datos y la necesidad de disponer de distintos sistemas de ficheros de alta capacidad y velocidad. Otros de los requisitos a la hora de dise˜ nar Spider era el permitir actualizaciones del sistema de ficheros con independencia de las 38 Cap´ ıtulo 1. Estudio del estado del arte en los sistemas de E/S Figura 1.15: Esquema de un blade de un Cray XC30. Las comunicaciones de rango 1 comunican los blades de un chasis, las de rango 2 los blades de un grupo, y las de rango 3 entre distintos grupos. Figura 1.16: Un blade f´ ısico de un XC30, cada disipador cubre un chip multicore, el de la izquierda est´ a sobre el chip Aries de comunicaciones. Al siguiente nivel del XC30 se le llama chasis. Caben 16 blades, o 64 nodos de c´ alculo en un chasis. La red de rango 1 se encarga de comunicar los blades del chasis, comunicando los blades usando una red todos a todos a trav´ es de la placa que los interconecta, por lo que no hacen falta cables para comunicarlos. En la figura 1.17 se muestra un grupo, que es el nivel superior al del chasis, estando formado por 6 chasis que ocupan 2 racks, teniendo en total 96 chips Aries (384 nodos de c´ alculo o 768 sockets). La red que comunica los chasis dentro de un grupo, de rango 2 en la terminolog´ ıa Aries, est´ a formada por cables pasivos de cobre y comunica cada blade con todos los que est´ an en el mismo slot en el resto de chasis. Por ejemplo, el chip Aries 1.4. Sistemas usados durante el desarrollo de esta tesis 39 Figura 1.17: Ejemplo de canales de comunicaci´ on entre dos blades dentro de un grupo en un Cray XC30. Cada caja azul es un blade, cada rect´ angulo es un chasis, y los seis rect´ angulos, la imagen completa, forman un grupo. Las l´ ıneas verdes indican comunicaciones dentro del chasis, y las negras entre elementos del grupo. que est´ a en la posici´ on 0 dentro de un chasis se conecta con los otros cinco chips Aries que est´ an en la misma posici´ on en los otros chasis, el que est´ a en la posici´ on 1 con todos los dem´ as de la posici´ on 1 de los otros cinco chasis del grupo y as´ ı sucesivamente. De esa forma la ruta m´ ınima entre dos nodos de un grupo (que no est´ en en el mismo chasis) va a tener siempre dos saltos, y la ruta m´ axima, cuando la ruta m´ ınima est´ e saturada, es de cuatro saltos, como se puede ver en la figura 1.17. Finalmente tenemos las comunicaciones entre grupos, que se realizan a trav´ es de la red de rango 3. Cada grupo tiene 240 puertos de fibra ´ optica para conectarse con otros grupos [87], por lo que el m´ aximo de grupos es de 241, pero se pueden hacer enlaces (bonding) entre las comunicaciones para tener m´ as ancho de banda. En cualquier caso deber´ an de ser sim´ etricas, es decir, tener las mismas conexiones con todos los grupos, y ser de todos a todos. Una de las ventajas de la conexi´ on Aries es que si detecta que un enlace est´ a saturado puede decidir usar otro, aunque sea m´ as costoso, para aumentar los anchos de banda entre dos nodos dados. Adem´ as el coste y la complejidad del cableado es mucho menor, ya que los tipos de conexiones m´ as abundantes son muy baratos, mientras que el n´ umero de conexiones “caras” es mucho menor. Adem´ as, en general Aries proporciona latencias mucho menores, ya que se suelen necesitar menos saltos para conectar dos nodos 40 Cap´ ıtulo 1. Estudio del estado del arte en los sistemas de E/S cualesquiera. Figura 1.18: Esquema de un sistema XC30, los cuadrados verdes son chasis, cada uno conteniendo 16 blades, los cuadrados azules son 2 racks, conteniendo 6 chasis en total. El sistema de ficheros est´ a conectado a trav´ es de una red infiniband FDR, llamada SION II, a los nodos LNET que enrutan los mensajes de E/S entre los clientes y los servidores de Lustre [30]. Hay en total 9 servidores LNET en EOS, cada uno con 4 conexiones infiniband, lo que da un total de 36 conexiones, una a cada uno de los switches de Spider II. Cada nodo de c´ alculo tiene una lista con un router primario y otro de respaldo para el caso de que falle el primario, que sirven para comunicarse con los OSSs y los MDSs. Otra ventaja de la red Aries respecto a la E/S es que, te´ oricamente, da igual donde se pongan los enrutadores LNET, siempre que est´ en repartidos entre los grupos, ya que el n´ umero de saltos para llegar a ellos es muy reducido [87]. Los enrutadores LNET usan un algoritmo de proyecci´ on donde cada uno tiene uno o m´ as OSS asignados y funcionan como proxies [43]. 1.4.5. Picasso Para las evaluaciones de rendimiento de la librer´ ıa FQbin se ha usado Picasso, un cluster ubicado en la Universidad de M´ alaga, formado por un conjunto heterog´ eneo de equipos de c´ alculo, que se pueden dividir en cuatro grupos de ordenadores. Por un lado 32 HP SL230G8 con 2 CPUs E5-2670, con un total 16 cores y 64 Gigabytes de RAM por nodo (en total 512 cores y 2 Terabytes de RAM). Por otro lado, 7 ordenadores HP DL980G7 con 8 procesadores E7-4870 con 80 cores y 2000 Gigabytes de RAM por nodo (en total 560 cores y 14 Terabytes de RAM). Adem´ as contiene 16 SL250G8 con 2 CPUs E5-2670 y dos tarjetas NVidia Tesla M2075, en total 256 cores, 1 terabyte de 1.4. Sistemas usados durante el desarrollo de esta tesis 41 Figura 1.19: Ordenador Picasso de la UMA, elementos conectados por infiniband. RAM y 32 tarjetas NVidia Tesla M2075. Todos estos equipos est´ an comunicados por una red de interconexi´ on infiniband FDR de 54 Gigabits por segundo con una topolog´ ıa fattree a trav´ es de 5 switches, 3 dispuestos como hojas, al que se conectan los servidores de c´ alculo y el sistema de almacenamiento, y los otros 2 ejerciendo de troncales para interconectar los otros 3. En la figura 1.19 se pueden ver los elementos anteriormente descritos, m´ as el sistema de almacenamiento y los servidores del sistema de ficheros. Adem´ as Picasso dispone de 41 nodos HP DL165G7, cada uno compuesto por 2 procesadores AMD Opteron 6176 (con 12 cores cada uno, 24 cores por nodo) y 96 Gigabytes de RAM. En total, estos equipos contienen 984 cores y 4 Terabytes de RAM conectados por ethernet Gigabit entre ellos y al resto de equipos. El sistema de almacenamiento es un Lustre compuesto por una cabina DDN con 2 controladoras SFA10000 redundantes y 240 discos de 3 Terabytes cada uno repartidos entre cinco cajones de discos, con una capacidad bruta de 720 Terabytes, que, al estar configurados en RAID-6, tienen una capacidad neta de unos 500 Terabytes. Adem´ as hay una cabina DDN EF3015 con 10 discos de 600 Gigabytes a 15000 rpm para los 42 Cap´ ıtulo 1. Estudio del estado del arte en los sistemas de E/S metadatos. Los sistemas de ficheros Lustre son servidos desde 2 MDSs y 4 OSSs, con 24 OST cada uno de ellos con 10 discos (en RAID-6 8+2), y cada OSS sirve 6 OSTs. Hay 2 MDS para servir los metadatos, y un MDT para cada uno de los sistemas de ficheros. 2Estudio de una interfaz de E/S en Chapel En este cap´ ıtulo presentamos la interfaz e implementaci´ on de una librer´ ıa para realizar E/S paralela desde el lenguaje Chapel. Para ello vamos a introducir las distribuciones en Chapel y la librer´ ıa GASNet que se usa para realizar las comunicaciones usando de forma ´ optima las distintas redes disponibles actualmente, entre otras las que ya hemos visto en la secci´ on 1.4. Adem´ as, se presenta una ampliaci´ on de las rutinas de comunicaci´ on de Chapel para mejorar, de manera transparente al programador, el rendimiento al realizar operaciones de movimiento de datos en arquitecturas con memoria distribuida. 2.1. Distribuciones en Chapel Una de las caracter´ ısticas m´ as importantes de Chapel es su soporte de distribuciones de datos. Adem´ as de las distribuciones de datos est´ andar (Block, Cyclic, Block-Cyclic) que est´ an disponibles en las librer´ ıas del compilador, Chapel tambi´ en permite al usuario definir distribuciones adicionales. Esto permite a programadores avanzados la creaci´ on de implementaciones de distribuciones de datos propias. Con estas distribuciones se permite trabajar con arrays de forma natural, a pesar de que estos pueden estar distribuidos en distintos nodos. Adem´ as Chapel permite especificar la forma en la que se distribuir´ a el array entre los nodos, la estrategia de recorrido paralelo, su disposici´ on espec´ ıfica dentro de la memoria de los nodos, as´ ı como otros detalles. Para facilitar la tarea de trabajar con arrays distribuidos, Chapel incorpora el concepto novedoso de mapa de dominio (domain map) [16], que es una receta mediante la que se instruye al compilador 43 44 Cap´ ıtulo 2. Estudio de una interfaz de E/S en Chapel la forma de mapear la vista global de los datos en la memoria del nodo. M´ as concretamente, los mapas de dominio especifican c´ omo los ´ ındices y los elementos del array se mapean en los locales (nodos), c´ omo se almacenan en memoria y c´ omo se implementan operaciones como recorridos, accesos y particionamientos. A los mapas de dominio se les puede poner nombre, ser manipulados y pasados a funciones. En particular, Chapel provee la palabra clave dmapped que permite mapear los ´ ındices del dominio a la arquitectura que designemos como objetivo usando una distribuci´ on espec´ ıfica. Un mapa de dominio se podr´ a llamar layout si solamente se va a realizar sobre un locale, o distribuci´ on si se va a realizar potencialmente sobre m´ ultiples locales. Un layout determina c´ omo se almacena e itera sobre dominios y arrays, mientras una distribuci´ on adem´ as debe decidir c´ omo mapear los ´ ındices del dominio a los locales. Un resultado interesante de esa organizaci´ on es que los arrays declarados sobre los dominios que han sido distribuidos usando la misma instancia de mapa de dominio est´ an alineados, y por tanto la localidad puede ser aprovechada mejor. En otras palabras, una distribuci´ on es una receta que Chapel usa para mapear los datos (y sus c´ alculos asociados) a los locales donde el programa se va a ejecutar [17]. Chapel incluye una librer´ ıa de mapas de dominio est´ andar para soportar las distribuciones m´ as comunes, a saber Block, Cyclic, Block-Cyclic. Adem´ as de las distribuciones que trae Chapel, ofrece soporte para distribuciones definidas por el usuario v´ ıa su DSI (Domain map Standard Interface) [16]. Esas distribuciones definidas por el usuario se desarrollan directamente en Chapel y pueden usar todas sus caracter´ ısticas (clases, iteradores,...), lo que simplifica la tarea del programador. 1use BlockDist; 2config const n=500; 3var Space = {1..n,1..n}; 4var Dom = Space dmapped Block(boundingBox=Space); 5var A: [Dom] real; 6... Figura 2.1: Declaraci´ on en Chapel de una matriz distribuida en bloques. En la figura 2.1 se puede ver una definici´ on de un conjunto de datos distribuidos. La distribuci´ on se hace en tiempo de ejecuci´ on, y depende del n´ umero de locales que se vayan a usar. La distribuci´ on de datos es transparente, por lo que la forma de acceder y procesar los datos ser´ a´ optima cuando se use correctamente. Para realizar la distribuci´ on, en la l´ ınea 1 indicamos que queremos cargar las definiciones de la distribuci´ on por bloques. En la l´ ınea 2 definimos una variable de configuraci´ on, n= 500, lo que significa que la podremos definir en la l´ ınea de comandos al invocar el programa. En la l´ ınea 3 definimos un dominio bidimensional de n×nque usaremos en la l´ ınea 4 para 2.1. Distribuciones en Chapel 45 crear un mapa de dominio. De esta forma, la bounding box definida por Space se reparte en porciones similares entre los locales. Finalmente, en la l´ ınea 5, se define un array, que almacenar´ a datos de tipo real usando el dominio Dom, por lo que estar´ a distribuido sobre los locales siguiendo una distribuci´ on en bloques. Por simplicidad, en la l´ ınea 4 coincide el dominio de la bounding box con los ´ ındices a almacenar, pero podr´ ıan ser distintos. En ese caso, s´ olo algunos ´ ındices de la bounding box contendr´ ıan elementos, como veremos en un ejemplo posterior en la figura 2.2. Para comprender mejor la forma en la que se distribuyen los datos en Chapel, y por tanto lo que son los mapas de dominio, los dominios y los arrays, vamos a ver primero una visi´ on m´ as detallada de los conceptos explicados en los p´ arrafos anteriores, y a continuaci´ on unos ejemplos. 2.1.1. Las distribuciones de datos en detalle Hemos introducido los conceptos de array, dominio y mapa de dominio, que son los conceptos clave para tener datos distribuidos y paralelismo en las operaciones realizadas sobre esos datos. Vamos ahora a ilustrar el entorno de trabajo de Chapel para implementar mapas de dominio, usando la distribuci´ on Block est´ andar como ejemplo. El m´ odulo de Chapel BlockDist define tres clases que indican: i) el mapa de dominio; ii) un dominio, mapeado usando los mapas de dominio; y iii) un array sobre ese dominio. En la figura 2.2 ilustramos un ejemplo de distribuci´ on de una array bidimensional sobre una malla de 2x2 locales implementado mediante 5 l´ ıneas de c´ odigo Chapel. El c´ odigo (arriba a la derecha) declara el mapa de dominio myDist, el dominio myDom y el array M. En tiempo de ejecuci´ on, estos tres componentes se representan usando descriptores globales, que son instancias de los descriptores de las respectivas clases, Block,BlockDom yBlockArr. Adem´ as existen descriptores locales, que son instancias de las clases LocBLock,LocBlockDom yLocBlockArr. Los descriptores locales son almacenados en cada locale y contienen el estado correspondiente al subconjunto del espacio de datos que esta almacenado en ese locale. El objeto de distribuci´ on de bloque myDist almacena el estado global de la distribuci´ on: el campo boundingBox determina c´ omo se realiza la descomposici´ on entre los locales; targetLocDom es el dominio que describe la red de locales que almacenar´ a los datos distribuidos; el array targetLocales[] identifica los locales reales en uso; y el array locDist[] apunta a los descriptores LocBlock locales. Esos objetos se almacenan en cada locale, y definen las fronteras de las distribuciones locales en sus campos myChunk. De forma similar, para el dominio en stride bidimensional distribuido, myDom, el objeto myDom mantiene el conjunto de ´ ındices en el campo whole=[3..10 by 3, 2..10 by 2] y apunta a los cuatro descriptores LocBlockDom almacenados en los 46 Cap´ ıtulo 2. Estudio de una interfaz de E/S en Chapel BlockArr rank = 2; eltType = real(64); stridable = true; dom locArr [0..1, 0..1] 1 BB=[1..10,1..10]; 2 D=[3..10 by 3, 2..10 by 2]; 3 myDist = new dmap(new Block(BB)); 4 const myDom = D dmapped myDist; 5 var M:[myDom] real(64) = 1..; Block rank = 2; boundingBox = [1..10, 1..10]; targetLocDom = [0..1, 0..1]; targetLocales = L0 L1 L2 L3 locDist[0..1, 0..1]; myDist GlobalLocal BlockDom rank = 2; stridable = true; dist whole = [3..10 by 3, 2..10 by 2]; locDoms[0..1, 0..1] myDom M LocBlockArrLocBlockDom L3 myBlock = [6..10 by 3, 6..10 by 2]; L2 myBlock = [6..10 by 3, 2..5 by 2]; L1 myBlock = [3..5 by 3, 6..10 by 2]; L0 myBlock = [3..5 by 3, 2..5 by 2]; LocBlock L3 myChunk = [6..+∞, 6..+∞] L2 myChunk = [6..+∞, -∞..5] L1 myChunk = [-∞..5, 6..+∞] L0 myChunk = [-∞..5, -∞..5] L3 myElems = 8.0 9.0 10.0 13.0 14.0 15.0 L2 myElems = 6.0 7.0 11.0 12.0 L1 myElems = 3.0 4.0 5.0 L0 myElems = 1.0 2.0 12345 678910 11 12 13 14 15 1 1 5 5 6 6 10 10 Figura 2.2: Ejemplo de distribuci´ on de bloques sobre cuatro locales. locales. Cada uno de ellos a su vez, almacena el conjunto de ´ ındices local, myBlock, calculados por la intersecci´ on entre el dominio global myDom.whole con el correspondiente myChunk local. Finalmente, el array Mse declara sobre myDom y se inicializa con la secuencia 1.. (es decir, 1, 2, 3, 4, ...). El descriptor del array global es una instancia BlockArr que identifica el tipo de datos que se almacenar´ an realmente (en el ejemplo son n´ umeros reales de 64 bits) y apunta a su descriptor de dominio global. Los elementos reales del array Mse almacenan localmente en cada uno de los arrays locales myElems. Los arrays myElems est´ an mapeados usando el layout est´ andar, DefaultRectangular, que ser´ a llamado array DR de ahora en adelante. Entre los arrays de Chapel, los arrays DR son los que representan la memoria f´ ısica de forma m´ as directa, es decir, un array DR se mapea en una ´ unica regi´ on contigua de memoria. Tambi´ en podemos ver en la figura 2.2 los cuatro locales, llamados L0, L1, L2 y L3. El n´ umero de locales se puede especificar en el momento de lanzar el programa a ejecuci´ on usando el par´ ametro -nl (number of locales). Por ejemplo, el ejecutable resultante de compilar un programa Chapel se puede lanzar en 4 locales mediante el siguiente comando: programa -nl 4. Podemos ver las cinco sentencias Chapel que est´ an en el programa del usuario en la esquina superior derecha de la figura 2.2, y debajo representamos gr´ aficamente los datos que se almacenan en cada locale concreto. El c´ odigo y la distribuci´ on resultante es la ´ unica parte de la figura 2.2 que es visible al programador. Es decir, los objetos globales de las clases Block,BlockDom yBlockArr as´ ı como las instancias de las clases locales LocBLock,LocBlockDom yLocBlockArr que se muestran a la izquierda en la figura, son transparentes al usuario y s´ olo se describen brevemente aqu´ ı con objeto de ilustrar las estructuras de datos que hay que manejar a la hora de optimizar las comunicaciones y la E/S en Chapel. 2.1. Distribuciones en Chapel 47 1 2 3 4 7 8 9 10 13 14 15 16 19 20 21 22 (0,0) (5,5) 25 26 27 28 31 32 33 34 5 6 11 12 17 18 23 24 29 30 35 36 1 2 3 4 7 8 9 10 5 6 11 12 13 14 15 16 19 20 21 22 17 18 23 24 25 26 27 28 31 32 33 34 29 30 35 36 locale 0: locale 1: locale 2: Figura 2.3: Distribuci´ on en bloques unidimensional de una matriz en tres locales. 1234 5678 910 11 12 13 14 15 16 (0,0) 1 2 5 6 3 4 7 8 910 13 14 11 12 15 16 (3,3) locale 0: locale 2: locale 1: locale 3: Figura 2.4: Distribuci´ on en bloques bidimensional de una matriz en cuatro locales. Por defecto, una distribuci´ on Block de Chapel reserva subregiones de datos en cada uno de los locales (o nodos). En el caso de definir un array de datos de una dimensi´ on, cada locale almacenar´ a un bloque consecutivo del array. Sin embargo, al declarar una matriz (un array bidimensional), el particionamiento de los datos se realizar´ a dependiendo del n´ umero de locales que se usen en tiempo de ejecuci´ on, y que se indica en la l´ ınea de comando con el argumento que se le pasa a nl. El comportamiento por defecto cuando el n´ umero de locales es primo, es el de configurar una malla unidimensional de locales, lo que resulta en una distribuci´ on por bloques de las filas de la matriz, tal y como 54 Cap´ ıtulo 2. Estudio de una interfaz de E/S en Chapel A B copia srcaddr(6,6,6) dstaddr(10,3,5) count[0]=4 count[1]=3 count[2]=2 stride[0] stride[1] stridelevels=2 srcstride[0]=11 srcstride[1]=11*12=132 dststride[0]=14 dststride[1]=14*15=210 14 16 15 12 11 13 Figura 2.7: Ejemplo de copia entre dos arrays tridimensionales usando dos niveles de salto (stridelevels=2), count={4,3,2},srcstride={11,11×12}ydststride ={14,14×15}. srcstride={11,11×12}ydststride={14,14×15}. El array count={4,3,2}indica en su primera componente que los bloques consecutivos tienen 4 elementos, que hay 3 bloques en cada plano y por ´ ultimo que vamos a usar 2 planos. Existen implementaciones concretas de las funciones gasnet_puts_bulk() y gasnet_gets_bulk() para muchas arquitecturas de forma que su funcionamiento sea lo m´ as ´ optimo posible. Por ejemplo, para las redes Gemini de Cray se utiliza una t´ ecnica de env´ ıo de doble-buffer en la que en el nodo local se van empaquetando los bloques no consecutivos en un buffer de env´ ıo al tiempo que se est´ a enviando el buffer que se prepar´ o con anterioridad. De esta forma se reduce el n´ umero de comunicaciones y se aprovecha mejor el HW de RDMA. En las redes Gemini, la implementaci´ on optimizada a´ un se considera beta, por lo que hay que usar variables de entorno para activar las optimizaciones de copia en stride. Espec´ ıficamente hay que definir las variables GASNET_VIS_AMPIPE=1, que activa la optimizaci´ on de doble-buffer explicada anteriormente, y GASNET_VIS_REMOTECONTIG=1, que habilita el empaquetamiento para datos que est´ an dispersos localmente, pero que en el nodo remoto ser´ an almacenados de forma continua. 2.2. Agregaci´ on de comunicaciones 55 2.2.2. Ampliaci´ on del lenguaje Chapel para usar funciones avanzadas de GASNet Las funciones de copia en stride se usar´ an intensivamente en la librer´ ıa de E/S que proponemos para Chapel. El objetivo es realizar la agregaci´ on de los datos en los distintos locales y reducir as´ ı el n´ umero de mensajes que se transfieren por la red. A cambio, el tama˜ no de los mensajes ser´ a mayor pero debido a que el inverso de la latencia es siempre menor que el ancho de banda, las comunicaciones agregadas son ventajosas cuando el volumen de datos es suficiente. Para usar las funciones de comunicaci´ on agregada en Chapel es necesario implementar en el run-time de Chapel las rutinas que hacen de interfaz entre los arrays de Chapel y los arrays de C que se usan en GASNet3. Las funciones que proponemos son chpl_comm_gets ychpl_comm_puts [105], que han sido dise˜ nadas para tener el mismo interfaz (los mismos argumentos) que las funciones de GASNet. La funcionalidad principal de estas funciones es convertir los arrays de Chapel a arrays de C y posteriormente llamar a la funci´ on GASNet correspondiente. stridelevels = 2, count = (3,2,4), srcstride = (12,16), dststride=(8,32) B[1..4, 1..4 by 3, 2..4] (dim3, dim2, dim1) 1 2 3 4 1 23 4 4 3 2 1 dim3 dim1 dim2 1 2 3 4 dim1 4 3 2 1 dim2 A[1..8 by 2, 1..3 by 2, 2..4] (dim3, dim2, dim1) 1234 dim3 5678 Copy = Figura 2.8: Ejemplo m´ as detallado de transferencia entre dos arrays tridimensionales con dos niveles de salto (stridelevels=2), count={3,2,4},srcstride={12,16}y dststride={8,32}. En la figura 2.8 podemos ver otro ejemplo m´ as detallado de una transferencia agregada en la que ahora se usar´ an las funciones chpl_comm_gets ochpl_comm_puts (dependiendo de si el comando se ejecuta en el destino o en el origen, respectivamente). Los argumentos a pasar a la funci´ on son los mismos que en las funciones equivalentes de 3La librer´ ıa GASNet est´ a implementada en C 56 Cap´ ıtulo 2. Estudio de una interfaz de E/S en Chapel GASNet y que fueron explicados en la secci´ on anterior. En este ejemplo se realiza una transferencia agregada desde el array B, definido como B[4][4][4], los datos del subdominio B[1..4, 1..4 by 3, 2..4]. El array destino es A[8][4][4], suponemos que est´ a en otro locale, y se escribir´ an s´ olo las posiciones A[1..8 by 2, 1..3 by 2,2..4]. En la figura, los datos que hay que leer de By copiar en Aest´ an identificados como cubos con fondo gris oscuro. La versi´ on en memoria de esos arrays contiene los elementos a transferir en 8 bloques de 3 elementos consecutivos, pero separados por strides diferentes en los arrays ByA. Con todo eso, los argumentos de llamada a las funciones chpl_comm_gets ochpl_comm_puts deben ser: stridelevels=2), count ={3,2,4},srcstride={12,16}ydststride={8,32}. En este ejemplo, dstaddr debe apuntar al elemento A[1,1,2] ysrcaddr aB[1,1,2]. El c´ odigo en la figura 2.9 muestra la primitiva chpl_comm_gets que ahora usa el compilador de Chapel. La primitiva chpl_comm_puts es an´ aloga. Como describiremos en la siguiente secci´ on, hemos modificado el compilador de Chapel para que inserte estas primitivas de forma autom´ atica cada vez que hay una asignaci´ on entre arrays distribuidos. Tambi´ en se deben insertar llamadas a funciones que calculen los argumentos que hay que pasar a dichas funciones. Las llamadas a array_get que vemos en los argumentos de la primitiva devuelven un manejador o puntero al array Chapel correspondiente. __primitive("chpl_comm_gets", __primitive("array_get",destr, buffer._value.getUnshiftedDataIndex(0)), __primitive("array_get",dststr,dststrides._value.getDataIndex(0)), srcnode, __primitive("array_get",src, privArr.locArr[locFrom.id].myElems._value.getUnshiftedDataIndex(0)), __primitive("array_get",srcstr,srcstrides._value.getDataIndex(0)), __primitive("array_get",cnt,count._value.getDataIndex(0)), stridelevels); Figura 2.9: Rutina del run-time de Chapel que genera comunicaciones agregadas. 2.2.3. Implementaci´ on de la agregaci´ on de datos El siguiente paso consiste en reducir en n´ umero de transferencias de datos para asignaciones de arrays siempre que estos usen la distribuci´ on Block o Cyclic. Aprovechando las caracter´ ısticas OOP (Object Oriented Programming) de Chapel, se sobrecarg´ o el operador =con una funci´ on especializada que proporciona la implementaci´ on optimizada. Como los arrays distribuidos por bloques se implementan usando arrays DR (Default 2.2. Agregaci´ on de comunicaciones 57 Rectangular) que son arrays que est´ an almacenados de forma local en un s´ olo nodo, hemos implementado primero la comunicaci´ on masiva (bulk) para realizar asignaciones entre dos arrays DR. Sobre eso hemos implementado la asignaci´ on entre un array DR y uno BD (Block Distributed), y a continuaci´ on la asignaci´ on entre dos arrays BD. Ilustramos esa organizaci´ on usando cuatro locales en la figura 2.10. A[D1] Array A B[D2] Array B 0 2 3 1 0 2 3 1 Array A Array B B[D2] 0 2 3 10 2 3 1 A[D1] Array A Array B A[D1] B[D2] 0 2 3 1 0 2 3 1 a) DR=DR b) DR=BD c) BD=BD Figura 2.10: A[D1]=B[D2] cuando A y/o B son arrays DR (Default Rectangular) o BD (Block Distributed). La figura 2.10 a) muestra la asignaci´ on A[D1]=B[D2] donde AyBson arrays DR, Areside en el locale 0 y Ben el locale 1. Aqu´ ı el programa, en tiempo de ejecuci´ on, tiene que calcular los argumentos y ejecutar una sola llamada a chpl_comm_gets (si la sentencia se ejecuta en el locale 0), o chpl_comm_puts (si la sentencia se ejecuta en el locale 1), tal y c´ omo se describe en la secci´ on 2.2.2. Si Bes un array del tipo BD (distribuido por bloques) sus elementos estar´ an almacenados en varios arrays de tipo DR. Estos arrays son los campos myElems de los descriptores locales LocBlockArr, que est´ an contenidos en B. Si Aes un array DR, entonces la operaci´ on A[D1]=B[D2] se compone de varias operaciones del tipo DR=DR. Por ejemplo, en la figura 2.10 b), B[D2] est´ a repartido entre cuatro locales, mientras Aest´ a almacenada completamente en el locale 0. En ese caso se divide A[D1] en cuatro regiones, de A[D10]aA[D13], de forma que A[D1i] es el destino de los elementos de B[D2] que est´ an almacenados en el locale i, con i∈ 0..3. Para cada i,src representa el trozo de B.locArr[i].myElems que corresponde aB[D2] en el locale i. La funci´ on bulk gets oputs que realiza la asignaci´ on en el locale 0 est´ a implementada internamente con un memcpy. Finalmente, si AyBson ambos arrays que siguen una distribuci´ on por bloques, entonces dividiremos el problema en varios del tipo DR=BD, ya que el array de destino BD es simplemente un conjunto de arrays DR. En la figura 2.10 c) el array Aest´ a almacenado en cuatro arrays DR locales. Por tanto la asignaci´ on A[D1]=B[D2] se convierte en: 1forall iin 0..3 do 2on A.locArr[i] do //en el local i ejecuta lo siguiente: 3A.locArr[i].myElems[dest] = B[src] //DR=BD 58 Cap´ ıtulo 2. Estudio de una interfaz de E/S en Chapel donde src ydest son las subregiones del dominio fuente y destino que indican que porciones del array van a ser asignados en cada locale. En este ejemplo A.locArr[0]. myElems necesita cuatro regiones de B, que est´ an almacenadas en cuatro locales. A. locArr[1].myElems s´ olo necesita informaci´ on de los locales 1 y 3. La asignaci´ on a A.locArr[3].myElems es realizada c´ omo una operaci´ on memcpy en este ejemplo. Optimizamos los movimientos de datos al realizar asignaciones entre distribuciones bloque y c´ ıclica, y entre c´ ıclica y bloque de una forma similar, convirti´ endolas en llamadas a chpl_comm_gets/puts. Todas esas redistribuciones de datos son casos especiales de la copia de una regi´ on rectangular de n-dimensiones a otra con una forma distinta. Por ejemplo, para realizar una redistribuci´ on de los datos de un bloque a una distribuci´ on c´ ıclica, cada locale fuente determinar´ a el trozo de datos que enviar´ a a cada locale de destino. Para cada trozo y para cada locale de destino, el locale fuente llamar´ a chpl_comm_puts con los argumentos apropiados. En general, esto conduce a una comunicaci´ on de todos a todos. El caso de la asignaci´ on entre c´ ıclica y bloques es similar, con la diferencia de que la asignaci´ on se pedir´ a desde los locales de destino llamando a la funci´ on chpl_comm_gets. 2.2.4. C´ alculo de las subregiones Como hemos visto en la secci´ on anterior, durante la asignaci´ on de las subregiones del array, hay que calcular distintas regiones dentro de cada trozo, como las subregiones D1iysrc. Para hacerlo en el caso de la asignaci´ on A[DA]=B[DB], primero necesitamos una correspondencia desde el dominio del destino, DA, al dominio de la fuente, DB. Esto se puede describir de manera formal de la siguiente forma. Sea: DA=[la1..ha1by sa1, ...,lan..hanby san] DB=[lb1..hb1by sb1, ...,lbn..hbnby sbn] donde nes el n´ umero de dimensiones (rank) de los dominios DA yDB. Un vector con todos los strides,sai, se puede obtener por A._value.dom.whole.stride y de forma similar para B. Sean [a1,a2,...,an]y[b1, b2,...,bn]las coordenadas de los dominios Aand Brespectivamente. Definimos la funci´ on biyectiva m:DB→DA de forma que (a1,a2,...,an)= m(b1,b2,...,bn)donde: ai=lai+sai×(bi−lbi)/sbi, i ∈1..n La funci´ on inversa correspondiente m−1:DA→DB, resulta en: (b1, b2, ..., bn) = m−1(a1, a2, ..., an) 2.2. Agregaci´ on de comunicaciones 59 1config const n=500; 2var D1 = [1..n,1..n]; 3var D2 = [1..2*n,1..2*n]; 5var A: [D1] dmapped Block(D1) real; 6var B: [D2] dmapped Block(D2) real; 8var DA = [101..200 by 2, 51..200 by 3]; 9var DB = [201..700 by 10, 301..600 by 6]; 11 A[DA] = B[DB]; Figura 2.11: C´ odigo Chapel para ilustrar las funciones de correspondencia. Array B 0 1 23 B[DB] Array A 0 1 2 3 A[DA] 51 150 198 101 159 199 1 1 500 500 1 1000 10001 301 595505499 691 501 491 201 153 161 3 250 251 250 251 500 501 500 501 Figura 2.12: Ilustrando la asignaci´ on A[DA]=B[DB] de la figura 2.11 cuando se usan 4 locales. donde: bi=lbi+sbi×(ai−lai)/sai, i ∈1..n Para ilustrar c´ omo se usan esas funciones de correspondencia, mym−1, vamos a basarnos en una versi´ on simplificada bidimensional del ejemplo que venimos estudiando, que podemos ver en la figura 2.11. Si ese c´ odigo se ejecuta en cuatro locales, la figura 2.12 ilustra la distribuci´ on y las subregiones que hay que mover, donde podemos ver que A[DA] est´ a en el locale 0 mientras B[DB] est´ a distribuida entre los cuatro locales. Cuando se ejecuta la asignaci´ on A[DA]=B[DB], se crean dos nuevos alias de los arrays originales, A’=A[DA] yB’=B[DB]. Ahora, para el caso de DR=BD, como ya hemos dicho, tenemos que ejecutar el siguiente bucle: 60 Cap´ ıtulo 2. Estudio de una interfaz de E/S en Chapel 1forall iin 0..3 do 2on A’ do // se ejecuta en el \textit{locale} que contiene A’ 3A’[Di] = B’.locArr[i].myElems //DR=DR De forma que el problema a resolver ahora es el c´ alculo de las subregiones Di,i∈ {0..3}. Esto se realiza usando la funci´ on de correspondencia m. Por ejemplo, para i=0, el subdominio de B almacenado en el locale 0 viene dado por B’.dom.locDoms[0].myBlock = [201..491 by 10, 301..499 by 6]. Entonces podemos obtener D0usando la funci´ on m para calcular las coordenadas de finalizaci´ on de la regi´ on de destino: [lxa1..hxa1by 2, lxa2..hxa2by 3]: lxa1= 101 + 2 ×(201 −201)/10 = 101 hxa1= 101 + 2 ×(491 −201)/10 = 159 lxa2= 51 + 3 ×(301 −301)/6 = 51 hxa2= 51 + 3 ×(499 −301)/6 = 150 De forma que obtenemos D0=[101..159 by 2, 51..150 by 3]. Y similarmente: D1= [101..159 by 2, 153..198 by 3] = M(201..491 by 10, 505..595 by 6) D2= [161..199 by 2, 51..150 by 3] = M(501..691 by 10, 301..499 by 6) D3= [161..199 by 2, 153..198 by 3] = M(501..691 by 10, 505..595 by 6) donde Mconvierte los l´ ımites alto y bajo de su rango de argumentos usando la funci´ on m. Como ya hemos dicho, si tanto Acomo Bson arrays distribuidos en bloques, podemos dividir el problema en varios del tipo DR=BD como hemos mostrado en la figura 2.10, pero ahora necesitamos la funci´ on m−1. As´ ı, si la funci´ on de asignaci´ on es A[DA]=B [DB], se crean dos nuevos alias de la misma forma que anteriormente: A’=A[DA] y B’=B[DB], y hay que ejecutar el siguiente c´ odigo: 1forall iin 0..3 do 2on A’.locArr[i] do // se ejecuta en el \textit{locale} i 3A’.locArr[i].myElems = B’[Di]//DR=BD 2.2. Agregaci´ on de comunicaciones 61 de forma que tenemos que encontrar las subregiones Di,i∈ {0..3}. Pero ahora tenemos que usar la funci´ on m−1para poder calcular los subdominios Dide B que corresponden a cada dominio local A’.dom.locDoms[i].myBlock. Para el ejemplo anterior, pero con DA=[101..400 by 2, 51..350 by 3] yDB=[201..500 by 2, 151..450 by 3] y cuatro locales, tenemos que A’.dom.locDoms[0].myBlock= [101..250 by 2, 51..250 by 3] yD0= [201..350 by 2, 150..350 by 3]. De forma similar para Diy para el resto de casos. Para profundizar en otros detalles de implementaci´ on el lector puede acceder al c´ odigo fuente, disponible en: http://chapel.cray.com. Esta implementaci´ on de la agregaci´ on es ´ util en multitud de casos. Tiene la ventaja a˜ nadida de no requerir de modificaciones en el c´ odigo fuente para usarla. Por ejemplo, la asignaci´ on A=B usar´ a la agregaci´ on si es necesario. Mostramos a continuaci´ on un ejemplo de uso de la agregaci´ on de comunicaciones en el algoritmo PARACR. Adem´ as, la usamos tambi´ en en la librer´ ıa de E/S paralela para el lenguaje Chapel que describiremos en la secci´ on 2.3. 2.2.5. Evaluaci´ on del rendimiento usando el algoritmo PARACR 0 200 400 600 800 1000 24816 32 Tiempo.en.segundos Num..Locales Titan:'PARACR(BlockToCyclic)'con'n=2^28 Sin.Agregación Con.Agregación Figura 2.13: Tiempo en segundos del algoritmo PARACR en Titan, sin/con agregaci´ on de datos. La Reducci´ on C´ ıclica Paralela (Parallel Cyclic Reduction), llamada PARACR en [61], es un algoritmo bien conocido para la resoluci´ on de sistemas tridiagonales de ecuaciones. Funciona en dos fases: sustituci´ on y soluci´ on. Para un problema de tama˜ no n, la fase de sustituci´ on est´ a compuesta de O(log n)pasos y en cada paso se realizan O(n) operaciones de tipo butterfly. Este algoritmo requiere de una distribuci´ on bloques en las primeras etapas de la computaci´ on, pero para explotar mejor la localidad es necesario 62 Cap´ ıtulo 2. Estudio de una interfaz de E/S en Chapel cambiar a la distribuci´ on c´ ıclica para las ´ ultimas etapas del algoritmo. Este cambio de distribuci´ on “Block to Cyclic” es inmediato en Chapel ya que basta con realizar la asignaci´ on A=B con Adistribuido c´ ıclicamente y Bdistribuido por bloques. Sin embargo, en la implementaci´ on original de Chapel con comunicaciones no agregadas el tiempo de redistribuci´ on consum´ ıa una gran parte del tiempo total de ejecuci´ on. En la figura 2.13 mostramos el tiempo de ejecuci´ on en segundos de la redistribuci´ on de bloques a c´ ıclica en varios nodos del supercomputador Titan (descrito en el ap´ endice A.1.2). La red de comunicaciones usada es Gemini y el n´ umero de locales va desde 2 hasta 32. En este c´ odigo hay cuatro arrays, cada uno de ellos de tama˜ no n= 228, que tienen que ser redistribuidos. La agregaci´ on aporta un aumento del rendimiento significativo, sobre todo con un bajo n´ umero de locales. Cuando el n´ umero de locales aumenta el tiempo de redistribuci´ on crece muy ligeramente, debido a que el n´ umero de mensajes bulk por locale es O(num.locales). Sin embargo, en el caso de no usar agregaci´ on el tiempo de redistribuci´ on baja al aumentar el n´ umero de locales. Esto es as´ ı porque el n´ umero total de mensajes punto-a-punto cuando no hay agregaci´ on es igual an(un mensaje por cada elemento del array), pero el n´ umero de mensajes por locale es O(n/#locales), por lo que al aumentar el n´ umero de locales disminuye el n´ umero de mensajes por locale. Con 2 locales, usar agregaci´ on proporciona una aceleraci´ on de 264x y con 32 locales la aceleraci´ on sigue siendo de m´ as de un orden de magnitud (16x). En [105] se muestran resultados similares para un algoritmo FFT radix-4. Las ventajas de la agregaci´ on se han evaluado adicionalmente con otros algoritmos como la transposici´ on de matrices, y algoritmos tipo “stencil” en los que hay que mover bloques de datos con frecuencia. En todos los casos, la agregaci´ on de mensajes ha resultado ser una optimizaci´ on de gran valor para reducir el tiempo total de ejecuci´ on. 2.3. Librer´ ıa de E/S en Chapel 63 2.3. Librer´ ıa de E/S en Chapel Como ya hemos visto anteriormente, uno de los retos de la HPC (High Performance Computing) es que el cuello de botella de las aplicaciones est´ a pasando de ser el procesamiento de los datos a ser la disponibilidad de los datos. Muchas aplicaciones HPC acceden de forma intensiva a datos, bien ley´ endolos y/o escribi´ endolos. Esto hace necesaria la optimizaci´ on de la E/S para obtener mejores rendimientos. Por razones de portabilidad y de representaci´ on de los datos, los sistemas HPC se despliegan a menudo con una amplia variedad de pilas software para manejar la E/S. En la capa de arriba, las aplicaciones cient´ ıficas realizan la E/S a trav´ es de librer´ ıas middleware como Parallel NetCDF [67], HDF [89] o MPI IO [120]. En cualquier caso, todas las librer´ ıas anteriores siguen dejando visibles al programador detalles de bajo nivel, que no son apropiados para una computaci´ on de alta productividad. Por otro lado, en el fondo de la pila de software usada para acceder a los datos, los sistemas de ficheros paralelos sirven directamente las peticiones de E/S dividiendo los bloques de los ficheros entre m´ ultiples dispositivos de almacenamiento. Obtener un buen rendimiento colectivo de E/S cuando se tienen muchos procesos sobre esas capas de software es una tarea compleja, que requiere no s´ olo conocer los patrones de acceso a los datos por parte de los procesos, sino adem´ as comprender c´ omo funciona la pila de software completa, especialmente el comportamiento del sistema de ficheros subyacente. Creemos que esta es una tarea muy pesada de la que es mejor abstraer al programador. En otras palabras, un lenguaje paralelo deber´ ıa proveer las construcciones de alto nivel adecuadas para realizar las operaciones de E/S, mientras el c´ odigo en tiempo de ejecuci´ on (runtime) ser´ ıa el encargado de optimizar el rendimiento y de tratar con el soporte de middleware para el correspondiente sistema de ficheros. Para ello, el runtime se puede apoyar en pistas (hints) sobre la localidad dadas por los usuarios (por ejemplo, la distribuci´ on de datos en los nodos). En esta tesis se ha realizado la primera implementaci´ on de operaciones E/S paralelas en el contexto del lenguaje de alta productividad Chapel. Actualmente Chapel soporta operaciones b´ asicas de E/S. Nuestras nuevas funciones de E/S paralelas se han implementado usando caracter´ ısticas est´ andar de Chapel, tomando las funciones b´ asicas de Chapel como capa intermedia y Lustre como el sistema de ficheros objetivo. Conforme evolucion´ o el trabajo el interfaz que presentamos se hizo m´ as simple, pero igual de potente, con lo que la productividad tambi´ en mejor´ o. Adem´ as se estudiaron distintos algoritmos para su implementaci´ on, las cuales se compararon con un c´ odigo equivalente en MPI IO, usando la librer´ ıa ROMIO, resultando el c´ odigo de Chapel no s´ olo mucho m´ as simple, si no adem´ as un poco m´ as r´ apido. A la hora de construir una interfaz de E/S hay que separar el dise˜ no de la interfaz de 70 Cap´ ıtulo 2. Estudio de una interfaz de E/S en Chapel cada locale escribe su porci´ on de los datos. Esto no significa que no estemos explotando el paralelismo en la escritura, si no que s´ olo aprovechamos el paralelismo que provee el sistema de ficheros, que viene dado por el n´ umero de OSTs en los que est´ an repartidos los datos que escribe cada locale. Si queremos explotar m´ as paralelismo, otra posibilidad es implementar una escritura colectiva en dos fases (o collective buffering en ingl´ es). La idea es recurrir a una soluci´ on intermedia entre que todos los locales escriban al mismo tiempo o que s´ olo uno de ellos pueda escribir en un momento dado. Para ello se selecciona un subconjunto de locales que ser´ an los responsables de recolectar los datos y transferirlos al sistema de ficheros. A estos locales con capacidad de realizar operaciones de E/S se les llama agregadores y todos ellos pueden realizar sus funciones de acceso al sistema de ficheros en paralelo. El n´ umero de agregadores es configurable y t´ ıpicamente inferior al n´ umero total de locales, lo que permite jugar con el compromiso entre el grado de paralelismo explotado y la contenci´ on en el acceso al sistema de ficheros. Es importante organizar la pol´ ıtica de agregaci´ on de datos que llevan a cabo los agregadores de forma que dos agregadores nunca tengan que acceder a un mismo stripe al mismo tiempo. Esta aproximaci´ on se estudia en [71], donde se usan 32 locales, 16 OSTs, 16 agregadores, y un fichero de 96 GiB en un Cray XT5. En estas condiciones, si los 32 locales escriben sus propios elementos s´ olo consiguen un ancho de banda de 420 MiB/seg, que se puede optimizar a 1629MiB/seg cuando s´ olo escriben los 16 agregadores y se recurre a un algoritmo con agregaci´ on de los datos. En otra prueba usando 1000 locales la velocidad pasa de 13 MiB/seg (escritura de todos los locales en paralelo) a 256 MiB/seg usando s´ olo los 16 agregadores de datos. Claramente, recurrir a agregadores conduce a una reducci´ on de los bloqueos y de la contenci´ on en el acceso al sistema de ficheros distribuido. La fase en la que los agregadores recolectan los datos de los dem´ as locales implica el realizar copias de subconjuntos de los datos entre distintos nodos. En esta fase aprovecharemos las funciones de agregaci´ on de datos presentadas en el cap´ ıtulo 2.2.2, de forma que las transferencias de datos discontinuos entre locales se implementen con una ´ unica operaci´ on. El run-time de Chapel y la librer´ ıa GASNet se encargar´ an autom´ aticamente del empaquetamiento de las transferencias para minimizar las comunicaciones que se realizan. Por tanto aprovecharemos mejor las caracter´ ısticas de la arquitectura de red subyacente. 2.3.3. Posibles fuentes de paralelismo Vamos a sistematizar la manera de abordar el problema de la escritura de un array distribuido en el sistema de almacenamiento secundario. Pretendemos realizar un estudio met´ odico de las distintas fuentes de paralelismo que se pueden explotar. Estas fuentes 2.3. Librer´ ıa de E/S en Chapel 71 OST 0 OST 1 OST 2 Locale 0 Locale 1 Locale 2 Locale 3 Locale 4 Locale 5 Locale 6 Locale 7 1 2 3 3 3 1 1 2 2 2 2 22 2 Locale 8 2 Figura 2.18: Arquitectura global de acceso al sistema de almacenamiento usando 3 agregadores. Se pueden ver tres puntos donde aplicar el paralelismo en el acceso. de paralelismo son ortogonales y se pueden explotar de forma combinada, as´ ı que tendremos que estudiar el rendimiento de cada combinaci´ on y las interrelaciones que se producen cuando explotamos varias fuentes de paralelismo al mismo tiempo. Fuentes de paralelismo en operaciones de escritura 1 Paralelismo en la E/S Varios agregadores escriben en paralelo, cada uno en un OST diferente. Paraleliza las comunicaciones con los OSTs. 2 Paralelismo en la agregaci´ on Varios locale escriben en paralelo en los buffers de los agregadores. Existe un agregador por cada OST. Paraleliza la operaci´ on de agregaci´ on. 3 Paralelismo entre E/S y agregaci´ on Explota doble buffering en los agregadores. Cada agregador simultan´ ea las operaciones de agregaci´ on y E/S. 4 Paralelismo en la escritura en OST Varios agregadores escriben en paralelo en un OST. El n´ umero de agregadores es un m´ ultiplo del n´ umero de OSTs. Tabla 2.1: Resumen de las distintas fuentes de paralelismo que explotaremos de forma combinada. En la tabla 2.1 enumeramos las 4 posibles fuentes de paralelismo que se pueden explotar de forma independiente o combinada. Para describir cada una de ellas nos apoyaremos en las figuras 2.18 y 2.19. En estas figuras esquematizamos la arquitectura de E/S para HPC que queremos explotar. Se muestran 9 locales que contienen los datos a almacenar en fichero y 3 OSTs encargados de almacenar los stripes del fichero. En la figura 2.18 se representan 3 agregadores que coinciden con los locales 0, 1 y 2, mientras 72 Cap´ ıtulo 2. Estudio de una interfaz de E/S en Chapel OST 0 OST 1 OST 2 Locale 0 Locale 1 Locale 2 Locale 3 Locale 4 Locale 5 Locale 6 Locale 7 1 2 3 3 3 3 3 3 3 3 1 1 2 2 2 2 222 4 4 4 Locale 8 3 2 4 4 4 Figura 2.19: Fuente adicional de paralelismo, indicada con un 4, que explota la escritura simult´ anea de varios agregadores en el mismo OST. que en la figura 2.19 los 9 locales tienen el papel de agregador. Se dibujan con flechas los flujos de datos que intervienen en el sistema, que son dos: i) por un lado hay que mover los datos desde los locales hacia los agregadores; y ii) posteriormente el contenido de los buffers de los agregadores se env´ ıa hacia los servidores del sistema de ficheros. En la figura 2.18 nos centramos en las tres primeras fuentes de paralelismo, indicadas por los n´ umeros 1, 2 y 3. El n´ umero 1 representa el paralelismo en las comunicaciones entre los agregadores y los OSTs del sistema de ficheros, es decir, los tres agregadores escriben al mismo tiempo, cada uno en un OST diferente. En esta configuraci´ on s´ olo hay un agregador por cada OST, que es responsable de agregar los stripes que se deben almacenar en ese OST. El n´ umero 2 identifica el paralelismo en los movimientos de datos desde los locales hacia los agregadores. En la figura se muestra como todos los locales agregan en el locale 0 todos los stripes que deben ser almacenados en el OST 0. Para evitar una excesiva complicaci´ on en la figura 2.18 se han obviado los flujos de datos que van desde todos los locales hasta los locales 1 y 2. Es decir, al tiempo que el locale 0 agrega todos los stripes que tienen como destino el OST 0, los locales 1 y 2 hacen lo propio con los stripes que tienen como destino los OSTs 1 y 2, respectivamente. El n´ umero 3 especifica el paralelismo que podemos explotar mediante una t´ ecnica de doble buffering que permita simultanear las operaciones de agregaci´ on y las de E/S. 2.3. Librer´ ıa de E/S en Chapel 73 Para ello se procede en fases en las que dos buffers se usan de forma alternativa: i) uno de ellos se escribe con los datos que los dem´ as locales van agregando en ´ el; ii) del otro buffer se leen los datos que se agregaron en una fase anterior, de forma que se puedan escribir a disco. Cuando ocurra que el buffer de lectura est´ a vac´ ıo y el de escritura ya est´ a lleno, se cambian las tornas y se vuelve a proceder como en la fase anterior. En la figura 2.19 se ilustra una cuarta fuente de paralelismo, representada con el n´ umero 4. En esta figura, tenemos 9 agregadores y vemos c´ omo varios de ellos pueden escribir en paralelo en un mismo OST. Por ejemplo, el OST 0 recibe datos al mismo tiempo de los locales 0, 3 y 6. El hecho de tener m´ as agregadores que OSTs permite reducir el tr´ afico de datos entre los locales ya que cada agregador es responsable de menos datos. A cambio, los OSTs est´ an sometidos a una mayor presi´ on en el proceso de escritura. De nuevo en la figura 2.19, bajo el n´ umero 2 s´ olo se muestran los flujos de datos del agregador que agrega los stripes que van al OST 0. Los flujos de los agregadores de los datos que pertenecen a los servidores OSTs 1 y 2 ser´ ıan equivalentes, pero yendo hacia los respectivos agregadores. Hacemos notar que los gr´ aficos anteriores simplifican la realidad en gran medida. En particular, los grandes sistemas HPC, como Titan, descrito en la secci´ on A.1.2 y EOS, descrito en la secci´ on 1.4.4, presentan una arquitectura de E/S mucho m´ as compleja. Por ejemplo, entre los locales y servidores del sistema de ficheros (OSS), existen distintas redes y por tanto es necesario usar routers (llamados LNET, explicados en la secci´ on 1.2.4.2) entre las distintas redes. Las cuatro fuentes de paralelismo descritas tienen caracter´ ısticas diferentes e independientes, y hasta cierto punto podemos decir que son ortogonales. Por ejemplo, podemos explotar al mismo tiempo las fuentes de paralelismo 1 y 2 (ir agregando en paralelo y al mismo tiempo escribir los datos disponibles en paralelo en los OSTs) o 1, 2 y 3 (simultaneando la agregaci´ on paralela con la escritura paralela). En realidad, la fuente de paralelismo que hemos identificado como 4 incluye de forma impl´ ıcita al paralelismo en la E/S (fuente de paralelismo n´ umero 1). La diferencia es que explotando la fuente 4 el n´ umero de agregadores es mayor y nos permite jugar con el compromiso entre la presi´ on en la red de comunicaciones entre locales y la presi´ on en la red de comunicaciones hacia la E/S. Cuando el n´ umero de agregadores es peque˜ no, como en la fuente 1, la presi´ on es mayor en la red de comunicaciones. Por el contrario, cuando tenemos m´ as agregadores, como en la fuente 4, ejercemos m´ as presi´ on sobre la red de E/S ya que hay m´ as agregadores escribiendo al mismo tiempo. Antes de pasar a evaluar las prestaciones que obtenemos de las distintas combinaciones, presentamos a continuaci´ on algunos detalles de implementaci´ on del c´ odigo Chapel que permite activar y desactivar independientemente cada una de las fuentes de paralelismo. 74 Cap´ ıtulo 2. Estudio de una interfaz de E/S en Chapel 2.3.4. Implementaci´ on en Chapel 1// Si Paralelismo4=false --> n´ umero de agregadores = n´ umero de OSTs 2if ( ! Paralelismo4 ) numAgreg=numOSTs; 3else numAgreg=(numLocales/numOSTs)*numOSTs; 5// Si Paralelismo1=true --> BSA = forall, y si no BSA = for 6BSA agregador in 0..#numAgreg do 7on Locales(agregador) do { 9if ( ! Paralelismo4 ) { 10 localesToWriteFrom=0; 11 localesToWriteTo=numLocales; 12 }else { 13 localesToWriteFrom=(agregador/numOSTs)*numOSTs; 14 localesToWriteTo=localesToWriteFrom+numOSTs-1; 15 } 16 // Si Paralelismo2=true --> BSL = forall, y si no BSL = for 17 BSL loc in localesToWriteFrom..localesToWriteTo do { 19 if (Paralelismo3) 20 cobegin { 21 GatherToBuffer(bufferA, bufferB, loc,...); 22 WriteBufferToDisk(bufferA, bufferB, file,...); 23 } 24 else {// Paralelismo3=false 25 GatherToBuffer(bufferA, null, loc,...); 26 WriteBufferToDisk(bufferA, null, file,...); 27 } 28 }//BSL - Bucle sobre los locales 29 }//BSA - Bucle sobre los agregadores Figura 2.20: Pseudoc´ odigo Chapel que implementa el c´ odigo que ejecutan los agregadores. En la figura 2.20 mostramos el pseudoc´ odigo del c´ odigo Chapel que ejecutan los agregadores. Ese algoritmo es configurable de forma que se puedan activar las distintas fuentes de paralelismo descritas anteriormente. En el c´ odigo esas fuentes de paralelismo est´ an identificadas mediante las variables de configuraci´ on4Paralelismo1 aParalelismo4. Las variables numOSTs ynumLocales contienen el n´ umero de OSTs y de locales con los que se est´ a ejecutando el programa. BSA yBSL son cons4En Chapel, las variables de configuraci´ on se pueden inicializar en tiempo de llamada al programa pasando los valores deseados como argumentos. Por ejemplo, el comando ./programa --Paralelismo1 =true lanza la ejecuci´ on del programa con la variable de configuraci´ on Paralelismo1 igual a true. 2.3. Librer´ ıa de E/S en Chapel 75 Locale 0 (aggregator) OST 0 OST 1 OST 2 Locale 1 (aggregator) Locale 2 (aggregator) Locale 3 1 1 1 Figura 2.21: Paralelismo 1: Paralelismo en la E/S cuando todos los agregadores escriben a la vez en los OSTs. tantes que el preprocesador substituye por forall ofor dependiendo de las variables Paralelismo1 yParalelismo2, respectivamente. De esta forma, el bucle que recorre los agregadores (BSA=Bucle sobre los agregadores)) y el que recorre los locales (BSL=Bucle sobre los locales) se pueden ejecutar en serie o en paralelo dependiendo de los argumentos que se pasen a la llamada del programa. Pasamos a detallar la implementaci´ on de esos bucles y como se explotan las 4 fuentes de paralelismo en las siguientes subsecciones. 2.3.4.1. Paralelismo 1: en la E/S En la figura 2.21, identificamos con el n´ umero 1 la fase en la que se realiza la transferencia de los agregadores a los OSTs. Todas las transferencias est´ an etiquetadas con un 1 ya que cuando Paralelismo1=true las escrituras en los OSTs se realizan en paralelo. En la figura 2.20, la l´ ınea 6 tiene dos posibles implementaciones: for agregador in 0..#numAgreg si Paralelismo1=false, o forall agregador in 0..#numAgreg si Paralelismo1=true. En la implementaci´ on secuencial, usando for, la variable agregador toma secuencialmente los valores 0, 1, ..., hasta numAgreg-1. Esto es as´ ı porque en Chapel el operador #n en una expresi´ on de rango del tipo 0..#n representa nn´ umeros consecutivos empezando en 0 (por tanto, hasta el n-1). En la alternativa con forall, la variable agregador toma los mismos valores, pero el bucle se ejecuta en paralelo. La sentencia de la l´ ınea 7 es muy importante ya que especifica en qu´ elocale se debe 76 Cap´ ıtulo 2. Estudio de una interfaz de E/S en Chapel ejecutar el c´ odigo que sigue. Con la expresi´ on on Locales(agregador)do{...} se obliga a que el bloque de c´ odigo que sigue al do se ejecute en el locale con identificador agregador. En otras palabras, las l´ ıneas 9 a 28 se ejecutan en los locales 0 a numAgreg -1, de forma secuencial o paralela dependiendo de la variable Paralelismo1. En definitiva, Paralelismo1 activa o desactiva si los agregadores trabajan o no en paralelo. La agregaci´ on es b´ asicamente una redistribuci´ on de los datos de forma que cada agregador debe colectar todos los stripes que se deben almacenar en el OST correspondiente. La distribuci´ on destino es realmente una distribuci´ on c´ ıclica por bloques entre los OSTs, y el tama˜ no de bloque de esa distribuci´ on es el tama˜ no del stripe. De esta manera, un agregador se convierte en la ´ unica fuente de datos para su OST correspondiente, y aunque la escritura en los OSTs se realice en paralelo, evitamos la p´ erdida de rendimiento debida a los locks (ver secci´ on 1.2.4) que deben ser adquiridos para acceder astripes compartidos por varios locales. Esto es porque cada stripe se escribe completamente desde un ´ unico agregador (no hay stripes compartidos entre varios agregadores). El inconveniente es que los agregadores deben invertir un tiempo en colectar del resto de locales todos los stripes que deben ser escritos en sus respectivos OSTs. Aunque no est´ a contemplado en el c´ odigo de la figura 2.20, nuestra implementaci´ on tambi´ en funciona cuando el n´ umero de locales es inferior al n´ umero de OSTs. En ese caso no s´ olo todos los locales son agregadores, sino que estos tienen que servir a varios OSTs. Por ejemplo, si tenemos s´ olo 2 locales y 4 OSTs, el locale 0 colectar´ ıa los stripes que se deben almacenar en los OSTs 0 y 1, mientras que el locale 1 agregar´ ıa los stripes con destino en los OSTs 2 y 3. 2.3.4.2. Paralelismo 2: en la agregaci´ on En la figura 2.22 se muestra de forma esquem´ atica los pasos del proceso de agregaci´ on paralela. Por simplicidad, esta figura asume que ´ unicamente el locale 0 act´ ua de agregador y por tanto mantiene los buffers donde agregar los datos de todos los locales (incluido ´ el mismo) para luego realizar la escritura en el OST 0. La fase indicada con el n´ umero 1 en la figura expresa que todos los locales est´ an escribiendo en su correspondiente regi´ on del buffer en paralelo. Para el locale 0, esa operaci´ on implicar´ ıa un memcopy, pero en la implementaci´ on real, los datos del locale 0 se leen directamente de su zona de almacenamiento de datos (una matriz en el ejemplo). Los stripes de los dem´ as locales son agregados mediante llamadas, en ´ ultima instancia, a chpl_comm__gets, presentadas en la secci´ on 2.2.2. La segunda fase est´ a indicada en la figura con un 2 y representa la escritura en disco del buffer. Hacemos notar, que as´ ı como el Paralelismo1 habilita que haya varios agregadores agregando en paralelo, Paralelismo2 habilita que el trabajo de agregaci´ on que hace cada agregador sea tambi´ en paralelo. 2.3. Librer´ ıa de E/S en Chapel 77 local submatrix buffers Locale 0 (aggregator) local submatrix Locale 2 OST 0 OST 1 local submatrix Locale 3 local submatrix Locale 1 OST 2 ①① ① ① ② Figura 2.22: Paralelismo 2: Paralelismo en la agregaci´ on de los datos haciendo que se reciban datos de todos los locales en paralelo. Volviendo al c´ odigo de la figura 2.20, el bucle que agrega los datos del resto de los locales,BSL, aparece en la l´ ınea 17. De nuevo, ese bucle tiene dos posibles implementaciones: for loc in localesToWriteFrom..localesToWriteTo do si Paralelismo2=false, o forall loc in localesToWriteFrom..localesToWriteTo do si Paralelismo2=true. Es decir, si no activamos el paralelismo 2, el agregador ir´ a solicitando los datos a los locales uno tras otro, de forma secuencial. En caso contrario, se crear´ an tantos threads como locales involucrados en la agregaci´ on de forma que aceleraremos el proceso de confecci´ on del buffer que hay que escribir al OST. El inconveniente es la gran presi´ on a la que someteremos a la red de comunicaciones ya que, en el caso peor, esta operaci´ on de agregaci´ on implica una fase de comunicaciones de todos los locales con todos los agregadores. El rango de iteraciones de este bucle BSL identifica los locales,loc, desde los que se deben agregar datos para un agregador determinado. Los locales involucrados dependen de si Paralelismo4 est´ a activado o no, as´ ı que posponemos la explicaci´ on de este punto a la subsecci´ on correspondiente. 2.3.4.3. Paralelismo 3: entre E/S y agregaci´ on En la figura 2.23 mostramos el funcionamiento de la t´ ecnica de doble buffer que permite solapar la E/S con las comunicaciones necesarias para la agregaci´ on de datos. En la 78 Cap´ ıtulo 2. Estudio de una interfaz de E/S en Chapel local submatrix local row buffer a Locale 0 (aggregator) local submatrix local row Locale 1 ① ① file ② buffer b ① Figura 2.23: Paralelismo 3: Paralelismo entre la E/S y la agregaci´ on gracias al uso de doble buffer. figura suponemos un ´ unico agregador, el locale 0, que adem´ as de su matriz local, reserva espacio en memoria para dos buffers, el ay el b. La t´ ecnica de doble buffer se implementa en una serie de pasos. En cada paso, se usa uno de los buffers para ir agregando datos de otros locales, al tiempo que el otro buffer (que ya tiene los datos agregados del paso anterior) se escribe en el OST. En el momento que capturamos en este ejemplo, el buffer btiene elementos de una agregaci´ on que tuvo lugar en un paso anterior. Por tanto el n´ umero 1 en la figura indica que al mismo tiempo se est´ a escribiendo el buffer ben el disco y se est´ an agregando datos en el buffer adesde los locales 0 y 1. El n´ umero 2 en la figura indica que en el paso siguiente se escribir´ a el buffer aen disco y, aunque no mostrado expl´ ıcitamente, el buffer bse estar´ a utilizando para agregar datos. Aunque la figura 2.23 representa un movimiento de datos del locale 0 al buffer a, la implementaci´ on real evita este movimiento innecesario dentro del locale 0. De vuelta al c´ odigo de la figura 2.20, el c´ odigo entre las l´ ıneas 19 y 27 contiene las instrucciones relativas al procedimiento que acabamos de explicar. La variable Paralelismo3 controla si las funciones GatherToBuffer yWriteBufferToDisk se ejecutan dentro de un cobegin o fuera de ´ el. En Chapel, la instrucci´ on cobegin crea un thread para ejecutar de forma paralela cada l´ ınea del cuerpo del cobegin. Por tanto, si Paralelismo3 est´ a activado, GatherToBuffer yWriteBufferToDisk se ejecutar´ an al mismo tiempo. En ese caso, entre otros par´ ametros, estas funciones reciben dos buffers (bufferA ybufferB) como argumentos, para implementar la t´ ecnica de doble buffer. Aunque no se muestra expl´ ıcitamente en el pseudoc´ odigo de la figura 2.20, se usan variables de sincronizaci´ on para la comunicaci´ on entre los dos threads creados en el cobegin. Esto permite la coordinaci´ on entre los dos threads, evitando por ejemplo que el thread que ejecuta GatherToBuffer intente reescribir un buffer que a´ un no ha escrito a disco el thread que ejecuta WriteBufferToDisk. 2.3. Librer´ ıa de E/S en Chapel 79 Locale 0 OST 0 OST 1 OST 2 Locale 1 Locale 2 Locale 3 Locale 4 Locale 5 Locale 6 Locale 7 Figura 2.24: Paralelismo 4: Paralelismo en la escritura en cada OST cuando hay m´ as agregadores que OSTs. Los agregadores 0 y 3 se reparten el trabajo de agregar todos los stripes que han de almacenarse en el OST 0. 2.3.4.4. Paralelismo 4: en la escritura en un OST Cuando el n´ umero de locales es mucho mayor que el n´ umero de OSTs, tambi´ en es mucho mayor que el n´ umero de agregadores, porque hasta ahora hemos configurado un agregador para responsabilizarse de todos los stripes que van a un OST dado. En ese caso, un peque˜ no n´ umero de locales (los agregadores) necesitan acceder al resto para recolectar todos los stripes que ese agregador tiene que enviar al OST que representa. Por tanto, cada agregador se convierte en un cuello de botella, concentrando todos los datos que van dirigidos al correspondiente OST. Para evitar esta situaci´ on, la activaci´ on de la variable Paralelismo4 permite repartir el trabajo de agregaci´ on entre un mayor n´ umero de locales. Sin p´ erdida de generalidad, la implementaci´ on que hemos validado configura tantos agregadores como el m´ ultiplo del n´ umero de OSTs que m´ as se acerque al n´ umero total de locales. Por ejemplo, en la figura 2.24 tenemos 8 locales y 3 OSTs. Si Paralelismo4 no est´ a activado, tendremos tambi´ en 3 agregadores, uno para colectar los stripes que van destinados a cada OST (ver figura 2.21). Sin embargo, en nuestro ejemplo, con Paralelismo4 activado, vamos a configurar 6 agregadores. En la figura 2.24 mostramos expl´ ıcitamente los agregadores 0 y 3, los cuales se reparten el trabajo de agregar todos los stripes que van al OST 0: El agregador 0 agrega de los locales 0 a 2 y el agregador 3, de los locales 3 a 7. De igual forma los agregadores 1 y 4 se reparten el trabajo de agregar los stripes que van al OST 1, y por ´ ultimo, los agregadores 2 y 5, concentran los stripes que van al OST 2. Por tanto, al activar Paralelismo4 reducimos la presi´ on en la red de comunicaciones y tambi´ en podemos explotar cierta localidad si los agre- 86 Cap´ ıtulo 2. Estudio de una interfaz de E/S en Chapel En cuanto al n´ umero de locales, en las siguientes subsecciones se muestra el ancho de banda en escritura cuando este n´ umero cambia entre 4 y 64, y entre 16 y 512 locales. Algunas combinaciones de fuentes de paralelismo se exploraron s´ olo entre 4 y 64 locales, o incluso menos, ya que se comprob´ o que no ten´ ıa sentido continuar evaluando esas combinaciones por reportar pobres resultados, como se discutir´ a posteriormente. Por ´ ultimo, tambi´ en es necesario decidir cu´ al debe ser el tama˜ no de los buffers donde se realizan las agregaciones. Se evaluaron los siguientes tama˜ nos de buffer: 4, 10, 16, 25, 32, 40, 48 y 100MiB. Aunque en principio parecer´ ıa intuitivo que a mayor tama˜ no del buffer mayor ancho de banda en escritura, esto no ha resultado ser as´ ı. Recordemos primero que los buffers se reusan en los agregadores durante todo el proceso de escritura del archivo: los buffers se llenan con stripes desde el resto de los locales y luego se vac´ ıan al OST, para volver a empezar la fase de agregaci´ on. Cuando Paralelismo3 est´ a activado, estas dos fases (agregaci´ on y E/S) se ejecutan en paralelo pero de forma sincronizada. Esto es, no se puede intercambiar el buffer de agregaci´ on y el de E/S hasta que las dos operaciones no est´ en completadas. En la pr´ actica, es imposible balancear ambas fases y la m´ as r´ apida tendr´ a que esperar a que la m´ as lenta termine. Pues bien, en nuestros experimentos hemos comprobado que la fase de agregaci´ on es m´ as lenta que la de E/S si los tama˜ nos de buffer son mayores de 32MiB. Es decir, para tama˜ nos de buffer muy grandes, el tiempo de comunicaciones entre todos los locales y todos los agregadores para que a cada uno le lleguen los stripes que almacenan los OSTs respectivos es mayor que el tiempo de env´ ıo del buffer de los agregadores a los OSTs. Es cierto que la red de comunicaciones entre locales es m´ as r´ apida que la de E/S. Pero tambi´ en lo es que el patr´ on de comunicaciones de agregaci´ on es casi “todos-a-todos” ya que se est´ a haciendo una redistribuci´ on de los stripes, lo que termina siendo m´ as costoso. Por tanto, para mantener el flujo de escritura en los OSTs de forma ininterrumpida y obtener por tanto el mayor ancho de banda en escritura, el tama˜ no del buffer debe ser, en esta arquitectura, de tama˜ no inferior a 32MiB. El resto de los experimentos que mostramos en la secci´ on 2.4.5 se han realizado con buffers de 10MiB. 2.4.4. Estimaci´ on del rendimiento ideal En [98] se anuncia que el rendimiento m´ aximo de Spider II es de 1 TByte/seg. Una versi´ on simplista del rendimiento nos podr´ ıa llevar a dividir el ancho de banda total entre el n´ umero de OSTs para obtener el rendimiento por OST, de forma que 1000TB/s entre 2016 OSTs es aproximadamente 500 MBytes por segundo. Explotando 16 OSTs en paralelo se deber´ ıa tener un valor m´ aximo de 8000 Mbytes por segundo. Sin embargo esta cota obtenida es poco realista, ya que hay m´ as factores que influyen en el rendimiento y los 16 OSTs pueden estar siendo usados concurrentemente por otros procesos. 2.4. Evaluaci´ on del nuevo interfaz de E/S 87 Para obtener una estimaci´ on m´ as realista de la cota m´ axima del rendimiento ideal en escritura, hemos medido el ancho de banda de la operaci´ on de escritura m´ as b´ asica en las condiciones m´ as favorables que hemos encontrado. M´ as precisamente, hemos usado el comando dd, un comando habitual en sistemas operativos UNIX para realizar operaciones de E/S. El procedimiento consisti´ o en enviar un comando dd a cada locale a trav´ es del sistema de colas de forma que cada uno de esos locales escriba en un archivo diferente al mismo tiempo. Los argumentos del comando fueron los siguientes: dd if=/dev/zero of=delete.‘hostname‘.bor bs=1M count=80 es decir, la fuente de datos son 0’s, el tama˜ no de bloque, bs, es 1MiB y se escriben 80 bloques en un fichero de nombre delete.EOS-XX.bor (la XX del nombre se substituye por el identificador del locale desde el que se est´ a escribiendo). El resultado es que cada locale crea un fichero de 80MiB. La prueba se ha ejecutado usando dos tama˜ nos de stripe diferentes (1MiB y 4MiB), 4 OSTs por fichero, y con distinto n´ umero de locales (1, 16, 32, 48 y 64). Los anchos de banda obtenidos cuando s´ olo escribe un locale estaban en el rango de entre 600-800 MBytes/s. Cuando se ejecutan 16 comandos dd en paralelo desde 16 locales el ancho de banda sube a unos 3,000 MBytes/s (esto es el resultado de dividir el volumen agregado de todos los ficheros –16×80MiB– dividido entre el tiempo que se consume en crear los 16 ficheros). Si se ejecuta desde 32 nodos, llegamos a unos 4800 MBytes/s. Sin embargo, ejecutar desde m´ as nodos ya no mejora el ancho de banda as´ ı que entendemos que en ese punto se alcanza el cuello de botella en la red que comunica los nodos de c´ alculo con los servidores del sistema de ficheros. En [141] y en [29] se comprueba que escribir en paralelo en ficheros independientes es m´ as r´ apido que escribir en paralelo en un ´ unico fichero, debido a los locks que tienen que gestionar la coherencia en la escritura. Por tanto, entendemos que los 4800MBytes/s representan una cota superior a la que deber´ ıamos intentar acercarnos en nuestra implementaci´ on paralela. 2.4.5. Evaluaci´ on de los algoritmos de escritura Hemos realizado varios conjuntos de experimentos para evaluar el rendimiento de los algoritmos de escritura que se proponen en la secci´ on 2.3.2. Dado que existen 16 posibles combinaciones, para ser sistem´ aticos presentamos los resultados de ancho de banda en dos grupos: Combinaciones (0,x,x,x), en la figura 2.25 Combinaciones (1,x,x,x), en la figura 2.26 88 Cap´ ıtulo 2. Estudio de una interfaz de E/S en Chapel 50 100 150 200 250 300 350 0 10 20 30 40 50 60 70 MB/sec nodes eos chapel with 8 osts and 0xxx (0,0,0,0) (0,1,0,0) (0,0,1,0) (0,1,1,0) (0,1,1,1) Número de locales EOS con 16 OSTs y configuraciones (0, x, x, x) Figura 2.25: Ancho de banda para las mejores combinaciones (0,x,x,x). En la figura 2.25 presentamos las medidas obtenidas con 16 OSTs, entre 4 y 64 locales y usando los algoritmos definidos por las tuplas (0,0,0,0), (0,1,0,0), (0,0,1,0), (0,1,1,0) y (0,1,1,1). En todas ellas Paralelismo1 est´ a desactivado. Otras tres combinaciones dentro de este grupo son (0,0,0,1), (0,1,0,1) y (0,0,1,1), pero no se muestran en la figura ya que obtienen resultados incluso peores. Como vemos y en general, los anchos de banda obtenidos en estas condiciones est´ an entre los 100MB/s y los 300MB/s. Esto es, m´ as de un orden de magnitud por debajo de la cota m´ axima estimada en la secci´ on anterior (4800MB/s). Como veremos a continuaci´ on, estos resultados son tambi´ en mucho peores que los que obtenemos con las combinaciones en las que Paralelismo1 est´ a activado. De entre las combinaciones que presentamos en la figura 2.25, la (0,0,0,0) no explota ninguna fuente de paralelismo y es una de las que peor rendimiento ofrecen. Al aumentar el n´ umero de locales, aumenta el tama˜ no del fichero y casi en la misma medida el tiempo total de escritura. Esto conduce a que el ancho de banda permanezca en torno a los 200MBs. Por el otro extremo, una de las mejores combinaciones es la (0,1,0,0) (s´ olo activamos el paralelismo en la agregaci´ on), especialmente cuando el n´ umero de locales 2.4. Evaluaci´ on del nuevo interfaz de E/S 89 500 1000 1500 2000 2500 3000 3500 4000 4500 0 50 100 150 200 250 300 350 400 450 500 550 MB/sec locales EOS Chapel con 16 OSTs (1,0,0,0) (1,1,1,0) (1,0,1,0) (1,1,0,0) (1,0,1,1) (1,1,1,1) EOS con 16 OSTs y configuraciones (1,x,x,x) Número de locales Figura 2.26: Ancho de banda para las mejores combinaciones (1,x,x,x). es superior a 16, pero a´ un as´ ı el ancho de banda obtenido es bastante pobre. El problema subyacente de estas combinaciones es que ninguna de ellas deja trabajar a los agregadores en paralelo. Por tanto, los OSTs que son alimentados desde cada agregador tambi´ en est´ an trabajando de forma secuencial. Es decir, es necesario activar el Paralelismo1 para explotar mejor el hecho de que un fichero se almacena de forma distribuida entre varios OSTs. En la figura 2.26 se muestran resultados con 16 OSTs, entre 16 y 512 locales y con Paralelismo1 activado. A partir de 128 locales se pueden distinguir claramente tres grupos distintos de acuerdo al rendimiento, el m´ as alto, unos 3800MB/s, en el caso de los algoritmos (1,1,x,0), uno intermedio, sobre los 3000MBs, con (1,0,x,0), y el m´ as bajo rendimiento, alrededor de 2200MB/s, con (1,x,x,1). Claramente, en comparaci´ on con las configuraciones (0,x,x,x) de la grafica anterior, activar el paralelismo entre agregadores dispara el rendimiento mas o menos un orden de magnitud. Esto confirma la ventaja de escribir en paralelo en los OSTs y nos acerca a los 4800MB/s que hemos medido al escribir en paralelo en ficheros independientes. Es decir, existen combinaciones de paralelismo que evitan la p´ erdida de rendimiento debida a los locks de Lustre y casi igualan 90 Cap´ ıtulo 2. Estudio de una interfaz de E/S en Chapel las prestaciones de la escritura en ficheros independientes. El rendimiento de 4800MB/s no se llega a alcanzar, ya que en esta ´ ultima alternativa ideal no tiene lugar la redistribuci´ on entre locales y agregadores, penalizaci´ on que si aparece cuando escribimos en un ´ unico fichero. Sin embargo, algunas de las combinaciones siguen provocando bloqueos en los OSTs. Nos referimos a las combinaciones que tienen Paralelismo4 activado: las del tipo (1,x,x,1). En contra de nuestras expectativas, esta posibilidad en lugar de estar entre las mejores est´ a entre las que peor redimiento reportan. De hecho, las combinaciones (1,0,0,1) y (1,1,0,1) no se muestran en la figura 2.26 ya que reportaron peores resultados y adem´ as estas combinaciones est´ an en cierta forma cubiertas por las combinaciones (1,0,1,1) y (1,1,1,1) que s´ ı est´ an presentes en la figura entre 16 y 512 locales. Los problemas de rendimiento con las combinaciones (1,x,x,1) radican en que hay varios agregadores escribiendo en el mismo OST (cosa que no ocurre en las combinaciones (1,x,x,0)). Recordemos, pese a que Paralelismo1 yParalelismo4 est´ en activados, las escrituras de stripes del mismo fichero en el mismo OST est´ an serializadas por locks de Lustre (ver secci´ on 1.2.4.1). Sin embargo, esper´ abamos que esta fuente de paralelismo s´ ı ayudase a reducir la carga en la red de comunicaciones entre locales, ya que al haber m´ as agregadores por OST, estos se reparten la carga de redistribuci´ on de los stripes que se han de almacenar en esos OSTs. Por ejemplo, con 512 locales, 16 OSTs y Paralelismo4 desactivado, s´ olo hay 16 agregadores colectando stripes potencialmente desde los otros 512 locales. Por el contrario, con Paralelismo4 activado, los 512 locales toman el papel de agregador, 32 de ellos agregan para un mismo OST, lo que significa que cada uno de estos ´ ultimos 32 agrega de 16 locales diferentes. Lo mismo ocurre para 256, 128, 64, 32 y 16 locales, ya que en todos esos casos cada agregador necesitar´ ıa agregar stripes de otros 16 locales. Por tanto, esper´ abamos que tuviese alg´ un impacto en el rendimiento el reducir la carga de los agregadores de tener que comunicarse con 512 locales a tener que comunicarse s´ olo con 16 locales. Adem´ as, las combinaciones (1,x,x,1) estudiadas presentan una ca´ ıda de rendimiento pronunciada entre 32 y 63 locales. Con Paralelismo4 activado, todas esas situaciones tienen en com´ un que vamos a tener dos o 3 agregadores por OST. Para confirmar que este era el hecho determinante en la caida de rendimiento, se ejecutaron con 8 OSTs las combinaciones (1,x,x,1), cuyos resultados entre 8 y 56 locales se muestran en la figura 2.27. En este caso es entre 16 y 23 locales cuando nuestro algoritmo de escritura configura 2 agregadores por OST y cuando se observan las caidas importantes en el rendimiento de escritura. Para la combinaci´ on (1,0,0,1) se midi´ o el ancho de banda desde 8 hasta 36 locales con paso 1 locale (es decir de 1 en 1). En la figura se aprecia que justo en los valores {16, 17, 18,..., 23}locales el ancho de banda es el menor y en torno a 250MBs. Estos resultados nos hacen concluir que tener 2 o 3 agregadores contendiendo en el acceso a un OST provoca m´ as p´ erdida de rendimiento que una configuraci´ on en 2.4. Evaluaci´ on del nuevo interfaz de E/S 91 0 500 1000 1500 2000 2500 3000 3500 4000 4500 10 20 30 40 50 MB/sec N mero de locales EOS con 8 OSTs y combinaciones (1,x,x,1) (1,0,0,1) (1,0,1,1) (1,1,0,1) (1,1,1,1) Número de locales Figura 2.27: Ancho de banda para las combinaciones (1,x,x,1) con 8 OSTs. la que el n´ umero de agregadores contendiendo es mayor. Sin embargo no hemos podido obtener m´ as informaci´ on que explique los motivos reales de ese hecho. Si desactivamos Paralelismo2, es decir no explotamos el paralelismo en agregaci´ on, vemos en la figura 2.26 que obtenemos un rendimiento intermedio en torno a 3000MB/s con las configuraciones (1,0,x,0). El mejor rendimiento, entorno a 3800MB/s, se obtiene cuando adem´ as del Paralelismo1 tambi´ en activamos el Paralelismo2, es decir para las combinaciones (1,1,x,0). Como vemos en la figura, especialmente a partir de 128 locales, activar o desactivar el Paralelismo3 tiene muy poco impacto tanto en las combinaciones (1,0,x,0) como en las (1,1,x,0). Recordemos que este paralelismo se encarga de solapar la fase de agregaci´ on con la de E/S. El motivo por el que el Paralelismo3 no afecta es porque aunque est´ e desactivado, los buffers y cach´ es intermedios en los nodos LNET y OSSs est´ an virtualmente permitiendo que este solape entre E/S y agregaci´ on tambi´ en est´ e teniendo lugar. Con todo esto, la recomendaci´ on es explotar tanto la fuente de paralelismo en E/S (los agregadores escribiendo al mismo tiempo en distintos OSTs) como la fuente de paralelismo en la agregaci´ on (cada agregador recolecta stripes de los distintos locales en 92 Cap´ ıtulo 2. Estudio de una interfaz de E/S en Chapel 0 500 1000 1500 2000 2500 3000 3500 4000 4500 0 20 40 60 80 100 120 MB/seg locales Chapel vs MPI-IO EOS 8 OSTs MPI-IO (1,1,0,0) (1,1,1,0) Número de locales Figura 2.28: Comparaci´ on entre los algoritmos (1,1,x,0) y MPI IO con 8 OSTs y entre 1 y 128 locales. paralelo). En esa situaci´ on evitamos los locks de Lustre y nos acercamos a los 4800MB/s que se midieron en la secci´ on 2.4.4 en circunstancias ideales. 2.4.6. Comparaci´ on con MPI IO La librer´ ıa MPI IO tambi´ en incorpora un mecanismo de agregaci´ on similar al que nosotros hemos implementado en Chapel [101]. En particular, la versi´ on de MPI IO m´ as avanzada que hemos podido ejecutar se basa en un algoritmo de agregaci´ on del tipo (1,0,1,0). En la figura 2.28 se compara el rendimiento de la implementaci´ on MPI IO, y las mejores combinaciones que hemos encontrado en la secci´ on anterior, (1,1,x,0). El c´ odigo MPI IO usado fue la versi´ on optimizada por Cray [142][29] que est´ a basada en ROMIO [119][121] [118], una implementaci´ on de alto rendimiento de MPI desarrollada y mantenida en el Argonne National Laboratory (ANL), a˜ nadiendo c´ odigo de la librer´ ıa Lustre ADIO de Sun Microsystems [71]. 2.5. Trabajos relacionados 93 El c´ odigo MPI usa una matriz bidimensional distribuida por bloques de tama˜ no variable dependiendo del n´ umero de locales, de forma que a cada locale le correspondan 2GB de datos igual que en la implementaci´ on Chapel. El sistema donde se realiz´ o la comparaci´ on es EOS y tanto en la versi´ on Chapel como la versi´ on MPI, escriben en un ´ unico archivo con tama˜ no de stripe de 1 MiB, distribuido en 8 OSTs. En la figura podemos comprobar c´ omo la implementaci´ on en Chapel reporta mayor ancho de banda que la implementaci´ on MPI cuando el n´ umero de locales es inferior a 128. El motivo es que la implementaci´ on MPI no explota el paralismo en la agregaci´ on (Paralelismo2) que como hemos visto en la secci´ on anterior suele conducir a mejoras significativas en el ancho de banda. Sin embargo, para 128 locales, el rendimiento entra la implementaci´ on Chapel y la versi´ on MPI es muy similar (entorno a 2400MB/s). Teniendo en cuenta que Chapel es un modelo de programaci´ on de m´ as alto nivel y m´ as productivo que MPI, creemos que los resultados obtenidos permiten confiar en que la E/S en Chapel puede ser competitiva en rendimiento. 2.5. Trabajos relacionados Algunos trabajos previos han estudiado el rendimiento de la E/S en sistemas masivamente paralelos [77][46][142], proporcionando una gu´ ıa sobre los factores que le afectan. Otros autores [76] han trabajado en la evaluaci´ on del rendimiento de sistemas masivamente paralelos, identificando los grandes retos presentes en esos sistemas. M´ as cercano a nuestra investigaci´ on encontramos [142] donde se estudia el rendimiento de una amplia variedad de interfaces de E/S paralelos en plataformas de gran tama˜ no basadas en Lustre, incluyendo POSIX IO, MPI IO y HDF5. Su estudio incide especialmente en c´ omo la distribuci´ on de los datos sobre los elementos de procesamiento puede afectar a las estrategias de E/S. En cualquier caso apuntan a que las rutinas de E/S deben tener en cuenta tanto el tama˜ no de stripe como el factor de stripes, y que el n´ umero de procesos realizando accesos a un OST debe estar limitado para evitar la contenci´ on. En general, si los usuarios quieren alcanzar el m´ aximo rendimiento tienen que definir cuidadosamente los par´ ametros usados en las llamadas de E/S, en particular tienen que ser conscientes de la distribuci´ on de datos usada en la aplicaci´ on en cuesti´ on, as´ ı como de la organizaci´ on y gesti´ on de los datos en el sistema de ficheros. En este trabajo nos distinguimos en que se provee al usuario con funciones de alto nivel de E/S. mientras que es el runtime, que conoce tanto la distribuci´ on de los datos entre los locales as´ ı como la distribuci´ on del sistema de ficheros subyacente, quien se encargue de la organizaci´ on eficiente de los accesos. En un trabajo m´ as reciente [29] se analiza el bajo rendimiento que se obtiene con 94 Cap´ ıtulo 2. Estudio de una interfaz de E/S en Chapel MPI IO (ROMIO) en entornos donde se usa el sistema de ficheros Lustre, y se propone una nueva librer´ ıa de nivel de usuario, llamada Y-Lib, para mejorar los resultados. All´ ı descubren que las implementaciones actuales optimizadas de ROMIO producen patrones de acceso en las comunicaciones de todos a todos (all-to-all), lo que provoca contenci´ on en las operaciones de acceso a los OSTs y OSSs. Cuando aparece esta situaci´ on en nuestros experimentos, encontramos resultados similares, como se ha explicado en la secci´ on 2.4.5. En [86] llegan a las mismas conclusiones que en [29], pero ninguno de esos trabajos, ni de los otros vistos, aborda el problema que afecta a la velocidad en MPI-IO cuando se producen accesos concurrentes a un mismo fichero en Lustre, y es el relacionado con el funcionamiento de los locks. Dicho funcionamiento se explica en [47], y en las secciones 1.2.4 y 2.4.5. En nuestro trabajo, proponemos una estrategia de agregaci´ on consciente de este problema que es la finalmente consigue mayor rendimiento. 2.6. Conclusiones En este cap´ ıtulo hemos presentado una nueva interfaz de acceso paralelo a la E/S y una implementaci´ on de esa interfaz optimizada para ser usada en Lustre, el sistema de ficheros paralelo m´ as usado en grandes sistemas HPC de la actualidad. El nuevo interfaz de E/S, implementado en el lenguaje de programaci´ on de alta productividad Chapel, se basa en una nueva clase o en una nueva distribuci´ on que permite mapear un array en el sistema de ficheros Lustre. Las escrituras/lecturas en/de regiones de ese array se convertir´ an mediante el compilador y el runtime de Chapel en escrituras/lecturas en/de fichero. Esto aumenta la productividad del programador, que se puede despreocupar de la implementaci´ on de la distribuci´ on de datos y de la comunicaci´ on eficiente de los mismos entre los locales (esta tarea queda delegada en el compilador y runtime). El programador tampoco tiene que preocuparse de la gesti´ on de bloques en el sistema de ficheros paralelo, ni de los accesos a los servidores de datos. Otra ventaja de nuestra implementaci´ on es que es independiente del n´ umero de nodos que se utilice para crear y escribir el fichero, y posteriormente tambi´ en es independiente del n´ umero de nodos que se use para leerlo. La conclusi´ on principal es que, cuando se quiere acceder a grandes vol´ umenes de datos se puede aprovechar eficientemente el sistema de ficheros paralelo si la implementaci´ on de la E/S consigue: i) escribir en paralelo en los OSTs desde los locales que hacen la funci´ on de agregadores; ii) agregar en paralelo los datos a escribir en cada OST; y iii) minimizar la posibilidad de contenci´ on, desde distintos locales, en los accesos a cada OST. Por otro lado, el solapamiento entre E/S y agregaci´ on ya suele estar implementado de forma impl´ ıcita mediante los buffers y caches de los nodos intermedios entre los clientes y los discos de almacenamiento. Por tanto, en nuestros resultados no aporta un 2.6. Conclusiones 95 beneficio sustancial una implementaci´ on expl´ ıcita de esta fuente de paralelismo. Otros p´ arametros del algoritmo de E/S que hemos evaluado y que tienen gran impacto en el ancho de banda de escritura deben ser estudiados y configurados adecuadamente: i) no tiene sentido aumentar el n´ umero de OSTs cuando ha dejado de ser el cuello de botella (en EOS, los LNET son el cuello de botella para m´ as de 16 OSTs); y ii) elegir con cuidado el tama˜ no del buffer que realiza la agregaci´ on y que ser´ a el que contiene los bloques que se env´ ıan en cada operaci´ on de E/S al OST. 102 Cap´ ıtulo 3. Almacenamiento optimizado de datos de secuenciaci´ on gen´ etica maci´ on a partir de esos datos [106], por lo que los costes totales no bajan tanto como es de esperar cuando se tienen en cuenta todos los factores implicados. Por otro lado, la cantidad de datos a acceder es tan grande que se comienza a encontrar un cuello de botella en la E/S en los sistemas HPC, provocando una p´ erdida de rendimiento, al estar los elementos de proceso esperando a los datos, debido al alto volumen de datos que usan los algoritmos gen´ eticos. En [80] se describe un planificador de tareas que tiene en cuenta la E/S a la hora de planificar dichas tareas, para evitar que al ejecutar muchas tareas en paralelo se produzcan esperas debido a la E/S. En ese trabajo se utilizan los sistemas de fichero NFS y Lustre, y la clave de la estrategia que se plantea tiene en cuenta qu´ e OSTs (servidor de datos) usa cada tarea para que no coincidan, evitar la contenci´ on y as´ ı poder tener un mejor rendimiento global. Hay que tener tambi´ en en cuenta que el problema de tener una cantidad muy grande de datos se produce si se almacenan los datos en bruto. Si se almacenan tras ser procesados ocupar´ an mucho menos, ya que las repeticiones de datos desaparecer´ an, y se almacenar´ a s´ olo lo que se haya determinado que es el genoma real del organismo al que se le realiz´ o la secuenciaci´ on. El problema de esto es que distintos algoritmos de secuenciaci´ on dan distintos resultados, por lo que la decisi´ on de borrar los datos brutos originales no es sencilla, ya que una mejora futura de los algoritmos permitir´ ıa obtener resultados m´ as precisos. Por esta raz´ on, hoy en d´ ıa se suelen almacenar los datos en bruto para que que en un futuro se pueda obtener m´ as informaci´ on a partir de ellos. 3.1.3. Compresi´ on de los datos Se han usado tambi´ en distintos m´ etodos para comprimir los datos obtenidos a partir de los secuenciadores gen´ eticos. Se puede encontrar un resumen en [131]. Como se ha mencionado, estos datos suelen estar organizados en nombres, secuencias y calidades. A veces se usa un campo extra para apuntar alguna informaci´ on adicional de las secuencias. Comprimiendo las calidades Comenzando por las calidades, hay que tener en cuenta que tienen ciertas propiedades que dependen del mecanismo usado para realizar la secuenciaci´ on, lo que permite aprovecharnos de eso para comprimir, e incluso para tener compresi´ on con p´ erdidas sin que afecte al procesamiento de los datos. Hay varias formas est´ andar de almacenar las calidades. Usualmente se almacenan como n´ umeros enteros separados por espacios, o como caracteres ASCII que representan los n´ umeros. El rango de valores que pueden tomar las calidades viene dado por la tecnolog´ ıa que se haya usado para realizar la secuenciaci´ on, y suele ser bastante peque˜ no, por ejemplo con la tecnolog´ ıa de Illumina var´ ıa entre 0 y 40 [21] y en el resto 3.1. El almacenamiento de secuencias gen´ eticas 103 de tecnolog´ ıas se suele usar la codificaci´ on usada por el programa PHRED [42] y variar entre 0 y 99 [70]. Un m´ etodo bastante habitual para mejorar la compresi´ on de la calidad es cuantificar el valor de las calidades en varios niveles. Por ejemplo, si las calidades toman valores entre 1 y 100, podr´ ıamos usar s´ olo tres niveles de cuantificaci´ on: mala (1), regular (50) y buena (100), de forma que los valores a comprimir sean m´ as simples, y se pueda alcanzar una compresi´ on mayor [96]. Esta compresi´ on es con p´ erdidas, ya que tras comprimir no se podr´ an recuperar los valores originales, pero el resultado suele ser suficiente para los algoritmos que se van a aplicar, y tiene el beneficio de facilitar la compresi´ on, ahorrando mucho espacio [96]. Comprimiendo las secuencias gen´ eticas Los m´ etodos para comprimir los datos de los genomas se pueden clasificar en tres tipos dependiendo de en qu´ e se base la compresi´ on [131]. El m´ as b´ asico se basa en la manipulaci´ on de bits, a continuaci´ on tenemos los basados en diccionarios y los m´ as sofisticados son los basados en referencias. Vamos a ver las ventajas e inconvenientes de cada uno de ellos. Los algoritmos basados en manipulaci´ on de bits se aprovechan de que las bases est´ an formadas por cuatro letras (ATCG), por tanto no es necesario usar ocho bits para almacenarlos, que es lo que ocupan cuando se usa la codificaci´ on ASCII para las letras. En este caso, con dos bits es suficiente: A →00, T →01, C →10, G →11. Simplemente con ese sistema las secuencias ocupan cuatro veces menos. En el caso de la calidad podemos adoptar una soluci´ on equivalente, si la calidad va de 1 a 100, el valor se puede almacenar en un s´ olo byte con una codificaci´ on binaria en vez de usar tres caracteres. Podemos encontrar ejemplos de este tipo de compresi´ on en [123]. Tenemos un problema si en el origen aparecen m´ as de las cuatro letras b´ asicas, o may´ usculas y min´ usculas. Adem´ as el fichero resultante es binario, y por tanto no se puede ver como texto. Los m´ etodos basados en diccionarios se basan en reemplazar secuencias por referencias a un diccionario, que generalmente se suele ir construyendo sobre la marcha. Los ejemplos m´ as t´ ıpicos son los basados en algoritmos como el LZW (Lempel-ZivWelch) [138]. Tambi´ en se puede usar una codificaci´ on Huffman, o cadenas de Markov. En estos casos se busca crear una codificaci´ on m´ ınima. Eso implica almacenar tambi´ en el diccionario adem´ as de los datos [75]. Finalmente los algoritmos basados en referencias son los que tienen un almacenamiento m´ as eficiente, siempre que ya exista un genoma de referencia para el organismo del que se han obtenido los datos, ya que el organismo secuenciado se deber´ ıa de parecer mucho al de referencia, y por tanto casi todas las secuencias estar´ an ya en ese genoma de referencia. En este caso s´ olo hay que almacenar punteros a los sitios donde aparece cada secuencia. Habr´ a caracter´ ısticas que diferencian unos organismos de otros, y esas partes ser´ an mucho m´ as dif´ ıciles de comprimir. Otra desventaja es que la velocidad de 104 Cap´ ıtulo 3. Almacenamiento optimizado de datos de secuenciaci´ on gen´ etica compresi´ on no se puede predecir. 3.2. FQbin: una librer´ ıa para el almacenamiento y acceso a datos de secuenciaci´ on gen´ etica Para reducir el impacto sobre el sistema de almacenamiento, que supone la enorme cantidad de informaci´ on que se genera con las tecnolog´ ıas de pr´ oxima generaci´ on de secuenciaci´ on (NGS, Next Generation Sequencing), proponemos FQbin [59]. Se trata de un formato nuevo y vers´ atil para la compresi´ on, almacenamiento y recuperaci´ on de datos obtenidos por secuenciaci´ on gen´ etica de alto rendimiento de ADN. Es un formato compatible con FASTA y FASTQ y como veremos, mejora el rendimiento de las propuestas existentes. Se basa en la librer´ ıa zlib y ofrece una compresi´ on de hasta 10x. El fichero comprimido puede leerse y descomprimirse en un tiempo menor al que se necesita con el formato FASTQ, que se puede considerar el formato de facto en este campo. Adem´ as, nuestra propuesta de formato proporciona acceso aleatorio a todas las secuencias que almacena [65]. Por otro lado, hemos comprobado que en sistemas distribuidos la ventaja del nuevo formato propuesto es mayor cuanto m´ as lenta sea la red. Se han realizado dos implementaciones de FQbin, una en C y otra en Chapel, y en este cap´ ıtulo comparamos el rendimiento y la productividad de cada una de las soluciones. 3.3. El formato de contenedor FQbin El formato FQbin unifica los campos FASTA, Quality y extras de cada secuencia comprimi´ endolos y a˜ nadiendo una cabecera para facilitar su acceso aleatorio. En disco se almacenar´ a la cabecera seguida por los datos de esa secuencia, todo comprimido. Se puede ver un esquema en la figura 3.2. El formato FQbin comienza con una cabecera de longitud variable, los cuatro primeros bytes (FILE HEADER LENGTH) indican la longitud de la cabecera del fichero, a continuaci´ on viene una cadena alfanum´ erica (STRING) con el identificador del formato, y dos campos, VERSION ySUBVERSION, con la versi´ on y subversi´ on, controlando la existencia de futuras versiones sin provocar problemas con los ficheros y librer´ ıas ya existentes. A continuaci´ on viene un conjunto de bloques de datos comprimidos (compressed blocks), conteniendo cada bloque un n´ umero configurable de secuencias, por defecto 10000, comprimidas como un s´ olo flujo zlib. Cada uno de los campos de secuencias (SEQUENCE RECORD) dentro del bloque comprimido comprende una cabecera de 3.3. El formato de contenedor FQbin 105 SEQUENCE HEADER SEQUENCE RECORD FILE HEADER compressed block 1 compressed block 2 compressed block 3 FILE HEADER FILE HEADER LENGTH-4 digit STRING VERSION SUBVERSION SEQNAME FASTA LENGTH SEQ HEADER LENGTH-4 digit QUAL LENGTH EXTRA LENGTH SEQUENCE HEADER SEQUENCE DATA QUAL DATA EXTRA DATA SEQ1 SEQ3 SEQ2 SEQ4 … … Figura 3.2: Formato FQbin. secuencia (SEQUENCE HEADER) y a continuaci´ on el resto de los datos (SEQUENCE DATA, QUAL DATA, EXTRA DATA). La cabecera de cada una de las secuencias, SEQUENCE HEADER, que desglosamos en la parte inferior de la figura 3.2, comienza con cuatro bytes que contienen la longitud de la cabecera (SEQUENCE HEADER LENGTH), que es variable. A continuaci´ on vienen cuatro campos de texto que contienen: i) el nombre de la secuencia (SEQNAME), que es el campo ID que sirve de identificaci´ on para realizar el acceso aleatorio; ii) la longitud de la secuencia (FASTA LENGTH); iii) la longitud de la calidad (QUAL LENGTH); y iv) la longitud de los extras (EXTRA LENGTH), si es que hay. Una vez que el bloque se llena, es decir, se llega al n´ umero de secuencias predeterminado, se cierra el flujo y se crea uno nuevo. De esta forma cuando hay que extraer un dato s´ olo hay que descomprimir el bloque en el que se encuentra y no todo el fichero, ahorrando tiempo y accesos al disco. Como se explicar´ a posteriormente en la secci´ on 3.8.4 esta separaci´ on en bloques sirve adem´ as como protecci´ on contra la corrupci´ on de los datos. Cada vez que se crea un bloque nuevo se comienza con un diccionario vac´ ıo en zlib. Como los datos a comprimir son secuencias gen´ eticas probamos a realizar una precarga del diccionario con datos de secuencias gen´ ericas, con la intenci´ on de no partir de un diccionario vac´ ıo cada vez que se comienza a comprimir, pero observamos que 106 Cap´ ıtulo 3. Almacenamiento optimizado de datos de secuenciaci´ on gen´ etica no se produc´ ıa ninguna mejora en la compresi´ on. El motivo es que los datos son semi aleatorios, y por tanto la precarga del diccionario no le supone una ventaja respecto a partir del diccionario vac´ ıo, ya que se consigue adaptar en muy poco tiempo a las cadenas que se va encontrando, y en caso de tener el diccionario lleno, tiene que ir olvidando los c´ odigos que contiene para aprender los nuevos, lo que hace que no se obtenga ninguna mejora precargando el diccionario. La compatibilidad con software actual y versiones anteriores est´ a garantizada gracias a que los contenidos de un fichero FQbin pueden ser enviados v´ ıa streaming a programas que aceptan datos FASTA/QUAL o FASTQ, evitando la necesidad de realizar una conversi´ on, lo que suele implicar realizar una copia de los datos en otro fichero. 3.4. Acceso aleatorio El contenedor FQbin y sus herramientas asociadas permiten un acceso aleatorio r´ apido a secuencias usando su ID, que es su nombre. Esto se consigue gracias a dos ficheros externos, uno que almacena un ´ ındice de los datos y otro fichero con un hash al ´ ındice. Estos ficheros se pueden regenerar en cualquier momento. El fichero de ´ ındices almacena la posici´ on de cada una de las secuencias en el fichero principal. Para que el acceso sea m´ as r´ apido se incluye el nombre de la secuencia, la posici´ on del bloque donde est´ a almacenado y la posici´ on dentro del bloque donde comienza la cabecera de la secuencia, de esta forma se puede realizar un seek al bloque donde est´ a la secuencia, y a continuaci´ on habr´ ıa que moverse descomprimiendo, realizando un gzseek, hasta la posici´ on donde est´ a la secuencia dentro del bloque. Esta estrategia permite un acceso aleatorio r´ apido a una secuencia, aunque no es realmente instant´ aneo, ya que ser´ ıa necesario descomprimir el fichero al menos hasta el punto donde se encuentra la secuencia, y, si hay muchos millones de bloques almacenados en el mismo fichero, llegar hasta el ´ ultimo descomprimiendo puede tardar unos segundos. La siguiente mejora que se plante´ o fue crear una tabla hash que acelere el acceso a la tabla de ´ ındices. Para ello se ordena el fichero de ´ ındices, se comprime por bloques y en la tabla hash se almacena el nombre de la primera y ´ ultima secuencia que almacena cada bloque y su posici´ on dentro del fichero de ´ ındices. Este fichero de hash ocupa un espacio m´ ınimo respecto al total, sobre un 0,0001 %, por ejemplo, unos 12 Kilobytes para unos 9.4 Gigabytes de datos. Cuando es necesario acceder a una secuencia concreta, primero se carga el fichero de hash, se mira en qu´ e bloque del fichero de ´ ındices se puede encontrar la secuencia buscada, se accede a ese bloque y se busca esa secuencia. Si no se encuentra es que no est´ a en el fichero, se da un error de “secuencia no encontrada”. Si se encuentra se lee la posici´ on del bloque dentro del fichero FQbin y la posici´ on de 3.5. Simplificaci´ on de los valores de la calidad 107 esa secuencia dentro del bloque, finalmente se accede con un seek a ese bloque, y dentro de este con un gzseek a la secuencia en cuesti´ on. De esta forma para acceder a una secuencia s´ olo se acceder´ an una porci´ on m´ ınima de los datos. Adem´ as la secuencia se obtendr´ a en el stout en formato FASTA o FASTQ, evitando el tener que pasar por conversores para ser usadas por otros programas. 3.5. Simplificaci´ on de los valores de la calidad Como los valores de la calidad suelen estar bastante repetidos en las secuencias que son ´ utiles [129], se ha a˜ nadido un paso de compresi´ on adicional a los valores de calidad (QV, Quality Values). Adem´ as es posible realizar una compresi´ on con p´ erdidas de los valores de calidad debido a que a la hora de usar esos datos no se suelen tener en cuenta todos los valores posibles [26], por lo que es posible hacer una cuantificaci´ on de los valores sin que influya en los resultados de los procesos a los que se sometan esas secuencias. La compresi´ on de los valores de calidad incluye los siguientes pasos, y s´ olo en los opcionales se pierde informaci´ on: 1. En caso de venir la calidad en formato num´ erico, como es habitual cuando el formato original es un fichero FASTA+Quality, se procede a convertir los valores num´ ericos en caracteres, tal y como se usan en el formato FASTQ. 2. Discretizaci´ on (opcional): el rango de valores de calidad (usualmente de 1 a 40 en secuenciadores NGS) se divide en intervalos de longitud personalizable, sustituyendo los valores que se encuentran dentro de cada intervalo por un s´ olo valor representativo del intervalo. La f´ ormula de discretizaci´ on es: QVdiscretizado[i] = trunc(QV [i]/long personalizable)×long personalizable. 3. Filtrado (opcional): debido a que usar el valor de calidad en el postprocesado de secuencias NGS es usualmente impr´ actico [20, 129], s´ olo los valores bajos de la calidad tienen inter´ es para cortar las bases de baja calidad. Suele considerarse buena calidad cuando el valor de calidad es mayor o igual a 20. En la implementaci´ on de FQbin el punto de corte es personalizable. Por tanto para las bases que se consideran de buena calidad su valor de calidad es reemplazado por el valor de corte, y para el resto de las bases el valor se pone a cero. 4. Simplificaci´ on de los valores de calidad repetidos: cuando todos los valores de calidad de las bases de una secuencia tienen el mismo valor se almacenan como un s´ olo valor y esto se indica en la cabecera, en el campo de longitud de la calidad. 108 Cap´ ıtulo 3. Almacenamiento optimizado de datos de secuenciaci´ on gen´ etica Es decir, al leer en la cabecera un 1 en la longitud de la calidad de la secuencia se sabe que se ha usado esta simplificaci´ on, y que todos los valores de calidad son iguales. En el momento de la descompresi´ on, al leer un s´ olo valor de calidad, este se repite tantas veces como indique la longitud de la secuencia. 3.6. Herramientas para manipular el formato FQbin Se han creado diversas herramientas para la manipulaci´ on del formato FQbin. Para crear un nuevo fichero FQbin se usa mk fqbin indic´ andole los ficheros de entrada y de salida. Para recrear el ´ ındice, en caso de que se pierda el fichero de ´ ındices, se usa idx fqbin. Para crear el fichero de hash se usa mk hash, para acceder a una secuencia de forma aleatoria se usa read fqbin. Por ejemplo, para leer todas las secuencias de un fichero FQbin y enviarlas al programa blast para ser comparadas contra una base de datos, en el ejemplo blast database, se puede usar el comando: iterate_fqbin -F file.fqbin | blastn -db blast_database En caso de que otro software no permita el uso de tuber´ ıas, se puede evitar el volcar el fichero FQbin en otro mediante el uso de tuber´ ıas con nombre (named pipes), esto se puede realizar, por ejemplo, como se muestra a continuaci´ on: # crea una tuber´ ıa con nombre (named pipe) mk_fifo nam_pipe1 # bowtie2 usa la tuber´ ıa bowtie2 index nam_pipe1 & # env´ ıa datos a la tuber´ ıa iterate_fqbin file.fqbin > nam_pipe1 # borra la tuber´ ıa rm nam_pipe1 3.7. Tests En el caso de los test de la implementaci´ on en C de nuestra librer´ ıa, se usaron tres clases distintas de secuencias: un conjunto de datos generados con un Illumina, el SRR314795 (21,908,723 reads, 3.9 GB y formato FASTQ), otro con un Roche 454/FLX+, 3.8. Resultados y discusi´ on 109 el SRR073389 (811,509 reads, 0.9 GB y formato FASTQ) y un fichero FASTA conteniendo las secuencias sin valores de calidad de los primeros 10 cromosomas humanos (1.6 GB y formato FASTA: AC 000133.1, AC 000134.1, AC 000135.1, AC 000136.1, AC 000137.1, AC 000138.1, AC 000139.1, AC 000140.1, AC 000141.1, AC 000142.1). FQbin se ejecut´ o comprimiendo los datos originales sin modificar, y se compar´ o con otros algoritmos de compresi´ on general como ZIP,GZIP (un envoltorio de ZLIB), BSC (un compresor de alto rendimiento sin p´ erdidas [49]) y DSRC (el mejor compresor publicado para datos con el formato FASTQ [27]). Posteriormente en otro conjunto de experimentos se a˜ nadi´ o a FQbin el filtrado de los valores de calidad que hemos explicado en la secci´ on 3.5, llamando a estos resultados FQbin+QV filters, como veremos en la siguiente secci´ on. Los test de las secciones 3.8.1 y 3.8.2 se han realizado usando un iMac quad-core a 2.8 GHz con 8 GB de RAM, usando un almacenamiento compartido por samba a trav´ es de una red ethernet de 1 Gbit/s o de 100 Mbit/s. Las medidas de tiempo se tomaron con el comando de Unix time. Para el resto de las secciones de este cap´ ıtulo se usa el supercomputador Picasso, descrito en la secci´ on 1.4.5. 3.8. Resultados y discusi´ on Se demostrar´ a que FQbin es un desarrollo robusto que provee un formato apropiado para el almacenamiento de datos generados con NGS, con un acceso casi instant´ aneo a cualquier secuencia. 3.8.1. Capacidad de compresi´ on Cuando se compara el factor de compresi´ on de FQbin con el que se consigue con otros compresores tanto gen´ ericos como espec´ ıficos, como se muestra en la figura 3.3A, se observa que las mayores diferencias se consiguen en los ficheros de m´ as tama˜ no (Illumina 3.9 GBytes). En particular, los compresores que m´ as comprimen son el BSC, FQbin y el DSRC en todos los casos. En la figura 3.3B representamos por otro lado el tiempo de compresi´ on para los mismos compresores. Podemos observar que BSC a pesar de ofrecer buenos factores de compresi´ on, incrementa notablemente el tiempo con el tama˜ no del fichero, lo que indica que es sensible al tama˜ no de los datos de entrada. DSRC es un compresor competitivo porque consigue mayores factores de compresi´ on en menos tiempo que FQbin, sin embargo no es capaz de manejar ficheros con formato FASTA (de ah´ ı que no haya resultados con el fichero Fasta en la figura). Si consideramos la optimizaci´ on FQbin+ QV filters, podemos observar que en ficheros de gran tama˜ no 110 Cap´ ıtulo 3. Almacenamiento optimizado de datos de secuenciaci´ on gen´ etica como Illumina es el que consigue mayor factor de compresi´ on, lo que confirma la naturaleza repetitiva de los datos de calidad, y las ventajas de realizar un filtrado [129]. Sin embargo, DSRC sigue siendo m´ as r´ apido cuando se compara con esta implementaci´ on de FQbin. Por otro lado, gZip que usa zlib, la misma librer´ ıa de compresi´ on que FQbin, no es mejor que ´ este comprimiendo. El motivo es que FQbin olvida todo lo aprendido al iniciar un nuevo bloque, y esto, debido a la alta aleatoriedad de los datos gen´ eticos, ayuda a comprimir mejor. En cuanto a velocidad, gZip es ligeramente m´ as r´ apido que FQbin debido a que este ´ ultimo tiene la sobrecarga de generar y comprimir los ´ ındices a los datos. En la bibliograf´ ıa encontramos SpeedGene, que dice tener un factor de compresi´ on que va desde 16 a varios cientos [102], pero no se puede comparar con FQbin, ya que SpeedGene comprime bas´ andose en un genoma de referencia para la realizaci´ on de estudios de asociaci´ on de genoma completo, mientras FQbin est´ a enfocado a almacenar los datos originales completos. 3.8.2. Lectura de los ficheros comprimidos Podemos ver en la figura 3.3 que el tama˜ no de un fichero FQbin es entre 6 y 10 veces menor que el de un FASTQ generado con un Illumina, pero su contenido, como se puede ver en la figura 3.2, es m´ as complejo, lo que podr´ ıa llevar a pensar en que leer todos los datos puede ser m´ as lento con el formato FQbin que con el FASTQ original. Usando los mismos ficheros con formato FASTQ que usamos anteriormente (454/FLX+ e Illumina), en la figura 3.4 representamos la ratio entre el tiempo que se tarda en leer todos los datos usando el formato FASTQ y el tiempo que tardamos el leer esos datos a partir del fichero comprimido con FQbin (aceleraci´ on), tanto en redes de 1 Gbit/s como 100MBit/s. FQbin lee el fichero 454/FLX+ completo en 26s en una red de 1 Gbit/s, siendo 2 veces m´ as r´ apido que el acceso directo al fichero FASTQ original, que tard´ o 63 segundos, y mucho m´ as r´ apido que el acceso con DSRC, que tambi´ en se midi´ o y tard´ o un total de 230 segundos. Hay que tener en cuenta que los ficheros comprimidos con otros programas como DSRC o zip, gZip y BSC deben ser convertidos primero en ficheros FASTQ, y a continuaci´ on estos ficheros FASTQ ya pueden ser procesados, por lo que con estos formatos siempre ser´ a m´ as lento el acceso que con un fichero FASTQ directamente. Con tama˜ no de ficheros mayores, como es el caso de Illumina, FQbin consigue mayores aceleraciones cuando se compara con FASTQ. Destaca el caso en el que 4 procesos concurrentes est´ an leyendo el mismo fichero (ver la barras 4 X), y en particular cuando ficheros grandes son le´ ıdos usando redes lentas de 100 Mbit/s. En ese caso, leer los datos de Illumina disminuy´ o el tiempo de lectura desde los 3329 segundos del fichero original 3.8. Resultados y discusi´ on 111 2 4 6 8 10 Compression factor BSC Zip gZip DSRC FQbin FQbin + QV filters 0 300 600 900 1200 1500 454/FLX+ Fasta Illumina Compression time (s) BSC Zip gZip DSRC FQbin FQbin + QV filters A B Figura 3.3: Comparaci´ on entre distintos algoritmos de compresi´ on. A: Factor de compresi´ on respecto al fichero sin comprimir calculado como el tama˜ no del fichero original dividido entre el tama˜ no del fichero comprimido (mayor es mejor). B: Tiempo que se tarda en realizar la compresi´ on (menor es mejor). en formato FASTQ a 1137 segundos en formato FQbin, consiguiendo una aceleraci´ on de 2.9x. En resumen, FQbin no s´ olo ahorra espacio de disco, sino que adem´ as ahorra tiempo en el acceso a los datos, acelerando la carga de las secuencias. Los compresores con los que se compara FQbin en la figura 3.3 requer´ ıan la descompresi´ on del fichero completo para acceder a una secuencia cualquiera, pero DSRC y FQbin permiten acceder a cualquier secuencia sin necesitar la descompresi´ on del fichero completo. La eficiencia en el acceso aleatorio se puede ver comparando el tiempo de acceder a la primera y ´ ultima secuencia del fichero, tal y como se muestra en la Tabla 3.1. Como era de esperar, en FASTQ el acceso secuencial al ´ ultimo elemento tarda m´ as cuanto m´ as grande es el fichero. Gracias al fichero hash que se a˜ nade a FQbin para acelerar los accesos, se obtiene un acceso pr´ acticamente instant´ aneo a cualquier secuencia contenida en FQbin, mejor incluso que con el formato DSRC. Por ejemplo, para el acceso al ´ ultimo elemento, FQbin obtiene una aceleraci´ on de 86x con el fichero 454/FLX+ y de 118 Cap´ ıtulo 3. Almacenamiento optimizado de datos de secuenciaci´ on gen´ etica una implementaci´ on distribuida. Para ello hemos de cambiar primero las estructuras de datos definidas en la figura 3.5, convirti´ endolas en distribuciones de Chapel (que se explicaron en la secci´ on 1.3.1). En la figura 3.7 vemos la declaraci´ on de variables y los arrays distribuidos para esta nueva versi´ on. A continuaci´ on s´ olo hay que cambiar algunos bucles para que la iteraci´ on de cada elemento se ejecute en el locale en el que est´ e contenido. Esto se realiza como se muestra en la figura 3.8. Hay que tener en cuenta que chunkArr es un array que est´ a distribuido entre los locales gracias a la distribuci´ on definida en la figura 3.7. En la figura 3.9 podemos ver los tiempos que tarda la librer´ ıa FQbin en comprimir de forma distribuida el fichero con formato FASTA y calidades que se us´ o como entrada en los experimentos del apartado previo. En el eje yse indica el tiempo de compresi´ on y escritura de los ficheros FQbin, y en el eje xse indica el n´ umero de cores. Se condujeron 8 experimentos diferentes: en cada uno se vari´ o el n´ umero de nodos (locales) usados, entre uno y ocho. En particular, en cada experimento se incrementaba el n´ umero de threads desde 1 hasta alcanzar el n´ umero m´ aximo de cores en ese escenario. Recordemos que en esta arquitectura hay 16 cores por locale, por lo tanto el n´ umero m´ aximo de cores en cada escenario viene dado por n´ umero de locales ×16. En la gr´ afica se representa hasta 80 cores. Inicialmente, los threads se reparten de forma uniforme entre los nodos Figura 3.6: Tiempos y aceleraci´ on de la implementaci´ on Chapel de FQbin para un locale. 3.9. Motivaci´ on para una nueva implementaci´ on en Chapel 119 var Space={0..#parallel_threads}; var chunkDom = Space dmapped new dmap(new Block(boundingBox=Space)); var chunkArr: [chunkDom] chunk; var indexForChunkDom = Space dmapped Block(boundingBox=Space); var indexDom = {0..#seqs_per_stream}; var indexArr: [indexForChunkDom][indexDom] index_el; Figura 3.7: Declaraci´ on en Chapel de las variables usadas en la librer´ ıa FQbin distribuida. // Bucle original forall iin 0..#parallel_threads do ... // Bucle anterior cambiado para que sea distribuido forall Lin Locales { on L { const indices = chunkArr.localSubdomain(); forall iin indices ... Figura 3.8: Cambios a efectuar en los bucles para que pasen a ejecutarse en todos los locales de manera distribuida. solicitados. Por ejemplo, 16 threads pueden corresponder a 1 s´ olo nodo con 16 cores, a 2 nodos con 8 cores cada uno, a 4 nodos con 4 cores, etc. Por ello se han usado distintos colores para distinguir el n´ umero de nodos. Un primer resultado que llama la atenci´ on en la figura 3.9 es que en la implementaci´ on Chapel distribuida de FQbin tiene un comportamiento similar a la implementaci´ on en un s´ olo locale (ver figura 3.6), pero s´ olo en configuraciones de hasta 8 threads por nodo. Como podemos ver en la figura, en configuraciones de m´ as de 8 threads por nodo los tiempos se degradan en cualquier escenario. La raz´ on de esta degradaci´ on est´ a relacionada con la planificaci´ on de los threads que dentro de cada nodo realiza GASNet. Hemos observado que GASNet asigna (pinning) cada thread a un core concreto dentro del nodo. Mientras el n´ umero de threads es menor que 8 asigna correctamente un thread por core. Sin embargo, a partir de 8 comienza a asignar m´ as de 1 thread por core (a pesar de que a´ un queden cores libres) lo que provoca la oversubscription del core y que empiecen a aparecer problemas de contenci´ on en ese core y por lo tanto el desbalanceo 120 Cap´ ıtulo 3. Almacenamiento optimizado de datos de secuenciaci´ on gen´ etica del trabajo entre esos threads y el resto, lo que degrada significativamente los tiempos. Efecto del env´ ıo y recepci´ on de mensajes con GASNet Aunque en principio las transferencias RDMA no usan la CPU ni el sistema operativo de ninguno de los nodos intervinientes, en la pr´ actica hay que preparar las transferencias, y sobre todo hacer un polling sobre la cola de terminaci´ on (completion queue), siendo la intenci´ on la de minimizar la latencia cuando termina una operaci´ on de RDMA. Se podr´ ıa realizar con una interrupci´ on, pero eso incrementar´ ıa la latencia, ya que habr´ ıa que preparar la interrupci´ on, realizar un cambio de contexto, afectando a las cach´ es y habr´ ıa que procesar la interrupci´ on, necesitando otro cambio de contexto. Por lo tanto, el realizar polling para disminuir la latencia de las comunicaciones y de la E/S est´ a documentado en diversos white papers de buenas pr´ acticas, como por ejemplo [126], [125], y [62]. Al implementar GASNet decidieron usar polling en el propio espacio de usuario, evitando el tener que realizar un cambio de contexto al pasar a modo kernel. Pero este polling continuo de la cola de terminaci´ on requiere de un thread espec´ ıfico que ocupa totalmente un core, lo que puede afectar al resto de threads de la aplicaci´ on si alguno de ellos es mapeado en el core que realiza el polling. Este detalle de la implementaci´ on interna de GASNet es el que explica los puntos at´ ıpicos que vemos en la figura 3.9. Figura 3.9: Tiempos de la implementaci´ on distribuida en Chapel de FQbin. Se representan por n´ umero de nodos, entre uno y ocho. En estas pruebas se ha usado GASNet. 3.10. Conclusiones 121 3.10. Conclusiones En ese cap´ ıtulo se ha presentado un formato de datos nuevo, FQbin, para el almacenamiento y acceso eficiente a datos obtenidos por secuenciaci´ on gen´ etica de ´ ultima generaci´ on. Se ha dise˜ nado una librer´ ıa para su f´ acil manipulaci´ on, ofreciendo las siguientes ventajas: 1. Provee el almacenamiento de secuencias y/o valores de calidad y/o extras en el mismo fichero, expandiendo su uso m´ as all´ a de las secuencias de nucle´ otidos. De hecho, la base de datos EuroPineDB [50] usa la librer´ ıa FQbin para recuperar contigs de grandes ficheros ACE. 2. Provee de funciones de recuperaci´ on directa soportando la E/S de datos FASTA, QUAL y FASTQ. 3. Permite que la informaci´ on comprimida sea f´ acilmente transferida por pipeline desde/hacia otros programas, asegurando la compatibilidad con el software actual sin necesidad de recodificar los datos. 4. El decremento en el tama˜ no de los ficheros se puede comparar favorablemente con el que se obtiene en otros algoritmos que representan el estado del arte. 5. Lee la misma cantidad de datos en menos tiempo que otros y provee un acceso casi instant´ aneo a cualquier secuencia por su nombre, independientemente de su posici´ on. Hemos comprobado que gracias a FQbin el acceso a secuencias concretas de forma aleatoria se acelera gracias al uso de ´ ındices y de tablas hash, como tambi´ en se aceleran la copia y manipulaci´ on en general de los ficheros completos. Pero mientras esper´ abamos que el acceso a todos los datos contenidos en un fichero FQbin fuera m´ as r´ apido que al fichero original, esto no es siempre cierto, ya que depende de la relaci´ on entre dos factores, la velocidad de la CPU, que es qui´ en deber´ a descomprimir los datos, y la velocidad del acceso a los datos en el almacenamiento secundario. Tras la definici´ on del formato, la implementaci´ on de la librer´ ıa en C y el estudio de su rendimiento, hemos usado el lenguaje de programaci´ on Chapel para convertir la librer´ ıa secuencial original en una paralela, consiguiendo una buena escalabilidad cuando se ejecuta dentro de un nodo (o locale). Cuando hemos adaptado esta versi´ on para permitir la ejecuci´ on distribuida en varios nodos hemos descubierto ciertas ineficiencias debidas a una incorrecta planificaci´ on interna de threads en la librer´ ıa GASNet (que es la que se encarga de las comunicaciones remotas de los mensajes entre nodos) as´ ı como al efecto del thread de polling que gestiona el env´ ıo y recepci´ on de mensajes en esa librer´ ıa. 122 Cap´ ıtulo 3. Almacenamiento optimizado de datos de secuenciaci´ on gen´ etica Para finalizar queremos comentar que nuestra soluci´ on [59] aparece referenciada en dos patentes, en el a˜ no 2014 la patente n´ umero 8.847.799 [66], y en 2015 la patente 8.976.049 [65], con t´ ıtulo ”Methods and systems for storing sequence read data”. Adem´ as se ha descargado la gema de ruby que contiene la librer´ ıa FQBin [58] m´ as de 2000 veces a d´ ıa de hoy, de las cuales m´ as de 400 descargas pertenecen a la ´ ultima versi´ on. 4Conclusiones En este trabajo hemos presentado dos contribuciones principales. Por un lado, proponemos una librer´ ıa de funciones de E/S de arrays distribuidos en el lenguaje de programaci´ on de alta productividad Chapel as´ ı como el estudio del rendimiento de dicha librer´ ıa cuando ´ esta se configura para explotar distintos grados de paralelismo en el sistema. Por otro lado, dise˜ namos una librer´ ıa para el almacenamiento y acceso a secuencias gen´ eticas que reduce el coste del almacenamiento y aumenta el rendimiento y la fiabilidad en los accesos. 4.1. Librer´ ıa para el almacenamiento distribuido de arrays El primer criterio de dise˜ no a la hora de crear una librer´ ıa de E/S paralela ha sido el de garantizar una interfaz de alta productividad. Tras analizar el estado del arte, hemos encontrado que el principal problema a la hora de proporcionar una interfaz amigable al programador viene derivado de la distribuci´ on de los datos, tanto entre los distintos nodos como dentro de cada nodo. Gracias a la facilidad de definir arrays distribuidos en Chapel, hemos conseguido crear una interfaz sencilla para realizar la E/S, ampliando la interfaz de las distribuciones de Chapel con operaciones de E/S. De esta forma, la informaci´ on relativa a la distribuci´ on de los datos es accesible desde la librer´ ıa de E/S al tiempo que se ocultan los detalles de implementaci´ on (paso de mensajes, gesti´ on de buffers, etc) al programador, lo cual impacta muy positivamente en un incremento de la productividad de nuestra soluci´ on. Otro criterio de dise˜ no ha consistido en mantener una alta flexibilidad en la librer´ ıa y la portabilidad de los datos almacenados en fichero. En concreto, nuestra soluci´ on per123 124 Cap´ ıtulo 4. Conclusiones mite realizar la escritura de un array distribuido en un ´ unico archivo, el cual puede ser le´ ıdo posteriormente a otro array que siga otra distribuci´ on diferente. En otras palabras, nuestra librer´ ıa desacopla el formato del fichero de la distribuci´ on usada en los arrays que almacena. En aras de la portabilidad hemos descartado la alternativa de generar varios ficheros para almacenar los resultados de salida de una aplicaci´ on. Esta posibilidad acelerar´ ıa la escritura en paralelo, pero complica sobremanera la posterior lectura si se usa una distribuci´ on de datos diferente a la usada en escritura, lo que podr´ ıa suceder simplemente cambiando el n´ umero de nodos usados por la aplicaci´ on. Por ´ ultimo, aunque no menos importante, el criterio de dise˜ no relacionado con las prestaciones ha requerido la mayor inversi´ on de tiempo en investigaci´ on y desarrollo. Para implementar eficientemente el interfaz anterior hemos buscado en las distintas fases de una operaci´ on de E/S todos los niveles en los que es posible explotar paralelismo. Hemos encontrado varios niveles que han resultado ser ortogonales y que por tanto se han podido combinar. Las distintas posibilidades se han validado experimentalmente para seleccionar la mejor soluci´ on. Tras realizar un an´ alisis exhaustivo hemos encontrado que el rendimiento ´ optimo se obtiene al paralelizar tanto las escrituras en disco como la agregaci´ on de datos. El solapar la escritura de los datos en disco y la agregaci´ on de los datos desde los clientes aumenta de forma no significativa el rendimiento. Sin embargo, intentar explotar m´ as paralelismo incrementando el n´ umero de agregadores por OST termina impactando negativamente en el rendimiento. Tambi´ en en relaci´ on con las prestaciones de la E/S, hemos comprobado que los overheads necesarios para realizar las comunicaciones son una fuente cada vez m´ as importante de p´ erdida de rendimiento. Para paliar estos overheads hemos propuesto una optimizaci´ on de las operaciones de copia de arrays en Chapel. Esta mejora consiste en agregar autom´ aticamente los datos para explotar mejor el ancho de banda de las comunicaciones lo que ha repercutido muy positivamente en la gesti´ on de los buffers usados en las operaciones de E/S. 4.1.1. Contribuciones Esta l´ ınea ha dado lugar a las siguientes publicaciones: Rafael Larrosa Jim´ enez, Rafael Asenjo Plaza, Angeles Gonz´ alez Navarro, and Emilio L´ opez Zapata. Implementing a chapel library for parallel io. In Actas de las XXI Jornadas de Paralelismo, pages 323–330, 2010. Rafael Larrosa, Rafael Asenjo, Angeles G. Navarro, and Bradford L. Chamberlain. A first implementation of parallel io in chapel for block data distribution. In PARCO, volume 22 of Advances in Parallel Computing, pages 447–454. IOS Press, 2011. 4.2. Librer´ ıa para el almacenamiento optimizado de datos de secuenciaci´ on gen´ etica125 Alberto Sanz, Rafael Asenjo, Juan Lopez, Rafael Larrosa, Angeles Navarro, Vassily Litvinov, Sung-Eun Choi, and Bradford L Chamberlain. Global data re-allocation via communication aggregation in chapel. In Computer Architecture and High Performance Computing (SBAC-PAD), 2012 IEEE 24th International Symposium on, pages 235–242. IEEE, 2012. Adem´ as, varias de las optimizaciones del runtime de Chapel surgidas durante el desarrollo de esta tesis se incorporaron al repositorio de dominio p´ ublico disponible en GitHub (https://github.com/chapel-lang/). En particular, estas mejoras est´ an disponibles desde las versiones 1.8 y 1.9 de Chapel. 4.2. Librer´ ıa para el almacenamiento optimizado de datos de secuenciaci´ on gen´ etica Tras estudiar los formatos existentes de almacenamiento de secuencias gen´ eticas as´ ı como los casos de uso, hemos propuesto un formato que a´ una alta capacidad de compresi´ on y r´ apido acceso aleatorio. Este nuevo formato llamado FQbin es funcionalmente competitivo con el est´ andar de facto, el formato FASTQ, ya que incluye todos los campos de datos opcionales de ´ este (como el campo de calidad o el campo con informaci´ on extra), pero lo aventaja en productividad y en ahorro de espacio en disco. Uno de los puntos m´ as importantes de la librer´ ıa FQbin desde el punto de vista de la productividad es que no hace falta modificar las aplicaciones ya existentes para usar el nuevo formato. Generalmente cuando almacenamos los datos en un formato pero la aplicaci´ on que los necesita est´ a preparada para leer otro formato diferente, la soluci´ on habitual consiste en realizar una conversi´ on de formatos. Esto suele implicar una descompresi´ on del archivo y el consiguiente consumo de recursos, tanto de espacio en almacenamiento secundario como de ancho de banda de E/S y de c´ alculo para la descompresi´ on. Nuestra implementaci´ on evita la creaci´ on del fichero descomprimido en disco mediante una estrategia basada en un stream de datos procesados por bloques y mantenidos temporalmente en memoria. El formato FQbin facilita ese proceso al tiempo que permite un r´ apido acceso aleatorio a las secuencias gen´ eticas. Adem´ as los ficheros con formato FQbin ocupan entre 5 y 10 veces menos espacio que los datos originales. El hecho de que las aplicaciones virtualmente procesen los datos en su versi´ on comprimida (no es necesario descomprimir todo el fichero para que sea procesado) repercute positivamente en varios aspectos ya que se reduce el impacto de la aplicaci´ on sobre la jerarqu´ ıa de memoria, la red de comunicaciones y los propios dispositivos de almacenamiento. 126 Cap´ ıtulo 4. Conclusiones Adem´ as de la implementaci´ on secuencial se ha realizado la implementaci´ on paralela en Chapel de la escritura de ficheros con formato FQbin. La escalabilidad observada es pr´ acticamente lineal cuando se usa memoria compartida. Sin embargo, una decisi´ on de dise˜ no de la capa de comunicaciones GASNet en la que se basa Chapel y que se escapa a nuestro control da lugar a p´ erdidas de escalabilidad en arquitecturas de memoria distribuida. 4.2.1. Contribuciones Esta l´ ınea ha dado lugar a la siguiente publicaci´ on: Dar´ ıo Guerrero-Fern´ andez, Rafael Larrosa, and M. Gonzalo Claros. FQbin, a compatible and optimized format for storing and managing sequence data. In International Work-Conference on Bioinformatics and Biomedical Engineering, IWBBIO 2013, Granada, Spain, March 18-20, 2013. Proceedings, pages 337–344. 2013. Adem´ as, la librer´ ıa FQBin [59] ha sido citada en dos patentes en los EEUU, en 2014 en la patente 8,847,799 [66] y en 2015 en la patente 8,976,049 [65], adem´ as la gema de ruby de la librer´ ıa FQBin, disponible en RubyGems (https://rubygems.org/gems/scbi fqbin/), ha tenido m´ as de 2000 descargas hasta ahora [58]. 4.3. Trabajos futuros Durante el transcurso del presente trabajo han aparecido oportunidades de mejorar las aportaciones propuestas. A continuaci´ on presentamos las que consideramos m´ as interesantes: La tecnolog´ ıa basada en burst buffers [83] ya se integra en algunos sistemas de almacenamiento pioneros y se prev´ e su uso tambi´ en en HPC, como por ejemplo en la plataforma Summit (ver secci´ on A.1.3). Este tipo de buffers asimilan las r´ afagas que se producen en la escritura, que posteriormente van pasando al almacenamiento f´ ısico. De esta forma el usuario percibe una velocidad de escritura superior. Pretendemos estudiar el impacto de esta tecnolog´ ıa en nuestra librer´ ıa de E/S, y en qu´ e medida puede ayudar a mejorar el rendimiento. En relaci´ on con la librer´ ıa FQbin queremos investigar como afectan los accesos de lectura al mismo fichero desde varios nodos de c´ alculo. Analizaremos los distintos puntos donde se puede realizar la descompresi´ on y como se alteran las comunicaciones realizadas dependiendo del lugar donde se produce la descompresi´ on 4.3. Trabajos futuros 127 (servidores de disco, del sistema de fichero, capa burst buffer o nodos de c´ alculo). Es decir, buscaremos el balance ´ optimo entre la carga de las CPUs, del sistema de almacenamiento y de las comunicaciones.