A Unified Coded Deep Neural Network Training Strategy Based on Generalized PolyDot Codes for Matrix Multiplication

Sanghamitra Dutta, Ziqian Bai, Haewon Jeong, Tze Meng Low, Pulkit Grover

I Introduction

Ever-increasing data and computing requirements increasingly require massively distributed and parallel processing. However, as the number of parallel processing units are scaling, the expected number of faults, errors or delays in computing are also scaling. Thus, one of the major challenges of large-scale computing today is ensuring “reliability at scale”. Coded computing has emerged as a promising solution to the various problems arising from the unreliability of processing nodes in parallel and distributed computing, such as straggling delays or “soft-errors”. It is, in fact, a significant step in a long line of work on noisy computing started by von Neumann in 1956, that has been followed upon by Algorithm-Based Fault Tolerance (ABFT), the predecessor of coded computing. Recent advances in coded computing have generated significant interest both within and outside information theory.

The problem of distributed matrix-matrix multiplication S=WX\bm{S}=\bm{W}\bm{X} under storage constraints, i.e., when each node is allowed to store a fixed 1K\frac{1}{K} fraction of each of the matrices W\bm{W} and X\bm{X}, has been of considerable interest in the coded computing community. In , the authors introduced Product codes that use two different Maximum Distance Separable (MDS) codes to encode the two matrices being multiplied. This work was followed upon by , where the authors proposed Polynomial codes, that use a polynomial-based encoding scheme for the storage-constrained matrix multiplication problem to achieve a lower recovery threshold, i.e., the number of processing nodes to wait for out of a total of PP nodes. In our prior work , we first demonstrated that the recovery threshold for the storage-constrained matrix multiplication problem can be reduced further (beyond Polynomial codes ) in scaling sense by using novel constructions called MatDot codes. When a fixed 1K\frac{1}{K} fraction of each matrix can be stored at each node, MatDot achieves a recovery threshold of 2K−12K-1 as compared to Polynomial codes which achieve a threshold of K2K^{2}, albeit at a higher communication cost. In fact, MatDot codes can be proved to be optimal for storage-constrained matrix-matrix multiplication using fundamental limits in . In the same work , we also proposed the PolyDot codes for matrix-matrix multiplication that interpolate between Polynomial codes (for low communication costs) and MatDot codes (for lowest recovery threshold), trading off recovery threshold and communication costs, as illustrated in Fig. 1.

In this work, we have two main contributions as follows: 1. Generalized PolyDot Codes: First we build upon our prior work on PolyDot codes and propose a new class of codes for distributed matrix-vector and matrix-matrix multiplication called Generalized PolyDot codes. The proposed Generalized PolyDot codes interpolate better between Polynomial codes and MatDot codes, improving the recovery threshold (see Theorems 1 and 2) by introducing the new idea of “garbage alignment” for this problem. Garbage alignment is essentially a clever substitution in the multivariate polynomial introduced in our prior work on PolyDot framework of matrix multiplication that aligns some of the unwanted coefficients and reduces the number of unknowns during polynomial interpolation, as we elaborate in Section III. The recovery threshold of PolyDot codes are improved by a factor of 22 using Generalized PolyDot codes by replacing the bijection-based substitution with garbage alignment. We note that, a concurrent work , that in fact appeared at the same venue as the original publication of this work , also achieve the same recovery threshold as Generalized PolyDot codes through slightly different routes.

2. Coded DNN Training Strategy: Our next contribution is that we develop a unified coded computing strategy, by appropriately utilizing Generalized PolyDot codes, for the training of model-parallelData parallel and model parallel are two different architectures for DNN training. In data parallelism , different nodes store and train a different replica of the entire DNN on different pieces of data, and a central parameter server combines inputs from all the nodes to train a central replica of the DNN. In model parallelism, different parts of a single DNN are parallelized across multiple nodes. Coding for data parallel training is examined in . Deep Neural Networks in presence of soft-errors. Soft-errors refer to undetected errors, e.g., bit-flips or gate errors in computation, that can corrupt the end resultIgnoring soft-errors entirely during training of DNNs can severely degrade the accuracy of training, as we experimentally observe in .. For this problem, we consider two kinds of error models: an adversarial model and a probabilistic model. Interestingly, as we show in Theorem 3, under the probabilistic model a real number (P,Q)(P,Q) MDS code can theoretically correct P−Q−1P-Q-1 errors with probability 11 as compared to ⌊P−Q2⌋\lfloor\frac{P-Q}{2}\rfloor for the more conventional, adversarial case. While ideas of correcting more errors than P−Q2\frac{P-Q}{2} have been prevalent in finite fields (see ), to the best of our knowledge this seems to be the first result of this nature for real number error correction using MDS codes (see for similar results on LDPC codes).

The problem of coded DNN training under soft-errors also has several additional challenges (see ), that are all addressed by our proposed unified strategy, as follows:

Prohibitive overhead of encoding matrix W\bm{W} at each iteration: Existing works on coded matrix multiplication (e.g. for computing WN×NXN×B\bm{W}_{N\times N}\bm{X}_{N\times B} where N≫BN\gg B) require encoding of both the matrices W\bm{W} and X\bm{X}. Because N≫BN\gg B, encoding W\bm{W} is computationally expensive (complexity Θ(N2P)\Theta(N^{2}P)) and in fact can be as expensive as the matrix multiplication itself (Θ(N2B)\Theta(N^{2}B)). Thus, for this problem, the coding techniques in would be helpful if W\bm{W} is known in advance and is fixed over a large number of computations so that the encoding cost is amortized. However, when training DNNs, because the parameter matrices update at every iteration, a naive extension of existing techniques would require encoding of parameter matrices at every iteration and thus introduce an undesirable additional overhead of Θ(N2P)\Theta(N^{2}P) at every iteration that can no longer be amortized. To address this, we carefully weave Generalized PolyDot codes into the operations of DNN training so that an initial encoding of the weight matrices is maintained across the updates at each iteration. To do so, at each iteration each node locally encodes much smaller matrices consisting of NBNB elements instead of the large matrix W\bm{W} of N2N^{2} elements, adding negligible overhead. In particular, for the case of B=1B=1, this simply reduces to encoding vectors instead of matrices which is much cheaper in terms of computational complexity.

Master node acting as a single point of failure: Because of our focus on soft-errors in this work, if we allow the architecture to use a master node, this node can often become a “single point of failure” (considered undesirable in parallel computing literature, e.g., ). Thus, we consider a completely decentralized setting, with no master node. In that spirit, our strategy allows encoding/decoding to be error-prone as well, along with all the other primary steps, namely, matrix multiplication, nonlinear activation, Hadamard product and update. We only introduce two verification steps (that check for decoding errors by exchanging some values among all nodes and comparing) that are extremely low complexityThe longer a computation, the more is the probability of soft-errors . In fact, the number of soft-errors that occur within a time interval is often modelled as a Poisson random variable with mean proportional to the length of the time interval. and hence, may be assumed to be error-free.

Nonlinear activation between layers: The nonlinear activation (e.g. sigmoid, ReLU) between layers also poses a difficulty for coded training because most coding techniques are linear. To circumvent this issue, we code the linear operations (matrix multiplication of complexity Θ(N2B)\Theta(N^{2}B)) at each layer separately. Matrix multiplication and update are the most critical and complexity-intensive steps in the training of DNNs as compared to other operations such as nonlinear activation or Hadamard product which are of complexity Θ(NB)\Theta(NB), and hence are also more likely to have errors. Moreover, as our implementation is decentralized, every node acts as a low-complexity functional replica of the master node, performing encoding/decoding/nonlinear activation/Hadamard product and helping us detect (and if possible correct) errors in all the primary steps, including the nonlinear activation step.

Overview of results in coded DNN Training: We show (in Theorems 4 and 6) that under both the adversarial and probabilistic error models, the coded DNN strategy using Generalized PolyDot codes improves the error tolerance in scaling sense over competing replication strategy and a preliminary MDS-code-based DNN training strategy. Moreover, to demonstrate the utility of the DNN training strategy, we also show in Theorems 5 and 7 that the additional overhead due to coding per iteration is negligible as compared to the computational complexity of the local matrix operations at each node as long as the number of processors P4=o(N)P^{4}=o(N). Our coding technique also extends to DNN training with regularization as is done more commonly in practice.

For a distributed matrix-matrix multiplication problem S=WX\bm{S}=\bm{W}\bm{X}, Polynomial codes use a horizontal splitting of the first matrix W\bm{W} and vertical splitting of the second matrix X\bm{X} into KK blocks each, thus satisfying the storage constraint (see Fig. 2). Alternately, MatDot codes use a vertical splitting of the first matrix W\bm{W} and a horizontal splitting of the second matrix X\bm{X} into KK blocks each. PolyDot and Generalized PolyDot codes both incorporate simultaneous vertical and horizontal splitting of both matrices W\bm{W} and X\bm{X} to interpolate between Polynomial codes and MatDot codes, while satisfying storage constraints. Interestingly, as it turns out, the problem of coded computing by splitting the matrix W\bm{W} both vertically and horizontally arises naturally in DNN training.

In the problem of DNN training, at each iteration in any particular layer, the same matrix W\bm{W} is required to be multiplied once with another matrix X\bm{X} (or vector x\bm{x}) from the right side in the feedforward stage, i.e., WX\bm{W}\bm{X}, and once with another matrix ΔT\bm{\Delta}^{T} (or vector δT\bm{\delta}^{T}) from the left side in the backpropagation stage, i.e., ΔTW\bm{\Delta}^{T}\bm{W}. One would like to use the same encoding on W\bm{W} (or sub-matrices of W\bm{W}) for both the matrix-matrix multiplication because the available storage is limited and storing two encoded sub-matrices of W\bm{W} for the two matrix-matrix multiplications is expensive. Now, suppose that we choose to use MatDot codes for WX\bm{W}\bm{X} and hence stick with vertical partitioning of the first matrix W\bm{W}. Then, we would have to use Polynomial codes for ΔTW\bm{\Delta}^{T}\bm{W} as W\bm{W} is the second matrix for this multiplication. This is also pictorially illustrated in Fig. 2. Thus, we would like to allow for both horizontal and vertical partitioning of the matrix W\bm{W} into a grid of m×nm\times n sub-matrices with mn=Kmn=K for the storage constraint. This would allow us to be able to interpolate between MatDot and Polynomial codes for the two matrix-matrix multiplications WX\bm{W}\bm{X} and ΔTW\bm{\Delta}^{T}\bm{W}, so as to achieve a good recovery threshold (and hence error tolerance) for both forward and backward matrix-matrix multiplications.

We note that a concurrent work , that in fact appeared at the same conference as the original publication of this work , proposes the coding scheme “Entangled Polynomial codes” that achieve the same recovery threshold as Generalized PolyDot codes through slightly different routes for the problem of distributed matrix multiplication under storage constraints allowing for both vertical and horizontal splitting. In this expanded version of we show that the Generalized PolyDot codes, that were originally proposed for matrix-vector products in , naturally extend to matrix-matrix products as well. This is because Generalized PolyDot codes are simply a clever substitution in a multivariate polynomial in our prior PolyDot framework for matrix-matrix multiplication (this work precedes both and ), so that some unwanted coefficients of the polynomial align with each other, reducing the degree of the polynomial and hence number of unknowns in polynomial interpolation, thereby improving the recovery threshold by a factor of 22. More importantly, our work also introduces a novel result in real-number error correction and is also the first line of work that considers the problem of coded DNN training.

I-B Organization.

The rest of the paper is organized as follows. We introduce our two problem formulations, namely, (1) coded matrix multiplication; and (2) coded DNN, in Section II. First, we address Problem 11, i.e., the coded matrix multiplication problem. For this problem, we introduce some motivating examples in Section III and then describe the Generalized PolyDot codes in detail in Section IV. Then moving on to Problem 22, we first elaborate upon some modelling assumptions, e.g., the two error models for the coded DNN problem in Section V, which leads to a novel result on real number error correction. This is followed by possible solutions for the coded DNN problem for B=1B=1 (existing strategies in Section VI, our proposed strategy in Section VII). We formally analyze the error tolerance and the computation and communication costs of our strategy in Section VIII. Next, we extend the proposed strategy for the case of B>1B>1 in Section IX. Finally, in Section X, we discuss an application of our strategy to more recent but closely related neural network architectures, namely, sparse autoencoders.

II System Models and Problem Formulations for the two problems

System Model: We assume that there is a centralized, reliable master node and PP memory-constrained worker nodes that may be unreliable. The master node allocates computational tasks to the worker nodes. The worker nodes perform their computations in parallel and send their outputs back to the master node. Outputs of some worker nodes may be modeled as erasures (e.g., because of straggling or faults). The master node gathers the outputs of the worker nodes (possibly a subset), and uses them to compute the final result. It is desirable that the computational overhead of the master node as well as the communication costs should be smaller than the local computational complexity of each worker node after parallelization.

Problem Formulation: Compute distributed matrix multiplication S=WX\bm{S}=\bm{W}\bm{X} using PP worker nodes prone to erasures such that each node can store only a fixed fraction 1K\frac{1}{K} of matrix W\bm{W} and 1K′\frac{1}{K^{\prime}} of matrix X\bm{X}. For this problem, our goal is to minimize the erasure recovery threshold, i.e., the number of nodes that the decoder has to wait for out of the total PP nodes, to be able to compute the entire final result. We are also interested in studying tradeoffs between communication costs and recovery threshold.

We assume that both KK and K′K^{\prime} are less than PP.

When K=K′K=K^{\prime}, our prior work on MatDot codes achieves the optimal recovery threshold of 2K−12K-1 under the storage constraint, albeit at a high communication cost. Here, we propose a more general encoding strategy for the problem of coded matrix multiplication under fixed storage constraints that also allows us to tradeoff between recovery threshold and communication costs, with MatDot codes and Polynomial codes being two special cases.

Note that no specific assumptions are made on the relative dimensions of W\bm{W} and X\bm{X} for this problem formulation. The matrices W\bm{W} and X\bm{X} do not necessarily have to be square matrices either.

II-B Problem 222 (Coded DNN).

Background: Before introducing our system model and problem formulation, we first introduce the main computational steps at each iteration of classical (error-free) DNN training using Stochastic Gradient Descent (SGD)As a first step in this direction of coded neural networks, we assume that the training is performed using vanilla SGD. As a future work, we plan to extend these coding ideas to other training algorithms such as momentum SGD, Adam etc.. For a more detailed introduction, we refer to the seminal work or Appendix A.

A DNN consists of LL weight matrices (also called parameter matrices), Wl\bm{W}^{l}, of dimensions Nl×Nl−1N_{l}\times N_{l-1} where ll denotes the layer-index and NlN_{l} denotes the number of neurons at layer ll, for l=1,2,…,Ll=1,2,\ldots,L. These weight matrices are updated at each iteration of training based on a “mini-batch” of BB data-points and their labels. Because we are primarily interested in large models that require parallelization across multiple nodes, we assume Nl≫BN_{l}\gg B for all ll. For simplicity of presentation, we also assume Nl=NN_{l}=N for all ll, and thus N≫BN\gg B. DNN training has 33 stages in each iteration, (i) the feedforward stage, (ii) the backpropagation stage, and (iii) the update stage. As the operations are similar and repeat across all the layers (see Appendix A), we limit our discussion to layer ll. The operations during feedforward stage (see Fig. 3a) can be summarized as:

[[Step O1]O1] Compute matrix-matrix product Sl=WlXl\bm{S}^{l}=\bm{W}^{l}\bm{X}^{l} where Xl\bm{X}^{l} is a matrix of dimension N×BN\times B.

[[Step C1]C1] Compute X(l+1)=f(Sl)\bm{X}^{(l+1)}=f(\bm{S}^{l}) where f(⋅)f(\cdot) is a nonlinear activation function applied element-wise.

At the last layer (l=Ll=L), the backpropagated error matrix is generated by accessing the true label matrix from memory and the estimated label matrix as output of last layer (see Fig. 3b). Then, the backpropagated error propagates from layer LL to 11 (see Fig. 3c), also updating the weight matrices at every layer alongside (see Fig. 3d). The operations for the backpropagation stage can be summarized as: From layer l=Ll=L to 11,

[[Step O2]O2] Compute matrix-matrix product [Cl]T=[Δl]TWl[\bm{C}^{l}]^{T}=[\bm{\Delta}^{l}]^{T}\bm{W}^{l} where [Δl]T[\bm{\Delta}^{l}]^{T} is a matrix of dimensions B×NB\times N.

[[Step C2]C2] Compute Hadamard product [Δl−1]T=[Cl]T∘g([Xl]T)[\bm{\Delta}^{l-1}]^{T}=[\bm{C}^{l}]^{T}\circ g([\bm{X}^{l}]^{T}) where g(⋅)g(\cdot) is another function applied element-wise (more specifically df(u)du=g(f(u))\frac{df(u)}{du}=g(f(u)) for the chosen nonlinear activation function f(u)f(u)) and the Hadamard product “∘\circ” between two matrices of the same dimensions is another matrix of those dimensions, such that its elements are element-wise products of the corresponding elements of the operands.

Finally, the step in the Update stage is as follows: For all layers ll,

[[Step O3]O3] Update matrix Wl\bm{W}^{l} as follows: Wl←Wl+ηΔl[Xl]T\bm{W}^{l}\leftarrow\bm{W}^{l}+\eta\bm{\Delta}^{l}[\bm{X}^{l}]^{T} where η\eta is the learning rate. Sometimes a regularization term is added with the loss function in DNN training (elaborated in Section A-C). For L2 regularization, the update rule is modified as: Wl←(1−ηλ)Wl+ηδl[xl]T\bm{W}^{l}\leftarrow(1-\eta\lambda)\bm{W}^{l}+\eta\bm{\delta}^{l}[\bm{x}^{l}]^{T} where η\eta is the learning rate and λ2\frac{\lambda}{2} is the regularization constant.

System model: We assume that there is a decentralized system of PP memory-constrained nodes that are prone to soft-errors during computation. Soft-errors cause the node to produce entirely garbage outputs. We introduce two error-models in Section V. Under Error-Model 11, which is an adversarial model, soft-errors only occur during the most computationally intensive operations, i.e., steps O1O1, O2O2 and O3O3 but the number of erroneous nodes is bounded. Under Error-Model 22, which is a probabilistic model, soft-errors can occur during the steps O1O1, O2O2, O3O3, C1C1, C2C2 as well as encoding/decoding. There is no upper bound on the number of erroneous nodes or any specific assumption on where they can occur, but when errors occur, the output of that node is assumed to have an additive continuous-valued random noise. We elaborate upon these two models in Section V.

There is no single reliable master node, and all nodes can be unreliable. After an initial error-free setup (a pre-processing stepWe assume that the cost of initial setup or pre-processing before the start of training is amortized across large number of iterations.), all the nodes may begin their computational tasks in parallel, and proceed with the iterations of training. These memory-constrained and unreliable nodes may also communicate with each other, as required by the DNN training algorithm, or even perform some functions of a master node, such as, gathering outputs from other nodes, encoding/decoding etc., while respecting their storage constraints. After completing the required number of iterations, the final computational results (in this case, the LL trained parameter matrices) remain stored locally across the multiple nodes in a distributed manner. The data set and their labels are stored in a separate reliable memory unit and are communicated to all the nodes at each iteration when they access it.

Problem formulation: Design an error-resilient DNN training strategy using PP nodes, such that:

Each node can store only a 1K\frac{1}{K} fraction of each weight matrix Wl\bm{W}^{l} for each layer. Thus, for each layer, there is a per-node storage of N2K+o(N2K)\frac{N^{2}}{K}+o\left(\frac{N^{2}}{K}\right) where the small additional storage of o(N2K)o\left(\frac{N^{2}}{K}\right) is to store additional quantities that are negligible in storage size as compared to the 1K\frac{1}{K} fraction of Wl\bm{W}^{l}, e.g., matrices Sl\bm{S}^{l}, Xl\bm{X}^{l}, Cl\bm{C}^{l} and Δl\bm{\Delta}^{l} which are all of dimensions N×BN\times B where B≪NB\ll N by our assumption. In particular, we assume that B=o(NK)B=o\left(\frac{N}{K}\right) to satisfy this storage constraint.

All additional overheads per node including the communication complexity as well as the computational complexity of encoding/decoding in an error-free iterationFor error-resilience, some operations such as error detection, encoding etc. are required to be performed at each iteration even though most iterations of training are actually error-free. The purpose of this assumption is only to ensure that the additional overheads introduced for error-resilience in these error-free iterations is negligible. In the few iterations where errors occur and are detected, all error-resilient strategies incur some extra costs, such as, possibly regenerating the erroneous nodes or reverting to the last checkpoint etc, that is not being compared here. should be negligible in scaling sense as compared to the computational complexity of the local matrix multiplications and updates, i.e., steps O1O1, O2O2 and O3O3 at each node after parallelization.

Our goal is to achieve maximum error tolerance in the steps of DNN training, i.e., maximize the number of errors that can be corrected during training in a single iteration under both the error models. Under Error model 11 (see Section V), the number of erroneous nodes are bounded and we require the number of errors that can be corrected in any step to be higher than the maximum number of errors that can occur in that step in the worst-case. Under Error model 22 (see Section V), which allows for unbounded number of errors but their values being drawn from a continuous-valued distribution, in addition to maximizing the number of errors we can correct, our goal is also to be able to detect the occurrence of errors even when they are too many to be corrected.

For this problem formulation, we also assume that N≫PN\gg P, the number of parallel nodes.

In practice, the nodes also perform “checkpointing” at large intervals, i.e., storing the entire DNN (the LL parameter matrices) at a reliable disk from which the values can be retrieved when errors cannot be corrected. However, checkpointing is very expensive, even though it is assumed to be error-free, as the nodes have to access the disk, and thus can only be performed at large time intervals.

Existing coded computing techniques require encoding of the matrices being multiplied. If we were to extend them naively to the problem of coded DNN training, then the matrix Wl\bm{W}^{l} has to be encoded afresh in each iteration because the matrix Wl\bm{W}^{l} is updated in each iteration during step O3O3. Encoding matrix Wl\bm{W}^{l} with non-sparse codes in each iteration has a huge computational cost, and is thus a major challenge for the problem of coded DNN training (violates the last criterion in problem formulation). Thus, one key contribution in this work is in proposing a unified strategy such that the matrix Wl\bm{W}^{l}, once initially encoded, remains encoded during updates in each iteration, obviating the need to encode afresh.

Matrix Partitioning Notations: Throughout this paper, matrices and vectors are denoted in bold font. When we block-partition a matrix A\bm{A} both row-wise and column-wise into m×nm\times n equal-sized blocks for any integers mm and nn, we let Ai,j\bm{A}_{i,j} denote the block with row index ii and column index jj, where i=0,1,…,m−1i=0,1,\dots,m-1 and j=0,1,…n−1j=0,1,\dots n-1. Similarly, when we partition a vector a\bm{a} into mm equal parts for any integer mm, the sub-vectors are denoted as a0,a1,…,an−1\bm{a}_{0},\bm{a}_{1},\dots,\bm{a}_{n-1} respectively. E.g., for m=n=2m=n=2, the partitioning is as follows:

Also, note that when a matrix A\bm{A} is split only horizontally or only vertically into mm blocks, we denote the sub-matrices as Ai,:\bm{A}_{i,:} or A:,i\bm{A}_{:,i} respectively for i=0,1,…,m−1i=0,1,\ldots,m-1.

III Motivating Example for Coded Matrix Multiplication

In this section, we introduce a motivating example to first understand both Polynomial codes and MatDot codes , that have been proposed for distributed matrix-matrix multiplication and then, we introduce the key idea of garbage alignment in the PolyDot framework for matrix-matrix multiplication . We choose K=K′=4K=K^{\prime}=4.

The main intuition behind these polynomial-based strategies is to carefully design two polynomials W~(v)\widetilde{\bm{W}}(v) and X~(v)\widetilde{\bm{X}}(v) (may also be multivariate) whose coefficients are sub-matrices of W\bm{W} and X\bm{X} respectively, such that, different sub-matrices of the final result S(=WX)\bm{S}(=\bm{W}\bm{X}) show up as coefficients of the product W~(v)X~(v)\widetilde{\bm{W}}(v)\widetilde{\bm{X}}(v). The pp-th processing node stores a unique evaluation of W~(v)\widetilde{\bm{W}}(v) and X~(v)\widetilde{\bm{X}}(v) at v=bpv=b_{p} for p=0,1,…,P−1p=0,1,\ldots,P-1, and then computes the product W~(bp)X~(bp)\widetilde{\bm{W}}(b_{p})\widetilde{\bm{X}}(b_{p}), essentially producing a unique evaluation of the polynomial S~(v):=W~(v)X~(v)\widetilde{\bm{S}}(v):=\widetilde{\bm{W}}(v)\widetilde{\bm{X}}(v) at v=bpv=b_{p}. From a sufficient number of unique evaluations of the polynomial S~(v)\widetilde{\bm{S}}(v), the decoder is able to interpolate back all its coefficients, which include the different sub-matrices of the result S\bm{S}. Thus, the recovery threshold, i.e., the number of nodes to wait for out of PP is (1+Degree(S~(v)))(1+Degree(\widetilde{\bm{S}}(v))), which is the total number of unknowns in the interpolation of the polynomial S~(v)\widetilde{\bm{S}}(v).

Matrix Multiplication using Polynomial Codes: In Polynomial codes , the first matrix W\bm{W} is split horizontally and the second matrix X\bm{X} is split vertically into K(=4)K(=4) blocks each, as follows:

Note that, the resultant matrix S=WX\bm{S}=\bm{W}\bm{X} therefore takes the following form:

For this coding strategy, two polynomials are chosen as follows:

The pp-th processing node stores a unique evaluation of W~(v)\widetilde{\bm{W}}(v) and X~(v)\widetilde{\bm{X}}(v) at v=bpv=b_{p} for p=0,1,…,P−1p=0,1,\ldots,P-1, and then computes the product W~(bp)X~(bp)=S~(bp)\widetilde{\bm{W}}(b_{p})\widetilde{\bm{X}}(b_{p})=\widetilde{\bm{S}}(b_{p}), which essentially produces a unique evaluation of the polynomial S~(v)=W~(v)X~(v)\widetilde{\bm{S}}(v)=\widetilde{\bm{W}}(v)\widetilde{\bm{X}}(v) at v=bpv=b_{p}. Observe the coefficients of the polynomial S~(v)\widetilde{\bm{S}}(v) as follows:

Interestingly, all the different sub-matrices of S\bm{S}, i.e., Si,j\bm{S}_{i,j} show up as coefficients of the polynomial S~(v)\widetilde{\bm{S}}(v). Because the degree of this polynomial S~(v)\widetilde{\bm{S}}(v) is 1515 (for K=4K=4), the decoder requires 1616 unique evaluations to interpolate the polynomial successfully. Thus, recovery threshold is 1616.

Matrix Multiplication using MatDot Codes: Contrary to Polynomial codes, in MatDot codes the first matrix W\bm{W} is split vertically and the second matrix X\bm{X} is split horizontally into KK blocks each as follows:

Now we carefully choose the two polynomials W~(v)\widetilde{\bm{W}}(v) and X~(v)\widetilde{\bm{X}}(v) as follows:

As before, the pp-th node stores a unique evaluation of W~(v)\widetilde{\bm{W}}(v) and X~(v)\widetilde{\bm{X}}(v) at v=bpv=b_{p} for p=0,1,…,P−1p=0,1,\ldots,P-1, and then computes the product W~(bp)X~(bp)=S~(bp)\widetilde{\bm{W}}(b_{p})\widetilde{\bm{X}}(b_{p})=\widetilde{\bm{S}}(b_{p}), essentially producing a unique evaluation of the polynomial S~(v):=W~(v)X~(v)\widetilde{\bm{S}}(v):=\widetilde{\bm{W}}(v)\widetilde{\bm{X}}(v). Now observe the coefficients of the polynomial S~(v)\widetilde{\bm{S}}(v) as follows:

Note that, the coefficient of vK−1(=v3)v^{K-1}(=v^{3}) gives the result S\bm{S}. All the other coefficients are practically of no use, and hence, are referred to as “garbage.” Because the degree of the polynomial S(v)\bm{S}(v) is 66, the decoder would require only 77 unique evaluations to interpolate the polynomial successfully. Thus, the recovery threshold is 77.

Matrix Multiplication using Generalized PolyDot Codes (with Garbage Alignment): Before introducing Generalized PolyDot Codes, we review the PolyDot framework for the matrix-matrix multiplication problem. The PolyDot framework splits both the matrices horizontally and vertically into 44 sub-matrices each, as follows:

The resultant matrix S\bm{S} therefore takes the form:

Now the PolyDot framework (see Fig. 4) encodes these sub-matrices of W\bm{W} into a polynomial in two variables with each variable corresponding to either the row or column dimension, as follows:

The sub-matrices of X\bm{X} are encoded as follows:

Now, observe the coefficients of the product of the two polynomials, i.e., S~(u,v,w)=W~(u,v)X~(v,w)\widetilde{\bm{S}}(u,v,w)=\widetilde{\bm{W}}(u,v)\widetilde{\bm{X}}(v,w) as follows:

The coefficient of uivwku^{i}vw^{k} corresponds to Si,k\bm{S}_{i,k}. The total number of unknowns or coefficients in this multivariate polynomial is 2×3×2=122\times 3\times 2=12. If we convert this multivariate polynomial into a polynomial of a single variable, e.g., using substitution v=u2,w=u6v=u^{2},w=u^{6}, so that there exists a bijection between the coefficients of the multivariate polynomial and the polynomial of a single variable, then we would require 1212 unique evaluations of the polynomial to be able to interpolate all its 1212 unknown coefficients, including the ones that contribute towards S\bm{S}.

In this work, one of our key observations is that even though the polynomial has 1212 coefficients, the number of coefficients that are useful to us is only 44, i.e., only the coefficients of uivwku^{i}vw^{k} for i,k=0,1i,k=0,1 while the others are garbage. Thus, we instead propose the following variable substitution: u=v2,w=v4u=v^{2},w=v^{4} in our previously proposed PolyDot framework . Observe the product now:

Note that the coefficient of v2i+4k+1v^{2i+4k+1} in S~(v)\widetilde{\bm{S}}(v) correspond exactly to the coefficient of uivwku^{i}vw^{k} in S~(u,v,w)\widetilde{\bm{S}}(u,v,w), i.e., Si,k\bm{S}_{i,k}. However, some of the garbage coefficients have now aligned with each other to reduce the total number of unknown garbage coefficients, e.g., coefficient of both uv2uv^{2} and ww in S~(u,v,w)\widetilde{\bm{S}}(u,v,w) now get added up and form the coefficient of v4v^{4} in S~(v)\widetilde{\bm{S}}(v). This is the key idea of garbage alignment that lies at the core of the design of Generalized PolyDot codes, as we discuss in details in the next section. The polynomial resulting after substitution is only of degree 88. Thus it only requires 99 unique evaluations to be able to interpolate all its coefficients. Thus, garbage alignment reduces the recovery threshold from 1212 to 99.

IV Generalized PolyDot codes for coded matrix multiplication

In this section, we describe our Generalized PolyDot code construction for the problem of coded matrix multiplication WN1×N0XN0×B\bm{W}_{N_{1}\times N_{0}}\bm{X}_{N_{0}\times B} using PP memory-constrained nodes, as discussed in Section II. We choose integers mm, nn and dd such that mn=Kmn=K and nd=K′nd=K^{\prime} respectively. Next, we partition the matrix W\bm{W} both horizontally and vertically into an m×nm\times n grid of smaller sub-matrices of dimensions N1m×N0n\frac{N_{1}}{m}\times\frac{N_{0}}{n} each. Note that, in this type of partitioning, each sub-matrix contains a 1K\frac{1}{K} fraction of WN1×N0\bm{W}_{N_{1}\times N_{0}}. Similarly, X\bm{X} is also partitioned into an n×dn\times d grid of sub-matrices of dimensions N0n×Bd\frac{N_{0}}{n}\times\frac{B}{d} each (1K′\frac{1}{K^{\prime}} of X\bm{X}). After this, we encode these sub-matrices of W\bm{W} and X\bm{X} by taking appropriate linear combinations and store encoded sub-matrices of W\bm{W} and X\bm{X} at each node that satisfy the storage constraints.

Theorem 1 states our achievability result for the problem of matrix-vector products where we are required to perform s=Wx\bm{s}=\bm{W}\bm{x} using PP nodes, such that every node can only store an N1m×N0n\frac{N_{1}}{m}\times\frac{N_{0}}{n} sub-matrix (1K\frac{1}{K} fraction) of W\bm{W} and an N0n×1\frac{N_{0}}{n}\times 1 sub-vector of x\bm{x}.

Generalized PolyDot codes for computing matrix-vector multiplication WN1×N0xN0×1\bm{W}_{N_{1}\times N_{0}}\bm{x}_{N_{0}\times 1} using PP nodes, each storing an N1m×N0n\frac{N_{1}}{m}\times\frac{N_{0}}{n} sub-matrix of W\bm{W} and an N0n×1\frac{N_{0}}{n}\times 1 sub-vector of x\bm{x}, has a recovery threshold of mn+n−1mn+n-1. Thus, it can tolerate at most P−mn−n+1P-mn-n+1 erasures.

Recall the PolyDot framework for matrix multiplication. We first block-partition W\bm{W} into m×nm\times n sub-matrix where Wi,j\bm{W}_{i,j} denotes the sub-matrix at location (i,j)(i,j) for i=0,1,…,m−1i=0,1,\ldots,m-1 and j=0,1,…,n−1j=0,1,\ldots,n-1. Let the pp-th node (p=0,1,…,P−1p=0,1,\ldots,P-1) store an encoded sub-matrix of W\bm{W}, which is a polynomial in uu and vv, as follows:

evaluated at some (u,v)=(ap,bp)(u,v)=(a_{p},b_{p}). The choice of these variables will be clarified later. The vector x\bm{x} is also partitioned into nn equal sub-vector denoted by x0,x1,…,xn−1\bm{x}_{0},\bm{x}_{1},\ldots,\bm{x}_{n-1}, each of dimensions N0n×1\frac{N_{0}}{n}\times 1. Each node stores an encoded sub-vector as follows:

evaluated at v=bpv=b_{p}. Now, each node computes the smaller matrix-vector multiplication W~(ap,bp)x~(bp)\widetilde{\bm{W}}(a_{p},b_{p})\widetilde{\bm{x}}(b_{p}) which effectively results in the evaluation, at (u,v)=(ap,bp)(u,v)=(a_{p},b_{p}), of the following polynomial:

even though the node is not explicitly evaluating it from all its coefficients. Observe that the coefficient of uivn−1u^{i}v^{n-1} for i=0,1,…,m−1i=0,1,\dots,m-1 turns out to be ∑j=0n−1Wi,jxj=si\sum_{j=0}^{n-1}\bm{W}_{i,j}\bm{x}_{j}=\bm{s}_{i}. This is obtained by fixing j′=jj^{\prime}=j. Thus, these mm coefficients constitute the mm sub-vectors of s(=Wx)\bm{s}(=\bm{W}\bm{x}). Therefore, s\bm{s} can be recovered by the decoder if it can interpolate these mm coefficients of the polynomial s~(u,v)\widetilde{\bm{s}}(u,v) in 3 from its evaluations. As an example, consider the case where m=n=2m=n=2.

More generally, we use the substitution u=vnu=v^{n} to convert s~(u,v)\widetilde{\bm{s}}(u,v) into a polynomial in a single variable. Therefore, ap=bpna_{p}=b_{p}^{n} and each bpb_{p} is unique for p=0,1,…,P−1p=0,1,\ldots,P-1. Some of the unwanted, garbage coefficients align with each other (e.g. uu and v2v^{2} in 4), but the coefficients of uivn−1u^{i}v^{n-1}, i.e., si\bm{s}_{i} remain unchanged and now correspond to the coefficients of vni+n−1v^{ni+n-1} for i=0,1,…,m−1i=0,1,\ldots,m-1. After the substitution u=vnu=v^{n} in s~(u,v)\widetilde{\bm{s}}(u,v), the resulting polynomial of single variable vv is of degree mn+n−2mn+n-2. Thus, the decoder needs to wait for mn+n−1mn+n-1 nodes, each providing a unique evaluation, to be able to interpolate all the mn+n−1mn+n-1 unknown coefficients. The recovery threshold is thus mn+n−1mn+n-1. ∎

Now, we extend the coding strategy to the problem of matrix-matrix multiplication.

Generalized PolyDot codes for computing matrix-matrix multiplication WN1×N0XN0×B\bm{W}_{N_{1}\times N_{0}}\bm{X}_{N_{0}\times B} using PP nodes, each storing an N1m×N0n\frac{N_{1}}{m}\times\frac{N_{0}}{n} sub-matrix of W\bm{W} and an N0n×Bd\frac{N_{0}}{n}\times\frac{B}{d} sub-matrix of X\bm{X} has a recovery threshold of mnd+n−1mnd+n-1. Thus, it can tolerate at most P−mnd−n+1P-mnd-n+1 erasures.

The matrix-matrix multiplication strategy is very similar to the matrix-vector case. The pp-th node (p=0,1,…,P−1p=0,1,\ldots,P-1) stores an encoded sub-matrix of W\bm{W} which is the same polynomial in uu and vv, as follows:

evaluated at (u,v)=(ap,bp)(u,v)=(a_{p},b_{p}). The matrix X\bm{X} is also block-partitioned into n×dn\times d sub-matrices, where the sub-matrix at location (j,k)(j,k) is denoted as Xj,k\bm{X}_{j,k}, for j=0,1,…,n−1j=0,1,\ldots,n-1 and k=0,1,…,d−1k=0,1,\ldots,d-1. As per the PolyDot framework , the pp-th node also stores an encoded sub-matrix of X\bm{X}, as a polynomial in (v,w)(v,w), as follows:

evaluated at (v,w)=(bp,cp)(v,w)=(b_{p},c_{p}). Next, each node computes the smaller matrix-matrix product: S~(ap,bp,cp)=W~(ap,bp)X~(bp,cp)\widetilde{\bm{S}}(a_{p},b_{p},c_{p})=\widetilde{\bm{W}}(a_{p},b_{p})\widetilde{\bm{X}}(b_{p},c_{p}) which effectively results in the evaluation, at (u,v,w)=(ap,bp,cp)(u,v,w)=(a_{p},b_{p},c_{p}), of the polynomial:

even though the node is not explicitly evaluating it from its coefficients. Now, fixing j′=jj^{\prime}=j, we observe that the coefficient of uivn−1wku^{i}v^{n-1}w^{k} for i=0,1,…,m−1i=0,1,\dots,m-1 and k=0,1,…,d−1k=0,1,\dots,d-1 turns out to be ∑j=0n−1Wi,jXj,k=Si,k\sum_{j=0}^{n-1}\bm{W}_{i,j}\bm{X}_{j,k}=\bm{S}_{i,k}. These mdmd coefficients constitute the m×dm\times d sub-matrices (or blocks) of S=WX\bm{S}=\bm{W}\bm{X}. Therefore, S\bm{S} can be recovered at the decoder if all these mdmd coefficients of the polynomial S~(u,v,w)\widetilde{\bm{S}}(u,v,w) can be interpolated from its evaluations at different nodes.

For garbage alignment, we propose the substitutions (u=vn,w=vmn)(u=v^{n},w=v^{mn}) to convert S~(u,v,w)\widetilde{\bm{S}}(u,v,w) into a polynomial of a single variable vv. Thus, ap=bpna_{p}=b_{p}^{n}, cp=bpmnc_{p}=b_{p}^{mn} and bpb_{p} is unique for p=0,1,…,P−1p=0,1,\ldots,P-1. The coefficient of uivn−1wku^{i}v^{n-1}w^{k} in S~(u,v,w)\widetilde{\bm{S}}(u,v,w) exactly correspond to the coefficient of vni+mnk+n−1v^{ni+mnk+n-1} in S~(v)=S~(u,v,w)∣u=vn,w=vmn\widetilde{\bm{S}}(v)=\widetilde{\bm{S}}(u,v,w)|_{u=v^{n},w=v^{m}n}, while some of the garbage terms align with each other, reducing the total number of unknowns. Observe the polynomial:

which is a polynomial in a single variable vv. Its degree is given by n(m−1)+mn(d−1)+n−1+n−1=mnd+n−2n(m-1)+mn(d-1)+n-1+n-1=mnd+n-2. Thus, the decoder needs to wait for mnd+n−1mnd+n-1 nodes, each producing a unique evaluation, to be able to interpolate all its mnd+n−1mnd+n-1 coefficients. ∎

Under erasures, the master node only waits for mnd+n−1mnd+n-1 nodes to finish and the decoding reduces to the problem of solving a linear system of equations (polynomial interpolation). If the outputs are corrupted by errors instead of erasures, the master node gathers outputs from all the PP nodes and then solves a sparse reconstruction problem to decode the correct output, using the techniques in , as we also discuss in Appendix B.

Now we discuss the communication and computation costs of Generalized PolyDot Codes (with mn=Kmn=K, nd=K′nd=K^{\prime}) in the centralized setup under erasures.

Computational complexity of encoding the sub-matrices at the master node: O(P(N1N0+N0B))\mathcal{O}(P(N_{1}N_{0}+N_{0}B)).

Total communication complexity of sending different encoded sub-matrices from the master node to each of the PP nodes: Θ(PN1N0K+PN0BK′)\Theta(\frac{PN_{1}N_{0}}{K}+\frac{PN_{0}B}{K^{\prime}}).

Computational complexity at each worker node for the matrix multiplication: Θ(N1N0Bmnd)\Theta(\frac{N_{1}N_{0}B}{mnd}).

Total communication complexity of gathering different outputs at the master node from the first mnd+n−1mnd+n-1 workers: Θ((mnd+n−1)N1Bmd)\Theta((mnd+n-1)\frac{N_{1}B}{md}).

Computational complexity of decoding at master node: O((mnd+n−1)3N1Bmd)\mathcal{O}((mnd+n-1)^{3}\frac{N_{1}B}{md}).

Tradeoff between communication cost and recovery threshold: In Fig. 5 we illustrate the tradeoff between recovery threshold and communication costs for Generalized PolyDot codes by varying mm, nn and dd. When we choose n=1,m=K,d=K′n=1,m=K,d=K^{\prime}, the Generalized PolyDot codes reduce to Polynomial codes with recovery threshold KK′KK^{\prime}. On the other hand, in the regime where K=K′K=K^{\prime}, the Generalized PolyDot codes reduce to MatDot codes when we choose m=d=1,n=Km=d=1,n=K, resulting in a recovery threshold of 2K−12K-1.

Now, we move on to our Problem Formulation 22, i.e., coded DNNs.

V Modelling Assumptions for coded DNNs with a Result on Real Number Error Correction

In this section, we will elaborate upon a few modeling assumptions for coded DNNs, such as, defining the two error models and communication complexity in decentralized settings. The error models introduced here lead to an interesting theoretical result on real number error correction, as stated in Theorem 3.

Let q\bm{q} be a Q×1Q\times 1 vector consisting of QQ real-valued symbols. The received output vector is as follows:

The subset A\mathcal{A} satisfies ∣A∣≤⌊P−Q2⌋|\mathcal{A}|\leq\lfloor\frac{P-Q}{2}\rfloor, with no specific assumptions on the locations or values of the errors and they may be chosen advarsarially.

Under the Probabilistic Error Model for channel coding, the decoder of a (P,Q)(P,Q) MDS Code can perform the following:

It can detect the occurrence of errors with probability 11, irrespective of the number of errors that occurred.

If the number of errors that occurred is less than or equal to P−Q−1P-Q-1, then all those errors can be corrected with probability 11, even without knowing in advance that how many errors actually occurred.

If the number of errors that occurred is more than P−Q−1P-Q-1, then the decoder is able to determine that the errors are too many to be corrected and declare a “decoding failure” with probability 11.

This result is interesting as it essentially means that in real number error correction, one can theoretically correct P−Q−1P-Q-1 errors with probability 11 which is more than the well-known adversarial error tolerance of ⌊P−Q2⌋\lfloor\frac{P-Q}{2}\rfloor for MDS coding. A detailed proof of this result is provided in Appendix B. Here we provide the main intuition.

Let us first consider the simple case of replication. Replication is essentially a (P,1)(P,1) MDS Code as one single real-valued symbol is replicated PP times. Under an adversarial error model, one could use a majority voting among the PP received values and thus correct upto ⌊P−12⌋\lfloor\frac{P-1}{2}\rfloor errors. However, under the probabilistic error model, the probability that two replicas are affected by the same value of error is . Thus, as long as not all the PP received values are equal, one can detect that errors have occurred. Moreover, if at least two received values match out of PP, it is most likely the original symbol unaffected by errors. Thus, one can correct P−2P-2 errors with probability 11. If no symbols match, then the decoder is able to declare a decoding failure.

This idea also extends to any (P,Q)(P,Q) MDS Code. For a (P,Q)(P,Q) MDS Code, the minimum Hamming distance between two codewords is dmin=P−Q+1d_{min}=P-Q+1. Thus, if the number of errors are within ⌊dmin−12⌋=⌊P−Q2⌋\lfloor\frac{d_{min}-1}{2}\rfloor=\lfloor\frac{P-Q}{2}\rfloor, the received output vector lies within a Hamming ball of radius ⌊dmin−12⌋\lfloor\frac{d_{min}-1}{2}\rfloor around the original codeword, and is thus closest in Hamming distance to the original codeword as compared to any other codeword. If one allows for more than ⌊dmin−12⌋\lfloor\frac{d_{min}-1}{2}\rfloor errors, the received output might fall within the ⌊dmin−12⌋\lfloor\frac{d_{min}-1}{2}\rfloor Hamming ball of another codeword, and hence may be decoded incorrectly. However, what Theorem 3 says is that if the error values are not adversarially chosen but allowed to be probabilistic, then even if we go slightly beyond ⌊P−Q2⌋\lfloor\frac{P-Q}{2}\rfloor, i.e., upto a Hamming radius of P−Q−1P-Q-1, the probability of the received output being closer to a different codeword is . In other words, for a given codeword, the set of all possible outputs that are closer to other codewords in a Hamming sense has probability measure in the space of all possible values that the output can take, i.e., the Hamming ball of radius P−Q−1P-Q-1.

Let H\bm{H} be the (P−Q)×P(P-Q)\times P sized parity check matrix of the MDS code, such that HGT=0\bm{H}\bm{G}^{T}=\bm{0}. We first propose the following decoding algorithm (see Algorithm 1) to produce e^\widehat{\bm{e}}, as an estimate of e\bm{e}, for a given channel output z\bm{z}. If an e^\widehat{\bm{e}} is obtained, the decoder can uniquely solve for q^\widehat{\bm{q}} from the linear set of equations GTq^=z−e^\bm{G}^{T}\widehat{\bm{q}}=\bm{z}-\widehat{\bm{e}}.

Note that, for He=He′\bm{H}\bm{e}=\bm{H}\bm{e}^{\prime} to hold, e=e′+h,\bm{e}=\bm{e}^{\prime}+\bm{h}, where h∈Null(H)\{0}\bm{h}\in Null(\bm{H})\backslash\{\bm{0}\}. Thus ∣A′∪A∣−∣A′∣|\mathcal{A}^{\prime}\cup\mathcal{A}|-|\mathcal{A}^{\prime}| indices of h\bm{h}, indexed in the set (A′∪A)\A′(\mathcal{A}^{\prime}\cup\mathcal{A})\backslash\mathcal{A}^{\prime} match exactly with e\bm{e} (see Fig. 6).

The key intuition behind this proof is that once a certain number of elements of a vector h\bm{h} are allowed to be chosen from iid Gaussian distributions, the probability of the vector h\bm{h} still lying in Null(H)Null(\bm{H}) becomes . To understand this better, observe that,

While we theoretically show that real number MDS coding can correct P−Q−1P-Q-1 errors, our proposed decoding algorithm requires sparse reconstruction for undetermined systems which is NP Hard . For practical purposes, one might consider using an L1-norm relaxation or other kinds of polynomial-time sparse reconstruction algorithms proposed in the compressed sensing literature , which are known to be reasonably accurate under various restrictions on matrix H\bm{H}.

Let us now understand what these models mean in the context of coded DNNs.

V-B Error Models for the coded DNN problem.

Recall from our problem formulation (Section II) that we are interested in model-parallel architectures that parallelize each layer across PP error-prone nodes (that can be reused across layers) because the nodes cannot locally store the entire matrix Wl\bm{W}^{l}. As the steps O1O1, O2O2 and O3O3 are the most computationally intensive (Θ(N2B)\Theta(N^{2}B)) steps at each layer, we restrict ourselves to schemes where these three steps for each layer are parallelized across the PP nodesThe steps C1C1 and C2C2 are lower in computational complexity, i.e., Θ(NB)\Theta(NB) which is lower in scaling sense as compared to Θ(N2B)\Theta(N^{2}B) and hence may or may not be parallelized across multiple nodes.. In such schemes, communication will be required after steps O1O1 and O2O2 as the partial computation outputs of steps O1O1 and O2O2 at one layer might be required at another node to compute the input X(l+1)\bm{X}^{(l+1)} or backpropagated error Δ(l−1)\bm{\Delta}^{(l-1)} for another layerNote that there could be alternate parallelization schemes where the different layers of the network are parallelized across different nodes instead of each layer being parallelized across all nodes. These schemes might have lower communication but some of the nodes stay idle and under-utilized, i.e., while computations are being performed in one layer, the nodes containing the other layers stay idle. It will be an interesting future work to explore the computation-communication tradeoffs among these alternate parallelization schemes.. We define Error Models 11 and 22, which are essentially realizations of the probabilistic and adversarial models for the coded DNN problem, with some additional assumptions.

Any node can have soft-errors but only during the steps O1O1, O2O2 and O3O3, which are the most computationally intensive operations in DNN training. Encoding, error-detection, decoding, nonlinear activation and Hadamard product are assumed to be error-free. These operations require negligible time and number of operationsThe shorter the computation, the lower is the probability of soft-errors. E.g., a Poisson process of soft errors makes the number of soft-errors have mean proportional to the interval length. because most of the time and resources are spent on steps O1O1, O2O2 and O3O3. There is no specific assumption on the locations of the erroneous nodes and they may be adversarial. However, the total number of erroneous nodes at any layer during O1O1, O2O2 and O3O3 are known to be bounded by t1t_{1}, t2t_{2} and t3t_{3} respectively. There is also no assumption on the distribution of the errors for this model, but all the output values of an erroneous node, e.g., all the values of the output matrix or vector are affected by errors.

Any node can have soft-errors during any primary operation such as encoding, decoding, nonlinear activation, Hadamard product as well as steps O1O1, O2O2 and O3O3, and there is no bound on the number of errors. The locations of the erroneous nodes can still be adversarial. The entire output of an erroneous node (all the values of the output matrix or vector) is assumed to be corrupted by additive iid Gaussian noise. However, under Error Model 22, we use verification steps to check for decoding errors that have very low complexity (compared to the primary steps), and hence those verification steps are assumed to be error-free.

Error Model 11 is a “worst-case” abstraction, which is useful when it is difficult to place probabilistic priors on errors. It essentially means that the number of errors that can occur in the longer steps is bounded. Hence, if we choose a strategy with higher error tolerance, we can correct all errors. On the other hand, Error Model 22 allows for errors in all primary operations and also has no upper bound on the number of errors. However, it makes one simplifying assumption. Specifically, the continuous distribution of noise simplifies our analyses by avoiding complicated probability distributions that arise in finite number of bits representations. We acknowledge that this simplification can lead to optimistic conclusions, e.g., it allows us to correct more errors than the adversarial model (see Theorem 3) with probability 11 and also detect the occurrence of errors (“garbage outputs”) with probability 11 (because the noise takes any specific value with probability zero; see Appendix B). This model is only accurate in the limit of large number of bits of precision. In practical implementations, our probability 11 results should be interpreted as holding with high probability (e.g. it is unlikely, but possible, that two erroneous nodes produce the exact same garbage output). Note that because both replication and coding can exploit Error Model 22 for error-correction and detection, it does not bias our results towards coding relative to replication.

V-C Error Tolerance Goals for the coded DNN Strategy.

We would like to be able to correct as many erroneous nodes as possible after the steps O1O1 and O2O2, because outputs are communicated to other nodes after these two steps.

Under Error Models 11, for any layer ll, the error tolerances are (tf,tb)(t_{f},t_{b}) if tft_{f} and tbt_{b} erroneous node outputs can be detected and corrected in the worst case immediately after steps O1O_{1} and after step O2O_{2} respectively. Similarly, under Error Model 22, the error tolerances are (tf,tb)(t_{f},t_{b}) if tft_{f} and tbt_{b} erroneous node outputs can be detected and corrected with probability 11 immediately after steps O1O_{1} and after step O2O_{2} respectively.

Our goal is to maximize the values of these error tolerances under both the error models.

For any coding strategy, the achievable tft_{f} and tbt_{b}’s depend on the number of nodes available (PP) and other parameters of the coding strategy, e.g., m,nm,n etc. as derived in Section VIII. Note that, after steps O1O1 and O2O2, we do not necessarily correct only the errors that occur during those steps. Depending on the coding strategy used, errors occurring in other steps could also get corrected after either O1O1 or O2O2. Under Error model 11, the values of the achievable tft_{f} and tbt_{b}’s will be required to be greater than appropriate functions of t1,t2t_{1},t_{2} and t3t_{3} based on the coding strategy being used, as we also elaborate in Section VIII.

V-D Communication Complexity.

In this work, we use standard definition of communication complexity for fully distributed and decentralized architectures, as mentioned in .

The communication cost of sending a message of NN items between two nodes will be modeled by α+βN\alpha+\beta N, in the absence of network conflicts. Here α\alpha and β\beta are two constants representing the message startup time and per data item transmission time respectively.

We assume that each node can communicate simultaneously to at most a constant number (say 22) of nodes. Thus, when one node has to broadcast the same NN values to PP other nodes, it usually initiates communication link (or startup) with the other nodes in the form of a spanning-tree in log⁡P\log{P} rounds and then starts communicating the NN values across this tree-like transmission network of the nodes. For more details, the reader is referred to . Following , the communication cost for this type of broadcast is given by: αlog⁡P+βN\alpha\log{P}+\beta N. Here, the first term arises because the communication link (or startup) between all the nodes is set up in Θ(log⁡2P)\Theta(\log_{2}{P}) rounds, and then the second term denotes the cost of sending the NN values across this tree-like network of nodes. When multiple nodes have to communicate with each other, the communication cost can be efficiently managed using collective communication protocols, as suggested in . For instance, when all nodes send their own, unique message of NN values to all other nodes, a communication protocol called All-Gather is used whose communication cost is αlog⁡P+2βPN\alpha\log{P}+2\beta PN. We will use Broadcast and All-Gather protocols to prove our results on communication complexities in Appendix D.

VI Applying Existing Strategies to the coded DNN Problem

The problem formulation 22 stated in Section II discusses our goals. Essentially, for B=1B=1, we are required to design a coded DNN training strategy, that we denote as C(N,K,P)\mathcal{C}(N,K,P), which performs distributed “post” and “pre” multiplication of the same matrix W\bm{W} with vectors x\bm{x} and δT\bm{\delta}^{T} respectively at each layer and a distributed update (W+ηδxT\bm{W}+\eta\bm{\delta}\bm{x}^{T}), along with all the other operations using PP memory-constrained nodes.

After these computations, all the replicas computing the same sub-vector, i.e., say si\bm{s}_{i}, exchange their computational outputs for error detection and correction. Note that the communication cost for exchanging outputs is only a function of the number of nodes, PP, and does not depend on NN. This is because the nodes can just exchange a single value of their computation result among each other instead of the entire sub-vector. Thus, this communication cost is much lower as compared to the computational cost of matrix multiplication in the regime N≪PN\ll P.

Under Error Model 11, any t=⌊P−mn2mn⌋t=\lfloor\frac{P-mn}{2mn}\rfloor errors can be tolerated in the worst case. However under Error Model 22, the probability of two outputs having exactly same error is . As long as an output occurs at least twice, it is almost surely the correct output. Thus, any t=Pmn−2t=\frac{P}{mn}-2 errors can be detected and corrected. Then, the correct sub-vectors (si\bm{s}_{i}’s) are communicated to the respective nodes that require it for generating their input for the next layer, and the sub-matrices stored in the erroneous nodes are regenerated by accessing other nodes known to be correct.

Additional Steps: At regular intervals, the system also checkpoints, i.e., sends the entire DNN to a disk for storage. This disk-storage, although time-intensive to retrieve from, can be assumed to be error-free. Under Error Model 22, if more than tt errors occur, then with probability 11, none of the outputs match. The system detects the occurrence of errors even though it is unable to correct them. So, it retrieves the DNN from the disk and reverts the computation to the last checkpoint.

A similar technique is applied for backpropagation. The the node with index (i,j)(i,j) accesses δiT\bm{\delta}^{T}_{i} and computes δiTWi,j\bm{\delta}^{T}_{i}\bm{W}_{i,j}. Finally the last node in every column aggregates and computes ∑i=0m−1δiTWi,j=cjT\sum_{i=0}^{m-1}\bm{\delta}^{T}_{i}\bm{W}_{i,j}=\bm{c}^{T}_{j} for j=0,1,…,n−1j=0,1,\dots,n-1. Error check occurs similarly. If errors can be corrected, then cjT\bm{c}^{T}_{j}’s are communicated to the respective nodes that require it to compute backpropagated error for the next layer, along with xj\bm{x}_{j}. Interestingly, after these operations, the node with index (i,j)(i,j) has xj\bm{x}_{j} and δiT\bm{\delta}^{T}_{i}, and is thus able to update itself as Wi,j←Wi,j+ηδixjT\bm{W}_{i,j}\leftarrow\bm{W}_{i,j}+\eta\bm{\delta}_{i}\bm{x}^{T}_{j} respectively.

The error tolerances for the replication strategy are tf=tb=⌊P−mn2mn⌋t_{f}=t_{b}=\lfloor\frac{P-mn}{2mn}\rfloor under Error Model 11 and tf=tb=Pmn−2t_{f}=t_{b}=\frac{P}{mn}-2 under Error Model 22, assuming mnmn divides PP.

where Gr\bm{G}_{r} and Gc\bm{G}_{c} are the generator matrices of the two systematic MDS codes for the row and column dimensions, ⊗\otimes denotes the Kronecker product and IN/m\bm{I}_{N/m} and IN/n\bm{I}_{N/n} denote identity matrices of the corresponding dimensions.

In step O1O1, only PfP_{f} nodes corresponding to the (Pfn,m)(\frac{P_{f}}{n},m) code are active. Similarly, in step O2O2, only PbP_{b} nodes are active corresponding to the (Pbm,n)(\frac{P_{b}}{m},n) code. Errors that happen in the update step O3O3 corrupt the updated sub-matrices and are detected and corrected the next time those sub-matrices are used to produce an output to be sent to another node, which could be either after step O1O1 or step O2O2 of the next iteration at that layer. Thus, errors of step O3O3 are corrected either after step O1O1 or step O2O2 at that layer, in the next iteration.

Under Error Model 11, in the worst case, the strategy thus requires tf≥t1+t3t_{f}\geq t_{1}+t_{3} and tb≥t2+t3t_{b}\geq t_{2}+t_{3} to be able to detect and correct all the errors. Under Error Model 22, when the number of errors are greater than tft_{f} or tbt_{b}, they can only be detected with probability 11 but cannot be corrected, as elaborated in .

The error tolerances for the MDS-code-based strategy are as follows: tf=⌊Pf−mn2n⌋, tb=⌊Pb−mn2m⌋t_{f}=\lfloor\frac{P_{f}-mn}{2n}\rfloor,\ t_{b}=\lfloor\frac{P_{b}-mn}{2m}\rfloor under Error Model 11 and tf=Pf−mn−nn, tb=Pb−mn−mmt_{f}=\frac{P_{f}-mn-n}{n},\ t_{b}=\frac{P_{b}-mn-m}{m} under Error Model 22.

VII Our Proposed coded DNN Training Strategy for mini-batch B=1𝐵1B=1

In this section, we introduce our proposed unified coded DNN training strategy. We propose an initial encoding scheme for W\bm{W} at each layer such that the same encoding allows us to perform coded “post” and “pre” multiplication of W\bm{W} with vectors x\bm{x} and δT\bm{\delta}^{T} respectively at each layer in every iteration. The key idea is that we encode W\bm{W} only for the first iteration. For all subsequent iterations, we encode and decode vectors (hence complexity o(N2K)o(\frac{N^{2}}{K}) as we show in Theorem 5) instead of matrices. As we will show, the encoded weight matrix W\bm{W} is able to update itself, maintaining its coded structure at very low additional overhead.

Initial Encoding of W\bm{W} (Pre-processing Step): Every node stores an Nm×Nn\frac{N}{m}\times\frac{N}{n} sub-matrix (or block) of W\bm{W} encoded using Generalized PolyDot. Recall from 1) that,

For p=0,1,…,P−1p=0,1,\ldots,P-1, node pp stores W~p:=W~(u,v)∣u=ap,v=bp\widetilde{\bm{W}}_{p}:=\widetilde{\bm{W}}(u,v)|_{u=a_{p},v=b_{p}}, i.e., the evaluation of W~(u,v)\widetilde{\bm{W}}(u,v) at (u,v)=(ap,bp)(u,v)=(a_{p},b_{p}), at the beginning of the training. This coded sub-matrix has N2K\frac{N^{2}}{K} entries. Thus every node stores a sub-matrix W~p=W~(ap,bp)\widetilde{\bm{W}}_{p}=\widetilde{\bm{W}}(a_{p},b_{p}) at the beginning of the training which has N2K\frac{N^{2}}{K} entries.

Encoding of matrix W\bm{W} is done only before the first iteration.

Feedforward stage: Assume that the entire input x\bm{x} to the layer is made available at every node by the previous layer (this assumption is justified at the end of this paragraph). Also assume that the updated W~p\widetilde{\bm{W}}_{p} of the previous iteration is available at every node (this assumption will be justified when we show that the encoded sub-matrices of W\bm{W} are able to update themselves, preserving their coded structure).

For p=0,1,…,P−1p=0,1,\ldots,P-1, node pp first block-partitions x\bm{x} into nn equal parts, and encodes them using the polynomial:

For p=0,1,…,P−1p=0,1,\dots,P-1, the pp-th node evaluates the polynomial x~(v)\widetilde{\bm{x}}(v) at v=bpv=b_{p}, yielding x~p:=x~(bp)\widetilde{\bm{x}}_{p}:=\widetilde{\bm{x}}(b_{p}). E.g., for n=2n=2, x\bm{x} is encoded as x~(v)=x0v+x1\widetilde{\bm{x}}(v)=\bm{x}_{0}v+\bm{x}_{1}.

Next, each node computes the matrix-vector product: s~p:=W~px~p\widetilde{\bm{s}}_{p}:=\widetilde{\bm{W}}_{p}\widetilde{\bm{x}}_{p}. The computation of s~p\widetilde{\bm{s}}_{p} at node pp is equivalent to the evaluation, at (u,v)=(ap,bp)(u,v)=(a_{p},b_{p}), of the following polynomial:

even though the node is not explicitly evaluating it by accessing all its coefficients separately. Now, fixing j′=jj^{\prime}=j, observe that the coefficient of uivn−1u^{i}v^{n-1} for i=0,1,…,m−1i=0,1,\dots,m-1 is ∑j=0n−1Wi,jxj=si\sum_{j=0}^{n-1}\bm{W}_{i,j}\bm{x}_{j}=\bm{s}_{i}. Thus, these mm coefficients constitute the mm sub-vectors of s=Wx\bm{s}=\bm{W}\bm{x}. Therefore, s\bm{s} can be recovered at any node if it can reconstruct all the coefficients of the polynomial s~(u,v)\widetilde{\bm{s}}(u,v) in 12, or rather just these mm coefficients. The matrix-vector product computed at the pp-th node results in the evaluation of this polynomial at (u,v)=(ap,bp)(u,v)=(a_{p},b_{p}). Every node then sends this product to every other nodeRecall that there is no single master node, and every node replicates the functioalities of the decoder. Using efficient all-to-all communication protocols popular in parallel computing, the communication cost of all nodes broadcasting its own sub-vector s~p\widetilde{\bm{s}}_{p} of length Nm\frac{N}{m} to all other nodes has a communication cost of αlog⁡P+2βNmP=Θ(NmP)\alpha\log{P}+2\beta\frac{N}{m}P=\Theta(\frac{N}{m}P) using All-Gather protocol. We are currently examining strategies to reduce this cost further. where some of these products may be erroneous. Now, if every node can still reconstruct the coefficients of uivn−1u^{i}v^{n-1} from these evaluations, then it can successfully decode s0,s1,…,sm−1\bm{s}_{0},\bm{s}_{1},\ldots,\bm{s}_{m-1}.

We use one of the substitutions u=vnu=v^{n} or v=umv=u^{m} (elaborated in Appendix C), to convert s~(u,v)\widetilde{\bm{s}}(u,v) into a polynomial in a single variable and then use standard decoding techniques to interpolate the coefficients of a polynomial in one variable from its evaluations at PP arbitrary points when some evaluations have an additive error. Once s\bm{s} is decoded at each node, the nonlinear function f(⋅)f(\cdot) is applied element-wise to generate the input for the next layer. This also makes x\bm{x} available at every node at the start of the next feedforward layer, justifying our assumption.

Regeneration: If the number of errors are few (≤\leq error tolerance), the nodes are not only able to decode the vectors correctly but also locate which nodes were erroneous (see Appendix C). Thus, the encoded W\bm{W} stored at those nodes are regeneratedThe encoded matrix at any node is the evaluation of a polynomial whose coefficients correspond to the original sub-matrices Wi,j\bm{W}_{i,j}. Thus, the number of nodes required by an error-prone node is the degree of this polynomial +1+1. Substituting u=vnu=v^{n} (alternatively, v=umv=u^{m}), this degree is mn−1mn-1, and thus an error-prone node needs to access mnmn correct nodes to regenerate itself. by accessing some of the nodes that are known to be correct and the algorithm proceeds forward.

It might appear that regeneration violates the storage constraint of each node, i.e., a storage of only a 1K\frac{1}{K} fraction of matrix W\bm{W}. However, note that the coded sub-matrix can also be computed element-wise rather than all at once, while adhering to the storage constraint. E.g. if the stored sub-matrix W~(ap,bp)\widetilde{\bm{W}}(a_{p},b_{p}) of size Nm×Nn\frac{N}{m}\times\frac{N}{n} is required to be regenerated, the node can access only the first element (location (0,0)(0,0)) of the stored sub-matrix of any mnmn nodes to regenerate the element (0,0)(0,0) of W~(ap,bp)\widetilde{\bm{W}}(a_{p},b_{p}). Then it deletes all the gathered values, and moves on to location (0,1)(0,1) and so on. This process only requires an additional storage of o(N2K)o\left(\frac{N^{2}}{K}\right).

Additional Steps (Under Error Model 22): Under Error Model 22, errors can occur in all the primary steps, and are unbounded. Similar to replication and MDS-code-based strategy, the DNN is checkpointed at a disk at regular intervals. If there are more errors than the error tolerance after steps O1O1 or O2O2, the nodes are unable to decode correctly. However, as the error is assumed to be additive and drawn from real-valued, continuous distributions, the occurrence of errors is still detected with probability 11 even though they cannot be located or corrected, and thus the entire DNN can again be restored from the last checkpoint.

To allow for decoding errors, we need to include one more verification step. This step is similar to the replication strategy where all nodes exchange some particular values of the decoded vector, i.e., say any PP pre-decided values of the decoded vector s\bm{s} and compare (additional communication overhead of αlog⁡P+2βP2\alpha\log{P}+2\beta P^{2} and computation overhead of Θ(P2)\Theta(P^{2}) as discussed in Appendix D). Again, it is unlikely that two nodes will have the exact same decoding error. If there is a disagreement at one or more nodes during this process, we assume that there has been errors during the decoding, and the entire neural network is restored from the last checkpoint. Because the total complexity of this verification step is low in scaling sense compared to encoding/decoding or communication (because it does not depend on NN), we assume that it is error-free since the probability of soft-errors occurring within such a small duration is negligible as compared to other computations of longer duration.

Under Error Model 22, errors can also occur in the step C1C1 or during encoding. If an error occurs during step C1C1, the vector x\bm{x} for the next layer is corrupted, which ultimately corrupts the encoded sub-vector of that node, i.e., x~p\widetilde{\bm{x}}_{p} for the next layer. Errors during encoding will also corrupt this sub-vector x~p\widetilde{\bm{x}}_{p}. Then, the error propagates into s~p\widetilde{\bm{s}}_{p} during the matrix-vector product s~p=W~px~p\widetilde{\bm{s}}_{p}=\widetilde{\bm{W}}_{p}\widetilde{\bm{x}}_{p} at the next layer and is finally detected and if possible corrected after step O1O1 of the next layer in the same iteration, when every node attempts to decode s\bm{s} from all its received s~p\widetilde{\bm{s}}_{p} sub-vectors, some of which may be erroneous.

Backpropagation stage: The backpropagation stage is very similar to the feedforward stage. The backpropagated error (transpose) δT\bm{\delta}^{T} is available at every node. Each node partitions the row-vector δT\bm{\delta}^{T} into mm equal parts and encodes them using the polynomial:

For p=0,1,…,P−1p=0,1,\dots,P-1, the pp-th node evaluates δ~T(u)\widetilde{\bm{\delta}}^{T}(u) at u=apu=a_{p}, yielding δ~pT:=δ~T(ap)\widetilde{\bm{\delta}}^{T}_{p}:=\widetilde{\bm{\delta}}^{T}(a_{p}). Next, it performs the computation c~pT:=δ~pTW~p\widetilde{\bm{c}}^{T}_{p}:=\widetilde{\bm{\delta}}^{T}_{p}\widetilde{\bm{W}}_{p} and sends the product to all other nodes, of which some products may be erroneous. Consider the polynomial:

The products computed at each node result in the evaluations of this polynomial c~T(u,v)\widetilde{\bm{c}}^{T}(u,v) at (u,v)=(ap,bp)(u,v)=(a_{p},b_{p}). Similar to feedforward stage, each node then decodes the coefficients of um−1vju^{m-1}v^{j} in the polynomial for j=0,1,…,n−1j=0,1,\dots,n-1, and thus reconstructs nn sub-vectors forming cT\bm{c}^{T}.

Additional Steps (Under Error Model 22): The additional steps of checkpointing and verification step to check for decoding errors is similar. Errors during step C2C2 or during encoding will corrupt the encoded sub-vector δ~pT\widetilde{\bm{\delta}}^{T}_{p}, which will eventually show up after the computation c~pT=δ~pTW~p\widetilde{\bm{c}}^{T}_{p}=\widetilde{\bm{\delta}}^{T}_{p}\widetilde{\bm{W}}_{p}, and will be corrected after step O2O2 when each node attempts to reconstruct cT\bm{c}^{T} from the outputs c~pT\widetilde{\bm{c}}^{T}_{p} of all the nodes, of which some may be erroneous.

Update stage: The key part is updating the coded W~p\widetilde{\bm{W}}_{p}. Observe that since x\bm{x} and δ\bm{\delta} are both available at each node, it can encode the vectors as ∑i=0m−1δiui\sum_{i=0}^{m-1}\bm{\delta}_{i}u^{i} and ∑j=0n−1xjvj\sum_{j=0}^{n-1}\bm{x}_{j}v^{j} at u=apu=a_{p} and v=bpv=b_{p} respectively, and then update itself as follows:

Thus, the update step preserves the coded nature of the weight matrix, with negligible additional overhead (see Theorem 5).

Update with Regularization: When training is performed using L2 regularization, recall that the original update rule is modified as: W←(1−ηλ)W+δxT\bm{W}\leftarrow(1-\eta\lambda)\bm{W}+\bm{\delta}\bm{x}^{T}. Adding the weight decay term only requires a minor change in the update step of coded DNN training. Now, the update equation in (14) can be modified as:

We do not need additional encoding or decoding for introducing the weight decay term. Instead, we only have to shrink the already encoded W\bm{W} matrix, W~p\widetilde{\bm{W}}_{p}, by (1−ηλ)(1-\eta\lambda) after each iteration. For a detailed discussion on regularization in DNN training, the reader is referred to Section A-C.

Errors during Update stage: Errors can occur in the update stage under both the error models. These errors corrupt the updated sub-matrix W~p\widetilde{\bm{W}}_{p} and then show up in the computation s~p=W~px~p\widetilde{\bm{s}}_{p}=\widetilde{\bm{W}}_{p}\widetilde{\bm{x}}_{p} at the same layer in the next iteration. So, these errors are finally detected and if possible corrected after step O1O1 at that layer in the next iteration when every node attempts to decode s\bm{s} from all its received s~p\widetilde{\bm{s}}_{p} sub-vectors, which may be erroneous.

VIII Results on Performance of our Proposed Strategy

We ignore integer effects here as we are primarily interested in error tolerance with scaling PP. However, strictly speaking, we need a floor function ⌊⋅⌋\lfloor\cdot\rfloor applied to all of the expressions for Error Model 11 above, and mnmn to divide PP for replication for both the error models.

Note that, for the proposed coded DNN training strategy, the errors in the update stage (step O3O3) are also corrected after step O1O1 at that layer, in the next iteration along with the other errors during step O1O1. Thus, under Error Model 11, we require tft_{f} to be greater than t1+t3t_{1}+t_{3} while tbt_{b} is only required to be greater than t2t_{2}. Because the computational complexities of steps O1,O2O1,O2 and O3O3 are similar, one might expect that t1,t2t_{1},t_{2} and t3t_{3} are nearly the same. Under Error Model 22, when the errors are more than tft_{f} or tbt_{b}, the strategy is able to detect errors with probability 11 but not correct them. However, even under Error Model 22, it is more desirable to have tf>tbt_{f}>t_{b} since more errors are likely to occur after step O1O1 as compared to step O2O2 since the total duration of computation that is covered is more.

The proof is provided in Appendix D. For this proof, we assume a pessimistic bound (Θ(P3)\Theta(P^{3})) on the decoding of a code of block length PP under errors, based on sparse reconstruction algorithms. We are currently examining the reduction of this complexity using other algorithms, which would also relax the condition of Theorem 5.

Note that we have used pessimistic bounds to characterize the communication and computational overheads of our strategy and in spite of that, we are able to show that the additional overheads are negligible in scaling sense as compared to the computational complexities of steps O1O1, O2O2 and O3O3 respectively. We are currently exploring the reduction of these additional overheads further using efficient implementation strategies.

IX Extension to mini-batch size B>1𝐵1B>1: coded matrix-matrix multiplication

So far, we have discussed DNN training using SGD for the special case of mini-batch size B=1B=1, i.e., when the DNN accesses only a single data point at a time. In this section, we extend our proposed strategy to the general case of mini-batch size B>1B>1, which leads to coded matrix-matrix multiplication.

Goal: The goal in this section is to design a coded DNN training strategy for mini-batch size B>1B>1, denoted by C(N,K,P,B)\mathcal{C}(N,K,P,B), using PP nodes such that every node can store only a 1K\frac{1}{K} fraction of the entries of W\bm{W} for each layer. We assume that B=o(NK)B=o(\frac{N}{K}) to ensure that the sizes of X\bm{X}, S(=WX)\bm{S}(=\bm{W}\bm{X}), ΔT\bm{\Delta}^{T} and CT(=ΔTW)\bm{C}^{T}(=\bm{\Delta}^{T}\bm{W}) are smaller in scaling sense than the size of 1K\frac{1}{K} fraction of W\bm{W}. Thus, every node has a total storage of LN2K+o(LN2K)\frac{LN^{2}}{K}+o(\frac{LN^{2}}{K}) where the small additional storage of o(LN2K)o(\frac{LN^{2}}{K}) is for storing X\bm{X}, S(=WX)\bm{S}(=\bm{W}\bm{X}), ΔT\bm{\Delta}^{T} and CT(=ΔTW)\bm{C}^{T}(=\bm{\Delta}^{T}\bm{W}) respectively for every layerIf we do not assume an upper bound on BB, then as BB increases, the allowed total storage per node would also be required to increase. Then, it may become possible to store more than 1K\frac{1}{K} fraction of W\bm{W} leading to alternative strategies altogether. This problem may be considered as a future work.. Similar to the coded DNN strategy for B=1B=1, the additional computation and communication complexities including encoding/decoding overheads in each iteration should be negligible in scaling sense as compared to the local computational complexity of the steps O1O1, O2O2 and O3O3 parallelized across each node, at any layer.

Thus, essentially we are required to perform distributed “post” and “pre” multiplication of the same matrix W\bm{W} of dimensions N×NN\times N with matrices X\bm{X} and ΔT\bm{\Delta}^{T} respectively and distributed update W+ηXΔT\bm{W}+\eta\bm{X}\bm{\Delta}^{T}, along with all the other operations of training. Similar to the case of B=1B=1, because outputs are only communicated to other nodes after steps O1O1 and O2O2 respectively, we aim to correct as many erroneous nodes as possible after these two steps, before moving to another layer.

The assumption that B=o(NK)B=o(\frac{N}{K}) is required only to satisfy the storage constraints. Irrespective of whether B=o(NK)B=o(\frac{N}{K}) or B=Ω(NK)B=\Omega(\frac{N}{K}), when considering the total computational or communication complexity, the matrix-matrix products and updates are still the most significant costs, and all other additional overheads including communication, encoding/decoding etc. add negligible overhead, as we will also formally show in Theorem 7, provided that the number of processing nodes are not too large, i.e., P4=o(N)P^{4}=o(N).

We modify the replication strategy of Section VI for B>1B>1. The matrix W\bm{W} is block-partitioned and stored in a manner similar to that for mini-batch B=1B=1 with Pmn\frac{P}{mn} replicas. However, instead of a single vector x\bm{x}, now the entire data matrix X\bm{X} is divided vertically into nn blocks, each of size Nn×B\frac{N}{n}\times B. Similarly, instead of vector δT\bm{\delta}^{T}, now the matrix ΔT\bm{\Delta}^{T} is divided into mm equal parts horizontally, each of size B×NmB\times\frac{N}{m}. All the operations that were performed on the ii-th part of x\bm{x} or δT\bm{\delta}^{T}, are now performed on the corresponding block of X\bm{X} or ΔT\bm{\Delta}^{T} respectively. The additional computational and communication complexity of this strategy is obviously higher than the case with B=1B=1, but they still satisfy the desired constraint that they should be smaller in scaling sense than the per-node computational complexity of the steps O1O1, O2O2 and O3O3.

The proposed MDS-code-based strategy can also be modified in a similar manner as replication. The matrix W\bm{W} is encoded and stored in a manner similar to that for the case of mini-batch size B=1B=1. However, the matrices X\bm{X} and ΔT\bm{\Delta}^{T} are now divided into nn and mm parts respectively, similar to replication, and the same operations are performed as in the case of mini-batch size B=1B=1. The additional computational and communication complexity of this strategy is still smaller in scaling sense than the per-node computational complexity of O1O1, O2O2 and O3O3.

IX-B Our Proposed coded DNN Training Strategy for mini-batch size B>1𝐵1B>1.

Here, we apply Generalized PolyDot codes for the matrix-matrix multiplication, as discussed previously in Theorem 2. Let us begin by reviewing the notations for matrix partitioning.

The weight matrix W\bm{W} is partitioned into a grid of m×nm\times n blocks as before. The main difference from the case of B=1B=1 is that the input to the layer X\bm{X} and backpropagated error Δ\bm{\Delta} are matrices now. The matrix X\bm{X} is also block partitioned into a grid of n×d1n\times d_{1} blocks, each sub-matrix (or block) being of size Nn×Bd1\frac{N}{n}\times\frac{B}{d_{1}}. This results in the matrix S=WX\bm{S}=\bm{W}\bm{X} to be partitioned into a grid of m×d1m\times d_{1} blocks, with each sub-matrix (or block) being of size Nm×Bd1\frac{N}{m}\times\frac{B}{d_{1}}, such that Si,k=∑j=0n−1Wi,jXj,k.\bm{S}_{i,k}=\sum_{j=0}^{n-1}\bm{W}_{i,j}\bm{X}_{j,k}.

Similarly, ΔT\bm{\Delta}^{T} is also partitioned into a grid of d2×md_{2}\times m blocks. This in turn results in CT=ΔTW\bm{C}^{T}=\bm{\Delta}^{T}\bm{W} to be partitioned into d2×nd_{2}\times n blocks, such that [CT]k,j=∑i=0m−1[ΔT]k,iWi,j.[\bm{C}^{T}]_{k,j}=\sum_{i=0}^{m-1}[\bm{\Delta}^{T}]_{k,i}\bm{W}_{i,j}.

Initial Encoding of W\bm{W} (Pre-processing step): The initial encoding of W\bm{W} is similar to that for the case of mini-batch size B=1B=1 (see 10 in Section VII).

Feedforward stage: Similar to the case of B=1B=1, we assume that the entire input X\bm{X} is made available at every node by the previous layerWe will discuss storage and communication tradeoffs in Section IX-C. Each node first partitions the matrix X\bm{X} into a grid of n×d1n\times d_{1} blocks, each sub-matrix being of size Nn×Bd1\frac{N}{n}\times\frac{B}{d_{1}}, and then encodes as follows:

For p=0,1,…,P−1p=0,1,\dots,P-1, the pp-th node evaluates the polynomial X~(v,w)\widetilde{X}(v,w) at (v,w)=(bp,cp)(v,w)=(b_{p},c_{p}), yielding X~p=X~(bp,cp)\widetilde{\bm{X}}_{p}=\widetilde{\bm{X}}(b_{p},c_{p}). E.g., for n=d1=2n=d_{1}=2, X\bm{X} is encoded as X~(v,w)=X0,0v+X0,1+X1,0wv+X1,1w\widetilde{\bm{X}}(v,w)=\bm{X}_{0,0}v+\bm{X}_{0,1}+\bm{X}_{1,0}wv+\bm{X}_{1,1}w. Next, each node computes the matrix-matrix product: S~p:=W~pX~p\widetilde{\bm{S}}_{p}:=\widetilde{\bm{W}}_{p}\widetilde{\bm{X}}_{p}. Computing S~p\widetilde{\bm{S}}_{p} is effectively resulting in the evaluation, at (u,v,w)=(ap,bp,cp)(u,v,w)=(a_{p},b_{p},c_{p}), of the following polynomial:

even though the node is not explicitly evaluating it by accessing all its coefficients separately.

Recall that following the basic idea of Generalized PolyDot codes, by fixing j′=jj^{\prime}=j, we can evaluate the coefficient of uivn−1wku^{i}v^{n-1}w^{k} for i=0,1,…,m−1i=0,1,\dots,m-1 and k=0,1,…,d1−1k=0,1,\dots,d_{1}-1, which turns out to be ∑j=0n−1Wi,jXj,k=Si,k\sum_{j=0}^{n-1}\bm{W}_{i,j}\bm{X}_{j,k}=\bm{S}_{i,k}. Thus, these md1md_{1} coefficients constitute the m×d1m\times d_{1} sub-matrices (or blocks) of the desired result S=WX\bm{S}=\bm{W}\bm{X}. Therefore, S\bm{S} can be recovered at any node if it can reconstruct all the coefficients, or rather these md1md_{1} coefficients, of the polynomial S~(u,v,w)\widetilde{\bm{S}}(u,v,w).

We use substitutions (u=vn,w=vmn)(u=v^{n},w=v^{mn}) or (v=um,w=umn)(v=u^{m},w=u^{mn}) as elaborated in Section IX-C to convert S~(u,v,w)\widetilde{\bm{S}}(u,v,w) into a polynomial of a single variable, and then use standard decoding techniques to interpolate the coefficients of a polynomial in a single variable from its evaluations at PP arbitrary points where some evaluations may be erroneous. Once S\bm{S} is decoded, a nonlinear activation function is applied element-wise to generate the input for the next layer. This also makes X\bm{X} available at each node at the start of the next feedforward layer, justifying our assumption.

The additional steps under Error Model 22 are similar to the case of B=1B=1.

Backpropagation stage: The backpropagated error (transpose) ΔT\bm{\Delta}^{T} is available at every node. Each node partitions ΔT\bm{\Delta}^{T} into a grid of d2×md_{2}\times m blocks, and encodes them using the polynomial:

For p=0,1,…P−1p=0,1,\dots P-1, the pp-th node evaluates Δ~T(w,u)\widetilde{\bm{\Delta}}^{T}(w,u) at (w,u)=(cp,ap)(w,u)=(c_{p},a_{p}), yielding Δ~pT=Δ~T(cp,ap)\widetilde{\bm{\Delta}}^{T}_{p}=\widetilde{\bm{\Delta}}^{T}(c_{p},a_{p}).

Next, it performs the computation C~pT:=Δ~pTW~pT\widetilde{\bm{C}}^{T}_{p}:=\widetilde{\bm{\Delta}}^{T}_{p}\widetilde{\bm{W}}^{T}_{p} and sends the product to all the other nodes, of which, some products may be erroneous. Consider the polynomial:

The products computed at each node effectively results in the evaluations of this polynomial C~(w,u,v)\widetilde{\bm{C}}(w,u,v) at (w,u,v)=(cp,ap,bp)(w,u,v)=(c_{p},a_{p},b_{p}). The coefficients of wkum−1vjw^{k}u^{m-1}v^{j} for k=0,1,…d2−1k=0,1,\dots d_{2}-1 and j=0,1,…,n−1j=0,1,\dots,n-1 in this polynomial actually correspond to the d2×nd_{2}\times n grid of sub-matrices of CT\bm{C}^{T}. Thus, if every node is able to decode these coefficients from the different evaluations of the polynomial at PP nodes, then every node can reconstruct CT\bm{C}^{T}. Observe that, since both CT\bm{C}^{T} and X\bm{X} are available at the node, the backpropagated error for the consecutive layer can be computed at each node by computing the Hadamard product CT∘g(XT)\bm{C}^{T}\circ g(\bm{X}^{T}).

Update stage: Since both X\bm{X} and ΔT\bm{\Delta^{T}} are available at each node, it can now encode the bb-th column of X\bm{X} and the bb-th row of ΔT\bm{\Delta^{T}} for b=0,1,…,B−1b=0,1,\dots,B-1, in a manner similar to that of mini-batch size B=1B=1.

We also let each individual vector δ(b)=[δ(b)0δ(b)1⋮δ(b)(m−1)]\bm{\delta}_{(b)}=\begin{bmatrix}\bm{\delta}_{(b)0}\\ \bm{\delta}_{(b)1}\\ \vdots\\ \bm{\delta}_{(b)(m-1)}\end{bmatrix} be partitioned into mm equal parts (similar to δ\bm{\delta} being partitioned into mm equal parts for the case of mini-batch size B=1B=1) and x(b)=[x(b)0x(b)1⋮x(b)(n−1)]\bm{x}_{(b)}=\begin{bmatrix}\bm{x}_{(b)0}\\ \bm{x}_{(b)1}\\ \vdots\\ \bm{x}_{(b)(n-1)}\end{bmatrix} be partitioned into nn equal parts (similar to x\bm{x}).

Then, the coded W~p\widetilde{\bm{W}}_{p} can be updated as:

Thus, the update step preserves the coded nature of the weight matrix W\bm{W}.

Update with regularization: As in the case of mini-batch size B=1B=1, the coded update can be easily extended to coded update with regularization as follows:

IX-C Comparison with existing strategies.

Here d1d_{1} and d2d_{2} are two integers that divide BB. Choosing higher values of d1d_{1} and d2d_{2} reduce the computational and communication complexity as discussed in Table IV, but also reduce the error tolerance.

The proof is provided in Appendix D. Now we include a table characterizing the storage, communication and computation costs of our proposed strategy in Table IV. Note that, the parameters d1d_{1} and d2d_{2} may be chosen accordingly to vary the communication cost as required.

The derivation of these terms is elaborated in Appendix D. Once again, we use pessimistic bounds for the additional overheads in our proposed strategy, and in spite of that, we are able to show that these additional overheads are negligible as compared to the complexities of the steps O1O1, O2O2 and O3O3. We are now examining strategies to reduce the overheads further.

X Coded Autoencoder

In this section, we show that our coded DNN technique can be easily extended to other commonly used architectures through the example of sparse autoencoders. Sparse autoencoders (see or Section A-D) are a specific type of DNN for learning sparse representations of given data in an unsupervised fashion. Sparse autoencoders usually have only one hidden layer, and training autoencoders with more than one hidden layer are treated as multiple one-hidden-layer autoencoders stacked together. Hence, we will only consider training an autoencoder with one hidden layer.

The major variation in numerical steps arises because of its different loss function E(W1(k),W2(k))E(\bm{W}^{1}(k),\bm{W}^{2}(k)). In mini-batch SGD, for calculating the gradient ∂E(W1(k),W2(k))∂Wi,jl(k)\frac{\partial E(\bm{W}^{1}(k),\bm{W}^{2}(k))}{\partial W^{l}_{i,j}(k)}, the loss function over one mini-batch of data is as follows:

where ρ^i\hat{\rho}_{i} is the sample sparsity of the second layer’s activation averaged over the mini-batch, i.e., ρ^i=1B∑b=0B−1Yi,b1(k)\hat{\rho}_{i}=\frac{1}{B}\sum_{b=0}^{B-1}Y^{1}_{i,b}(k) where Y1(k)=X2(k)=f(S1(k))=f(W1(k)X1(k))\bm{Y}^{1}(k)=\bm{X}^{2}(k)=f(\bm{S}^{1}(k))=f(\bm{W}^{1}(k)\bm{X}^{1}(k)) applied element-wise and X1(k)∈RN0×B\bm{X}^{1}(k)\in\mathcal{R}^{N_{0}\times B} is a matrix whose columns represent the BB data points chosen at the kk-th iteration.

As the difference is only in the loss function, the feedforward stage follow the same procedures given in Section IX. Now let us examine coded backpropagation and update stages. First, notice that updating W2\bm{W}^{2} at the final layer is simply an update with L22 regularization, explained in Section A-D. Then the major difference arises in updating W1\bm{W}^{1} due to the last term in the loss function. Updating W1\bm{W}^{1} follows:

If we compared 22 with the update rule of a generic DNN, where

we observe that the only additional term is Qρ^(k)\bm{Q}_{\bm{\hat{\rho}}}(k). The derivation of these autoencoder equations are given in Appendix A-D. The only additional term in Δauto1\bm{\Delta}_{\textnormal{auto}}^{1} compared to Δ1\bm{\Delta}^{1} is Qρ^\bm{Q}_{\bm{\hat{\rho}}}. Hence, we only have to analyze how this term can be incorporated into our coded DNN framework. Qρ^\bm{Q}_{\bm{\hat{\rho}}} is a matrix with the following form:

To compute this, we first need to obtain ρ^0,⋯ρ^N1−1\hat{\rho}_{0},\cdots\hat{\rho}_{N_{1}-1}. As each node already generates the entire input for the next layer during the feedforward stage, i.e., X2=Y1=f(W1(k)X1(k))\bm{X}^{2}=\bm{Y}^{1}=f(\bm{W}^{1}(k)\bm{X}^{1}(k)), the ρ^i\hat{\rho}_{i}’s can be computed at every node with computation complexity O(N1B)O(N_{1}B). Then, computing Q(ρ^i)Q(\hat{\rho}_{i})’s takes computation complexity of O(N1)O(N_{1}). After we complete coded multiplication, (C~2)pT=(Δ~2)pTW~p2(\widetilde{\bm{C}}^{2})^{T}_{p}=(\widetilde{\bm{\Delta}}^{2})^{T}_{p}\widetilde{\bm{W}}^{2}_{p} and decode (C2)T(\bm{C}^{2})^{T} at each node, the nodes can compute Δauto1\bm{\Delta}_{\textnormal{auto}}^{1} with complexity O(N1B)O(N_{1}B). Then it can also encode the computed Δauto1\bm{\Delta}_{\textnormal{auto}}^{1} for coded update of W1\bm{W}^{1}.

To summarize, the key steps are as follows:

Feedforward stage: Exactly same as before

Backpropagation and Update stage at layer 22: The W2\bm{W}^{2} update is only backpropagation with regularization.

Backpropagation stage at layer 11: The W1\bm{W}^{1} update is following equations (21) and (22).

Computing ρ^\bm{\hat{\rho}}: complexity is Θ(N1B)\Theta(N_{1}B).

Computing Qρ^\bm{Q}_{\bm{\hat{\rho}}}: complexity is Θ(N1)\Theta(N_{1}).

Computing Δauto1\bm{\Delta}_{\textnormal{auto}}^{1}: Computing coded matrix-matrix product (C~2)pT=(Δ~2)pTW~p2(\widetilde{\bm{C}}^{2})^{T}_{p}=(\widetilde{\bm{\Delta}}^{2})^{T}_{p}\widetilde{\bm{W}}^{2}_{p} which takes Θ(N1N2Bd2K)\Theta\left(\frac{N_{1}N_{2}B}{d_{2}K}\right) complexity.

Decode (C2)T(\bm{C}^{2})^{T} at every node: complexity Θ(P3N1Bd2n)\Theta(P^{3}\frac{N_{1}B}{d_{2}n}).

Update stage at layer 11: Encode Δauto1\bm{\Delta}_{\textnormal{auto}}^{1} for the update.

XI Discussion and Conclusions

To summarize, in this work we first proposed a novel coded computing technique for distributed matrix multiplications called Generalized PolyDot and then used it to design a unified coded computing strategy for DNN training. Our proposed coding strategy (and the concurrent strategy of ) advances on the existing coded computing techniques for distributed matrix-matrix multiplication, improving the recovery threshold and error tolerances. Lastly, we also show how our unified strategy can be adapted to specific applications, such as, autoencoders.

The problem of reliable computing using unreliable elements was first posed in 1956 by von Neumann, speculating that the efficiency and reliability of the human brain is obtained by allowing for low power but error-prone components with redundancy for error-resilience. This is also evident from the influence of McCulloch-Pitts model of a neuron in his work. It is often speculated that the error-prone nature of brain’s hardware actually helps it be more efficient: rather than making individual components more reliable using higher power/resources, it might be more efficient to accept component-level errors, and utilize sophisticated error-correction mechanisms for overall reliability of the computationSee also the “Efficient Coding Hypothesis” of Barlow and an application in .. It is thus surprising that this problem of training neural networks under errors has still remained open, even as massive artificial neural networks are being trained on increasingly low-cost and unreliable processing units. We believe that this work (and our prior work ) might be a significant step in the design of biologically inspired neural networks with error resilience that could hold the key to significant improvements in efficiency and reduction of energy consumption during neural network training. Thus, these results could be of broader scientific interest to communities like High Performance Computing (HPC), neuroscience as well as neuromorphic computing.

References

Appendix A DNN Background

We provide a background on DNNs for unfamiliar readers. We also follow the standard description used in DNN literature , so familiar readers can merely skim this part.

DNN Operations: We first explain the training of a DNN with l=1,2,…,Ll=1,2,\ldots,L layers (excluding the input layer which can be thought of as the layer ) using Stochastic Gradient Descent (SGD) with batch size B=1B=1. Later, we will explain how this can be extended to mini-batch SGD with mini-batch size B>1B>1. A Deep Neural Network (DNN) essentially consists of LL weight matrices (also called parameter matrices), one for each layer, that represent the connections between the ll-th and (l−1)(l-1)-th layer for l=1,2,…,Ll=1,2,\ldots,L. At the ll-th layer, NlN_{l} denotes the number of neurons. Thus, for layer ll, the weight matrix to be trained is of dimension Nl×Nl−1N_{l}\times N_{l-1}. At the ll-th layer (l=1,…,Ll=1,\ldots,L), NlN_{l} denotes the number the neurons, (i.e., the row-dimension of the weight matrix).

We use the index kk to denote the iteration number of the training. At the kk-th iteration, the neural network is trained based on a single data point using three stages: a feedforward stage, followed by a backpropagation stage and an update stage. We use the following notations:

N1,…,NL:N_{1},\dots,N_{L}: Number of neurons in layers 1,2,…,L1,2,\dots,L. We also introduce the notation N0N_{0} to denote the dimension of the original data vector, which serves as the input to the first layer.

Wi,jl(k):W_{i,j}^{l}(k): At iteration kk, the weight of the connection from neuron jj on layer l−1l-1 to neuron ii on layer ll for i=0,1,…,Nl−1i=0,1,\dots,N_{l}-1 and j=0,1,…,Nl−1−1j=0,1,\dots,N_{l-1}-1. Note that, the weights actually form a matrix Wl(k)\bm{W}^{l}(k) of dimension Nl×Nl−1N_{l}\times N_{l-1} for layer ll.

xl(k)∈RNl−1:\bm{x}^{l}(k)\in\mathcal{R}^{N_{l-1}}: The input of layer ll at the kk-th iteration. Note that, for the first layer, x1(k)\bm{x}^{1}(k) becomes the data point used for the kk-th iteration of training.

sl(k):\bm{s}^{l}(k): The summed output of the neurons of layer ll before a nonlinear function f(⋅)f(\cdot) is applied on it, at the kk-th iteration. Note that, for i=0,1,…,Nl−1i=0,1,\dots,N_{l}-1, the scalar sil(k)s_{i}^{l}(k) is the i−i-th entry of the vector sl(k)\bm{s}^{l}(k), i.e. the summed output of neuron ii on layer ll.

y^l(k)∈RNl:\hat{\bm{y}}^{l}(k)\in\mathcal{R}^{N_{l}}: The output of layer ll at the kk-th iteration after the application of the nonlinear function. Note that, the output of the last layer LL, i.e., y^L(k)\hat{\bm{y}}^{L}(k) is the final estimated label generated by the neural network that can be compared with the true label.

We now make some observations that explains the functional connectivity across the layers:

Input for any layer is the output of the previous layer except of course for the first layer whose input is the actual data vector itself:

At each layer, the input for that layer is summed with appropriate weights of that layer (entries Wi,jlW_{i,j}^{l}) to produce the summed output of each neuron given by:

The final output of each layer is given by a nonlinear function applied on the summed output of each neuron as below:

Observe that y^L(k)\hat{y}^{L}(k) of the last layer denotes the estimated output or label of the DNN and is to be compared with the true label vector y(k)\bm{y}(k) for the corresponding data point x1(k)\bm{x}^{1}(k).

Key Idea of Training: Let us assume we have a large data set χ\chi consisting of several data points and their labels. The goal of training a DNN is to find the weight matrices W1\bm{W}^{1} to WL\bm{W}^{L} that minimize the empirical loss function defined as follows:

Here bb denotes the index of the data point in χ\chi and ϵ2(b,W1,W2,…,WL)\epsilon^{2}(b,\bm{W}^{1},\bm{W}^{2},\ldots,\bm{W}^{L}) is the loss for the bb-th data point and its label when using the weight matrices (also called parameter matrices) W1,W2,…,WL\bm{W}^{1},\bm{W}^{2},\ldots,\bm{W}^{L} in the neural network. We clarify this using the following example. Let (x1,y)(\bm{x}^{1},\bm{y}) be a particular data point and its true label vector and suppose the network consists of only a single layer followed by an element-wise nonlinear activation function f(⋅)f(\cdot). Then, the estimated label y^L\hat{\bm{y}}^{L} (here L=1L=1) is given by y^L=f(W1x1)\hat{\bm{y}}^{L}=f(\bm{W}^{1}\bm{x}^{1}) applied element-wise. Therefore, the empirical loss function is as follows:

Other commonly used loss functions for ϵ2(⋅)\epsilon^{2}(\cdot) are hinge loss, logistic loss, or cross-entropy loss. The technique does not depend on the specific choice of the loss function.

Gradient Descent (GD) is one way to iteratively minimize this loss function using the following update rule:

Here kk denotes the index of the iteration, ll denotes the layer index of the neural network, Wl(k)\bm{W}^{l}(k) is the value of the weight matrix (or parameter matrix) of layer ll at the kk-th iteration, and Wi,jl(k)W_{i,j}^{l}(k) denotes the (i,j)(i,j)-th scalar element of the matrix Wl(k)\bm{W}^{l}(k). However, as is evident from this update rule, that GD would require access to all the data points to compute the gradient at each iteration which is computationally expensive.

Stochastic Gradient Descent: To remedy this, one often uses Stochastic Gradient Descent (SGD) instead of full Gradient Descent (GD) where the gradient with respect to only a single data point is used at each iteration and the update rule is replaced as:

Here bb denotes the index of the data point which is accessed at the kk-th iteration, and this depends on the iteration index kk. For the ease of explanation of SGD, we can simply rewrite ϵ2(b,W1(k),…,WL(k))\epsilon^{2}(b,\bm{W}^{1}(k),\ldots,\bm{W}^{L}(k)) only as a function of the iteration index kk alone, i.e., ϵ2(k)=ϵ2(b,W1(k),…,WL(k))\epsilon^{2}(k)=\epsilon^{2}(b,\bm{W}^{1}(k),\ldots,\bm{W}^{L}(k)). Thus, the update rule is as follows:

We will also include the case of mini-batch SGD with batch-size BB later in Section A-B.

Derivation of Backpropagation and Update: Recall that, at the kk-th iteration, the weights of every layer ll of the DNN are to be updated as follows:

The backpropagation (and update) helps us to compute the errors and updates in a recursive form, so that the update of any layer ll depends only on the backpropagated error vector of its succeeding layer, i.e. layer l+1l+1 and not on the layers {l+1,l+2,…,L}\{l+1,l+2,\dots,L\}. Here, we provide the backpropagation and update rules. Let us define the backpropagated error vector as δl(k)\bm{\delta}^{l}(k):

For example, one might again consider the special case of the L2 loss function given by:

For the L2 loss function, the error at the last layer is given by:

Now, observe that δil(k)\delta_{i}^{l}(k) can be calculated from δl+1(k)\bm{\delta}^{l+1}(k) as we derive here in Lemma 3.

During the training of a neural network using backpropagation, the backpropagated error vector δil(k)\delta_{i}^{l}(k) for any layer ll can be expressed as a function of the backpropagated error vector of the previous layer as given by:

Now, using the fact that sil(k)=∑j=0Nl−1−1Wi,jl(k)xjl(k)s_{i}^{l}(k)=\sum_{j=0}^{N_{l-1}-1}W_{i,j}^{l}(k)x^{l}_{j}(k), we have

Thus, the update rule is derived as follows:

Note that, the update rule does not depend on any particular choice of loss function.

We assume that a DNN with LL layers (excluding the input layer) is being trained using backpropagation with Stochastic Gradient Descent (SGD)As a first step in this direction of coded neural networks, we assume that the training is performed using vanilla SGD. As a future work, we plan to extend these coding ideas to other training algorithms such as momentum SGD, Adam etc. with a mini-batch size of B=1B=1 . The DNN thus consists of LL weight matrices (see Fig. 3), one for each layer, that represent the connections between the ll-th and (l−1)(l-1)-th layer for l=1,2,…,Ll=1,2,\ldots,L. At the ll-th layer, NlN_{l} denotes the number of neurons. Thus, the weight matrix to be trained is of dimension Nl×Nl−1N_{l}\times N_{l-1}. For simplicity of presentation, we assume that Nl=NN_{l}=N for all layers. In every iteration, the DNN (i.e. the LL weight matrices) is trained based on a single data point and its true label through three stages, namely, feedforward, backpropagation and update, as shown in Fig. 3. At the beginning of every iteration, the first layer accesses the data vector (input for layer 11) and starts the feedforward stage which propagates from layer l=1l=1 to l=Ll=L. For a layer ll, let us denote the weight matrix, input for the layer and backpropagated error for iteration kk by Wl(k)\bm{W}^{l}(k), xl(k)\bm{x}^{l}(k) and δl(k)\bm{\delta}^{l}(k) respectively. The operations performed in layer ll during feedforward stage (see Fig. 3a) can be summarized as:

[[Step O1]O1] Compute matrix-vector product sl(k)=Wl(k)xl(k)\bm{s}^{l}(k)=\bm{W}^{l}(k)\bm{x}^{l}(k).

[[Step C1]C1] Compute input for layer (l+1)(l+1) given by x(l+1)(k)=f(sl(k))\bm{x}^{(l+1)}(k)=f(\bm{s}^{l}(k)) where f(⋅)f(\cdot) is a nonlinear activation function applied elementwise.

At the last layer (l=Ll=L), the backpropagated error vector is generated by assessing the true label and the estimated label, f(sL(k))f(\bm{s}^{L}(k)), which is output of last layer (see Fig. 3b). Then, the backpropagated error propagates from layer LL to 11 (see Fig. 3c), also updating the weight matrices at every layer alongside (see Fig. 3d). The operations for the backpropagation stage can be summarized as:

[[Step O2]O2] Compute matrix-vector product [cl(k)]T=[δl(k)]TWl(k)[\bm{c}^{l}(k)]^{T}=[\bm{\delta}^{l}(k)]^{T}\bm{W}^{l}(k).

[[Step C2]C2] Compute backpropagated error vector for layer (l−1)(l-1) given by [δ(l−1)(k)]T=[cl(k)]TDl(k)[\bm{\delta}^{(l-1)}(k)]^{T}=[\bm{c}^{l}(k)]^{T}\bm{D}^{l}(k) where Dl(k)\bm{D}^{l}(k) is a diagonal matrix whose ii-th diagonal element depends only on the ii-th value of xl(k)\bm{x}^{l}(k). More specifically, Dl(k)\bm{D}^{l}(k) is a diagonal matrix whose ii-th diagonal element is a function g(⋅)g(\cdot) of the ii-th element of xl(k)\bm{x}^{l}(k), such that, g(f(u))=f′(u)g(f(u))=f^{\prime}(u) for the chosen nonlinear activation function f(⋅)f(\cdot) in the feedforward stage. This is equivalent to computing the Hadamard product: [δ(l−1)(k)]T=[cl(k)]T∘g([xl(k)]T)[\bm{\delta}^{(l-1)}(k)]^{T}=[\bm{c}^{l}(k)]^{T}\circ g([\bm{x}^{l}(k)]^{T}).

Finally, the step in the Update stage is as follows:

[[Step O3]O3] Update as: Wl(k+1)←Wl(k)+ηδl(k)[xl(k)]T\bm{W}^{l}(k+1)\leftarrow\bm{W}^{l}(k)+\eta\bm{\delta}^{l}(k)[\bm{x}^{l}(k)]^{T} where η\eta is the learning rate. Sometimes a regularization term is added with the loss function in DNN training (elaborated in Section A-C). For L2 regularization, the update rule is modified as: Wl(k+1)←(1−ηλ)Wl(k)+ηδl(k)[xl(k)]T\bm{W}^{l}(k+1)\leftarrow(1-\eta\lambda)\bm{W}^{l}(k)+\eta\bm{\delta}^{l}(k)[\bm{x}^{l}(k)]^{T} where η\eta is the learning rate and λ2\frac{\lambda}{2} is the regularization constant.

Important Note: As the primary, computationally intensive operations in the DNN training remains roughly the same across layers, we simplify our notations in the main part of the paper. For any layer, we denote the feedforward input xl(k)\bm{x}^{l}(k), the weight matrix Wl(k)\bm{W}^{l}(k) and the backpropagated error δl(k)\bm{\delta}^{l}(k) as x\bm{x}, W\bm{W} and δ\bm{\delta} respectively. We also assume N0=N1=⋯=NL=NN_{0}=N_{1}=\dots=N_{L}=N, and thus the dimension of W\bm{W} is N×NN\times N.

A-B Extension to mini-batch SGD with batch size B>1𝐵1B>1.

Note that while SGD is faster in terms of computations as compared to GD which requires access to the whole data set at each iteration, it is usually less accurate and has noise. So, often one uses mini-batch SGD with a batch size of B>1B>1 which is somewhat intermediate between SGD (batch size B=1B=1) and full GD. In mini-batch SGD, the gradient is computed over a mini-batch of BB data points together than only a single data point, though BB is usually much much smaller than the total size of the data set ∣χ∣|\chi|. Thus, the update rule is given by:

Here, BB denotes the size of the mini-batch, and bb denotes the index of the data point in the subset of BB data points that are chosen for the particular iteration kk. Again, it is more convenient to represent the whole term 1B∑b=0B−1ϵ2(b,W1(k),…,WL(k))\frac{1}{B}\sum_{b=0}^{B-1}\epsilon^{2}(b,\bm{W}^{1}(k),\ldots,\bm{W}^{L}(k)) as only a function of the iteration index kk. Thus, we have, ϵ2(k)=1B∑b=0B−1ϵ2(b,W1(k),…,WL(k))\epsilon^{2}(k)=\frac{1}{B}\sum_{b=0}^{B-1}\epsilon^{2}(b,\bm{W}^{1}(k),\ldots,\bm{W}^{L}(k)).

As we are operating on BB data points instead of a single one during the training using mini-batch SGD, we introduce some matrix notations here as an extension to the vector notations for the case of mini-batch B=1B=1.

N0,N1,…,NL:N_{0},N_{1},\dots,N_{L}: (Stays same as before)

Xl(k)∈RNl−1×B:\bm{X}^{l}(k)\in\mathcal{R}^{N_{l-1}\times B}: The input of layer ll at the kk-th iteration. Note that, this is now a matrix with BB columns instead of a single column vector. Note that, for the first layer, X1(k)\bm{X}^{1}(k) consists of the set of BB data points used for the kk-th iteration of training, all arranged in consecutive columns.

Sl(k)∈RNl×B:\bm{S}^{l}(k)\in\mathcal{R}^{N_{l}\times B}: The summed output of the neurons of layer ll before a nonlinear function f(⋅)f(\cdot) is applied on it, at the kk-th iteration. Note that, during the feedforward stage, Sl(k)=Wl(k)Xl(k)\bm{S}^{l}(k)=\bm{W}^{l}(k)\bm{X}^{l}(k).

Y^l(k)∈RNl×B:\hat{\bm{Y}}^{l}(k)\in\mathcal{R}^{N_{l}\times B}: The output of layer ll at the kk-th iteration. Note that, Y^l(k)=f(Sl(k))\hat{\bm{Y}}^{l}(k)=f(\bm{S}^{l}(k)), applied element-wise and the input to the next layer, i.e., Xl+1(k)=Y^l(k)\bm{X}^{l+1}(k)=\hat{\bm{Y}}^{l}(k).

The derivation of the backpropagation and update rule is similar to the case of mini-batch B=1B=1. The backpropagated error is also now a matrix ∈RNl×B\in\mathcal{R}^{N_{l}}\times B denoted by Δl(k)\bm{\Delta}^{l}(k) whose (i,b)(i,b)-th entry is defined as:

Using an analysis similar to the previous case, the recursion for backpropagation can be derived as follows:

Or, in matrix notations, this reduces to the following steps:

Matrix Multiplication: [Cl(k)]T=[Δl(k)]TWl(k).[\bm{C}^{l}(k)]^{T}=[\bm{\Delta}^{l}(k)]^{T}\bm{W}^{l}(k).

Hadamard product: [Δl−1(k)]T=e(0)[Cl(k)]TD0l(k)+e(1)[Cl(k)]TD1l(k)+⋯+e(B−1)[Cl(k)]TDB−1l(k)[\bm{\Delta}^{l-1}(k)]^{T}=\bm{e}_{(0)}[\bm{C}^{l}(k)]^{T}\bm{D}^{l}_{0}(k)+\bm{e}_{(1)}[\bm{C}^{l}(k)]^{T}\bm{D}^{l}_{1}(k)+\dots+\bm{e}_{(B-1)}[\bm{C}^{l}(k)]^{T}\bm{D}^{l}_{B-1}(k) where Dbl(k)\bm{D}^{l}_{b}(k) is a diagonal matrix that only depends on the bb-th column of Xl(k)\bm{X}^{l}(k), i.e., whose ii-th diagonal element depends on only the element at location (i,b)(i,b) of the matrix Xl(k)\bm{X}^{l}(k). More specifically, Dbl(k)\bm{D}^{l}_{b}(k) is a diagonal matrix whose ii-th diagonal element is a function g(⋅)g(\cdot) of the element at location (i,b)(i,b) of the matrix Xl(k)\bm{X}^{l}(k), such that g(f(u))=f′(u)g(f(u))=f^{\prime}(u) for the chosen nonlinear activation function f(⋅)f(\cdot) in the feedforward stage. Also, e(b)\bm{e}_{(b)} is a unit vector of dimension 1×B1\times B whose bb-th entry is 11 and all others are . To clarify, observe that,

Note that the aforementioned step can also be rewritten as a Hadamard product between two matrices as follows:

the Hadamard product “∘\circ” between two matrices of the same dimension is defined as a matrix of the same dimension as the operands, with each element given by the element-wise product of the corresponding elements of the operands.

Important Note: Similar to the case of B=1B=1, we omit the iteration index kk and layer index ll in the main text of the paper because we are mainly interested in only the matrix operations that occur at each iteration across each layer and these operations are similar across all iterations and layers.

A-C Regularization.

The loss function discussed in 27 is often used in practice with an additional regularization term. Thus, the empirical loss function takes the form:

where R(⋅)R(\cdot) is a regularization function of the matrices W1,⋯ ,WL\bm{W}^{1},\cdots,\bm{W}^{L}. Similar to above, the goal of training DNN is finding the weight matrices W1\bm{W}^{1} to WL\bm{W}^{L} that minimize the loss function in 46.

One of the commonly used regularization techniques in deep learning is adding penalty norms on weights, such as L2 norm and L1 norm. For L2 regularization (weight decay), loss function in (46) can be written as:

It is desirable to use different λ\lambda for the different layers, but for the simplicity we assume that we use the same λ\lambda across all layers. As we already described, SGD is usually preferred for these kinds of optimization problems as compared to the batch Gradient Descent. Thus, the SGD update rule takes the form:

where bb is the index of the data point used in the kk-th iteration. With the weight decay regularization term, the update equation in (40) can thus be rewritten as:

Note that by adding the weight decay term, we have to shrink Wi,jW_{i,j}’s by a constant factor (1−ηλ)(1-\eta\lambda) at each iteration. A similar expression can also be derived for the case of mini-batch B>1B>1.

A-D Training Autoencoders.

Autoencoders are a type of neural networks that are used to learn generative models of data in an unsupervised manner. Autoencoders usually consist of one hidden layer that has smaller number of neurons than input or output layers. The goal of autoencoders is to reconstruct the input data at at the output layer, and the outputs of hidden layer can be considered as compressed representation of data. Autoencoders can have more than one hidden layer, but they are usually treated as single-hidden-layer autoencoders stacked together, and they are trained one by one in a greedy fashion (See Fig 9). There are many variants of autoencoders, such as denoising autoencoders , variational autoencoders , or sparse autoencoders . In this paper, we will focus on sparse autoencoders.

In sparse autoencoders, the number of hidden layer neurons can be bigger than the number of neurons in input or output layers. Obviously, an identity function can perfectly reconstruct the input at the output layer in this case. However, what we are interested in is finding a sparse representation of the data, similar to sparse coding problem. Hence, we add additional term in the loss function to enforce sparsity of the outputs of the hidden layer. Since the autoencoder network has only 22 layers, the loss function EE is defined as below:

where KL denotes Kullback-Leibler divergence, ρ\rho is desired sparsity and ρ^i\hat{\rho}_{i} is the sample sparsity of the second layer’s activation averaged over the entire data set, i.e., ρ^i=1∣χ∣∑b=0∣χ∣−1Yi,b1\hat{\rho}_{i}=\frac{1}{|\chi|}\sum_{b=0}^{|\chi|-1}Y^{1}_{i,b} where Y1=f(W1X1)\bm{Y}^{1}=f(\bm{W}^{1}\bm{X}^{1}) applied element-wise and X1∈RN0×∣χ∣\bm{X}^{1}\in\mathcal{R}^{N_{0}\times|\chi|} is a matrix whose columns represent all the data points.

As full batch GD is expensive, one tends to use mini-batch SGD where the update depends on the gradient calculated on a mini-batch of BB data points instead of the whole data set. In mini-batch SGD, for calculating the gradient ∂E(W1(k),W2(k))∂Wi,jl(k)\frac{\partial E(\bm{W}^{1}(k),\bm{W}^{2}(k))}{\partial W^{l}_{i,j}(k)}, it easier to redefine the loss function E(W1(k),W2(k))E(\bm{W}^{1}(k),\bm{W}^{2}(k)) over one mini-batch of data, as follows:

where ρ^i\hat{\rho}_{i} is the sample sparsity of the second layer’s activation averaged over the mini-batch, i.e., ρ^i=1B∑b=0B−1Yi,b1(k)\hat{\rho}_{i}=\frac{1}{B}\sum_{b=0}^{B-1}Y^{1}_{i,b}(k) where Y1(k)=f(W1(k)X1(k))\bm{Y}^{1}(k)=f(\bm{W}^{1}(k)\bm{X}^{1}(k)) applied element-wise and X1(k)∈RN0×B\bm{X}^{1}(k)\in\mathcal{R}^{N_{0}\times B} is a matrix whose columns represent the BB data points chosen at the kk-th iteration.

Since the additional term in the loss function only affects W1\bm{W}^{1}, the update of W2\bm{W}^{2} remains the same and rewrite the update of W1\bm{W}^{1} as:

We now derive Δauto1(k)\bm{\Delta}^{1}_{\textnormal{auto}}(k) for the backpropagation of autoencoders in the following lemma.

We will omit kk here for the sake of simplicity. Let us first derive the partial derivative of the second term in the loss function first.

Now we obtain the derivative of the first two terms in the loss function.

Appendix B Proof of Theorem 3 and a discussion on error models and decoding techniques

In this appendix, we provide a proof of Theorem 3 along with discussion on the error models and decoding techniques. For completeness, we restate some of the descriptions already presented in Section V.

Recall the channel coding scenario. Let q\bm{q} be a Q×1Q\times 1 vector consisting of QQ real-valued symbols. The received output vector is as follows:

Here G\bm{G} is the generator matrix of a (P,Q)(P,Q) real number MDS Code and e\bm{e} is the P×1P\times 1 error vector that corrupts the codeword GTq\bm{G}^{T}\bm{q}. The locations of the codeword that are affected by errors is a subset A⊆{0,1,…,P−1}\mathcal{A}\subseteq\{0,1,\ldots,P-1\}. The elements of e\bm{e} indexed in A\mathcal{A} denote the corresponding values of additive error while the rest are .

Adversarial Error Model: The subset A\mathcal{A} satisfies ∣A∣≤⌊P−Q2⌋|\mathcal{A}|\leq\lfloor\frac{P-Q}{2}\rfloor, with no specific assumptions on the locations or values of the errors and they may be chosen advarsarially.

B-B Formal Proof of Theorem 3.

Under the Probabilistic Error Model for channel coding, the decoder of a (P,Q)(P,Q) MDS Code can perform the following:

It can detect the occurrence of errors with probability 11, irrespective of the number of errors that occurred.

If the number of errors that occurred is less than or equal to P−Q−1P-Q-1, then all those errors can be corrected with probability 11, even without knowing in advance that how many errors actually occurred.

If the number of errors that occurred is more than P−Q−1P-Q-1, then the decoder is able to determine that the errors are too many to be corrected and declare a “decoding failure” with probability 11.

Let H\bm{H} be the (P−Q)×P(P-Q)\times P sized parity check matrix of the MDS code, such that HGT=0\bm{H}\bm{G}^{T}=\bm{0}. We first propose the following decoding algorithm (Algorithm 1 restated) to produce e^\widehat{\bm{e}}, as an estimate of e\bm{e}, for a given channel output z\bm{z}. If an e^\widehat{\bm{e}} is obtained, the decoder can uniquely solve for q^\widehat{\bm{q}} from the linear set of equations GTq^=z−e^\bm{G}^{T}\widehat{\bm{q}}=\bm{z}-\widehat{\bm{e}}.

Claim 1: Under the Probabilistic Error Model,

Claim 22: Under the Probabilistic Error Model,

Given that ∣∣e∣∣0≤P−Q−1||\bm{e}||_{0}\leq P-Q-1, there is at least one vector, which is the true e\bm{e}, which lies in the search-space of the decoding algorithm and hence the declaration of a decoding failure does not arise. Thus, it is sufficient to show:

Claim 33: Under the Probabilistic Error Model,

Essentially, to prove both Claims 22 and 33, it is sufficient to show that

Next, we examine a term in 62 as follows:

where EA\mathcal{E}_{\mathcal{A}} is the set of eA\bm{e}_{\mathcal{A}} for which there exists another e′≠e\bm{e}^{\prime}\neq\bm{e} such that He=He′\bm{H}\bm{e}=\bm{H}\bm{e}^{\prime}, ∣∣e′∣∣0≤∣A∣||\bm{e}^{\prime}||_{0}\leq|\mathcal{A}| and ∣∣e′∣∣0≤P−Q−1||\bm{e}^{\prime}||_{0}\leq P-Q-1. Note that, this set EA\mathcal{E}_{\mathcal{A}} is a super-set of the set of eA\bm{e}_{\mathcal{A}}s for which e′≠e\bm{e}^{\prime}\neq\bm{e} because the existence of another such e′\bm{e}^{\prime} does not necessarily mean that the decoding algorithm always picks it over the correct e\bm{e}. Examining a term in 63,

where EA,A′⊆EA\mathcal{E}_{\mathcal{A},\mathcal{A}^{\prime}}\subseteq\mathcal{E}_{\mathcal{A}} is the set of all eA\bm{e}_{\mathcal{A}}s for which ∃e′ with non-zero indices A′\exists\bm{e}^{\prime}\text{ with non-zero indices }\mathcal{A}^{\prime} such that He=He′\bm{H}\bm{e}=\bm{H}\bm{e}^{\prime}, ∣∣e′∣∣0≤∣A∣||\bm{e}^{\prime}||_{0}\leq|\mathcal{A}| and ∣∣e′∣∣0≤P−Q−1||\bm{e}^{\prime}||_{0}\leq P-Q-1. Note that,

Observe that when ∣A′∣|\mathcal{A}^{\prime}| is greater than min⁡{∣A∣,P−Q−1}\min\{|\mathcal{A}|,P-Q-1\} for a given eA\bm{e}_{\mathcal{A}}, there is no possible e′\bm{e}^{\prime} that satisfies all the aforementioned conditions, and thus EA,A′\mathcal{E}_{\mathcal{A},\mathcal{A}^{\prime}} is an empty set. Let us see what happens when ∣A′∣≤min⁡{∣A∣,P−Q−1}|\mathcal{A}^{\prime}|\leq\min\{|\mathcal{A}|,P-Q-1\}.

First we will show that EA,A′\mathcal{E}_{\mathcal{A},\mathcal{A}^{\prime}} is also an empty set when ∣A∪A′∣=∣A′∣|\mathcal{A}\cup\mathcal{A}^{\prime}|=|\mathcal{A}^{\prime}|. For He=He′\bm{H}\bm{e}=\bm{H}\bm{e}^{\prime} to hold,

where h∈Null(H)\{0}\bm{h}\in Null(\bm{H})\backslash\{\bm{0}\}. As H\bm{H} is also the generator matrix of a (P,P−Q)(P,P-Q) MDS Code (parity check of a (P,Q)(P,Q) MDS Code), for any h∈Null(H)\{0}\bm{h}\in Null(\bm{H})\backslash\{\bm{0}\}, the number of non-zeros is ∣∣h∣∣0≥P−Q+1||\bm{h}||_{0}\geq P-Q+1. On the other hand, ∣A′∣≤P−Q−1|\mathcal{A}^{\prime}|\leq P-Q-1 which is less than P−Q+1P-Q+1. Thus

for any such e′\bm{e}^{\prime} (with non-zero elements A′\mathcal{A}^{\prime}) to exist, for a given eA\bm{e}_{\mathcal{A}}.

Therefore, if there exists such an e′\bm{e}^{\prime} for which 65 holds, then ∣A∪A′∣−∣A′∣(>0)|\mathcal{A}\cup\mathcal{A}^{\prime}|-|\mathcal{A}^{\prime}|(>0) entries of h(=e−e′)\bm{h}(=\bm{e}-\bm{e}^{\prime}), indexed by the set (A∪A′)\A′(\mathcal{A}\cup\mathcal{A}^{\prime})\backslash\mathcal{A}^{\prime} should match exactly with the true error vector e\bm{e}. For particular A\mathcal{A} and A′\mathcal{A}^{\prime}, satisfying ∣A∪A′∣>∣A′∣|\mathcal{A}\cup\mathcal{A}^{\prime}|>|\mathcal{A}^{\prime}| and ∣A′∣≤min⁡{∣A∣,P−Q−1}|\mathcal{A}^{\prime}|\leq\min\{|\mathcal{A}|,P-Q-1\}, let us denote the set of indices (A∪A′)\A′(\mathcal{A}\cup\mathcal{A}^{\prime})\backslash\mathcal{A}^{\prime} as A∗\mathcal{A}^{*}. Observe that,

where HA∗\bm{H}_{\mathcal{A}^{*}}, HA′\bm{H}_{\mathcal{A}^{\prime}} and HA∪A′\bm{H}_{\mathcal{A}\cup\mathcal{A}^{\prime}} denotes the sub-matrices of the matrix H\bm{H} consisting of the columns indexed in sets A∗\mathcal{A}^{*}, A′\mathcal{A}^{\prime} and A∪A′\mathcal{A}\cup\mathcal{A}^{\prime} respectively. Similarly the sub-scripted h\bm{h} denotes sub-vectors of h\bm{h} consisting of the elements indexed in the sub-script.

Thus, the set EA,A′\mathcal{E}_{\mathcal{A},\mathcal{A}^{\prime}} can be redefined as the set of all eA\bm{e}_{\mathcal{A}}s for which there exists an h=e′−e\bm{h}=\bm{e}^{\prime}-\bm{e} (and hence an hA′\bm{h}_{\mathcal{A}^{\prime}}) such that [HA∗HA′][eA∗hA′]=0\begin{bmatrix}\bm{H}_{\mathcal{A}^{*}}&\bm{H}_{\mathcal{A}^{\prime}}\end{bmatrix}\begin{bmatrix}\bm{e}_{\mathcal{A}^{*}}\\ \bm{h}_{\mathcal{A}^{\prime}}\end{bmatrix}=\bm{0} where A∗=(A∪A′)\A′\mathcal{A}^{*}=(\mathcal{A}\cup\mathcal{A}^{\prime})\backslash\mathcal{A}^{\prime}. Returning to where we left off in 64,

where eA\A∗\bm{e}_{\mathcal{A}\backslash\mathcal{A}^{*}} denotes the elements of e\bm{e} that are independent and outside of those in set A∗\mathcal{A}^{*}, and the set EA,A′,eA\A∗∗\mathcal{E}^{*}_{\mathcal{A},\mathcal{A}^{\prime},\bm{e}_{\mathcal{A}\backslash\mathcal{A}^{*}}} is defined as the set of all eA∗\bm{e}_{\mathcal{A}^{*}}s for a given eA\A∗\bm{e}_{\mathcal{A}\backslash\mathcal{A}^{*}}, such that there exists an hA′\bm{h}_{\mathcal{A}^{\prime}} satisfying

B-C Comparison with Adversarial Error Model.

For completion of the discussion on error models and also to facilitate both the proofs of Theorems 4 and 6, we now also formally show that the number of errors that can be corrected using a (P,Q)(P,Q) MDS Code under the adversarial model is ⌊P−Q2⌋\lfloor\frac{P-Q}{2}\rfloor. The result is well-known in coding theory for finite fields. This is an extension for real numbers (also see ).

Under the adversarial error model for channel coding, the decoder of a (P,Q)(P,Q) MDS Code can correct at most ⌊P−Q2⌋\lfloor\frac{P-Q}{2}\rfloor errors in the worst case.

This statement is also equivalent to the following: A real-valued polynomial with QQ coefficients can be interpolated correctly from its evaluations at PP distinct values, if at most ⌊P−Q2⌋\lfloor\frac{P-Q}{2}\rfloor evaluations are erroneous under the adversarial model.

Let z\bm{z} denote the vector of length PP that denotes the corrupted codeword with additive error e\bm{e}. Thus,

where G\bm{G} is the generator matrix of any real number (P,Q)(P,Q) MDS code. The location of errors, i.e., A\mathcal{A} represents the non-zero indices of e\bm{e}.

As the matrix GT\bm{G}^{T} is full-rank, there exists a full-rank annihilating matrix or parity check matrix H\bm{H} of dimension (P−Q)×P(P-Q)\times P such that HGT=0\bm{H}\bm{G}^{T}=0, leading to He=Hz:=z~.\bm{H}\bm{e}=\bm{H}\bm{z}:=\widetilde{\bm{z}}. Following the steps of , we will show that if ∣A∣≤⌊P−Q2⌋|\mathcal{A}|\leq\lfloor\frac{P-Q}{2}\rfloor, then the solution e\bm{e} to this linear system of equations He=z~\bm{H}\bm{e}=\widetilde{\bm{z}} is unique. Therefore, given z\bm{z} and GT\bm{G}^{T}, the vector q\bm{q} can also be reconstructed uniquely as GT\bm{G}^{T} is full-rank.

Suppose there exists another e′≠e\bm{e}^{\prime}\neq\bm{e} with non-zero indices A′\mathcal{A}^{\prime} such that ∣A′∣≤⌊P−Q2⌋|\mathcal{A}^{\prime}|\leq\lfloor\frac{P-Q}{2}\rfloor and He=He′\bm{H}\bm{e}=\bm{H}\bm{e}^{\prime}. Then, e−e′=h\bm{e}-\bm{e}^{\prime}=\bm{h} where h∈Null(H)\{0}\bm{h}\in Null(\bm{H})\backslash\{\bm{0}\}. The number of non-zero elements in h=e−e′\bm{h}=\bm{e}-\bm{e}^{\prime} satisfies:

This is a contradiction. Any h\bm{h} lying in Null(H)\{0}Null(\bm{H})\backslash\{\bm{0}\} has at most P−Q+1P-Q+1 non-zero elements. This is because H\bm{H} is also the generator of a (P,P−Q)(P,P-Q) MDS code, and thus any (P−Q)(P-Q) columns are linearly dependent.

For the special case of polynomial interpolation, the matrix GT\bm{G}^{T} becomes a Vandermonde matrix. Consider a polynomial g(u)=q0+q1u+⋯+qQ−1uQ−1g(u)=q_{0}+q_{1}u+\dots+q_{Q-1}u^{Q-1} with QQ unique coefficients, evaluated at PP distinct values a0,a1,…,aP−1a_{0},a_{1},\dots,a_{P-1} respectively. Let these evaluations be represented as γ0,γ1,…,γP−1\gamma_{0},\gamma_{1},\dots,\gamma_{P-1} respectively. Thus, we have,

B-D Decoding Methods (Polynomial Time).

The realistic decoding technique under either of the error models would be to first compute H z=z~\bm{H}\ \bm{z}=\widetilde{\bm{z}}. If z~=0\widetilde{\bm{z}}=\bm{0}, one can declare that there are no errors with probability 11. Otherwise, one can use standard sparse reconstruction algorithms to find a solution for the undetermined linear system of equations He=z~\bm{H}\bm{e}=\widetilde{\bm{z}} with minimum non-zero elements for e\bm{e}. Under Error Model 11, if a solution is found with number of non-zero entries at most ⌊P−Q2⌋\lfloor\frac{P-Q}{2}\rfloor, the error is detected and corrected. Under Error Model 22, there is no upper bound on the number of errors. If a solution is found with number of non-zero entries at most P−Q−1P-Q-1, it is declared as the correct e\bm{e}. Otherwise, the node declares a decoding failure even though error is detected, and reverts to the last checkpoint.

Appendix C Proofs of Theorems 4 and 6

Here we provide details of the decoding technique for both feedforward and backpropagation. As we discussed before, the problem of decoding reduces to the problem of interpolating the coefficients of a polynomial in two variables from its values evaluated at (ap,bp)(a_{p},b_{p}) for p=0,1,…,P−1p=0,1,\ldots,P-1, where some of the values might be prone to additive undetected errors (or erasures in the case of stragglers).

Observe that the total number of unique coefficients that the polynomial s~(u,v)\widetilde{\bm{s}}(u,v) has is m(2n−1)m(2n-1) (similarly c~T(u,v)\widetilde{\bm{c}}^{T}(u,v) has n(2m−1)n(2m-1) unique coefficients). Thus, at first glance it might appear that one would require at the very least m(2n−1)+2tm(2n-1)+2t (or n(2m−1)+2tn(2m-1)+2t for backpropagation) evaluations to be able to correct any tt errors during the matrix-vector product of the feedforward or backpropagation stages respectively (under Error Model 11).

However, we actually do not require all the unique m(2n−1)m(2n-1) (or n(2m−1)n(2m-1)) coefficients. We only need the coefficients of uivn−1u^{i}v^{n-1} (or um−1vju^{m-1}v^{j}) which means we effectively require the values of only mm (or nn) coefficients respectively. One technique for solving this interpolation problem is to convert it into a polynomial in a single variable, such that we can reconstruct the required coefficients from its evaluation at PP distinct values at PP different nodes, where some evaluations may be erroneous.

The degree of this polynomial is mn+n−2mn+n-2 which means it has mn+n−1mn+n-1 unique coefficients now. The number of unique coefficients reduce since some of the unique coefficients in s~(u,v)\widetilde{\bm{s}}(u,v) now align together with the substitution u=vnu=v^{n}. However, interestingly, the coefficients of our interest, i.e. the coefficients of uivn−1u^{i}v^{n-1} in s~(u,v)\widetilde{\bm{s}}(u,v) exactly correspond to the coefficients of vin+n−1v^{in+n-1} in s~(v)\widetilde{\bm{s}}(v). For interpolation of a polynomial with mn+n−1mn+n-1 unique coefficients from its evaluations at PP distinct values, one can correct at most ⌊P−mn−n+12⌋\lfloor\frac{P-mn-n+1}{2}\rfloor errors under Error Model 11 (see Lemma 5) and P−mn−nP-mn-n under Error Model 22 (see Theorem 3).

For the backpropagation however, none of the coefficients of c~T(u,v)\widetilde{\bm{c}}^{T}(u,v) align in c~T(v)\widetilde{\bm{c}}^{T}(v) with the substitution u=vnu=v^{n} and c~T(v)\widetilde{\bm{c}}^{T}(v) still has n(2m−1)n(2m-1) unique coefficients but is now a polynomial of one variable. Thus it can correct at most ⌊P−2mn+n2⌋\lfloor\frac{P-2mn+n}{2}\rfloor errors under Error Model 11 (see Lemma 5) and P−2mn+n−1P-2mn+n-1 under Error Model 22 (see Theorem 3).

This is a polynomial of degree mn+m−2mn+m-2, and thus has mn+m−1mn+m-1 unique coefficients. Thus, during backpropagation, one can correct ⌊P−mn−m+12⌋\lfloor\frac{P-mn-m+1}{2}\rfloor errors under Error Model 11 and P−mn−mP-mn-m under Error Model 22. However, for s~(u)\widetilde{\bm{s}}(u) with the substitution v=umv=u^{m}, the number of unique coefficients is m(2n−1)m(2n-1). Thus, the number of errors that can be corrected in feedforward stage is ⌊P−2mn+m2⌋\lfloor\frac{P-2mn+m}{2}\rfloor errors under Error Model 11 (see Lemma 5) and P−2mn+m−1P-2mn+m-1 under Error Model 22 (see Theorem 3).

Now, for the feedforward stage, one uses a (Pfn,m)(\frac{P_{f}}{n},m) MDS code, and is thus able to correct any ⌊Pfn−m2⌋\lfloor\frac{\frac{P_{f}}{n}-m}{2}\rfloor errors under Error Model 11 and Pfn−m−1\frac{P_{f}}{n}-m-1 under Error Model 22.

And for the backpropagation stage, one uses a (Pbm,n)(\frac{P_{b}}{m},n) MDS code. Therefore, one can correct any ⌊Pbm−n2⌋\lfloor\frac{\frac{P_{b}}{m}-n}{2}\rfloor errors under Error Model 11 and Pbm−n−1\frac{P_{b}}{m}-n-1 under Error Model 22.

Under Error Model 22, the probability of two outputs having exactly same error is as the errors are drawn from continuous distributions (recall Remark 7). Thus, as long as an output occurs at least twice, it is almost surely the correct output. Thus, any Pmn−2\frac{P}{mn}-2 errors can be detected and corrected both in the feedforward and backpropagation stages respectively. ∎

Now we move on to the proof of Theorem 6.

The proof is very similar to the previous case. We again consider the four cases separately.

The degree of this polynomial is mnd1+n−2mnd_{1}+n-2 which means that it has mnd1+n−1mnd_{1}+n-1 unique coefficients. The problem of interpolation of the coefficients of this polynomial under erroneous evaluations is equivalent to the problem of decoding of a (P,mnd1+n−1)(P,mnd_{1}+n-1) MDS code. Thus, one can correct ⌊P−mnd1−n+12⌋\lfloor\frac{P-mnd_{1}-n+1}{2}\rfloor errors under Error Model 11 and P−mnd1−nP-mnd_{1}-n under Error Model 22. Now let us see what happens for backpropagation.

The degree of this polynomial is mnd2+mn−n−1mnd_{2}+mn-n-1 which means there are mnd2+mn−nmnd_{2}+mn-n unique coefficients. The problem of interpolation of the coefficients of this polynomial under erroneous evaluations is equivalent to the problem of decoding of a (P,mnd2+mn−n)(P,mnd_{2}+mn-n) MDS code. Thus, one can correct ⌊P−mnd2−mn+n2⌋\lfloor\frac{P-mnd_{2}-mn+n}{2}\rfloor errors under Error Model 11 and P−mnd2−mn+n−1P-mnd_{2}-mn+n-1 errors under Error Model 22.

The degree of this polynomial is mnd1+mn−m−1mnd_{1}+mn-m-1 which means that it has mnd1+mn−mmnd_{1}+mn-m unique coefficients. Thus, one can correct ⌊P−mnd1−mn+m2⌋\lfloor\frac{P-mnd_{1}-mn+m}{2}\rfloor errors under Error Model 11 and P−mnd1−mn+m−1P-mnd_{1}-mn+m-1 under Error Model 22. Now let us see what happens for backpropagation.

The degree of this polynomial is mnd2+n−2mnd_{2}+n-2 which means there are mnd2+n−1mnd_{2}+n-1 unique coefficients. Thus, one can correct ⌊P−mnd2−n+12⌋\lfloor\frac{P-mnd_{2}-n+1}{2}\rfloor errors under Error Model 11 and P−mnd2−nP-mnd_{2}-n under Error Model 22.

The expressions for the MDS-based technique and replication follows from the proof of Theorem 4. ∎

Similar ratios can also be derived for tbt_{b} under Error Model 11 in the limit of large PP, as well as for both tft_{f} and tbt_{b} under Error Model 22.

Appendix D Complexity Analysis and Proofs of Theorem 5 and Theorem 7

Before proceeding to the proofs, we remind the reader about Remark 8. In this work, we assumed that each node can multi-cast simultaneously to at most a constant number (say 22) of nodes. In order to multi-cast the same NN values from one node to PP other nodes, the communication complexity is αlog⁡2P+βN=Θ(N)\alpha\log_{2}{P}+\beta N=\Theta(N), assuming a node cannot communicate to all the PP other nodes simultaneously but uses a spanning-tree type multi-cast mechanism. It might also be worth mentioning that when all PP nodes have to broadcast their own set of NN values to all other nodes except itself, we assume that the communication complexity is αlog⁡2P+βNP=Θ(NP)\alpha\log_{2}{P}+\beta NP=\Theta(NP) using an efficient all-to-all communication protocol called All-Gather . Currently, we are examining strategies to reduce this communication cost even further using improved implementation strategies.

For completeness, we include two of the communication protocols here. The reader can refer to for more details.

Broadcast Communication Protocol: In a cluster of PP nodes, when one node sends a data (say a vector x\bm{x}) of NN values to all other nodes, it is called a broadcast. The communication cost of Broadcast is αlog⁡P+βN\alpha\log{P}+\beta N.

All-Gather Communication Protocol: In a cluster of PP nodes, when every node sends its own, unique data (say vector xp\bm{x}_{p} for node index p=0,1,…,P−1p=0,1,\ldots,P-1) of NN values to all other nodes, it is called an All-Gather. The communication cost of All-Gather is αlog⁡P+2βPN\alpha\log{P}+2\beta PN.

Let us understand the computational and communication complexity per node for all the steps during DNN training in the decentralized implementation.

Matrix-vector products (steps O1O1 and O2O2): The computational complexity is Θ(N2K)\Theta(\frac{N^{2}}{K}) for both feedforward (step O1O1) and backpropagation (step O2O2).

Multi-casting the partial result to all other nodes for decoding: Feedforward stage: For node p=0,1,…,P−1p=0,1,\ldots,P-1, a sub-vector s~p\widetilde{\bm{s}}_{p} of length Nm\frac{N}{m} is to be multi-casted from the pp-th node to all other nodes. The total communication complexity of this step is αlog⁡P+2βNmP=Θ(NmP)\alpha\log{P}+2\beta\frac{N}{m}P=\Theta(\frac{N}{m}P) using All-Gather. Backpropagation stage: For node p=0,1,…,P−1p=0,1,\ldots,P-1, a sub-vector c~pT\widetilde{\bm{c}}^{T}_{p} of length Nn\frac{N}{n} is to be multi-casted from the pp-th node to all other nodes. The communication complexity of this step is αlog⁡P+2βNnP=Θ(NnP)\alpha\log{P}+2\beta\frac{N}{n}P=\Theta(\frac{N}{n}P) using All-Gather.

Decoding at every node (Requires solving sparse reconstruction problems ): Feedforward stage: O(NmP3)\mathcal{O}(\frac{N}{m}P^{3}). Backpropagation stage: O(NnP3)\mathcal{O}(\frac{N}{n}P^{3}). Note that, this step is actually the most dominant among the complexities of all the other steps apart from the steps O1O1, O2O2 and O3O3. We are using a pessimistic upper bound for this decoding complexity which can reduce further based on appropriate choice of decoding algorithm.

Crosscheck for decoding errors: In this step, every node generates a vector of length PP denoting which nodes were found erroneous and then it communicates this vector of PP values to every other node for crosscheck. Thus, the complexity is αlog⁡P+2βP2=Θ(P2)\alpha\log{P}+2\beta P^{2}=\Theta(P^{2}) for both feedforward and backpropagation stage, using All-Gather.

Post-processing on decoded result: Feedforward stage: Each node generates the vector x\bm{x} for next layer by applying nonlinear function f(⋅)f(\cdot) element-wise on vector s\bm{s}. The complexity is thus Θ(N)\Theta(N). Backpropagation stage: Each node multiplies the resulting vector cT\bm{c}^{T} with diagonal matrix D\bm{D}. The complexity is thus Θ(N)\Theta(N).

Encoding of vectors: Feedforward stage: The vector x\bm{x} is divided into nn equal parts of length Nn\frac{N}{n} and a linear combination of all of them is taken. The computational complexity is thus Θ(nNn)=Θ(N)\Theta\left(n\frac{N}{n}\right)=\Theta(N). Backpropagation stage: The vector δT\bm{\delta}^{T} is divided into mm equal parts of length Nm\frac{N}{m}. The computational complexity is thus Θ(mNm)=Θ(N)\Theta\left(m\frac{N}{m}\right)=\Theta(N).

Update stage (step O3O3): It involves a rank-11 update to the stored encoded sub-matrix W~p\widetilde{\bm{W}}_{p} and is thus of computational complexity Θ(N2K)\Theta(\frac{N^{2}}{K}).

Now, we can compute the total computational and communication complexity of all the steps except the individual matrix vector products at each node (step O1O1) for the feedforward stage.

Here the last line follows from the condition of the theorem that P4=o(N)P^{4}=o(N).

Similarly, for the backpropagation stage, the total complexity of all additional steps except the matrix-vector product (step O2O2) is given by,

Here the last line follows from the condition of the theorem that P4=o(N)P^{4}=o(N).

Thus, it is proved that the individual matrix-vector products (steps O1O1 and O2O2) and the update (step O3O3) are the most computationally intensive steps and the additional steps including encoding/decoding add negligible overhead. ∎

Let us now understand the computational and communication complexity of the additional steps for the case of B>1B>1 in the decentralized implementation. Again, we look at the complexities of the various steps as follows:

Matrix-matrix products: Feedforward stage (step O1O1): The multiplication of sub-matrix W~p\widetilde{\bm{W}}_{p} of size Nm×Nn\frac{N}{m}\times\frac{N}{n} with sub-matrix X~p\widetilde{\bm{X}}_{p} of size Nn×Bd1\frac{N}{n}\times\frac{B}{d_{1}} is performed at each node, which is of computational complexity Θ(N2BKd1)\Theta(\frac{N^{2}B}{Kd_{1}}) where K=mnK=mn. Backpropagation stage (step O2O2): The multiplication of sub-matrix Δ~pT\widetilde{\bm{\Delta}}^{T}_{p} of size Bd2×Nm\frac{B}{d_{2}}\times\frac{N}{m} with sub-matrix W~p\widetilde{\bm{W}}_{p} of size Nm×Nn\frac{N}{m}\times\frac{N}{n} is performed at each node, which is of computational complexity Θ(N2BKd2)\Theta(\frac{N^{2}B}{Kd_{2}}) where K=mnK=mn.

Broad-casting the partial result to all other nodes for decoding: Feedforward stage: For p=0,1,…,P−1p=0,1,\ldots,P-1, the node with index pp multi-casts the resulting sub-matrix S~p(=W~pX~p)\widetilde{\bm{S}}_{p}(=\widetilde{\bm{W}}_{p}\widetilde{\bm{X}}_{p}) of dimension Nm×Bd1\frac{N}{m}\times\frac{B}{d_{1}} to all other nodes. Thus, the communication complexity is given by: αlog⁡P+2βNBmd1P=Θ(NBmd1P)\alpha\log{P}+2\beta\frac{NB}{md_{1}}P=\Theta(\frac{NB}{md_{1}}P) using All-Gather. Backpropagation stage: For p=0,1,…,P−1p=0,1,\ldots,P-1, the node with index pp multi-casts the resulting sub-matrix C~pT(=Δ~pTW~p)\widetilde{\bm{C}}^{T}_{p}(=\widetilde{\bm{\Delta}}^{T}_{p}\widetilde{\bm{W}}_{p}) of dimension Bd2×Nn\frac{B}{d_{2}}\times\frac{N}{n} to all other nodes. Thus, the communication complexity is given by: αlog⁡P+2βNBnd2P=Θ(NBnd2P)\alpha\log{P}+2\beta\frac{NB}{nd_{2}}P=\Theta(\frac{NB}{nd_{2}}P) using All-Gather.

Decoding at every node (Requires solving a sparse reconstruction problem ): Feedforward stage: O(NBmd1P3)\mathcal{O}(\frac{NB}{md_{1}}P^{3}). Backpropagation stage: O(NBnd2P3)\mathcal{O}(\frac{NB}{nd_{2}}P^{3}). Note that, this step is actually the most dominant among the complexities of all the other steps apart from the matrix-matrix products (steps O1O1 and O2O2) and the update (step O3O3). We are using a pessimistic upper bound for this decoding complexity which can reduce based on appropriate choice of decoding algorithm.

Crosscheck for decoding errors: For both feedforward and backpropagation, this step is similar to that for the case of B=1B=1. Each node sends a vector of PP values to every other node. Thus, the communication complexity is αlog⁡P+2βP2=Θ(P2)\alpha\log{P}+2\beta P^{2}=\Theta(P^{2}).

Post-processing on decoded result: Feedforward stage: Each node generates the matrix X\bm{X} for the next layer by applying nonlinear function f(⋅)f(\cdot) element-wise on the matrix S\bm{S} of dimensions N×BN\times B. The complexity is thus Θ(NB)\Theta(NB). Backpropagation stage: Each node performs the Hadamard product of the matrix C\bm{C} with the matrix f′(S)f^{\prime}(\bm{S}) applied element-wise, where both the matrices are of dimension N×BN\times B. The complexity is thus Θ(NB)\Theta(NB).

Encoding of sub-matrices: Feedforward stage: The matrix X\bm{X} is divided into an n×d1n\times d_{1} grid of sub-matrices of dimensions Nn×Bd1\frac{N}{n}\times\frac{B}{d_{1}} and a linear combination of all of these nd1nd_{1} matrices is computed at each node. The computational complexity is thus Θ(NBnd1nd1)=Θ(NB)\Theta\left(\frac{NB}{nd_{1}}nd_{1}\right)=\Theta(NB). Backpropagation stage: The matrix ΔT\bm{\Delta}^{T} is divided into a d2×md_{2}\times m grid of sub-matrices of dimensions Bd2×Nm\frac{B}{d_{2}}\times\frac{N}{m} and a linear combination of all of these md2md_{2} matrices is computed at each node. The computational complexity is thus Θ(NBmd2md2)=Θ(NB)\Theta\left(\frac{NB}{md_{2}}md_{2}\right)=\Theta(NB).

Update stage (step O3O3): This involves performing BB consecutive rank-11 updates to the stored, encoded sub-matrix W~p\widetilde{\bm{W}}_{p}. Thus, the complexity is Θ(N2BK)\Theta(\frac{N^{2}B}{K}).

Now, we can compute the total complexity of all the steps except the individual matrix-matrix products (step O1O1) at each node for the feedforward stage.

Here the last line follows from the condition of the theorem that P4=o(N)P^{4}=o(N). Thus the ratio of the complexity of all the additional steps to that of the matrix-matrix products (step O1O1) tends to as K,NK,N and PP scale.

Similarly, for the backpropagation stage, the total complexity of all additional steps except the matrix-matrix product (step O2O2) is given by,

Here the last line follows from the condition of the theorem that P4=o(N)P^{4}=o(N). Thus the ratio of the complexity of all the additional steps to that of the matrix-matrix products (step O2O2) tends to as K,NK,N and PP scale.

Thus, it is proved that the individual matrix-matrix products and the update are the most computationally intensive steps and the additional steps including encoding or decoding add negligible overhead. ∎