Improving Distributed Gradient Descent Using Reed-Solomon Codes
Wael Halbawi, Navid Azizan-Ruhi, Fariborz Salehi, Babak Hassibi
Introduction
With the size of today’s datasets, due to high computation and/or memory requirements, it is virtually impossible to run large-scale learning tasks on a single machine; and even if that is possible, the learning process can be extremely slow due to its sequential nature. Therefore, it is highly desirable or, even necessary, to run the tasks in a distributed fashion on multiple machines/cores. For this reason, parallel and distributed computing has attracted a lot of attention in recent years from the machine learning, and other, communities .
When a task is divided among a number of machines, the “computation time” is clearly reduced significantly, since the task is being processed in parallel rather than sequentially. However, the taskmaster has to wait for all the machines in order to be able to recover the exact desired computation. Therefore, in the face of substantial or heterogeneous delays, distributed computing may suffer from being slow, which defeats the purpose of the exercise. Several approaches have been proposed to tackle this problem. One naive yet common way, especially when the task consists of many iterations, is to not wait for all machines, and ignore the straggling machines. One may hope that in this way on average the taskmaster receives enough information from everyone; however, it is clear that the performance of the learning algorithm may be significantly impacted in many cases because of lost updates. An alternative and more appropriate way to resolve this issue, is to introduce some redundancy in the computation of the machines, in order to efficiently trade off computation time for less wait time, and to be able to recover the correct update using only a few machines. But the great challenge here is to design a clever scheme for distributing the task among the machines, such that the computation can be recovered using a few machines, independent of which machines they are.
Over the past few decades, coding theory has been developed to address similar challenges in other domains, and has had enormous success in many applications such as mobile communication, storage, data transmission, and broadcast systems. Despite the existence of a great set of tools developed in coding theory which can be used in many machine learning problems, researchers had not looked at this area until very recently . This work is aimed at bridging the gap between distributed machine learning and coding theory, by introducing a carefully designed coding scheme for efficiently distributing a learning task among a number of machines.
More specifically, we consider gradient-based methods for additively separable cost functions, which are the most common way of training any model in machine learning, and use coding to cleverly distribute each gradient iteration across machines in an efficient way. To that end, we propose a deterministic construction based on Reed-Solomon codes accompanied with an efficient decoder, which is used to recover the full gradient update from a fixed number of returning machines. Furthermore, we provide a new delay model based on heavy-tail distributions that also incorporates the time required for decoding. We analyze this model theoretically and use it to optimally pick our scheme’s parameters. We compare the performance of our method on the MNIST dataset with other approaches, namely: 1) Ignoring the straggling machines , 2) Waiting for all the machines, and 3) GradientCoding as proposed by Tandon et al. . Our numerical results show that, for the same training time, our scheme achieves better test errors.
As mentioned earlier, coding theory in machine learning is a relatively new area. We summarize the recent related work here. Lee et al. recently employed a coding-theoretic method in two specific distributed tasks, namely matrix multiplication and data shuffling. They showed significant speed-ups are possible in those two tasks by using coding. Dutta et al. proposed a method that speeds up distributed matrix multiplication by sparsifying the inner products computed at each machine. A coded MapReduce framework was introduced by Li et al in which is used to facilitate data shuffling in distributed computing. The closest work to our framework is the work of Tandon et al. , which aims at mitigating the effect of stragglers in distributed gradient descent using Maximum-Distance Separable (MDS) codes. However, no analysis of computation time was provided. Furthermore, in their framework, along with the above-mentioned works, the decoding was assumed to be performed offline which might be impractical in certain settings.
2 Statement of Contributions
In this work, we make the following three main contributions.
We construct a deterministic coding scheme for efficiently distributing gradient descent over a given number of machines. Our scheme is optimal in the sense that it can recover the gradient from the smallest possible number of returning machines, , given a prespecified computational effort per machine.
We provide an efficient online decoder, with time complexity for recovering the gradient from any machines, which is faster than the best known method , .
We analyze the total computation time, and provide a method for finding the optimal coding parameters. We consider heavy-tailed delays, which have been widely observed in CPU job runtimes in practice .
The rest of the paper is organized as follow. In Section 2, we describe the problem setup and explain the design objectives in detail. Section 3, provides the construction of our coding scheme, using the idea of balanced Reed-Solomon codes. Our efficient online decoder is presented in Section 4. We then characterize the total computation time, and describe the optimal choice of coding parameters, in Section 5. Finally, we provide our numerical results in Section 6, and conclude in Section 7.
Preliminaries
As it will be explained in detail, and are to be chosen in such a way that the total computation time is minimized. For a fixed and , we want to be able to recover the gradient using the linear combinations received from the fastest machines at the master (or equivalently tolerate stragglers). Note that we do not assume any prior knowledge about the stragglers, i.e., we shall design a scheme that enables master to recover the gradient from any set of machines. It is known that for any fixed and , an upper-bound on the number of stragglers that any scheme can tolerate is:
The scheme proposed in this work achieves this bound. A coding scheme designed to tolerate stragglers consists of an encoding matrix , and a collection of decoding vectors . The matrix should satisfy:
Each row of contains exactly nonzero entries.
The linear space generated by any rows of contains the all-one vector of length , .
The values of these nonzero entries prescribe the linear combination sent by . In other words, the coded partial gradient sent from to is given by
When this holds for any set of indices of size , it means that the gradient can be recovered from the set of machines that return fastest. ,
2 Computational Trade-offs
In a distributed scheme that does not employ redundancy, the taskmaster has to wait for all the workers to finish in order to compute the full gradient. However, in the scheme outlined above, the taskmaster needs to wait for the fastest machines to recover the full gradient. Clearly, this requires more computation by each machine. Note that in the uncoded setting, the amount of computation that each worker does is of the total work, whereas in the coded setting each machine performs a fraction of the total work. From (2), we know that if a scheme can tolerate stragglers, the fraction of computation that each worker does is . Therefore, the computation load of each worker increases by a factor of . As will be explained further in Section 5, there is a sweet spot for (and consequently ) that minimizes the expected total time that the master waits in order to recover the full gradient update.
It is worth noting that it is often assumed that the decoding vectors are precomputed for all possible combinations of returning machines, and the decoding cost is not taken into account in the total computation time. In a practical system, however, it is not very reasonable to compute and store all the decoding vectors, especially as there are such vectors, which grows quickly with . In this work, we introduce an online algorithm for computing the decoding vectors on the fly, for the indices of the workers that respond first. The approach is based on the idea of inverting Vandermonde matrices, which can be done very efficiently. In the sequel, we show how to construct an encoding matrix for any and , such that the system is resilient to stragglers, along with an efficient algorithm for computing the decoding vectors .
Code Construction
The basic building block of our encoding scheme is a matrix , where each row is of weight , which serves as a mask for the matrix , where is the number of data partitions that is assigned to every machine. Each column of will be chosen as a codeword from a suitable Reed–Solomon Code over the complex field, with support dictated by the corresponding column in . Whereas the authors of choose the rows of as codewords from a suitable MDS code, this approach does not immediately work when is not equal to .
We will utilize techniques from to construct the matrix (and then ). For that, we present the following definition.
A matrix is column (row)-balanced if for fixed row (column) weight, the weights of any two columns (rows) differ by at most 1.
While our scheme can handle the case where is not an integer, in this section we illustrate the case where it is. The general result is described in the appendix.
Ultimately, we are interested in a matrix with row weight that prescribes a mask for the encoding matrix . As an example, let , and . Then, is given by
where each column is of weight . The following algorithm produces a balanced mask matrix. For a fixed column weight , each row has weight either or .
2 Reed–Solomon Codes
It is well-known that any rows of form an invertible matrix, which implies that specifying any evaluations of a polynomial of degree at most characterizes it. In particular, fixing evaluations of the polynomial to zero characterizes uniquely up to scaling. This property will give us the ability to construct from .
3 Building the Encoding Matrix from the Mask Matrix
Once a mask matrix has been determined using Algorithm 4, the encoding matrix can be built by picking appropriate codewords from . Consider in (5) and the following polynomials
The constant is chosen such that the constant term of , i.e. , is equal to . The evaluations of on are collected in the vector which sits as the column of . The validity of this process can be confirmed using (6), and is generalized in Algorithm 2.
Once the matrix is specified, the corresponding decoding vectors required for computing the gradient at the taskmaster have to be characterized.
Efficient Online Decoding
We exploit the fact that is constructed using Reed–Solomon codewords and show that each decoding vector can be computed in time. Recall that the taskmaster should be able to compute the gradient from any surviving machines, indexed by , according to (4). The column of is determined by a polynomial where . We can write as , where and is the vector of coefficients of , and is the matrix given in (6). Now consider , the coded partial gradients received from . The rows of corresponding to are given by
We require a vector such that . This is equivalent to finding a vector such that
Indeed, the matrix in the above product is a Vandermonde matrix defined by distinct elements and so it is invertible in time , which facilitates the online computation of the decoding vectors. This is an improvement compared to previous works where the decoding time is usually . A careful inspection of inverses of Vandermonde matrices built from an root of unity allows us to compute the required decoding vector in a space efficient manner. This is demonstrated in the next subsection.
Note that is nothing but the first row of the inverse of , which can be built from a set of polynomials . Let the column of be and associate it with . The condition , where is the elementary basis vector of length , implies that should vanish on . Specifically,
The first row of is given by , where is the constant term of . Indeed, we have , which can be computed in closed form according to the following formula,
By choosing as a primitive root of unity, one is guaranteed that there are only distinct values of . This observation proposes that the master should precompute and store the set , and then compute each by utilizing lookup operations. The following algorithm outlines this procedure.
Analysis of Total Computation time
In this section, we provide a theoretical model which can be used to optimize the choice of parameters that define the encoding scheme. For this purpose, we model the response time of a single computing machine as
where the quantity can be thought of the fundamental delay of the machine, i.e. the minimum time required for a machine to return in perfect conditions. Previous works model the return time of a machine as a shifted exponential random variable. We propose using this approach since the heavy-tailed nature of CPU job runtime has been observed in practice .
Let denote the expected time of computing the gradient using the first machines. As a result we have
where is the ordered statistic of , and is the time required at the taskmaster for decoding. Here we assume is large and define as the fraction of the dataset assigned to each machine. For this value of , the number of machines required for successful recovery of the gradient is given by
The expected value of the order statistic of the Pareto distribution with parameter will converge as grows, i.e.,
Using this result, we can approximate , for ,
where we assume that the taskmaster uses Algorithm 3 for decoding. If we assume is the time required for one FLOP, the total decoding time is given by . Since is bounded from above by the memory of each machine, one can find the optimal computation time, subject to memory constraints, by minimizing with respect to .
In the schemes where the decoding vectors are computed offline, the quantity does not appear in the total computation time . Therefore, for large values of , we can write:
This function can be minimized with respect to by standard calculus to give
Note that this quantity is valid (less than one) if and only if one has . It has been observed in practice that the parameter is close to one. Therefore, this assumption holds because is assumed to be large.
Numerical Results
To demonstrate the effectiveness of the scheme, we performed numerical simulations using MATLAB. We train a softmax regression model on a distributed cluster composed of machines to classify 10000 handwritten digits from the MNIST database while artificially introducing delay as a random variable sampled from a Pareto distribution according to (16) with parameters and . Similar to , knowledge of the entire gradient allows us to employ accelerated gradient methods such as the one proposed by Nesterov . Details of the experiment are given in the accompanying description of Figure 2. We compare several schemes by running each of them on the same dataset for a fixed amount of time (in seconds) and then measuring the test error. The results depicted in Figure 2 demonstrate that the scheme proposed in this paper outperforms the four other schemes.
Conclusion
We presented a straggler mitigation scheme that facilitates the implementation of distributed gradient descent in a computing cluster. For a fixed per-machine computational effort, the taskmaster recovers the full gradient from the least number of machines theoretically required, which is done via an algorithm that is efficient in both space and time. Furthermore, we propose a theoretical delay model based on heavy-tailed distributions and incorporates the decoding time, which allows us to minimize the expected running time of the algorithm.
References
Appendix
To lighten notation, we prove correctness for . The general case follows immediately.
Let and be integers where . The row weights of matrix produced by Algorithm 1 for are
Proof. The nonzero entries in column of are given by
In case , each element in , after reducing modulo , appears the same number of times. As a result, those indices correspond to columns of equal weight, namely . Hence, the two cases of are identical along with their corresponding index sets.
In the case where , each of the first elements, after reducing modulo , appears the same number of times. As a result, the nonzero entries corresponding to those indices are distributed evenly amongst the rows, each of which is of weight . The remaining indices contribute an additional nonzero entry to their respective rows, those indexed by . Finally, we have that the first rows are of weight , while the remaining ones are of weight .
Now consider the case when is not necessarily equal to zero. This amounts to shifting (cyclically) the entries in each column by positions downwards. As a result, the rows themselves are shifted by the same amount allowing to conclude us the following.
Let and be integers where . The row weights of matrix produced by Algorithm 1 are
2 General construction
The matrices and are constructed using Algorithm 1. Each column of has weight and each column of has weight . Note that according to (2), we require in order to tolerate a positive number of stragglers.
3 Correctness of Algorithm 4
According to the algorithm, the condition implies that leading to , which is constructed using Algorithm 1.
Moving on to the general case, the matrix given by
where each matrix is row-balanced. The particular choice of in aligns the “heavy” rows of with the “light" rows of , and vice-versa. The algorithm works because the choice of parameters equates the number of heavy rows of to the number of light rows of . The following lemma is useful in two ways.
.
and conclude that .
We have shown the concatenation of a “heavy” row of along with a “light” row of results in one that is of weight . It remains to show that the concatenation of and results of rows of this type only.
We will assume that holds. From Proposition 2, we have and . We will show that the two quantities are in fact equal. Indeed, we can express as
Hence and by the choice of , the “light" rows of align with the “heavy" rows of , and vice-versa. Furthermore, Lemma 1 guarantees that each row of is of weight . The same holds for the remaining rows, using the fact that when both and are non-integers.
4 Proof of Proposition 1
From , the expected value of the ordered statistic of the Pareto distribution is:
where is the gamma function given by . We now assume that is large and make the standard approximation
Furthermore, (2) implies that the number of machines we wait for is , for some which leads to
By letting , the first two terms in the product converge to and , respectively, which yields
Offline Decoding
For illustrative purposes, we plot the function from (21) for a given set of parameters and indicate the optimal point. This plot is given in Figure 3.