Near-Optimal Straggler Mitigation for Distributed Gradient Methods

Songze Li, Seyed Mohammadreza Mousavi Kalan, A. Salman Avestimehr, Mahdi Soltanolkotabi

I Introduction

Gradient descent (GD) serves as a working-horse for modern inferential learning tasks spanning computer vision to recommendation engines. In these learning tasks one is interested in fitting models to a training data set of mm training examples {xj}j=1m\{\bm{x}_{j}\}_{j=1}^{m} (usually consisting of input-output pairs). The fitting problem often consists of finding a mapping that minimizes the empirical risk

In order to scale GD to handle massive amount of training data, developing parallel/distributed implementations of gradient descent over multiple cores or GPUs on a single machine, or multiple machines in computing clusters is of significant importance .

In this paper we consider a distributed computing model consisting of a master node and nn workers as depicted in Fig. 1. Each worker ii stores and processes a subset of rir_{i} training examples locally, and then generates a message zi\bm{z}_{i} based on computing partial gradients using the local training data, and then sends this message to the master node. The master collects the messages from the workers, and uses these messages to compute the total gradient and update the model via (1). If each worker processes a disjoint subset of the examples, the master needs to gather all partial gradients from all the workers. Therefore, when different workers compute and communicate at different speeds, the run-time of each iteration of distributed GD is limited by the slowest worker (or straggler). This phenomenon known as the straggler effect, significantly delays the execution of distributed computing tasks when some workers compute or communicate much slower than others. For example, it was shown in that over a wide range of production jobs, stragglers can prolong the completion time by 34% at median.

We focus on straggler mitigation in the above distributed GD framework. To formulate the problem, we first define two key performance metrics that respectively characterize how much local processing is needed at each worker, and how many workers the master needs to wait for before it can compute the gradient. In particular, we define the computational load, denoted by rr, as the number of training examples each worker processes locally, and the recovery threshold, denoted by KK, as the average number of workers from whom the master collects the results before it can recover the gradient. As a function of the computational load rr, the recovery threshold KK decreases as rr increases. For example, when r=mnr=\frac{m}{n} such that each worker processes a disjoint subset of the examples, KK attains its maximum of nn. One the other hand, if each worker processes all examples, i.e., r=mr=m, the master only needs to wait for one of them to return the result, achieving the minimum K=1K=1. For an arbitrary computational load mn≤r≤m\frac{m}{n}\leq r\leq m, we aim to characterize the minimum recovery threshold across all computing schemes, denoted by K∗(r)K^{*}(r), which provides the maximum robustness to the straggler effect. Moreover, due to the high communication overhead to transfer the results to the master (especially for a high-dimensional model vector w\bm{w}), we are also interested in characterizing the minimum communication load, denoted by L∗(r)L^{*}(r), which is defined as the (normalized) size of the messages received at the master before it can recover the gradient.

To reduce the effect of stragglers in this paper we propose a distributed computing scheme, named “Batched Coupon’s Collector” (BCC). We will show that this scheme achieves the recovery threshold

where HnH_{n} denotes the nnth harmonic number. We also prove a simple lower bound on the minimum recovery threshold demonstrating that

Thus, our proposed BCC scheme achieves the minimum recovery threshold to within a logarithmic factor, that is,

We will also demonstrate that the BCC scheme achieves the minimum communication load to within a logarithmic factor, that is,

The basic idea of the proposed BCC scheme is to obtain the “coverage” of the computed partial gradients at the master. Specifically, we first partition the entire training dataset into mr\frac{m}{r} batches of size rr, and then each worker independently and randomly selects a batch to process. As a result, the process of collecting messages at the master emulates the coupon collecting process in the well-known coupon collector’s problem (see, e.g., ), which requires to collect a total of mr\frac{m}{r} different types of coupons using nn independent trials. Since the examples in different batches are disjoint, we can compress the computed partial gradients at each worker by simply summing them up, and send the summation to the master. Utilizing the algebraic property of the overall computation, the proposed BCC scheme attains the minimum communication load from each worker.

Beyond the theoretical analysis, we also implement the proposed BCC scheme on Amazon EC2 clusters, and empirically demonstrate performance gain over the state-of-the-art straggler mitigation schemes. In particular, we run a baseline uncoded scheme where the training examples are uniformly distributed across the workers without any redundant data placement, the cyclic repetition scheme in designed to combat the stragglers for the worst-case scenario, and the proposed BCC scheme, on clusters consisting of 5050 and 100100 worker nodes respectively. We observe that the BCC scheme speeds up the job execution by up to 85.4% compared with the uncoded scheme, and by up to 69.9% compared with the cyclic repetition scheme.

Finally, we generalize the BCC scheme to accelerate distributed GD in heterogeneous clusters, in which each worker may be assigned different number of training examples according to its computation and communication capabilities. In particular, we derive analytically lower and upper bounds on the minimum job execution time, by developing and analyzing a generalized BCC scheme for heterogeneous clusters. We have also numerically evaluated the performance of the proposed generalized BCC scheme. In particular, compared with a baseline strategy where the dataset is distributed without repetition, and the number of examples a worker processes is proportional to its processing speed, we numerically demonstrate a 29.2829.28% reduction in average computation time.

For the aforementioned distributed GD problem, a simple data placement strategy is that each worker selects rr out of the mm examples uniformly at random. Under this data placement, each worker processes each of the selected examples, and communicates the computed partial gradient individually to the master. Following the arguments of the coupon’s collector problem, this simple randomized computing scheme achieves a recovery threshold

Similar to the proposed BCC scheme, this randomized scheme achieves the minimum recovery threshold to within a logarithmic factor. However, since each worker communicates rr times more messages, the communication load has increased to

Recently a few interesting papers utilize coding theory to mitigate the effect of stragglers in distributed GD. In particular, a cyclic repetition (CR) scheme was proposed in to randomly generate a coding matrix, which specifies the data placement and how to encode the computed partial gradients across workers for communication. Furthermore, in and , the same performance was achieved using deterministic constructions of Reed-Solomon (RS) codes and cyclic MDS (CM) codes. These coding schemes can tolerate r−1r-1 stragglers in the worst case when the computational load is rr. More specifically, when the number of examples is equal to the number of workers (m=nm=n)When m>nm>n, we can partition the dataset into nn groups, and view each group of mn\frac{m}{n} training examples as a “super example”., the above coded schemes achieve the recovery threshold

In all of these coded schemes, each worker encodes the computed partial gradients by generating a linear combination of them, and communicates the single coded message to the master. This yields a communication load

While the above simple randomized scheme and the coding theory-inspired schemes are effective in reducing the recovery threshold and the communication load respectively, the proposed BCC scheme achieves the best of both. In Fig. 2, we numerically compare the recovery threshold of the randomized scheme, the CR scheme in , and the proposed BCC scheme, and demonstrate the performance gain of BCC. To summarize, the proposed BCC schemes has the following advantages

Simplicity: Unlike the computing schemes that rely on delicate code designs for data placement and communication, the BBC scheme is rather simple to implement, and has little coding overhead.

Reliability: The BCC scheme simultaneously achieves near minimal recovery threshold and communication load, enabling good straggler mitigation and fast job execution.

Universality: In contrast to the coding theory-inspired schemes like CR, the proposed BCC scheme does not require any prior knowledge about the number of stragglers in the cluster, which may not be available or vary across the iterations.

Scalability: The data placement in the BCC scheme is performed in a completely decentralized manner. This allows the BCC scheme to seamlessly scale up to larger clusters with minimum overhead for reshuffling the data.

Finally, we highlight some recent developments of utilizing coding theory to speedup a broad calss of distributed computing tasks. In , maximum distance separable (MDS) error-correcting codes were applied to speedup distributed linear algebra operations (e.g., matrix multiplications). In particular, MDS codes were utilized to generate redundant coded computing tasks, providing robustness to missing results from slow workers. The proposed coded computing scheme in was further generalized in , where it was shown that the solution of is a single operating point on a more general tradeoff between computation latency and communication load. Other than dealing with stragglers, coding theory was also shown to be an effective tool to alleviate communication bottlenecks in distributed computing. In , for a general MapReduce framework implemented on a distributed computing cluster, an optimal tradeoff between the local computation on individual workers and the communication between workers was characterized, exploiting coded multicasting opportunities created by carefully designing redundant computations across workers.

II Problem Formulation

We focus on a data-distributed implementation of the gradient descent updates in (1). In particular, as shown in Fig. 1 of Section I, we employ a distributed computing system that consists of a master node, and nn worker nodes (denoted by Worker 11, Worker 2,…,2,\ldots, Worker nn). Worker ii, stores and processes locally a subset of ri≤mr_{i}\leq m training examples. We use Gi⊆{1,…,m}{\cal G}_{i}\subseteq\{1,\ldots,m\} to denote the set of the indices of the examples processed by Worker ii. In the ttth iteration, Worker ii computes a partial gradient gj(wt)\bm{g}_{j}(\bm{w}_{t}) with respect to the current weight vector wt\bm{w}_{t}, for each j∈Gij\in{\cal G}_{i}. Ideally we would like the workers to process as few examples as possible. This leads us to the following definition for characterizing the computational load of distributed GD schemes.

We define the computational load, denoted by rr, as the maximum number of training examples processed by a single worker across the cluster, i.e., r:=max⁡i=1,…,nrir:=\underset{i=1,\ldots,n}{\max}r_{i}.

The assignment of the training examples to the workers, or the data distribution, can be represented by a bipartite graph G{\bf G} that contains a set of data vertices {d1,d2,…,dm}\{d_{1},d_{2},\ldots,d_{m}\}, and a set of worker vertices {k1,k2,…,kn}\{k_{1},k_{2},\ldots,k_{n}\}. There is an edge connecting djd_{j} and kik_{i} if Worker ii computes gj\bm{g}_{j} locally, or in other words, jj belongs to Gi{\cal G}_{i}. Since each data point needs to be processed by some worker, we require that N(k1)∪…∪N(kn)={d1,…,dm}{\cal N}(k_{1})\cup\ldots\cup{\cal N}(k_{n})=\{d_{1},\ldots,d_{m}\}, where N(ki){\cal N}(k_{i}) denotes the neighboring set of kik_{i}. After Worker ii, i=1,…,ni=1,\ldots,n, finishes its local computations, it communicates a function of the local computation results to the master node. More specifically, as shown in Fig. 1 Worker ii communicates to the master a message zi\bm{z}_{i} of the form

Let W⊆{1,…,n}{\cal W}\subseteq\{1,\ldots,n\} denote the index of the subset of workers whose messages are received at the master. After receiving these messages, the master node calculates the complete gradient (based on all training data) by using a decoding function ψ\psi. More specifically,

In order for the master to be able to calculate the complete gradient from the received messages it needs to wait for a sufficient number of workers. We quantify this and a related parameter more precisely below.

We define the communication load, denoted by LL, as the average aggregated size of the messages the master receives from the workers with indices in W{\cal W}, normalized by the size of a partial gradient computed from a single example.

We say that a pair (r,K)(r,K) is achievable if for a computational load rr, there exists a distributed computing scheme, such that the master recovers the gradient after receiving messages from on average KK or less workers.

We define the minimum recovery threshold, denoted by K∗(r)K^{*}(r), as

We also define the minimum communication load, denoted by L∗(r)L^{*}(r), in a similar manner.

In the next section, we propose and analyze a computing scheme for distributed GD over a homogeneous cluster, and show that it simultaneously achieves a near optimal recovery threshold and communication load (up to a logarithmic factor).

III The Batched Coupon’s Collector (BCC) Scheme

In this section, we consider homogeneous workers with identical computation and communication capabilities. As a result, each worker processes the same number of training examples, and we have r1=r2=⋯=rn=rr_{1}=r_{2}=\cdots=r_{n}=r. We note that in this case for the entire dataset to be stored and processed across the cluster, we must have mr≤n\frac{m}{r}\leq n. For this setting, we propose the following scheme which we shall refer to as “batched coupon’s collector” (BCC).

The key idea of the proposed BCC scheme is to obtain the “coverage” of the computed partial gradients at the master. As indicated by the name of the scheme, BCC is composed of two steps: “batching” and “coupon collecting”. In the first step, the training examples are partitioned into batches, which are selected randomly by the workers for local processing. In the second step, the processing results from the data batches are collected at the master, emulating the process of the well-known coupon’s collector problem. Next, we describe in detail the proposed BCC scheme.

Data Distribution. For a given computational load rr, as illustrated in Fig. 3, we first evenly partition the entire data set into ⌈mr⌉\lceil\frac{m}{r}\rceil data batches, and denote the index sets of the examples in these batches by B1,B2,…,B⌈mr⌉{\cal B}_{1},{\cal B}_{2},\ldots,{\cal B}_{\lceil\frac{m}{r}\rceil}. Each of the batches contains rr examples (with the last batch possibly being zero-padded). Each worker node independently picks one of the data batches uniformly at random for local processing. We denote index set of the data points selected by Worker ii as Bσi{\cal B}_{\sigma_{i}}, i.e. Gi=Bσi{\cal G}_{i}={\cal B}_{\sigma_{i}}.

Communication. After computing the partial gradient gj\bm{g}_{j} for all j∈Bσij\in{\cal B}_{\sigma_{i}}, Worker ii computes a single message by summing them up i.e.,

Data Aggregation at the Master. When the master node receives the message from a worker, it discards the message if the master has received the result from processing the same batch before, and keeps the message otherwise. The master keeps collecting messages until the processing results from all data batches are received. Finally, the master reduces the kept messages to the final result by simply computing their summation.

We would like to note that the above BCC scheme is fully decentralized and coordination-free. Each worker selects its data batch independently of the other workers, and performs local computation and communication in a completely asynchronous manner. There is no need for any feedback from the master to the workers or between the workers. All these features make this scheme very simple to implement in practical scenarios.

III-B Near Optimal Performance Guarantees for BCC

In this subsection, we theoretically analyze the BCC scheme, whose performance provides an upper bound on the minimum recovery threshold of the distributed GD problem, as well as an upper bound on the minimum communication load. To start, we state the main results of this paper in the following theorem, which characterizes the minimum recovery threshold and the minimum communication load to within a logarithmic factor.

For a distributed gradient descent problem of training mm data examples distributedly over nn worker nodes, we have

for sufficiently large nn, where K∗(r)K^{*}(r) and L∗(r)L^{*}(r) are the minimum recovery threshold and the minimum communication load respectively, KBCC(r)K_{\textup{BCC}}(r) and LBCC(r)L_{\textup{BCC}}(r) are the recovery threshold and the communication load achieved by the BCC scheme, and Ht=∑k=1t1kH_{t}=\sum_{k=1}^{t}\frac{1}{k} is the tt-th harmonic number.

Given that H⌈mr⌉≈⌈mr⌉log⁡(⌈mr⌉)H_{\lceil\tfrac{m}{r}\rceil}\approx\lceil\tfrac{m}{r}\rceil\log(\lceil\tfrac{m}{r}\rceil), the results of Theorem 1 imply that for the homogeneous setting, the proposed BCC scheme simultaneously achieves a near minimal recovery threshold and communication load ( up to a logarithmic factor). \hfill□\hfill\square

As we mentioned before, other coding-based approaches mostly focus on the worst-case scenario, resulting in a high recovery threshold e.g. KCR=m−r+1K_{\textup{CR}}=m-r+1.This is assuming m=nm=n. We would like to point out that although designed for the worst-case, the fractional scheme proposed in may finish when the master collects results from less than m−r+1m-r+1 workers. However, it only applies to the case where rr divides mm. In contrast, instead of focusing on worst-case scenarios, our proposed scheme aims at achieving the “coverage” of the partial computation results at the master, by collecting the computation of a much smaller number of workers (on average). As numerically demonstrated in Fig. 2 in Section I, The BCC scheme brings down the recovery threshold from m−r+1m-r+1 to roughly mrlog⁡mr\frac{m}{r}\log\frac{m}{r}. \hfill□\hfill\square

In the coded computing schemes proposed in , a linear combination of the locally computed partial gradients is carefully designed at each worker, such that the final gradient can be recovered at the master with minimum message sizes communicated by the workers. In the BCC scheme, each worker also communicates a message of minimum size, which is created by summing up the local partial gradients. As a result, BCC achieves a much smaller recovery threshold and hence can substantially reduces the total amount of network traffic. This is especially true when the dimension of the gradient is large, leading to significant speed-ups in the job execution. \hfill□\hfill\square

The coded schemes in are designed to make the system robust to a fixed number of stragglers. Specifically, for a cluster with ss stragglers, a code can be designed such that the master can proceed after receiving m−sm-s messages, no matter which ss workers are slow. However, it is often difficult to predict the number of stragglers in a cluster, and it can change across iterations of the GD algorithm, which makes the optimal selection of this parameter for the coding schemes in practically challenging. In contrast, our proposed BCC scheme is universal, i.e., it does not require any prior knowledge about the stragglers in the cluster, and still promises a near-optimal straggler mitigation. \hfill□\hfill\square

The lower bound mr\frac{m}{r} in (13) and (14) is straightforward. They correspond to the best-case scenario where all workers the master hears from before completing the task, have mutually disjoint training examples. The upper bound in (13) and (14) is simultaneously achieved by the above described BCC scheme. To see this, we view the process of collecting messages at the master node as the classic coupon collector’s problem (see e.g., ), in which given a collection of NN types of coupons, we need to draw uniformly at random, one coupon at a time with replacement, until we collect all types of coupons. In this case, we have ⌈mr⌉\lceil\frac{m}{r}\rceil batches of training examples, from which each worker independently selects one uniformly at random to process. It is clear that the process of collecting messages at the master is equivalent to collecting coupons of N=⌈mr⌉N=\lceil\frac{m}{r}\rceil types. As we know that the expected numbers of draws to collect all NN different types of coupons is NHNNH_{N}, we use N=⌈mr⌉N=\lceil\frac{m}{r}\rceil and reach the upper bound on the minimum recovery threshold. To characterize the communication load of the BCC scheme, we first note that since each worker communicates the summation of its computed partial gradients, the message size from each worker is the same as the size of the gradient computed from a single example. As a result, a communication load of 11 is accumulated from each surviving worker, and the BCC scheme achieves a communication load that is the same as the achieved recovery threshold.∎

Beyond the theoretical analysis, we also implement the proposed BCC scheme for distributed GD over Amazon EC2 clusters. In the next section, we describe the implementation details, and compare its empirical performance with two baseline schemes.

III-C Empirical Evaluations of BCC

In this subsection, we present the results of experiments performed over Amazon EC2 clusters. In particular, we compare the performance of our proposed BCC scheme, with the following two schemes.

uncoded scheme: In this case, there is no repetition in data among the workers and the master has to wait for all the workers to finish their computations.

cyclic repetition scheme of : In this case, each worker processes rr training examples and in every iteration, the master waits for the fastest m−r+1m-r+1 workers to finish their computations.

We train a logistic regression model using Nesterov’s accelerated gradient method. We compare the performance of the BCC, the uncoded and the cyclic repetition schemes on this task. We use Python as our programming language and MPI4py for message passing across EC2 instances. In our implementation, we load the assigned training examples onto the workers before the algorithms start. We measure the total running time via Time.time(), by subtracting the starting time of the iterations from the completion time at the master. In the ttth iteration, the master communicates the latest model wt\bm{w}_{t} to all the workers using Isend(), and each worker receives the updated model using Irecv(). In the cyclic repetition scheme, each worker sends the master a linear combination of the computed partial gradients, whose coefficients are specified by the coding scheme in . In the BCC and uncoded schemes the workers simply send the summation of the partial gradients back to the master. When the master receives enough messages from the workers, it computes the overall gradient and updates the model.

We run Nesterov’s accelerated gradient descent distributedly for 100 iterations, using the aforementioned three schemes. We compare their performance in the following two scenarios:

scenario one: We use 5151 t2.micro instances, with one master and n=50n=50 workers. We have m=50m=50 data batches, each of which contains 100100 data points generated according to the aforementioned model.

scenario two: We use 101101 t2.micro instances, with one master and n=100n=100 workers. We have m=100m=100 data batches, each of which contains 100100 data points.

III-C2 Results

For the uncoded scheme, each worker processes r=mnr=\frac{m}{n} data batches. For the cyclic repetition and the BCC schemes, we select the computational load rr based on the memory constraints of the instances so as to minimize the total running times.

We plot the total running times of the three schemes in both scenarios in Fig. 4. We also list the breakdowns of the running times for scenario one in Table I and scenario two in Table II respectively. Within each iteration, we measure the computation time as the maximum computation time among the workers whose results are received by the master before the iteration ends. After the last iteration, we add the computation times in all iterations to reach the total computation time. The communication time is computed as the difference between the total running time and the computation time.Due to the asynchronous nature of the distributed GD, we cannot exactly characterize the time spent on computation and communication (e..g., often both are happening at the same time). The numbers listed in Tables I and II provide approximations of the time breakdowns.

We draw the following conclusions from these results.

As we observe in Fig. 4, in scenario one, the BCC scheme speeds up the job execution by 85.4% over the uncoded scheme, and 69.9% over the cyclic repetition scheme. In scenario two, the BCC scheme speeds up the job execution by 73.0% over the uncoded scheme, and 69.7% over the cyclic repetition scheme. In scenario one, we observe the master waiting for on average 1111 workers to finish their computations, compared with 4141 workers for the cyclic repetition scheme and all 5050 workers for the uncoded scheme. In scenario two, we observe the master waiting for on average 2525 workers to finish their computations, compared with 9191 workers for the cyclic repetition scheme and all 100100 workers for the uncoded scheme.

As we note in Fig. 4, the performance gains of both cyclic repetition and BCC schemes over the uncoded scheme become smaller with increasing number of workers. This is because that as the number of workers increases, in order to optimize the total running time, we need to also increase the computational load rr at each worker to maintain a low recovery threshold. However, due to the memory constraints at the worker instances, we cannot increase rr beyond the value 1010 to fully optimize the run-time performance.

We observe from Table I and Table II that having a smaller recovery threshold benefits both the computation time and the communication time. While the BCC scheme and the cyclic repetition scheme have the same computational load at each worker, the computation time of BCC is much shorter since it needs to wait for a smaller number of workers to finish. On the other hand, lower recovery threshold of BCC yields a lower communication load that is directly proportional to the communication time. As a result, since in all experiments the communication time dominates the computation time, the total running time of each scheme is approximately proportional to its recovery threshold.

IV Extension to Heterogeneous Clusters

For distributed GD in heterogeneous clusters, workers have different computational and communication capabilities. In this case, the above proposed BCC scheme is in general sub-optimal due to its oblivion of network heterogeneity. In this section, we extend the above BCC scheme to tackle distributed DC over heterogeneous clusters. We also theoretically demonstrate that the extended BCC scheme provides an approximate characterization of the minimum job execution time.

In the heterogeneous setting, we consider an uncoded communication scheme where after processing the local training examples, each worker communicates each of its locally computed partial gradients separately to the master. That is, Worker ii, i=1,…,ni=1,\ldots,n, communicates zi={gj:j∈Gi}\bm{z}_{i}=\{\bm{g}_{j}:j\in{\cal G}_{i}\} to the master. Under this communication scheme, the master computes the final gradient as soon as it collects the partial gradients computed from all examples. When this occurs, we say that coverage is achieved at the master node.

We assume that the time required for Workers to process the local examples and deliver the partial gradients are independent from each other. We assume that this time interval, denoted by TiT_{i} for Worker ii, is a random variable with a shift-exponential distribution, i.e.,

Here, μi≥0\mu_{i}\geq 0 and ai≥0a_{i}\geq 0 are the fixed straggler and shift parameters of Worker ii.

In this case, the total job execution time, or the time to achieve coverage at the master is given by

We are interested in characterizing the minimum average execution time in a heterogeneous cluster, which can be formulated as the following optimization problem.

In the rest of this section, we develop lower and upper bounds on the optimal value of P1{\cal P}_{1}.

To start, we first define the waiting time for the master to receive at least ss partial gradients (possibly with repetitions)

We also consider the following optimization problem

It is intuitive that once we fix the work loads at the worker, i.e., (r1,r2,…,rn)(r_{1},r_{2},\ldots,r_{n}), the time for the master to receive ss results T^s\hat{T}_{s} should increase as ss increases. We formally state this phenomenon in the following lemma.

Consider an arbitrary dataset placement G{\bf G} where Worker ii processes ∣Gi∣=ri|{\cal G}_{i}|=r_{i} training examples, for any 0≤s1,s2≤∑i=1nri0\leq s_{1},s_{2}\leq\sum_{i=1}^{n}r_{i}, such that s1≤s2s_{1}\leq s_{2}, we have

To tackle the distributed GD problem over heterogeneous cluster, we generalize the above BCC scheme, and characterize the completion time of the generalized scheme using the optimal value of the above problem P2{\cal P}_{2}. The characterized completion time serves as an upper bound on the minimum average coverage time. Next, we state this result in the following theorem.

For a distributed gradient descent problem of training mm data examples distributedly over nn heterogeneous worker nodes, where the computation and communication time at Worker ii has an exponential tail with a straggler parameter μi\mu_{i} and a shift parameter aia_{i}, the minimum average time to achieve coverage is bounded as

where c=2+log⁡(a+Hn/μ)log⁡mc=2+\frac{\log(a+H_{n}/\mu)}{\log m}, a=max⁡(a1,…,an)a=\max(a_{1},\ldots,a_{n}), μ=min⁡(μ1,…,μn)\mu=\min(\mu_{1},\ldots,\mu_{n}).

The proof of Theorem 2 is deferred to the appendix.

The upper bound on the average coverage time is achieved by a generalized BCC scheme, in which given the optimal data assignments (r1∗,…,rn∗)(r_{1}^{*},\ldots,r_{n}^{*}) for P2{\cal P}_{2} with s ⁣= ⁣⌊cmlog⁡m⌋s\!=\!\lfloor cm\log m\rfloor, Worker ii independently selects ri∗r_{i}^{*} examples uniformly at random. We emphasize that similar to the BCC data distribution policy in the homogeneous setting, the main advantages of the generalized BCC lies in its simplicity and decentralized nature. That is, each node selects the training examples randomly and independently from the other nodes, and we do not need to enforce a global plan for the data distribution. This also provides a scalable design so that when a new worker is added to the cluster, according to the updated dataset assignments computed from P2{\cal P}_{2} with n+1n+1 workers and s=⌊cmlog⁡m⌋s=\lfloor cm\log m\rfloor, each worker can individually adjust its workload by randomly adding or dropping some training examples, without needing to coordinate with the master or other workers. \hfill□\hfill\square

IV-C Numerical Results

We numerically evaluate the performance of the generalized BCC scheme in heterogeneous clusters, using the proposed random data assignment. In this case, we compute the optimal assignment (r1∗,…,rn∗)(r_{1}^{*},\ldots,r_{n}^{*}) to minimize the average time for the master to collect ⌊mlog⁡m⌋\lfloor m\log m\rfloor partial gradients. In comparison, we also consider a “load balancing” (LB) assignment strategy where the mm data points are distributed across the cluster based on workers’ processing speeds, i.e., ri=μi∑μimr_{i}=\frac{\mu_{i}}{\sum\mu_{i}}m.

We consider the computation task of processing m=500m=500 examples over a heterogeneous cluster of n=100n=100 workers. All workers have the same shift parameter ai=20a_{i}=20, for all i=1,…,ni=1,\ldots,n. The straggling parameter μi=1\mu_{i}=1 for 9595 workers, and μi=20\mu_{i}=20 for the remaining 55 workers. As shown in Fig. 5, the computation of the LB assignment is long since the master needs to wait for every worker to finish. However, utilizing the proposed random assignment, the master can terminate the computation once it has achieved coverage, which significantly alleviates the straggler effect. As a result, the generalized BCC scheme reduces the average computation time by 29.28%29.28\% compared with the LB scheme.

V Conclusion

We propose a distributed computing scheme, named batched coupon’s collector (BCC), which effectively mitigates the straggler effect in distributed gradient descent algorithms. We theoretically illustrate that the BCC scheme is robust to the maximum number of stragglers to within a logarithmic factor. We also empirically demonstrate the performance gain of BCC over baseline straggler mitigation strategies on EC2 clusters. Finally, we generalize the BCC scheme to minimize the job execution time over heterogeneous clusters.

References

Appendix Proof of Theorem 2

Before starting the formal proof of Theorem 2, we first state a result for the coupon collector’s problem that will become useful later. We denote the random variable that represents the minimum number of coupons one needs to collect before obtaining all mm types of coupons as M(M≥m)M(M\geq m), and present an upper bound on the tail probability in the following lemma.

Pr(M≥(1+ϵ)mlog⁡m)≤1mϵ\textup{Pr}(M\geq(1+\epsilon)m\log m)\leq\frac{1}{m^{\epsilon}}, for any ϵ≥0\epsilon\geq 0.

We prove Theorem 2 in two steps. In the first step, we propose a generalized BCC scheme, for which no batching operation is performed on the dataset, and the workers simply sample the examples to process uniformly at random. In the second step, we analyze the average execution time of the generalized BCC scheme. To start, we obtain an estimate of the number of partial gradients the master receives before coverage is achieved (analogous to the recovery threshold in the homogeneous setting). Then, conditioned on the value of this number, we derive an upper bound on the average coverage time, which is obviously also an upper bound on the minimum coverage time over all schemes.

where cc is specified in the statement of Theorem 2. Assume the optimal task assignment is given by

First, we consider a relaxed data distribution strategy G1{\bf G}_{1} in which Worker ii, independently, and uniformly at random selects ri∗r_{i}^{*} data points with replacement, and processes them locally. That is, G1{\bf G}_{1} allows each worker to process an example more than once. It is obvious that

We note that when using the data distribution G1{\bf G}_{1} the process of receiving partial gradients at the master mimics the process of collecting coupons in the coupon collector’s problem. We define a random variable W(W≥m)W(W\geq m) as the minimum number of partial gradients (possibly with repetition) the master receives before it reaches coverage. We note that WW is statistically equivalent to the minimum number of coupons one needs to collect in the coupon collector’s problem. In what follows, we only consider the case where the coverage can be achieving using the messages sent by all nn nodes (or the computation can be successfully executed), i.e., W≤∑i=1nrn∗W\leq\sum_{i=1}^{n}r_{n}^{*}.

Taking expectation conditioned on the value of WW, we have

where step (a) is due to Lemma 2, and step (b) results from Lemma 1 in Section IV. In step (c), Tˉ1,Tˉ2,…,Tˉn\bar{T}_{1},\bar{T}_{2},\ldots,\bar{T}_{n} are i.i.d. random variables with the shift-exponential distribution

for all i=1,…,ni=1,\ldots,n, where μ=min⁡(μ1,…,μn)\mu=\min(\mu_{1},\ldots,\mu_{n}), a=max⁡(a1,…,an)a=\max(a_{1},\ldots,a_{n}), and r∗=max⁡(r1∗,…,rn∗)r^{*}=\max(r^{*}_{1},\ldots,r^{*}_{n}). Step (d) is because that we choose c=2+log⁡(a+Hn/μ)log⁡mc=2+\frac{\log(a+H_{n}/\mu)}{\log m}.