Coded convolution for parallel and distributed computing within a deadline

Sanghamitra Dutta, Viveck Cadambe, Pulkit Grover

I Introduction

The operation of convolution has widespread applications in mathematics, physics, statistics and signal processing. Convolution plays a key role in solving inhomogeneous differential equations to determine the response of a system under initial conditions . Convolution of any arbitrary input function with the pre-specified impulse response of the system, also known as the Green’s function , provides the response of the system to any arbitrary input function. Convolution is also useful for filtering or extraction of features, in particular for Convolutional Neural Networks (CNNs), which are frequently used in time-critical systems like autonomous vehicles. In , Dally points out that application of machine learning tools in low-latency inference problems is primarily bottlenecked by the computation time.

In this paper, we propose a novel coded convolution for time-critical distributed systems prone to straggling and delays, for fast and reliable computation within a target deadline. We also propose a novel analysis technique based on a new metric – the exponent of the asymptotic probability of failure to meet a target deadline – for comparison of various delay tolerant techniques in distributed systems. Our contributions are two fold:-

We go beyond distributed matrix-vector products, and explore the problem of distributed convolution (building on ) in a straggler prone time-critical distributed scenario. We propose a novel strategy of splitting the vectors to be convolved and coding them, so as to perform fast and reliable convolution within specified deadlines. Moreover our strategy also allows for encoding and decoding online as compared to existing work on coded matrix-vector products which require encoding to be performed as a pre-processing step because of its high computational complexity.

We introduce a new analysis technique to compare the performance of various fault/delay tolerant strategies in distributed systems. Under a more generalized shifted Weibull computation time model (that also encompasses shifted exponential), we demonstrate the utility of the proposed deadline-driven analysis technique for the coded convolution problem.

The motivation for our work arises from the increasing drive towards distributed and parallel computing with ever-increasing data dimensions. Parallelization reduces the per processor computation time, as the task to be performed at each processor is substantially reduced. However, the benefits of parallelization are often limited by “stragglers”: a few slow processors that delay the entire computation. These processors can cause the probabilistic distribution of computation times to have a long tail, as pointed out in the influential paper of Dean and Barroso. For time-critical applications, stragglers hinder the completion of tasks within a specified deadline, since one has to wait for all the processors to finish their individual computations.

The use of replication, , is a natural strategy for dealing with stragglers in distributed systems since the computation only requires a subset of the processors to finish. Error correcting codes have been found to offer advantages over replication in various situations. The use of simple check-sum based codes for correcting errors in linear transforms dates back to the ideas of algorithmic fault tolerance and its extensions in . A strategy of coding multiple parallel convolutions over reals for error tolerance is proposed in . An alternative technique for error-resilient convolutions, that limits itself to finite-field vectors and requires coding both vectors into residue polynomials, was proposed in . Recently, the use of erasure codes, e.g., MDS code based techniques , has been explored for speeding up matrix-vector products in distributed systems. In , we introduce a novel class of codes called – Short-Dot codes – that compute multiple straggler-tolerant short dot products towards computing a large matrix-vector product. The most common metric for comparing the performance of various strategies in existing literature is expected computation time that often uses a shifted exponential computation time model. However, this method of analysis becomes unwieldy when extended to higher order moments. Moreover, expected time does not capture the probability of failure to meet a specific deadline which could be important in time-critical applications.

II Assumptions and Problem Formulation

Before proceeding further, note that the computational complexity of convolving two vectors of any lengths, say m1m_{1} and m2m_{2} using Fast Fourier Transform (FFT) is Θ((m1+m2−1)(log⁡(m1+m2−1)+1))\Theta\left((m_{1}+m_{2}-1)(\log{(m_{1}+m_{2}-1)}+1)\right) . When one of m1m_{1} or m2m_{2} is much smaller than the other, say log⁡m1≪log⁡m2\log{m_{1}}\ll\log{m_{2}}, then the computational complexity can be reduced even further using overlap methods like overlap add . When using overlap methods, the longer vector is divided into smaller pieces of length comparable to the smaller one and each piece is convolved separately. The outputs are then combined together. In this paper, for conceptual simplicity we make some assumptions on the computational complexity of convolution, under different scenarios.

Scenario 1: We assume that when m1m_{1} and m2m_{2} are comparable, the computational complexity of convolution using FFT is given by C(m1+m2)(log⁡(m1+m2))C(m_{1}+m_{2})(\log{(m_{1}+m_{2})}) where CC is a constant.

Scenario 2: We assume that when the length of one vector is sufficiently smaller than another, specifically log⁡(m2)=o(log⁡(m1))\log(m_{2})=o(\log(m_{1})), then the computational complexity using overlap methods is 2Cm1(log⁡(2m2)+1)2Cm_{1}(\log{(2m_{2})}+1) where CC is a constant.

We consider the problem of convolution of a vector a\bm{a} of length N1N_{1} with an unknown input vector x\bm{x} of length N2N_{2}, using PP parallel processing units. The computation goal is to convolve the two vectors reliably in a distributed and parallelized fashion, in the presence of stragglers, so as to prevent failure to meet a specified deadline.

Assumption: For the ease of theoretical analysis, we assume that 2N1N2P≤min⁡{N1,N2}2\sqrt{\frac{N_{1}N_{2}}{P}}\leq\min\{N_{1},N_{2}\}. This ensures that N1N_{1} and N2N_{2} are not too far from each other for relatively small PP.

Note that, the operation of convolution can be represented as the multiplication of a Toeplitz matrix, with a vector, requiring a naive computational complexity of Θ(N1N2)\Theta(N_{1}N_{2}). One might wonder if we should parallelize this matrix vector product and, in doing so, use techniques from that make matrix-vector products resilient to straggling. However, using such a parallelization scheme over the PP given processors, the computational complexity per processor cannot be reduced to below Θ(N1N2P)\Theta(\frac{N_{1}N_{2}}{P}). If PP is not large enough, this is still substantially larger than C(N1+N2)log⁡(N1+N2)C(N_{1}+N_{2})\log{(N_{1}+N_{2})}, the computational complexity of convolution using FFT on a single processor under Scenario 1.

III The Naive Uncoded Strategy

An uncoded strategy to perform the convolution is to divide both the vectors a\bm{a} (length N1N_{1}) and x\bm{x} (length N2N_{2}) into smaller parts of lengths s1s_{1} and s2s_{2} respectively. Let the parts of the vector a\bm{a} be denoted by a1,a2,…,aN1s1\bm{a}_{1},\bm{a}_{2},\dots,\bm{a}_{\frac{N_{1}}{s_{1}}}, each of length s1s_{1}. Similarly, the parts of vector x\bm{x} are denoted by x1,x2,…,xN2s2\bm{x}_{1},\bm{x}_{2},\dots,\bm{x}_{\frac{N_{2}}{s_{2}}}. We first consider the case where the entire convolution is computed in parallel using all the given PP processors, in one shot. Consider the following algorithm:

Thus, we must have N1s1N2s2=P\frac{N_{1}}{s_{1}}\frac{N_{2}}{s_{2}}=P. Depending on the lengths s1s_{1} and s2s_{2}, the convolutions can be performed under Scenario 1 or 2. First consider the case where we perform the convolutions under Scenario 1, choosing s1s_{1} and s2s_{2} of comparable lengths. Then, for minimizing per processor computational complexity, we have the following optimization problem:

A quick differentiation reveals that the above minimization is achieved when s1=s2=N1N2/Ps_{1}=s_{2}=\sqrt{N_{1}N_{2}/P} (say ss). Thus, the computational complexity required for convolution under Scenario 1 would be C\big{(}2\sqrt{\frac{N_{1}N_{2}}{P}}\big{)}\log{\Big{(}2\sqrt{\frac{N_{1}N_{2}}{P}}\Big{)}}. We show in Appendix A that operating under Scenario 2 requires equal or higher per processor complexity than the aforementioned strategy. We also show that, when we use the PP processors several times instead of operating in a single shot, the computational complexity only increases.

Thus, for uncoded convolution, we provably show that the minimum computational complexity per processor in order sense is attained when s1=s2=N1N2/Ps_{1}=s_{2}=\sqrt{N_{1}N_{2}/P} (say s=N1N2/Ps=\sqrt{N_{1}N_{2}/P}). Now, we can perform PP convolutions of ss-length vectors, convolving each of the N1s\frac{N_{1}}{s} parts of a\bm{a} with N2s\frac{N_{2}}{s} parts of x\bm{x} in parallel. A convolution of length s=N1N2/Ps=\sqrt{N_{1}N_{2}/P} in each processor has a computational complexity of 2CN1N2Plog⁡2N1N2P2C\sqrt{\frac{N_{1}N_{2}}{P}}\log{2\sqrt{\frac{N_{1}N_{2}}{P}}}, where CC is a constant which is independent of the length of the convolution. Fig. 1 shows the uncoded convolution strategy.

Reconstruction: First let us discuss how to reconstruct each a∗xj\bm{a}\ast\bm{x}_{j} of length N1+s2N_{1}+s_{2} from successfully computed convolution outputs {ai∗xj∣i=1,2,…,N1/s1}\{\bm{a}_{i}\ast\bm{x}_{j}|i=1,2,\dots,N_{1}/s_{1}\} each of size s1+s2s_{1}+s_{2}. Consider the following algorithm:

Note that this uncoded strategy requires all the PP processors to finish. If some processors straggle, the entire computation is delayed, and may fail to finish within the specified deadline. In the next sections, we first analyze a replication strategy that introduces redundancy and does not require all the processors to finish, and then propose our coded convolution strategy.

IV The Replication Strategy

In the previous section, we have shown that the strategy of minimizing the per processor computational complexity in convolution is to divide each vector into equal sized portions, and convolve in parallel processors, all in one shot. Now, let us consider a (P,r)(P,r) replication strategy, where every sub task has rr replicas or copies. We basically divide the task of convolution into P/rP/r equal parts, and each part has rr replicas, so as to use all the given PP processors. Then, the length of the vectors to be convolved at each processor is given by s=N1N2rPs=\sqrt{\frac{N_{1}N_{2}r}{P}}.

For the convolution of a vector of length N1N_{1} with another vector of length N2N_{2} using PP processors, a (P,r)(P,r) replication strategy needs K=P−Ps2N1N2+1K=P-\frac{Ps^{2}}{N_{1}N_{2}}+1 short convolutions of length s=N1N2rPs=\sqrt{\frac{N_{1}N_{2}r}{P}} to finish in the worst case.

Observe that there are Pr\frac{P}{r} tasks that need to finish, each having rr replicas. The worst case wait arises when the replicas of all but one of the Pr\frac{P}{r} tasks have completed before any one copy of the last task has finished. Thus,

V The Coded Convolution Strategy

We now describe our proposed coded convolution strategy which ensures that computation outputs from a subset of the PP processors are sufficient to perform the overall convolution.

Let us choose a length ss such that N1N2P<s\sqrt{\frac{N_{1}N_{2}}{P}}<s. Now, we can divide both the vectors a\bm{a} and x\bm{x} into parts of length ss, as shown in Fig. 3. The vector x\bm{x} has N2s\frac{N_{2}}{s} parts of length ss, namely x1,…,xN2s\bm{x}_{1},\dots,\bm{x}_{\frac{N_{2}}{s}}. Similar to uncoded strategy, for every convolution a∗xj\bm{a}\ast\bm{x}_{j} we assign Ps/N2Ps/N_{2} processors . The N1s\frac{N_{1}}{s} parts of the vector a\bm{a} are encoded using a \big{(}\frac{Ps}{N_{2}},\frac{N_{1}}{s}\big{)} MDS code to produce PsN2\frac{Ps}{N_{2}} vectors of length ss, as discussed below:

Let \bm{G}_{\big{(}\frac{N_{1}}{s}\times\frac{Ps}{N_{2}}\big{)}} be the generator matrix of a (PsN2,N1s)\left(\frac{Ps}{N_{2}},\frac{N_{1}}{s}\right) MDS code on the real field, i.e., a matrix such that every N1s×N1s\frac{N_{1}}{s}\times\frac{N_{1}}{s} square sub-matrix is invertible . For example, an \big{(}\frac{N_{1}}{s}\times\frac{Ps}{N_{2}}\big{)} Vandermonde matrix satisfies this property. Now we encode the ai\bm{a}_{i}s as follows:

Here GLT\bm{G}^{T}_{L} denotes a N1s×N1s\frac{N_{1}}{s}\times\frac{N_{1}}{s} square sub-matrix of GT\bm{G}^{T}, consisting of the rows indexed in LL. Since GLT\bm{G}^{T}_{L} is invertible, we have

The decoded outputs can be combined similar to the uncoded strategy. Now, we state some results for Coded Convolution.

For convolution of a vector of length N1N_{1} with another vector of length N2N_{2} using PP processors, an (N1,N2,P,s)(N_{1},N_{2},P,s) Coded Convolution decoder needs K=P−PsN2+N1sK=P-\frac{Ps}{N_{2}}+\frac{N_{1}}{s} short convolutions of length s=N1N2Ps=\sqrt{\frac{N_{1}N_{2}}{P}} to finish in the worst case.

Note that each part of x\bm{x} of length ss is convolved with PsN2\frac{Ps}{N_{2}} vectors, so that any N1s\frac{N_{1}}{s} are sufficient. Now, in the worst case, all but one part of x\bm{x} finishes convolution with all PsN2\frac{Ps}{N_{2}} vectors before the last one finishes with N1s\frac{N_{1}}{s} vectors. Thus, the number of processors required is at most,

V-B Comparison with replication strategy:

We first show that the worst-case number of processors required using Coded Convolution is less than that required using replication for the same per-processor computational complexity.

For convolution of vector of length N1N_{1} with an unknown input vector of length N2N_{2} using PP processors, an (N1,N2,P,s)(N_{1},N_{2},P,s) Coded Convolution decoder waits for fewer processors in the worst case as compared to that of a (P,r)(P,r) replication strategy with equal per processor computational complexity.

Recall from Section IV that in a (P,r)(P,r) replication strategy, the length of the vectors convolved at each processor is given by N1N2rP\sqrt{\frac{N_{1}N_{2}r}{P}}. For equal per processor task, this length should be equal to ss in Coded Convolution. Using the fact that s=N1N2rP>N1N2Ps=\sqrt{\frac{N_{1}N_{2}r}{P}}>\sqrt{\frac{N_{1}N_{2}}{P}}, we get,

Comment: To justify our focus on the worst case KK, we also discuss in Section VI how the exponent of the probability of failure to meet a deadline in the asymptotic regime depends on the worst case KK.

V-C Computational Complexity of Encoding and Reconstruction:

Observe that, the computational complexity of each processor is (2Cs)log⁡(2s)(2Cs)\log{(2s)} from Scenario 1. To be able to encode and reconstruct online, it is desirable that the encoding complexity and reconstruction complexity, i.e., the total complexity of decoding and subsequent additions are both negligible compared to the per processor complexity, since otherwise the encoder/decoder becomes the primary bottleneck during successful completion of tasks within a specified deadline.

For an (N1,N2,P,s)(N_{1},N_{2},P,s) Coded Convolution using a Vandermonde Encoding Matrix, the ratio of the total encoding complexity and reconstruction complexity to the per processor complexity tends to , as N1,N2,P→∞N_{1},N_{2},P\to\infty, when P(\log{P})^{2}=o\Big{(}\log{\sqrt{\frac{N_{1}N_{2}}{P}}}\Big{)}.

A proof based on , is provided in Appendix B.

VI Asymptotic Analysis of Exponents

We now perform an asymptotic analysis of the exponent of the probability of failure to meet a deadline in the limit of the deadline diverging to infinity. Let Fs(t)F_{s}(t) denote the probability that a convolution of two vectors of length ss is computed in a single processor within time tt, based on a shifted Weibull model. Thus

Here μ>0\mu>0 is a straggling parameter, α>0\alpha>0 denotes the exponent of the delay distribution and CC denotes the constant of convolution that is independent of the length of the vectors. For shifted exponential model, α=1\alpha=1.

Let Psf(t)P^{f}_{s}(t) be the probability of failure to finish the convolution of two vectors of length ss at or before time tt.

From Theorem 3, when P(\log{P})^{2}=o\big{(}\log{\sqrt{N_{1}N_{2}/P}}\big{)}, the decoding complexity is negligible compared to the per processor complexity. Thus, only the straggling in the parallel processors is the bottleneck for failure to meet a deadline. We now compute the failure-exponent for the coded strategy, in the limit of t→∞t\to\infty. For the coded strategy, at most K=P−PsN2+N1sK=P-\frac{Ps}{N_{2}}+\frac{N_{1}}{s} processors need to finish computation in the worst case. An upper bound on the failure probability is thus given by the probability that at most K−1K-1 processors have finished the convolution of two ss length vectors. Thus,

Here c(P,K)c(P,K) denotes a function of PP and KK, that is independent of tt. The last line follows since for tt large enough, (Fs(t))(1−Fs(t))≫1\frac{(F_{s}(t))}{(1-F_{s}(t))}\gg 1, and thus for the purpose of analysis of the failure exponent for large tt, the binomial summation is of the same order as the largest term, as we show here:

Bounding the leading co-efficient in the failure exponent,

. Definition: Let us define a function ϵ(s)\epsilon(s) as follows:

Note that ϵ(s)\epsilon(s) is an upper bound on the leading co-efficient of the failure exponent for an (N1,N2,P,s)(N_{1},N_{2},P,s) Coded Convolution. Now we compare coding with the uncoded strategy.

For the uncoded strategy in Section III, the failure event occurs exactly (not upper bound) when at most P−1P-1 processors have failed to finish. The failure exponent can thus be computed exactly using steps similar to (13). For the uncoded strategy, the length of the vectors in each processor is s=N1N2/Ps=\sqrt{N_{1}N_{2}/P} and the number of processors to wait for is PP (also given by K=P−PsN2+N1sK=P-\frac{Ps}{N_{2}}+\frac{N_{1}}{s} from Theorem 2). The leading co-efficient in the failure exponent is exactly given by,

For N1=Θ(N2)N_{1}=\Theta(N_{2}) and fixed Weibull parameter α\alpha, there exists an (N1,N2,P,s)(N_{1},N_{2},P,s) Coded Convolution strategy with ss strictly greater than N1N2/P\sqrt{N_{1}N_{2}/P} such that the ratio of the exponents of the asymptotic failure probability of the coded strategy to the naive uncoded strategy scales as Ω(P)\Omega(\sqrt{P}) and thus diverges to infinity for large PP.

From (14), the leading co-efficient in the failure exponent for the uncoded strategy is \epsilon\big{(}\sqrt{N_{1}N_{2}/P}\big{)}. Now let us choose a coded strategy with s=2N1N2/Ps=2\sqrt{N_{1}N_{2}/P}. Putting this value of ss in ϵ(s)\epsilon(s), we get ϵ(s)=−μα(32PN1N2+1)(4CN1N2Plog⁡4N1N2P)α\epsilon(s)=-\frac{\mu^{\alpha}(\frac{3}{2}\sqrt{\frac{PN_{1}}{N_{2}}}+1)}{\left(4C\sqrt{\frac{N_{1}N_{2}}{P}}\log{4\sqrt{\frac{N_{1}N_{2}}{P}}}\right)^{\alpha}}. Thus comparing the ratio,

Since from the conditions of the theorem, N1=Θ(N2)N_{1}=\Theta(N_{2}), the ratio scales as Ω(P)\Omega(\sqrt{P}) and diverges for large PP. ∎

Comment: The best choice of ss would be an integer in the range \big{(}\sqrt{N_{1}N_{2}/P},\min\{N_{1},N_{2}\}\big{]} that maximizes ∣ϵ(s)∣\left|\epsilon(s)\right|.

We show in Appendix C that the exponent of the probability of failure to meet a deadline in the asymptotic regime depends on the worst case KK for replication. In this paper, we include simulations in Fig. 4 showing that coded convolution outperforms replication in terms of the exponent of the probability of failure to meet a deadline. Moreover the exponents observed from the simulations are quite close to those calculated theoretically using the worst case KK.

Simulation Results: We consider the convolution of a given vector of length N1=212N_{1}=2^{12} with an unknown vector of length N2=211N_{2}=2^{11}, using P=8P=8 parallel processors, based on an exponential model in MATLAB. Our results in Fig. 4 show that coded convolution has the fastest decay of failure exponent for large deadlines. The slopes obtained from simulations are found to be quite close to the theoretically calculated values.

VII Discussion

Thus, the proposed strategy outperforms existing strategies in terms of the exponent of the probability of failure to meet a deadline, in the limit of large deadlines. It might be observed that our strategy allows for online encoding and reconstruction as compared to existing works on coded matrix-vector products where the encoding cost can be amortized only if the matrix is pre-specified.

Acknowledgment

This work was supported in part by Systems on Nanoscale Information fabriCs (SONIC), one of the six SRC STARnet Centers, sponsored by MARCO and DARPA. The support of NSF Awards 1350314, 1464336 and 1553248 is also acknowledged. Sanghamitra Dutta also received Prabhu and Poonam Goel Graduate Fellowship.

The authors would like to thank Franz Franchetti, Tze Meng Low and Doru Thom Popovici for their useful feedback regarding this research. Praveen Venkatesh, Haewon Jeong and Yaoqing Yang are also thanked for their comments and suggestions.

References

Appendix A Optimal Uncoded Strategy

We prove that performing convolution under Scenario 2, results in equal or higher per processor complexity. Without loss of generality, assume that s2≤s1s_{2}\leq s_{1}. Then, under Scenario 2, the complexity of the convolution is given by C(2s1)(log⁡(2s2)+1)C(2s_{1})(\log{(2s_{2})}+1). Now, we are to solve

This function is concave in s1s_{1}. Let us examine its derivative.

The derivative is positive in the range of s1s_{1}, i.e. 1≤s1≤min⁡{N1,N1N2P}1\leq s_{1}\leq\min\{N_{1},\frac{N_{1}N_{2}}{P}\}. Thus, the minimum is attained for the least value of s1s_{1}. However, as the product of s1s_{1} and s2s_{2} is constant, s1s_{1} has to be greater or equal to the geometric mean, since otherwise s2>s1s_{2}>s_{1} which violates our assumption. Thus, the minimum value of s1s_{1} is N1N2P\sqrt{\frac{N_{1}N_{2}}{P}}, which results in the same time complexity as convolution without overlap methods.

Now, we consider performing convolutions using several serial uses of the given PP parallel processors, say rr uses. Then, effectively we have rPrP processors, and the length of the convolution vector in each serial use of a single processor is N1N2rP\sqrt{\frac{N_{1}N_{2}}{rP}}. As this operation is to be done rr times, the computational complexity may be written as r(2N1N2rP)log⁡(2N1N2rP)r(2\sqrt{\frac{N_{1}N_{2}}{rP}})\log{(2\sqrt{\frac{N_{1}N_{2}}{rP}})}, which increases roughly with a factor of r\sqrt{r}, compared to convolution in a single shot. Thus, use of the processors serially, only worsens the computational complexity per processor, while still requiring all of them to finish.

Appendix B Encoding and Decoding Complexity

Theorem 3: For an (N1,N2,P,s)(N_{1},N_{2},P,s) Coded Convolution using a Vandermonde Encoding Matrix, the ratio of the total encoding complexity and reconstruction complexity to the per processor complexity tends to , as N1,N2,P→∞N_{1},N_{2},P\to\infty, when P(\log{P})^{2}=o\Big{(}\log{\sqrt{\frac{N_{1}N_{2}}{P}}}\Big{)}.

Encoding For encoding, the first vector a\bm{a} is split into smaller pieces of length ss, as given by {a1,a2,…,aN1s}\{\bm{a}_{1},\bm{a}_{2},\dots,\bm{a}_{\frac{N_{1}}{s}}\}. Now we encode the ai\bm{a}_{i}s as follows:

Here GN1s,PsN2\bm{G}_{\frac{N_{1}}{s},\frac{Ps}{N_{2}}} is chosen as a Vandermonde matrix and {g1,g2,…,gPsN2}\{g_{1},g_{2},\dots,g_{\frac{Ps}{N_{2}}}\} are all distinct reals. Observe that the encoding process thus becomes the evaluation of polynomials of degree N1s−1\frac{N_{1}}{s}-1( Note that N1s<PsN2\frac{N_{1}}{s}<\frac{Ps}{N_{2}}) at PsN2\frac{Ps}{N_{2}} points, repeated ss times for ss length vector ai\bm{a}_{i}s. From , , it is known that the evaluation of a polynomial of degree PsN2−1\frac{Ps}{N_{2}}-1 at PsN2\frac{Ps}{N_{2}} arbitrary points can be performed in O(PsN2(log⁡PsN2)2)\mathcal{O}(\frac{Ps}{N_{2}}(\log{\frac{Ps}{N_{2}}})^{2}). As this process is repeated ss times, the total encoding complexity is given by O(s(PsN2)(log⁡PsN2)2)\mathcal{O}\left(s\left(\frac{Ps}{N_{2}}\right)(\log{\frac{Ps}{N_{2}}})^{2}\right). Now observe that,

Note that, (21) holds from the choice of N1N2P≤s\sqrt{\frac{N_{1}N_{2}}{P}}\leq s which implies PsN2≤P\frac{Ps}{N_{2}}\leq P. Thus, the encoding complexity is O(sP(log⁡P)2)\mathcal{O}(sP(\log{P})^{2}).

Reconstruction The reconstruction complexity consists of the complexity of decoding and subsequent additions. First we compute the complexity of decoding.

For decoding, convolution outputs of length 2s2s arrive from each parallel processor. For reconstructing each of the N2s\frac{N_{2}}{s} parts of vector x\bm{x}, any N1s\frac{N_{1}}{s} convolution outputs are sufficient. Recall that for any xj\bm{x}_{j},

If G\bm{G} is chosen as a Vandermonde Matrix, the reconstruction problem takes the form,

The reconstruction of {ai∗xj}\{\bm{a}_{i}\ast\bm{x}_{j}\} for i=1,2,…,N1si=1,2,\dots,\frac{N_{1}}{s} thus reduces to the interpolation of a polynomial of degree N1s−1\frac{N_{1}}{s}-1 from its value at N1s\frac{N_{1}}{s} known, arbitrary points, to be repeated 2s2s times, which is the length of the convolution output. From , , it is known that the interpolation of a polynomial of degree N1s−1\frac{N_{1}}{s}-1 from its value at N1s\frac{N_{1}}{s} arbitrary points can be performed in O(N1s(log⁡N1s)2)\mathcal{O}(\frac{N_{1}}{s}(\log{\frac{N_{1}}{s}})^{2}). As this step is repeated 2s2s times for N2s\frac{N_{2}}{s} parts of vector x\bm{x}, the total decoding complexity is given by O((2s)(N1N2s2)(log⁡N1s)2)\mathcal{O}\left((2s)\left(\frac{N_{1}N_{2}}{s^{2}}\right)(\log{\frac{N_{1}}{s}})^{2}\right). Now observe that,

Note that, (24) holds from the choice of N1N2P≤s\sqrt{\frac{N_{1}N_{2}}{P}}\leq s which also implies N1s≤N1N2s2≤P\frac{N_{1}}{s}\leq\frac{N_{1}N_{2}}{s^{2}}\leq P.

Now, let us consider the complexity of the subsequent additions. Note that, we have a total of N1sN2s\frac{N_{1}}{s}\frac{N_{2}}{s} vectors to be added with appropriate shifts. As each vector is of length 2s2s, the complexity of adding any one vector is O(2s)\mathcal{O}(2s). Thus, the total complexity of the subsequent additions is given by \mathcal{O}\big{(}(2s)\frac{N_{1}}{s}\frac{N_{2}}{s}\big{)}. Now observe that,

Note that, (25) also holds from the choice of N1N2P≤s\sqrt{\frac{N_{1}N_{2}}{P}}\leq s.

Thus, from (21), (24) and (25), we can bound the total encoding and reconstruction complexity as follows:

Here DD is a constant. And (27) follows from the conditions of the theorem that P(log⁡P)2=o(2slog⁡N1N2P)P(\log{P})^{2}=o\left(2s\log{\sqrt{\frac{N_{1}N_{2}}{P}}}\right) and from the choice of N1N2P≤s\sqrt{\frac{N_{1}N_{2}}{P}}\leq s. Thus the total reconstruction complexity, i.e., the complexity of decoding and subsequent additions is o(2slog⁡(2s))o\left(2s\log{(2s)}\right). Recall that the per processor complexity for performing convolutions of two vectors of length ss is given by C(2s)log⁡(2s)C(2s)\log{(2s)}. Thus, the ratio of the decoding complexity to the per-processor time complexity tends to as N1,N2,P→∞N_{1},N_{2},P\to\infty. ∎

Appendix C Exponent for Repetition

In Section VI, we perform an asymptotic analysis of the exponent of the probability of failure to meet a deadline in the limit of the deadline diverging to infinity. Recall from Section VI that an upper bound for the leading coefficient of the failure exponent for Coded Convolution is given by

Here KK is the number of processors one needs to wait for in the worst case for Coded Convolution. We also showed that for uncoded strategy, this upper bound is actually exactly equal to the leading co-efficient of the failure exponent (putting worst case K=PK=P) as the deadline diverges to infinity. Now, we show that even for repetition strategy this upper bound is actually exactly equal to the upper bound using worst case KK as the deadline diverges to infinity.

Theorem 4: For a (P,r)(P,r) repetition strategy, the leading coefficient of the exponent in the probability of failure to meet a deadline converges as

where K=P−r+1K=P-r+1 is the number of processors to wait for in the worst case for a (P,r)(P,r) repetition strategy.

Recall that in a (P,r)(P,r) repetition strategy, we basically divide the operation of convolution into P/rP/r equal tasks, and each task has rr repetitions, so as to use all the given PP processors. Let us introduce the notation Ti,jT_{i,j} to represent the random variable corresponding to the different computation times on different processors. Here ii is the task index which varies from 11 to P/rP/r for the P/rP/r different parts of the convolution task. Also let jj be the replica index for each task which varies from 11 to rr for rr replicas of each task. Thus Ti,jT_{i,j} represents the computation time of the jj-th copy of the ii-th task.

Consider any task with index ii. Let Ti′T^{\prime}_{i} denote the completion time of the particular task. The task finishes when any one of its rr replicas finish. Thus, the distribution of Ti′T^{\prime}_{i} is given by

Now let us try to find the distribution of Ti′T^{\prime}_{i}. Recall that each Ti,jT_{i,j} is a random variable denoting the time required to perform a convolution of length ss at a processor where s=N1N2rPs=\sqrt{\frac{N_{1}N_{2}r}{P}}. Its computational complexity is given by 2Cslog⁡(2s)2Cs\log(2s). From shifted Weibull model assumptions, the c.d.f. of each Ti,jT_{i,j} is given by

Since the {Ti,j}\{T_{i,j}\}s are i.i.d., we thus have

Here c1(P,K)=(P/rP/r−1)c_{1}(P,K)=\binom{P/r}{P/r-1}, i.e., a function of PP and K(=P−r+1)K(=P-r+1) only and independent of tt. Similarly, c2(P,K)=∑l=0Pr−1(P/rl)c_{2}(P,K)=\sum_{l=0}^{\frac{P}{r}-1}\binom{P/r}{l} is also a function of PP and K(=P−r+1)K(=P-r+1) only and independent of tt. Thus,

Bounding the leading co-efficient in the failure exponent,

Appendix D Heuristic analysis of exponents

Statement: There exists an (N1,N2,P,s)(N_{1},N_{2},P,s) Coded Convolution strategy that outperforms the naive uncoded strategy in terms of asymptotic probability of failure if P(log⁡P)2=o(log⁡N1N2P)P(\log{P})^{2}=o\left(\log{\sqrt{\frac{N_{1}N_{2}}{P}}}\right) and α<2PN1N2(1+1log⁡2N1N2P)−1\alpha<2\sqrt{\frac{PN_{1}}{N_{2}}}\left(1+\frac{1}{\log{2\sqrt{\frac{N_{1}N_{2}}{P}}}}\right)^{-1}.

Justification: Consider the function given by

A higher value of E(s)E(s) implies faster decay of the failure exponent. Note that ss only takes values in the range N1N2P≤s≤min⁡{N1,N2}\sqrt{\frac{N_{1}N_{2}}{P}}\leq s\leq\min\{N_{1},N_{2}\}. Here s=N1N2Ps=\sqrt{\frac{N_{1}N_{2}}{P}} corresponds to the uncoded strategy, and any ss strictly greater than N1N2P\sqrt{\frac{N_{1}N_{2}}{P}} lies in the coded regime.

Now, we compare the exponents for uncoded and coded strategies.

Consider the derivative of E(s)E(s) in this range, as given by,

At s=N1N2Ps=\sqrt{\frac{N_{1}N_{2}}{P}} (Uncoded), the derivative is given by,

If E′(s)∣N1N2P>0E^{\prime}(s)|_{\sqrt{\frac{N_{1}N_{2}}{P}}}>0, then E(s)E(s) is strictly increasing, as ss increases beyond N1N2P\sqrt{\frac{N_{1}N_{2}}{P}}, i.e., into the coded regime. Thus, E(s)E(s) attains a maximum in the regime s>N1N2Ps>\sqrt{\frac{N_{1}N_{2}}{P}}, i.e., in the coded regime. Thus, it is sufficient that

Thus there exists an ss in the coded regime for which the asymptotic failure probability is lesser than that in the uncoded regime.

Note that, α=1\alpha=1, i.e., the shifted exponential also belongs to this regime. We also show some plots for the special case of N1=4NN_{1}=4N, N2=NN_{2}=N and P=log⁡NP=\sqrt{\log{N}}, showing that choosing an ss in the coded regime outperforms uncoded strategy for the chosen values of α\alpha in Fig. 5.