scieee AI-readable full text Open interactive document viewer

Performance Analysis of Distributed GPU-Accelerated Task-Based Workflows

Simitsis, Alkis; Carvalho, Marcos N. L.; Queralt, Anna; Romero, Oscar; Tatu, Cristian Cătălin; Badia, Rosa M

Abstract

We present an empirical approach to identify the key factors affecting the execution performance of task-based workflows on aHigh Performance Computing (HPC) infrastructure composed of heterogeneous CPU-GPU clusters. Our results reveal that the execution performance in distributed GPU-accelerated task-based workflows highly depends on several interrelated factors regarding the task algorithm, dataset, resources, and system employed. In addition, our analysis identifies key correlations among these factors, presents novel observations, and offers guidelines toward designing an automated method to handle task-based workflows in modern, high-compute capacity, CPU-GPU engines.

Full text

Performance Analysis of Distributed GPU-Accelerated Task-Based Workflows Marcos N. L. Carvalho∗ [email protected] UPC, Barcelona, Spain NKUA & Athena RC, Athens, Greece Anna Queralt [email protected] Universitat Politècnica de Catalunya Barcelona, Spain Oscar Romero oscar[email protected] Universitat Politècnica de Catalunya Barcelona, Spain Alkis Simitsis [email protected] Athena Research Center Athens, Greece Cristian Tatu [email protected] Barcelona Supercomputing Center Barcelona, Spain Rosa M. Badia [email protected] Barcelona Supercomputing Center Barcelona, Spain ABSTRACT We present an empirical approach to identify the key factors affecting the execution performance of task-based workflows on a High Performance Computing (HPC) infrastructure composed of heterogeneous CPU-GPU clusters. Our results reveal that the execution performance in distributed GPU-accelerated task-based workflows highly depends on several interrelated factors regarding the task algorithm, dataset, resources, and system employed. In addition, our analysis identifies key correlations among these factors, presents novel observations, and offers guidelines toward designing an automated method to handle task-based workflows in modern, high-compute capacity, CPU-GPU engines. 1 INTRODUCTION Data Science (DS) pipelines are essential for many software systems today. Such pipelines are composed of multiple processing stages that perform different tasks to move data from one stage to a next stage [ 7 ]. The relationship between these stages creates complex workflows, where data preparation, training and evaluation of Machine Learning (ML) models are just examples of processing that is frequently done over big datasets. To process large amounts of data efficiently, distributed and parallel applications are required and task-based workflows provide a highlevel programming abstraction to develop such applications [ 55 ]. High Performance Computing (HPC) clusters are being widely used to process DS workloads because of the massive parallelism enabled by distributed architectures [ 54 ]. In addition, improvements in the hardware industry have increased the popularity of accelerators, such as graphic processing units (GPUs), making modern distributed infrastructures even more heterogeneous [ 2 ]. This evolution has led to distributed GPU-accelerated task-based workflows [ 6 ], which provide massive processing power to users by leveraging both task-level parallelism of distributed CPUs and thread-level parallelism of GPUs. In this context, tasks are distributed and processed in parallel in CPUs and each task, internally, has its threads parallelized by GPUs. Considering the performance of such workflows, we face three core challenges: (i) Ad hoc design: Lacking concrete guidelines ∗ The author pursues a joint PhD degree under the auspices of DEDS (No 955895), a Horizon 2020 MSCA ITN, and he is co-affiliated with Universitat Politècnica de Catalunya in Spain, Athena Research Center and National and Kapodistrian University of Athens (a degree awarding institute for Athena R.C.) in Greece. ©2024 Copyright held by the owner/author(s). Published in Proceedings of the 27th International Conference on Extending Database Technology (EDBT), 25th March-28th March, 2024, ISBN 978-3-89318-095-0 on OpenProceedings.org. Distribution of this paper is permitted under the terms of the Creative Commons license CC-by-nc-nd 4.0. and optimization heuristics, the workflow developer often resorts to intuition to assign tasks to GPU accelerators. And oftentimes, due to the huge design space of execution parameters (a.k.a. factors), the developer might exhaustively rerun representative workloads in order to identify efficient execution settings [ 10 ]. In practice, even expert developers spend a considerable amount of time searching for appropriate, but often sub-optimal, settings. (ii) Resource wastage: Finding fitting settings results into efficient resource usage. On the other hand, a poor combination of execution parameters leads to load imbalance and waste of resources. For example, a non desirable situation would be to keep the CPUs busy while the GPUs stay idle during workflow processing. (iii) Complexity: Previous studies of the problem have focused on individual factors that could affect execution performance, such as block size, dataset input size, etc. In our study, we argue and show that this is not sufficient, as in real-world scenarios the workflow execution performance is typically affected by a combination of factors. To motivate the discussion, consider the following example that attempts to improve the performance of a task-based workflow using GPU accelerators. The workflow represents a CPUbased and a GPU-accelerated version of a distributed implementation of the K-means algorithm 1 . In this experiment, we consider single task (on 1 CPU core and 1 GPU device) and parallel tasks (on all CPU cores and GPU devices) execution. Generally, task processing involves data computation, represented by the task user code, and data movement (we detail these in Section 3.3). Figure 1 shows the execution performance at different task processing stages: (i) For a single task, a 5.69x speedup of GPU over CPU is achieved when only the parallel fraction of a task user code is considered; this is the fraction of the task user code where threads can be parallelized on a GPU. (ii) The speedup, again for a single task, is reduced to 1.24x when the total execution of the task user code is considered; this includes the CPU-GPU communication overhead and the serial (i.e., non-parallel) fraction of the task as well. (iii) Interestingly, for parallel tasks execution (i.e. all tasks being distributed and processed in parallel), we observe that the overall performance of GPUs is worse than CPUs (-1.20x speedup). Arguably, as a partial analysis of the performance of GPUs vs. CPUs in distributed task-based workflows may produce misleading or incomplete results, we need to conduct a more thorough and principled analysis. And this exactly is the goal of our work. To efficiently run distributed task-based workflows, we need to identify and characterize the factors affecting the performance at 1 The experiment involves a 10 GB dataset processed by 256 tasks distributed on a cluster with 128 CPU cores and 32 GPU devices. For more details, see Section 4. Experiments & Analyses Paper Series ISSN: 2367-2005 690 10.48786/edbt.2024.59 -1.20x speedup 1.24x speedup 5.69x speedup Figure 1: Performance of distributed K-means at different processing stages on CPUs and GPUs all task processing stages, as those presented in Figure 1. Earlier work has identified a number of factors affecting the performance of heterogeneous CPU-GPU processing in a multitude of scenarios, as illustrated in Figure 2. Different performance limitations arise depending on the infrastructure used, such as single, server machines and compute clusters. On a single machine, the main limitations are CPU-GPU data transfer bottleneck and device speedup. Given a task user code, a potential mitigation technique to overcome CPU-GPU communication would be to increase computation on GPUs using techniques such as staged pipeline and zero-copy [ 60 ]. The more computation, the higher the device speedup. The amount of computation in a task heavily depends on its arithmetic intensity, i.e., the number of math operations divided by number of bytes read [ 10 , 45 ], and the task granularity, i.e., the amount of data to process [ 78 ]. However, tasks may involve parallel fraction processing that exploits thread-level parallelism on GPUs, and serial fraction processing executed on CPUs. Hence, the parallel fraction processing within tasks is a key factor to consider in order to fully utilize GPUs, as it defines the abundance of thread-level parallelism [ 22 ]. Earlier work has studied the issues of CPU-GPU communication and parallel fraction independently, but has not considered both problems in tandem, exactly as it happens in real settings. On compute clusters, tasks are distributed and processed in parallel on multiple nodes, thus, scaling out task-level parallelism [ 72 ] and providing memory robustness [ 62 ] to GPUs by breaking the input dataset into chunks. Such an infrastructure adds non-negligible overheads, including storage I/O [ 27 , 38 , 69 , 70], network I/O [6, 26, 34, 78], and task scheduling [2, 25]. Regardless of the infrastructure used, task granularity is frequently considered as a key performance factor. However, increasing task granularity alone is not always the best strategy to improve the performance of distributed task-based workflows, as it increases thread-level parallelism at the cost of reducing tasklevel parallelism, which in turn leads to load imbalance between CPUs and GPUs. Our empirical analysis reveals that the performance of distributed GPU-accelerated task-based workflows is a result of many interrelated factors involving, beside task granularity, CPU-GPU data transfer, task parallel fraction, storage I/O, network I/O, task scheduling, etc. We argue that focusing on a single factor separately (e.g. increasing task granularity to maximize device speedup) as the current related work suggests, leads to sub-optimal designs. On the other hand, combining multiple performance factors (e.g. multi-level parallelism) to optimize task-based workflows in heterogeneous CPU-GPU, distributed environments involves a non trivial design complexity, which makes this a quite challenging problem. The literature misses a detailed performance analysis to characterize the performance of task-based workflows and study how to balance thread-level and task-level parallelism to maximize resource utilization. Our contributions. We present a systematic performance analysis of task-based workflows and make the following contributions: • A systematic analysis of thread-level parallelism; our results reveal that the gains provided by GPUs are highly affected by the parallel fraction processing within tasks. • A systematic analysis of task-level parallelism; our results reveal that depending on the task granularity, scheduling policy, and storage architecture used, system overheads can be significantly high, dominating the total execution time and eliminating the potential gain of GPU processing. • A method to identify what are the factors affecting the performance of task-based workflows, how these relate to each other and to the algorithms, datasets, resources, and distributed execution system employed; our results indicate that certain execution parameters are highly correlated with the performance of task-based workflows and can be considered as key factors. • We present novel observations and offer guidelines toward designing an automated method to handle task-based workflows in modern, high compute capacity, CPU-GPU engines. Structure. The rest of the paper is organized as follows. Section 2 presents the related work. Sections 3 and 4 describe background concepts and the method used in our analysis, respectively. Section 5 presents our experimental analysis, and Section 6 concludes the paper. 2 RELATED WORK The literature about heterogeneous CPU-GPU processing is wide and several approaches have been proposed to improve the performance of CPU-GPU processing. In Figure 2, we present a classification of the state of the art on CPU-GPU processing and highlight in red the related work that is within the scope of our analysis. In particular, the related work on heterogeneous CPU-GPU has focused on frameworks [ 27 , 31 , 71 ], libraries [ 20 , 56 , 66 ], schedulers [ 10 , 25 , 30 ], cost-based models [ 10 , 53 , 62 ], programming models [ 3 , 34 , 69 ], and benchmarks [ 14 , 78 ]. Previous performance analysis studies focused on specific applications like deep learning models, ETL processes, image classification, and query performance [ 9 , 10 , 29 , 32 , 42 , 63 ]. Our study targets generic taskbased workflows in CPU-GPU environments, follows a bottomup analysis from thread-level to task-level parallelism, and considers a multiplicity of factors affecting workflow performance. Previous works have also explored different ways to use GPUs, for example, as the primary processor [ 27 , 62 ], as an accelerator [ 73 , 78 ] or in the context of heterogeneous CPU-GPU processing [ 32 , 59 , 69 , 71 ]. Regarding the application field, several papers have demonstrated interesting results by using GPUs to accelerate database query processing [ 10 , 32 , 62 , 71 ] and dataintensive analytics applications with task-based workflows [ 2 , 3 , 9, 29, 42, 78], dataflows [15, 57], and graph processing [39, 76]. 691 CPU-GPU Processing GPU Usage GPU Integration Application Level of Analysis Primary Processor [27, 62] Accelerator [73, 78] Heterog. CPU-GPU [32, 59, 69, 71] Integrated [33, 35, 75] Dedicated Database [10, 32, 62, 71] Analytics Instruction [10, 64, 69] Task [32, 62, 71] DAG [27, 39] Infrastructure Data-intensive Applications Single Machine Cluster Task-based Workflows [2, 3, 9, 29, 42, 78] Dataflows [15, 57] Graph Processing [39, 76] CPU-GPU Data Transfer [11, 32, 33, 36, 59, 60, 71] Device Speedup [9, 16] Storage I/O [27, 38, 69, 70] Network I/O [6, 26, 34, 78] Task Scheduling [2, 25] Figure 2: A taxonomy for CPU-GPU processing (the scope of our study is highlighted in red text) Our focus is on analyzing the impact of GPUs on task-based workflows, which are nowadays commonly used in various applications, such as data science pipelines. Different levels of analysis can be found in the literature, spanning instruction-level optimization [ 10 , 64 , 69 ], task-level [ 32 , 62 , 71 ] and distributed acyclic graph (DAG) evaluation [27, 39]. Some approaches investigate using integrated GPUs [ 33 , 35 , 75 ]. However, the majority of previous work considers dedicated GPUs due to the higher processing capacity compared to integrated GPUs [ 60 ]. Depending on the infrastructure used, i.e. single machine or cluster, different limitations may arise for dedicated GPUs. Limitations on single server machine. One of the main limitations reported in the literature about the performance of heterogeneous CPU-GPU processing is the data transfer bottleneck [ 11 , 32 , 33 , 36 , 59 , 60 , 71 ]. This bottleneck has been studied extensively and several techniques have been proposed to alleviate it. The techniques include overlapping transfer with execution, data compression, approximation, caching, data locality, single-pass algorithms, heterogeneous execution and faster system bus [ 60 ]. Some studies have also proposed increasing the amount of data to process (i.e., task granularity) as much as the GPU memory supports in order to increase the throughput and the device speedup [9, 16]. However, recent publications have demonstrated that this strategy is not always the best option to improve the performance of heterogeneous CPU-GPU environments [ 32 , 61 ]. Additionally, the parallel fraction processing within a task is a factor that can limit performance as much as the CPU-GPU communication. A theoretical analysis of the parallel fraction is done in [ 53 ], but there is no empirical study about it in the literature. Moreover, no previous work has evaluated the impact that these limitations all together may cause in the performance of GPU-accelerated tasks. Our study reveals that it is worth using GPUs when parallel fraction processing time is large enough to overcome not only CPU-GPU processing time, but also serial fraction processing time. Limitations on cluster deployment. Distributed environments enable scaling out task processing over several nodes at the cost of adding potential bottlenecks, such as data storage I/O [ 27 , 38 , 69 , 70 ], network I/O [ 6 , 26 , 34 , 78 ], and task scheduling overheads [ 2 , 25 ]. Although these are relevant limitations to consider when assessing the performance of parallel task processing, there is no study on how to tune task granularity to balance task-level and thread-level parallelism and smooth such overheads. In addition, there is no work showing how the scheduling policy used and the storage architecture (i.e. local or shared disk) may affect these overheads. Our study shows that such factors must be considered in tandem in order to distribute tasks efficiently on available CPU cores while at the same time utilizing GPUs. To the best of our knowledge, ours is the first study to relate multiple factors to the execution performance of distributed GPU-accelerated task-based workflows. In the past, other works have considered multiple factors but not all together and not for distributed GPU-accelerated workflows [2, 6]. 3 PRELIMINARIES Several systems facilitate the development and execution of parallel applications on distributed infrastructures (e.g, clusters, clouds, containerized platforms) such as Apache Spark [ 74 ] and COMPSs [ 4 ]. They all share a similar flow. The developer submits an application (e.g., matrix multiplication, K-means, neural network, etc.) to such a system, which then generates a Directed Acyclic Graph (DAG) based on the data dependencies between its tasks (e.g., math operations, such as dot product, blocked matrix multiplication, etc.). The system handles parallelism and distribution of tasks by scheduling and allocating resources to execute the tasks in heterogeneous processors (e.g., CPU or GPU) of a cluster. Furthermore, the required data to process the tasks is accessed through a storage architecture. Although HPC architectures typically decouple processing from storage using shared disks, the option of using local disks is also available. An overview of a typical distributed task-based workflow processing system is illustrated in Figure 3. Next, we describe the core components of such a system and set the basis of our analysis. 3.1 DAG creation In distributed GPU-accelerated task-based workflows, an application is broken into small manageable parts defined as tasks. Each task takes an input data, perform some calculation over it, and 692 Code Sequential Parallel Distributed System (1) Code submission (2) DAG creation Cluster Resources Processing Storage Node 0 core core core core CPU ... GPU Node N core core core core CPU GPU (3) Task scheduling and resource allocation DAG Local Disk Disk ... Disk Node 0 Node N Shared Disk Disk Network (4) Task execution (5) Data access PCIe PCIe Runtime System Task parallelization Task dependency Figure 3: An overview of a distributed task-based workflow processing system return an output to be used by the subsequent dependent tasks, if any. When an application code is sent to a processing system (e.g., COMPSs), the data dependencies between tasks are automatically identified and an execution DAG is generated, composing an execution workflow. In this DAG, the vertices represent tasks and the edges represent dependencies. The shape of the DAG reveals key characteristics of task processing. In particular, the width of the DAG shows the degree of parallelism (i.e., #parallel tasks generated) and the height of the DAG shows the degree of task dependency. 3.2 Task scheduling The system runtime schedules the tasks of a DAG that are free of dependencies on the available resources. In general, various scheduling policies are typically available prioritizing aspects, such as task generation order and data locality. Depending on the scheduler selected, the parallel task execution time may vary due to scheduling overheads, which may be low (e.g., for prioritizing task generation) or high (e.g., when considering data locality). 3.3 Task execution Figure 4 shows an abstract illustration of the typical processing stages of a task, namely the task user code that relates to data computation, and serialization and deserialization that relate to data movement. The task user code might be serial (i.e., runs single-threaded) and/or parallel (i.e., multi-threaded). A serial task only contains serial code. A partially parallel task contains both serial and parallel code fractions, and a fully parallel task only parallel code. Besides data computation, the user code also includes inter/intra-CPU, inter/intra-GPU, and CPU-GPU communication functionality. For the scope of this work, here we only consider CPU-GPU communication; i.e., the overhead to move data between CPUs and GPUs over a bus. There are two levels of parallelism in task execution: threadlevel parallelism and task-level parallelism. Thread-level parallelism. This takes place in the parallel fraction of a task user code. Although thread parallelism is feasible in CPUs (e.g., internally multi-threaded or parallelized with OpenMP [ 46 ]), several frameworks recommend assigning a task to a single CPU to avoid CPU over-subscription and hence, leverage full task parallelism in CPUs. Our initial micro-benchmarks corroborate this practice, but we plan to investigate this further in future work. Hence, we consider that serial tasks are assigned to CPUs, and partially or fully parallel tasks to GPUs, which accelerate them with thread parallelism. Note, that a GPU-accelerated task is not assigned directly to a GPU; it is first processed in a Deserialization Serial fraction Serialization Task user code Parallel tasks (a) Serial task Deserialization Serial fraction Serialization Parallel fraction CPU-GPU Comm. Task user code Parallel tasks (b) Partially parallel task Deserialization SerializationParallel fraction CPU-GPU Comm. Task user code Parallel tasks (c) Fully parallel task Figure 4: Abstract task processing stages CPU core to deserialize data, then a GPU executes the parallel code, and the result is sent to a CPU core to serialize the output. Task-level parallelism. Multiple tasks can be executed in parallel. The degree of task parallelism is affected by task dependencies and available resources. Independent tasks are processed in parallel on available CPU cores in the cluster. Tasks with dependencies are processed sequentially, as soon as the execution of their dependencies has finished. Task parallelism is limited by the amount of available GPU devices in the cluster. Since cluster nodes typically have fewer GPU devices than CPU cores, GPU-accelerated tasks have lower degree of task parallelism. For example, in a cluster containing 128 CPU cores and 32 GPU devices, we can execute in parallel a maximum of 128 CPU-based tasks and only 32 GPU-accelerated tasks. Still, although GPU-accelerated tasks reap reduced benefits from task parallelism, they achieve considerable speedups due to thread parallelism. And this creates a critical trade-off between thread-level and task-level parallelism. 3.4 Data access Some tasks may need to read from and write data to a storage device (e.g., disk or memory), which results to storage I/O and data transfer overheads e.g., due to data deserialization and serialization, respectively. The overhead varies depending on the storage architecture. For example, with local disks the data is accessed 693 Grid G(k x l), 8 blocks (4x2) Task 1 Task 2 Task 3 Task 4 Input dataset D(i x j), 64 elements (8 x 8) Row-wise chunking Blocks B(m x n), 8 elements (2x4) Grid G(k x l), 8 blocks (4x2) Task 5 Task 6 Task 7 Task 8 Task 1 Task 2 Task 3 Task 4 Hybrid row-wise and column-wise chunking Figure 5: Example data partitioning and task parallelization locally without having a communication overhead. With shared disks, processing and storage are decoupled and the data is accessed from a shared file system accessible from all nodes of the network, which however (from a performance perspective) adds considerable communication overhead due to network latency, resource contention, etc. 3.5 Programming model Without loss of generality, in our analysis, we consider that the input data is in the form of a matrix. This is in accordance with typical applications in parallel programming and popular libraries, such as CUDA [44] and dislib [17]. The processing system splits the input dataset (here, matrix) into blocks. And the blocks are typically organized in grids (e.g., ds_array in dislib, grid in CUDA, etc.). Figure 5 shows an example partitioning scheme for an input dataset comprising 64 elements in an 8x8 matrix. Assuming a block size of 8 elements organized in 4 columns and 2 rows, the grid will contain 8 blocks. Depending on the application (e.g., an algorithm), a chunking policy determines how to organize the blocks of a grid and assign them to tasks (task parallelism). We refer to task granularity as the blocks of the grid assigned to a single task. Blocks sent to GPU for processing are further chunked to get parallelized at a thread level (thread parallelism) using specialized libraries (e.g., CuPy for CUDA [20]). Note that most systems do not provide a direct way to control how many tasks are spawned. Still we can achieve this indirectly by tuning the block size, which is determined by the developer. In our experiments, we set the task granularity to one block per grid to control the exact number of tasks get spawned (see also Section 4.4.4). Therefore, large blocks result into coarse-grained tasks and small blocks result into fine-grained tasks. Formally, let 𝑫𝑖×𝑗 be the input dataset with 𝑖 rows and 𝑗 columns, 𝑩𝑚×𝑛 be a block, and 𝑮𝑘×𝑙 be a grid with 𝑘 blocks per row and 𝑙 blocks per column. The size of 𝑫𝑖×𝑗 (i.e., 𝑖×𝑗 ) gives the total number of elements of the input dataset to be processed, the size of 𝑮𝑘×𝑙 (i.e., 𝑘×𝑙 ) gives the grid dimension (i.e., the total number of blocks per grid) and the size of 𝑩𝑚×𝑛 (i.e., 𝑚×𝑛 ) gives the block dimension (i.e. the total number of elements per block). The relationship between these sizes is described by Eq. (1): 𝑖=𝑘×𝑚, 𝑗 =𝑙×𝑛, (1) To define a grid, the block dimension should be specified in the application. Thus, assuming that 𝑫𝑖×𝑗 and 𝑩𝑚×𝑛 are given, Eq. (1) can be rewritten as Eq. (2): 𝑘= 𝑖 𝑚,𝑙 = 𝑗 𝑛,(2) Note that 𝑘 and 𝑙 are inversely proportional to 𝑚 and 𝑛 , respectively. This relationship reveals a trade-off between task-level and thread-level parallelism. The larger the block dimension, the smaller the grid dimension, which increases the task granularity (and, hence, thread parallelism, as now the block contains more elements to process) and reduces the number of generated tasks (and, hence, task parallelism). On the contrary, the smaller the block dimension, the larger the grid dimension, which minimizes the degree of thread parallelism and maximizes the degree of task parallelism. Consequently, these relationships are key to achieve a balanced degree of task-level and thread-level parallelism. However, there are two constraints. First, the memory size of a block must fit in the memory of a processor to avoid undesirable effects; e.g., out-of-memory (OOM) errors. Second, the block dimension ( 𝑘×𝑙 ) cannot be greater than the input dataset dimension ( 𝑖×𝑗 ). 4 OUR ANALYSIS METHOD Increasing task granularity is not the only strategy to improve performance of distributed GPU-accelerated task-based workflows. We argue that additional factors should also be considered. To validate our hypothesis we conduct a performance analysis adopting the systematic method proposed by Jain [28]. 4.1 Evaluated algorithms We consider task-based workflows that represent data-intensive applications and run on a heterogeneous CPU-GPU cluster with dedicated GPUs. A task-based performance analysis relates to how the low-level tasks in each algorithm are processed by a distributed system. From this aspect, the algorithm’s logic is orthogonal to how the underlying tasks implementing it (we don’t measure how these tasks are produced) are being processed (this is what we measure in our analysis). In this setting, the parallel fraction of the user code is of paramount importance to task-based analysis (see also Section 5). Hence, our criterion to choose representative algorithms is the level of parallelism employed in each algorithm. To this end, we consider two families of algorithms based on how they perform data computation in the task user code employing serial and/or parallel processing fractions (see Section 3.3): (a) Fully parallelizable algorithms involve fully parallel tasks whose user code can be fully parallelized in a GPU device (see Figure 4c). And (b) Partially parallelizable algorithms involve partially parallel tasks whose user code contains both serial and parallel fractions (see Figure 4b). In our analysis, we employ two representative algorithms of each family: an algorithm from family (a) that involves only parallel task execution vs. one from family (b) that has a low ratio of parallel / serial code in its task execution. For the first case, we use Matrix Multiplication (Matmul), which fulfills our criterion and is also a fundamental operation in many ML/DL techniques, including LLMs, PCA, SVD, linear regression, etc. For the second case, we use K-means, another widely used and easy to understand algorithm, whose ratio of parallel / serial code is low and hinders its parallel execution. Note that any other fully parallelizable (in the first case) or partially parallelizable (in the second case) algorithm would also be applicable to our analysis. 694 For these algorithms, we investigate when it is beneficial to use CPU and/or GPU acceleration taking advantage of thread and task parallelism yielded by key factors related to the algorithms, datasets, resources, and the distributed system employed. 4.2 Evaluated metrics First, we consider metrics related to task execution according to its processing stages: task user code and (de-)serialization. These can be used to perform a head-to-head comparison between CPUbased and GPU-accelerated tasks leaving aside other overheads, such as storage I/O, network I/O, and task scheduling. Task user code metrics. We report these metrics aggregated per task type; i.e., tasks running the same code are aggregated together. • Serial fraction execution time: average time per task to process the serial fraction of the user code. • Parallel fraction execution time: average time per task to process the parallel fraction of the user code. • CPU-GPU communication time: average time per task to move data from CPU to GPU or vice versa. • User code execution time: average time per task user code; i.e., summary of serial fraction, parallel fraction, and CPU-GPU communication times. Data movement overheads. We report the data serialization and deserialization times grouped by CPU core for all task types. • Deserialization time: average time per CPU core to read data from storage (e.g., disk) and load it to main memory. • Serialization time: average time per CPU core to write data from main memory to storage. Task level metrics. These metrics are computed per DAG level, i.e., for tasks that are in the same level in the DAG. • Parallel task execution time: average time per algorithm iteration to run in parallel tasks placed in the same level in the DAG, considering all overheads related to data movement. Figure 6 shows example DAGs generated by PyCOMPSs [ 68 ] for the two algorithms we consider: K-means and Matrix Multiplication (Matmul). For the task user code metrics we compute the average time per task of the same task type in the algorithm. Hence, for Matmul there are two tasks matmul_func (blue nodes in Figure 6b) and add_func (white nodes in Figure 6b). The data movement metrics are computed per CPU core involving all tasks (blue and white). And the parallel task execution time is computed per each level in the DAG (e.g., all blue nodes, the four white, the two white, etc.). This classification of metrics help us scrutinize thread and task parallelism in Section 5. 4.3 Evaluated factors Each variable that affects measured performance (i.e., the outcome of an experiment) and has several alternatives is called afactor [ 28 ]. In our analysis, we measure the metrics detailed earlier by varying the factors explained in this section. Table 1 presents a list of factors that affect the performance of task-based workflows. The list originates from experimentation, our experience, and our interviews with the team developed the distributed system COMPSs [ 4 ]. Some of these factors have been studied before (see Section 2), but never all together. We classify the factors into four dimensions, namely task algorithm, dataset, resources, and system employed. Each factor further affects a number of parameters. For example, when we vary the block dimension, then parameters such as block size, main 1 d3v1 (2) d4v1 2 d8v1 (2) d4v1 3 d12v1 (2) d4v1 4 d16v1 (2) d4v1 sync 5 d5v2 d9v2 d13v2 d17v2 d18v2 6 d20v1 (2) d21v1 7 d24v1 (2) d21v1 8 d27v1 (2) d21v1 9 d30v1 (2) d21v1 sync 10 d22v2 d25v2 d28v2 d31v2 d32v2 11 d34v1 (2) d35v1 12 d38v1 (2) d35v1 13 d41v1 (2) d35v1 14 d44v1 (2) d35v1 sync 15 d36v2 d39v2 d42v2 d45v2 d46v2 (a) DAG for K-means, grid dimension 4x1, 3 iterations main barrier 1 d1v1 d1v1 2 d3v1 d4v1 3 d6v1 d7v1 4 d9v1 d10v1 8 d1v1 d3v1 9 d3v1 d16v1 10 d6v1 d18v1 11 d9v1 d20v1 15 d1v1 d6v1 16 d3v1 d26v1 17 d6v1 d28v1 18 d9v1 d30v1 22 d1v1 d9v1 23 d3v1 d36v1 24 d6v1 d38v1 25 d9v1 d40v1 29 d4v1 d1v1 30 d16v1 d4v1 31 d26v1 d7v1 32 d36v1 d10v1 36 d4v1 d3v1 37 d16v1 d16v1 38 d26v1 d18v1 39 d36v1 d20v1 43 d4v1 d6v1 44 d16v1 d26v1 45 d26v1 d28v1 46 d36v1 d30v1 50 d4v1 d9v1 51 d16v1 d36v1 52 d26v1 d38v1 53 d36v1 d40v1 57 d7v1 d1v1 58 d18v1 d4v1 59 d28v1 d7v1 60 d38v1 d10v1 64 d7v1 d3v1 65 d18v1 d16v1 66 d28v1 d18v1 67 d38v1 d20v1 71 d7v1 d6v1 72 d18v1 d26v1 73 d28v1 d28v1 74 d38v1 d30v1 78 d7v1 d9v1 79 d18v1 d36v1 80 d28v1 d38v1 81 d38v1 d40v1 85 d10v1 d1v1 86 d20v1 d4v1 87 d30v1 d7v1 88 d40v1 d10v1 92 d10v1 d3v1 93 d20v1 d16v1 94 d30v1 d18v1 95 d40v1 d20v1 99 d10v1 d6v1 100 d20v1 d26v1 101 d30v1 d28v1 102 d40v1 d30v1 106 d10v1 d9v1 107 d20v1 d36v1 108 d30v1 d38v1 109 d40v1 d40v1 barrier 5 d2v2 d5v2 6 d8v2 d11v2 7 d12v2 d13v2 12 d15v2 d17v2 13 d19v2 d21v2 14 d22v2 d23v2 19 d25v2 d27v2 20 d29v2 d31v2 21 d32v2 d33v2 26 d35v2 d37v2 27 d39v2 d41v2 28 d42v2 d43v2 33 d45v2 d46v2 34 d47v2 d48v2 35 d49v2 d50v2 40 d52v2 d53v2 41 d54v2 d55v2 42 d56v2 d57v2 47 d59v2 d60v2 48 d61v2 d62v2 49 d63v2 d64v2 54 d66v2 d67v2 55 d68v2 d69v2 56 d70v2 d71v2 61 d73v2 d74v2 62 d75v2 d76v2 63 d77v2 d78v2 68 d80v2 d81v2 69 d82v2 d83v2 70 d84v2 d85v2 75 d87v2 d88v2 76 d89v2 d90v2 77 d91v2 d92v2 82 d94v2 d95v2 83 d96v2 d97v2 84 d98v2 d99v2 89 d101v2 d102v2 90 d103v2 d104v2 91 d105v2 d106v2 96 d108v2 d109v2 97 d110v2 d111v2 98 d112v2 d113v2 103 d115v2 d116v2 104 d117v2 d118v2 105 d119v2 d120v2 110 d122v2 d123v2 111 d124v2 d125v2 112 d126v2 d127v2 dislib.data.array._matmul_func dislib.data.array._add_func (b) DAG for Matmul, grid dimension 4x4 Figure 6: DAGs of partially & fully parallelizable algorithms grid dimension, and DAG shape are affected. We also consider an algorithm-specific parameter to study the computational complexity of tasks, e.g., the number of clusters in K-means. It is worth noting, that some presumably relevant parameters are determined from the factors and therefore, they are not directly included in our experiments; e.g., task dependencies may be extracted from the input algorithm, the block dimension determines the DAG shape, etc. Finally, we consider as future work other resource parameters, such as #GPU devices, RAM and GPU memory size, CPU-GPU bus throughput, and disk throughput. Our experiments indicate that the factors reported here are sufficient to detect relevant trends. 4.4 Experimental setup 4.4.1 Cluster configuration. We employed the Minotauro System, an HPC cluster hosted at BSC [ 13 ]. We used 8 out of the 38 nodes totally available. Each node has 16 CPU cores (Intel Xeon E5-2630 with 128 GB of RAM) and 4 GPU devices (NVIDIA K80, each with 12 GB of memory, using PCIe 3.0 for CPU-GPU interconnect). Hence, at most 128 tasks can be parallelized by CPUs and 32 by GPUs. 4.4.2 Distributed system. Our tested distribution system is PyCOMPSs (v.3.0) [ 68 ]. We experimented with two of its scheduling policies [12], one considering the task generation order and another one based on data locality. We also experimented with both local and shared disk storage running General Parallel File System (GPFS), which is commonly used in HPC environments. 4.4.3 Scripts. The software bits to run our analysis can be found in our public code repo 2 . We employ Python’s performance counter [ 51 ] to measure serial fraction execution, parallel 2 https://github.com/mnlcarv/Performance-Analysis-of-Distributed-GPUAccelerated-Task-Based-Workflows.git 695 Table 1: Factors and parameters Dimension Factors Parameters Task algorithm a) block dimension∗ ∥†‡§, b) computational complexity∥, c) parallel fraction∥, and d) algorithm-specific parameter∥ a) block size, grid dimension, and DAG shape b) - c) - d) - Dataset e) dataset dimension∗∥ †‡§ e) dataset size Resources f) processor type (i.e., CPU or GPU)∥, and g) storage architecture† f) maximum #CPU cores available depending on the processor type g) - System h) scheduling policy‡§ h) - System functions affected: Device Speedup (∥); Storage I/O (†); Network I/O (‡); CPU-GPU Data Transfer (∗); Task Scheduling (§) fraction execution (only for CPU-based tasks), CPU-GPU communication, and parallel task execution times. Since GPU executions run asynchronously with respect to CPU executions, we used CUDA events [ 21 ] to measure the parallel fraction execution for GPU-accelerated tasks. We used Paraver [ 48 ] to collect data deserialization and serialization times from traces automatically generated by PyCOMPSs runtime. 4.4.4 Algorithms. We used the K-means and Matmul implementations from dislib library (version 0.6.4) [ 17 , 23 ], a distributed version of the scikit-learn library [ 47 ]. The algorithms used have the following GPU-accelerated tasks and computational complexities: Matmul: matmul_func: 𝑂(𝑁3) and add_func: 𝑂(𝑁) , where 𝑁 is the order of the block (these two tasks share a dependency as well). K-means: partial_sum: 𝑂(𝑀𝑁𝐾2) , where 𝑀 and 𝑁 are the number of rows (i.e., samples) and columns (i.e., features) in a block, respectively, and 𝐾is the number of clusters. The Matmul tasks matmul_func and add_func have fully parallel user code. The K-means task partial_sum has partially parallel user code. The two algorithms use different chunking strategies. Matmul chunks the datasets into rows and columns, while K-means chunks the dataset into rows. As result, they yield differently shaped DAGs as shown in Figure 6. The Matmul DAG is wide and shallow, which implies a high level of task parallelism. The K-means DAG is narrow and deep, resulting in a low degree of task parallelism and a high level of task dependencies. As discussed in Section 3.5, to control the number of tasks generated in all executions we assign exactly one block per grid per task. Matmul implementation guarantees such task granularity by default, but in K-means we enforce it by setting the number of grid columns to 1. 4.4.5 Datasets. We generated synthetic datasets in the form of NumPy array with random float64 (double precision) values. To ensure reproducibility across multiple executions, we used a fixed random state value. We varied the dataset sizes to ensure that each algorithm tested reaches the GPU memory limits. As we discussed, varying the block size leads to different grid dimensions and hence, different (but predicable) number of tasks and task granularities. Smaller blocks generate a large number of finegrained tasks, whereas larger blocks generate a small number of coarse-grained tasks. And this allows us to stress the cluster resources for various scenarios of thread and task parallelism. Specifically, we used the following sizing scenarios: • Matmul: two datasets: 8 GB, 32K x 32K (1024M elements) and 32 GB, 64K x 64K (4B elements), and five grid dimensions: 1x1, 2x2, 4x4, 8x8, 16x16. • K-means: two datasets: 10 GB, 12.5M samples x 100 features (1250M elements) and 100 GB, 125M samples x 100 features (12.5B elements), and nine grid dimensions: 1x1, 2x1, 4x1, 8x1, 16x1, 32x1, 64x1, 128x1, 256x1. 5 EXPERIMENTS Our analysis spans four main experiments: (a) an end-to-end performance analysis, (b) profiling task user code processing, (c) profiling parallel task processing, and (d) conducting a correlation analysis of all the factors considered. Section 5.1 presents an end-to-end performance evaluation considering all system functions affected, i.e., CPU-GPU data transfer, limited device speedup, storage I/O, network I/O and scheduling overhead. Besides the costly CPU-GPU communication, task serial fraction may also cause a significant limitation in device speedup. As expected, data (de-)serialization dominates storage I/O and represents a critical bottleneck in distributed environments. Next, we follow a bottom-up approach to analyze performance from thread-level to task-level parallelism. Section 5.2 discusses the gains obtained by GPUs due to thread parallelism by testing different task workloads with different computational complexities. For a focused evaluation, we study the overheads caused by CPU-GPU communication and serial fraction processing. Section 5.3 extends the previous analysis by considering storage I/O, network I/O, and scheduling overheads caused by task distribution. For this experiment, we compare the execution times between CPUs and GPUs over different combination of storage architectures and scheduling policies. Section 5.4 presents a correlation analysis of all relevant factors, summarizing our findings and discussing relevant trends identified in our study. Finally, a discussion about the generalizability of our approach is provided in Section 5.5. Unless otherwise stated, in the experiments we use the experimental setup presented in Section 4.4 with the 8 GB and 10 GB datasets for Matmul and K-means, respectively, shared disk as storage architecture, and task generation order as a scheduling policy. We ran each experiment six times and discarded the first run to avoid fluctuations due to warm up processing, such as loading required modules, compile the GPU kernel, etc. 5.1 End-to-end analysis Figure 7 presents an end-to-end performance analysis for Matmul (Figure 7a) and K-means (Figure 7b). The analysis shows how factors, such as task computational complexity of each algorithm, block dimension (represented by block size and grid dimension in X-axis), and processor type (CPU or GPU) affect performance metrics, such as GPU speedup over CPU (top charts), execution times (bottom charts), parallel fraction (P. Frac - green line), CPU-GPU communication and serial fraction (blue line), and data serialization/deserialization (orange line). The left charts refer to the 8 GB and 10 GB datasets, and the right charts to the larger datasets 32 GB and 100 GB for Matmul and K-means, respectively. (blue line), and data serialization/deserialitzation (red line). Previous studies allude to CPU-GPU communication and (lack of) memory being limiting factors in GPU acceleration [ 60 ]. Our 696 GPU OOM GPU OOM GPU OOM GPU OOM (a) Matmul, 8GB left and 32GB right GPU OOM GPU OOM (b) K-means, 10GB left and 100GB right Figure 7: End-to-end performance analysis experiments corroborate this hypothesis. Speedups obtained in the parallel fraction scale with the block size. However, Figure 7 shows that GPUs’ memory limitations eventually lead to out-ofmemory errors for large task granularities. Similarly, the user code speedup is affected by the block size. Consider for example Matmul on the 8 GB dataset. In fine-grained tasks, communication dominates the parallel fraction computation and the user code speedup decreases on average about 35% compared to the parallel fraction speedup. The decrease is smaller for coarsegrained tasks (e.g., 20% in block size 2048 MB), as computation dominates communication. The same trend holds for both algorithms in all datasets tested. This result is aligned with a typical performance optimization strategy that suggests increasing the volume of data to process for improving GPU throughput [ 32 ]. However, Figure 7 reveals that this strategy is not always effective when considering serial fraction processing and data (de- )serialization. 5.1.1 Serial fraction processing. Figure 7b shows that user code speedups do not change significantly with the block size. In partially parallelizable algorithms, these speedups suffer from serial processing and CPU-GPU communication that dominate the parallel fraction for all block sizes. Moreover, both serial and parallel fractions scale with the block size in similar proportions. Hence, the gain obtained by GPUs in the parallel fraction are 0 10 20 GPU Speedup over CPU matmul_func Usr. Code add_func 32 128 512 2048 8192 Block size MB 10 3 10 1 101 103 Average Time per Task (s) P. Frac. CPU P. Frac. GPU CPU-GPU Comm. 32 128 512 2048 Block size MB GPU OOM GPU OOM Figure 8: Task computational complexity in Matmul diminished by the cost of serial processing, resulting in marginal speedups for all block sizes. 5.1.2 Data (de-)serialization. Figure 7 shows that the speedups of parallel tasks are largely affected by data (de-)serialization. An excess of fine-grained tasks saturates the nodes with more tasks than available CPU cores and the disk with an abundance of read/write processes. On the other hand, a small number of coarse-grained tasks do not fully utilize the available CPU cores in the cluster and increase the data volume per read/write processes, thus, increasing the cost of (de-)serialization that cannot be parallelized. Figure 7 shows that the maximum GPU speedup is obtained when the maximum parallelism in data (de)-serialization is achieved; i.e., when the number of tasks is equal to the number of available CPU cores. It is also worth noting that for small block sizes the parallel task GPU speedup over CPU is negative due to the relatively considerable communication and data movement overheads. For larger block sizes, these overheads are amortized with the increased task parallelism and hence, GPU speedup turns positive as we reach the maximum task parallelism (32 tasks in our settings). 5.1.3 Dataset size. Figures 7 shows that the same trends hold for both smaller (left charts) and larger (right charts) datasets. In fact, as the dataset size increases, there is also an increase in GPU speedup (on avg ∼5x) for parallel fraction and user code, whilst for parallel task the difference is negligible. Note that for the large datasets (32 GB for Matmul and 100 GB for K-means) our analysis is limited by the available GPU memory, which is 12 GB in our settings. We explain this in more detail in Section 5.3. This limitation does not allow testing block sizes larger than 2 MB (4x4 grid size) in Matmul and 6 MB (16x1 grid size) in K-means due to the GPU memory size needed to handle larger blocks along with the intermediate results produced in each algorithm. Summary. Higher task granularity does not always achieve higher GPU speedups over CPU, mainly due to serial processing, CPU-GPU communication, and data (de-)serialization overheads. Our findings are as follows. • Observation O1: User code speedups are not affected significantly by block size when parallel processing gains are diminished by the serial processing and CPU-GPU communication costs. 697 0 2 4 6 8 GPU Speedup over CPU 10 clusters Usr. Code 100 clusters 1000 clusters 3978 156 313 625 1250 2500 5000 10000 Block size MB 10 3 10 1 101 103 Average Time per Task (s) P. Frac. CPU S. Frac. GPU P. Frac. GPU S. Frac. GPU CPU-GPU Comm. 39 78 156 313 625 1250 2500 5000 10000 Block size MB 39 78 156 313 625 1250 2500 5000 10000 Block size MB GPU OOM GPU OOM CPU GPU OOM GPU OOM GPU OOM CPU GPU OOM (a) Varying clusters, dataset size 10GB, K-means (b) Varying data skew, Matmul 2GB, K-means 1GB, 10 clusters Figure 9: The effect of (a) algorithm-specific parameter (#clusters) in K-means and (b) data skew in Matmul and K-means • Observation O2: Parallel task speedups do not increase significantly for coarse-grained tasks, but can significantly improve when data (de-)serialization is fully parallelized using all available CPU cores. 5.2 Profiling task user code processing In this experiment, we investigate how task algorithm factors affect the task user code metrics. In particular, we explore the effect of (a) the computational complexities of matmul_func and add_func tasks in Matmul, (b) the algorithm-specific parameter, which in this experiment is the number of clusters (#clusters) in K-means, and (c) data skew in Matmul and K-means. 5.2.1 Computational complexity. We use Matmul that comprises two types of tasks with different computational complexities, namely matmul_func and add_func. Figure 8 shows how factors such as task computational complexity, block dimension (represented by block size in X-axis), and processor type affect performance metrics such as the user code GPU speedup over CPU (top charts), and execution times (bottom charts) related to the parallel fraction (green line) and CPU-GPU communication (purple line). Note that for the maximum task granularity (8192 MB) the matrix is multiplied in a single matmul_func task and no add_func task is needed (the reason we skip this value in Figure 8). Figure 8 shows that the parallel fraction execution time in matmul_func dominates CPU-GPU communication times in most cases. Consequently, speedups scale with block size and increase as high as 21x. However, Figure 8 shows that this pattern does not repeat in add_func because communication dominates parallel fraction computation in all block sizes, resulting in performance degradation of GPUs compared to CPUs. This happens because the computational complexity of add_func is two orders of magnitude less than the computational complexity of matmul_func. Therefore, the parallel fraction in add_func is too small to benefit from the massive thread parallelism provided by GPUs. 5.2.2 Algorithm-specific parameter. K-means has a single task, partial_sum, whose computational complexity is affected by an algorithm-specific parameter: #clusters. Figure 9a shows how factors such as task computational complexity with 10, 100, and 1000 clusters, respectively, block dimension (represented by block size in X-axis), parallel fraction and processor type affect performance metrics such as the user code speedup of GPU over CPU (top charts) and execution times (bottom charts) related to the parallel fraction (green line), serial fraction (yellow line), and CPU-GPU communication (purple line). For 10 clusters, the computational complexity is so low that the parallel fraction execution time is less than the serial fraction and CPU-GPU communication times, resulting in marginal speedups (no more than 1.5x). For 100 clusters, the computational complexity is higher, making the parallel fraction execution time greater than the CPU-GPU communication time, but still less than the serial fraction execution time. In this case, the speedups are increased in about two times over the scenario with 10 clusters. Finally, for 1000 clusters, the parallel fraction keeps dominating the CPU-GPU communication, but the gap between the serial and parallel fraction execution times is reduced. Therefore, the speedups are up to 7x higher than in the scenario with 10 clusters. Note that the speedups keep scaling with #clusters until GPU’s memory capacity is reached (OOM). Figure 9a shows that the speedups do not scale with the block size. This happens because the effect of #clusters dominates the effect of block dimension in the computational complexity of partial_sum task. And this is expected as #clusters has a quadratic impact in the computational complexity of partial_sum, while block dimension has a linear impact (see also Section 4.4.4). 5.2.3 Data Skew. Next, we investigate how data skew may affect our analysis. We generated two skewed datasets of size: 2 GB (16K x 16K, 256M elements) for Matmul, and 1 GB (1.25M samples x 100 features) for K-means. For doing so, we adapted the uniform distribution of the NumPy random routine [ 43 ] to move 50% of the elements to certain regions of the distribution forcing groups of elements in the dataset. Figure 9b compares the CPU (top charts) and GPU (bottom charts) task user code execution time in Matmul (left charts) and K-means (right charts) for the uniform (0% skew) and skewed datasets. We observe that data skew does not affect the task user code execution time. Our analysis also indicates similar performance for the fine-grained stages of parallel and serial fraction, CPU-GPU communication, and parallel task (the respective charts are omitted due to space considerations). This result is expected as the algorithms tested do not process differently tasks involving uniform or skewed data. In general, we believe that data skew might have an impact in certain pipelines that exploit data distribution in a special manner; how this could affect processing in GPUs is an excellent topic for future work. 698