scieee AI-readable full text Open interactive document viewer

Ignis: An efficient and scalable multi-language Big Data framework

Piñeiro Pomar, César Alfredo; Martínez Castaño, Rodrigo; Pichel Campos, Juan Carlos

Abstract

Most of the relevant Big Data processing frameworks (e.g., Apache Hadoop, Apache Spark) only support JVM (Java Virtual Machine) languages by default. In order to support non-JVM languages, subprocesses are created and connected to the framework using system pipes. With this technique, the impossibility of managing the data at thread level arises together with an important loss in the performance. To address this problem we introduce Ignis, a new Big Data framework that benefits from an elegant way to create multi-language executors managed through an RPC system. As a consequence, the new system is able to execute natively applications implemented using non-JVM languages. In addition, Ignis allows users to combine in the same application the benefits of implementing each computational task in the best suited programming language without additional overhead. The system runs completely inside Docker containers, isolating the execution environment from the physical machine. A comparison with Apache Spark shows the advantages of our proposal in terms of performance and scalability

Full text

Rúa Jenaro de la Fuente, s/n –Campus Vida – Universidade de Santiago de Compostela -15782 Santiago de Compostela – citius.usc.es ©2020 Elsevier B.V. This manuscript version is made available under the CC-BY-NC-ND 4.0 license Ignis: An efficient and scalable multi-language Big Data framework César Piñeiro, Rodrigo Martínez-Castaño and Juan C. Pichel Version: accepted article César Piñeiro, Rodrigo Martínez-Castaño and Juan C. Pichel (2020) Ignis: An efficient and scalable multi-language Big Data framework. Future Generation Computer Systems, 105, 705 - 716. Doi: https://doi.org/10.1016/j.future.2019.12.052 How to cite: Copyright information: Ignis: an efficient and scalable multi-language Big Data framework⋆ César Piñeiroa,Rodrigo Martínez-Castañoaand Juan C. Pichela,∗ aCiTIUS, Universidade de Santiago de Compostela, 15782 Santiago de Compostela, Spain ARTICLE INFO Keywords: Big Data; Multi-language; Performance; Scalability; Container Abstract Most of the relevant Big Data processing frameworks (e.g., Apache Hadoop, Apache Spark) only support JVM (Java Virtual Machine) languages by default. In order to support non-JVM languages, subprocesses are created and connected to the framework using system pipes. With this technique, the impossibility of managing the data at thread level arises together with an important loss in the performance. To address this problem we introduce Ignis, a new Big Data framework that benefits from an elegant way to create multi-language executors managed through an RPC system. As a consequence, the new system is able to execute natively applications implemented using non-JVM languages. In addition, Ignis allows users to combine in the same application the benefits of implementing each computational task in the best suited programming language without additional overhead. The system runs completely inside Docker containers, isolating the execution environment from the physical machine. A comparison with Apache Spark shows the advantages of our proposal in terms of performance and scalability. 1. Introduction Nowadays we are living in the Big Data era, which demands processing in a efficient way huge amounts of data coming from many different sources. In the past, clusters were reserved for HPC (High-Performance Computing) tasks, but with the arrival of Big Data it became necessary to design new frameworks and technologies for the execution on these systems. The de facto standards for parallel processing of Big Data are Apache Hadoop [31] and Apache Spark [35] engines. To fully exploit the capabilities of these frameworks programmers should implement their applications in languages based on the Java Virtual Machine (JVM) (fundamentally, Java and Scala) or other high-level languages such as Python. However, developers build applications in the programming languages that best suit their needs and, many times, those languages are not natively supported by the corresponding Big Data framework. 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-tosource compilers [27] 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 [9]. Note that in the latter case, the main code that performs the calls should be anyway implemented in Python, Java or Scala. To overcome the above limitations of the Big Data processing engines we introduce Ignis, a new framework that allows users to execute their applications written in multiple ⋆This work has been supported by MICINN (RTI2018-093336-BC21), Xunta de Galicia (ED431G/08 and ED431C-2018/19) and European Regional Development Fund (ERDF). ∗Corresponding author [email protected] (C. Piñeiro); [email protected] (R. Martínez-Castaño); [email protected] (J.C. Pichel) ORCID(s): 0000-0001-9505-6493 (J.C. Pichel) 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. Currently Ignis supports Python, Java, C and C++, but thanks to its modular architecture, adding support for new languages is a straightforward process. Several efforts have been done to bridge the gap between HPC and Big Data technologies [11,12,19,20,24,33]. 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. We can summarize the key contributions of this paper as follows: – To the best of our knowledge, Ignis is the first Big Data framework with native multi-language support including both JVM and non-JVM-based languages. In this way, users might combine in the same application the benefits of implementing each computational task in the best suited programming language. – Ignis outperforms 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. As a consequence, Ignis facilitates the convergence of HPC and Big Data since data and compute-intensive tasks can be executed efficiently in the same framework. – Contrary to the developers community belief, we show that Python is not natively supported by Spark since data transfers between the JVM and external processes degrade noticeably the overall performance. – 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 C. Piñeiro et al.: Preprint submitted to Elsevier Page 1 of 14 Ignis: an efficient and scalable multi-language Big Data framework 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 [10] 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. The remainder of this paper is organized as follows. Section 2gives the background and discusses some related work. Section 3describes the architecture and the modules of Ignis. The different ways to storage data in Ignis are explained in Section 4. Section 5details the Ignis API. Section 6shows the experimental results. Finally, the main conclusions derived from this work are explained. 2. Background & Related Work 2.1. Big Data Frameworks MapReduce [8] is a programming model introduced by Google for processing and generating large data sets on a huge number of computing nodes. A MapReduce program execution is divided into two main phases: map and reduce. The input and output of a MapReduce computation is a list of key-value pairs. Users only need to focus on implementing map and reduce functions. 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. In this way, map tasks achieve data parallelism, while reduce tasks perform parallel reduction. Currently, several processing frameworks support this programming model such as Hadoop [31], Spark [35], Flink [5] and Tez [28]. In particular, Apache Hadoop is the most successful opensource implementation of the MapReduce programming model. Hadoop consists of three main layers: a data storage layer (HDFS), a resource manager layer (YARN), 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. Apache Spark was designed to overcome some of the Hadoop limitations, especially when considering iterative jobs. It supports both in-memory and on-disk computations in a fault tolerant manner by introducing the idea of Resilient Distributed Datasets (RDDs). An RDD represents a readonly collection of objects partitioned across the cluster nodes that can be rebuilt if a partition is lost. Users can explicitly cache an RDD in memory across machines and reuse it in multiple MapReduce-like parallel operations. By using RDDs, programmers can perform iterative operations on their data without writing intermediary results to disk. Apart from running interactively using Python, Scala and R, Spark can also be linked into applications in either Java, Python, Scala or R. 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 [15], 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 [9]. 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. 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.2. Bridging the Gap between HPC and Big Data There is unanimity in the research community about the importance of approaching HPC and Big Data worlds. We can find in the literature several works dealing with this challenge from different points of view. We can categorized those works depending on if the path of convergence goes from Big Data to HPC or in the opposite direction. In this way, the first category includes solutions that try to boost C. Piñeiro et al.: Preprint submitted to Elsevier Page 2 of 14 Ignis: an efficient and scalable multi-language Big Data framework the performance of well-established Big Data technologies when running on HPC systems. For example, taking advantage of the fast interconnection networks such as Infiniband [19,21], introducing new solutions for data analytics based on classical HPC programming languages [24], or implementing highly optimized libraries to run on HPC systems but using JVM-based languages [11]. In the second category we find works that extend HPC technologies in order to support Big Data tasks. For instance, extending MPI to improve the processing and communication of large numbers of key-value pairs [20]. Note that, unlike previous works, our framework covers both categories since its final goal is that applications belonging to HPC and Big Data worlds can be executed efficiently in a unique framework. Finally, there is an interesting paper that do not fit in the previous categories that analyzes and compares in depth the Big Data and HPC software stacks [12]. It identifies several architecture layers where there is synergy, which can facilitate the integration of both stacks. 2.3. Apache Thrift Apache Thrift [1] 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 [16] 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.4. Docker Docker [10] 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 scriptDriver Backend Manager Executor 1 Executor n ... Driver container Executor containers Docker Resource Manager Figure 1: Scheme of the Ignis architecture. ing 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 [17,32], showing that performance differences between non-containerized and 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 [2] 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 [7] propose a method to allow Docker containers to take advantage of a Infiniband network when running HPC applications. Finally, in a recent work [29], the authors demonstrate how Docker containers can be integrated with HPC environments and run MPI applications with cloudenabled schedulers. 3. 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 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 C. Piñeiro et al.: Preprint submitted to Elsevier Page 3 of 14 Ignis: 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: Docker Resource Manager architecture. 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 operations defined by the BACKEND. Next, a more detailed description of each module is provided. 3.1. Docker Resource Manager The DOCKER RESOURCE MANAGER 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. 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 phys1https://www.consul.io 1"image":"test", 2"resources": { 3"cores": 2, 4"memory":"2GB", 5"swap":"0", 6"volumes": [{ 7"size":"50MiB", 8"mode":"rw-ro", 9"type":"tmpfs" 10 }, { 11 "mode":"rw-rw", 12 "type":"glusterfs" 13 }], 14 "devices": [], 15 # RPC and data ports 16 "ports": [2013, 1963] 17 }, 18 "opts": { 19 "prefered_hosts": ["node1","node2"], 20 "swappiness": 0 21 }, 22 "events": { 23 "on_exit": { 24 "restart":false, 25 "destroy":true 26 } 27 }, 28 "args": ["manager","2013","6000"] Figure 3: Example of a Docker Resource Manager request. ical/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. Tasks (in the context of our resource manager) represent containers with assigned resources. They can be included into an existing task group. Otherwise, a new task group would be automatically created for the new task. This feature allows the user to take actions on all the tasks under the same group at the same time. When launching a task, the following parameters can be provided: number of virtual cores, memory, Docker image, arguments for the image, on-exit behaviors, ports to be opened and volumes. A preference node list can be also supplied in such a way that the first node of the list with enough available resources will be selected to execute the container on. Otherwise, the selection process will follow a round-robin scheduling. In Consul, tasks are registered as services and a health check endpoint is provided so Consul can be used to observe the status of the existing tasks. An example of task request is shown in Figure 3. The master is responsible for checking the request and finding a slave node for it. If succeeds, it will interact with the chosen slave node in order to launch the task with the required resources. Currently, there are three types of supported volumes: local, in-memory and distributed. Local volumes are disk images mounted as loop devices and stored in local disks. When volumes are not in use, they are compressed. In-memory volumes are local volumes whose content is copied to a tmpfs, so memory is used instead of disk during the execution of a task. Finally, distributed volumes are symbolic links to directories within a distributed filesystem such as GlusterFS [13]. C. Piñeiro et al.: Preprint submitted to Elsevier Page 4 of 14 Ignis: an efficient and scalable multi-language Big Data framework Job Cluster 1 Worker 1 ... ... Worker n Cluster n ... Task n Task 1 ... DOCKER CONTAINER Node 1 DOCKER CONTAINER Node 2 DOCKER CONTAINER Node n Figure 4: Job hierarchy in the Ignis framework. Volumes can be created with custom permissions (ro, rw) for the first task that mounts it and, optionally, for the same volume group. Group-shared volumes will be available for all the tasks with mounted volumes in the same group. Local volumes will only be available for the group if tasks are in the same host. All the containers receive the IP address of their host. The requested ports are mapped to random available ports in the host machine. 3.2. Driver Module The DRIVER MODULE is a user API through which users can access all the available functionalities of the Ignis framework. Driver code can be programmed in any of the supported languages (Java, Python and C/C++). This module does not perform any heavy computation and uses Thrift RPC to delegate its work to the BACKEND MODULE. The DRIVER MODULE was designed as a mere interface so the logic has not to be reimplemented for every programming language. Figure 4shows the hierarchy of the components of a job in Ignis. A cluster is a group of Docker containers distributed in multiple computing nodes. In order to build multilanguage applications, at least one worker has to be created for each programming language. Each worker contains tasks, which are the parallel operations that compose a job. Tasks are instantiated as executor processes inside Docker containers. Figure 5shows an example of a driver implemented in two different programming languages, Python and C++, for the well-known Wordcount application. Note that both drivers show a similar behavior and syntax. After initializing the Ignis framework, a cluster of Docker containers is configured and built (lines 4-10, driver code). Several parameters such as Docker image, number of containers, number of cores and memory per container are established. Wordcount has two phases. The first stage consists of a map operation that takes as input a text file and is tokenized into key-value pairs (𝑤𝑜𝑟𝑑, 1). This task in the example of Figure 5was implemented in Python. As a consequence it is necessary to create a Python worker in the cluster (line 13, driver code). Note that the worker is mandatory and it is not related to the programming language in which the driver code is written. Once the worker is created, the map task is defined indicating the file and class paths where the function to be applied can be found (line 16, driver code). We must highlight that for languages that support source code serialization as Python or Java, references to functions can be used instead. Resulting key-value pairs are stored in words. In the following phase of the Wordcount job all the keys are grouped together and the values for similar keys are added up to find the occurrences for a particular word. To illustrate the multi-language support of Ignis that reduction task in the example was implemented in C++. As a consequence, it is necessary to create a C++ worker (line 19, driver code). Data can be shared between different workers using the importData function (line 22, driver code). It is worth noting that this function is only necessary when a task requires data generated by a previous one that belongs to a different worker. The driver code ends storing the final results in a file. Lazy evaluation is performed so tasks in the driver code are only executed when a result is required explicitly. In Figure 5the trigger that causes the tasks to be launched is writing the final result to file (line 26, driver code). This approach is also followed by Spark where RDDs are computed lazily the first time they are used in an action [34]. 3.3. Backend Module The BACKEND MODULE and the DRIVER MODULE are executed in the same Docker container (see Figure 1). The BACKEND MODULE has the services which defines the logic of the DRIVER MODULE. For instance, reduceByKey consists of three steps: searching and grouping the keys, and accumulating values with the same key. These operations are defined in the BACKEND MODULE, but their implementations for a particular programming language are done in the EXECUTOR MODULE. The BACKEND MODULE is also in charge of making the requests to the DOCKER RESOURCE MANAGER with the aim of building the cluster following the properties specified in the driver code. Tasks are stored by the BACKEND MODULE, which are instantiated as executor processes inside the cluster containers. Therefore, each task uses several executors to apply in parallel the same operation to multiple data items. Performance optimizations can be done such as placing executors with high interaction together in the same host. The distribution and exchange of data among executors is also handled by this module. Data exchange is an asynchronous operation required by functions such as sort,shuffle or reduce. The BACKEND MODULE determines the data to be sent, and the addresses of the source and destination executors for each communication. Once the exchange service is initiated all the messages are sent from the executor processes following the directives of the BACKEND MODULE. The service is the responsible of building connections among executors and allocating a buffer to write and read C. Piñeiro et al.: Preprint submitted to Elsevier Page 5 of 14 Ignis: an efficient and scalable multi-language Big Data framework 1# Initialization of the framework 2ignis.Ignis.start() 3# Resources/Configuration of the cluster 4prop = ignis.IProperties() 5prop["ignis.executor.image"]="wordcount" 6prop["ignis.executor.instances"]="2" 7prop["ignis.executor.cores"]="4" 8prop["ignis.executor.memory"]="2GB" 9# Construction of the cluster 10 cluster = ignis.ICluster(prop) 11 12 # Initialization of a Python Worker in the cluster 13 worker_python = ignis.IWorker(cluster, "python") 14 # Task 1 - Python: tokenize text into pairs ('word', 1) 15 text = worker_python.readFile("text.txt") 16 words = text.flatMap("wordcount/wd.py:Split") 17 18 # Initialization of a C++ Worker in the cluster 19 worker_cpp = ignis.IWorker(cluster, "cpp") 20 # Task 2 - C++: reduce pairs with same word and obtain totals 21 # Transfer data from Task 1 - Python 22 words_cpp = worker_cpp.importData(words) 23 count = words_cpp.reduceBykey("wordcount/libwd.so:Join") 24 25 # Print results to file 26 count.saveAsFile("wordcount.txt") 27 28 # Stop the framework 29 ignis.Ignis.stop() 1// Initialization of the framework 2Ignis::start(); 3// Resources/Configuration of the cluster 4prop = std::make_shared<IProperties>(); 5(*prop)["ignis.executor.image"]="wordcount"; 6(*prop)["ignis.executor.instances"]="2"; 7(*prop)["ignis.executor.cores"]="4"; 8(*prop)["ignis.executor.memory"]="2GB"; 9// Construction of the cluster 10 cluster = make_shared<ICluster>(prop); 11 12 // Initialization of a Python Worker in the cluster 13 worker_python = make_shared<IWorker>(cluster, "python"); 14 // Task 1 - Python: tokenize text into pairs ('word', 1) 15 text = worker_python.readFile("text.txt"); 16 words = text.flatMap("wordcount/wd.py:Split"); 17 18 // Initialization of a C++ Worker in the cluster 19 worker_cpp = make_shared<IWorker>(cluster, "cpp"); 20 // Task 2 - C++: reduce pairs with same word and obtain totals 21 // Transfer data from Task 1 - Python 22 words_cpp = worker_cpp.importData(words); 23 count = words_cpp.reduceBykey("wordcount/libwd.so:Join"); 24 25 // Print results to file 26 count.saveAsFile("wordcount.txt"); 27 28 // Stop the framework 29 Ignis::stop(); Figure 5: Wordcount driver code example using Python (left) and its equivalent C++ code (right). data in the connections. Once all the messages are sent, the exchange service stops. Note that for efficiency reasons executors in the same host exchange data through the shared memory. Finally, we must highlight that Ignis is able to recover after a failure of a cluster node or some of the executors. The BACKEND MODULE is able to follow the task trace of the affected executors in such a way that only those executors are reallocated and recomputed. It means that even if the input data of an executor is lost, only the tasks necessary to recompute that portion of data are performed. This process is done without the intervention of the user. An experimental evaluation of the fault tolerance mechanisms of Ignis is shown in Section 6. 3.4. Manager Module The MANAGER MODULE is the connection point between the BACKEND MODULE and the executor processes (see Figure 1). Any command from the BACKEND MODULE to the executor processes goes through this module, which includes launching the executors in the containers. In the scenario of no answer from an executor process, the MANAGER MODULE kills the corresponding executor and informs the BACKEND MODULE in order to start the recovery process. 3.5. Executor Module The operations defined in the BACKEND MODULE are implemented for each supported programming language in the EXECUTOR MODULE. Therefore, adding a new language to Ignis only requires to implement those operations in the corresponding language. Most of the driver functions such as map or reduceByKey are meta-functions. That is, generic functions that require another function to perform an internal operation. For example, flatMap in the codes of Figure 5apply the Split function to the input text (line 16, driver code). Those functions in Ignis are defined using a common interface based on the number of input and output parameters. Figure 6shows as example the Split and Join functions of the Wordcount application in Python and C++. Split tokenizes a text into key-value pairs (𝑤𝑜𝑟𝑑, 1). Therefore, it is a typical flat operation that produces an arbitrary number (zero or more) values for each input value. The difference with a map operation is that it must produce one output value for each input. In the Ignis API, IFlatFunction represents a flat operation, where the call method always contains the implementation of the function. Join takes two count values for the same word and accumulates them. Consequently, multiple values are reduced iteratively to a single value. In this case, IFunction2 represents an operation that takes two arguments and generates only one result. Note that in the C++ code data types of the input and output parameters should be specified. The context object is also passed as argument to Split and Join (Python - lines 4 and 10, C++ - lines 4 and 19), which allows the user to access more information such as the properties defined in the driver code. On the other hand, some applications may need to perform certain specific tasks for each executor process. For instance, opening a database connection, reading a particular file or preparing the environment before processing. To facilitate this task all the interfaces to the functions in the Ignis C. Piñeiro et al.: Preprint submitted to Elsevier Page 6 of 14 Ignis: an efficient and scalable multi-language Big Data framework 1# Tokenize text into pairs ('word', 1) 2class Split(IFlatFunction): 3 4def call(self, elem, context): 5return [(word, 1) for word in elem.split()] 6 7# Reduce pairs with same word and obtain totals 8class Join(IFunction2): 9 10 def call(self, elem1, elem2, context): 11 return elem1 + elem2 1// Tokenize text into pairs ('word', 1) 2class Split : public api::function::IFlatFunction<string, pair<string,int64_t>> 3{ 4api::Iterable<pair<string,int64_t>> call(string& line, api::IContext& context) { 5// Vector to accumulate pairs 6vector<pair<string,int64_t>> pairs; 7// String stream tokenizer 8stringstream words(line); 9// Iterator 10 istream_iterator<string> begin(words), end; 11 12 for(it = begin; it != end; it++) pairs.emplace_back(*it, 1); 13 return(pairs); } 14 }; 15 16 // Reduce pairs with same word and obtain totals 17 class Join : public api::function::IFunction2<int64_t, int64_t, int64_t> 18 { 19 size_t call(int64_t& count1, int64_t& count2, api::IContext& context) { 20 return count1 + count2; } 21 }; Figure 6: Split and Join functions for the Wordcount application using Python (left) and its equivalent C++ code (right). API provide a before method, which is called at the beginning of the executor, and an after method, which is called once at the end of the processing. This functionality is similar to the setup and cleanup methods in Hadoop, but it is not supported by Spark. 4. Data Storage Ignis provides several options for data storage thanks to its common storage interface. In this way, it is possible to create different representations of the data without modifying the operations. Users can choose the type of storage that best fits their application taking into account the limitations of their particular execution environment. Only it is necessary to specify the selection as a property in the driver code. Next we explain the different storage options supported by Ignis: –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 [16], 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 [26]. 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 inmemory storage when using a Python worker and the number of cores per container is greater than one (see the Wordcount driver code of Figure 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: inmemory, 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. 5. Ignis API To use Ignis, developers should write a driver program that implements their application at high-level (see the example of Figure 5). As we explained in Section 3.5, some of the driver functions such as map require another one to execute. In that case it is also necessary to implement those C. Piñeiro et al.: Preprint submitted to Elsevier Page 7 of 14 Ignis: an efficient and scalable multi-language Big Data framework functions. We must highlight that although the Ignis code uses a sequential notation, 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 (𝑘𝑒𝑦, 𝑣𝑎𝑙𝑢𝑒)pair after applying a function to the 𝑣𝑎𝑙𝑢𝑒 argument with the aim of calculating the 𝑘𝑒𝑦.values does not need an extra function, it only returns 𝑣𝑎𝑙𝑢𝑒 from a (𝑘𝑒𝑦, 𝑣𝑎𝑙𝑢𝑒)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. –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 𝑛elements, where 𝑛is chosen by the user. takeSample returns a random sample of 𝑛elements 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. 6. 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. 6.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 [30] 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 2Do not confuse with the data mining benchmark suite NUMineBench. C. Piñeiro et al.: Preprint submitted to Elsevier Page 8 of 14