Optimization Techniques For A Distributed In-Memory Computing Platform By Leveraging SSD Part 2
Aug 17, 2023
3.1. Cluster Environment
Figure 1 shows our testbed cluster consisting of one name node (master) and four data nodes (slaves). In the name node (master), we configured the NameNode and Secondary NameNode of Hadoop (HDFS) and the Driver Node (master node) of Spark. In each data node, we run the DataNode of Hadoop (HDFS) and Worker Node of Spark. The name node and data node machines have the same H/W environments (3.4 GHz Xeon E3-1240V3 QuadCore Processor with hyper-threading), except for the amount of main memory (8 GB for the name node and 4 GB for each data node).
Namename is the Master node in the Hadoop architecture, responsible for managing and monitoring the file system of the entire Hadoop cluster. The Namename node is also one of the critical nodes of the entire Hadoop cluster, and its performance and reliability will directly affect the operating efficiency and availability of the entire Hadoop cluster.
There are many indicators related to the Namename node, one of the most important indicators is memory. The Namename node requires a lot of memory to store and manage the namespace of the entire HDFS file system, which includes metadata information of files and directories, such as file names, permissions, timestamps, file sizes, and so on.
The memory of the Namename node not only determines the number of files it can manage and the size of the file system but also affects the performance and reliability of the Hadoop cluster. If the Namename node has insufficient memory, it will not be able to respond quickly to client requests, resulting in reduced throughput of the entire Hadoop cluster. In addition, if the Namename node fails, the metadata information it stores may be lost, making the entire HDFS file system unavailable.
Therefore, in the Hadoop cluster, the memory of the Namename node is crucial. It is recommended that administrators select the appropriate Namename node hardware configuration based on specific business needs, and regularly monitor the performance and availability of Namename nodes to ensure that they can provide efficient and reliable services for the entire Hadoop cluster. It can be seen that we need to improve our memory. Cistanche can significantly improve memory because meat paste is a traditional Chinese medicinal material with many unique effects, one of which is to improve memory. The efficacy of minced meat comes from various active ingredients, including carboxylic acid, polysaccharides, flavonoids, etc. These ingredients can promote brain health through various channels.

Click know supplements to boost memory
We used two SSDs as storage spaces where a 120 GB SATA3 SSD is used for the operating system, and a 512 GB SATA3 SSD is equipped for the HDFS, respectively. In addition, the 512 GB SATA3 SSD can be effectively leveraged for expanding the bandwidth of insufficient main memory to cache the RDDs of Spark. All nodes including the name node and data node are connected with a 1 Gb Ethernet switch, as seen in Figure 1. Table 2 shows the summary of hardware and software configurations in each data node of our testbed cluster.


3.2. Spark JVM Heap
A Spark job runs as a Java process on the Java Virtual Machine (JVM), and Spark exploits Scala, a functional language extended from Java. The worker process of Spark also runs on the JVM of each data node, so that on each data node, the worker process has the JVM heap in the main memory as depicted in Figure 2. When Spark submits a job, the worker process that has the JVM heap executes the job as distributed tasks.

We can customize the ratio of the JVM heap size of a Spark worker through the configuration file spark-defaults. conf in the spark/conf/ directory. In the spark defaults.conf file, the value of spark.executor.memory is the JVM heap size where the default is 512 MB that each worker node can utilize in the data node. In addition, the value of spark.storage.safetyFraction is fixed as 0.9, which means that Spark can use up to 90% of the JVM heap size (also known as safety area). This is to prevent the JVM from generating OOM (out of memory) errors due to the lack of available main memory during the task processing.
In this safety area, the overall JVM heap space is divided into three sub-regions: unroll, storage, and shuffle spaces, as shown in Figure 2. The unroll space is used for unrolling data blocks in memory. When an RDD is cached on other storage media such as an SSD or HDD not on the main memory, the RDD should be serialized. Then, when Spark reads this RDD back to the memory, the RDD has to be unrolled. The storage space is used for caching an RDD. If the storage space is not sufficient for caching the RDD, some RDDs can be evicted from this space based on the LRU (least recently used) policy, or they can be cached on other storage media, such as an SSD. The shuffle space is used for shuffling the intermediate data. This shuffle space can play an important role in iterative applications such as machine learning since it can substantially affect the overall job completion time.
In the default Spark configuration, the storage and shuffle spaces of the JVM heap have capacity fraction ratios of 0.6 and 0.2, respectively (i.e., 60% of the safety area for the storage and 20% for the shuffle). The unroll space takes 20% of the storage space by default. The capacity of these three spaces of the JVM heap can be set by a spark. storage.unrollFraction, spark.storage.memoryFraction, and spark.shuffle.memoryFraction. For example, in our testbed cluster, we can set the spark.executor.memory as 2.6 GB of the 4 GB memory of the worker node, which means that the JVM heap size is set to a maximum of 2.6 GB. Then, the actual capacities of storage space and shuffle space are 2.6 GB × 0.9 × 0.6 = 1.4 GB and 2.6 GB × 0.9 × 0.2 = 0.46 GB, respectively. Accordingly, the unroll space takes 1.4 GB × 0.2 = 0.28 GB.
3.3. RDD Caching Policy
The Spark platform offers diverse RDD caching options involving main memory and disks. The default option is the MEMORY_ONLY, where the RDD is maintained in the storage space described in Section 3.2 as a non-serialized Java object. If this storage space is insufficient for holding all RDDs, some of them would be evicted from the main memory based on a pre-defined cache replacement policy. However, whenever a noncached RDD is required for task processing, this RDD should be re-created based on the lineage information which can result in substantial performance degradation in this MEMORY_ONLY caching policy.
Besides the MEMORY_ONLY option, Spark provides alternative MEMORY_AND_DISK, DISK_ONLY, and OFF_HEAP options. The MEMORY_AND_DISK option stores RDDs in the non-volatile disk when the storage space is not enough to store all required RDDs. Disks can consist of HDDs or SSDs; however, normal spindle disks have relatively poor read/write throughput, so the overall execution time can be longer than that of the MEMORY_ONLY caching option. To address this problem, we can effectively leverage SSDs, which can potentially reduce the overall job completion time compared to the normal HDD-based approach.

The DISK_ONLY option stores RDDs only in non-volatile storage devices such as HDDs or SSDs, i.e., not in the main memory. A cluster that does not have a sufficient amount of available memory can achieve good performance with this option. In this case, since RDD is stored only in disk media, the shuffle space can be extended instead of using the storage space of the memory. As a result, when running an application such as PageRank, which generates a relatively large amount of shuffle data, we can observe better performance than in the MEMORY_ONLY case.
The OFF_HEAP option enables Spark to use off-heap space, which is outside of the Java garbage collector’s management. Thus, if we use off-heap space, we must deal with complicated memory operations such as allocation/deallocation and serialization/deserialization. Therefore, for practical purposes, we do not use the OFF_HEAP configuration.
3.4. Optimization Methodology
As we discussed in Sections 3.2 and 3.3, our optimization methods include (1) the configuration of the Spark JVM heap and (2) the experimental options of the RDD caching policy as follows:
1. Spark JVM heap configuration: We investigated the effects of changing the capacity fraction ratios of the shuffle and storage spaces. The ratio of shuffle and storage space is 60%:30%, 50%:40%, and 20%:60%, respectively. The “20%:60%” shuffle and storage ratio is the default value in the Spark setup. We choose “60%:30%” to contrast the result with sufficient shuffle space and configure “50%:40%” to show the performance in a balanced way.
2. RDD caching policy: We also examined the effects of different RDD caching policies. We compared the performance of various policies such as OFF_HEAP, MEMORY_ONLY, MEMORY_AND_DISK, and DISK_ONLY, where DISK denotes the SSD in this experiment.
Table 3 shows a total of 12 different experimental configurations based on the RDD caching policies and Spark JVM capacity fraction ratios. In the experiment configurations labeled with “_1” (for example, “N_1”), we set 60% of the Spark JVM heap for shuffling and 30% for storage spaces. With those labeled “_2”, we set 50% of the Spark JVM heap for shuffling and 40% for storage. Finally, for those labeled with “_3”, we set 20% of the Spark JVM heap for shuffling and 60% for storage, as can be seen from the “Option”, “Shuffle”, and “Storage” columns in Table 3. Note that our testbed cluster’s maximum memory size of the executor is 2.7 GB, i.e., each worker node has 2.7 GB as the Spark JVM heap size.

In terms of the RDD caching policy, the “N” option is not to cache the RDD, “M” option is to cache the RDD on the memory only, “M&S” option is to cache the RDD on the memory and SSD together, and finally, “S” option is for caching the RDD on the SSD only.
Through our experiments, we propose optimizing strategies that can achieve the best performance from the cluster which has insufficient memory amounts by carefully adjusting the Spark JVM heap configuration and employing an effective RDD caching policy, as we will see in Section 4.
4. Experimental Results and Analysis
4.1. 500 MB PageRank Experiments
4.1.1. Results with Changing JVM Heap Configurations
Figure 3 shows the experimental results of each stage in the PageRank workload by changing the JVM heap sizes. In the Distinct stage, Spark reads the input data and distinguishes the URL and links. As we can see from the results of the Distinct0 stage, the overall execution time decreases by changing the JVM heap sizes from _1 and _2 to _3 options, mainly due to the garbage collection (GC). For example, the GC time takes 25 s, 24 s, and 16 s in M&S_1, M&S_2, and M&S_3, respectively. Therefore, in the Distinct0 stage, as we increase the amount of storage space, we can improve the overall performance by reducing the GC time. On the other hand, in the Distinct1 stage, the overall execution time increases as we change the options from _1 and _2 to _3. This is mainly because of the shuffle spill. When we checked the Spark web UI, the shuffle data were spilled onto the disk because of the lack of shuffle memory space. For example, the sizes of shuffle spill data on disk in M&S_1, M&S_2, and M&S_3 are 0, 220 MB, and 376 MB respectively. When the shuffle spill occurs, the CPU overheads for spilling the data onto the disk increase because the data need to be serialized.

After the Distinct stages, there are iterative flatMap stages to obtain ranks. FlatMap stages generate a lot of shuffle data, which can make our cluster lack the necessary shuffle memory space. Therefore, as the available amount of shuffle space decreases (in order from options _1, _2, and _3), the more shuffle spill can occur, which can potentially affect the overall job execution time (e.g., M&S option flatMap2 stage _1: 37 s, _2: 40 s, _3: 49 s). However, when the data are cached on memory only (i.e., M_1, M_2, and M_3), they show another pattern. The main reason for this behavior is that the Spark scheduler schedules the tasks unevenly because there is a lack of memory storage space for caching the RDD on options _1 and _2. If a worker does not have RDDs, he is excluded from the scheduling pool. Therefore, the other workers have to handle additional tasks with GC overheads which can affect the whole job execution time.
4.1.2. Results with Changing RDD Caching Options
First of all, Distinct stages are not affected by changing the RDD caching policy but only by memory usage. The stages that are affected by the RDD caching option are flatMap stages since during the shuffle phase, cached RDDs are used again.

In Figure 4, the graph is normalized by the N_1 option that does not cache the RDD and _1 memory configuration to check the performance difference. When comparing only the graphs of _1, in the order of M_1, M&S_1, and S_1, there is a 32% performance degradation in M_1 and 30% and 20% performance improvements with M&S_1 and S_1, respectively. With the M_1 option, the reason for relatively poor performance is that RDDs are cached unevenly due to the lack of storage, which will result in uneven scheduling as we previously mentioned. This means the JVM heap space is insufficient for shuffling the data and saving the RDDs.

To address this problem, we spread the RDDs to cache both on the memory and SSD, which can improve the performance as shown with the M&S_1 option. Caching the RDD on memory improves the access speed for the RDD, and caching the RDD on the SSD can avoid the shuffle spill by effectively extending the available shuffle space in memory. With the S_1 option that has shown a 20% performance improvement, the RDD is cached on the SSD only. The shuffle spill is reduced by caching the RDD on the SSD. However, it has achieved a lower performance improvement than M&S_1, where the RDD is mainly cached on memory and reused from the memory.
In the default configuration of Spark, which is option _3, we can see that in the order of M_3, M&S_3, S_3, and N_3, overall performance decreases. In the default configuration, the storage of the JVM heap is enough for the RDD to be cached with balance. Therefore, the overall performance mainly depends on the performance of the memory device used. However, we can still see the best performance with the M&S_1 option since we can effectively reduce the GC time and shuffle spill by caching the RDD both on the memory and SSD.
4.2. 1 GB PageRank Performance
We experimented with the PageRank workload by increasing the size of data from 500 MB to 1 GB. Figure 5 shows the different behaviors of the system compared to PageRank for the 500 MB dataset. We can see some failed jobs that could not succeed in completing the job until the take6 stage (e.g., N_1, N_2, M_1, M_2, M_3, M&S_3). Among these failed jobs, there are ones that failed in the flatMap2 stage, which are N_1, N_2, and M_1. The reason for the job failure is the lack of storage memory. The GC occurs when the RDD is cached on insufficient memory. Because of this GC overhead, the Spark executor receives an ExecutorLostFailure exception.
M_2, M_3, and M&S_3 could proceed with the processing until the flatMap2 stage; however, after this, failure occurs. M&S_3 operates similarly to the M_3 until the flatMap2 stage because when the M&S_3 option is used, there is sufficient memory to cache the RDD. After the flatMap2, the OutOfMemory error occurs due to the lack of shuffle memory space in the flatMap3 stage.
4.2.1. Results with Changing JVM Heap Configuration
The Distinct0 stage shows very similar results to the 500 MB dataset, and the overall performance improves in the order of options _1, _2, and _3. This is because the GC time is reduced to 78 s, 59 s, and 28 s, respectively
On the other hand, in the Distinct1 stage, it showed different results for the 500 MB dataset. In the 500 MB dataset experiment, we can see the performance gain by increasing the shuffle space of memory. However, in the 1GB dataset experiment, the executor memory of the worker node cannot accommodate the large size of data. Therefore, the shuffle space of memory becomes relatively insufficient. For example, the amounts of shuffle spill for options _1, _2, and _3 are 575.5 MB, 813.8 MB, and 843.4 MB, respectively, and the GC time takes 33 s, 10 s, and 8 s, respectively. As we mentioned before, when a shuffle spill occurs, the RDD needs to be serialized so that the CPU computations can increase, which can result in overall performance degradation.

4.2.2. Results of Changing RDD Caching Policy
For analyzing execution time by changing the RDD caching policy, as we can see from Figure 6, we exclude distinct stages from Figure 5. This is because we do not need to analyze distinct stages since there are no changes caused by changing the RDD caching policy.
Interestingly, there are no changes with various JVM heap configurations in flatMap stages as opposed to the case of the 500 MB dataset. The reason for this is that the shuffle spill occurs in all configurations because there is insufficient memory. The overall execution time by changing the RDD caching option increases in the order of M&S, S, N, and M. (M&S is the fastest option.) In option N, the ExecutorLostFailure error occurs because there is insufficient memory space. In option M, when the RDD is cached on memory, GC overhead occurs because there is insufficient memory space. Even if the RDD is cached on memory, the job fails because of the ExecutorLostFailure error that occurs when the shuffle memory space is insufficient (OutOfMemory).
In such low available memory situations, M&S and S options can be effective alternatives. In the M&S_1 option, we increase the accessibility of the RDD by caching the RDD using both the memory and SSD. As a result, there is a performance improvement for the same reason as the 500 MB dataset. In addition, there is sufficient shuffle memory space due to caching the RDD on the SSD. As seen in Figure 6, the M&S_1 option becomes the fastest option in this experiment (M&S_1:0.6, S_1:0.63, T_1 0.64).

4.3. TC Experiment Analysis
Figure 7 shows the results of TC (transitive closure) experiments that use the input data involving 50,000 edges and 25,000 vertexes generated at random. The iteration number is 10. Through the iterations, the number of tasks is doubled at each iteration, and therefore, the size of the RDD increases and the amounts of shuffle read and write also increase. In the last iteration, the number of tasks becomes 4096. As there are more iteration stages, there is a larger effect on the total job execution time, and the last iteration stage is the biggest, consisting of many tasks that can lower the overall performance.

As we can see from Figure 7, the performance improves in the order of _3, _2, and _1 on the M, M&S, and S options, which means that obtaining a sufficient shuffle memory of the JVM heap is helpful. With the M option, the performance of option _1 is 18% faster than option _3, whereas, in the M&S option, the performance of option _1 is 3% faster than _3. In the S option, the performance of _1 is 2% faster than _3.
When we focus on changing the RDD caching option, the performance of option S_1 is 42% faster than N_1, and it is also 31% faster than M_1. The reason for the performance gain of job execution time is dependent on the last iteration stage. The key factor affecting the last iteration stage is the shuffle read blocked time. Shuffle read blocked time occurs when the RDD executed in the previous stage is read from another worker node through the network because of the lack of executor memory.
Even if each task has a performance gain of about 1–2 s through solving the shuffle read blocked time, we can achieve a significant performance gain because, in the last state, the number of tasks is quite large (i.e., 4096). In addition, one of the main factors affecting job execution time is the counting stage, which counts how many edges the TC matrix has at the last job.
With option N, because there are no RDDs cached in the counting stage, Spark reads the shuffle data executed from the previous stage, which takes 60 s. In addition, in option M, the RDD is not cached on memory because of the lack of executor memory. As a result, it also takes 60 s. However, in the M&S and S options, the RDD can be cached on the memory and SSD, so that it takes only 2 s in the count stage.
4.4. TeraSort Experiment Analysis
Figure 8 shows the experimental results of the TeraSort benchmark that uses a 10 GB dataset by changing the JVM heap configuration and RDD caching option. This graph is normalized by option N_1. We can see that all of the job execution times are similar; the difference between them is less than 5%. In the TeraSort workload, there were no performance improvements or degradations by changing the configurations and options. In the sorting stage, there are a few shuffles through the network. However, the sizes of shuffle read and shuffle write are 25 MB each, which is quite small compared to PageRank and TC. Therefore, the JVM heap configuration and RDD caching option do not affect the performance. Furthermore, the TeraSort workload does not consist of iterative jobs as in the transitive closure, so there is no benefit from RDD caching in the previous stage.

4.5. K-Means Clustering Experiment Analysis
The normalized job completion time of k-means clustering for the 1.5 GB dataset is shown in Figure 9. The purpose of k-means clustering is to find the k clusters in the dataset based on the distance measurement (e.g., Euclidean distance). In this workload, the algorithm reduces the SSE (sum of squared error) [24] by iterating the distance calculation between the k center points and each data point. In this experiment, we iterate this process eight times. The amount of data to be shuffled is minimal because the data needed from the previous stage are the information about the center points and SSE in each stage. In our k-means clustering workload, the maximum amount of shuffle read/write data is 1.0 MB, and the minimum is 0.8 MB. The shuffle spill does not occur here because the shuffle space is sufficient in all settings. In the experiments with no caching options, there is no difference among options _1, _2, and _3, because these settings do not cache any RDD, and in all three settings, the shuffle space is sufficient.

When caching the RDDs in the main memory or the memory and SSD, the more the storage space for the RDD, the more the performance in the job execution time improves because more RDDs can be cached on the storage space. When comparing the memory_only option and memory_and_SSD option, the memory_and_SSD option showed better performance improvement. This is because, in the memory_only option, the storage space is insufficient even in the M_3 option. In addition, caching the RDDs on the SSD solves such a lack of storage memory. Memory_and_SSD options improved performance by 10% on average as compared with the memory_only option.
Note that the k-means clustering workload shows an opposite performance tendency from PageRank and transitive closure workloads because of the difference in the amount of shuffle data. We will discuss this in more detail in the following subsection.
5. Discussion and Summary
5.1. Discussion
We analyzed the main factors of potential performance degradation problems based on the feature of workload and processing stages. Our extensive experimental results are summarized concerning applying the performance optimization techniques of the Spark platform for various workloads as follows:
• The performance degradation by Java garbage collection: In the PageRank workload with the 500 MB dataset and 1GB dataset, the GC occurs when there is insufficient storage space of the JVM heap to store the RDD. In the Distinct0 stage which reads the input file from the HDFS and caches it into the RDD, the GC occurs. We expand the storage space of the JVM heap through the configuration to solve this GC problem. We can improve performance to reduce GC because the storage space of the JVM heap can be expanded. In Figures 3 and 5, with the same RDD caching option, the _3 configuration shows the best performance in the Distinct0 stage. In addition, in PageRank with the 1 GB dataset, some options fail in the flatMap stage due to the lack of memory. The GC overhead increases so much that the stage fails or goes into an infinite loop. Thus, we construct the cluster with SSDs to solve this problem. It shows a performance improvement and succeeds in the job that failed using only memory, as seen in Figure 6, M&S_1 and S_1.
• The performance degradation by shuffle spill: In the PageRank workload with the 500 MB dataset and 1 GB dataset, in the flatMap stage, we can see that option M&S_1 shows the best performance because it has the least amount of shuffle spill (Figure 4: M&S_1 is 30% faster than N_1; Figure 6: M&S_1 is 40% faster than N_3). PageRank has many shuffle tasks. Thus, when the shuffle space of the JVM heap is insufficient for shuffling the data through the network, the shuffle spill occurs. Therefore, to reduce the shuffle spill, expanding the shuffle space of the JVM heap becomes the key factor of performance improvement.
Furthermore, we can improve the performance by storing the RDD both on the memory and SSD. This can make the executor expand the shuffle memory of the JVM heap to reduce shuffle spill. If there are more iterations, the performance from the flatMap stage would be the key point of the performance improvement. In the 1GB dataset experiment, the job execution time of S_3 is the best option, because RDDs are cached on the SSD only, and there is sufficient heap memory on the executors. Thus, in the S_3 option, the Distinct stages are faster than any other option. However, if the iteration number increases, the flatMap stage affects the job execution time. Thus, the M&S_1 option can achieve great performance in this case. Through these analyses, we can identify that shuffle has a key effect on job completion time. Thus, we have to expand the shuffle memory of the JVM heap and cache the RDD both in the memory and SSD to obtain enough shuffle memory space for preventing shuffle spill.
• The performance degradation by shuffle read blocked time: There is shuffle read blocked time on the TC workload. It occurs when there are a lot of tasks in the stage and each task needs to read the previous RDD through the network. As a result of the TC experiment (Figure 7), the M&S option is faster than option M. In the same RDD caching option, expanding the shuffle space of the JVM heap is faster than expanding the storage space. The reason for improved performance is that by expanding the shuffle space of the JVM heap, the shuffle read blocked time decreases in each task.
5.2. Summary: Which Is the Best Way?
In the comprehensive experimental results, there is not a single best setup to boost all workloads, since each of these workloads has different characteristics, even in its job lifetime. However, we can still propose how to optimize the configurations of a distributed in-memory computing platform by considering the variety of target workloads as follows:
• Spark JVM heap configuration—shuffle area vs. storage area: According to the experimental results of four different workloads, we can observe the performance differences depending on workload characteristics. For example, PageRank is a typical example of having a large amount of shuffle data so allocating more memory to the shuffle portion improves the overall performance. However, in the case of k-means clustering, the more we allocate to storage memory, as opposed to shuffle memory, the less execution time is required. Therefore, if we can adjust the JVM memory allocation percentage dynamically according to the workload characteristics, we can optimize the total execution time. The Hadoop YARN [25] enables us to assign jobs to different types of clusters (configurations) so that we can apply this idea to a large-sized Hadoop cluster to meet memory characteristics for various types of jobs.
• RDD caching policy—memory vs. SSD: In most cases, the SSD-backed memory caching shows the best performance unless all of the RDDs can fit within the actual main memory. Therefore, the SSD-assisted memory caching policy can be a viable choice for challenging workloads requiring substantial amounts of main memory that cannot be met by any single node in a cluster.
6. Conclusions
In this paper, we have investigated the main factors for the performance degradation of the Spark system running on top of a commodity-server-based computing cluster with insufficient available main memories. After experimentation and analysis, we presented alternatives that can improve the overall performance.
Java garbage collection occurs when the storage space of the JVM heap is insufficient due to a lack of physical memory. Java GC makes tasks wait for garbage collection so that the overall job completion time increases. The shuffle spill occurs when the shuffle space of the JVM heap is insufficient during the shuffle phase. Shuffle spill increases the CPU overhead to perform serialization for spilling intermediate shuffle data to the disk due to the lack of shuffle space. In the TC workload experiment, shuffle read blocked time makes the task wait for reading shuffle data through the network because of the lack of shuffle space. All of these factors can potentially increase the overall job completion time which can seriously affect the performance of the Spark system.
To address these problems, we construct a cluster with an SSD and cache the RDD both on the memory and SSD separately by utilizing the SSD to supplement the storage space of the memory. In addition, we adjust the JVM heap configuration for expanding the shuffle space. As a result, we could achieve a 30% performance improvement for the PageRank workload and a 42% performance improvement for the TC workload. We have identified that the shuffle spill can be a key factor of performance degradation and showed through experimentation that in workloads consisting of several iterations and shuffling, expanding the shuffle space can provide significant performance gains. In addition, we found that different memory usage patterns of jobs can affect the total execution time depending on the storage/shuffle memory percentage allocation in the JVM. According to the performance analysis of PageRank and k-means clustering, memory allocation in the JVM that is well-tuned to the workload characteristics can significantly improve job completion time.
Integrating these findings into the Spark platform would be one of our future works. For example, if workloads can be characterized in terms of the amounts of shuffle data, an optimized configuration can be automatically applied to accelerate the processing of target workloads. Therefore, in heterogeneous server configurations, developing a workload memory usage-aware scheduling system can improve the overall performance of a Spark-based cluster.
Author Contributions:
Conceptualization, J.L. (Jaehwan Lee); methodology, J.L. (Jaehwan Lee) and J.C.; software, J.C. and J.L. (Jaehyun Lee); validation, J.C., J.L. (Jaehyun Lee) and J.L. (Jaehwan Lee); investigation, J.L. (Jaehwan Lee) and J.-S.K.; resources, J.L. (Jaehwan Lee) and J.-S.K.; data curation, J.C. and J.L. (Jaehyun Lee); writing—original draft preparation, J.C. and J.L. (Jaehyun Lee); writing— review and editing, J.L. (Jaehwan Lee) and J.-S.K.; visualization, J.L. (Jaehyun Lee); supervision, J.L. (Jaehwan Lee) and J.-S.K.; project administration, J.L. (Jaehwan Lee) and J.-S.K.; funding acquisition, J.L. (Jaehwan Lee). All authors have read and agreed to the published version of the manuscript.

Funding:
This research was supported by the Basic Science Research Program (NRF-2020R1F1A1072696) through the National Research Foundation of Korea (NRF) funded by the Ministry of Science and ICT, GRRC program of Gyeonggi Province (No. GRRC-KAU-2017-B01, “Study on the Video and Space Convergence Platform for 360VR Services”), and ITRC (Information Technology Research Center) support program (IITP-2021-2018-0-01423).
Institutional Review Board Statement:
Not applicable.
Informed Consent Statement:
Not applicable.
Data Availability Statement:
Available upon request.
Conflicts of Interest:
The authors declare no conflict of interest.
References
1. Dean, J.; Ghemawat, S. MapReduce: Simplified data processing on large clusters. Commun. ACM 2008, 51, 107–113. [CrossRef]
2. The Apache Hadoop Project: Open-Source Software for Reliable, Scalable, Distributed Computing. Available online: https: //hadoop.apache.org/ (accessed on 10 September 2021).
3. Shvachko, K.; Kuang, H.; Radia, S.; Chansler, R. The Hadoop distributed file system. In Proceedings of the 2010 IEEE 26th symposium on mass storage systems and technologies (MSST), Incline Village, NV, USA, 3–7 May 2010; pp. 1–10.
4. Zaharia, M.; Chowdhury, M.; Franklin, M.J.; Shenker, S.; Stoica, I. Spark: Cluster computing with working sets. HotCloud 2010, 10, 95.
5. Ousterhout, K.; Rasti, R.; Ratnasamy, S.; Shenker, S.; Chun, B.G. Making sense of performance in data analytics frameworks. In Proceedings of the 12th USENIX Symposium on Networked Systems Design and Implementation (NSDI), Oakland, CA, USA, 4–6 May 2015; pp. 293–307.
6. Xing, W.; Ghorbani, A. Weighted PageRank algorithm. In Proceedings of the IEEE Second Annual Conference on Communication Networks and Services Research, Fredericton, NB, Canada, 21 May 2004; pp. 305–314.
7. Chakradhar, S.T.; Agrawal, V.D.; Rothweiler, S.G. A transitive closure algorithm for test generation. IEEE Trans. Comput.-Aided Des. Integr. Circuits Syst. 1993, 12, 1015–1028. [CrossRef]
8. O’Malley, O. Terabyte Sort on Apache Hadoop. Yahoo. May 2008. pp. 1–3. Available online: http://sortbenchmark.org/ YahooHadoop.pdf (accessed on 10 September 2021).
9. K-Means Clustering. Available online: https://en.wikipedia.org/wiki/K-means_clustering (accessed on 10 September 2021).
10. Zaharia, M.; Chowdhury, M.; Das, T.; Dave, A.; Ma, J.; McCauly, M.; Franklin, M.J.; Shenker, S.; Stoica, I. Resilient distributed datasets: A fault-tolerant abstraction for in-memory cluster computing. In Proceedings of the 9th USENIX Symposium on Networked Systems Design and Implementation (NSDI), San Jose, CA, USA, 25–27 April 2012; pp. 15–28.
11. Davidson, A.; Or, A. Optimizing Shuffle Performance in Spark; Technical Report; Berkeley-Department of Electrical Engineering and Computer Sciences, University of California: Berkeley, CA, USA, 2013.
12. Nicolae, B.; Costa, C.H.A.; Misale, C.; Katrinis, K.; Park, Y. Leveraging Adaptive I/O to Optimize Collective Data Shuffling Patterns for Big Data Analytics. IEEE Trans. Parallel Distrib. Syst. 2017, 28, 1663–1674. [CrossRef]
13. Zhang, H.; Cho, B.; Seyfe, E.; Ching, A.; Freedman, M.J. Riffle: Optimized Shuffle Service for Large-Scale Data Analytics. In Proceedings of the Thirteenth EuroSys Conference; EuroSys ’18; Association for Computing Machinery: New York, NY, USA, 2018. [CrossRef]
For more information:1950477648nn@gmail.com






