Queuing Network Models to Predict the Completion Time of the Map Phase of MapReduce Jobs
Daniel A. Menascé, Shouvik Bardhan · Int. CMG Conference · 2012
Big Data processing is generally defined as a situation when the size of the data itself becomes part of the computational problem. This has made divide-and-conquer type algorithms implemented in clusters of multi-core CPUs in Hadoop/MapReduce environments an important data processing tool for many organizations. Jobs of various kinds, which consists of a number of automatically parallelized tasks, are scheduled on distributed nodes based on the capacity of the machines. A key challenge in provisioning such jobs in a Hadoop/MapReduce cluster is to be able to predict their completion times based on various job characteristics. Standard makespan computations that ignore the contention on compute nodes significantly underestimate a job’s completion time. This paper proposes a mathematically sound model based on closed Queuing Networks for predicting the execution time of the map phase of a MapReduce job. The model captures contention at compute nodes and parallelism gains due to increased number of slots available to map tasks. Experiments validated the model on a single as well as a 2-node Hadoop environment. We ran experiments for different input split sizes and different map slot sizes to validate our model.