Efficient Straggler Replication in Large-scale Parallel Computing
Da Wang, Gauri Joshi, Gregory Wornell
Introduction
In cloud computing, large-scale sharing of computing resources provides users with great flexiblity and scalability. Computing frameworks such as MapReduce and Apache Spark are developed to harness these benefits. These frameworks employ massive parallelization by dividing a large job into many tasks that can be executed parallely on different machines. These frameworks can be used to run optimization and machine learning algorithms that can be easily divided into independent parallel tasks, for example alternating direction method of multipliers (ADMM) and Markov Chain Monte-Carlo (MCMC) .
The execution time of a task on a machine is subject to stochastic variations due to co-hosting, virtualization and other hardware and network variations . Thus, a key challenge in executing a job that consists of a large number of parallel tasks is the latency in waiting for the slowest tasks, or the “stragglers” to finish. As pointed out in [5, Table 1], the latency of executing many parallel tasks could be significantly larger ( ms) than the median latency of a single task ( ms).
In this work we provide a mathematical framework to analyze how replication of straggling tasks affects the latency and the cost of computing resources, and propose better scheduling policy designs.
The idea of replicating tasks in parallel computing has been recognized by system designers , and first adopted at a large scale via the “backup tasks” in MapReduce . A line of systems work and references therein further developed this idea. For example, Apache Spark implements “speculative execution” to allow relaunching slow running tasks.
While task replication has been studied in systems literature and also adopted in practice, there is not much work on mathematical analysis of replication strategies. Replication strategies are analyzed in , mainly for the single task case. In this paper we consider task replication for a job consisting of a large number of tasks, which corresponds more closely to today’s large-scale cloud computing frameworks.
The use of redundancy to reduce latency has also attracted attention in other contexts such as cloud storage and networking . Most of these works that consider queueing focus on the case of one task. Waiting for many tasks is harder to analyze as indicated by fork-join queue analysis.
2 Our contributions
In this work we propose a framework to analyze strategies for replicating straggling tasks of a large computing job. In particular we consider three parameters of a straggler replication strategy: 1) the fraction of tasks declared as stragglers, 2) number of replicas for each straggling tasks, and 3) whether the original copy should be killed or kept running. We characterize how these parameters impact the trade-off between latency and computing cost. Our characterizations allow us to identify regimes with the surprising property that replicating a small fraction of tasks drastically reduces latency while saving computing cost. These insights allow one to apply optimization to search for scheduling policies based on one’s sensitivity to computing latency and computing cost.
The rest of the paper is organized as follows. In Section 2 we introduce notation, formulate the problem, and define performance metrics used in the paper. In Section 3 we provide an analysis of single-fork task replication policies and defer all proofs to Appendix A. Then in Section 4 we describe an algorithm that finds a good scheduling policy for execution time distributions that are not analytically tractable (e.g., empirical distributions from real-world traces). In Section 5 we conclude with a discussion of the implications and future perspectives.
Problem Formulation
Lower-case letters (e.g., ) denote a particular value of the corresponding random variable, which is denoted in upper-case letters (e.g., ). We denote the cumulative distribution function (c.d.f.) of by . Its complement, the tail distribution is denoted by . We denote the upper end point of by
For i.i.d. random variables , we define as the -th order statistic, i.e., the -th smallest of the random variables.
2 System Model
We consider a job consisting of parallel tasks, where is largeAnalysis of real-world trace data shows that it is common for a job to contain hundreds or even thousands of tasks . and each task is assigned to a different machine. We use the probability distribution to model the random variation in machine response time due to factors such as congestion, queueing, virtualization, and competing jobs being run on the same machines, and assume this execution time distribution is independent and identically distributed (i.i.d.) across machines. The identical assumption of implies that tasks in this job are assigned to machines with processing power proportional to task size, with the simplest case being a group of homogeneous tasks are assigned to a group of homogeneous machines. The independent assumption of could be satisfied when machine response times fluctuate independently over time, or when each new task (or new replica) is assigned to a new machine that is not previously used to run tasks of the current job. Note that we treat the variability that captures as an exogenous factor from a user’s perspective—in general a user renting machines from a cloud computing service has little or no control over other jobs that share the resources.A system designer may be able to influence this variability by adjusting the resource sharing among different jobs, another interesting direction that is beyond the scope of this work.
3 Scheduling Policy
A scheduling policy or scheduler assigns one or more replicas of each task to different machines, possibly at different time instants. In this work, we assume the scheduler receives instantaneous feedback notifying it when a machine finishes its assigned task, and there is no intermediate feedback indicating the status of processing of a task. We focus our attention on a set of policies called single-fork policies, defined as follows.
A single-fork scheduling policy launches all tasks at time . It waits until tasks finish. For each of the remaining straggling tasks, it chooses one of the following two actions:
When the earliest replica of a task finishes, all the other remaining replicas of the same task are terminated.
Note that in both scenarios there are a total of replicas running after the forking point. Fig. 1 illustrates these two cases of keeping or killing the original copy of a task. For simplicity of notation we assume that is such that is an integer. We note that corresponds to running tasks in parallel and waiting for all to finish, which is the baseline case without any replication or killing any original tasks.
Although we focus on single-fork policies in this paper, the analysis can be generalized to multi-fork policies, where new replicas of straggling tasks are launched at multiple times during the execution of the job [24, Section 6.4]. Forking multiple times can achieve a better latency-cost trade-off, but could be undesirable in practice due to additional delay and complexity in obtaining new and killing existing replicas.
4 Performance Metrics
We now define the latency and cost metrics used to compare straggler replication policies and understand when and how replication is useful.
where is the time when at least one replica of task finishes. More specifically, suppose the scheduler launches replicas of each of the tasks at times for , then
where are i.i.d., drawn from the execution time distribution .
where is given in 3 and .
Fig. 2 illustrates the execution of a job with two tasks, and evaluation of the corresponding latency and cost . Given two tasks, we launch two replicas of task 1 and , and two replicas of task 2 at and . The task execution times are , , , and . Machine finishes the task first at time , and the second replica running on is terminated before it finishes executing. Similarly, machine finishes task at time , and the replica running on is terminated. Thus the latency of the job is . The cost is the sum of all running times normalized by , i.e., .
Single-fork policy analysis
For a computing job with tasks, and task execution time distribution , the latency and cost metrics as are
where is the residual execution time of a straggling tasks after launching replicas. Its tail distribution is given by
Using Theorem 1 we can determine the single-fork policy parameters and that give the best latency-cost trade-off for a given service time distribution . To decide whether to kill or to keep the original copy of the straggling task, we are essentially comparing the additional time needed for the original time to finish and the completion time for a new copy. In Lemma 1 we identify when killing the original task is better than keeping the original task and vice versa.
For a given , killing the original task gives lower latency and cost than keeping it running if
Conversely, if the inequality in (8) is reversed for all , then keeping the original task is better.
The proof is given in Appendix A. For a class of distributions called ‘new-longer-than-used’ distributions , (8) is true for any . An example of such distributions is the shifted-exponential distribution for which we analyze the latency-cost trade-off in Section 3.2 below.
2 Single-fork scheduling with analytical execution time distributions
In this section we evaluate the latency-cost trade-off in Theorem 1 for two execution time distributions: Shifted exponential and Pareto. The shifted exponential distribution has an exponential tail, while Pareto distribution has a heavy tail.
Consider that the task execution time distribution is a shifted exponential distribution . Its tail distribution function is given by
The shifted exponential distribution has an exponentially decaying tail. It is lower bounded by a constant , aiming to capture the delay due to machine start-up or task initialization. Due to this constant , the shifted exponential distribution satisfies (8) for any . Thus, it is always better to keep the original straggling task, and launch additional replicas if necessary.
For a computing job with tasks, if the execution time distribution of tasks are i.i.d. , then as , the latency and cost metrics are
2.2 Pareto execution time
The tail distribution function of the Pareto distribution is
The Pareto distribution has a heavy-tail that decays polynomially. It has been observed to fit task execution time distributions in data centers .
For a computing job with tasks, if the execution time distribution of tasks are i.i.d. , then as , the latency and cost metrics are
Case 2: Keeping the original task The tail distribution of
For a computing job with tasks, if the execution time distribution of each task is , then the expected latency satisfies
In Figures 6(a) and 6(b) we plot the expected latency and cost as varies, for different values of . The black dot is the baseline case (), where no replication is used and we simply wait for the original copies of all tasks to finish. Note that and keeping the original copy is also equivalent to the baseline case, and thus not plotted in the figures. The diminishing return of increasing in terms of latency reduction is clearly demonstrated. In addition, we observe that a small amount of replication (small and ) can reduce latency significantly in comparison with the baseline case. But as increases further, the latency may increase (as observed for ) because of the second term in 5.
Fig. 6(c) shows the latency versus the computing cost for different values of , with varying along each curve. Depending upon the latency requirement and limit on the cost, one can choose an appropriate operating point on this trade-off curve. This plot again demonstrates the non-intuitive phenomenon that it is possible to reduce latency (from to about for and cases) and computing cost simultaneously.
Empirical execution time distributions
In practice, it may be difficult to fit the empirical behavior of the task execution time to a well-characterized distribution, thus making the latency-cost analysis using the framework presented in Section 3 difficult. In this section we propose an algorithm to estimate the latency and cost from the empirical distribution of task execution time. This enables users to evaluate the latency-cost trade-off of various replication strategy using execution trace directly, instead of a fitted execution time distribution. Applying our algorithm to the Google Cluster Trace data , we show that it is possible to improve upon the performance of the default replication policy in MapReduce-style frameworks.
To estimate the latency and cost from empirical execution time samples, we apply the bootstrapping method that uses the empirical distribution as an approximation of the true distribution.
2 Demonstration using Google Cluster Trace
The Google Cluster Trace data gives timestamps of events such as SCHEDULE, EVICT, FINISH, FAIL, KILL etc. for each of the tasks of computing jobs that are run on Google’s cluster machines. In this section we apply Algorithm 1 to two jobs in the Google Cluster Trace, and study the latency-cost trade-offs for these real-world task service distributions.
In our demonstration we only consider tasks with SCHEDULE and FINISH times, as we would like to obtain samples that represent a normal execution (not killed or evicted). In a few rare cases, a task is associated with multiple SCHEDULE and FINISH events due to duplicate execution. For these we choose to keep the first occurrences in each event category.
We choose two jobs (Job ID 6252284914 and 6252315810) with different numbers of tasks. For each task in a job, we obtain the task execution time by calculating the time difference between SCHEDULE and FINISH. The normalized histograms of the task execution times of the two jobs are shown in Fig. 7(a) and Fig. 7(b) respectively. Both the distributions have straggling tasks whose execution time is significantly longer than average. To emphasize the importance of such stragglers, we modify the trace for Job 6252315810 by removing the 3 samples with execution time longer than 1400 seconds, leading to the execution time distribution shown in Fig. 7(c).
For the tail-shortened trace histogram in Fig. 7(c), killing the original copy increases the latency, because it is too “impatient”—the original copy is likely to finish before a new copy of the task. On the other hand, if we keep the original copy, adding a small amount of redundancy can reduce latency and computing cost simultaneously, as shown in Fig. 10(a). Lastly, Fig. 10 indicates that killing and replicating tasks can lead to a worse performance trade-off, so one needs to apply replication with care.
3 Scheduling policy selection
For example, a latency-sensitive user may choose to define the optimal scheduling policy via the following constrained optimization problem:
where is the baseline scheduling policy without replication and the maximum allowed number of copies for a task. On the other hand, a cost-sensitive user may choose to define the optimal scheduling policy via the following optimization problem:
Concluding remarks
Replication of the slowest tasks of a computing job (straggling tasks) has been observed to be highly effective in practice to speed-up job completion. In this paper we provide a theoretical framework to understand the effect of straggler replication on the job completion latency, and the additional computing time spent on running the replicas. Our latency-cost analysis gives the insight that the scaling of job completion latency with the number of tasks depends on the tail of the per-task execution time. We identify regimes where replicating a small fraction of stragglers can drastically reduce latency and computing cost simultaneously. With the guidance from this asymptotic analysis, we propose a bootstrapping-based algorithm to estimate the latency and cost from empirical traces of execution time. The effectiveness of this algorithm is demonstrated on the Google Cluster Trace data, where we show that careful choice of the replication strategy can improve the latency-cost trade-off as compared to the default option in MapReduce.
2 Future Directions
Generalizations of this straggler replication model include considering heterogeneous servers, dependencies between tasks (some tasks need to complete in order to begin others), and taking into account queueing delay of tasks as considered in for the single task case. Another direction is to analyze approximate computing, where we need only a subset of the tasks of a job to complete, a relevant model for information retrieval and machine learning jobs. This idea is developed in the context of coded distributed storage in . We also aim to develop an algorithm that learns the task execution time distribution online, and use it to decide when and how many replicas to launch. This has an exploration-exploitation trade-off, similar to the multi-arm bandit problems studied in reinforcement learning .
More broadly, our analysis framework can be applied to other systems with stochastically varying components, for example, in crowdsourcing, each worker may take a variable amount of time to complete a task .
Appendix A Appendix
Given , if and , where , then for , the order statistic is asymptotically normal,
where is the p.d.f. corresponds to and denotes convergence in probability as .
Extreme value theory (EVT) is an asymptotic theory of extremes, i.e., minima and maxima. It shows that if a distribution belongs to one of three families of distributions Theorem 5), then its maxima can be well characterized asymptotically as given by Theorem 6, which is also referred to as the Fisher-Tippett-Gnedenko Theorem (Theorem 1.1.3 in ).
A distribution function has one of the following domains of attraction if it satisfies the conditions of the extreme value distribution if and only if
where , the upper end point of the distribution .
as and is a non-degenerate distribution. The extreme value distribution and the values of and depend on the domain of attraction (and hence the tail behavior) of given by Theorem 5.
where is called the Gumbel distribution.
where is called the Fréchet distribution.
where is called the reversed-Weibull distribution.
Based on Theorem 6, we can derive the expected value of extreme values, as shown in Lemma 2.
We can also characterize the limit distribution of the sample extreme analogously via Theorem 6 by
It is worth noting that the distribution function for may be in a different domain of attraction from that of .
A.2 Proofs of Single Fork Analysis
As the number of tasks by Theorem 4 we have . Hence,
When we keep the original copy, the residual execution time of a straggling task is
When we kill the original copy, new copies of the straggling task are launched at the forking point. Thus the residual execution time is
Based on Theorem 5, for we have
By Theorem 6 and Theorem 5, the maximum of shifted exponential belongs to the Gumbel family with
Note that the first term does not include because for large the original task would have run for at least seconds. Thus the tail distribution of is given by
By Theorem 6 and Theorem 5 similar to the relaunching case we have
Before showing the detailed proof of Theorem 3], we state in Lemma 3 how the domain of attraction of relates to that of .
Given a single fork policy with ,
The proof follows directly from 7 and Theorem 5, and hence is omitted here.
and then the result holds again following (14).