scieee AI-readable full text Open interactive document viewer

New technologies to bridge the gap between High Performance Computing (HPC) and Big Data

Piñeiro Pomar, César Alfredo

Abstract

The unification of HPC and Big Data has received increasing attention in the last years. It is a common belief that exascale computing and Big Data are closely associated since HPC requires processing large-scale data from scientific instruments and simulations. But, at the same time, it was observed that tools and cultures of HPC and Big Data communities differ significantly. One of the most important issues in the path to the convergence is caused by the differences in their software stacks. This thesis will address the research challenge of bridging the gap between Big Data and HPC worlds. With this goal in mind, a set of tools and technologies will be developed and integrated into a new unified Big Data-HPC framework that will allow the execution of scientific multi-language applications on both environments using containers.

Full text

INTERNATIONAL DOCTORAL SCHOOL OF THE USC César Alfredo Piñeiro Pomar PhD Thesis New technologies to bridge the gap between High Performance Computing (HPC) and Big Data Santiago de Compostela, 2022 Doctoral Programme in Information Technology Research TESE DE DOUTORAMENTO NEW TECHNOLOGIES TO BRIDGE THE GAP BETWEEN HIGH PERFORMANCE COMPUTING (HPC) AND BIG DATA César Alfredo Piñeiro Pomar ESCOLA DE DOUTORAMENTO INTERNACIONAL DA UNIVERSIDADE DE SANTIAGO DE COMPOSTELA PROGRAMA DE DOUTORAMENTO EN INVESTIGACIÓN EN TECNOLOXÍAS DA INFORMACIÓN SANTIAGO DE COMPOSTELA 2022 Declaración do autor da tese Don César Alfredo Piñeiro Pomar Título da tese: New technologies to bridge the gap between High Performance Computing (HPC) and Big Data. Presento a miña tese, seguindo o procedemento adecuado ao Regulamento, e declaro que: 1. A tese abarca os resultados da elaboración do meu traballo. 2. De ser o caso, na tese faise referencia ás colaboracións que tivo este traballo. 3. Confirmo que a tese non incorre en ningún tipo de plaxio doutros autores nin de traballos presentados por min para a obtención doutros títulos. 4. A tese é a versión definitiva presentada para a súa defensa e coincide a versión impresa coa presentada en formato electrónico. E comprométome a presentar o Compromiso Documental de Supervisión no caso de que o orixinal non estea na Escola. En Santiago de Compostela, 15 de outubro de 2022 Asdo. César Alfredo Piñeiro Pomar Autorización do director/titor da tese New technologies to bridge the gap between High Performance Computing (HPC) and Big Data Don Juan Carlos Pichel Campos , Profesor Titular da Área de Arquitectura e Tecnoloxía de Computadores da Universidade de Santiago de Compostela INFORMA: Que a presente tese correspóndese co traballo realizado por Don César Alfredo Piñeiro Pomar, baixo a miña dirección/titorización, e autorizo a súa presentación, considerando que reúne os requisitos esixidos no Regulamento de Estudos de Doutoramento da USC, e que como director desta non incorre nas causas de abstención establecidas na Lei 40/2015. De acordo co indicado no Regulamento de Estudos de Doutoramento, declara tamén que a presente tese de doutoramento é idónea para ser defendida en base á modalidade de Monográfica con reproducción de publicaciones, nas que a participación do doutorando foi decisiva para a súa elaboración e que as publicacións se axústan ó Plan de Investigación. En Santiago de Compostela, 15 de outubro de 2022 Asdo. Juan Carlos Pichel Campos Director/a tese Á miña nai... cho debo todo Where we’re going, we don’t need roads Dr. Emmett Lathrop Brown AGRADECEMENTOS En primeiro lugar debo facer unha especial mención, de inmensa gratitude, ao meu titor e director, Dr. Juan Carlos Pichel, polo seu apoio incondicional ao longo de todos estes anos, pola súa xenerosidade e pola paciencia á hora de corrixir os meus textos. Para min sempre será un referente, un espello onde mirarme cada día e pensar que vou ser capaz de acadar que se sinta orgulloso de este seu alumno. En segundo lugar, e no eido persoal, ao meu añorado avó, José Pomar, emigrante con mínimos estudos, traballador incansable, de mentalidade aberta, ávido de coñecemento, que hoxe , onde queira que estea, sinto o se apoio. Dende pequeniño, del aprendín a ”estrebillar”, a procurar solucións, a razoar, a pensar no porqué das cousas. E a tentar ser sempre unha boa persoa. Non podo esquecerme de bos amigos e compañeiros que me acompañaron na carreira, Cristian e Álex e, como non, a amigos incondicionais como Damián, Carolina, Carlos, Cristina e Miryan que sempre estiveron ao meu carón, tanto nas boas coma nas malas. Por último, mais non menos importantes na miña vida persoal, José Luis, Fernando e Chus, que me enriqueceron e axudaron a ser que son. E a pesar da distancia tampouco me quero esquecer as miñas primas Bea e Natalia. This work has received financial support from European Commission RIA - H2020 (HPCEUROPA3 - INFRAIA-2016-1-730897), MICINN, Spain (RTI2018-093336-B-C21, PLEC2021007662), Xunta de Galicia, Spain (ED481A-2019/137, ED431G/08, ED431G-2019/04 and ED431C2018/19) and the European Regional Development Fund (ERDF). característica salientable de Ignis é que está completamente desenvolvido dentro de contedores Docker, o que illa a contorna de execución do sistema físico e evita problemas de dependencias. Finalmente, co obxectivo de facilitar a súa adopción pola comunidade Big Data, a API da ferramenta está inspirada na API de Spark, de forma que os códigos de Ignis son facilmente comprensibles polos usuarios familiarizados con Spark. A avaliación experimental levouse a cabo usando catro tipos diferentes de aplicacións coa intención de ter unha mostra representativa dos algoritmos máis usados en Big Data e na computación científica. En primeiro lugar, considerouse unha aplicación que simula o algoritmo empregado para o minado de Bitcoins, formado por varias operacións map encadeadas. A segunda foi o K-Means, un algoritmo iterativo de machine learning que foi implementado usando o modelo de programación MapReduce. En terceiro lugar, unha operación de ordenación (Sort), que é un dos núcleos computacionais máis importantes en moitas aplicacións paralelas. E por último, o Gradiente Conxugado, un método moi coñecido no ámbito HPC para resolver sistemas de ecuacións lineais. Os resultados obtidos permitíronnos facer unha comparación entre Ignis e Spark. Por exemplo, Ignis foi 2.2×máis rápido que Spark para aplicacións que constan de múltiples operacións tipo map, mostrando así o seu bo rendemento en termos de escalabilidade tanto forte coma débil. Outro exemplo ilustrativo foi o K-Means. Neste caso Ignis foi en promedio entre 1.6×e 2.1× máis rápido que Spark. Resultados de rendemento similares foron observados para o Sort e o Gradiente Conxugado. Por outro lado, en contra da crenza da comunidade de desenvolvedores, demostramos que Python non está soportado de forma nativa por Spark, xa que as transferencias de datos entre a JVM e os procesos externos degradan notablemente o rendemento xeral das aplicacións escritas nesa linguaxe. Como se explicou previamente, fixéronse moitos esforzos no eido da investigación para diminuir a fenda entre as tecnoloxías HPC e Big Data. Non obstante, Ignis segue un camiño diferente na procura dunha converxencia real, debido a que ningún dos traballos anteriores tiñan como obxectivo final a creación dun novo motor de procesamento común para aplicacións Big Data e HPC. A pesares de que Ignis fai importantes contribucións na busca da converxencia, ten algunhas limitacións que lle impiden ser un framework universal para a execución de aplicacións HPC e Big Data. A desvantaxe máis importante está relacionada coa forma en que se realizan as comunicacións. En particular, Ignis está restrinxido ao uso de sockets TCP para a comunicación entre nodos. Isto causaría un importante problema de escalabilidade cando se utiliza un alto número de nodos de computación nos actuais sistemas HPC, xa que Ignis non podería aproveitar as redes máis avanzadas (por exemplo, Infiniband ou Slingshot). No Capitulo 3 introdúcese IgnisHPC como solución as limitacións presentes en Ignis. IgnisHPC herda algunhas características de Ignis, como a tolerancia a fallos, pero redeseñouse xv completamente co obxectivo de mellorar o rendemento, a escalabilidade e a produtividade. IgnisHPC usa MPI como tecnoloxía principal, o que lle permite soportar moitos modelos de comunicacións e arquitecturas de rede. Ademáis, as aplicacións e librarías MPI poden ser executadas directamente en IgnisHPC de forma eficiente. Grazas a isto a maioría das aplicacións científicas HPC non teñen que ser portadas a unha nova API ou modelo de programación. Por exemplo, para executar o resolutor alxebraico AMG en IgnisHPC só é preciso engadir unhas 40 liñas de codigo, cando toda a aplicación suma máis de 65.000. Esta característica a día de hoxe non está presente en ningún outro framework de computación. Sumado a isto, os códigos MPI poden ser combinados con operacións que seguen o modelo MapReduce. A API de IgnisHPC foi extendida para soportar, entre outros, algoritmos para o procesamento de grafos. Os resultados experimentais mostraron os beneficios en termos de rendemento e produtividade. O estudo realizado demostrou que IgnisHPC supera claramente a Spark e Ignis ao considerar aplicacións que representan diferentes tipos de patróns no eido Big Data. En termos de rendemento, IgnisHPC é dende 1.1×a 3.9×máis rápido que Spark, e dende 1.1×a 1.3×máis rápido que Ignis. Pero tamén hai outros beneficios, como os relacionados co consumo de memoria, xa que IgnisHPC permite a procesamento de conxuntos de datos extremadamente grandes (por exemplo, o benchmark Terasort). Tamén comprobouse que executar aplicacións MPI (e híbridas MPI+OpenMP) desde IgnisHPC é fácil e tan eficiente como executalas de forma nativa. En particular, as diferenzas de rendemento entre MPI nativo e IgnisHPC son realmente pequenas, sempre inferiores ao 1,7%. Finalmente, no Capítulo 4 preséntase a aplicación BigSeqKit. Esta supón unha contribución dobre con respecto a IgnisHPC. Por un lado, BigSeqKit foi a primeira aplicación científica completa implementada en IgnisHPC. Trátase dun conxunto de ferramentas paralelas baseadas no software SeqKit para o procesado de arquivos FASTA e FASTQ centrada na escalabilidade e no rendemento. Hai que ter en conta que manipular estes arquivos de forma eficiente é esencial para analizar e interpretar os datos en calquer fluxo de traballo de xenómica. Por outro lado, IgnisHPC foi extendido para soportar a linguaxe de programación Go. Na actualidade non existe ningún framework Big Data-HPC que teña soporte nativo para esta linguaxe. Os resultados experimentais demostraron os beneficios da nosa proposta, logrando aumentos de velocidade entre 11×e 387×con respecto a SeqKit. Temos que ter en conta que SeqKit só pode executarse nun único nodo, de xeito que procesos que requiren horas en SeqKit tardan unicamente uns poucos minutos en BigSeqKit. Polo tanto, podemos concluír que os obxectivos da tese cumpríronse con éxito. As tecnoloxías e solucións desenvolvidas supoñen un importante avance cara á converxencia real dos mundos HPC e Big Data. xvi Contents 1 Introduction 1 1.1 High Performance Computing (HPC) . . . . . . . . . . . . . . . . . . . . . . 2 1.2 BigData...................................... 5 1.3 HPCvs.BigData................................. 7 1.4 Bridging the gap between HPC and Big Data . . . . . . . . . . . . . . . . . . 8 1.5 Objectives..................................... 10 1.6 Methodology ................................... 11 1.7 Publications.................................... 12 1.8 Dissertationstructure............................... 15 2 Towards an efficient and scalable multi-language Big Data framework 17 2.1 Background.................................... 18 2.2 Architecture of the Ignis framework . . . . . . . . . . . . . . . . . . . . . . . 20 2.3 Datastorage.................................... 27 2.4 IgnisAPI ..................................... 28 2.5 Experimentalresults ............................... 30 2.6 Finalremarks ................................... 38 3 Improving the interoperability between HPC and Big Data languages and programming models 39 3.1 Background.................................... 40 3.2 IgnisHPC ..................................... 41 3.3 Programming applications for IgnisHPC . . . . . . . . . . . . . . . . . . . . . 48 3.4 MPIonIgnisHPC................................. 52 3.5 Experimental evaluation . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 57 3.6 Relatedwork ................................... 68 3.7 Finalremarks ................................... 70 4 BigSeqKit: a Big Data approach to process FASTA/FASTQ files at scale 71 4.1 Approach ..................................... 73 Contents 4.2 Performanceresults................................ 77 4.3 Finalremarks ................................... 80 5 Conclusions 81 5.1 Futurework.................................... 84 Bibliography 86 A Getting Started with IgnisHPC 95 A.1 IgnisHPCusage.................................. 95 A.2 BigSeqKit User’s Guide . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 101 List of Figures 105 List of Tables 108 xix 1 Introduction Since the beginning of computing, humans have been constantly looking for ways to create more powerful computers. Computers enabled a fast and systematic way to perform mathematical operations, leaving free time for humans to focus on other significant tasks. As a consequence, computers were forced to evolve with the aim of solving more and more complex operations. But this was far from simple, it has been a long way from performing sums and subtractions to predicting the weather. In this way, single computers were soon overloaded with the amount of work that they had to carry out. However, just like humans, computers can also collaborate to achieve a common goal, and where a single computer is not enough to perform a task, a group of them can provide the necessary resources to do it. This procedure, where several calculations are executed simultaneously, is known as parallel computing and its origin is almost as ancient as the first computers. At the beginning, parallel computing was restricted to single machines, which contained several CPUs that share resources such as memory and I/O devices. Those systems could execute instructions in parallel improving the overall performance. Note that the machine had to be designed for a specific number of processors, whose maximum number was limited by the stateof-the-art technology at that time. It was the appearance of computer networks that showed the true potential of parallel computing. Although the first networks were quite basic, nowadays, thousands or even millions of computers can collaborate to solve a common problem. Those systems are known as clusters of computers or supercomputers. On the other hand, the increasing diversity of resources, networks and machines available caused that programmers could no longer deal with all the hardware and software details to parallelize an application. This resulted in the appearance of standards for the construction of parallel applications, which provided a higher level of abstraction. Consequently, the development of parallel applications was simplified, increasing noticeably the portability and performance. This fact also facilitated the construction of complex applications that deal with important research challenges such as climate modeling, protein folding, or drug discovery. Chapter 1. Introduction Figure 1.1: Frontier is the current No. 1 system in the TOP500 supercomputer list (June 2022). 1.1 HIGH PERFORMANCE COMPUTING (HPC) High Performance Computing (HPC) refers to any form of computing in which the processing density or the size of the problems addressed require a higher capacity than that provided by a standard computer, in order to achieve the objective under given constraints. For this reason, HPC occurs on clusters and supercomputers. The focus of HPC is always on the most complex and demanding computational tasks. Thus, we can find applications in all areas of industry and research that make use of parallel systems. Some significant examples are fluid simulation [94], molecular dynamics [91], astrophysics [24], bioinformatics and medicine [74], weather forecasting [30] or quantum computing [32], among many others. Among all the components of a HPC infrastructure, we can highlight the computing nodes, the network and the storage. The basic computing unit is the CPU, which currently consists of several cores (multicore). Computing nodes could contain one or several of these CPUs. For instance, each node of the Frontier supercomputer, current No. 1 in the TOP5001list, includes one AMD EPYC 7A53s CPU with 64 cores (Figure 1.1). In addition to CPUs, computing nodes can boost their performance including hardware accelerators such as GPUs and FPGAs. In this way, a node of the Frontier supercomputer also has 4 GPUs (AMD Radeon Instinct MI250X). Unfortunately, because accelerators require new restrictions and programming tools (e.g. CUDA [84]), taking advantage of this architectural heterogeneity comes at a significant programming cost. The network also plays an important role in large scale computing systems for the efficient communication when multiple nodes are involved. In fact, it is one of the most critical components because its impact on the performance of applications increases with the system size. 1https://www.top500.org [accessed September, 2022] 2 Chapter 1. Introduction Ethernet is usually the technology of choice for the most basic clusters, but due to its high latency and transfer rate, most performance-demanding clusters implement more complex network technologies (e.g. Infiniband [77], Slingshot [28], TofuD [3], etc.). The Frontier supercomputer, for example, has a 200 Gb/s Slingshot-11 interconnect. Finally, storage also takes a major role in HPC systems as it is the origin and destination of the computed data in most cases. Data must be available to all computing nodes and require the use of Network File System (NFS) or some Distributed File System (DFS). In the first case, data are stored on a single node and the others access them through the network. An NFS-based storage solution is a popular choice for HPC because it is simple to configure and administer. It performs reasonably well for small to medium clusters. However, it has important performance issues when the number of I/O operations is high. On the other hand, a DFS is a file system that is distributed on multiple locations. That is, data are distributed among all the nodes. In this way, a single file may be spread across different servers (striping) or even replicated (mirroring). These parameters can be configured to maximize the reading and/or writing speed depending on the type of computations to be performed in the HPC system. The most popular distributed file system is Lustre [16], which is used in the Frontier supercomputer. Other examples are GlusterFS [41] and IBM Spectrum Scale (formerly GPFS) [81]. Although the increased scale of the largest computing systems has stimulated the search for more scalable algorithms and the use of libraries that provide new levels of scalability, HPC applications have changed little over the years [44]. Adoption of new programming models and languages has been conservative: applications are written in Fortran, C, or C++ using Message Passing Interface (MPI) to support inter-node parallel execution and OpenMP or other alternatives to exploit intra-node parallelism. In particular, MPI is the most widely used and dominant programming model in HPC. In MPI, processes make explicit calls to library routines defined by the MPI standard to communicate data between two or more processes. These routines include both point-to-point (two party) and collective (many party) communication. Before MPI, there were many competing message passing libraries, where each vendor implemented their own library for their hardware. Different libraries made the application development very hard and differences between their APIs completely remove portability. So MPI was an attempt to define common standard messaging interfaces for all vendors. MPI Standard is defined by the MPI Forum and, since its creation, it has been updated to the present day. MPI has remained stable since its first release, and new versions only fix bugs and extend the existing functionality. The most important updates have added support for parallel I/O, asynchronous messages, dynamic process creation and recently, with the release of version 4.0, support for long messages (>2GB). Among the different MPI implementations, the most successful ones are MPICH [72] and 3 Chapter 1. Introduction Open-MPI [75]. In particular, MPICH is a more mature, efficient and robust project, but even so, both are implementing the latest version of the MPI standard. In particular, MPICH 4.1 already implements most of the MPI-4 standard while Open-MPI 5.0 is still in its early stages. Other libraries such as ScaMPI and MPI/Pro that could have competed in their time, can be considered obsolete and are no longer available. The standard is only defined for C and Fortran, so any other language will be dependent on the library that implements it. For example, MPICH has support for C++ and Open-MPI for Java. However, they do not implement the functionality in these languages, they call the C API (bindings). This is the most common way MPI is exported to other programming languages. A popular binding is mpi4py [26], a library to provide Python bindings for MPI. The library package follows the MPI specification and is inspired by the object-oriented interface of MPI-2 (MPI standard for C++ that was discontinued and removed in MPI-3). The library is compatible with MPICH and Open-MPI, and although it only supports the MPI-2 standard, it can be used with modern versions of both. It also offers an additional layer with the MPI functionality adapted to native Python types. This layer uses a serializer that converts objects into bytes so that they can be encoded with the standard functions. Additionally, other libraries such as Numpy, which are also implemented in C, can efficiently make use of mpi4py. The overall performance is significantly lower than C but, in terms of scalability, it offers a satisfactory performance. Java has also been a language with multiple MPI bindings. The first was javampi2, an objectoriented Java interface to the standard MPI-1.1. The library used Java Native Interface (JNI) for the invocation of native C methods and although functional, it was far from efficient. The purpose of the library was to create a Java MPI development environment, but it was abandoned in 2003. Although the project disappeared, the API became the de-facto standard for MPI development in Java. MPJ-express [88] was a library that implemented this standard. MPJ-express was a native implementation of the MPI standard in pure Java. The reference manual only shows the implementation of the most known MPI primitives. However, the library has not been updated since 2014 and seems to have been abandoned. Currently, the only way to use MPI in Java is Open-MPI, which includes support for this language in its latest version and, like javampi, uses JNI. Finally, some programming languages with a relatively short lifetime have already available libraries for using MPI. They are usually small software projects created by the community, but sufficient to develop common applications. This is the case, for example, of Julia, which ended its development in 2020. One of its objectives is the implementation of parallel applications, so its developers have included an MPI library [17]. 2https://code.google.com/archive/p/javampi/ [accessed September, 2022] 4 Chapter 1. Introduction 1.2 BIG DATA Nowadays we are living in the Big Data era, which demands processing in a efficient way huge amounts of data. In the past, clusters were reserved for HPC tasks, but with the arrival of Big Data it became necessary to design new frameworks and technologies for the execution on these systems. These data come from all type of sources: sensors used to obtain information on the climate, publications in social networks, blogs, digital images and video, etc. One of the characteristics of this amount of information is the fact that, in many cases, is not structured. As a consequence, there is not a common representation for the data, so classic programming models and computing engines have problems to deal with that limitation. Google tried to overcome that issue introducing MapReduce [29], a new programming model for processing and generating large data sets on a huge number of computing nodes. Since then, the MapReduce paradigm has been closely linked to Big Data processing. MapReduce simplifies the computation process by dividing the work into three main phases: map,shuffle and reduce. The input and output of a MapReduce computation is a list of keyvalue pairs. Users only need to focus on implementing map and reduce functions, while shuffle is performed automatically between phases. In the map phase, map workers take as input a list of key-value pairs and generate a set of intermediate output key-value pairs, which are stored in the intermediate storage (i.e., files or in-memory buffers). The reduce function processes each intermediate key and its associated list of values to produce a final dataset of key-value pairs. This way, map tasks achieve data parallelism, while reducing tasks perform parallel reduction. MapReduce can handle large jobs, such as sorting a petabyte of data in a few hours, as long as enough servers are available. Furthermore, MapReduce is not only focused on parallelism, it offers the possibility to recover from failures during the data processing. If a mapper or a reducer fail, the job can be rescheduled and restarted from a previous phase. Apache Hadoop [96] was the first open-source implementation of the MapReduce programming model. It was widely adopted by both industry and academia, thanks to that simple yet powerful programming model that hides the complexity of parallel task execution and faulttolerance from the users. Hadoop consists of three main layers: a data storage layer (HDFS), a resource manager layer (YARN [93]), and a data processing layer (Hadoop MapReduce Framework). HDFS is a block-oriented distributed file system based on the idea that the most efficient data processing pattern is a write-once, read-many-times pattern. For this reason, Hadoop shows good performance with embarrassingly parallel applications requiring a single MapReduce execution (assuming intermediate results between map and reduce phases are not huge), and even for applications requiring a small number of sequential MapReduce executions. However, most applications do not fit the Hadoop model and require a more general data orchestration. Apache Spark [103] was designed to overcome some of the Hadoop limitations, especially when considering iterative jobs. Spark only replaces the MapReduce Framework, while the 5 Chapter 1. Introduction – Objectives: The first step, as specified in the previous section, involves specifying the objectives to be met. The objectives will determine the tools and solutions that will be integrated into the new Big Data-HPC framework. – A review and analysis of the state-of-the-art research on new technologies and algorithms that can replace and improve the current solutions used by the existing frameworks, as detailed in Chapter 1. This process includes a study of the current design and the identification of strengths and weaknesses. The weaknesses will determine the scientific challenges to be addressed. – Design and implementation of a new computing engine from scratch. Different technologies and solutions will be developed and analyzed. Performance, scalability and productivity will be the most important metrics in order to guide this process. The designs and technologies considered in each stage will be explained in detail in Chapters 2 and 3. – Validation: Once the proposed solutions have been implemented, several HPC and Big Data scientific applications from different areas of knowledge will be used as case studies. The validation process consists of a performance and usability comparison with respect to the current state of the art Big Data frameworks. Chapters 2 and 3 include the corresponding experimental results. On the other hand, in Chapter 4 we validate our proposal using a real bioinformatics application. 1.7 PUBLICATIONS The research developed in this thesis has led to the following contributions: Journals: César Piñeiroa, Rodrigo Martínez-Castañoaand Juan C. Pichela.Ignis: an Efficient and Scalable Multi-language Big Data Framework. Future Generation Computer Systems, Volume 105, Pages 705-716, 2020. ISSN: 0167-739X. The publication is available at https://doi.org/10.1016/j. future.2019.12.052. aCentro Singular de Investigación en Tecnoloxías Intelixentes, Universidade de Santiago de Compostela. –Quality indicators: Journal Impact Factor: 7.187 (JCR 2020), 1.260 (SJR 2020). Journal ranked in JCR 2020, in Computer science, Theory & Methods (Q1, First decile, 7/110). –PhD candidate contribution: Research conceptualization, design and implementation of the framework and test, evaluation and verification, and partially manuscript writing. –Reproduction rights: Elsevier allows its inclusion as part of a doctoral thesis without express permission (see www.elsevier.com/about/policies/copyright#Author-rights or Figure 1.3). 12 Chapter 1. Introduction Figure 1.3: Licensing information of the previous Elsevier publication (reference [79]). César Piñeiroaand Juan C. Pichela.A Unified Framework to Improve the Interoperability between HPC and Big Data Languages and Programming Models. Future Generation Computer Systems, Volume 134, Pages 123-139, 2022. ISSN: 0167-739X. The publication is available at https:// doi.org/10.1016/j.future.2022.04.002. aCentro Singular de Investigación en Tecnoloxías Intelixentes, Universidade de Santiago de Compostela. –Quality indicators: Journal Impact Factor: 7.307 (JCR 2021), 2.233 (SJR 2021). Journal ranked in JCR 2021, in Computer science, Theory & Methods (Q1, First decile, 10/109). –PhD candidate contribution: Research conceptualization, design and implementation of the framework and test, evaluation and verification, and partially manuscript writing. –Reproduction rights: : This article has been published as Open Access (see Figure 1.4). In any case, Elsevier allows its inclusion as part of a doctoral thesis without express permission (see www.elsevier.com/about/policies/copyright#Author-rights). Figure 1.4: The previous paper, reference [80], has been published as Open Access. 13 Chapter 1. Introduction Conferences: César Piñeiroa, Rodrigo Martínez-Castañoaand Juan C. Pichela.Towards a Big Data Multilanguage Framework using Docker Containers. Proceedings of the XXIX Jornadas de Paralelismo – SARTECO, 2018. ISBN: 978-84-09-04334-7. aCentro Singular de Investigación en Tecnoloxías Intelixentes, Universidade de Santiago de Compostela. –PhD candidate contribution: Research conceptualization, design and implementation of the framework and test, evaluation and verification, and partially manuscript writing. –Reproduction rights: This research was a position paper and only contains preliminary ideas used afterwards in [79]. Contents of this paper were not included in this document. César Piñeiroa.On the road to a unified Big Data and HPC framework. Proceedings of the 35th IEEE International Parallel and Distributed Processing Symposium (IPDPS), 2021. ISBN: 978-1-66543577-2. The publication is available at https://doi.org/10.1109/IPDPSW52791.2021.00160. aCentro Singular de Investigación en Tecnoloxías Intelixentes, Universidade de Santiago de Compostela. –Quality indicators: Conference ranked as Class 2 in GII-GRIN-SCIE (GGS) Conference Rating. –PhD candidate contribution: Research conceptualization, design and implementation of the framework and test, evaluation and verification, and partially manuscript writing. –Reproduction rights:: The contents of this poster are a preliminary introduction to [80], so only contents of that paper were included in the thesis. Other journal publications not related to this thesis: Carlos Eiras-Francoa, David Martínez-Regoa, Leslie Kanthanb, César Piñeiroc, Antonio Bahamonded, Bertha Guijarro-Berdiñasa, and Amparo Alonso-Betanzosa.Fast Distributed kNN Graph Construction Using Auto-tuned Locality-sensitive Hashing. ACM Transactions on Intelligent Systems and Technology, Volume 11, Issue 6, 2020. ISSN: 2157-6904. The publication is available at https://doi.org/10.1145/3408889. aResearch Center on Information and Communication Technologies (CITIC), Universidade da Coruña bDepartment of Maths and Computer Science, University College London cCentro Singular de Investigación en Tecnoloxías Intelixentes, Universidade de Santiago de Compostela dDepartment of Computer Science, Universidad de Oviedo –Quality indicators: Journal Impact Factor: 4.654 (JCR 2020). Journal ranked in JCR 2020, in Computer Science, Information Systems (Q1, 36/161). 14 Chapter 1. Introduction César Piñeiroa, José M. Abuínaand Juan C. Pichela.VeryFastTree: Speeding up the Estimation of Phylogenies for Large Alignments through Parallelization and Vectorization Strategies. Bioinformatics, Volume 36, pages 4658-4659, 2020. ISSN: 1367-4803. The publication is available at https://doi.org/10.1093/bioinformatics/btaa582. aCentro Singular de Investigación en Tecnoloxías Intelixentes, Universidade de Santiago de Compostela. –Quality indicators: Journal Impact Factor: 6.937 (JCR 2020). Journal ranked in JCR 2021, in Mathematical & Computational Biology (Q1, First decile, 3/58). 1.8 DISSERTATION STRUCTURE The structure of this PhD dissertation is organized in five chapters and one appendix. This document covers in detail the state-of-the-art, the research challenges and objectives of the thesis, the design and development of the proposed solutions, a discussion and analysis of the experimental results, and finally the conclusions derived from this work. Specifically, we structured this document as follows: – An introduction about Big Data and HPC together with a review of the state-of-the-art is presented in Chapter 1. In particular, we perform a comparison between HPC and Big Data ecosystems, highlighting the reasons that difficult their convergence. The scientific challenges and objectives of this PhD dissertation are also detailed. – Chapter 2 is focused on the development from scratch of a new Big Data platform called Ignis. The necessity of this new framework is motivated, and its architecture details and API are described. A performance comparison with Spark is provided running applications that represent the most typical algorithmic patterns in Big Data and scientific computing. In this chapter we have focused on objectives O1 and O4. – Chapter 3 presents IgnisHPC, a new computing framework that removes some of the important limitations and performance issues of Ignis. The chapter explains in detail how MPI is used as backbone technology in IgnisHPC, which has two important consequences. First, the framework supports many different communication models and network architectures, covering the vast majority of Big Data and/or HPC clusters. And second, MPI applications and libraries can be directly executed in IgnisHPC. The chapter ends with a thorough experimental evaluation using several Big Data and HPC (MPI and MPI+OpenMP) applications. Note that IgnisHPC fulfills all the research objectives of the thesis. 15 Chapter 1. Introduction – Chapter 4 describes BigSeqkit, a bioinformatics application implemented using the IgnisHPC platform. Performance results in terms of execution time and scalability are presented. It is important to highlight that this chapter complements objective O1 because a new programming language was included in IgnisHPC to develop BigSeqkit. – The concluding remarks of this thesis are detailed in Chapter 5. Moreover, this chapter includes a discussion about future lines of work. – Finally, Appendix A shows how to install, configure and execute applications using IgnisHPC. In addition, a BigSeqKit user’s guide is also presented. 16 2 Towards an efficient and scalable multi-language Big Data framework In the Introduction of the thesis, we described the current state of the art in Big Data frameworks. Most of them, e.g., Apache Hadoop and Apache Spark, only support natively JVM (Java Virtual Machine) languages. However, developers build applications in the programming languages that best suit their needs. For instance, most of the existent scientific applications are developed in languages like C, C++ or Fortran. In that case, it is necessary to port the source codes, which requires a huge effort, use source-to-source compilers [78] or take advantage of the mechanisms provided by the frameworks to call external processes based on system pipes with the corresponding degradation in the performance [31]. Note that in the latter case, the main code that performs the calls should be anyway implemented in Python, Java or Scala. In this chapter, we introduce Ignis, a new framework that allows users to execute their applications written in multiple programming languages without additional overhead. Our framework uses a multi-language RPC (Remote Procedure Call) approach to create a native executor for each language in order to avoid data transfers. In this way, data is handled by the executor in the most efficient way for each programming language. Ignis supports Python, Java, C and C++, but thanks to its modular architecture, adding support for new languages is a straightforward process. To facilitate its adoption from the Big Data community, the Ignis API was inspired by the Spark API in such a way that Ignis codes are easily understandable by users who are familiar with Spark. As explained in Section 1.4, several efforts have been done to bridge the gap between HPC and Big Data technologies. However, Ignis follows a different path in order to reach this convergence since none of the previous works have as their final goal searching for a unique processing engine for Big Data and HPC applications. The remainder of the chapter is organized as follows. Section 2.1 gives some background on several technologies required by Ignis. Section 2.2 describes the architecture and the modules of Ignis. The different ways to storage data in Ignis are explained in Section 2.3. Section 2.4 details the Ignis API. Section 2.5 shows the experimental results. Finally, the final remarks about our proposal are explained. Chapter 2. Towards an efficient and scalable multi-language Big Data framework The contents of this chapter were extracted from the following publication: César Piñeiroa, Rodrigo Martínez-Castañoaand Juan C. Pichela.Ignis: an efficient and scalable multi-language Big Data framework. Future Generation Computer Systems, Volume 105, Pages 705-716, 2020. ISSN: 0167-739X. The publication is available at https://doi.org/10.1016/ j.future.2019.12.052. aCentro Singular de Investigación en Tecnoloxías Intelixentes, Universidade de Santiago de Compostela, Spain. 2.1 BACKGROUND 2.1.1 Big Data frameworks It is well known that Hadoop and Spark are the de facto standards for Big Data processing. Jobs on both frameworks are composed by a driver and a number of executors. The driver is a high-level process that controls the workflow, and executors are a set of independent processes distributed on a cluster that run the work in parallel. The vast majority of the applications developed for these frameworks are written in Java, Scala or Python. However, there are some situations where it is not reasonable to port an application to one of the previous languages (e.g., performance issues). Note that Python is a non-JVM language, so it is not natively supported by Spark. However, Spark takes advantage of Jython [51], which allows a driver implemented in Python to be executed within the JVM. This causes a misunderstanding in the Spark users community. By contrast, executors are directly executed with the available Python interpreter. Hadoop and Spark use system pipes to share data outside of the Java Virtual Machine in order to run non-JVM codes, which introduces an additional overhead that negatively affects performance. Both Hadoop and Spark deal with non-JVM codes in a similar way. Next the process followed by Spark is explained: – The user application (non-JVM code) must read data from the standard input and write results to the standard output. Therefore, input and output must be represented in a string format. – An executor containing the input data is created inside a JVM. – Each executor launches a subprocess with the user application, connecting the JVM to the subprocess using pipes. – The executor converts each object stored inside an RDD to its string representation, writing the result to the standard input of the subprocess. At the same time, another thread is reading from the standard output of the subprocess to generate the resulting RDD with the output data. 18 Chapter 2. Towards an efficient and scalable multi-language Big Data framework As we mentioned, the above process requires that input and output data must be represented as a string. As a consequence, the user application is responsible to parse this string data to a native format supported by the considered non-JVM language. Note that this process becomes complicated when more complex data structures such as trees or maps are considered. Another important drawback is that non-JVM processes are not allowed to access Spark and Hadoop functions such as the context. 2.1.2 Apache Thrift Apache Thrift [10] is an RPC (Remote Procedure Call) system whose main functionality is the invocation of remote methods between different programming languages. Apache Thrift has its own IDL (Interface Definition Language) to define multiple services. In the first place, each service exports series of functions with parameters, returns and exceptions. Second, Thrift generates the corresponding skeleton interface and stub class for the user selected language. A client uses the stab to make remote calls and the server defines the methods implementing the skeleton. Thrift is composed of a set of protocols and transports. Protocols define how data types are serialized, while transports indicate the medium through which the data is sent. There is also the possibility of use intermediate transports such as Zlib [56] to apply compression in streaming when data is sent. Spark and Ignis use Thrift for the communication between modules, but Ignis uses a modified version to add data transfer without defining an IDL. 2.1.3 Docker Docker [68] containers provide the benefits of virtualization (isolation, flexibility, portability, agility, etc.) without penalizing the I/O performance considerably. Docker makes use of resource isolation characteristics of the Linux kernel, so independent containers can be executed on the same host using different assigned resources without interfering among them. Containers supply a virtual environment with their own space of processes and networks. The containers are built with stacked layers. When a container is in execution, a new writable layer is created over a set of read-only layers which define a Docker image. The Docker images are always built from a base image, usually a root filesystem coming from a GNU/Linux distribution. These images can be easily distributed via the official registry, with our own registry or with tarballs. Images can be built with a custom scripting language (dockerfiles) or by saving the state of a running container. Big data processing engines such as Spark have been successfully integrated in a Docker environment [58, 100], showing that performance differences between non-containerized and 19 Chapter 2. Towards an efficient and scalable multi-language Big Data framework Driver Backend Manager Executor 1 Executor n ... Driver container Executor containers Docker Resource Manager Figure 2.1: Scheme of the Ignis architecture. containerized versions are small. However, while Docker is widely used in the industry, its adoption in the HPC world is not very common. There are important efforts in the research community to deal with some of its limitations. For example, authors in [12] developed a secure way of running Docker containers on a HPC environment avoiding the privilege scalation problems. To improve the bandwidth and throughput of HPC jobs when using containers, other researchers [22] propose a method to allow Docker containers to take advantage of a Infiniband network when running HPC applications. Finally, in a recent work [82], the authors demonstrate how Docker containers can be integrated with HPC environments and run MPI applications with cloud-enabled schedulers. 2.2 ARCHITECTURE OF THE IGNIS FRAMEWORK Ignis is divided into four independent main modules which run inside Docker containers: Backend, Manager, Driver and Executor (see Figure 2.1). They are coded in different languages, using Apache Thrift for the inter-module communications. We can summarize the interactions and main goals of the different Ignis modules as follows. The Docker Resource Manager is responsible of launching and destroying containers and also of assigning the required resources to them. Since it can be used outside of the context of Ignis, the interaction is performed through an HTTP API instead. The Driver and the Backend share the same container. Users access all the available features of Ignis through a user API defined by the Driver. Note that this module is only an interface to the Backend, where services that define the logic of the API operations are specified. The Backend is responsible of interacting with the Docker Resource Manager to build a cluster of containers according to the instructions specified in the driver user code. This module also handles the distribution and exchange of data among executors through the Manager. In case some data is lost due to a failure of a cluster node or some of the executors, the Backend is able to recompute the corresponding portion of data without requiring a costly replication. Finally, the Executor Module implements for the supported programming languages the 20 Chapter 2. Towards an efficient and scalable multi-language Big Data framework Slave API Consul Register / deregister node Services health check Master API Register / deregister service Check resources Docker Launch / destroy container Figure 2.2: Docker Resource Manager (Ancoris) architecture. operations defined by the Backend. Next, a more detailed description of each module is provided. 2.2.1 Docker resource manager The Docker Resource Manager, named Ancoris, executes itself inside Docker containers and is composed of two types of instances: masters and slaves. Master instances manage the available resources in a cluster and they are responsible of launching the containers with the assigned resources through the slave interfaces. Slaves expose the resources of their host when they are deployed. Both client requests and internal calls from masters to slaves are done through HTTP APIs. The Docker Resource Manager uses Consul1, which is a distributed, highly available system that provides a framework for discovering and configuring services within a cluster. Among its main functionalities are service discovery (find new providers of a given service), health checking (check the status of the registered services) and hierarchical key/value storage. The basic communications between Consul, masters and slaves are illustrated in Figure 2.2. When a slave is initialized, the configured resources of that machine are registered in the key-value store provided by Consul. The resources are defined in a granular way so a client can request different types of devices of the same family (e.g., hard disks, GPUs) and even choose a particular device. One important attribute is the CPU normalizer, whose goal is to represent the relative performance of a physical/hyperthread core on a cluster of heterogeneous nodes. For instance, a less powerful CPU could specify a normalizer factor of 1.0, whereas another node with a CPU twice as powerful would specify 2.0. In this way, the second node would double the number of virtual cores. 1https://www.consul.io [accessed September, 2022] 21 Chapter 2. Towards an efficient and scalable multi-language Big Data framework –In-Memory: This is the Ignis default option and provides the fastest performance since all data is stored in memory. –Raw memory storage: A serialized representation in memory that allows to remove extra space introduced by objects at the expense of an additional overhead. Serialization protocols are the same used for data exchange between executors. The memory buffer is compressed by Zlib [56], which has 9 compression levels that can be changed when defining the properties in the driver code. By default, level 6 is applied. The disadvantage of serializing data and applying compression is that elements are not indexed and must be accessed sequentially. –Index Raw memory storage: This is a special Python storage method to overcome the Global Interpreter Lock (GIL) limitations, which causes only one thread to execute at a time [76]. This type of storage uses a binary representation, uncompressed, with a table to index the data elements. Data and the index table are stored in shared memory, so they can be accessed by multiple Python processes. This storage replaces the default in-memory storage when using a Python worker and the number of cores per container is greater than one (see the Wordcount driver code of Figure 2.5). –Raw Disk storage: This is a variation of the raw memory method where the buffer uses a file mapping approach. In this way, while a portion of the data is in memory, the remaining is stored in disk. This storage option allows to work with large amounts of data that can not be completely stored in memory. Data in Ignis is ephemeral. It means that when data previously computed is used as input of a new task, it is discarded after use. But, what happens when multiple tasks have the same input data?. In this case, the best solution in terms of efficiency is not discarding the data to avoid extra computations for each task. To deal with this issue, Ignis provides two mechanisms to alter the persistence of data: cache and persist. In particular, the cache method hints that data should be kept in memory after the first time is computed, because it will be reused. The persist method allows users to choose a different storage option for the data: in-memory, raw memory storage or raw disk storage. When data is no longer needed, it must be explicitly removed by users with uncache or unpersist indistinctly. 2.4 IGNIS API To use Ignis, developers should write a driver program that implements their application at high-level (see the example of Figure 2.5). As we explained in Section 2.2.5, some of the driver functions such as map require another one to execute. In that case it is also necessary to implement those functions. We must highlight that although the Ignis code uses a sequential notation, 28 Chapter 2. Towards an efficient and scalable multi-language Big Data framework operations on data are performed in parallel. In order to facilitate the adoption from the Big Data community, the Ignis API was inspired by the Spark API in such a way that Ignis codes are easily understandable by users who are familiar with Spark. Next we provide details about the current functions supported by Ignis: –Managing files: Ignis is able to load and save data from/to any location in a (distributed) file system. Currently data can only be read from text files using readFile. Data can be saved to a file in plain text (saveAsFile) or JSON format (saveAsJSON). It is possible to save distributed data into one single file or several files (one file per executor that owns a portion of the data). –Map functions: The common characteristic to functions belonging to this type is that they apply the same function to each element in the data. As a result of the transformation, the output could be of different size with respect to the input. The available functions are: map, flatMap,filter,keyBy and values. The first two are very similar. While map applies a one-to-one transformation to the input data, flatMap generates an arbitrary number of results. filter is a map operation that pick elements from the input data matching a predicate. keyBy returns a (key,value)pair after applying a function to the value argument with the aim of calculating the key.values does not need an extra function, it only returns value from a (key,value)pair. –Reduce functions: The reduction method aggregates all the elements in the input data using a function. aggregation is a sort of reduction where the type of the input and output data is different. Two functions are necessary, the first one is applied to each element in a data partition, and the second one combines the partial results obtained for each partition. reduceByKey and aggregateByKey are variations where the operation is performed only among elements with the same key in such a way that the final result is a set of unique pairs with values calculated using reduce or aggregate operations, respectively. –Sort functions: In order to sort elements Ignis provides two functions: sort and sortBy. The first method uses the natural order and does not need any additional function. sortBy allows to use a user-defined function to specify the order of the elements. If the result of applying that function to two elements is true, then the first element should precede the second one. Both methods support ascending and descending order. –Shuffle functions: The shuffle method balances the number of elements to be processed by each executor, keeping the same order. This operation is useful to preserve performance after an operation that greatly unbalances the data. It should be invoked by the user. The function importData, which allows different workers to share data, performs an internal shuffle operation if the number of executors for each worker is different. 29 Chapter 2. Towards an efficient and scalable multi-language Big Data framework –Other functions: Ignis implements several operations that return a value to the driver code, but they do not modify or generate new stored data. Spark refers to this type of operations as actions. In particular, Ignis supports count,take,takeSample and collect. The most basic operation is count that returns the number of elements of a stored data collection. collect returns a collection with all the elements stored in the executors of a task. take applies a collect operation but obtains only the first nelements, where nis chosen by the user. takeSample returns a random sample of nelements from the distributed data, with or without replacement. Finally, parallelize distributes the elements of a collection among the executors to form a distributed dataset. In this case new stored data is created. As we mentioned, the above functions exchange data with the driver. With the aim of improving overall performance, Ignis also implements optimized versions for scenarios where data to be exchanged is of small size. 2.5 EXPERIMENTAL RESULTS In this section we evaluate Ignis using several applications in terms of performance, scalability and fault tolerance. A comparison with Spark is also provided. 2.5.1 Hardware Platform and Software The experiments shown in this section were carried out on an 8-node cluster, where each node consists of: – CPU: 2 x Intel Xeon E5-2630v4 (2.2Ghz, 10 cores) – Memory: 384 GB of RAM – Storage: 8 x 4TB 7.2k SATA – Network: 2 x 10GbE It is a Linux cluster running CentOS 7 (kernel 3.10.0), Docker 18.09.1-ce and Spark 2.2.0 (with YARN [93] as cluster manager). GlusterFS on XFS was used by Ignis as distributed file system. In order to illustrate the benefits of our proposal we have considered four applications: Minebench2,K-Means,Sort and Conjugate Gradient. The first three represent different types of application patterns for which Spark is considered the best performing Big Data framework with respect to other approaches such as Hadoop. For instance, Minebench can be considered a chain of map operations, while K-Means uses an iterative MapReduce model. For completeness we have also analyzed the Conjugate Gradient, which is one of the most relevant iterative solvers 2Do not confuse with the data mining benchmark suite NU-MineBench. 30 Chapter 2. Towards an efficient and scalable multi-language Big Data framework for systems of linear equations in the HPC world. Using these applications we will compare the performance of Ignis and Spark. On the one hand, we will demonstrate that, while Python is natively supported by Ignis, an important overhead is caused by using system pipes for data transferring among executors in Spark. As a consequence, Spark degrades both performance and scalability since Python is treated as an external language. On the other hand, we will show the benefits of our approach when running applications coded using several programming languages. Instead of considering only speedup to assess the scalability, we will consider a better alternative based on plotting the raw execution time when using different number of cores on a cluster [46]: – Strong scaling: for a fixed problem, a straight line with slope -1 indicates good scalability, whereas any upward curvature away from that line indicates limited scalability. – Weak scaling: for a sequence of problems with a fixed amount of work per core, a horizontal straight line indicates good scalability, whereas any upward trend of that line indicates limited scalability. 2.5.2 Minebench Minebench3performs the calculation of SHA-256 hashes imitating the Proof-of-Work algorithm used in the Bitcoin protocol [73]. This algorithm has two phases which are implemented using two chained map operations. The first map is data-intensive, while the second is a computeintensive task. In particular, in the first stage a set of Bitcoin transactions are grouped together forming a block proposal. A binary Merkle tree [69] is calculated for those transactions and its Merkle root hash is added to a block header in addition to other attributes such as protocol version, hash of the previous block header, the current timestamp and the difficulty, through which is calculated the network target. The network target determines the threshold under which the hash of the proposed block header must be considered valid. The second stage calculates the hash of the block header iteratively while the condition is not met. In order to obtain different results, another attribute is present in the block header: the nonce. It is an integer that is changed in every iteration so the resulting hash also changes with it. Two different implementations of the application were considered in the tests. The first one was programmed using only Python. In the second version, two different programming languages were used: Python and C++. In this case, the first phase of the application uses Python since it requires handling data, while the compute-intensive tasks use C++ to achieve the best results in terms of performance. 3Publicly available at: https://github.com/brunneis/minebench [accessed September, 2022] 31 Chapter 2. Towards an efficient and scalable multi-language Big Data framework 1 10 100 101 102 103 Cores Execution Time (seconds) Spark Ignis Spark (least−squares fit) Ignis (least−squares fit) Ideal Scalability (a) Python, strong scaling 1 10 100 100 101 102 103 Cores Execution Time (seconds) Spark Ignis Spark (least−squares fit) Ignis (least−squares fit) Ideal Scalability (b) Python & C++, strong scaling 1 10 100 103 104 Cores Execution Time (seconds) Spark Ignis Spark (least−squares fit) Ignis (least−squares fit) Ideal Scalability (c) Python, weak scaling 1 10 100 102 103 104 Cores Execution Time (seconds) Spark Ignis Spark (least−squares fit) Ignis (least−squares fit) Ideal Scalability (d) Python & C++, weak scaling Figure 2.7: Study of the scalability of Ignis and Apache Spark running the Minebench application. Axis are in log scale. Figure 2.7 shows the scalability results obtained by Ignis and Spark when running the Proofof-Work algorithm using both implementations (i.e., only Python and Python & C++) on our cluster. The strong scaling tests were obtained using a 120MB input file containing 300k blocks, while the weak scaling experiments start from a 120MB input file (one core) to reach 4.2GB and 9,600k blocks (32 cores). A least-squares fit is also provided to estimate the scalability results using up to 256 cores. In addition, graphs show a line corresponding to the ideal scalability. First, we analyze the strong scaling results. When considering the Python application (Figure 2.7(a)), Spark and Ignis obtain similar results with a small number of cores. However, as the number of executors increases, the cost of starting JVMs and transferring data through system pipes to the Python processes degrades the Spark global performance. Therefore, running Python applications on Spark prevents from reaching high levels of parallelism. For example, 32 Chapter 2. Towards an efficient and scalable multi-language Big Data framework 1 2 4 8 16 32 64 102 103 104 Cores Execution Time (seconds) Spark (MLlib) Spark (Python) Ignis (Python) Ignis (C++) 1 2 4 8 16 32 64 102 103 104 Cores Execution Time (seconds) Spark (MLib) Spark (Python) Ignis (Python) Ignis (C++) Figure 2.8: Study of the scalability of Ignis and Apache Spark running the K‐Means application: K=12 (left) and K=81 (right). Axis are in log scale. Ignis is 1.3×faster than Spark using 32 cores. This behavior is even more clear running the multi-language application (Figure 2.7(b)). Since Spark sends data from Python to C++ processes through the JVM, the number of pipe operations increases. In this way, the overhead is greater than the one observed for the Python implementation. In this case, Ignis is 2.2×faster than Spark using 32 cores. The least-square fits point out that execution times diverge and, as a consequence, the performance difference between Ignis and Spark will increase when considering more cores. We can conclude that our framework shows a very good behavior in terms of strong scalability, always close to the ideal case (red line). On the other hand, the multi-language implementation clearly outperforms the Python benchmark. Weak scaling results are displayed in Figures 2.7(c) and 2.7(d). In those tests we keep constant the amount of work performed for each core. The overhead detected in the results considering the Python code demonstrates that Python is not natively supported in Spark like Java or Scala. Note how the Spark scalability is getting away from the ideal case (horizontal red line) as the number of cores increases. It can be observed that Ignis shows a better behavior. When considering the multi-language code, the overhead of Spark becomes even more noticeable. That overhead is due to the exchange of data between processes. Pipes use the hard disk for those exchanges, causing an important bottleneck in the performance. However, Ignis still keeps a good scalability. 2.5.3 K-means K-Means is a classical machine learning algorithm for data clustering, and it is a good example of an iterative MapReduce application pattern. The goal of this algorithm is to classify a given data set through a certain number of clusters (Kclusters). Each cluster has a centroid. The algorithm 33 Chapter 2. Towards an efficient and scalable multi-language Big Data framework 1 2 4 8 16 32 64 0 0.5 1 1.5 2 2.5 Cores Speedup Ignis (Python) with respect to Spark (Python) K = 12 K = 81 1 2 4 8 16 32 64 0 0.5 1 1.5 2 2.5 3 3.5 Cores Speedup Ignis (C++) with respect to Spark (MLlib) K = 12 K = 81 Figure 2.9: Speedup of Ignis with respect to Apache Spark running the K‐Means application. works iteratively in such a way that every iteration each data point is assigned to the nearest centroid, and the Kcentroids are recalculated as barycenters of the clusters resulting from the previous assignment step. These operations are compute-intensive tasks. The algorithm iterates until the centroids do not change their location or other criterion is fulfilled. The final goal of the algorithm is that points within the same cluster are as similar as possible (i.e., high intra-class similarity), while points from different clusters are as dissimilar as possible (i.e., low inter-class similarity). Spark provides its own implementation of K-means in MLlib (Machine Learning Library) [67]. We have implemented a pure Python version for Spark and Ignis of the algorithm described in [13]. A Python-C++ multi-language version of that algorithm was also executed on the Ignis framework. In that version compute-intensive tasks were implemented in C++, while the K-means algorithm was specified in Python. For fair comparison we also included the Spark performance results for the MLlib version. Experiments were conducted using the NUS-WIDE dataset [21], which contains 269,648 images with 500 attributes per image. Figure 2.8 shows the strong scalability of the K-Means application. Results were obtained after 10 iterations using different number of clusters: K=12 (left graph) and K=81 (right graph). As we noted above, Spark (Python) and Ignis (Python) execute the same K-means code on both platforms. Tests point out that the continuous exchange of data between the Python processes and the JVM degrades the Spark performance when running an iterative application like K-Means. Performance differences between Ignis and Python grow as the parallelism increases. On the other hand, the Spark implementation that uses MLlib outperforms Spark (Python) and Ignis (Python) even though the driver code was also programmed in Python. This behavior indicates that the KMeans method in a Python Spark driver code is just a wrapper of the most 34 Chapter 2. Towards an efficient and scalable multi-language Big Data framework 1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 Iteration 0 2 4 6 8 10 12 14 16 18 20 Iteration Time (seconds) Spark Ignis Figure 2.10: Iteration times for K‐Means in presence of a failure. One node failed at the start of the 10th iteration. efficient Scala implementation of the algorithm. In this way, the overhead caused by the pipe operations disappears. In any case, the Python-C++ multi-language code on Ignis is the fastest implementation. To summarize the benefits of using Ignis with respect to Spark when running iterative MapReduce applications we show in Figure 2.9 the previous K-Means results expressed in terms of speedup. The left graph illustrates how Python is considered a non-native language by Spark. To do so, the speedup between the pure Python K-Means versions running on Spark and Ignis for each number of cores is calculated. According to the results, Ignis is on average 1.94×and 1.25×faster than Spark with K=12 and K=81, respectively. Speedups up to 2.4×were reached. The right graph displays the speedup between the best performing Spark version (MLlib) and the Ignis multi-language code for a particular number of cores. In this case, Ignis is on average 1.58×and 2.06×faster than Spark with K=12 and K=81, respectively. Therefore, Ignis is able to handle efficiently the combination of different programming languages in one iterative MapReduce application, extracting the maximum performance from the considered parallel system. On the other hand, one of the main features of Ignis is that it provides fault tolerance efficiently. As we pointed out in Section 2.2.3, if some data is lost, Ignis has enough information about how it was derived. In this way, only those operations needed to recompute the corresponding portion of data are performed. Therefore, lost data can be recovered without requiring costly replication. Next we will evaluate the cost of reconstructing a data partition after a node failure in the K-Means application. Figure 2.10 compares the running times for 15 iterations of K-Means, with one where a node fails at the start of the 10th iteration. We also included the Spark times for the same 35 Chapter 2. Towards an efficient and scalable multi-language Big Data framework 1 2 4 8 16 32 64 102 103 104 105 Cores Execution Time (seconds) Spark Ignis (Python) Ignis (C++) 1 2 4 8 16 32 64 0 10 20 30 40 50 60 70 80 Cores Speedup Spark Ignis (Python) Ignis (C++) 1.14x faster 1.72x faster Figure 2.11: Study of the scalability of Ignis and Apache Spark sorting 35GB of text data. scenario. Tests were performed considering the Python implementation using 16 cores on a 4 nodes cluster. Until the end of the 9th iteration, the Ignis iteration times were about 3 seconds. In the 10th iteration, one of the nodes is killed, resulting in the loss of the tasks running on that machine and the data partitions stored there. As a consequence, Ignis reruns these tasks on other machine nodes, rebuilding the data, which increases the iteration time to 15s. Once the lost data is reconstructed, the Ignis iteration time goes back down to 3s. Spark behaves similarly, increasing the 10th iteration time from about 7s to 19s, going back to the normal iteration time afterwards. Note that with a checkpoint-based fault recovery mechanism, recovery would likely require rerunning at least several iterations, depending on the frequency of checkpoints. 2.5.4 Sort Benchmark Sorting is a very common and useful data-intensive operation, thus it is supported by Spark and Ignis. Elements in Ignis are sorted by means of the SampleSort algorithm where elements are distributed by a regular sampling among the executors [61]. This task requires that executors exchange data. In both platforms the user can define a comparison function to be used together with sort, but the algorithm itself can not be modified. In contrast to the map operation, Spark does not allow to implement that user-defined function in a programming language different than the one used in the driver. That limitation does not appear in Ignis since any combination of programming languages for the driver and tasks is permitted. In this way, a Ignis driver written in Python could use a Java or C++ comparison function. Performance tests were carried out sorting in ascending order 35GB of text data, which contains about two billion lines. Each line has from 10 to 30 bytes of random text. We have 36 Chapter 2. Towards an efficient and scalable multi-language Big Data framework 1 2 4 8 16 32 64 103 104 Cores Execution Time (seconds) Spark Ignis 1 2 4 8 16 32 64 0 0.5 1 1.5 2 2.5 3 3.5 4 4.5 5 Cores Speedup Spark Ignis 1.43x faster Figure 2.12: Study of the scalability of Ignis and Apache Spark running the CG application. Axis are in log scale. compared the sort built-in capability of Spark with Ignis. Two different implementations in Ignis were analyzed. In the first one, the code was completely programmed in Python, while the second has a Python driver and the sort operation uses a user-defined C++ function for comparison purposes. Results are displayed in Figure 2.11. The left graph shows the execution times up to 64 cores. We can observe that while the behavior of the strong scalability is similar for all the cases, Ignis outperforms Spark in terms of execution time even for the Python implementation. Using a Python-C++ multi-language application is again the best option. Speedups, which are displayed in the right graph, were calculated using as reference the Spark execution with one core. Results confirm the previous observations regarding computing times. For instance, Ignis sorts data 1.14×and 1.72×faster than Spark using 64 cores when considering the Python and the multi-language versions, respectively. 2.5.5 Conjugate Gradient (CG) The Conjugate Gradient (CG) is an iterative method for solving sparse, large systems of linear equations. It is considered one of the most important computational kernels lying at the heart of many HPC scientific and engineering applications. A CG iteration involves the following compute-intensive linear algebra operations [14]: one sparse matrix-vector product, three vector updates, and two inner products. We have implemented the CG method for Spark and Ignis using Python. BLAS routines were carried out using the Intel MKL library. Our tests were conducted considering as input a 100k×100ksparse matrix with 900M nonzeros. The method was executed until 100 iterations were reached. The experimental results are shown in Figure 2.12. The left graph illustrates the strong scalability of CG using up to 64 cores. Although the scalability results are not very good, 37 Chapter 3. Improving the interoperability between HPC and Big Data languages and programming models –Mesos+Singularity: The same benefits commented above apply to the combination of Mesos and Singularity. –Nomad: It combines in the same framework a resource and a scheduler manager. Due to its lack of dependencies, it is the best option to install in a cluster from scratch. Moreover, it has better support for devices like GPUs than Mesos, which allows to create heterogeneous execution environments. It is important to highlight that the cost of deploying Docker containers is very low. For instance, considering Nomad, it is possible to deploy thousands of Docker containers in just a few seconds2. 3.2.4 Driver module The Driver module is a user API that allows access to all IgnisHPC functionalities. The driver program describes the high-level control flow of the application, and it can be programmed in any of the supported languages (Java, Python and C/C++). The Driver was designed as a Thrift RPC interface to the Backend so the logic has not to be re-implemented for every programming language. More details about the driver API and how to implement an application in IgnisHPC are provided in Section 3.3. 3.2.5 Backend module The Backend module contains the services that define the Driver’s logic. For instance, the reduceByKey function requires searching and grouping the keys. These operations are defined in the Backend, but they are implemented in the Executor module for a specific programming language. The Backend module is also responsible of sending requests to the resource manager in accordance with the instructions specified in the driver code. These instructions are lazily executed, so the Backend registers the function calls as a task dependency graph. When a task that represents an action in the driver code is created (e.g., a call to count), all the tasks in its dependency graph are executed. An example is shown in Figure 3.3, where the Action Task depends on Task 3, and Task 3 depends on Task 1 and 2. Note that a task dependency is only computed if the task was never executed or if its result was not explicitly cached. The Executor and Container tasks are always executed as final dependencies. These tasks check that executors and containers are running and ready to be used. Finally, IgnisHPC is able to recover after a failure of a cluster node or some of the executors. Affected tasks are traced by the Backend in such a way that only their executors are reallocated 2https://www.hashicorp.com/c1m [accessed September, 2022] 44 Chapter 3. Improving the interoperability between HPC and Big Data languages and programming models Container Task Executor Task Task 1 Task 3 Action Task Task 2 isCached isNotCached Figure 3.3: Example of a task dependency graph. and recomputed. If the affected tasks are cached, the recovery process will be faster since it is not necessary to recalculate their dependencies. Note that this process is automatic, but users can tune the recovery process using the persistence functions in the driver code (see Section 3.3). 3.2.6 Executor module The Executor module implements the operations defined by the Backend, where each supported programming language has its own implementation. In order to add support for a new language in IgnisHPC, a minimum implementation only requires programming the context class. The executor context allows the API functions to interact with the rest of the IgnisHPC system. In this way, among the functionalities of the context we find the exchange of user variables between driver and executors or providing the executors access to the MPI communicators. As we explained previously, IgnisHPC uses MPI (that is, MPI communicators) to perform all the communications related to the Executor module. IgnisHPC constructs three types of communicators for data transfers (see Figure 3.4): –Base communicator: for each worker there is a communicator that includes all its executors. This communicator always exists. If one executor is lost, the communicator is destroyed and a new communicator is created including a new executor. To that end, the capability of linking dynamically a single process to a communicator introduced in MPI-3 was of special importance. Without this feature all processes would have to be launched at the same time, and in case a process died, it could not be replaced causing the job to fail. –Driver communicator: it joins a base communicator to the driver process. It is created when the driver and a worker exchange data. –Inter-worker communicator: it is created joining the base communicators of two workers. It is used to send data from one worker to another and it will be destroyed as soon as one of 45 Chapter 3. Improving the interoperability between HPC and Big Data languages and programming models Base Communicator Executor11 Executor1n Executorm1 ... Worker with mExecutors Containers Driver Communicator Worker Driver Driver Container Inter-worker Communicator Worker1 Worker2 ... Executormn ... Figure 3.4: MPI communicators in IgnisHPC. the two workers stops its execution. This communicator is only created between workers that execute the operation ImportData. The base communicator for each worker is accessible to programmers by the executor context. It means that IgnisHPC functions can be implemented using that communicator and MPI primitives (e.g. gather, scatter, broadcast, reduce, etc.). As a result, IgnisHPC supports the execution of pure MPI applications with minimal modifications in the original code. A detailed explanation is provided in Section 3.4. Another benefit of using MPI for data transfers is the performance improvement of iterative applications. When using Big Data frameworks such as Ignis and Spark, an iterative application requires the driver to perform an evaluation task per iteration to obtain the final result. Each evaluation has three steps: stopping the executors, analysis of the partial results by the driver and restarting the executors. Note that starting and stopping the executors is very costly in terms of performance. IgnisHPC avoids the driver evaluations because executors share the partial results of each iteration using their MPI base communicator. Therefore, it is not necessary to stop them because they do not need to wait for the driver. This has even a bigger impact on performance for those applications with many short iterations. 3.2.7 Submitter module The Submitter is an IgnisHPC module consisting of a set of scripts and utilities for configuring and launching jobs. As we commented previously, Ignis had no module to launch tasks, and the Driver module was launched manually using Ancoris (see Section 3.2.3). The Submitter is a container, which can be accessed by SSH. There users can set up jobs in a similar environment where the IgnisHPC applications will be executed. The main utility of the Submitter module is the ignis-submit script that, like spark-submit, allows users to launch IgnisHPC jobs in the cluster. The script only requires as mandatory argu46 Chapter 3. Improving the interoperability between HPC and Big Data languages and programming models 1# Python driver 2ignis-submit ignishpc/python python3 mydriver.py 3 4# C++ driver 5ignis-submit --name myapp --properties ignis.driver.memory=2GB ignishpc/cpp ./mydriver 0 -g 2 Figure 3.5: Job submission examples. ments a Docker image and the driver program. There are also the following optional parameters: –name: a job name can be specified. –arguments: after the driver program name, all the parameters will be considered as driver arguments. –attach mode: by default, jobs are launched in unattached mode. That is, ignis-submit launches the job and exits. Attach mode allows users to control the job as if the driver runs locally, so output is printed in real time, and it is possible to manually kill the job. –properties: users can change the default properties before launching the Driver module. Executor properties can be redefined later but Driver properties are set only by ignis-submit. Figure 3.5 shows two job submission examples. The first one is a Python basic submission with only its base image and the driver application. The second one deals with a C++ driver and optional parameters. In particular, --name sets the job name, --properties changes the driver default memory to 2 GB, and 0 -g 2 are considered arguments of mydriver. Note that ignishpc/cpp is the C++ base image. 3.2.8 Data storage In the same way as Ignis, IgnisHPC provides multiple options for data storage. Users can choose a type of storage according to their particular execution environment. Storage must be defined as a property before the worker creation. IgnisHPC supports the following storage options: –In-Memory: it is the best performer since all data is stored in memory. Memory consumption could be an issue, so it is not suitable for all kinds of jobs. –Raw memory: data is stored in a memory buffer using a serialized binary format. Extra memory consumption is minimal and the buffer is compressed by Zlib, which has nine compression levels. Level six is applied by default, but it can be changed when the worker properties are defined. –Disk: similar to raw memory but the buffer is stored as a POSIX file. Performance is much lower but it allows to work with large amounts of data that cannot be completely stored in memory. 47 Chapter 3. Improving the interoperability between HPC and Big Data languages and programming models Type Functions Conversion map, filter, flatmap, keyBy, mapPartitions, keys, values, mapValues, etc. Group groupBy, groupByKey Sort sort, sortBy, sortByKey Reduce reduce, treeReduce, aggregate, treeAggregate, fold, reduceByKey, aggregateByKey, etc. I/O collect, top, take, saveAsObjectFile, saveAsTextFile, saveAsJsonFile, etc SQL union, join, distinct Math sample, sampleByKey, takeSample, count, max, min, countByKey, countByValue Balancing repartition, partitionBy Persistence persist, cache, unpersist, uncache Table 3.1: Example of some IDataFrame functions supported by IgnisHPC. IgnisHPC behaves very similar to Spark in terms of data locality. Data is split into several partitions which are assigned to executors. Each partition is assigned to a single executor. In case another executor needs that partition, it will be sent using MPI. By default, IgnisHPC stores all partitions in memory to achieve the best possible performance. If there is not available memory, Docker sends some data to the container swap automatically. The user can modify this policy by changing the swap size or removing it. In case the data size exceeds the swap or the swap performance is insufficient, IgnisHPC allows users to set disk as primary storage, so partitions will always be stored on disk. In this case partitions will only be loaded into memory at processing time. On the other hand, there is an important difference in how memory is handled by Ignis and IgnisHPC. Since Ignis was just a prototype, for simplicity in the implementation, it assigns one data partition to each executor. In this way, if it is necessary to increase the partition size, a realloc operation is performed in such a way that the complete partition is copied to a different memory location. The consequence is a noticeable increment in the memory consumption. This restricts Ignis to work with smaller input datasets. IgnisHPC overcomes that limitation supporting several data partitions per executor. Note that an executor can spawn several threads to process the data partitions in parallel. 3.3 PROGRAMMING APPLICATIONS FOR IGNISHPC IgnisHPC requires a minimal driver code that implements the application at high-level. To facilitate the adoption from the Big Data community, the IgnisHPC API is inspired by the Spark API in such a way that IgnisHPC codes are easily understandable by users who are familiar with Spark. In comparison to Ignis, we have extended the API to cover the most important primitives 48 Chapter 3. Improving the interoperability between HPC and Big Data languages and programming models required by Big Data applications. For instance, IgnisHPC includes functions such as join and union for graph processing. The IgnisHPC driver API is composed by six main classes: –Ignis starts and stops the driver environment. –IProperties defines the execution environment properties. –ICluster represents a group of executors containers. It is possible, for example, to execute remote commands (execute,executeScript) and send files (sendFile,sendCompressedFile) to them. –IWorker represents a group of processes of the same programming language. This class includes functions to read files (textFile,partitionJsonFile, etc.), import data partitions from another worker (importData), send data from the driver (parallelize) and execute external codes (loadLibrary,call,voidCall). As we explain later, the former routines allow IgnisHPC to execute MPI applications within the framework. –IDataFrame contains all the functions of the MapReduce paradigm, similarly to Spark RDD. A function can be a transformation that generates another IDataFrame or an action that generates a final result (see Table 3.1). With respect to Ignis, besides the support for new API functions, IgnisHPC has increased the overall performance for some types of routines (e.g., Group and Sort) thanks to its complete redesign using MPI. –ISource is an auxiliary class used by meta-functions such as map in the driver. This class acts as a wrapper for the input parameters. It is also used to store variables and send them to the executors. Those variables can be obtained by the executors using the context. Note that all the API operations that move data between executors perform an internal shuffling operation (e.g., parallelize,collect,partitionBy, etc.). 3.3.1 An example: Transitive Closure With the goal of illustrating how applications are programmed in IgnisHPC, Figure 3.6 shows an example of a driver implemented in Python for computing the Transitive Closure of a graph. This algorithm finds out if a vertex xis reachable from another vertex yfor all vertex pairs (x,y) in the graph. Note that an equivalent driver code could be implemented in any of the IgnisHPC supported languages (C/C++ and Java) using a similar syntax. First, the IgnisHPC framework is initialized (line 6). Properties are created to configure and build a cluster (lines 8 to 14). Note that the properties definition is optional, and IgnisHPC could read them from a default configuration file. Moreover, IgnisHPC introduce the possibility of overwrite the default values when a job is submitted like Spark using the new Submitter module. Therefore, the Docker image, the number of containers, the number of cores and the memory per container could be defined out of the driver code. Computing the Transitive Closure has two phases. To illustrate the multi-language support in IgnisHPC, each phase was implemented in 49 Chapter 3. Improving the interoperability between HPC and Big Data languages and programming models 1#!/usr/bin/python 2 3import ignis 4 5# Initialization of the framework 6ignis.Ignis.start() 7# Resources/Configuration of the cluster 8prop = ignis.IProperties() 9prop["ignis.executor.image"] = "ignishpc/full" 10 prop["ignis.executor.instances"] = "2" 11 prop["ignis.executor.cores"] = "4" 12 prop["ignis.executor.memory"] = "2GB" 13 # Construction of the cluster 14 cluster = ignis.ICluster(prop) 15 # Initialization of a Python Worker 16 worker_python = ignis.IWorker(cluster, "python") 17 # Task 1: Tokenize text into pairs (x, y) 18 # The edges are stored in reversed order 19 tc = worker_python.textFile("graph.dat") 20 edges = tc.map(lambda x_y: x_y.split(" ")[::-1]) 21 # Initialization of a C++ Worker 22 worker_cpp = ignis.IWorker(cluster, "cpp") 23 # Transfer data from Task 2 - Python 24 tc2 = worker_cpp.importData(tc).cache() 25 26 oldCount = 0 27 nextCount = tc.count() 28 while True: 29 oldCount = nextCount 30 # Task 2: Perform the join,(y, (z, x)) pairs, 31 # to obtain the new (x, z) paths. 32 new_edges = tc2.join(edges). 33 map("libexample.so:Reverse2") 34 tc2 = tc2.union(new_edges).distinct().cache() 35 nextCount = tc2.count() 36 if nextCount == oldCount: 37 break 38 # Show result 39 print("TC has %i edges" % tc2.count()) 40 # Stop the framework 41 ignis.Ignis.stop() Figure 3.6: Transitive Closure driver code in Python. a different programming language. The first one uses a Python executor and the second a C++ executor. The first stage consists of a map operation that takes as input a text file and creates pair values that represent edges in the graph (line 20). As a consequence, it is necessary to previously create a Python worker in the cluster (line 16). It is important to note that creating the worker is mandatory and is not related to the driver programming language. On the other hand, if the worker and the driver code are in the same language, lambda functions can be used (line 20). The following phase of the algorithm is iterative: edges are joined into a path until there are no new paths. Since this phase is implemented in C++, a C++ worker should be created (line 22). Data is shared between workers using the importData function (line 24). The driver code ends printing the results. The framework must be stopped before the driver ends (line 42) to stop the backend. Unlike Ignis, IgnisHPC automatically detects when the driver process ends 50 Chapter 3. Improving the interoperability between HPC and Big Data languages and programming models 1#include<ignis/executor/api/function/IFunction.h> 2 3using namespace ignis::executor::api; 4typedef std::pair<int64_t, int64_t> ipair; 5typedef std::pair<int64_t, ipair> ipair2; 6 7// Reverse pairs (z, (x, y)) to (z, (y, x)) 8class Reverse2 : public function::IFunction< 9ipair2, ipair2>{ 10 ipair2 call(ipair2& x_y, IContext& context){ 11 return {x_y.second, x_y.first}; 12 }; 13 14 ignis_export(Reverse2, Reverse2) Figure 3.7: Function in C++ used by a map operation for the Transitive Closure application. and stops it. However, this is not a good practice. In IgnisHPC a lazy evaluation is performed when a result is not required explicitly. In the example, the trigger that causes the tasks to be launched are the calls to the count function. This approach is also followed by Spark where RDDs are computed lazily the first time they are used in an action. Most of the driver functions are meta-functions. That is, generic functions that require another one to perform an internal operation. This is the case of map in the example of Figure 3.6 (lines 20 and 33-34). To implement those functions we should use the executor API provided by IgnisHPC. Basically, it defines a simple interface based on the number of required input parameters. Figure 3.7 shows an example corresponding to the C++ function used by map in the driver code. Since map takes one parameter and also returns one parameter, IFunction is used. In case there are two input parameters (e.g., reduce), IFunction2 would be used, and so on. Note that if the function does not return any value (e.g. foreach), functions have the same name but with the Void prefix. 3.3.2 Text lambda functions As explained previously, lambda functions need that driver and executor codes were implemented in the same language because native code serialization is required. We refer to code serialization as the process by which a function or set of instructions are converted into bytes to be sent and executed in a different environment. Note that although Python and Java are able to serialize code both languages face compatibility problems. On top of that, C++ does not allow any type of code serialization. To overcome these limitations IgnisHPC implements its own multi-language lambdas without source code serialization, named text lambdas. In this way, IgnisHPC allows to define lambdas as text, using the executor language syntax. The executor will transform the lambda text into source code to be used as a meta-function parameter. It is important to highlight that the language of the driver code is indifferent. 51 Chapter 3. Improving the interoperability between HPC and Big Data languages and programming models 1// Python lambda 2data.reduce("lambda a, b: a + b") 3 4// C++ lambda 5v = 2 6data.reduce( 7ISource("[](int a, int b, IContext &c){ \ 8return (a + b) * c.var<int>(\"v\"); \ 9}").addParam("v",v) 10 ) Figure 3.8: Examples of text lambda for a Python and a C++ executor. Figure 3.8 shows an example of a text lambda that accumulates all elements (line 2) used by areduce function. It uses Python syntax so must be evaluated by a Python executor. Another example is shown in line 7. It defines a text lambda that captures the value of a variable, which is read from the context. This lambda function will be compiled and loaded by a C ++ executor. Performance is not affected by using text lambda functions but it can add some overhead to the compilation process, especially for C++. In the same way, using the mechanism that allows to execute text lambda functions, IgnisHPC can send a complete job or application to the executors. This is possible thanks to loadLibrary, which can be used to execute a full source code as an IgnisHPC library. It has again a small impact on the compilation time. More details about loadLibrary are provided in Section 3.4.2 using MPI applications as use case. 3.4 MPI ON IGNISHPC Our first prototype, Ignis, was limited to perform inter-process communications using only TCP sockets (different computing nodes) or shared memory (same node). However, IgnisHPC was completely redesigned to use MPI as backbone technology. As a consequence, all communications are internally implemented by MPI routines. As we explained previously, this change makes it possible for IgnisHPC to support more communication models and network architectures. In addition, an important advantage of our approach is that, once IgnisHPC has configured the MPI communications, users can combine in the same application pure MPI libraries using the IgnisHPC communicators together with standard Big Data functions (map,reduce,collect, etc.). 3.4.1 Integration of MPI into a Big Data environment MPI was not designed to run on Docker containers. As a consequence, there are several problems that should be addressed. First, by default and to preserve an isolation runtime environment, Docker creates a private virtual network between the host and the containers. If two containers 52 Chapter 3. Improving the interoperability between HPC and Big Data languages and programming models are launched on the same host, we can execute MPI processes in the same way that they were two real nodes of a cluster. But if we launched those containers on different hosts, the communication is impossible since they belong to different networks. We can find in the literature several works that deal with this issue (see Section 3.6). For instance, some approaches opt for launching the container on the host network or creating a virtual network between the cluster nodes [27]. However, these configurations are difficult to handle and implement by resource managers in Big Data environments. Second, there are important differences in how ports are handled by a Big Data or an HPC environment. For instance, ports are considered a resource in a Big Data environment because there are services that require an exclusive port, which is not the case in HPC. MPI needs ports to establish connections between processes but restricted to a range. However, ports provided by resource managers in a Big Data environment are usually random and not consecutive. Finally, IgnisHPC can internally spawn several threads per MPI process (executor) to increase the performance when processing and/or communicating data. All these threads use communicators to exchange data in parallel. Every time a communicator is created, MPI assigns a virtual interface to it. However, a virtual interface can only be used by one communicator at the same time, so parallel communications require the use of multiple virtual interfaces. In the most recent MPICH version, which is the MPI implementation used by IgnisHPC, virtual interfaces are assigned sequentially. Since IgnisHPC creates and destroys communicators dynamically, it is not possible to assure that threads can always exchange data in parallel using communicators with different assigned virtual interfaces. The consequence is a degradation in the performance. To overcome the above problems, IgnisHPC applies the following changes to MPICH: – Containers: MPICH has been designed to work on a local network. Docker containers can be joined to a network but only within the same node (internal network). Although resource managers can export a service from the internal network to the local network, this causes a problem when MPI is executed containerized. MPICH uses a service to store the network addresses of the launched MPI processes, but when using containers, each MPI process stores the values corresponding to the internal network which are not valid outside the node. For this reason, it is necessary to modify MPICH in order to store the correct network values that correspond to the local network. – Ports: now MPICH uses a list of ports provided by the resource manager instead of a range. – Multithreading: MPICH was modified to assure that all threads use a different virtual connection. In this way, communications between threads can always be performed in parallel. 3.4.2 Running MPI applications in IgnisHPC 53 Chapter 3. Improving the interoperability between HPC and Big Data languages and programming models 1 2 4 8 16 32 64 128 160 Cores 102 103 104 105 Time (sec) - log scale Spark IgnisHPC (Python) IgnisHPC (C++) (a) Times (strong scaling) 1 2 4 8 16 32 64 128 160 Cores 0 20 40 60 80 100 120 Speedup Spark IgnisHPC (Python) IgnisHPC (C++) (b) Speedup Figure 3.16: Study of the scalability of IgnisHPC and Apache Spark running the TeraSort application. 1 2 4 8 16 32 64 128 160 240 Cores 101 102 103 104 Time (sec) - log scale Spark (MLlib) Ignis (C++) Ignis (Python & C++) IgnisHPC (Python & C++) (a) Times (strong scaling) 1 2 4 8 16 32 64 128 160 240 Cores 0 20 40 60 80 Speedup Spark (MLlib) Ignis (C++) Ignis (Python & C++) IgnisHPC (Python & C++) (b) Speedup Figure 3.17: Study of the scalability of IgnisHPC, Ignis and Apache Spark running the K‐Means application. the NUS-WIDE dataset [21], which contains 269,648 images with 500 attributes per image. In the tests the results were obtained after 10 iterations and K=81. –PageRank (PR). It is a graph algorithm which ranks elements by counting the number and quality of links. To evaluate the PR algorithm on IgnisHPC and Spark we used the LiveJournal graph from the SNAP repository [59], which contains 4.8M vertices and about 69M edges. –Transitive Closure (TC). One of the most basic questions that arises when analyzing a complex graph Gis whether one vertex xcan reach another vertex yvia a directed path. A way to store this information is to construct another graph, such that there is an edge (x,y)in the new graph if and only if there is a path from xto yin the input graph. This new graph is called the Transitive Closure of G. Since computing the TC is very costly, we used a small graph with 75 vertices and 200 edges in our tests. 3.5.2.1 Analysis and discussion We now present the performance results from our evaluation of Spark, Ignis and IgnisHPC considering all the Big Data applications detailed above. Speedups were calculated using as reference the Spark sequential time. All experiments have been executed ten times and their average result is reported. In addition, the relative standard deviation (RSD), also known as the 60 Chapter 3. Improving the interoperability between HPC and Big Data languages and programming models coefficient of variation, is calculated as 100 × σ /x. Low RSD values point out that the data is tightly clustered around the mean. Figures 3.14 and 3.15 show the scalability study of the Minebench (MB) application using two different implementations. In the first one, MB was programmed using only Python, while in the second one, the two chained map operations are implemented using Python (data-intensive task) and C++ (compute-intensive task), respectively. Results show that IgnisHPC exhibits very good strong-scaling behavior for both implementations. On the contrary, the Spark scalability is impacted for the cost of starting JVMs and transferring data through system pipes to the Python processes, causing an important degradation in the overall performance [79]. This scenario is even more clear for the multi-language implementation in Figure 3.15(a) since Spark sends data from Python to C++ processes through the JVM, increasing the number of pipe operations. As a consequence, the Spark strong scalability is really poor. On the other hand, IgnisHPC weak scales very well for both code versions (Figures 3.14(c) and 3.15(c)), which is not surprising since there is not much communication in the MB application. As a consequence, IgnisHPC is able to extract all the existent parallelism. Finally, we must highlight that IgnisHPC clearly outperforms Ignis both in terms of strong and weak scalability, especially when considering the multi-language application. The RSD for the Minebench experiments ranges from 0.3% to 4.9%. Note that for the Python implementation (Figure 3.14), the performance differences between IgnisHPC and Ignis are only caused by architectural improvements in the framework since no MPI operations are carried out. Executors read the input data from a file and exchange partial results using the shared memory. However, if we consider the multi-language implementation of MB (Figure 3.15), performance improvements are also due to the use of MPI for the communication between Workers. Performance results of the TeraSort (TS) application running on IgnisHPC and Spark frameworks are displayed in Figure 3.16. Ignis results are not shown because the memory consumption of sorting 1 TB of data is too high for our cluster. As we explained in Section 3.2.8, Ignis assigns one data partition to each executor. For TS those partitions are very large. Every time an element is added to a partition, the complete partition is copied to a different memory location (realloc operation), which causes a boost in the memory requirements. This restricts Ignis to work with smaller input datasets. We avoid this limitation since IgnisHPC was designed to allow several partitions per worker. Two different TS implementations in IgnisHPC were analyzed: a pure Python code and a multi-language Python-C++ code. In the latter case, the sort operation uses a user-defined C++ function for comparison purposes. For all the cases IgnisHPC outperforms Spark, especially when considering the multi-language implementation. In this way, for instance, IgnisHPC is 116×faster than sequential Spark when using 160 cores, while Spark reaches a speedup of only 66×(Figure 3.16(b)). It allows IgnisHPC to sort 1 TB of data in barely 5 minutes. The RSD for the TS experiments ranges from 1.3% to 6%. 61 Chapter 3. Improving the interoperability between HPC and Big Data languages and programming models 1 2 4 8 16 32 64 128 192 240 Cores 102 103 104 Time (sec) - log scale Spark IgnisHPC (a) Times (strong scaling) 1 2 4 8 16 32 64 128 192 240 Cores 0 5 10 15 Speedup Spark IgnisHPC (b) Speedup Figure 3.18: Study of the scalability of IgnisHPC and Apache Spark running the PageRank application. 1 2 4 8 16 32 64 128 192 240 Cores 103 104 Time (sec) - log scale Spark IgnisHPC (a) Times (strong scaling) 1 2 4 8 16 32 64 128 192 240 Cores 0 1 2 3 4 5 Speedup Spark IgnisHPC (b) Speedup Figure 3.19: Study of the scalability of IgnisHPC and Apache Spark running the Transitive Closure application. Strong scaling results of K-Means (KM) are shown in Figure 3.17. We used as reference the Spark implementation of this algorithm included in MLlib (Machine Learning Library) [67]. For Ignis, a pure C++ and a Python-C++ implementations of KM were analyzed. For IgnisHPC, we only show the results for the multi-language Python-C++ code because the numbers obtained by a pure C++ application are very similar. We can observe that Ignis was able to beat Spark when considering the C++ code. However, an important degradation in the scalability was detected for the multi-language implementation as the parallelism increases. This problem was explained in Section 3.2.6 and is related to the way Ignis handles iterative applications. Ignis starts and stops the executors each iteration because the driver must compute the partial results, which has an important impact on the performance. To deal with this, IgnisHPC takes advantage of MPI in such a way that executors compute the partial results and share them without intervention of the driver. For this reason IgnisHPC exhibits a very good strong scalability even for multi-language iterative applications, decreasing noticeably the execution times with respect to Spark and Ignis. It can be observed in Figure 3.17(b) that IgnisHPC is about two and four times faster than Spark and Ignis (multi-language code) when using all the cores in the cluster, respectively. The RSD for the KM experiments ranges from 1.1% to 4.8%. In Big Data analytics many problems require processing graphs. For this reason it is essential for a Big Data framework as IgnisHPC to include primitives to support this kind of applications. 62 Chapter 3. Improving the interoperability between HPC and Big Data languages and programming models Application No. times faster than Spark No. times faster than Ignis Minebench 3.87×[Python & C++] 1.23×[Python & C++] 1.26×[Python] 1.08×[Python] TeraSort 1.76×[C++] – 1.35×[Python] K-Means 1.94×[Python & C++] 1.28×[Python & C++] PageRank 1.10×[Python] – Transitive Closure 1.12×[Python] – Table 3.3: Summary of the IgnisHPC performance results for all the Big Data applications considering the maximum number of cores and the best Ignis and Spark implementation (in case there is more than one). Between brackets the programming language/s used in the IgnisHPC implementation. Point-to-point Collective Application Blocking Non-blocking Blocking Non-blocking LULESH –Isend,Irecv Allreduce,Barrier – AMG Send,Recv Isend,Irecv, Irsend Allreduce,Barrier,Bcast, Reduce,Alltoall, Allgather(v),Gather(v), Scan,Scatter(v) – MiniAMR Send,Recv Isend,Irecv Allreduce,Barrier,Bcast, Alltoall – MiniVite Sendrecv Isend,Irecv Allreduce,Barrier,Bcast, Reduce,Alltoall(v), Exscan Ialltoall MSAProbs Send,Recv Isend,Irecv Allreduce,Barrier,Bcast – Table 3.4: MPI calls used for communications in the HPC applications. In addition, since the size of the graphs to be processed is often very large, a good scalability is essential. With this goal in mind we have evaluated two well-known graph algorithms in Spark and IgnisHPC: PageRank (PR) and Transitive Closure (TC). Performance results are shown in Figures 3.18 and 3.19, respectively. The algorithms in Spark were implemented using GraphX [97]. Note that Ignis does not support several operations that are basic for this type of applications such as join and union (see Table 3.3), so it cannot be evaluated. Despite the fact that GraphX is a highly tuned API for graph processing, IgnisHPC is capable of outperforming Spark in both cases. The RSD for the PR experiments ranges from 0.6% to 2.7%, while for the TC varies from 0.1% to 1.7%. 3.5.2.2 Analysis and discussion Table 3.3 summarizes the performance gains obtained by IgnisHPC with respect to Spark and Ignis when running all the considered Big Data applications. Results were obtained using the maximum number of cores available and taking into account the best Ignis and Spark implementation (if there is more than one). In this way, for example, two implementations are available for Minebench, pure Python and multi-language Python-C++. In that case we used as reference for Spark the Python code, while for Ignis the best performing implementation was the 63 Chapter 3. Improving the interoperability between HPC and Big Data languages and programming models multi-language one (see the values in Figures 3.14(a) and 3.15(a) when using 240 cores). According to the results, IgnisHPC is from 1.10×to 3.87×faster than Spark. The good behavior of IgnisHPC is particularly relevant when considering multi-language applications. At the same time, IgnisHPC is a step forward with respect to Ignis in terms of performance. In this case, IgnisHPC is from 1.08×to 1.28×faster than Ignis. However, there are additional benefits. First, the memory consumption in IgnisHPC was optimized allowing multiple partitions per executor, which allows to work with extremely large datasets. That is the reason why Ignis is not able to execute TeraSort in our cluster. And second, the IgnisHPC API was extended to support, among others, graph processing algorithms such as PageRank and Transitive Closure. 3.5.3 HPC applications For many years MPI has been the dominant parallel programming model in the HPC area. As we explained in Section 3.4, thanks to its architectural design, one of the most important features of IgnisHPC is its ability to execute native MPI applications within the framework just adding a few lines of code. To evaluate the benefits of our approach we are interested in two key areas: performance (with respect to the native MPI execution) and productivity (additional Source Lines Of Code - SLOC). In this way, we have selected five HPC applications coming from different scientific fields that represent a variety of MPI communication patterns. Table 3.4 summarizes the most important MPI calls (point-to-point and collective operations) used in the applications. Note that during a specific run, an application may use only a subset of these communications. All the codes were implemented using C/C++. For some of them we have considered hybrid implementations (MPI+OpenMP) to demonstrate that is also possible to efficiently execute this type of applications in IgnisHPC without additional effort. Next we provide some information about the selected HPC applications used in the experimental evaluation: –LULESH (Livermore Unstructured Lagrange Explicit Shock Hydrodynamics). It is a shock hydrodynamics code developed at Lawrence Livermore National Lab (LLNL) [55]. It has been ported to a number of programming models: MPI, OpenMP, MPI+OpenMP, CUDA, etc. In this Chapter we have considered the hybrid MPI+OpenMP implementation, which uses MPI between nodes and OpenMP for cores on a node. Performance tests were run on 8 nodes with a problem size of 703on each node, which corresponds to the most representative problem size [54]. –AMG. It is a parallel algebraic multigrid solver for linear systems arising from problems on unstructured grids. It is part of the Exascale Computing Project (ECP) proxy applications suite4and was derived directly from the BoomerAMG [48] solver. AMG is an SPMD ap4http://proxyapps.exascaleproject.org/ecp-proxy-apps-suite 64 Chapter 3. Improving the interoperability between HPC and Big Data languages and programming models (a) 1 2 4 8 16 20 Cores per node 102 103 Time (sec) - log scale IgnisHPC MPI (b) 1 2 4 8 16 20 Cores per node 0 1 2 3 4 5 Speedup IgnisHPC MPI Figure 3.20: Study of the scalability of LULESH (8 nodes). plication with about 65,000 lines of code which uses OpenMP threading within MPI tasks. Parallelism is achieved by simply subdividing the grid into logical P×Q×R(in 3D) chunks of equal size. AMG is a highly synchronous and memory-access bound code. The scalability tests were obtained with a fixed local problem grid size per MPI process of 100×100×100 points. –miniAMR. It is a proxy app for adaptive mesh refinement (AMR), which is a frequently used technique for efficiently solving partial differential equations (PDEs) [86]. It applies a stencil calculation on a unit cube computational domain, which is divided into blocks. This application also belongs to the ECP proxy app collection and was implemented using MPI. We used blocks with dimensions 8×8×8 and a maximum of 4 refinement levels. The test case we considered is that of an expanding sphere, which closely mimics an explosion. Blocks are refined along the boundary of the expanding sphere. –miniVite. It implements a parallel Louvain method for community detection, which is one of the most important graph kernels used in scientific and social networking applications for discovering higher order structures within a graph [38]. It is also included in ECP proxy app collection. miniVite was programmed using MPI and OpenMP. As input we used a graph with 10% of the vertices of the well-known friendster social network graph [99]. It consists of 6.6M vertices and 24.2M edges. –MSAProbs. One basic step in many bioinformatics analyses is the multiple sequence alignment (MSA). MSAProbs [62] is a state-of-the-art tool to compute protein MSA based on hidden Markov models. In this Chapter we have considered its MPI+OpenMP parallel implementation [42]. The input dataset PF07085 [71] used in the tests consists of 975 sequences with an average length of 512. 3.5.3.1 Analysis and discussion Next we carry out the analysis of the execution of the MPI-based HPC applications within IgnisHPC. As mentioned previously, we will focus on two aspects. First, the performance dif65 Chapter 3. Improving the interoperability between HPC and Big Data languages and programming models (a) 1 2 4 8 12 Nodes - log scale 102 103 Time (sec) - log scale IgnisHPC MPI (b) 1 10 100 240 Cores - log scale 102 103 Time (sec) - log scale IgnisHPC MPI Figure 3.21: Study of the weak scalability of AMG (20 threads/cores per node) (a) and miniAMR (b). (a) 1 2 4 8 12 Nodes 103 104 Time (sec) - log scale IgnisHPC MPI (b) 1 2 4 8 12 Nodes 0 2 4 6 8 10 Speedup IgnisHPC MPI Figure 3.22: Study of the scalability of miniVite (20 threads/cores per node). (a) 1 2 4 8 12 Nodes 102 103 104 Time (sec) - log scale IgnisHPC MPI (b) 1 2 4 8 12 Nodes 0 2 4 6 8 Speedup IgnisHPC MPI Figure 3.23: Study of the scalability of MSAProbs (20 threads/cores per node). ferences between running the HPC applications on the cluster as native MPI tasks or using IgnisHPC. It is important to highlight that is out of the scope of this thesis to analyze the particular behavior of each MPI application in terms of performance and scalability. This was extensively explained in the references provided in the description paragraphs of Section 3.5.3. The second key aspect is productivity. In our case we measured the source lines of code (SLOC) of the applications. This metric is very important since scientists will only adopt IgnisHPC to execute MPI applications if porting them requires little effort. Performance. Figure 3.20 shows the strong scaling results of LULESH. Note that this application uses a hybrid MPI+OpenMP approach. Performance differences between native MPI and 66 Chapter 3. Improving the interoperability between HPC and Big Data languages and programming models Application SLOC MPI SLOC IgnisHPC LULESH 5,918 5,993 (+75) AMG 65,154 65,197 (+43) MiniAMR 9,958 9,987 (+39) MiniVite 3,264 3,324 (+60) MSAProbs 6,045 6,062 (+17) Table 3.5: SLOC of the HPC applications. IgnisHPC are really small, always lower than 1.7%. Running LULESH from IgnisHPC shows the same scalability trend than the MPI native execution. Weak scalability tests were run to evaluate AMG (MPI + OpenMP) and miniAMR (MPI). Results are shown in Figures 3.21(a) and 3.21(b). In both cases also, the performance of IgnisHPC comes close to that of its counterpart, the native MPI implementation. The maximum performance difference is only about 1.4% for both applications. In this way, for instance, the execution times of miniAMR using all the cores in the cluster were 520 and 517 seconds with native MPI and IgnisHPC, respectively. miniVite scalability results are displayed in Figure 3.22. The behavior replicates the observations commented previously for the other HPC applications. That is, running an MPI application within the IgnisHPC framework achieves very similar performance with respect to the native execution. In this particular case, the maximum difference drops to only 0.2%. The same scalability analysis applied to MSAProbs (Figure 3.23) produces a maximum performance difference of 0.4%. So we conclude that running MPI (and MPI+OpenMP) applications from IgnisHPC is as efficient as executing them natively. Productivity. As we explained in Section 3.4, running MPI applications in IgnisHPC requires some minimal modifications to the original source code and adding a few lines to call the application from the driver code. It is important to highlight that most of these extra lines are devoted to parsing the arguments of the MPI application. In any case, this is a very simple and repetitive code as shown in the example of Figure 3.10, which can be considered as boilerplate. We measure the SLOC using SLOCcount [95] of each original MPI application and its counterpart adapted to IgnisHPC (see Table 3.5). The number of extra lines, between brackets in the table, ranges only from 17 to 75. This demonstrates that integrating MPI applications and libraries in IgnisHPC is a straightforward process, which is very important for the HPC community since it is not necessary to port MPI codes to a new API or programming model. Therefore, IgnisHPC fulfills its design goal of unifying in a single framework the benefits of HPC and Big Data applications. 67 Chapter 3. Improving the interoperability between HPC and Big Data languages and programming models 3.6 RELATED WORK 3.6.1 HPC and containers HPC workloads tend to be monolithic in nature so that each component and its dependencies must be present for running or compiling the code. In addition, if an update of any component is required, all modules will be affected. For this reason, the most common difficulty faced by endusers when creating and implementing scientific software is the installation and configuration of a framework with thousands of dependencies. Containers are a good way to self-contain an application and its dependencies in a controlled environment. Containers do not interfere with each other and allow to be deleted or updated without leaving any trace on the physical machine. They are an alternative to virtual machines while maintaining a similar level of isolation and showing a superior performance that in some cases is almost identical to the one obtained when executing natively on a real machine [2, 23]. IgnisHPC can be seen as an MPI application, so it can be efficiently executed inside containers as it was proven in several works. For example, running MPI applications on a containerized cluster using Docker on a cluster [15] or in the Cloud [102], or using Shifter [83] instead. Other works deal with the orchestration of Docker containers in an HPC environment. For instance, Higgins et al. [49] implemented a script based on SSH for the creation of an MPI environment inside a Docker container. Unlike IgnisHPC, this approach requires root privileges to modify the hosts configuration. Another paper introduces Scylla [82], a framework for deploying MPI jobs within Docker Containers using Apache Mesos. As explained in Section 3.1, Apache Mesos requires an orchestration framework such as Marathon or Singularity to be used. However, the authors, instead of considering a third party framework, implemented an ad-hoc solution. Their approach also requires root privileges. We must highlight that IgnisHPC supports all the functionalities included in Scylla without needing root permissions and is not limited to work with Apache Mesos. 3.6.2 Spark-MPI As a general-purpose framework, Spark has been widely used for many scientific applications and algorithms. However, there are examples from different areas such as linear algebra [39], genomics [1] or even data science [87] where Spark does not obtain the expected performance. One way to approach Big Data and HPC worlds is trying to boost the performance of wellestablished Big Data technologies when running on HPC systems. Although several related works were shown in Section 1.4, next we will analyze in more detail Spark-MPI [65, 66], which is the most similar to IgnisHPC. 68 Chapter 3. Improving the interoperability between HPC and Big Data languages and programming models Spark-MPI is a hybrid platform that combines Spark and MPI taking advantage of the MPI Exascale Process Management Interface (PMIx). A Spark-MPI application consists of a driver launched by Spark and a set of processes launched by MPI. These MPI processes connect to the Spark driver as workers, and they will be able to execute both RDD and MPI functions. Note that using MPI routines in Spark-MPI only makes sense when the data is processed using mapPartitions. That is the only way that Spark provides to work on a complete partition instead of on each element of the partition. The authors showed the benefits of Spark-MPI with deep learning algorithms and ptychographic and tomographic applications. However, SparkMPI has several limitations that we detailed next: – Spark-MPI requires a hybrid environment for Spark and MPI, which will be configured with their respective resource managers. For example, Mesos or Yarn for Spark and Hydra or Slurm for MPI. It is hard to find a system configured this way since resource managers cannot share the available hardware resources. Moreover, a Spark-MPI job should queue for both Spark and MPI queuing systems, and the requested resources may not be granted at the same time. On the other hand, IgnisHPC is a single application so the previous problems do not apply. – Due to its particular architecture and how Spark executors are launched, Spark-MPI loses the fault tolerance system provided by Spark. As a consequence, after any failure in the executors, all the job is lost. On the contrary, as we demonstrated in [79], Ignis and IgnisHPC are able to recover after a failure of a cluster node or some of the executors. If some data is lost, IgnisHPC has enough information about how it was derived in such a way that only those operations needed to recompute the corresponding portion of data are performed. – Spark-MPI can only be used with Python codes (PySpark). Although communications are performed using MPI, Spark-MPI would experience the same degradation in the performance observed for Spark when executing Python applications (see Section 3.5.2.1). Spark suffers performance issues since it requires sharing data outside the JVM through system pipes. Therefore, Spark-MPI performance results would be similar to those obtained by Spark when running Python codes (see, for example, the Minebench results in Figure 3.14). On the other hand, although the execution of pure MPI codes in Spark-MPI is discussed, the vast majority of MPI applications are implemented in C/C++ for which Spark and Spark-MPI does not have native support. – Spark-MPI is limited to access one partition at the same time per executor. Partition size is restricted to a maximum of 2 GB, which is related to the use of JVMs in Spark. Therefore, additional executors should be created in case more data is necessary, degrading the overall I/O performance. This restriction only applies to IgnisHPC for Java applications, but not for Python and C/C++. 69 Chapter 4. BigSeqKit: a Big Data approach to process FASTA/FASTQ files at scale –shuffle: sequences shuffling can be implemented using the IgnisHPC API function partitionByRandom. 4.1.4 Another implementation details In order to parallelize and integrate the SeqKit routines into IgnisHPC it is necessary to start considering the sequence parser. It takes a stream of characters in FASTA/Q format and generates a data structure with the sequence representation. In SeqKit, this stream can be represented by a file or the standard input. In BigSeqKit, this stream is implemented using the IgnisHPC iterators, which grant the users access to the data partitions. In this way, BigSeqKit will read the data from a file and split it in multiple partitions, which facilitates their parallel processing. As a result, the SeqKit command arguments that affect file processing will have no effect in BigSeqKit. For example, the --two-pass option, which reads a file multiple times instead of storing all the sequences in memory, does not make sense in BigSeqKit. Another important advantage of using IgnisHPC is how memory is handled. Users can choose a type of storage according to their particular case. For instance, if an input file is too large to be kept completely in the server memory, it could be stored compressed in memory or in disk. Performance would be lower, but it could be successfully processed. That scenario is not considered by SeqKit that simply would raise an ”out of memory” error. In particular, BigSeqKit supports the following IgnisHPC storage options: –In-Memory: it is the best performer since all data is stored in memory. It is the default option. –Raw memory: data is stored in a memory buffer using a serialized binary format. Extra memory consumption is minimal and the buffer is compressed by Zlib. –Disk: similar to raw memory but the buffer is stored as a POSIX file. Although the performance is significantly worse, it enables working with vast amounts of data that cannot be entirely kept in memory. On the other hand, rmdup,common and pair commands in SeqKit use hash functions to check duplicates. It is well-known that hash functions can produce the same result for different values. This event is commonly known as a hash collision. However, SeqKit does not check for collisions, so it is possible to generate incorrect results. BigSeqKit uses hashes to group sequences but then checks for collisions by comparing the real values. 76 Chapter 4. BigSeqKit: a Big Data approach to process FASTA/FASTQ files at scale 1 2 4 6 8 16 32 (2 nodes) 64 (4 nodes) 128 (8 nodes) Cores 0 20 40 60 Speedup D 1 D 2 D 3 1 2 4 6 8 16 32 (2 nodes) 64 (4 nodes) 128 (8 nodes) Cores 0 100 200 300 400 Speedup D1 D2 D3 Max. speedup seqkit D1 = 1.9, D2 = 7.2, D3 = 7.2 1 2 4 6 8 16 32 (2 nodes) 64 (4 nodes) 128 (8 nodes) Cores 0 5 10 15 Speedup D 1 D 2 D 3 1 2 4 6 8 16 32 (2 nodes) 64 (4 nodes) 128 (8 nodes) Cores 0 10 20 30 Speedup D1 D2 D3 in seqkit: Out of Memory 1 2 4 6 8 16 32 (2 nodes) 64 (4 nodes) 128 (8 nodes) Cores 0 5 10 15 20 25 Speedup D 1 D 2 D 3 1 2 4 6 8 16 32 (2 nodes) 64 (4 nodes) 128 (8 nodes) Cores 0 20 40 60 80 Speedup D 1 D 2 D 3 Figure 4.3: Speedups obtained by BigSeqKit with respect to the same commands executed using SeqKit. Note that grep was partially parallelized in SeqKit. 4.2 PERFORMANCE RESULTS In this section we analyze the performance results obtained by BigSeqKit. Experiments were performed using up to 8 computing nodes of the Marconi1004supercomputer installed at CINECA (Italy). Each node contains one IBM Power9 AC922 @3.1GHz processor and 256 GB of memory. The main characteristics of the three FASTA/FASTQ files considered in the evaluation are the following: –D1(Homo_sapiens.GRCh38.dna_sm.toplevel.fa - FASTA - 59.7 GB): Human genome obtained from http://www.ensembl.org. Number of sequences: 639, Minimum length: 970, Average length: 98.8M, Maximum length: 248.9M. –D2(SRR642648_1.filt.fq - FASTQ - 24.1 GB): DNA sequences obtained from https://www.internationalgenome.org [35]. Number of sequences: 98.7M, Minimum length: 100, Average length: 100, Maximum length: 100. –D3(ERR4667750.fq - FASTQ - 79.1 GB): DNA sequences obtained from https://www.internationalgenome.org [35]. Number of sequences: 318.1M, Minimum length: 101, Average length: 101, Maximum length: 101. Figure 4.3 shows the results in terms of the speedups obtained by BigSeqKit with respect to the sequential execution (1 core) of the same SeqKit command. Each result was computed as the 4https://www.hpc.cineca.it/content/hardware [accessed September, 2022] 77 Chapter 4. BigSeqKit: a Big Data approach to process FASTA/FASTQ files at scale 1 2 4 6 8 16 32 (2 nodes) 64 (4 nodes) 128 (8 nodes) D1 SeqKit 74.1 – – – – – – – – BigSeqKit 99.6 51.1 26.5 18.5 14.2 7.9 7 6.6 6.5 [11.2×] D2 SeqKit 55.3 – – – – – – – – BigSeqKit 70.2 36.5 18.9 13.3 10 6.8 5.5 5.2 5.1 [10.8×] D3 SeqKit 116.8 – – – – – – – – BigSeqKit 140.7 71.3 36.5 25.5 18.8 10.9 10 9.1 8.5 [13.7×] Table 4.2: Execution times (seconds): seq command. Highlighted in blue, fastest time and number of times faster than sequential SeqKit. median of five experiments. As we mentioned previously, to improve the usability and facilitate the adoption of BigSeqKit, it implements the same command interface than SeqKit. As example to illustrate the benefits of our tool, we show results for the following utilities (see Table 4.1 for a complete list of commands): seq transforms sequences (extract ID, filter by length, etc.) and removes gaps, grep searches sequences by ID/name/sequence, rmdup removes duplicated sequences, replace replaces a name/sequence using a regular expression, concatenate concatenates sequences with same ID from multiple files and sort sorts sequences by ID/name/sequence/length. According to the results, several conclusions can be made. BigSeqKit outperforms SeqKit for all the commands considered. Speedups higher than one are always obtained using two or more cores on a single node. In particular, BigSeqKit is from 8.1×(seq -D2) to 49.1×(grep - D1) faster than SeqKit (using 16 cores). Note that grep was partially parallelized in SeqKit, but its best speedup only reaches 7.2 (D2and D3). On the other hand, BigSeqKit can be executed using several nodes on a cluster. This is not the case of SeqKit, whose scalability is limited to use a few threads on a single node. As a consequence, BigSeqKit is from 10.8×(seq -D2) to 387.3×(grep -D1) faster than SeqKit when considering 8 nodes. It means that, for example, SeqKit takes about 5 hours to find 10 short sequences in D1, while BigSeqKit only requires 1.4 minutes. An additional advantage of using BigSeqKit is that if the memory consumed by a particular file is too high, it can be split into several nodes. This is the case of using sort to process our largest file, D3. Both BigSeqKit and SeqKit consume too much memory that prevent the execution of this command on a single node of our cluster. BigSeqKit deals with that problem using two or more nodes to sort the input file. In this way, BigSeqKit is able to decrease the time to sort D3from 20 minutes (2 cores - 2 nodes) to barely 1 minute (128 cores - 8 nodes). Tables from 4.2 to 4.7 display the execution times of BigSeqKit and SeqKit when running seq,grep,rmdup,replace,concat and sort utilities, respectively. It is important to highlight that speedups shown in Figure 4.3 were computed using the times included in these tables. Results confirm the benefits of using BigSeqKit with respect to SeqKit in terms of performance 78 Chapter 4. BigSeqKit: a Big Data approach to process FASTA/FASTQ files at scale 1 2 4 6 8 16 32 (2 nodes) 64 (4 nodes) 128 (8 nodes) D1 SeqKit 32,650 23,926 17,273 out of memory out of memory out of memory – – – BigSeqKit 10,140 5,100 2,587 1,721 1,285 665 330 163 84 [387×] D2 SeqKit 21,762 15,240 10,550 6,731 5,645 3,025 – – – BigSeqKit 7,509 3,800 1,929 1,300 984 503 259 139 73 [299×] D3 SeqKit 67,804 49,592 33,221 21,152 17,771 9,394 – – – BigSeqKit 20,568 10,336 5,224 3,486 2,677 1,376 692 374 201 [337×] Table 4.3: Execution times (seconds): grep command. Highlighted in blue, fastest time and number of times faster than sequential SeqKit. 1 2 4 6 8 16 32 (2 nodes) 64 (4 nodes) 128 (8 nodes) D1 SeqKit 522.6 – – – – – – – – BigSeqKit 778.5 408.2 212.7 150.4 110 57.8 32.3 17.3 9.8 [53.3×] D2 SeqKit 383.5 – – – – – – – – BigSeqKit 523.5 273.7 146.6 102.2 77.2 41.2 22.9 12.9 7.5 [51.1×] D3 SeqKit 918.1 – – – – – – – – BigSeqKit 1,200.5 630.2 330.8 230.5 174.5 92.3 49.5 26.8 14.9 [61.6×] Table 4.4: Execution times (seconds): rmdup command. Highlighted in blue, fastest time and number of times faster than sequential SeqKit. 1 2 4 6 8 16 32 (2 nodes) 64 (4 nodes) 128 (8 nodes) D1 SeqKit 135 – – – – – – – – BigSeqKit 152.4 79.2 38.9 26.5 21 10.7 6.3 6 [22.5×]6.4 D2 SeqKit 303.8 – – – – – – – – BigSeqKit 320.1 165.3 86.2 56.8 45.6 24.6 13.3 7.5 5.5 [55.2×] D3 SeqKit 843.2 – – – – – – – – BigSeqKit 861.2 444.2 230.9 155.1 117.8 61.9 32.9 18.2 10.6 [79.6×] Table 4.5: Execution times (seconds): replace command. Highlighted in blue, fastest time and number of times faster than sequential SeqKit. 1 2 4 6 8 16 32 (2 nodes) 64 (4 nodes) 128 (8 nodes) D1 SeqKit 150.1 – – – – – – – – BigSeqKit 182.4 99.4 51.9 35.4 28 16.9 12.7 8.4 8.3 [18.1×] D2 SeqKit 221.5 – – – – – – – – BigSeqKit 260.4 140.5 76.5 50.9 42.5 24.1 19.2 15.2 12.6 [17.6×] D3 SeqKit 620.5 – – – – – – – – BigSeqKit 700.9 375.5 199 131 106.4 61.5 50.1 39.2 30.5 [20.3×] Table 4.6: Execution times (seconds): concat command. Highlighted in blue, fastest time and number of times faster than sequential SeqKit. 79 Chapter 4. BigSeqKit: a Big Data approach to process FASTA/FASTQ files at scale 1 2 4 6 8 16 32 (2 nodes) 64 (4 nodes) 128 (8 nodes) D1 SeqKit 143.2 – – – – – – – – BigSeqKit 165.4 87 44 29.2 22 14.1 12.1 11.5 11.4 [12.6×] D2 SeqKit 1,127.6 – – – – – – – – BigSeqKit 1,260.5 659.5 348.5 240.9 186.5 102.6 62.4 44.5 38.4 [29.4×] D3 SeqKit out of memory – – – – – – – – BigSeqKit out of memory 1,239∗649.7∗446.9∗345.7∗185.9∗111.1 75.1 64.8 Table 4.7: Execution times (seconds): sort command. Times with an asterisk were obtained using two nodes. Highlighted in blue, fastest time and number of times faster than sequential SeqKit. and scalability. Note that even considering a single node (up to 16 cores), BigSeqKit always outperforms SeqKit when using more than 1 core. Tables also highlight the fastest times achieved by BigSeqKit and how many times faster are those with respect to sequential SeqKit. Independently of the command considered, BigSeqKit is able to reduce the processing time at least one order of magnitude. For example, the maximum processing time observed for BigSeqKit is only 201 seconds (grep command - D3, see Table 4.3), while the corresponding SeqKit time is 9,394 seconds. In other words, it means that SeqKit takes 2.6 hours to perform the same operation. 4.3 FINAL REMARKS In this chapter we showed the potential of IgnisHPC to boost the performance and scalability of a bioinformatics toolkit to process FASTA and FASTQ files. Current state of the art tools such as SeqKit are not ready for processing and manipulating very large files because all of them are mainly based on sequential processing. To that end, we have presented BigSeqKit, which parallelizes and optimizes the SeqKit routines using IgnisHPC. Since SeqKit was programmed in Go, IgnisHPC was extended to support that language. As a consequence, IgnisHPC is nowadays the first parallel computing framework that supports Go. Regarding the experimental results, BigSeqKit clearly outperforms SeqKit, being from 11× to 387×faster when using 8 nodes of a cluster. Even when only one computing node is considered, BigSeqKit obtains speedups from 8×to 48×. It means, for instance, that processes that require hours in SeqKit take just a few minutes in BigSeqKit. As future work we plan to add also the remainder SeqKit commands not included in BigSeqKit:sliding,faidx,sana,fx2tab,tab2fx,convert,amplicon,fish,split,split2, restart and mutate. Note that all of them are independent routines, so their implementation using IgnisHPC will be straightforward. 80 5 Conclusions The unification of High Performance Computing and Big Data has received increasing attention in the last years. It is a common belief that exascale computing and Big Data are closely associated since HPC requires processing large-scale data from scientific instruments and simulations. But, at the same time, it was observed that tools and cultures of HPC and Big Data communities differ significantly [47]. As it was explained, one of the most important sources of divergence comes from the differences between their software ecosystems. In this way, HPC applications have traditionally been based on MPI to support inter-node parallel execution, and based on OpenMP or other alternatives to exploit intra-node parallelism. However, Big Data programming models are based on interfaces like Hadoop or Spark. In addition to different programming models, programming languages also differ between both communities: being Fortran and C/C++ the most common languages in HPC applications, and Java, Scala, or Python being the most common languages in Big Data applications. This divergence between programming models and languages sets out a convergence problem, not only related to the interoperability of the applications but also to the interoperability between data formats from different programming languages [11]. In this scenario, we need to consider how to build end-to-end workflows where, for example, simulations can be MPI applications written in Fortran or C/C++, and the analytics codes can be written in Java or Python (maybe parallelized by a Big Data framework). In this thesis, we address the research challenge of bridging the gap between Big Data and HPC ecosystems by unifying the development, combination and execution of HPC and Big Data parallel tasks using different languages and programming models in the same computing framework. This means a different path in order to reach the convergence with respect to other existent works in the literature. Those works try to boost the performance of well-established Big Data technologies when running on HPC systems, or extend HPC technologies to support Big Data tasks. None of them search for a unique processing engine for Big Data and HPC workloads. Chapter 2 presented our first attempt to build an efficient and scalable multi-language framework, named Ignis. The main conclusions and contributions derived from that work were the following: Chapter 5. Conclusions – To the best of our knowledge, Ignis was the first Big Data framework with native multilanguage support including both JVM and non-JVM-based languages. It supports C/C++, Python and Java. In this way, users might combine in the same application the benefits of implementing each computational task in the best suited programming language. – Ignis outperformed the state-of-the-art framework Spark in terms of performance and scalability running applications that represent the most typical algorithmic patterns in Big Data and scientific computing. For example, Ignis was 2.2×faster than Spark for an application that chains several Map operations, showing also a very good behavior in terms of strong and weak scaling. We can find another illustrative example when considering the K-Means machine learning algorithm. K-Means is a good example of an iterative MapReduce application pattern. In this case, Ignis was on average from 1.6×to 2.1×faster than Spark. Similar performance results were also observed for Sort and the Conjugate Gradient. Note that the latter is considered one of the most important computational kernels in HPC applications. – Contrary to the developers community belief, we demonstrated that Python is not natively supported by Spark since data transfers between the JVM and external processes degrade noticeably the overall performance. – Ignis deals efficiently with fault tolerance in such a way that if some data is lost, only those operations needed to recompute the corresponding portion of data are performed. – Ignis provides a simple and powerful user API to develop applications. To facilitate the adoption from the Big Data community, the Ignis API was inspired by the Spark API in such a way that Ignis codes are easily understandable by users who are familiar with Spark. – Finally, Ignis is fully developed inside Docker containers, which isolates the execution environment from the physical system and avoids dependency problems. It includes a custom resource manager to assign hardware resources and launch containers. However, despite Ignis makes significant contributions to the convergence, it has some limitations that prevent it to be a universal framework for HPC and Big Data applications. The most important drawback is related to how communications are performed. In particular, Ignis is restricted to use TCP sockets for inter-node communication. This would cause an important problem of scalability when using a high number of computing nodes in the current HPC systems, since Ignis could not take advantage of the most advanced networks (e.g., Infiniband or Slingshot). 82 Chapter 5. Conclusions To solve the above and other issues regarding Ignis, we introduced IgnisHPC in Chapter 3. IgnisHPC inherits some characteristics from Ignis such as fault tolerance, but it was completely redesigned with the aim of improving performance, scalability and productivity. We can summarize the main contributions of IgnisHPC as follows: – Just like Ignis, IgnisHPC supports natively both JVM and non-JVM-based languages. Therefore, applications can be implemented using one or several programming languages following an API inspired by Spark’s one. – IgnisHPC uses MPI as backbone technology, which allows the framework to support many communication models and network architectures. In this way, it covers the characteristics of the vast majority of Big Data and HPC clusters. In addition, MPI applications and libraries can be directly executed in an efficient way in IgnisHPC. In this way, most of the HPC scientific applications, which in many cases contain tens of thousands of lines of code, do not have to be ported to a new API or programming model. For instance, to execute the parallel algebraic multigrid solver AMG in IgnisHPC only requires about 40 extra lines of code, while the whole application totalized more than 65,000. It is important to highlight that, to the best of our knowledge, nowadays there is no other computing framework with that feature. – In IgnisHPC, MPI codes can be easily combined with typical MapReduce operations to create hybrid applications. As a result, developers can use different programming models inside the same application according to the type of the computational tasks (dataintensive or compute-intensive). – A thorough experimental evaluation was carried out to demonstrate the benefits of IgnisHPC in terms of performance and productivity. The study showed that IgnisHPC clearly outperforms Spark and Ignis when considering different types of Big Data application patterns. In terms of performance, IgnisHPC is from 1.1×to 3.9×faster than Spark, and from 1.1×to 1.3×faster than Ignis. However, there are also another benefits. For instance, the memory consumption in IgnisHPC was optimized. In this way, unlike Ignis, IgnisHPC is capable of processing extremely large datasets (e.g., TeraSort benchmark). In addition, the IgnisHPC API was extended to support, among others, graph processing algorithms. – We also proved that running MPI (and hybrid MPI+OpenMP) applications from IgnisHPC is easy and as efficient as executing them natively. In particular, performance differences between native MPI and IgnisHPC are really small, always lower than 1.7%. – IgnisHPC is also fully containerized to avoid dependencies and supports some of the most well-known resource and scheduler managers (e.g., Mesos, Nomad, etc.). 83 Chapter 5. Conclusions Finally, Chapter 4 introduced BigSeqKit, whose contribution regarding IgnisHPC is twofold. On one hand, BigSeqKit was the first complete scientific application implemented in IgnisHPC. It is a parallel toolkit based on SeqKit to process FASTA and FASTQ files at scale specially designed for speed and scalability. Note that manipulating these files efficiently is essential to analyze and interpret data in any genomics pipeline. On the other, IgnisHPC was extended to support Go programming language. There is no other Big Data-HPC framework that natively supports this language. The experimental results on a cluster demonstrated the benefits of our proposal, achieving speedups from 11×to 387×with respect to SeqKit. Note that SeqKit only can be executed on a single node. As a consequence, processes that require several hours in SeqKit take only a few minutes in BigSeqKit. Therefore, we can conclude that the objectives of the thesis detailed in Section 1.5 have been successfully fulfilled. The technologies and solutions developed mean an important step forward toward the real convergence of HPC and Big Data worlds. 5.1 FUTURE WORK We can find several research directions to continue bridging the gap between HPC and Big Data technologies: – One interesting line of research is related to checkpointing. Some experimental technologies can freeze a running container (or an individual application) and checkpoint its state to disk. In this way, a container (or group of containers) can be restored to the exact execution point before it was frozen using the stored information. This has enormous implications for IgnisHPC. First, IgnisHPC, just like Spark, already has a recovery system, but they share the same weakness, the driver. If the driver fails, the complete job will fail, and all progress that has not been written to disk will be lost. In this way, if the framework performed a checkpoint on the driver, this limitation could be overcome. Additionally, calls to MPI applications and libraries from IgnisHPC are not covered by the fault tolerant mechanism. As a consequence, any error that causes a crash in the MPI processes would require the MPI task to be started over again. However, if we checkpoint the MPI processes, we could restart them from those points. – Although IgnisHPC currently supports several programming languages (Python, C/C++, Java and Go), it would be interesting to include support for R, which is a programming language and environment with a focus on statistical analysis. As such, due to its poor efficiency, R is not utilized to write libraries for compute-intensive applications, but libraries implemented in other languages are exported for using from R. The strength of R is to offer a simple programming environment accessible to the entire scientific community. IgnisHPC could implement an R interface using one of the already supported 84 Chapter 5. Conclusions languages and export its functionality. In this way, libraries such as BigSeqKit could also be exported for the R users community. – Another interesting direction for improvement is the GUI. Currently the user interaction with IgnisHPC is minimal. Users interact with the framework to launch jobs and the only feedback that the users receive are job logs stored as plain text files. Nowadays, other frameworks like Spark or Hadoop have a web interface where the job information can be found. As mentioned in previous chapters, a job is composed by tasks, so it would be interesting to display real-time information about tasks and the system resources usage. As a result, the task tracking would allow an easier debugging process and it could be very helpful to identify inefficient phases. 85 Bibliography [78] C. Piñeiro, J. M. Abuín, and J. C. Pichel. Perldoop2: A Big Data-Oriented Source-toSource Perl-Java Compiler. In Proc. of the IEEE Intl. Conf. on Big Data Intelligence and Computing (DataCom), pages 933–940, 2017. [79] César Piñeiro, Rodrigo Martínez-Castaño, and Juan C. Pichel. Ignis: An efficient and scalable multi-language Big Data framework. Future Generation Computer Systems, 105:705–716, 2020. [80] César Piñeiro and Juan C. Pichel. A unified framework to improve the interoperability between HPC and Big Data languages and programming models. Future Generation Computer Systems, 134:123–139, 2022. [81] Dino Quintero, Luis Bolinches, Puneet Chaudhary, Willard Davis, Steve Duersch, Carlos Henrique Fachim, Andrei Socoliuc, Olaf Weiser, et al. IBM Spectrum Scale (formerly GPFS). IBM Redbooks, 2017. [82] Pankaj Saha, Angel Beltre, and Madhusudhan Govindaraju. Scylla: a Mesos Framework for Container Based MPI Jobs. CoRR, abs/1905.08386, 2019. [83] Pankaj Saha, Angel Beltre, Piotr Uminski, and Madhusudhan Govindaraju. Evaluation of Docker Containers for Scientific Workloads in the Cloud. In Proc. of the Practice and Experience on Advanced Research Computing, 2018. [84] Jason Sanders and Edward Kandrot. CUDA by example: an introduction to generalpurpose GPU programming. Addison-Wesley Professional, 2010. [85] Peter Sanders, Sebastian Lamm, Lorenz Hübschle-Schneider, Emanuel Schrade, and Carsten Dachsbacher. Efficient parallel random sampling—vectorized, cache-efficient, and online. ACM Transactions on Mathematical Software (TOMS), 44(3):1–14, 2018. [86] Aparna Sasidharan and Marc Snir. MiniAMR - A miniapp for Adaptive Mesh Refinement. Technical report, 2016. [87] Manvi Saxena, Shweta Jha, Saba Khan, John Rodgers, Peggy Lindner, and Edgar Gabriel. Comparison of MPI and Spark for Data Science Applications. In IEEE Int. Parallel and Distributed Processing Symposium Workshops (IPDPSW), pages 682–690, 2020. [88] Aamir Shafi, Bryan Carpenter, and Mark Baker. Nested parallelism for multi-core HPC systems using java. J. Parallel Distributed Comput., 69(6):532–545, 2009. [89] Wei Shen, Shuai Le, Yan Li, and Fuquan Hu. SeqKit: a cross-platform and ultrafast toolkit for FASTA/Q file manipulation. PloS One, 11(10):e0163962, 2016. 92 Bibliography [90] Juwei Shi and et al. Clash of the Titans: MapReduce vs. Spark for Large Scale Data Analytics. Proc. VLDB Endowment, 8(13):2110–2121, sep 2015. [91] Vladimir V Stegailov, Nikita D Orekhov, and Grigory S Smirnov. Hpc hardware efficiency for quantum and classical molecular dynamics. In International conference on parallel computing technologies, pages 469–473. Springer, 2015. [92] Giselle van Dongen and Dirk Van den Poel. Evaluation of stream processing frameworks. IEEE Transactions on Parallel and Distributed Systems, 31(8):1845–1858, 2020. [93] Vinod Kumar Vavilapalli and et al. Apache Hadoop YARN: Yet Another Resource Negotiator. In Proc. of the 4th Annual Symposium on Cloud Computing, pages 5:1–5:16. ACM, 2013. [94] Yong-Xian Wang, Li-Lun Zhang, Wei Liu, Yong-Gang Che, Chuan-Fu Xu, Zheng-Hua Wang, and Yu Zhuang. Efficient parallel implementation of large scale 3d structured grid cfd applications on the tianhe-1a supercomputer. Computers & Fluids, 80:244–250, 2013. [95] D. Wheeler. SLOCCount. http://www.dwheeler.com/sloccount. [Online; accessed November, 2021]. [96] Tom White. Hadoop: The Definitive Guide. O’Reilly Media, Inc., 4th edition, 2015. [97] Reynold S. Xin, Joseph E. Gonzalez, Michael J. Franklin, and Ion Stoica. GraphX: A Resilient Distributed Graph System on Spark. In 1st International Workshop on Graph Data Management Experiences and Systems. ACM, 2013. [98] Luna Xu, Min Li, and Ali R. Butt. GERBIL: MPI+YARN. In 2015 15th IEEE/ACM International Symposium on Cluster, Cloud and Grid Computing, pages 627–636, 2015. [99] Jaewon Yang and Jure Leskovec. Defining and Evaluating Network Communities based on Ground-truth, 2012. [100] K. Ye and Y. Ji. Performance Tuning and Modeling for Big Data Applications in Docker Containers. In Int. Conference on Networking, Architecture, and Storage (NAS), pages 1–6, 2017. [101] Andy B. Yoo, Morris A. Jette, and Mark Grondona. Slurm: Simple linux utility for resource management. In JSSPP, 2003. [102] Andrew J. Younge, Kevin Pedretti, Ryan E. Grant, and Ron Brightwell. A Tale of Two Systems: Using Containers to Deploy HPC Applications on Supercomputers and Clouds. 93 Bibliography In IEEE Int. Conference on Cloud Computing Technology and Science (CloudCom), pages 74–81, 2017. [103] Matei Zaharia, Mosharaf Chowdhury, Michael J. Franklin, Scott Shenker, and Ion Stoica. Spark: Cluster Computing with Working Sets. In Proc. of the 2nd USENIX Conf. on Hot Topics in Cloud Computing (HotCloud), pages 10–10, 2010. [104] Matei Zaharia and et al. Resilient Distributed Datasets: A Fault-tolerant Abstraction for In-memory Cluster Computing. In Proceedings of the 9th USENIX Conference on Networked Systems Design and Implementation, pages 2–2. USENIX Association, 2012. 94 A Getting Started with IgnisHPC In previous chapters, we have introduced some application examples using the Ignis/IgnisHPC API. For example, Chapters 2 and 3 showed implementations of Wordcount and the Transitive Closure algorithm, respectively. In addition, an implementation of a hybrid application with MPI using the LULESH library when considering IgnisHPC was also explained. In this appendix, we will introduce a guide to help users in order to install and use our framework. In particular, Section A.1 shows the installation and configuration steps required to run a simple job on IgnisHPC. Section A.2 explains how to execute the BigSeqKit application step by step from scratch. A.1 IGNISHPC USAGE In this section, we summarize the minimum steps to execute a simple job in IgnisHPC. This information is also available as the IgnisHPC online documentation for users1. IgnisHPC is a modularized Docker framework consisting of multiple source code repositories. The framework is Open Source, so all repositories can be found in GitHub2. A.1.1 Requirements As we mentioned, IgnisHPC is a dockerized framework, so all the system modules are executed inside Docker containers. On the other hand, IgnisHPC external dependencies such as schedulers or image storage can be installed independently, but IgnisHPC includes a dockerized version of them. Therefore, the minimum requirements to run IgnisHPC are: – Docker: It must be installed and accessible. We recommend using the newest version available. Please refer to https://docs.docker.com/get-docker/ for instructions. – Python3: Available by default in most Linux distributions. It is used to execute a deploy script and simplify the installation of IgnisHPC and its dependencies. 1https://ignishpc.readthedocs.io [accessed September, 2022] 2https://github.com/ignishpc [accessed September, 2022] Appendix A. Getting Started with IgnisHPC – Pip: The deploy script is available as a pip package, although it can be downloaded from the source code repository. In any case, using pip is the easiest way to install IgnisHPC. – Git (optional): The git binary is required for building IgnisHPC images from repositories. A.1.2 Installation IgnisHPC can be installed just using the following command: $ pip install ignishpc Once the command is executed, we can check that we have the ignis-deploy command available in the path. This process must be carried out in all the computing nodes where IgnisHPC is going to be executed. A.1.3 Creating IgnisHPC containers Images of IgnisHPC containers are not available for download. They must be built in your runtime environment. This allows the creation of a custom development environment with complete isolation from other framework installations even on the same machines. In this example we will use a local Docker repository and ignishpc as the base path for all images. $ ignis-deploy registry start --default This command will launch a Docker repository that will be available on port 5000. –default parameter indicates that other calls to ignis-deploy should use this repository as the default source. Once executed, we will receive the following warning message: info: add '{"insecure-registries" : [ "myhost:5000" ]}' to /etc/docker/daemon.json and restart docker daemon service use myhost:5000 to refer the registry The Docker registry requires a certificate for validation. We can create one locally or add ”insecure −registries” : [”myhost : 5000”]to /etc/docker/daemon.json with the aim of forcing Docker to use an insecure one. Do not forget to restart the Docker daemon to reload the configuration. Once the registration is available, we can proceed with the creation of the IgnisHPC images: $ ignis-deploy images build --full --sources \ https://github.com/ignishpc/dockerfiles.git \ 96 Appendix A. Getting Started with IgnisHPC https://github.com/ignishpc/backend.git \ https://github.com/ignishpc/core-cpp.git \ https://github.com/ignishpc/core-python.git \ https://github.com/ignishpc/core-go.git The first two repositories are essential for the construction of the base images. The command can be executed in several phases, it is not necessary to specify all the repositories in the same execution. The only restriction is that if you use the –full parameter, which creates an extra image with all the core repositories, you must have all the cores together. This allows users to create an image that can run Python, C/C++ and Go codes in the same container. Finally, IgnisHPC files can be extracted from an image with: $ docker run --rm -v $(pwd):/target <ignis-image> ignis-export-all /target The result will not be executable but can be used in the application development. A.1.4 Deploying IgnisHPC Containers Once all images are created, it is necessary to deploy the containers. IgnisHPC jobs are launched using the submitter module, but to use a cluster it requires a resource and scheduler manager such as Nomad or Mesos. Alternatively, IgnisHPC can also be launched as a Slurm job in an HPC cluster. In that case Docker is replaced by Singularity, so a different submitter must be used. Next we show examples deploying the containers locally (no manager is needed), and using Nomad, Mesos and Slurm resource managers. –Docker (only local): Submitter using an http endpoint: $ ignis-deploy submitter start --dfs <working-directory-path> \ --scheduler docker tcp://myhost:2375 Submitter using a Unix-socket: $ ignis-deploy submitter start --dfs <working-directory-path> \ --scheduler docker /var/run/docker.sock \ --mount /var/run/docker.sock /var/run/docker.sock –Nomad: Master node: $ ignis-deploy nomad start --password 1234 --volumes <working-directory-path\*> 97 Appendix A. Getting Started with IgnisHPC Worker nodes: $ ignis-deploy nomad start --password 1234 --join myhost1 --default-registry myhost1:5000 Submitter: $ ignis-deploy submitter start --dfs <working-directory-path/*> \ --scheduler nomad http://myhostX:4646 *The working directory must be available on all nodes via NFS (Network File System) or a DFS (Distributed File System). (Only required for working with files) –Mesos: Zookeeper is required by Mesos: $ ignis-deploy zookeeper start --password 1234 Master node: $ ignis-deploy mesos start -q 1 --name master -zk zk://master:2281 \ --service [marathon | singularity] --port-service 8888 Worker nodes: $ ignis-deploy mesos start --name nodoX -zk zk://master:2281 \ --port-service 8888 --default-registry master:5000 Submitter: $ ignis-deploy submitter start --dfs <working-directory-path/*> \ --scheduler [marathon | singularity] http://master:8888 * The working directory must be available on all nodes via NFS (Network File System) or a DFS (Distributed File System). (Only required for working with files) –Slurm: The ignis-slurm submitter can be obtained from ignishpc/slurm-submitter with: $ docker run --rm -v $(pwd):/target ignishpc/slurm-submitter ignis-export /target This submitter will allow you to launch IgnisHPC on a cluster as a non-root user and without docker. IgnisHPC Docker images can be converted to singulairty image files with: 98 Appendix A. Getting Started with IgnisHPC $ ignis-deploy images singularity [--host] ignishpc/full ignis_full.sif The basic syntax of ignis-slurm is the same as the later shown ignis-submit, but a first parameter with job-time must be passed to be requested to Slurm. The time can be specified in any format supported by Slurm. For example, a 10 minute job should start with: $ ignis-slurm 00:10:00 .... In addition, help text can be displayed using: $ ignis-slurm --help A.1.4.1 Launching the first job The first step to launch a job is to connect to the submiter container. The default password is ignis, but we can change it inside the container or choose one when launching the submitter: $ ssh root@myhost -p 2222 The code we will use as an example is the classic Wordcount application, which can be seen below: 1#!/usr/bin/python 2 3import ignis 4 5# Initialization of the framework 6ignis.Ignis.start() 7# Resources/Configuration of the cluster 8prop = ignis.IProperties() 9prop["ignis.executor.image"] = "ignishpc/python" 10 prop["ignis.executor.instances"] = "1" 11 prop["ignis.executor.cores"] = "2" 12 prop["ignis.executor.memory"] = "1GB" 13 # Construction of the cluster 14 cluster = ignis.ICluster(prop) 15 16 # Initialization of a Python Worker in the cluster 17 worker = ignis.IWorker(cluster, "python") 18 # Task 1 - Tokenize text into pairs ('word', 1) 19 text = worker.textFile("text.txt") 20 words = text.flatmap(lambda line: [(word, 1) for word in line.split()]) 21 # Task 2 - Reduce pairs with same word and obtain totals 22 count = words.toPair().reduceByKey(lambda a, b: a + b) 23 # Print results to file 24 count.saveAsTextFile("wordcount.txt") 25 26 # Stop the framework 27 ignis.Ignis.stop() Figure A.1: IgnisHPC minimal job example with wordcount 99 Appendix A. Getting Started with IgnisHPC In order to run it, we need to create a file containing the Figure A.1 (text.txt) and store it in the working directory. By default the submitter sets the working directory to /media/dfs. All relative paths used in the source code are resolved using this working directory, so /media/dfs/text.txt is an alias of text.txt. Finally, we can execute our code using the submitter: $ ignis-submit ignishpc/python python3 driver.py or: $ ignis-submit ignishpc/python ./driver.py When the execution has finished, we can see the result of the execution in wordcount.txt located in the working directory. If we want to check the execution logs, we must navigate to the scheduler web or use docker log in case of using docker directly. A.1.5 Launching without a container The ignis-submit can also be used outside the submitter container, for example where permanent containers are not allowed: $ docker run --rm -v $(pwd):/target ignishpc/submitter ignis-export /target This command will create an ignis folder in the current directory with everything needed to run the submitter. The ignis-deploy command configures the submitter container, but when there is no container, we must set the configuration manually. The submitter needs a dfs and a scheduler, as ignis-deploy showed, these can be defined as environment variables or in ignis/etc/ignis.conf property file. # set current directory as job directory (ignis.dfs.id in ignis.conf) export IGNIS_DFS_ID=$(pwd) # set docker as scheduler (ignis.scheduler.type in ignis.conf) export IGNIS_SCHEDULER_TYPE=docker # set where docker is available (ignis.scheduler.url in ignis.conf) export IGNIS_SCHEDULER_URL=/var/run/docker.sock The above example could be launched as follows: $ ./ignis/bin/ignis-submit ignishpc/python ./driver.py 100 Appendix A. Getting Started with IgnisHPC A.2 BIGSEQKIT USER'S GUIDE 1package main 2 3import ( 4"ignis/driver/api" 5"bigseqkit" 6) 7// Auxiliary function for error checking, driver aborts if err is not nil 8func check[T any](e T, err error) T {if err != nil {panic(err) } else {return e }} 9 10 func main() { 11 // Initialization of the framework 12 check(0, api.Ignis.Start()) 13 // Stop the framework when main ends 14 defer api.Ignis.Stop() 15 // Resources/Configuration of the cluster 16 prop := check(api.NewIProperties()) 17 check(prop.Set("ignis.executor.image","ignishpc/full")) 18 check(prop.Set("ignis.executor.instances","2")) 19 check(prop.Set("ignis.executor.cores","4")) 20 check(prop.Set("ignis.executor.memory","6GB")) 21 // Construction of the cluster 22 cluster := check(api.NewIClusterProps(prop)) 23 // Initialization of a Go Worker 24 worker := check(api.NewIWorkerDefault(cluster, "go")) 25 // Sequence reading 26 seqs := check(bigseqkit.ReadFASTA("mysequences.fa", worker)) 27 // Deletion of duplicated sequences 28 rm_opts := &bigseqkit.SeqKitRmDupOptions{} 29 rm_opts.BySeq(true) 30 u_seqs := check(bigseqkit.RmDup(seqs, rm_opts)) 31 // Sort by ID 32 sort_opts := &bigseqkit.SeqKitSortOptions{} 33 u_sorted_seqs := check(bigseqkit.Sort(u_seqs, sort_opts)) 34 // Save the result 35 check(0, u_sorted_seqs.SaveAsTextFile("result.fa")) 36 } Figure A.2: BigSeqKit example using rmdup and sort commands. Figure A.2 shows an example of a IgnisHPC driver code implemented in Go that executes rmdup and sort commands. BigSeqKit has been created as a library, so it only needs to be imported to be used within IgnisHPC (line 5). Instead of commands from terminal like SeqKit,BigSeqKit utilities are functions that can be called from a driver code. Note that their names and arguments are exactly the same than those included in SeqKit, which can be found in https://bioinf. shenwei.me/seqkit/usage. Functions in BigSeqKit do not use files as input, they use DataFrames instead, an abstract representation of parallel data used by IgnisHPC (similar to RDDs in Spark [103]). Parameters are grouped in a data structure where each field represents the long names of a parameter. Note that BigSeqKit functions can be linked (like system pipes using ”|”), so the DataFrame generated by one can be used as input to another. In this way, integrate BigSeqKit routines in a more complex code is really easy. 101 List of Tables Tab. 3.1 Example of some IDataFrame functions supported by IgnisHPC. . . . . . . . 48 Tab. 3.2 Operations used in each Big Data application. Minebench (MB), Terasort (TS), K-means (KM), PageRank (PR) and Transitive Closure (TC). Operators annotated with I are specific only to IgnisHPC. . . . . . . . . . . . . . . . . . . . 58 Tab. 3.3 Summary of the IgnisHPC performance results for all the Big Data applications considering the maximum number of cores and the best Ignis and Spark implementation (in case there is more than one). Between brackets the programming language/s used in the IgnisHPC implementation. . . . . . . . . . 63 Tab. 3.4 MPI calls used for communications in the HPC applications. . . . . . . . . . 63 Tab. 3.5 SLOC of the HPC applications. . . . . . . . . . . . . . . . . . . . . . . . . . 67 Tab. 4.1 List of commands included in both BigSeqKit and SeqKit. .......... 73 Tab. 4.2 Execution times (seconds): seq command. Highlighted in blue, fastest time and number of times faster than sequential SeqKit................ 78 Tab. 4.3 Execution times (seconds): grep command. Highlighted in blue, fastest time and number of times faster than sequential SeqKit................ 79 Tab. 4.4 Execution times (seconds): rmdup command. Highlighted in blue, fastest time and number of times faster than sequential SeqKit................ 79 Tab. 4.5 Execution times (seconds): replace command. Highlighted in blue, fastest time and number of times faster than sequential SeqKit. ............ 79 Tab. 4.6 Execution times (seconds): concat command. Highlighted in blue, fastest time and number of times faster than sequential SeqKit. ............ 79 List of Tables Tab. 4.7 Execution times (seconds): sort command. Times with an asterisk were obtained using two nodes. Highlighted in blue, fastest time and number of times faster than sequential SeqKit. .......................... 80 109 The unification of HPC and Big Data has received increasing attention in the last years. It is a common belief that exascale computing and Big Data are closely associated since HPC requires processing large-scale data from scientific instruments and simulations. But, at the same time, it was observed that tools and cultures of HPC and Big Data communities differ significantly. One of the most important issues in the path to the convergence is caused by the differences in their software stacks. This thesis will address the research challenge of bridging the gap between Big Data and HPC worlds. With this goal in mind, a set of tools and technologies will be developed and integrated into a new unified Big Data-HPC framework that will allow the execution of scientific multi-language applications on both environments using containers.