"Short-Dot": Computing Large Linear Transforms Distributedly Using Coded Short Dot Products

Sanghamitra Dutta, Viveck Cadambe, Pulkit Grover

Introduction

This work proposes a coding-theory inspired computation technique for speeding up computing linear transforms of high-dimensional data by distributing it across multiple processing units that compute shorter dot products. Our main focus is on addressing the “straggler effect,” i.e., the problem of delays caused by a few slow processors that bottleneck the entire computation. To address this problem, we provide techniques (building on ) that introduce redundancy in the computation by designing a novel error-correction mechanism that allows the size of individual dot products computed at each processor to be shorter than the length of the input. Shorter dot products offer advantages in computation, storage and communication in distributed linear transforms.

The problem of computing linear transforms of high-dimensional vectors is “the" critical step in several machine learning and signal processing applications. Dimensionality reduction techniques such as Principal Component Analysis (PCA), Linear Discriminant Analysis (LDA), taking random projections, require the computation of short and fat linear transforms on high-dimensional data. Linear transforms are the building blocks of solutions to various machine learning problems, e.g., regression and classification etc., and are also used in acquiring and pre-processing the data through Fourier transforms, wavelet transforms, filtering, etc. Fast and reliable computation of linear transforms are thus a necessity for low-latency inference . Due to saturation of Moore’s law, increasing speed of computing in a single processor is becoming difficult, forcing practitioners to adopt parallel processing to speed up computing for ever increasing data dimensions and sizes.

Classical approaches of computing linear transforms across parallel processors, e.g., Block-Striped Decomposition , Fox’s method , and Cannon’s method , rely on dividing the computational task equally among all available processorsStrassen’s algorithm and its generalizations offer a recursive approach to faster matrix multiplications over multiple processors, but they are often not preferred because of their high communication cost . without any redundant computation. The fusion node collects the outputs from each processors to complete the computation and thus has to wait for all the processors to finish. In almost all distributed systems, a few slow or faulty processors – called “stragglers” – are observed to delay the entire computation. This unpredictable latency in distributed systems is attributed to factors such as network latency, shared resources, maintenance activities, and power limitations. In order to combat with stragglers, cloud computing frameworks like Hadoop employ various straggler detection techniques and usually reset the task allotted to stragglers. Forward error-correction techniques offer an alternative approach to deal with this “straggler effect” by introducing redundancy in the computational tasks across different processors. The fusion node now requires outputs from only a subset of all the processors to successfully finish. In this context, the use of preliminary erasure codes dates back to the ideas of algorithmic fault tolerance . Recently optimized Repetition and Maximum Distance Separable (MDS) codes have been explored to speed up computations.

We consider the problem of computing Ax\bm{A}\bm{x} where A(M×N)\bm{A}_{(M\times N)} is a given matrix and x(N×1)\bm{x}_{(N\times 1)} is a vector that is input to the computation (M≪N)(M\ll N). In contrast with , which also uses codes to compute linear transforms in parallel, we allow the size of individual dot products computed at each processor to be smaller than NN, the length of the input.

Why might one be interested in computing short dot products while performing an overall large linear transform? One reason is straightforward: the computation time depends on the length of the dot-products computed. Processors are also inherently memory limited, which limits the size of dot products that can be computed. In some distributed and cloud computing systems, the computation time is dominated by the time taken to communicate x\bm{x} to the processors. In systems where multi-casting x\bm{x} is not possible or is inefficient, it may be faster to communicate a subset of the co-ordinates of x\bm{x} to each processor. In such systems, we anticipate that communicating shorter vectors, each formed by these subsets of coordinates of x\bm{x}, is likely to result in substantial speedups over schemes that require the entire x\bm{x} vector (in particular when multi-casting is difficult).Another interesting example comes from recent work on designing processing units that exclusively compute dot-products using analog components . These devices are prone to errors and increased delays in convergence when designed for larger dot products. In Sections 4 and 6, we show both theoretically (under model assumptions inspired from that admit simplified expected time analysis while being a crude approximation of experimental observations) and experimentally that the speed-up using Short-Dot can be increased beyond that obtained using the strategy proposed in , in straggler-prone environments.

To summarize, our main contributions are:

To compute Ax\bm{A}\bm{x} for a given matrix A(M×N)\bm{A}_{(M\times N)}, we instead compute Fx\bm{F}\bm{x} where we construct F(P×N)\bm{F}_{(P\times N)} (total no. of processors PP > Required no. of dot-products MM) such that each NN-length row of F\bm{F} has at most N(P−K+M)/PN(P-K+M)/P non-zero elements. Because the locations of zeros in a row of F\bm{F} are known by design, this reduces the complexity of computing dot-products of rows of F\bm{F} with x\bm{x}. Here KK parameterizes the resilience to stragglers: any KK of the PP dot products of rows of F\bm{F} with x\bm{x} are sufficient to recover Ax\bm{A}\bm{x}, i.e., any KK rows of F\bm{F} can be linearly combined to generate the rows of A\bm{A}.

We provide fundamental limits on the trade-off between the length of the dot-products and the straggler resilience (number of processors to wait for) for any such strategy in Section 3. This suggests a lower bound on the length of task allotted per processor. Our limits show that Short-Dot is near-optimal.

Assuming exponential tails of service-times at each server (used in ), we derive the expected computation time required by our strategy and compare it to uncoded parallel processing, repetition strategy and MDS coding (see Fig. 2) based linear computation. We also explicitly show a regime (M=Plog⁡PM=\frac{P}{\log P}) where Short-Dot outperforms all its competing strategies in expected computation time, by a factor of log⁡(P)log⁡(log⁡P)\frac{\log(P)}{\log(\log P)}, that diverges to infinity for large PP. In general, Short-Dot is found to be universally faster than all its competing strategies over the entire range of M≤PM\leq P. When MM is linear in PP, Short-Dot offers speed-up by a factor of Ω(log⁡(P))\Omega(\log(P)) over uncoded, parallel processing and repetition. When MM is sub-linear in PP, Short-Dot out-performs repetition or MDS coding based linear computations by a factor of Ω(PMlog⁡(P/M))\Omega\left(\frac{P}{M\log(P/M)}\right) .

We also provide experimental results showing that Short-Dot is faster than existing strategies.

In a concurrent work, in , Tandon et al. consider a coded computation problem similar to ours for the special case where MM, the number of NN-length dot-products to be computed, is 11 and the given matrix AM×N\bm{A}_{M\times N} (in this case just a single row vector) is [1,1,…,1]1×N[1,1,\dots,1]_{1\times N}. We note that for M=1M=1, the gain of using coded strategies over replication-based strategies is bounded even as NN and P→∞P\to\infty for s=Θ(N/P)s=\Theta(N/P) . Our paper differs from in that we consider the more general case M≥1M\geq 1, and observe that the gains over replication can be unbounded with this scaling in the regime s=Θ(MN/P)s=\Theta(MN/P). For M>1M>1, the number of operations per processor using our strategy is lower than an application of for the same worst-case straggler resilience. To see this, note that a straightforward extension of the strategy proposed in that encodes each row of A\bm{A} separately for M(>1)M(>1) rows would require MM dot-products of length N(P−K+1)P\frac{N(P-K+1)}{P} at each processor while using a “joint” encoding across rows, Short-Dot only requires a single dot-product of length N(P−K+M)P\frac{N(P-K+M)}{P} (note that N(P−K+M)P<M×N(P−K+1)P\frac{N(P-K+M)}{P}<M\times\frac{N(P-K+1)}{P}) at each processor, while still requiring the same number of processors (any KK out of PP) to finish. Further, we also provide a tighter converse for M>1M>1 that proves that Short-Dot is near-optimal. It is worth noting that also introduces the notion of partial stragglers, which is outside the scope of our paper.

For multiple dot-products, an alternative repetition-based strategy is to compute MM dot products P/MP/M times in parallel at different processors. Now we only have to wait for at least one processor corresponding to each of the MM vectors to finish (see Fig. 2). Improving upon repetition, it is shown in that an (P,M)(P,M)-MDS code allows constructing PP coded vectors such that any MM of PP dot-products can be used to reconstruct all the MM original vectors (see Fig. 2). This strategy is shown, both experimentally and theoretically, to perform better than repetition and uncoded strategies.

Can we go beyond MDS codes? MDS codes-based strategies require NN-length dot-products to be computed on each processor. Short-Dot instead constructs PP vectors of sparsity ss (less than NN), such that the dot product of x\bm{x} with any K (≥M)K\ (\geq M) out of these PP short vectors is sufficient to compute the dot-product of x\bm{x} with all the MM given vectors (see Fig. 2). Compared to MDS Codes, Short-Dot is more flexible as it waits for some more processors (since K≥MK\geq M), but each processor computes a shorter dot product. Short-Dot also effectively reduces the communication cost since only a shorter portion of the input vector is to be communicated to each processor. We also propose Short-MDS, an extension of the MDS codes-based strategy in to create short dot-products of length ss, through block partitioning, and compare it with Short-Dot. In regimes where Ns\frac{N}{s} is an integer, Short-MDS may be viewed as a special case of Short-Dot. But when Ns\frac{N}{s} is not an integer, Short-MDS has to wait for more processors in worst case than Short-Dot for the same sparsity ss, as discussed in Remark 2 in Section 2.

Our coded parallelization strategy: Short-Dot

Given row vectors {a1T,a2T,…,aMT}\{\bm{a}_{1}^{T},\bm{a}_{2}^{T},\ldots,\bm{a}_{M}^{T}\}, there exists a P×NP\times N matrix F\bm{F} such that a linear combination of any K(>M)K(>M) rows of the matrix is sufficient to generate the row vectors and each row of F\bm{F} has sparsity at most s=NP(P−K+M)s=\frac{N}{P}(P-K+M), provided PP divides NN.

In the next lemma, we show how the row sparsity of F\bm{F} can be constrained to be at most NP(P−K+M)\frac{N}{P}(P-K+M) by appropriately choosing the appended vectors z1,…,zK−M\bm{z}_{1},\ldots,\bm{z}_{K-M}.

Proof: We select a sparsity pattern that we want to enforce on F\bm{F} and then show that there exists a choice of the appended vectors z1,…,zK−M\bm{z}_{1},\ldots,\bm{z}_{K-M} such that the pattern can be enforced. Sparsity Pattern enforced on F\bm{F}: This is illustrated in Fig. 4. First, we construct a P×PP\times P “unit block” with a cyclic structure of nonzero entries, where (K−M)(K-M) zeros in each row and column are arranged as shown in Fig. 4. Each row and column have at most sc=P−K+Ms_{c}=P-K+M non-zero entries. This unit block is replicated horizontally N/PN/P times to form an P×NP\times N matrix with at most scs_{c} non-zero entries in each column, and and at most s=Nsr/Ps=Ns_{r}/P non-zero entries in each row. We now show how choice of z1,…,zK−M\bm{z}_{1},\ldots,\bm{z}_{K-M} can enforce this pattern on F\bm{F}.

Here Aj=[a1(j) a2(j)…aM(j)]T\bm{A}_{j}=[a_{1}(j)\ a_{2}(j)\ldots a_{M}(j)]^{T} is actually the jj-th column of given matrix A\bm{A} and z=[z1(j), … zK−M(j)]T\bm{z}=[z_{1}(j),\ \ldots\ z_{K-M}(j)]^{T} depends on the choice of the appended vectors. Thus,

From Lemmas 1 and 2, for a given M×NM\times N matrix A\bm{A}, there always exists a P×NP\times N matrix F\bm{F} such that a linear combination of any KK columns of F\bm{F} is sufficient to generate our given vectors and each row of F\bm{F} has sparsity at most s=NP(P−K+M)s=\frac{N}{P}(P-K+M). This proves the theorem. ■\blacksquare

Remark 1: Relaxed conditions on matrix B\bm{B}

It has been stated in Lemmas 11 and 22 that all square sub-matrices of B\bm{B} need to be invertible. A matrix with i.i.d. Gaussian entries can be shown to satisfy this property with probability 11. In fact the condition on B\bm{B} in Lemmas 11 and 22 can be relaxed, as evident from the proof. For matrix BP×K\bm{B}_{P\times K} we only need two conditions. (1) All K×KK\times K square sub-matrices are invertible. (2) All (K−M)×(K−M)(K-M)\times(K-M) square sub-matrices in the last K−MK-M columns of B\bm{B} are invertible. A Vandermonde Matrix satisfies both these properties and thus can be used for encoding in Short-Dot.

With this insight in mind, we now formally state our computation strategy:

Remark 2: Short-MDS - a special case of Short-Dot

An extension of the MDS codes-based strategy proposed in , that we call Short-MDS can be designed to achieve row-sparsity ss. First block-partition the matrix of NN columns, into ⌈N/s⌉\left\lceil{N/s}\right\rceil sub-matrices of size M×sM\times s, and also divide the total processors PP equally into ⌈N/s⌉\left\lceil{N/s}\right\rceil parts. Now, each sub-matrix can be encoded using a (P⌈N/s⌉,M)(\frac{P}{\left\lceil{N/s}\right\rceil},M) MDS code. In the worst case, including all integer effects, this strategy requires K=P−⌊P⌈N/s⌉⌋+MK=P-\left\lfloor{\frac{P}{\left\lceil{N/s}\right\rceil}}\right\rfloor+M processors to finish. In comparison, Short-Dot requires K=P−⌊PsN⌋+MK=P-\left\lfloor{\frac{Ps}{N}}\right\rfloor+M processors to finish. In the regime where, ss exactly divides NN, Short-MDS can be viewed as a special case of Short-Dot, as both the expressions match. However, in the regime where ss does not exactly divide NN, Short-MDS requires more processors to finish in the worst case than Short-Dot. Short-Dot is a generalized framework that can achieve a wider variety of pre-specified sparsity patterns as required by the application. In Table 1, we compare the lengths of the dot-products and straggler resilience KK, i.e., the number of processors to wait for in worst case, for different strategies.

Limits on trade-off between the length of dot-products and parameter K

In this section, we derive fundamental trade-offs between the length of the dot-products computed at each individual processor and the number of processors to wait for, i.e., KK, which parametrizes the resilience to stragglers. First we derive an information-theoretic limit in Theorem 2 that holds for any matrix A\bm{A}, such that each column has at least one non-zero entryNote that choice of such a class of matrix A\bm{A} is reasonable, since if say the jj-th column of A\bm{A} consists entirely of zeros, then the jj-th column and its corresponding entry in unknown vector x\bm{x} can simply be omitted from the problem.. In Theorem 3, we show how this bound can be tightened further, so that in the limit of large number of columns of matrix A\bm{A}, Short-Dot is near-optimal.

Let AM×N\bm{A}_{M\times N} be any matrix such that each column has at least one non-zero element. For any matrix FP×N\bm{F}_{P\times N} satisfying the property that the span of its any KK rows contains the span of the MM rows of AM×N\bm{A}_{M\times N}, the average sparsity sˉ\bar{s} over the rows of FP×N\bm{F}_{P\times N} must satisfy \bar{s}\geq\frac{N}{P}\big{(}P-K+1\big{)}.

Proof: We claim that KK is strictly greater than the maximum number of zeros that can occur in any column of the matrix F\bm{F}. If not, suppose the jj-th column of F\bm{F} has more than KK zeros. Then there exists a choice of KK rows of F\bm{F} such that any linear combination of these rows will always be at the jj-th column index. However, since the jj-th column of A\bm{A} has at least one non-zero entry, say at row ii, it is not possible to generate the ii-th row of A\bm{A} by linearly combining these chosen KK rows of F\bm{F}. Thus,

Here the last line follows since maximum value is always greater than average. Note that if sˉ\bar{s} is the average sparsity over the rows of FP×N\bm{F}_{P\times N}, then the average number of zeros over the columns of FP×N\bm{F}_{P\times N} can be written as (N−sˉ)PN\frac{(N-\bar{s})P}{N}. Thus, from (4),

A slight re-arrangement establishes the lower bound in Theorem 2. ■\blacksquare Recall that, Short-Dot achieves a column sparsity of at most (P−K+M)(P-K+M) while a hard lower bound is (P−K+1)(P-K+1) from this proof. The bound is tight for M=1M=1. The bound on average row-sparsity s\geq\frac{N}{P}\big{(}P-K+1\big{)} is also tight only for M=1M=1 (implicitly assuming PP divides NN, since P≪NP\ll N). Now we tighten this bound further for M>1M>1.

Let M>1M>1. Then there exists a matrix AM×N\bm{A}_{M\times N}, such that any FP×N\bm{F}_{P\times N} satisfying the property that any KK rows of FP×N\bm{F}_{P\times N} can span all the rows of AM×N\bm{A}_{M\times N}, must also satisfy the following property:

The average sparsity over the rows of FP×N\bm{F}_{P\times N} is lower bounded as

Moreover, if NN is sufficiently large, such that M2(PK−M+1)=o(N)M^{2}\binom{P}{K-M+1}=o(N), then the average sparsity over the rows of FP×N\bm{F}_{P\times N} is lower bounded as

Note that the second term in the lower bound in (6) does not depend on NN. Thus, if NN is sufficiently larger than PP and MM, the second term in the lower bound becomes negligible compared to the first term, and the first term is precisely what Short-Dot can achieve. Thus, from this lower bound, we can conclude that when NN is large, Short-Dot is near optimal.

Before proceeding with the proof, we give a basic intuition on the proof technique. We basically divide the columns of FP×N\bm{F}_{P\times N} into two groups, one with at most (K−M)(K-M) zeros, and other with more than (K−M)(K-M) zeros. Then we show that there exist matrices AM×N\bm{A}_{M\times N} such that the number of columns in the latter group, i.e., with more than (K−M)(K-M) zeros is bounded, and this in turn bounds the average sparsity. Now we formally prove the theorem.

Proof: Let us denote the number of columns of FP×N\bm{F}_{P\times N} with more than (K−M)(K-M) zeros as λ\lambda. We will show later in Lemma 3 that λ<M(PK−M+1)\lambda<M\binom{P}{K-M+1}. Now, compute the average number of zeros over the columns of F\bm{F}. The columns of F\bm{F} can be divided into two groups : λ\lambda columns with greater than (K−M)(K-M) zeros and (N−λ)(N-\lambda) columns with at-most (K−M)(K-M) zeros. Recall from (3), that if A\bm{A} is chosen such that every column has at-least one non-zero entry, then the maximum number of zeros in any column of F\bm{F} is upper bounded by (K−1)(K-1). Thus, the group of λ\lambda columns can have at-most K−1K-1 zeros each. Thus,

If sˉ\bar{s} is the average sparsity of each row of F\bm{F}, then the average zeros of each column of F\bm{F} is given by (N−sˉ)PN\frac{(N-\bar{s})P}{N}. Thus,

After slight re-arrangement, the average sparsity of each row of F\bm{F} can be bounded as:

Thus, the first part of the theorem, i.e., (6) is proved. Using the condition that M2(PK−M+1)=o(N)M^{2}\binom{P}{K-M+1}=o(N) in (6), we can also obtain (7). Thus,

Thus, the theorem is proved. ■\blacksquare

Lemma 3: Let M>1M>1. Then there exists a matrix AM×N\bm{A}_{M\times N}, such that any FP×N\bm{F}_{P\times N} satisfying the property that any KK rows of FP×N\bm{F}_{P\times N} can span all the rows of AM×N\bm{A}_{M\times N}, must also satisfy the following property: The number of columns (λ)(\lambda) with more than K−MK-M zeros is upper bounded as λ<M(PK−M+1)\lambda<M\binom{P}{K-M+1}.

Proof: Assume, λ≥M(PK−M+1)\lambda\geq M\binom{P}{K-M+1}. Now, a column with more than (K−M)(K-M) zeros will have at least (K−M+1)(K-M+1) zeros. There can be at most (PK−M+1)\binom{P}{K-M+1} different patterns in which (K−M+1)(K-M+1) zeros can occur in a column of length PP. Every column with more than (K−M+1)(K-M+1) zeros also has one of these (PK−M+1)\binom{P}{K-M+1} column sparsity pattern, just with more zeros. From a pigeon-hole argument, at least one of these sparsity patterns of (K−M+1)(K-M+1) zeros will surely occur in λ(PK−M+1)\frac{\lambda}{\binom{P}{K-M+1}} columns or more. Let us consider the sub-matrix of F\bm{F}, of size P×λ(PK−M+1)P\times\frac{\lambda}{\binom{P}{K-M+1}}, consisting of only the columns of F\bm{F} having (K−M+1)(K-M+1) zeros in the same locations, i.e., with similar sparsity pattern. Any KK rows of this sub-matrix of F\bm{F} should generate all the rows of a corresponding M×λ(PK−M+1)M\times\frac{\lambda}{\binom{P}{K-M+1}} sub-matrix of the given A\bm{A}, consisting of the same columns of A\bm{A} as picked in this sub-matrix of F\bm{F}.

There always exists a fully dense matrix A\bm{A} such any M×λ(PK−M+1)M\times\frac{\lambda}{\binom{P}{K-M+1}} sub-matrix of A\bm{A} is full-rank, since A\bm{A} can be arbitrary. This sub-matrix of A\bm{A} is of rank min⁡{M,λ(PK−M+1)}=M\min\{M,\frac{\lambda}{\binom{P}{K-M+1}}\}=M(from assumption). Any KK rows of the sub-matrix of F\bm{F}, should generate MM linearly independent rows of this sub-matrix of A\bm{A}. But since the sub-matrix of F\bm{F} has (K−M+1)(K-M+1) rows consisting of all zeros, there is a choice of KK rows, such that all these zero rows are chosen, and we are only left with at most M−1M-1 non-zero rows to generate MM linearly independent rows of A\bm{A}. This is a contradiction. Thus, we must have λ>M(PK−M+1)\lambda>M\binom{P}{K-M+1}. ■\blacksquare

Analysis of expected computation time for exponential tail models

We now provide a probabilistic analysis of the computation time required by Short-Dot and compare it with uncoded parallel processing, repetition and MDS coding based linear computation scheme as shown in Fig. 5. We follow the shifted-exponential computation time model as described in . Although the shifted exponential distribution may only be a crude approximation of the delay of real systems, we use the shifted exponential model since it is analytically tractable and allows for a fair comparison with the strategy proposed in . We assume that the time required by a processor to compute a single dot-product of length NN be distributed as:

Here, μ(>0)\mu(>0) is the “straggling parameter” that determines the unpredictable latency in computation time. Intuitively, the shifted exponential model states that for a task of size NN, there is a minimum time offset proportional to NN such that the probability of completion of the task before that time is . The probability of task completion is maximum at the time-offset and then decays with an exponential tail after that. This nature of the model might be attributed to the fact that while a processor is most likely to finish its task of size NN at a time proportional to NN, but an unpredictable latency due to queuing and various other factors causes an exponential tail. For an ss length dot product, we simply replace NN by ss in (12), as suggested in . The analysis of expected computation time requires closed form expressions of the KK-th statistic which is simplistic for exponential tails. However a more thorough empirical study is necessary to establish any chosen model for straggling in a particular environment.

The expected computation time for Short-Dot is the expected value of the KK-th order statistic of these PP iid exponential random variables, which is given by:

Here, (13) uses the fact that the expected value of the KK-th statistic of PP iid exponential random variables with parameter 11 is ∑i=1P1i−∑i=1P−K1i≈log⁡(P)−log⁡(P−K)\sum_{i=1}^{P}\frac{1}{i}-\sum_{i=1}^{P-K}\frac{1}{i}\approx\log(P)-\log(P-K) . The expected computation time in the RHS of (13) is minimized when P−K=Θ(M)P-K=\Theta(M). This minimal expected time is O(MNP)\mathcal{O}(\frac{MN}{P}) for MM linear in PP and is O(MNlog⁡(P/M)P)\mathcal{O}\left(\frac{MN\log(P/M)}{P}\right) for MM sub-linear in PP.

A detailed analysis of the expected computation time for the competing strategies, i.e., uncoded strategy, repetition and MDS coding strategy is provided in the Appendix. Table 2 shows the order-sense expected computation time in the regimes where MM is linear and sub-linear in PP.

Note that in the regime where MM is linear in PP, Short-Dot outperforms Uncoded Strategy by a factor diverging to infinity for large PP. Similarly, in the regime where MM is sub-linear in PP, Short-Dot outperforms MDS coding strategy by a factor that diverges to infinity for large PP. Thus Short-Dot universally outperforms all its competing strategies over the entire range of MM.

Now we explicitly provide a regime, where the speed-ups from Short-Dot diverges to infinity for large PP, in comparison to all three competing strategies - MDS Coding, Repetition or Uncoded strategies.

Suppose MM scales as Plog⁡P\frac{P}{\log{P}}. Then, Short-Dot with K=P−M2K=P-\frac{M}{2} has an expected computation time (scaled by NN) as IE[TSD]N=O(log⁡(log⁡P)log⁡P)\frac{{\rm I\kern-2.10002ptE}[T_{SD}]}{N}=O(\frac{\log(\log{P})}{\log{P}}) that decays to as P→∞P\to\infty. In contrast, the expected computation time (scaled by NN) for MDS coding, repetition and uncoded strategies scale as Ω(1)\Omega(1) and thus do not decay to as P→∞P\to\infty.

Proof: For the proof of this theorem, we simply substitute the values of MM and KK in the expressions of expected computation time as follows. We let M=Plog⁡PM=\frac{P}{\log P} for all the strategies. For uncoded strategy, we thus obtain,

For MDS Coding based linear computation, we obtain,

Now, we consider the Short-Dot strategy with K=P−M2=P−P2log⁡PK=P-\frac{M}{2}=P-\frac{P}{2\log{P}}. Note that the inequality K>MK>M is satisfied for log⁡P>32\log{P}>\frac{3}{2}. Now let us calculate the expected computation time for Short-Dot.

Thus, the speed-up offered by Short-Dot in this regime is log⁡Plog⁡(log⁡P)\frac{\log{P}}{\log(\log P)}, and thus diverges to infinity for large PP, as illustrated in Fig. 6.

Encoding and Decoding Complexity

Even though encoding is a pre-processing step (since A\bm{A} is assumed to be given in advance), we include a complexity analysis for the sake of completeness. Recall from Section 2 that we first choose an appropriate matrix B\bm{B} of dimension P×KP\times K, such that every K×KK\times K square sub-matrix is invertible and all (K−M)×(K−M)(K-M)\times(K-M) sub-matrices in the last (K−M)(K-M) columns are invertible. Now, for each of the NN columns of the given matrix A\bm{A}, we perform the following.

modulo𝑗1…𝑗𝐾𝑀1𝑃1\ \ \ \ U\leftarrow(\{(j-1),\ldots,(j+K-M-1)\}\mod P)+1 ⊳\rhd The set of (K−M)(K-M) indices that are 0 for the jj-th column of F\bm{F} Set \ \ \ \bm{B}^{U}\leftarrow\text{Rows of }\bm{B}\text{ indexed byU} Solve for z:     (Bcols M+1:KU)[ z ]=− Bcols 1:MU[Aj]\bm{z}:\ \ \ \ \ (\bm{B}^{U}_{cols\ M+1:K})[\ \bm{z}\ ]=-\ \bm{B}^{U}_{cols\ 1:M}[\bm{A}_{j}] ⊳\rhd z(K−M)×1\bm{z}_{(K-M)\times 1} is a row vector. Set      Fj=B[AjT ∣zT ]T\ \ \ \ \ \bm{F}_{j}=\bm{B}[\bm{A}_{j}^{T}\ |\bm{z}^{T}\ ]^{T} ⊳\rhd Fj\bm{F}_{j} is a column vector ( jj-th col of F\bm{F}) For each of the NN columns, the encoding requires a matrix inversion of size (K−M)×(K−M)(K-M)\times(K-M) to solve a linear system of equations, a matrix-vector product of size (K−M)×M(K-M)\times M and another matrix vector product of size P×KP\times K. The naive encoding complexity is therefore O(N((K−M)3+(K−M)M+PK))\mathcal{O}(N((K-M)^{3}+(K-M)M+PK)). Note that effectively there are only N/PN/P different column sparsity patterns for this particular design discussed in this paper. Thus, there are effectively N/PN/P unique BU\bm{B}^{U}s , and thus N/PN/P unique matrix inversions can suffice for all the NN columns, as sparsity pattern is repeated. Thus, the complexity can be reduced to O(NP(K−M)3+(K−M)MN+PKN)=O(NP(K−M)3+2PKN)\mathcal{O}(\frac{N}{P}(K-M)^{3}+(K-M)MN+PKN)=\mathcal{O}(\frac{N}{P}(K-M)^{3}+2PKN)

This is higher than MDS coding based linear computation that has an encoding complexity of O(NMP))\mathcal{O}(NMP)), but it is only a one-time cost that provides savings in online steps (as discussed earlier in this section).

The encoding complexity can be reduced further for special choices of the matrix B\bm{B}. Let us choose B\bm{B} to be a Vandermonde matrix as given by

Here UU denotes a set of (K−M)(K-M) indices ∈{1,2,…,P}\in\{1,2,\dots,P\}. The matrix-vector product Bcols 1:MU[Aj]\bm{B}^{U}_{cols\ 1:M}[\bm{A}_{j}] is equivalent to the evaluation of a polynomial of degree (K−1)(K-1) with the KK co-efficients as [AjT0(K−M)×1][\bm{A}_{j}^{T}\bm{0}_{(K-M)\times 1}] at (K−M)(K-M) arbitrary points given by {hl∣l∈U}\{h_{l}|l\in U\}. Once this product is obtained, the linear system of equations reduces to the interpolation of the (K−M)(K-M) unknown co-efficients of a polynomial of degree (K−M−1)(K-M-1) (which is z\bm{z}), from its value at (K−M)(K-M) arbitrary points as given by {hl∣l∈U}\{h_{l}|l\in U\}. Once z\bm{z} is obtained, we perform the following operation.

This step is equivalent to the evaluation of a polynomial of degree (K−1)(K-1) at PP points given by {hl∣l=1,2,…P}\{h_{l}|l=1,2,\dots P\}. Thus we decompose our encoding problem for each column of A\bm{A} into a bunch of polynomial evaluation and interpolation problems, all of degree less than PP. Now, from , , we know that both the interpolation and the evaluation of a polynomial of degree less than PP, at PP arbitrary points is O(Plog⁡2(P))\mathcal{O}(P\log^{2}(P)). Thus, the complexity of encoding is O(NPlog⁡2(P))\mathcal{O}(NP\log^{2}(P)).

2 Decoding Complexity:

During decoding, we get KK dot-products from the first KK processors out of PP. We then perform the following operations.

We solve a system of KK linear equations in KK variables and use only MM values of the obtained solution vector. Thus, effectively we do a single matrix inversion of size K×KK\times K followed by a matrix-vector product of size K×MK\times M. The decoding complexity of Short-Dot is thus O(K3+KM)\mathcal{O}(K^{3}+KM) which does not depend on NN when M,K≪NM,K\ll N. This is nearly the same as O(M3+M2)\mathcal{O}(M^{3}+M^{2}) complexity of MDS coding based linear computation.

Similar to encoding, using Vandermonde matrices can reduce the decoding complexity further. As already discussed, we choose the encoding matrix B\bm{B} as a Vandermonde matrix as described in (18). The decoding problem consists of solving a a system of KK linear equations in KK variables.

Here VV is a set of KK indices ∈{1,2,…,P}\in\{1,2,\dots,P\}. The problem of finding w\bm{w} is equivalent to the interpolation of the co-efficients of a polynomial of degree (K−1)(K-1), from its values at KK arbitrary points given by {hl∣l∈V}\{h_{l}|l\in V\}. Again, from , , the interpolation of a polynomial of degree (K−1)(K-1), at KK arbitrary points can be done in O(Klog⁡2(K))\mathcal{O}(K\log^{2}(K)), which thus becomes the decoding complexity.

Experimental Results

We perform experiments on computing clusters at CMU to test the computational time. We use HTCondor to schedule jobs simultaneously among the PP processors. We compare the time required to classify 1000010000 handwritten digits of the MNIST database, assuming we are given a trained 11-layer Neural Network. We separately trained the Neural network using training samples, to form a matrix of weights, denoted by A10×785\bm{A}_{10\times 785}. For testing, the multiplication of this given 10×78510\times 785 matrix, with the test data matrix X785×10000\bm{X}_{785\times 10000} is considered. The total number of processors was 2020.

Assuming that A10×785\bm{A}_{10\times 785} is encoded into F20×785\bm{F}_{20\times 785} in a pre-processing step, we store the rows of F\bm{F} in each processor apriori. Now portions of the data matrix X\bm{X} of size s×10000s\times 10000 are sent to each of the PP parallel processors as input. We also send a C-program to compute dot-products of length s=NP(P−K+M)s=\frac{N}{P}(P-K+M) with appropriate rows of F\bm{F} using command condor-submit. Each processor outputs the value of one dot-product. The computation time reported in Fig. 7 includes the total time required to communicate inputs to each processor, compute the dot-products in parallel, fetch the required outputs, decode and classify all the 1000010000 test-images, based on 3535 experimental runs.

Key Observations: (See Table 3 for detailed results). Computation time varies based on nature of straggling, at the particular instant of the experimental run. Short-Dot outperforms both MDS and Uncoded, in mean computation time. Uncoded is faster than MDS since per-processor computation time for MDS is larger, and it increases the straggling, even though MDS waits for only for 1010 out of 2020 processors. However, note that Uncoded has more variability than both MDS and Short-Dot, and its maximum time observed during the experiment is much greater than both MDS and Short-Dot. The classification accuracy was 85.98%85.98\% on test data.

The experimental times are quite high due to some limitations of the experimental platform used. The time includes some overhead to start the cluster, and communicate data in the form of text files to all the processors, and also collect the output data files back from all the processors. The read time also depends on the size of file to be read. Currently, we are looking at performing these experiments in alternate distributed computing platforms, with better communication protocols.

Discussion

The major advantage of using Short-Dot codes over the MDS coding strategy in is that the length of the pre-stored vectors (rows of F\bm{F}) as well as the communicated input (portions of x\bm{x}) is shorter than NN. It is thus applicable when processing units have limitations of memory, and it is not possible to pre-store the long vectors of length NN. Short-dot also has advantages over in systems where the principle bottlenecks in computation time is in communicating the input x\bm{x} to all the processors, and it may not be feasible to broadcast (multicast) x\bm{x} to all processors at the same time. Thus, it is also useful in applications where communication costs are predominant over computation costs.

2 Errors instead of erasures:

While we focus on the problem of erasures in this paper, Short-Dot can also be used to correct errors. Consider the scenario when instead of straggling or failures, some processors return entirely faulty or garbage outputs, in a distributed system and we do not know which of the outputs are erroneous. We argue from coding theoretic arguments that Short-Dot codes designed to tolerate (P−K)(P-K) stragglers, can also correct ⌊(P−K)2⌋\lfloor\frac{(P-K)}{2}\rfloor errors. First observe that if the code can tolerate (P−K)(P-K) stragglers, then the Hamming Distance between any two code-words should at least be (P−K+1)(P-K+1). Hence, the number of errors that can be corrected is ⌊(Hamming Distance -1)2⌋\lfloor\frac{(\text{Hamming Distance -1})}{2}\rfloor which is ⌊(P−K)2⌋\lfloor\frac{(P-K)}{2}\rfloor. The same result can also be derived by recasting the decoding problem as a sparse reconstruction problem, and borrowing ideas from standard compressive sensing literature which also yields a concrete, decoding algorithm. The problem reduces to an l0l_{0} minimization problem, which can be relaxed into an l1l_{1} minimization, or solved using alternate sparse reconstruction techniques, under certain constraints on the encoding matrix B\bm{B}.

3 More dot-products than processors

While we have presented the case of M<PM<P here, Short-Dot easily generalizes to the case where M≥PM\geq P. The matrix can be divided horizontally into several chunks along the row dimension (shorter matrices) and Short-Dot can be applied on each of those chunks one after another. Moreover if rows with same sparsity pattern are grouped together and stored in the same processor initially, then the communication cost is also significantly reduced during the online computations, since only some elements of the unknown vector x\bm{x} are sent to a particular processor.

Acknowledgments: Systems on Nanoscale Information fabriCs (SONIC), one of the six SRC STARnet Centers, sponsored by MARCO and DARPA. We also acknowledge NSF Awards 1350314, 1464336 and 1553248. S Dutta also received Prabhu and Poonam Goel Graduate Fellowship.

References

Appendix

We now provide a probabilistic analysis of the computational time required by Short-Dot and compare it with uncoded parallel processing, repetition and MDS code based linear computation as shown in Fig. 5. We assume that the time required by a processor to compute a single dot-product follows an exponential distribution and is independent of other parallel processors.

Let us assume, the time required to compute a single dot-product of length NN, follow the distribution:-

Here, μ(>0)\mu(>0) is a straggling parameter, that determines the “unpredictable latency” in computation time. We also assume, that if the length of the dot-product is ss where ss is the sparsity of the vector, the probability distribution of the computational time varies as:-

Now we derive the expected computation time using our proposed strategy and compare it with existing strategies in the regimes where the number of dot-products MM is linear and sub-linear in PP.

Table 2 shows the order-sense expected computation time in the regimes where MM is linear and sub-linear in PP.

The computation time over each of the PP processors behaves as iid exponential random variables following the distribution:-

Now, the expected computation time is the expected value of the KK-th order statistic of these PP iid exponential random variables, which is given by:-

Here we use the result (from ) that the K−K- th order statistic of PP exponential random variables that are iid as ∼exp⁡(−T) ∀ T ≥0\sim\exp(-T)\ \forall\ T\ \geq 0 is given by

For large PP and K<PK<P, we can approximate the following:

Note that the expected computation time is minimized when K=P−Θ(M)K=P-\Theta(M), and is given by:-

If M=Θ(P)M=\Theta(P), the expected time is O(MNP)\mathcal{O}(\frac{MN}{P}). If M=o(P)M=o(P), the expected time is O(MNlog⁡(P/M)P)\mathcal{O}\left(\frac{MN\log(P/M)}{P}\right). Note that s=(P−K+M)NPs=\frac{(P-K+M)N}{P} is actually an upper bound on the length of each dot-product achieved using Short-Dot. Thus the expression obtained in (27) is an upper bound for the actual expected computation time. Thus we use O(.)\mathcal{O}(.) instead of Θ(.)\Theta(.).

2 Existing Strategies

For one single processor to compute all MM dot-products of length NN, the computation time is distributed as

Thus, the expected computation time can be easily derived to be

Now, consider an uncoded strategy where the computation is simply divided into PP dot-products and sent to PP processors. We assume that each processor is sent only one dot-product at a time. We wait for all the processors to finish computation. Note that integer effects arise when MM does not exactly divide PP. Some rows can be divided among ⌈PM⌉\left\lceil{\frac{P}{M}}\right\rceil processors, while the remaining are divided among ⌊PM⌋\left\lfloor{\frac{P}{M}}\right\rfloor processors. Let m1m_{1} and m2m_{2} denote the number of rows that get ⌈PM⌉\left\lceil{\frac{P}{M}}\right\rceil processors and ⌊PM⌋\left\lfloor{\frac{P}{M}}\right\rfloor processors respectively. Clearly the values can be obtained by solving:-

Now, we have two groups of exponential variables: one group consisting of m1⌈PM⌉m_{1}\left\lceil{\frac{P}{M}}\right\rceil iid exponential random variables of task size N⌈PM⌉\frac{N}{\left\lceil{\frac{P}{M}}\right\rceil} , and another group consisting of m2⌊PM⌋m_{2}\left\lfloor{\frac{P}{M}}\right\rfloor iid exponential random variables of task size N⌊PM⌋\frac{N}{\left\lfloor{\frac{P}{M}}\right\rfloor}. The two groups are independent of each other. Note that we assume that NN is large compared to PP and is divisible by P,⌊PM⌋,⌊PM⌋P,\left\lfloor{\frac{P}{M}}\right\rfloor,\left\lfloor{\frac{P}{M}}\right\rfloor, so that the integer effects with respect to NN do not appear and the plots can be scaled with respect to NN for ease of understanding.

The expected computation time is thus given by the expectation of the maximum of all these P=m1⌈PM⌉+m2⌊PM⌋P=m_{1}\left\lceil{\frac{P}{M}}\right\rceil+m_{2}\left\lfloor{\frac{P}{M}}\right\rfloor exponential random variables.

This expression is numerically computed using MATLAB and plotted in the plot of theoretical computation time in Fig. 5. When MM divides PP exactly, the expressions are simpler. The computation time for each processor is distributed as

The expected computation time is the maximum of PP such independent and identically distributed random variables, as given by:-

The expected time for this uncoded strategy is Θ(MNlog⁡(P)P)\Theta\left(\frac{MN\log(P)}{P}\right) regardless of whether MM is linear or sub-linear in PP. Our strategy Short-Dot thus offers a speed-up of Ω(log⁡(P))\Omega(\log(P)) in expected computation time when MM is linear in PP, as mentioned in (27), and thus outperforms by a factor that diverges to infinity for large PP.

When a (P,M)(P,M) repetition strategy is used, we separate the matrix into MM rows and repeat each row P/MP/M times, so as to obtain a total of PP tasks. Note that integer effects arise when MM does not exactly divide PP. Some rows are repeated ⌈PM⌉\left\lceil{\frac{P}{M}}\right\rceil times, while the remaining are repeated ⌊PM⌋\left\lfloor{\frac{P}{M}}\right\rfloor times. Let m1m_{1} and m2m_{2} denote the number of rows that are repeated ⌈PM⌉\left\lceil{\frac{P}{M}}\right\rceil times and ⌊PM⌋\left\lfloor{\frac{P}{M}}\right\rfloor times respectively. Clearly the values can be obtained by solving:-

Now, the minimum of ⌈PM⌉\left\lceil{\frac{P}{M}}\right\rceil (or similarly ⌊PM⌋\left\lfloor{\frac{P}{M}}\right\rfloor ) iid exponential random variables is also exponential with parameter scaled by ⌈PM⌉\left\lceil{\frac{P}{M}}\right\rceil (or similarly ⌊PM⌋\left\lfloor{\frac{P}{M}}\right\rfloor ). The expected computation time is thus given by the expectation of the maximum of m1m_{1} independent exponential variables with parameter scaled by ⌈PM⌉\left\lceil{\frac{P}{M}}\right\rceil and m2m_{2} independent exponential variables with parameter scaled by ⌊PM⌋\left\lfloor{\frac{P}{M}}\right\rfloor.

This expression is computed using MATLAB in the plot of theoretical expected computation time (Fig. 5). When MM exactly divides PP, the analysis is simpler, and the two types of exponential distributions are identical. Following an analysis similar to , it simplifies to the expectation of the maximum of MM iid exponential random variables, each of which is the minimum of P/MP/M iid exponential random variables.

When MM is linear in PP, the expected computation time is Θ(MNPlog⁡(P))\Theta(\frac{MN}{P}\log(P)) while our strategy achieves O(N)\mathcal{O}(N) in this regime. When MM is sub-linear in PP, the expected computation time is Θ(N)\Theta(N) while Short-Dot achieves O(MNlog⁡(P/M)P)\mathcal{O}\left(\frac{MN\log(P/M)}{P}\right) that offers speed-up by a factor diverging to infinity.

The matrix is separated into MM rows and coded into PP rows using a (P,M)(P,M) MDS code. Thus, each processor effectively computes a dot-product of length NN. We have to wait for any MM processors to finish. Assuming the computation of each processor is independent, following an analysis similar to , we obtain that,

When MM is linear in PP, the expected computation time is Θ(N)\Theta(N) as compared to our strategy that achieves O(MN/P)\mathcal{O}(MN/P). However, in the regime where MM is sub-linear in PP, the expected computation time is also Θ(N)\Theta(N) while our strategy achieves O(MNlog⁡(P/M)P)\mathcal{O}\left(\frac{MN\log(P/M)}{P}\right), and thus outperforms MDS codes by a factor that diverges to infinity for large PP.