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 nn worker nodes (or servers), in which the data is partitioned into kk 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 ss machines fail to perform their work. The storage overhead of the system, which is denoted by dd, 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 ×5\times 5 slower than the typical worker (×8\times 8 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 n−sn-s non-straggling nodes. The work of established the fundamental bound d≥s+1d\geq s+1, provided a deterministic construction which achieves it with equality when s+1∣ns+1|n, and a randomized one which applies to all ss and nn. Subsequently, deterministic constructions were also obtained by and . These works have focused on the scenario where ss 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 ss, but no error bound is guaranteed if this number exceeds ss.

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 ϵ\epsilon 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 N(w)\textbf{N}(\bm{w}) be the matrix

If aa and B satisfy the EC condition, then for all r∈[t]r\in[t] we have vr=∇LS(w(r))\bm{v}_{r}=\nabla L_{\mathcal{S}}(\bm{w}^{(r)}).

For a given r∈[t]r\in[t], let B′\textbf{B}^{\prime} be the matrix whose ii’th row bi′\bm{b}_{i}^{\prime} equals bi\bm{b}_{i} if i∈Kri\in\mathcal{K}_{r}, and zero otherwise. By the definition of C in Algorithm 1 it follows that C=B′⋅N(w(r))\textbf{C}=\textbf{B}^{\prime}\cdot\textbf{N}(w^{(r)}), and since supp⁡(a(Kr))⊆Kr\operatorname{supp}(a(\mathcal{K}_{r}))\subseteq\mathcal{K}_{r} it follows that a(Kr)B′=a(Kr)Ba(\mathcal{K}_{r})\textbf{B}^{\prime}=a(\mathcal{K}_{r})\textbf{B}. Therefore, we have

The next lemma bounds the deviance of vr\bm{v}_{r} from the gradient of the empirical risk at the current model w(r)\bm{w}^{(r)} by using the function ϵ\epsilon and the spectral norm ∥⋅∥spec\lVert\cdot\rVert_{\text{spec}} of N(w)\textbf{N}(\bm{w}). Recall that for a matrix P the spectral norm is defined as ∥P∥spec≜max⁡x:∥x∥2=1∥Px∥2\lVert\textbf{P}\rVert_{\text{spec}}\triangleq\max_{\bm{x}:\lVert\bm{x}\rVert_{2}=1}\lVert\textbf{P}\bm{x}\rVert_{2}.

For a function ϵ\epsilon as above, if aa and B satisfy the ϵ\epsilon-AC condition, then d2(vr,∇LS(w(r)))≤ϵ(∣Krc∣)⋅∥N(w(r))∥specd_{2}(\bm{v}_{r},\nabla L_{\mathcal{S}}(\bm{w}^{(r)}))\leq\epsilon(|\mathcal{K}_{r}^{c}|)\cdot\lVert\textbf{N}(\bm{w}^{(r)})\rVert_{\text{spec}}.

Due to Lemma 4 and Lemma 5, in the remainder of this paper we focus on constructing aa and B that satisfy either the EC condition (Section V) or the ϵ\epsilon-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 C\mathcal{C} is called cyclic if the cyclic shift of every codeword is also a codeword, namely,

C⊥\mathcal{C}^{\bot} is an [n,n−κ][n,n-\kappa] MDS code, and hence its minimum Hamming weight is κ+1\kappa+1.

For any subset K⊆[n]\mathcal{K}\subseteq[n] of size n−κ+1n-\kappa+1 there exists a codeword in C\mathcal{C} whose support (i.e., the set of nonzero indices) is K\mathcal{K}.

The reverse code CR≜{(cn,…,c1)∣(c1,…,cn)∈C}\mathcal{C}^{R}\triangleq\{(c_{n},\ldots,c_{1})|(c_{1},\ldots,c_{n})\in\mathcal{C}\} is an [n,κ][n,\kappa] MDS code.

Further, the structure of RR may also imply a lower bound on the distance of CC.

An infinite sequence of dd regular graphs on nn vertices satisfies that λ≥2d−1−on(1)\lambda\geq 2\sqrt{d-1}-o_{n}(1), where on(1)o_{n}(1) is an expression which tends to zero as nn tends to infinity.

Constant degree regular graphs (i.e., families of graphs with fixed degree dd that does not depend on nn) for which λ\lambda is small in comparison with dd are largely referred to as expanders. In particular, graphs which attain the above bound asymptotically (i.e., λ≤2d−1\lambda\leq 2\sqrt{d-1}) 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.

∥b∥0=s+1\lVert\bm{b}\rVert_{0}=s+1 for every row b\bm{b} of B.

Every row of B is a codeword in CR\mathcal{C}^{R}.

The column span of B is the code C\mathcal{C}.

To prove B1 and B2, observe that B is of the following form, where c1≜(β1,…,βs+1,0,…,0)\bm{c}_{1}\triangleq(\beta_{1},\ldots,\beta_{s+1},0,\ldots,0).

To prove B3, notice that the leftmost n−sn-s 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 n−sn-s, and since dim⁡C=n−s\dim\mathcal{C}=n-s, the claim follows.

. The above aa 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 nn and ss there exist explicit complex valued aa and B that satisfy the EC-condition with optimal d=s+1d=s+1. The respective encoding (i.e., constructing B) and decoding (i.e., constructing a(K)a(\mathcal{K}) given K\mathcal{K}) complexities are O(s(n−s))O(s(n-s)) and O(slog⁡2s+nlog⁡n)O(s\log^{2}s+n\log n), respectively. In addition, for any given nn and ss such that n≠s mod 2n\neq s\bmod 2 there exist explicit real valued aa and B that satisfy the EC-condition with optimal d=s+1d=s+1. The encoding and decoding complexities are O(min⁡{slog⁡2s,nlog⁡n})O(\min\{s\log^{2}s,n\log n\}) and O(gs+s(n−s))O(g_{s}+s(n-s)), where gsg_{s} 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 n≠s mod 2n\neq s\bmod 2.

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 nn and ss such that n≠s mod 2n\neq s\bmod 2, define the following BCH codes over the reals. In both cases denote ω≜e2πi/n\omega\triangleq e^{2\pi i/n}.

According to Lemma 9, it is clear that C1\mathcal{C}_{1} and C2\mathcal{C}_{2} are cyclic. According to the BCH bound (Theorem 10), it is also clear that the minimum distance of C1\mathcal{C}_{1} is at least ∣R1∣+1=s+1|\mathcal{R}_{1}|+1=s+1, and the minimum distance of C2\mathcal{C}_{2} is at least ∣R2∣+1=s+1|\mathcal{R}_{2}|+1=s+1. Hence, to prove that C1\mathcal{C}_{1} and C2\mathcal{C}_{2} are MDS codes, it is shown that their code dimensions are n−sn-s.

Since the sets R1\mathcal{R}_{1} and R2\mathcal{R}_{2} are closed under conjugation (i.e., rr is in Ri\mathcal{R}_{i} if and only if the conjugate of rr is in Ri\mathcal{R}_{i}) it follows that the polynomials p1(x)≜∏r∈R1(x−r)p_{1}(x)\triangleq\prod_{r\in\mathcal{R}_{1}}(x-r) and p2(x)≜∏r∈R2(x−r)p_{2}(x)\triangleq\prod_{r\in\mathcal{R}_{2}}(x-r) have real coefficients. Hence, by the definition of BCH codes it follows that

and hence, dim⁡C1≥n−s\dim\mathcal{C}_{1}\geq n-s and dim⁡C2≥n−s\dim\mathcal{C}_{2}\geq n-s. Let d(C1)d(\mathcal{C}_{1}) and d(C2)d(\mathcal{C}_{2}) be the minimum distances of C1\mathcal{C}_{1} and C2\mathcal{C}_{2}, respectively, and notice that by the Singleton bound (Sec. 4.1) it follows that

Algorithms for computing the matrix B and the vector a(K)a(\mathcal{K}) for the codes in this subsection are given in Appendix B. The algorithm for construction B outperforms previous works whenever s=o(n)s=o(n), and the algorithm for computing a(K)a(\mathcal{K}) outperforms previous works for a smaller yet wide range of ss values.

VI Approximate Gradient Coding from Expander Graphs

Recall that in order to retrieve that exact gradient, one must have d≥s+1d\geq s+1, 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 nn nodes, which is constructed by the master before dispersing the data, and setting aa to be some simple function.

The resulting error function ϵ(s)\epsilon(s) depends on the parameters of the graph, whereas the resulting storage overhead dd 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 ϵ\epsilon which is smaller than (3) for any ss.

For any K⊆[n]\mathcal{K}\subseteq[n] we have uK∈⟨v2,…,vn⟩\bm{u}_{\mathcal{K}}\in{\left\langle{\bm{v}_{2},\ldots,\bm{v}_{n}}\right\rangle}.

Notice that the eigenvalues of B are μi≜λid\mu_{i}\triangleq\frac{\lambda_{i}}{d}, and hence μ≜max⁡{∣μ2∣,∣μn∣}\mu\triangleq\max\{|\mu_{2}|,|\mu_{n}|\} equals λd\frac{\lambda}{d}. Further, the eigenvectors are identical to those of AG\textbf{A}_{G}. Therefore, it follows from Corollary 21 that

and since {v2,…,vn}\{\bm{v}_{2},\ldots,\bm{v}_{n}\} are orthonormal, it follows that

The above aa and B satisfy the ϵ\epsilon-AC condition for ϵ(s)=λdnsn−s\epsilon(s)=\frac{\lambda}{d}\sqrt{\frac{ns}{n-s}}. The storage overhead of this scheme equals the degree dd of the underlying regular graph GG.

It is evident that in order to obtain small deviation ϵ(s)\epsilon(s), it is essential to have a small λ\lambda and a large dd. However, most constructions of expanders have focused in the case were dd is constant (i.e., d=O(1)d=O(1)). On one hand, constant dd serves our purpose well in terms of storage overhead, since it implies a constant storage overhead. On the other hand, a constant dd does not allow λ/d\lambda/d to tend to zero as nn tends to infinity due to Theorem 11.

It is readily verified that our scheme outperforms the trivial one by a multiplicative factor of λd\frac{\lambda}{d}, 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 nn there exists an 88-regular graph on nn nodes with λ≤52\lambda\leq 5\sqrt{2}. For example, by using these graphs with the parameters n=500n=500, d=8d=8, we have an improvement factor of λd=528≈0.883\frac{\lambda}{d}=\frac{5\sqrt{2}}{8}\approx 0.883.

Several additional examples for Ramanujan graphs, which attain an improvement factor ≤2d−1d≈2d\leq\frac{2\sqrt{d-1}}{d}\approx\frac{2}{\sqrt{d}} but are harder to construct, are as follows.

Let pp and qq be distinct primes such that p=1 mod 4p=1\bmod 4, q=1 mod 4q=1\bmod 4, and such that the Legendre symbol (pq)\left(\frac{p}{q}\right) is 11 (i.e., pp is a quadratic residue modulo qq). Then, there exist a non-bipartite Ramanujan graph on n=q(q2−1)2n=\frac{q(q^{2}-1)}{2} nodes and constant degree p+1p+1.

If p=13p=13 and q=17q=17 then n=2448n=2448, d=14d=14, and 2d−1d≈0.534\frac{2\sqrt{d-1}}{d}\approx 0.534.

If p=5p=5 and q=29q=29 then n=12180n=12180, d=6d=6, and 2d−1d≈0.816\frac{2\sqrt{d-1}}{d}\approx 0.816.

Restricting dd to be a constant (i.e., not to grow with nn) is detrimental to the improvement factor λd\frac{\lambda}{d} 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 nn) to produce a graph GG with the parameters (n,d,λ)=(2m,m−1,mlog⁡3m)(n,d,\lambda)=(2^{m},m-1,\sqrt{m\log^{3}m}). For this family of graphs, the term λd\frac{\lambda}{d} goes to zero as nn goes to infinity.

The above approximation scheme can be used with a bipartite graph GG as well. However, bipartite graphs satisfy that λ=d\lambda=d, and hence the resulting error function ϵ(s)=nsn−s\epsilon(s)=\sqrt{\frac{ns}{n-s}} 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 2n2n nodes can be employed in a slightly different fashion, and obtain ϵ(s)=σ2dnsn−s\epsilon(s)=\frac{\sigma_{2}}{d}\sqrt{\frac{ns}{n-s}} for some σ2<d\sigma_{2}<d that is defined shortly.

Let G=(L∪R,E)G=(\mathcal{L}\cup\mathcal{R},\mathcal{E}) be a dd-regular (in both sides) and connected bipartite graph on 2n2n nodes (and hence ∣L∣=∣R∣=n|\mathcal{L}|=|\mathcal{R}|=n), with an adjacency matrix

for some n×nn\times n real matrix C, and eigenvalues λ1=d≥λ2≥⋯≥λ2n=−d\lambda_{1}=d\geq\lambda_{2}\geq\cdots\geq\lambda_{2n}=-d. Let {(ui,vi,σi)}i=1n\{(\bm{u}_{i},\bm{v}_{i},\sigma_{i})\}_{i=1}^{n} be the set of triples of left-singular vectors, right-singular vectors, and singular values of C, as explained above, where σ1≥⋯≥σn≥0\sigma_{1}\geq\cdots\geq\sigma_{n}\geq 0. The next well-known lemma presents the connection between the singular values of C and the eigenvalues of AG\textbf{A}_{G}. Its proof is a combination of a few simple exercises, and is given for completeness.

{λi}i=12n={σi}i=1n∪{−σi}i=1n\{\lambda_{i}\}_{i=1}^{2n}=\{\sigma_{i}\}_{i=1}^{n}\cup\{-\sigma_{i}\}_{i=1}^{n}.

and hence σi∈{λi}i=12n\sigma_{i}\in\{\lambda_{i}\}_{i=1}^{2n}. It is an easy exercise to verify that the eigenvalues of a bipartite graph are symmetric around zero, and hence it follows that −σi∈{λi}i=12n-\sigma_{i}\in\{\lambda_{i}\}_{i=1}^{2n} as well.

Therefore, it follows that λi2∈{σi2}i=1n\lambda_{i}^{2}\in\{\sigma_{i}^{2}\}_{i=1}^{n}, and thus λi∈{σi}i=1n∪{−σi}i=1n\lambda_{i}\in\{\sigma_{i}\}_{i=1}^{n}\cup\{-\sigma_{i}\}_{i=1}^{n}. ∎

From Lemma 27, the following corollaries are easy to prove.

Since σ1≥…≥σn\sigma_{1}\geq\ldots\geq\sigma_{n}, it follows from Lemma 27 that λ2=σ2\lambda_{2}=\sigma_{2}. Further, since GG is connected, it follows that λ2<d\lambda_{2}<d, and hence σ2<d\sigma_{2}<d.∎

By setting B≜1dC\textbf{B}\triangleq\frac{1}{d}\textbf{C}, we have the following lemma, which may be seen as the equivalent of Lemma 22 to the bipartite case.

By Corollary 28, and since {vi}i=2n\{\bm{v}_{i}\}_{i=2}^{n} is an orthonormal set, it follows that

Applying the above lemma on several constructions of bipartite expanders provides the following examples.

Let pp and qq be distinct primes such that p=1 mod 4p=1\bmod 4, q=1 mod 4q=1\bmod 4, and such that the Legendre symbol (pq)\left(\frac{p}{q}\right) is −1-1. Then, there exists a bipartite Ramanujan graph on 2n=q(q2−1)2n=q(q^{2}-1) 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 p=5p=5 and q=13q=13 then n=21842=1092n=\frac{2184}{2}=1092, d=6d=6, and σ2d≤256≈0.745\frac{\sigma_{2}}{d}\leq\frac{2\sqrt{5}}{6}\approx 0.745.

If p=13p=13 and q=5q=5 then n=1202=60n=\frac{120}{2}=60, d=14d=14, and σ2d≤21314≈0.515\frac{\sigma_{2}}{d}\leq\frac{2\sqrt{13}}{14}\approx 0.515.

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 A\mathcal{A} in a graph GG, let N(A)\mathcal{N}(\mathcal{A}) be the set of vertices which has at least one neighbor in A\mathcal{A}, and for a vertex VV let d(V)d(V) denote its degree.

Let G=(W∪P,E)G=(\mathcal{W}\cup\mathcal{P},\mathcal{E}) be a bipartite graph where ∣W∣=∣P∣=n|\mathcal{W}|=|\mathcal{P}|=n for some nn and d(W)≤dd(W)\leq d for every W∈WW\in\mathcal{W} and some dd. Then, for every r∈{1,2,…,⌊nd⌋}r\in\{1,2,\ldots,\lfloor\frac{n}{d}\rfloor\} there exists a set Qr⊆P\mathcal{Q}_{r}\subseteq\mathcal{P} of size rr such that ∣N(Qr)∣≤d∣Qr∣|\mathcal{N}(\mathcal{Q}_{r})|\leq d|\mathcal{Q}_{r}|.

Let P={P1,…,Pn}\mathcal{P}=\{P_{1},\ldots,P_{n}\} and without loss of generality assume that d(P1)≤d(P2)≤…≤d(Pn)d(P_{1})\leq d(P_{2})\leq\ldots\leq d(P_{n}). Due to this ordering of P1,…,PnP_{1},\ldots,P_{n}, and due to the bounded degree of vertices in W\mathcal{W}, it follows that d≥avgn≥avgn−1≥avg1d\geq\text{avg}_{n}\geq\text{avg}_{n-1}\geq\text{avg}_{1}, where avgj≜1j∑i=1jd(Pi)\text{avg}_{j}\triangleq\frac{1}{j}\sum_{i=1}^{j}d(P_{i}) for every j∈[n]j\in[n]. Therefore, the average degree of Qr≜{P1,…,Pr}\mathcal{Q}_{r}\triangleq\{P_{1},\ldots,P_{r}\} is at most dd, which implies that ∑j=1rd(Pi)≤rd\sum_{j=1}^{r}d(P_{i})\leq rd. Since the size of N(Qr)\mathcal{N}(\mathcal{Q}_{r}) is at most the sum of degrees in Qr\mathcal{Q}_{r}, it follows that ∣N(Qr)∣≤∑j=1rd(Pi)≤rd=d∣Qr∣|\mathcal{N}(\mathcal{Q}_{r})|\leq\sum_{j=1}^{r}d(P_{i})\leq rd=d|\mathcal{Q}_{r}|. ∎

Associate a bipartite graph G=(W∪P,E)G=(\mathcal{W}\cup\mathcal{P},\mathcal{E}) with B as follows. Consider left vertices W={W1,W2,…,Wn}\mathcal{W}=\{W_{1},W_{2},\ldots,W_{n}\} corresponding to workers, and right vertices P={P1,…,Pn}\mathcal{P}=\{P_{1},\ldots,P_{n}\} corresponding to data parts. We draw an edge {Wi,Pj}∈E\{W_{i},P_{j}\}\in\mathcal{E} between worker WiW_{i} and data part PjP_{j} if (B)i,j≠0(\textbf{B})_{i,j}\neq 0. According to Lemma 31, for any ss such that d<s<nd<s<n there exists a set Q⊆P\mathcal{Q}\subseteq\mathcal{P} of size ⌊sd⌋\lfloor\frac{s}{d}\rfloor such that ∣N(Q)∣≤d⋅⌊sd⌋≤s|\mathcal{N}(\mathcal{Q})|\leq d\cdot\lfloor\frac{s}{d}\rfloor\leq s.

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 λ≤2d\lambda\leq 2\sqrt{d}. Thus, using our proposed scheme, we get an aa and B that satisfy the ϵ\epsilon-AC condition for ϵ(s)≤2ds1−(s/n)\epsilon(s)\leq\frac{2}{\sqrt{d}}\sqrt{\frac{s}{1-(s/n)}} which tends to 2sd2\sqrt{\frac{s}{d}} as s/n→0s/n\to 0.

VI-C A few remarks about convergence

Let XiX_{i} be a Bernoulli(q)\text{Bernoulli}(q) random variable that equals 11 if the ii’th server responded, and otherwise. Further, if no server responded then v=0\bm{v}=0, and hence we have that v=z⋅B⋅N(w)\bm{v}=\bm{z}\cdot\textbf{B}\cdot\textbf{N}(\bm{w}), where z\bm{z} 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 B(K:)\textbf{B}(\mathcal{K}:) is the submatrix of B which consists of the rows that are indexed by K\mathcal{K}. 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 B(K,:)\textbf{B}(\mathcal{K},:).

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 n−sn-s 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 26,20026,200 training samples, and a model dimension of 241,915241,915 (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 c1/(r+c2)c_{1}/(r+c_{2}), which is typical for SGD, where c1c_{1} and c2c_{2} 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 n=30,50n=30,50 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 O(n)O(n) time, in contrast to O(n3)O(n^{3}) 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 a(K)a(\mathcal{K}) given a set K\mathcal{K} 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 c1⊤\bm{c}^{\top}_{1} of B is a codeword in an [n,n−s][n,n-s] Reed-Solomon code whose evaluation points are all roots of unity of order nn, denoted {αi}i=0n−1\{\alpha_{i}\}_{i=0}^{n-1}. Hence, to find c1\bm{c}_{1}, one can define the polynomial m(x)≜∏j=s+1n−1(x−αj)m(x)\triangleq\prod_{j=s+1}^{n-1}(x-\alpha_{j}) and evaluate it over α0,α1,…,αs\alpha_{0},\alpha_{1},\ldots,\alpha_{s}, which is possible in O(s(n−s))O(s(n-s)) operations. This compares favorably with the respective O(n2log⁡2(n))O(n^{2}\log^{2}(n)) of (Sec. 5.1.1) for any ss.

Let C\mathcal{C} be the code from Subsection V-A, and notice that it is a GRS code whose column multipliers are all equal to 11. By Lemma 35, it follows that the generator matrix of C⊥\textbf{C}^{\bot} is V⋅D\textbf{V}\cdot\textbf{D}, where

Solving the equation fVKc=−xKc′⋅DKc−1\bm{f}\textbf{V}_{\mathcal{K}^{c}}=-\bm{x}^{\prime}_{\mathcal{K}^{c}}\cdot\textbf{D}_{\mathcal{K}^{c}}^{-1} amounts to an interpolation problem, i.e., finding a degree (at most) s−1s-1 polynomial which passes through ss given points. This is possible in O(slog⁡2s)O(s\log^{2}s) operations by .

Given f\bm{f}, computing the product fV\bm{f}\textbf{V} reduces to evaluation of a degree (at most) nn polynomial on all roots of unity of order nn. This is possible in O(nlog⁡n)O(n\log n) operations by utilizing the famous Fast Fourier Transform (FFT) .

Hence, the total complexity of Algorithm 2 is O(slog⁡2s+nlog⁡n)O(s\log^{2}s+n\log n). ∎

Appendix B

where {αi}i=1s\{\alpha_{i}\}_{i=1}^{s} are the roots of the code.

As in Appendix A, to compute the matrix B it suffices to find the codeword c1\bm{c}_{1}, which in this case is a lowest weight codeword in a BCH code. It is readily verified that c1\bm{c}_{1} may be given by the coefficients of the generator polynomial g(x)≜∏i=1s(x−αi)g(x)\triangleq\prod_{i=1}^{s}(x-\alpha_{i}) (denoted by p1p_{1} and p2p_{2} in Eq. (V-B)) . Finding these coefficients is possible by evaluating g(x)g(x) in s+1s+1 arbitrary and pairwise distinct roots of unity of order nn, and solving an interpolation problem in O(slog⁡2s)O(s\log^{2}s) by . However, evaluating g(x)g(x) at s+1s+1 points requires O(nlog⁡n)O(n\log n) operations using the FFT algorithm. Hence, the overall complexity of computing B is O(min⁡{slog⁡2s,nlog⁡n})O(\min\{s\log^{2}s,n\log n\}), an improvement over whenever s=o(n)s=o(n).

In Algorithm 3, for a given K⊆[n]\mathcal{K}\subseteq[n], let VK\textbf{V}_{\mathcal{K}} be the matrix of columns of V that are indexed by K\mathcal{K}. The complexity of this algorithm outperforms whenever s=o(log⁡2n)s=o(\log^{2}n), and outperforms whenever s=o(n2/3)s=o(n^{2/3}).

As in the complexity analysis of Algorithm 2, the complexity is clearly O(γs+s(n−s))O(\gamma_{s}+s(n-s)), where γs\gamma_{s} 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 O(s3)O(s^{3}) 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 WjW_{j} to the master node contains p′p^{\prime} complex numbers (i.e., 2p′2p^{\prime} real numbers) rather than pp complex numbers (i.e., 2p2p real numbers), and it is clearly optimal.

Appendix D

We operate under the convention that if all servers are stragglers (i.e., K=∅\mathcal{K}=\varnothing), 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 YY above, we define a random variable