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 (140140 ms) than the median latency of a single task (11 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., xx) denote a particular value of the corresponding random variable, which is denoted in upper-case letters (e.g., XX). We denote the cumulative distribution function (c.d.f.) of XX by FX(x)F_{X}(x). Its complement, the tail distribution is denoted by FˉX(x)≜1−FX(x)\bar{F}_{X}(x)\triangleq 1-F_{X}(x). We denote the upper end point of FXF_{X} by

For i.i.d. random variables X1,X2,⋯ ,XnX_{1},X_{2},\cdots,X_{n}, we define Xj:nX_{j:n} as the jj-th order statistic, i.e., the jj-th smallest of the nn random variables.

2 System Model

We consider a job consisting of nn parallel tasks, where nn 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 FXF_{X} 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 FXF_{X} 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 FXF_{X} 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 FXF_{X} 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 π(p,r)\pi\left(p,r\right) launches all nn tasks at time . It waits until (1−p)n(1-p)n tasks finish. For each of the remaining pnpn 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 r+1r+1 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 pp is such that pnpn is an integer. We note that p=0p=0 corresponds to running nn 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 TiT_{i} is the time when at least one replica of task ii finishes. More specifically, suppose the scheduler launches rr replicas of each of the nn tasks at times ti,jt_{i,j} for j=0,1,2,…rj=0,1,2,\dots r, then

where Xi,jX_{i,j} are i.i.d., drawn from the execution time distribution FXF_{X}.

where TiT_{i} is given in 3 and (x)+=max⁡(0,x)\left({x}\right)^{+}=\max(0,x).

Fig. 2 illustrates the execution of a job with two tasks, and evaluation of the corresponding latency TT and cost CC. Given two tasks, we launch two replicas of task 1 t1,1=0t_{1,1}=0 and t1,2=2t_{1,2}=2, and two replicas of task 2 at t2,1=0t_{2,1}=0 and t2,2=5t_{2,2}=5. The task execution times are X1,1=8X_{1,1}=8, X1,2=7X_{1,2}=7, X2,1=11X_{2,1}=11, and X2,2=5X_{2,2}=5. Machine M1M_{1} finishes the task first at time t=8t=8, T1=8T_{1}=8 and the second replica running on M2M_{2} is terminated before it finishes executing. Similarly, machine M4M_{4} finishes task 22 at time T2=10T_{2}=10, and the replica running on M3M_{3} is terminated. Thus the latency of the job is T=max⁡{T1,T2}=10T=\max\left\{{T_{1},T_{2}}\right\}=10. The cost is the sum of all running times normalized by nn, i.e., C=(8+6+10+5)/2=14.5C=(8+6+10+5)/2=14.5.

Single-fork policy analysis

For a computing job with nn tasks, and task execution time distribution FXF_{X}, the latency and cost metrics as n→∞n\rightarrow\infty are

where YY is the residual execution time of a straggling tasks after launching replicas. Its tail distribution FˉY\bar{F}_{Y} is given by

Using Theorem 1 we can determine the single-fork policy parameters pp and rr that give the best latency-cost trade-off for a given service time distribution FXF_{X}. 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 0<p≤10<p\leq 1, killing the original task gives lower latency and cost than keeping it running if

Conversely, if the inequality in (8) is reversed for all x≥0x\geq 0, 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 0<p≤10<p\leq 1. 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 FXF_{X} is a shifted exponential distribution ShiftedExp(Δ,μ)\textsf{ShiftedExp}\left({\Delta},{\mu}\right). Its tail distribution function is given by

The shifted exponential distribution has an exponentially decaying tail. It is lower bounded by a constant Δ\Delta, aiming to capture the delay due to machine start-up or task initialization. Due to this constant Δ\Delta, the shifted exponential distribution satisfies (8) for any 0<p≤10<p\leq 1. Thus, it is always better to keep the original straggling task, and launch additional replicas if necessary.

For a computing job with nn tasks, if the execution time distribution of tasks are i.i.d. ShiftedExp(Δ,μ)\textsf{ShiftedExp}\left({\Delta},{\mu}\right), then as n→∞n\rightarrow\infty, the latency and cost metrics are

2.2 Pareto execution time

The tail distribution function of the Pareto distribution Pareto(α,xm)\textsf{Pareto}\left({\alpha},{x_{m}}\right) 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 nn tasks, if the execution time distribution of tasks are i.i.d. Pareto(α,xm)\textsf{Pareto}\left({\alpha},{x_{m}}\right), then as n→∞n\rightarrow\infty, the latency and cost metrics are

Case 2: Keeping the original task The tail distribution of YY

For a computing job with nn tasks, if the execution time distribution of each task is Pareto(α,xm)\textsf{Pareto}\left({\alpha},{x_{m}}\right), then the expected latency satisfies

In Figures 6(a) and 6(b) we plot the expected latency and cost as pp varies, for different values of rr. The black dot is the baseline case (p=0p=0), where no replication is used and we simply wait for the original copies of all nn tasks to finish. Note that r=0r=0 and keeping the original copy is also equivalent to the baseline case, and thus not plotted in the figures. The diminishing return of increasing rr in terms of latency reduction is clearly demonstrated. In addition, we observe that a small amount of replication (small pp and rr) can reduce latency significantly in comparison with the baseline case. But as pp increases further, the latency may increase (as observed for r=0r=0) because of the second term in 5.

Fig. 6(c) shows the latency versus the computing cost for different values of rr, with pp 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 7070 to about 1515 for r=1r=1 and r=2r=2 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 π0\pi_{0} is the baseline scheduling policy without replication and rmax⁡r_{\max} 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 FXF_{X} 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 X1,X2,…,Xn∼ i.i.d.FXX_{1},X_{2},\allowbreak\ldots,\allowbreak X_{n}\stackrel{{\scriptstyle~{}{i.i.d.}}}{{\sim}}F_{X}, if 0<p<10<p<1 and 0<f(xp)<∞0<f(x_{p})<\infty, where xp=FX−1(p)x_{p}=F_{X}^{-1}(p), then for k=np+o(n)k=np+o\left({\sqrt{n}}\right), the kthk^{th} order statistic is asymptotically normal,

where f(⋅)f(\cdot) is the p.d.f. corresponds to FXF_{X} and →P\stackrel{{\scriptstyle P}}{{\rightarrow}} denotes convergence in probability as n→∞n\rightarrow\infty.

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 FXF_{X} has one of the following domains of attraction if it satisfies the conditions of the extreme value distribution G(x)G(x) if and only if

where ω(x)=sup⁡{x:FX(x)<1}\omega\left({x}\right)=\sup\{x:F_{X}(x)<1\}, the upper end point of the distribution FXF_{X}.

as n→∞n\rightarrow\infty and G(⋅)G(\cdot) is a non-degenerate distribution. The extreme value distribution G(x)G(x) and the values of ana_{n} and bnb_{n} depend on the domain of attraction (and hence the tail behavior) of FXF_{X} given by Theorem 5.

where Λ(x)\Lambda(x) is called the Gumbel distribution.

where Φξ(x)\Phi_{\xi}(x) is called the Fréchet distribution.

where Ψξ(x)\Psi_{\xi}(x) 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 X1:n{X}_{{1}:{n}} analogously via Theorem 6 by

It is worth noting that the distribution function for −X-X may be in a different domain of attraction from that of XX.

A.2 Proofs of Single Fork Analysis

As the number of tasks n→∞n\rightarrow\infty by Theorem 4 we have T(1)→FX−1(1−p)T^{(1)}\rightarrow F_{X}^{-1}(1-p). Hence,

When we keep the original copy, the residual execution time of a straggling task is

When we kill the original copy, r+1r+1 new copies of the straggling task are launched at the forking point. Thus the residual execution time is

Based on Theorem 5, for η(y)=1/((r+1)μ)\eta(y)=1/((r+1)\mu) 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 Δ\Delta because for large nn the original task would have run for at least Δ\Delta seconds. Thus the tail distribution of YY 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 FYF_{Y} relates to that of FXF_{X}.

Given a single fork policy π(p,r;n)\pi\left(p,r;n\right) with 0<p<10<p<1,

The proof follows directly from 7 and Theorem 5, and hence is omitted here.

and then the result holds again following (14).

References