Gradient Coding from Cyclic MDS Codes and Expander Graphs
Netanel Raviv, Itzhak Tamo, Rashish Tandon, Alexandros G. Dimakis
I Introduction
Data intensive machine learning tasks have become ubiquitous in many real-world applications, and with the increasing size of training data, distributed methods have gained increasing popularity. However, the performance of distributed methods (in synchronous settings) is strongly dictated by stragglers, i.e., nodes that are slow to respond or unavailable. In this paper, we focus on coding theoretic (and graph theoretic) techniques for mitigating stragglers in distributed synchronous gradient descent.
A coding theoretic framework for straggler mitigation called gradient coding was first introduced in . It consists of a system with one master and worker nodes (or servers), in which the data is partitioned into parts, and one or more parts is assigned to each one of the workers. In turn, each worker computes the partial gradient on each of its assigned parts, linearly combines the results according to some predetermined vector of coefficients, and sends this linear combination back to the master node. Choosing the coefficients at each node judiciously, one can guarantee that the master node is capable of reconstructing the full gradient even if any machines fail to perform their work. The storage overhead of the system, which is denoted by , refers to the amount of redundant computations, or alternatively, to the number of data parts that are sent to each node (see example in Fig. 1).
The importance of straggler mitigation was demonstrated in a series of recent studies (e.g., and ). In particular, it was demonstrated in that stragglers may run up to slower than the typical worker ( in ) on Amazon EC2, especially for the cheaper virtual machines; such erratic behavior is unpredictable and can significantly delay training. One can, of course, use more expensive instances but the goal here is to use coding theoretic methods to provide reliability out of cheap unreliable workers, overall reducing the cost of training.
By and large, the purpose of gradient coding is to enable the master node to compute the exact gradient out of the responses of any non-straggling nodes. The work of established the fundamental bound , provided a deterministic construction which achieves it with equality when , and a randomized one which applies to all and . Subsequently, deterministic constructions were also obtained by and . These works have focused on the scenario where is known prior to the construction of the system. Furthermore, the exact computation of the full gradient is guaranteed if the number of stragglers is at most , but no error bound is guaranteed if this number exceeds .
The contribution of this work is twofold. For the computation of the exact gradient we employ tools from classic coding theory, namely, cyclic MDS codes, in order to obtain a deterministic construction which compares favourably with existing solutions; both in the applicable range of parameters, and in the complexity of the involved algorithms. Some of these gains are a direct application of well known properties of these codes.
This paper is organized as follows. Related work regarding gradient coding (and coded computation in general) is listed in Section II. A framework which encapsulates all the results in this paper is given in Section III. Necessary mathematical notions from coding theory and graph theory are given in Section IV. The former notions are used to obtain an algorithm for exact gradient computation in Section V, and the latter ones are used to obtain an algorithm for the approximate gradient in Section VI. Experimental results are given in Section VII.
II Related Work
The work of Lee et al. initiated the use of coding theoretic methods for mitigating stragglers in large-scale learning. This work is focused on linear regression and therefore can exploit more structure compared to the general gradient coding problem that we study here. The work by Li et al. , investigates a generalized view of the coding ideas in , showing that their solution is a single operating point in a general scheme of trading off latency of computation to the load of communication.
Further closely related work has shown how coding can be used for distributed MapReduce, as well as a similar communication and computation tradeoff . We also mention the work of which addresses straggler mitigation in linear regression by using a different approach, that is not mutually exclusive with gradient coding. In their work, the data is coded rather than replicated, and the nodes perform their computation on coded data.
The work by generalizes previous work for linear models but can also be applied to general models to yield explicit gradient coding constructions. Our results regarding the exact gradient are closely related to the work by which was obtained independently from our work. In , similar coding theoretic tools were employed in a fundamentally different fashion. Both and are comparable in parameters to the randomized construction of and are outperformed by us in a wide range of parameters. A detailed comparison of the theoretical asymptotic behaviour is given in the sequel.
None of the aforementioned works studies approximate gradient computations. However, we note that subsequent to this work, two unpublished manuscripts study a similar approximation setting and obtain related results albeit using randomized as opposed to deterministic approaches. Furthermore, the exact setting was also discussed subsequent to this work in and . In it was shown that network communication can be reduced by increasing the replication factor, and respective bounds were given. The work of discussed coded polynomial computation with low overhead, and applies to gradient coding whenever the gradient at hand is a polynomial.
III Framework
Notice that the error term in Definition 3 is a function of the number of stragglers since it is not expected to decrease if more stragglers are present. The conditions which are given in Definition 2 and Definition 3 guarantee the exact and approximate computation by the following lemmas. In the upcoming proofs, let be the matrix
If and B satisfy the EC condition, then for all we have .
For a given , let be the matrix whose ’th row equals if , and zero otherwise. By the definition of C in Algorithm 1 it follows that , and since it follows that . Therefore, we have
The next lemma bounds the deviance of from the gradient of the empirical risk at the current model by using the function and the spectral norm of . Recall that for a matrix P the spectral norm is defined as .
For a function as above, if and B satisfy the -AC condition, then .
Due to Lemma 4 and Lemma 5, in the remainder of this paper we focus on constructing and B that satisfy either the EC condition (Section V) or the -AC condition (Section VI).
IV Mathematical Notions
This section provides a brief overview on the mathematical notions that are essential for the suggested schemes. The exact computation (Sec. V) requires notions from coding theory, and the approximate one (Sec. VI) requires notions from graph theory. The coding theoretic material in this section is taken from , which focuses on finite fields, and yet the given results extend verbatim to the real or complex case (see also , Sec. 8.4).
[22, Sec. 8] A code is called cyclic if the cyclic shift of every codeword is also a codeword, namely,
is an MDS code, and hence its minimum Hamming weight is .
For any subset of size there exists a codeword in whose support (i.e., the set of nonzero indices) is .
The reverse code is an MDS code.
Further, the structure of may also imply a lower bound on the distance of .
An infinite sequence of regular graphs on vertices satisfies that , where is an expression which tends to zero as tends to infinity.
Constant degree regular graphs (i.e., families of graphs with fixed degree that does not depend on ) for which is small in comparison with are largely referred to as expanders. In particular, graphs which attain the above bound asymptotically (i.e., ) are called Ramanujan graphs, and several efficient constructions are known . Since explicit constructions of Ramanujan graphs are often rather intricate, one may resort to choosing a random regular graph and verify its expansion by computing the respective eigenvalues. This process that is known to produce a good expander with high probability [10, Theorem 7.10], and is used in our experiments.
V Exact Gradient Coding from Cyclic MDS Codes
The matrix B satisfies the following properties.
for every row of B.
Every row of B is a codeword in .
The column span of B is the code .
To prove B1 and B2, observe that B is of the following form, where .
To prove B3, notice that the leftmost columns of B have leading coefficients in different positions, and hence they are linearly independent. Thus, the dimension of the column span of B is at least , and since , the claim follows.
. The above and B satisfy the EC condition (Definition 2).
In the remainder of this section, two cyclic MDS codes over the complex numbers and the real numbers are suggested, from which the construction in Theorem 13 can be obtained. These constructions are taken from (Sec. II.B), and are given with a few adjustments to our case. The contributions of these codes is summarized in the following theorem.
For any given and there exist explicit complex valued and B that satisfy the EC-condition with optimal . The respective encoding (i.e., constructing B) and decoding (i.e., constructing given ) complexities are and , respectively. In addition, for any given and such that there exist explicit real valued and B that satisfy the EC-condition with optimal . The encoding and decoding complexities are and , where is the complexity of inverting a generalized Vandermonde matrix.
Note that the use of complex rather than real matrix B may potentially double the required bandwidth, since every complex number contains two real numbers. A simple manipulation of Algorithm 1 which resolves this issue is given in Appendix C. This optimal bandwidth is also attained by the scheme in the next section, which uses a smaller number of multiplication operations. However, it is applicable only if .
V-B Cyclic-MDS Codes Over the Real Numbers
If one wishes to abstain from using complex numbers, e.g., in order to reduce bandwidth, we suggest the following construction, which provides a cyclic MDS code over the reals. This construction relies on (Property 3), with an additional specialized property.
For a given and such that , define the following BCH codes over the reals. In both cases denote .
According to Lemma 9, it is clear that and are cyclic. According to the BCH bound (Theorem 10), it is also clear that the minimum distance of is at least , and the minimum distance of is at least . Hence, to prove that and are MDS codes, it is shown that their code dimensions are .
Since the sets and are closed under conjugation (i.e., is in if and only if the conjugate of is in ) it follows that the polynomials and have real coefficients. Hence, by the definition of BCH codes it follows that
and hence, and . Let and be the minimum distances of and , respectively, and notice that by the Singleton bound (Sec. 4.1) it follows that
Algorithms for computing the matrix B and the vector for the codes in this subsection are given in Appendix B. The algorithm for construction B outperforms previous works whenever , and the algorithm for computing outperforms previous works for a smaller yet wide range of values.
VI Approximate Gradient Coding from Expander Graphs
Recall that in order to retrieve that exact gradient, one must have , an undesirable overhead in many cases. To break this barrier, we relax the requirement to retrieve the gradient exactly, and settle for an approximation of it. Note that trading the exact gradient for an approximate one is a necessity in many variants of gradient descent (such as the acclaimed stochastic gradient descent [21, Sec. 14.3]), and hence our techniques are aligned with common practices in machine learning.
We show that this can be outperformed by setting B to be a normalized adjacency matrix of a connected regular graph on nodes, which is constructed by the master before dispersing the data, and setting to be some simple function.
The resulting error function depends on the parameters of the graph, whereas the resulting storage overhead is given by its degree (i.e., the fixed number of neighbors of each node). The error function is given below for a general connected and regular graph, and particular examples with their resulting errors are given in the sequel. In particular, it is shown that taking the graph to be an expander graph provides an error term which is smaller than (3) for any .
For any we have .
Notice that the eigenvalues of B are , and hence equals . Further, the eigenvectors are identical to those of . Therefore, it follows from Corollary 21 that
and since are orthonormal, it follows that
The above and B satisfy the -AC condition for . The storage overhead of this scheme equals the degree of the underlying regular graph .
It is evident that in order to obtain small deviation , it is essential to have a small and a large . However, most constructions of expanders have focused in the case were is constant (i.e., ). On one hand, constant serves our purpose well in terms of storage overhead, since it implies a constant storage overhead. On the other hand, a constant does not allow to tend to zero as tends to infinity due to Theorem 11.
It is readily verified that our scheme outperforms the trivial one by a multiplicative factor of , which is less than one for every non-bipartite graph; bipartite graphs are used in a slightly different fashion in the sequel. We conclude the discussion with several examples, the first of which uses Margulis graphs (, Sec. 8), that are rather easy to construct.
For any integer there exists an -regular graph on nodes with . For example, by using these graphs with the parameters , , we have an improvement factor of .
Several additional examples for Ramanujan graphs, which attain an improvement factor but are harder to construct, are as follows.
Let and be distinct primes such that , , and such that the Legendre symbol is (i.e., is a quadratic residue modulo ). Then, there exist a non-bipartite Ramanujan graph on nodes and constant degree .
If and then , , and .
If and then , , and .
Restricting to be a constant (i.e., not to grow with ) is detrimental to the improvement factor due to Theorem 11, but allows lower storage overhead. If one wishes a smaller (i.e., better) improvement factor at the price of higher overhead, the following is useful.
There exists a polynomial algorithm (in ) to produce a graph with the parameters . For this family of graphs, the term goes to zero as goes to infinity.
The above approximation scheme can be used with a bipartite graph as well. However, bipartite graphs satisfy that , and hence the resulting error function is identical to the error function of the trivial scheme (3), which requires lower overhead, and hence no gain is attained. However, in what follows it is shown that bipartite graphs on nodes can be employed in a slightly different fashion, and obtain for some that is defined shortly.
Let be a -regular (in both sides) and connected bipartite graph on nodes (and hence ), with an adjacency matrix
for some real matrix C, and eigenvalues . Let be the set of triples of left-singular vectors, right-singular vectors, and singular values of C, as explained above, where . The next well-known lemma presents the connection between the singular values of C and the eigenvalues of . Its proof is a combination of a few simple exercises, and is given for completeness.
.
and hence . It is an easy exercise to verify that the eigenvalues of a bipartite graph are symmetric around zero, and hence it follows that as well.
Therefore, it follows that , and thus . ∎
From Lemma 27, the following corollaries are easy to prove.
Since , it follows from Lemma 27 that . Further, since is connected, it follows that , and hence .∎
By setting , we have the following lemma, which may be seen as the equivalent of Lemma 22 to the bipartite case.
By Corollary 28, and since is an orthonormal set, it follows that
Applying the above lemma on several constructions of bipartite expanders provides the following examples.
Let and be distinct primes such that , , and such that the Legendre symbol is . Then, there exists a bipartite Ramanujan graph on nodes with {\color[rgb]{0,0,0}\definecolor[named]{pgfstrokecolor}{rgb}{0,0,0}\pgfsys@color@gray@stroke{0}\pgfsys@color@gray@fill{0}\sigma_{2}}\leq 2\sqrt{p}.
If and then , , and .
If and then , , and .
VI-B Lower bound.
Finally, we have the following lower bound on the approximation error of any Approximate Computation (AC) scheme, that establishes asymptotic optimality of our scheme, up to constants, when used with Ramanujan graphs. In what follows, for any set of vertices in a graph , let be the set of vertices which has at least one neighbor in , and for a vertex let denote its degree.
Let be a bipartite graph where for some and for every and some . Then, for every there exists a set of size such that .
Let and without loss of generality assume that . Due to this ordering of , and due to the bounded degree of vertices in , it follows that , where for every . Therefore, the average degree of is at most , which implies that . Since the size of is at most the sum of degrees in , it follows that . ∎
Associate a bipartite graph with B as follows. Consider left vertices corresponding to workers, and right vertices corresponding to data parts. We draw an edge between worker and data part if . According to Lemma 31, for any such that there exists a set of size such that .
The above lemma establishes the asymptotic optimality (up to constants) of our scheme, when used with Ramanujan graphs. Recall that for Ramanujan graphs we have . Thus, using our proposed scheme, we get an and B that satisfy the -AC condition for which tends to as .
VI-C A few remarks about convergence
Let be a random variable that equals if the ’th server responded, and otherwise. Further, if no server responded then , and hence we have that , where is a random variable such that
and thus convergence of the algorithm can be guaranteed as a special case of the SGD algorithm, by adjusting the learning rate properly.
By performing some technical calculations, that are given in full in Appendix D, we have that
Therefore, Theorem 3 and the analysis in Section VI show that our algorithm provides better bounds than the trivial algorithm, and thus it is likely to converge faster in many cases.
VII Experimental Results
In this section, we present experimental results of our proposed approximate gradient coding scheme (Sec. VI).
where is the submatrix of B which consists of the rows that are indexed by . Note that even though we have no additional theoretical guarantees for the optimal decoder, it is always possible to compute it in cubic time, e.g., by the singular value decomposition of .
VII-B Generalization Error
In this section, our approximate gradient coding (AGC) scheme is compared to other approaches. We compare against gradient coding from (GC), as well as the trivial scheme (IS), where the data is divided equally among all workers, but the master only uses the first gradients.
We measured the performance of our coding schemes in terms of the area under the curve (AUC) on a validation set for a logistic regression problem, on a real dataset. The dataset we used was the Amazon Employee dataset from Kaggle. We used training samples, and a model dimension of (after one-shot encoding with interaction terms), and used gradient descent to train the logistic regression. For GC we used a constant learning rate, chosen using cross-validation. For AGC and IS we used a learning rate of , which is typical for SGD, where and were also chosen via cross-validation.
All our methods were implemented in python using MPI4py (similar to ). We ran our experiments using t2.micro worker instance types on Amazon EC2 and a c3.8xlarge master instance type. The results for are given in Fig. 3 and Fig. 4, in which AGC corresponds to our approximation schemes with the optimal decoder, whereas AGC (Linear), termed AGCL is our full proposed approximation scheme.
We observe that both these approaches are only slightly worse than GC, which utilizes the full gradient, and are quite better than the IS approach. Compared to each other, AGC and AGCL seem equivalent, however AGC was marginally better. That being said, AGCL can be faster since computing the linear decoder only requires time, in contrast to time for the optimal decoder.
Acknowledgments
This research has been supported by NSF Grants CCF 1422549, 1618689, DMS 1723052, ARO YIP W911NF- 14-1-0258 and research gifts by Google, Western Digital and NVIDIA. The work of Rashish Tandon was done while he was at UT Austin, prior to joining apple. The work of Itzhak Tamo and Netanel Raviv was supported in part ISF Grant 1030/15 and NSF-BSF Grant 2015814. The work of Netanel Raviv was supported in part by the postdoctoral fellowship of the Center for the Mathematics of Information (CMI) in the California Institute of Technology, and in part by the Lester-Deutsch postdoctoral fellowship. The authors express their gratitude to Prof. Roi Livni for his valuable input.
References
Appendix A
In this section, efficient algorithms for encoding (i.e., computing the matrix B) and decoding (i.e., computing the vector given a set of non-stragglers) are given for the scheme in Section V-A.
Since the matrix B is circulant, it suffices to compute only its leftmost column. Further, the leftmost column of B is a codeword in an Reed-Solomon code whose evaluation points are all roots of unity of order , denoted . Hence, to find , one can define the polynomial and evaluate it over , which is possible in operations. This compares favorably with the respective of (Sec. 5.1.1) for any .
Let be the code from Subsection V-A, and notice that it is a GRS code whose column multipliers are all equal to . By Lemma 35, it follows that the generator matrix of is , where
Solving the equation amounts to an interpolation problem, i.e., finding a degree (at most) polynomial which passes through given points. This is possible in operations by .
Given , computing the product reduces to evaluation of a degree (at most) polynomial on all roots of unity of order . This is possible in operations by utilizing the famous Fast Fourier Transform (FFT) .
Hence, the total complexity of Algorithm 2 is . ∎
Appendix B
where are the roots of the code.
As in Appendix A, to compute the matrix B it suffices to find the codeword , which in this case is a lowest weight codeword in a BCH code. It is readily verified that may be given by the coefficients of the generator polynomial (denoted by and in Eq. (V-B)) . Finding these coefficients is possible by evaluating in arbitrary and pairwise distinct roots of unity of order , and solving an interpolation problem in by . However, evaluating at points requires operations using the FFT algorithm. Hence, the overall complexity of computing B is , an improvement over whenever .
In Algorithm 3, for a given , let be the matrix of columns of V that are indexed by . The complexity of this algorithm outperforms whenever , and outperforms whenever .
As in the complexity analysis of Algorithm 2, the complexity is clearly , where is the complexity of inverting a generalized Vandermonde matrix. While explicit formulas for inverting a generalized Vandermonde matrix were discussed in , for simplicity, one may employ the trivial algorithm. ∎
Appendix C Bandwidth reduction in Subsection V-A
and hence the correctness of Algorithm 1 under the scheme from Subsection V-A is preserved. Under this framework, the transmission from to the master node contains complex numbers (i.e., real numbers) rather than complex numbers (i.e., real numbers), and it is clearly optimal.
Appendix D
We operate under the convention that if all servers are stragglers (i.e., ), then the algorithm outputs the vector as the approximation of the gradient. To analyze the trivial algorithm under this convention, define the random variable
To conduct a similar analysis for the scheme in Section VI we observe that
In the spirit of the definition of above, we define a random variable