Fundamental Limits of Coded Linear Transform

Sinong Wang, Jiashang Liu, Ness Shroff, Pengyu Yang

I Introduction

Classical approaches of distributed linear transforms rely on dividing the input matrix A\mathbf{A} equally among all available worker nodes, and the master node has to collect the results from all workers to output y\mathbf{y}. As a result, a major performance bottleneck is the latency incurred in waiting for a few slow or faulty processors – called “stragglers” to finish their tasks . Recently, forward error correction and other coding techniques have shown to provide an effective way to deal with the “straggler” in the distributed computation tasks . By exploiting coding redundancy, the vector y\mathbf{y} is recoverable even if all workers have not finished their computation, thus reducing the delay due to straggle nodes. One main result in this paper is the development of optimal codes to deal with stragglers in computing distributed linear transforms, which provides an optimum recovery threshold and optimum computation redundancy.

What is the minimum recovery threshold and computation load for coded distributed linear transforms? Can we find an optimum computation strategy that achieves both optimality?

To that end, there have been works that have investigated the problems of reducing the recovery threshold and computational load for coded distributed linear transforms by leveraging ideas from coding theory. The first work proposes a maximum distance separable (MDS) code-based linear transform. We illustrate this approach using the following example. Consider a distributed system with 33 workers, the coding scheme first splits the matrix A\mathbf{A} into two submatrices, i.e., A=[A1;A2]\mathbf{A}=[\mathbf{A}_{1};\mathbf{A}_{2}]. Then each worker computes A1x\mathbf{A}_{1}\mathbf{x}, A2x\mathbf{A}_{2}\mathbf{x} and (A1+A2)x(\mathbf{A}_{1}+\mathbf{A}_{2})\mathbf{x}. The master node can compute Ax\mathbf{A}\mathbf{x} as soon as any 22 out of the 33 workers finish, and therefore achieves a recovery threshold of 22 and a computation load of 44. More specifically, it is shown that MDS code achieves a recovery threshold of Θ(m)\Theta(m) and computation load O(mn)O(mn). An improved scheme proposed in , referred to as the short MDS code, can offer a larger recovery threshold but with a constant reduction of computation load by imposing some sparsity of the encoded submatrices. More recently, the work in designs a type of polynomial code. It achieves the recovery threshold of nn that is independent of number of workers. However, it incurs a computation load of mnmn. In other words, each worker of polynomial code has to access the information of entire data matrix A\mathbf{A}.

In this paper, we show that, given number of partitions nn and number of stragglers ss, the minimum recovery threshold is nn and the minimum computation load is n(s+1)n(s+1). We design a novel coded computation strategy, referred to as the ss-diagonal code. It achieves both optimum recovery threshold and optimum computation load, which significantly improves the state-of-the-art. We also design a hybrid decoding algorithm between the peeling decoding and Gaussian elimination. We show that the decoding time is nearly linear time of the output dimension O(r)O(r).

Moreover, we define a new probabilistic recovery threshold performance metric that allows the coded computation scheme to provide the decodability in a large percentage (with high probability) of straggling scenarios. Based on this new metric, we propose two kinds of random codes, pp-Bernoulli code and (d1,d2)(d_{1},d_{2})-cross code, which can probabilistically guarantee the same recovery threshold but with much less computation load.

In the pp-Bernoulli code, each worker stores a weighted linear combination of several submatrices such that each Ai\mathbf{A}_{i} is chosen with probability pp independently. We show that, if p=2log⁡(n)/np=2\log(n)/n, the pp-Bernoulli code achieves the probabilistic recovery threshold of nn and a constant communication and computation load O(nlog⁡(n))O(n\log(n)) in terms of the number of stragglers. In the (d1,d2)(d_{1},d_{2})-cross code, each worker first randomly and uniformly chooses d1d_{1} submatrices, and each submatrix is randomly and uniformly assigned to d2d_{2} workers; then each worker stores a linear combination of the assigned submatrices and computes the corresponding linear transform. We show that, when the number of stragglers s=poly(log⁡(n))s=\text{poly}(\log(n)), the (2,3)(2,3)-cross code, or when s=O(nα),α<1s=O(n^{\alpha}),\alpha<1, the (2,2/(1−α))(2,2/(1-\alpha))-cross code achieves the probabilistic recovery threshold of nn and a much lower communication and computation load O(n)O(n). We compare our proposed codes and the existing codes in Fig. 1.

We finally implement and benchmark the constructed code in this paper on the Ohio Super Computing Center , and empirically demonstrate its performance gain compared with existing strategies.

II Problem Formulation

Given the above system model, we can formulate the coded distributed linear transform problem based on the following definitions.

This is a general definition of a large class of coded computation schemes. For example, in the polynomial code , the coding matrix M\mathbf{M} is the Vandermonde matrix. In the MDS type of codes , it is a specific form of corresponding generator matrix.

(Computation load) The computation load of strategy M\mathbf{M} is defined as l(M)=∥M∥0l(\mathbf{M})=\|\mathbf{M}\|_{0}, the number of the nonzero elements of coding matrix.

III Diagonal Code and its Optimality

In this section, we first describe the optimum recovery threshold and optimum computation load. Then we introduce the ss-diagonal code that exactly matches such lower bounds.

We first establish the information theoretical lower bound of the recovery threshold.

The proof is similar with the one in that uses the cut-set type argument around the master node. This bound implies that the optimum recovery threshold is independent of the number of workers. We next establish the lower bound of the computation load, i.e., the density of coding matrix M\mathbf{M}.

The polynomial code achieves the optimum recover threshold of nn. However, the coding matrix (Vandermonde matrix) is fully dense, i.e., l(M)=nml(\mathbf{M})=nm, and far beyond the above bound. The short MDS code can reduce the computation load but sacrifices the recovery threshold. Therefore, a natural question that arises is, can we design a code achieves both optimum recovery threshold and optimum computation load? We will answer this question in the sequel of this paper.

III-B Diagonal Code

Now we present the ss-diagonal code that achieves both the optimum recovery threshold and optimum computation load for any given parameter values of nn, mm and ss.

(ss-diagonal code) Given parameters mm,nn and ss, the ss-diagonal code is defined as

where each coefficient mijm_{ij} is chosen from a finite set SS independently and uniformly at random.

The reason we name this code as the ss-diagonal code is that the nonzero positions of the coding matrix M\mathbf{M} show the following diagonal structure.

where ∗* indicates nonzero entries of M\mathbf{M}. Before we analyze the recovery threshold of the diagonal code, the following example provides an instantiation of the above construction for n=4n=4, s=1s=1 and m=5m=5.

Example 1: (11-Diagonal code) Consider a distributed linear transform task Ax\mathbf{Ax} using m=5m=5 workers. We evenly divide the matrix A\mathbf{A} along the row side into 44 submatrices: AT=[A1T,A2T,A3T,A4T]\mathbf{A}^{T}=[\mathbf{A}_{1}^{T},\mathbf{A}_{2}^{T},\mathbf{A}_{3}^{T},\mathbf{A}_{4}^{T}]. Given this notation, we need to compute the following 44 uncoded components.

Based on the definition of 11-diagonal code, each worker stores the following submatrices,

Suppose that the first worker is straggler and master node receives results from worker {2,3,4,5}\{2,3,4,5\}. According to the above coded computation strategy, we have

The coefficient matrix is an upper diagonal matrix, which is invertible since the elements in the main diagonal are nonzero. Then we can recover the uncoded components {Aix}\{\mathbf{A}_{i}\mathbf{x}\} by direct inversion of the above coefficient matrix. The decodability for the other 44 possible scenarios can be proved similarly. Therefore, this code achieves the optimum recovery threshold of 44. Besides, the total number of partial linear transforms is 88, which also matches the lower bound of computation load.

III-C Optimality of Diagonal Code

The optimality of computation load of the ss-diagonal code can be easily obtained by counting the number of nonzero elements of coding matrix M\mathbf{M}. In terms of recovery threshold, we have the following result of the diagonal code

Let finite set SS satisfy ∣S∣≥2n2Cmn|S|\geq 2n^{2}C^{n}_{m}, then the s−s-diagonal code achieves the recovery threshold of nn.

For each subset U⊆[m]U\subseteq[m] with ∣U∣=n|U|=n, let MU\mathbf{M}^{U} be a n×nn\times n submatrices consisting rows of M\mathbf{M} index by UU. To prove ss-diagonal code achieves the optimum recovery threshold of nn, we need to show all the n×nn\times n submatrices MU\mathbf{M}^{U} are full rank. We first define the following bipartite graph model between the set of mm workers [m][m] and the set of nn data partitions [n][n].

Based on the above definition, the coding matrix M\mathbf{M} of ss-diagonal code can be obtained by assigning each intermediate xijx_{ij} of Edmonds matrix M(x)\mathbf{M}(\mathbf{x}) a value from set SS independently and uniformly at random. Given the subset U⊆[m]U\subseteq[m] with ∣U∣=n|U|=n, define GD(U,V2)G^{D}(U,V_{2}) as a subgraph of GD(V1,V2)G^{D}(V_{1},V_{2}) and MU(x)\mathbf{M}^{U}(\mathbf{x}) as the corresponding Edmonds matrix. Then the probability that matrix MU\mathbf{M}^{U} is full rank is equal to the probability that the determinant of the Edmonds matrix MU(x)\mathbf{M}^{U}(\mathbf{x}) is nonzero at the given value x\mathbf{x}. The following technical lemma provides a simple lower bound of the such an event.

A classic result in graph theory is that a bipartite graph GD(U,V2)G^{D}(U,V_{2}) contains a perfect matching if and only if the determinant of Edmonds matrix, i.e., ∣MU(x)∣|\mathbf{M}^{U}(\mathbf{x})|, is a nonzero polynomial. Combining this result with Schwartz-Zeppel Lemma, we can finally reduce the analysis of the full rank probability of the submatrix MU\mathbf{M}^{U} to the probability that the subgraph GD(U,V2)G^{D}(U,V_{2}) contains a perfect matching.

The next technical lemma shows that for all subsets U⊆[m]U\subseteq[m], the subgraph GD(U,V2)G^{D}(U,V_{2}) exactly contains a perfect matching. Then utilizing the union bound, we conclude that, with probability at least 1/21/2, all the n×nn\times n submatrices of M\mathbf{M} are full rank. Since we can generate the coding matrix M\mathbf{M} offline, with a few rounds of trials (22 in average), we can find a coding matrix with all n×nn\times n submatrices being full rank.

Construct bipartite graph GD(V1,V2)G^{D}(V_{1},V_{2}) from Definition 5. For each U⊆[m]U\subseteq[m] with ∣U∣=n|U|=n, the subgraph GD(U,V2)G^{D}(U,V_{2}) contains a perfect matching.

The proof is mainly based on Hall’s marriage theorem. One special case of ss-diagonal code is that all the nonzero elements of matrix M\mathbf{M} can be equal to 11, when we are required to resist only one straggler.

Given the parameters nn and m=n+1m=n+1, define the 11-diagonal code:

It achieves the optimum computation load of 2n2n and optimum recovery threshold of nn.

The ss-diagonal code achieves both optimum recovery threshold of nn and computation load of n(s+1)n(s+1). This result implies that the the average computation load, i.e., n(s+1)/(n+s)n(s+1)/(n+s), will increase when number of stragglers increases. In the extreme case, the coding matrix M\mathbf{M} will degenerate to a fully dense matrix. The reason behind this phenomenon mainly derives from the pessimistic definition of the recovery threshold, which requires the coded computation scheme M\mathbf{M} to resist any ss stragglers. In the next section, we show that, if we relax requirement in Definition 2 to the probabilistic scenario, we can design a code with a computation load independent of the number of stragglers ss.

IV Random Code: “Break” the Limits

In practice, the straggles randomly occur in each worker, and the probability that a specific straggling configuration happen is in low probability. We first demonstrate the main idea through the following motivating examples.

Example 2: Consider a distributed linear transform tasks with n=20n=20 data partitions, s=5s=5 stragglers and m=25m=25 workers. The recovery threshold of 2020 implies that all C2520=53130C_{25}^{20}=53130 square 20×2020\times 20 submatrices are full rank. Suppose that a worker being a straggler is identically and independently Bernoulli random variable with probability 10%10\%. Then, the probability that workers {1,2,3,4,5}\{1,2,3,4,5\} being straggler is 10−510^{-5}. if there exists a scheme that can guarantee the master to decode the results in all configurations except straggling configurations {1,2,3,4,5}\{1,2,3,4,5\}, we can argue that this scheme achieves a recovery threshold of 2020 with probability 1−10−51-10^{-5}.

Based on the above two examples, the computation load can be reduced when we allow the coded computation strategy fails in some specific scenarios. Formally, we have the following definition of the probabilistic recovery threshold.

Instead of guaranteeing the recoverability of all the straggling configurations, new definition provides a probabilistic relaxation such that a small percentage of all the straggling configurations, i.e., vanishes as n→∞n\rightarrow\infty, are allowed to be unrecoverable. In the sequel, we show that, under such a relaxation, one can construct a coded computation scheme that achieves probabilistic recovery threshold of nn and a constant (regarding parameter ss) computation load.

IV-B Construction of Random Code

Based on our previous analysis of the recovery threshold of s−s-diagonal code, we show that, for any subset U⊆[m]U\subseteq[m] with ∣U∣=n|U|=n, the probability that MU\mathbf{M}^{U} is full rank is lower bounded by the probability (multiplied by 1−o(1)1-o(1)) that the corresponding subgraph G(U,V2)G(U,V_{2}) contains a perfect matching. This technical path motivates us to utilize the random proposal graph to construct the coded computation scheme. The first one is the following pp-Bernoulli code, which is constructed from the ER random bipartite graph model .

(pp-Bernoulli code) Given parameters m,nm,n, construct the coding matrix M\mathbf{M} as follows:

where tijt_{ij} is picked independently and uniformly from the finite set SS.

For any parameters m,nm,n and ss, if p=2log⁡(n)/np=2\log(n)/n, the pp-Bernoulli code achieves the probabilistic recovery threshold of nn.

This result implies that each worker of pp-Bernoulli code requires accessing 2log⁡(n)2\log(n) submatrices in average, which is independent of number of stragglers. Note that the existing works in distributed functional computation proposes a random sparse code that also utilizes Bernoulli random variable to construct the coding matrix. There exist two key differences: (i) the elements of our matrix can be integer valued, while the random sparse code adopts real-valued matrix; (ii) the density of pp-Bernoulli code is 2log⁡(n)2\log(n), while the density of the random sparse code is an unknown constant.

The second random code is the following (d1,d2)(d_{1},d_{2})-cross code.

((d1,d2)(d_{1},d_{2})-cross code) Given parameters m,nm,n, construct the coding matrix M\mathbf{M} as follows: (1) Each row (column) independently and uniformly choose d1d_{1} (d2d_{2}) nonzero positions; (2) For those nonzero positions, assign the value independently and uniformly from the finite set SS.

It is easy to see that the computation load of (d1,d2)(d_{1},d_{2})-cross code is upper bounded by d1m+d2nd_{1}m+d_{2}n. The next theorem shows that a constant choice of d1d_{1} and d2d_{2} can guarantee the probabilistic recovery threshold of nn.

For any parameters m,nm,n and ss, if s=poly(log⁡(n))s=\text{poly}(\log(n)), the (2,3)(2,3)-cross code achieves the probabilistic recovery threshold of nn. If s=Θ(nα)s=\Theta(n^{\alpha}), α<1\alpha<1, the (2,2/(1−α))(2,2/(1-\alpha))-cross code achieves the probabilistic recovery threshold of nn.

The proof of this theorem is based on analyzing the existence of perfect matching in a random bipartite graph constructed as following: (i) each node in the left partition randomly and uniformly connects to d1d_{1} nodes in the opposite class; (ii) each node in the right partition randomly and uniformly connects to ll nodes in the opposite class, where ll is chosen under a specific degree distribution. This random graph model can be regarded as a generalization of Walkup’s 22-out bipartite graph model . The main technical difficulty in this case derives from the intrinsic complicated statistical model of node degree of right partition.

One straightforward construction is using the LDPC codes or rateless codes such as LT and raptor code . These codes can not only guarantee the probabilistic recovery threshold of Θ(n)\Theta(n), but also provide a low computation load, i.e., O(nlog⁡(1/ϵ))O(n\log(1/\epsilon)) from raptor code. However, the desired performance of these codes requires the nn sufficiently large, i.e., n>105n>10^{5}, which is impossible in distributed computation. For example, when n=20n=20, the recovery threshold of LT code is roughly equal to 3030, which is much larger than the lower bound.

IV-C Numerical Results of Random Code

We examine the performance of proposed pp-Bernoulli code and (d1,d2)(d_{1},d_{2})-cross code in terms of convergence speed of full rank probability and computation load. In Fig. 4, we plot the percentage of full rank n×nn\times n square submatrix and the average computation load l(M)/ml(\mathbf{M})/m of each scheme, based on 10001000 experimental run. Each column of the (2,2.5)(2,2.5)-code independently and randomly chooses 22 or 33 nonzero elements with equal probability. It can be observed that the full rank probability of both pp-Bernoulli code and (d1,d2)(d_{1},d_{2})-cross code converges to 11 quite fast. The (d1,d2)(d_{1},d_{2})-cross code exhibits even faster convergence and much less computation load compared to the pp-Bernoulli code. For example, when n=20,s=4n=20,s=4, the (2,2)(2,2)-cross code achieves the full rank probability of 0.860.86 and average computation load 3.43.4. This provides the evidence that (d1,d2)(d_{1},d_{2})-cross code is useful in practice. Moreover, in the practical use of random code, one can run multiple rounds of trails to find a “best” coding matrix with even higher full rank probability and lower computation load.

V Fast Decoding Algorithm

(rooting step) If rank(MU)=n(\mathbf{M}^{U})=n, then for any k0∈{1,2,…,n}k_{0}\in\{1,2,\ldots,n\}, we can recover a particular block Aix\mathbf{A}_{i}\mathbf{x} with column index k0k_{0} in matrix MU\mathbf{M}^{U} via the following linear combination.

Moreover, we can show that, for the ss-diagonal code, when the position of blocks recovered from rooting step (13) is specifically designed, the Algorithm 1 requires at most ss rooting steps. This result provides an O(rs)O(rs) time decoding algorithm for ss-diagonal code.

(Nearly linear time decoding of ss-diagonal code) For ss-diagonal code constructed in Definition 4 and any nn received results indexed by U={i1,…,in}⊆[n+s]U=\{i_{1},\ldots,i_{n}\}\subseteq[n+s] and 1≤i1≤⋯≤in≤n+s1\leq i_{1}\leq\cdots\leq i_{n}\leq n+s. Let the index k∈[n]k\in[n] be ik≤n<ik+1i_{k}\leq n<i_{k+1}. Recover the blocks indexed by [n]\{i1,…,ik}[n]\backslash\{i_{1},\ldots,i_{k}\} from rooting steps (13). Then the Algorithm 1 uses the at most ss times rooting steps and the total decoding time is O(rs)O(rs).

VI Experimental Results

In this section, we present the experimental results at OSC . We compare our proposed coding schemes including s−s-diagonal code and (d1,d2)(d_{1},d_{2})-cross code against the following existing schemes in both single matrix vector multiplication and gradient descent: (i) uncoded scheme: the input matrix is divided uniformly across all workers without replication and the master waits for all workers to send their results; (ii) sparse MDS code : the generator matrix is a sparse random Bernoulli matrix with average computation overhead Θ(log⁡(n))\Theta(\log(n)). (iii) polynomial code : coded matrix multiplication scheme with optimum recovery threshold and nearly linear decoding time; (iv) short dot code : append the dummy vectors to data matrix A\mathbf{A} before applying the MDS code, which provides some sparsity of encoded data matrix with cost of increased recovery threshold. (v) LT code : rateless code widely used in broadcast communication. It achieves average computation overhead Θ(log⁡(n))\Theta(\log(n)) and a nearly linear decoding time using peeling decoder. To simulate straggler effects in large-scale system, we randomly pick ss workers that are running a background thread.

We first use a matrix with r=t=1048576r=t=1048576 and nnz(A)=89239674(\mathbf{A})=89239674 from data sets , and evenly divide this matrix into n=12n=12 and 2020 partitions. In Fig. 5, we report the job completion time under s=2s=2 and s=4s=4, based on 2020 experimental runs. It can be observed that both (2,2)(2,2)-cross code outperforms uncoded scheme (in 50% the time), LT code (in 80% the time), sparse MDS code (in 60% the time), polynomial code (in 20% the time) and our ss-diagonal code. Moreover, we have the following observations: (i)LT code performs better than sparse MDS code when number of stragglers ss is small, and worse when ss increases. (ii) the uncoded scheme is faster than the polynomial code, because the density of encoded data matrix is greatly increased, which leads to increased computation time per worker and additional I/O contention at the master node.

We further compare our proposed ss-diagonal code with (2,2)(2,2)-cross code versus the number of stragglers ss. As shown in Fig. 5, when the number of stragglers ss increases, the job completion time of ss-diagonal code increases while the (2,2)(2,2)-cross code does not change. If the number of stragglers is smaller than 22, the ss-diagonal code performs better than (2,2)(2,2)-cross code. Another interesting observation is that the irregularity of work load can decrease the I/O contention. For example, when s=2s=2, the computation load of 22-diagonal code is similar as (2,2)(2,2)-cross code, which is equal to 3636 in the case of n=12n=12. However, the (2,2)(2,2)-cross code costs less time due to the unbalanced worker load.

VI-B Coded Gradient Descent

We first describe the gradient-based distributed algorithm to solve the following linear regression problem.

In the uncoded gradient descent, each worker ii first stores a submatrix Ai\mathbf{A}_{i}; then during iteration tt, each worker ii first computes a vector Ai⊺Aixt\mathbf{A}_{i}^{\intercal}\mathbf{A}_{i}\mathbf{x}_{t} and returns it to the master node; the master node updates the weight vector xt\mathbf{x}_{t} according to the above gradient descent step and assigns the new weight vector xt+1\mathbf{x}_{t+1} to each worker. The algorithm terminates until the gradient vanish, i.e., ∥A⊺(Axt−b)∥≤ϵ\|\mathbf{A}^{\intercal}(\mathbf{Ax}_{t}-\mathbf{b})\|\leq\epsilon. In the coded gradient descent, each worker ii first stores several submatrices according to the coding matrix M\mathbf{M}; during iteration tt, each worker ii computes a linear combinations,

where mijm_{ij} is the element of the coding matrix M\mathbf{M}. The master node collects a subset of results, decodes the full gradient A⊺(Axt−b)\mathbf{A}^{\intercal}(\mathbf{Ax}_{t}-\mathbf{b}) and updates the weight vector xt\mathbf{x}_{t}. Then continue to the next round.

We use a data from LIBSVM dataset repository with r=19264097r=19264097 samples and t=1163024t=1163024 features. We evenly divide the data matrix A\mathbf{A} into n=12,20n=12,20 submatrices. In Fig. 6 , we plot the magnitude of scaled gradient ∥ηA⊺(Ax−b)∥2\|\eta\mathbf{A}^{\intercal}(\mathbf{Ax}-\mathbf{b})\|^{2} versus the running times of the above seven different schemes under n=12,20n=12,20 and s=2,4s=2,4. Among all experiments, we can see that (2,2)(2,2)-cross code converges at least 20%20\% faster than sparse MDS code, 22 times faster than both uncoded scheme and LT code and at least 44 times faster than short dot and polynomial code. The (2,2)(2,2)-cross code performs similar with ss-diagonal code when s=2s=2 and converges 30%30\% faster than ss-diagonal code when s=4s=4.

References

-C Proof of Theorem 2

Given the parameter mm and nn, considering any coded computation scheme that can resist ss stragglers. We first define the following bipartite graph model between the mm workers [m][m] and nn data partitions [n][n], where we connect node i∈[m]i\in[m] and node j∈[n]j\in[n] if mij≠0m_{ij}\neq 0 (worker ii has access to data block Aj\mathbf{A}_{j}). The degree of worker node i∈[m]i\in[m] is ∥Mi∥0\|\mathbf{M}_{i}\|_{0}. We show that the degree of node j∈[n]j\in[n] must be at least s+1s+1. Suppose that it is less than s+1s+1 and all its neighbors are stragglers. In this case, there exists no worker that is non-straggler and also access to Aj\mathbf{A}_{j}(or the corresponding submatrix of coding matrix is rank deficient). It is contradictory to the assumption that it can resist ss stragglers.

Based on the above argument and the fact that, the sum of degrees of one partition is equal to the sum of degrees in another partition in the bipartite graph, we have

-D Proof of Lemma 2

A direct application of Hall’s theorem is: given a bipartite graph GD(V1,V2)G^{D}(V_{1},V_{2}), for each U⊆[m]U\subseteq[m] with ∣U∣=n|U|=n, each subgraph GD(U,V2)G^{D}(U,V_{2}) contains a prefect matching if and only if every subset set S⊆V1S\subseteq V_{1} such that ∣N(S)∣<∣S∣|N(S)|<|S|, where the neighboring set N(S)N(S) is defined as N(S)={y∣x,y are connected for some x∈S}N(S)=\{y|x,y\text{ are connected for some }x\in S\}. This result is equivalent to the following condition: for each subset I⊆[m]I\subseteq[m],

where supp(Mi)(\mathbf{M_{i}}) is defined as the support set: supp(Mi)={j∣mij≠0,j∈[n]}(\mathbf{M_{i}})=\{j|m_{ij}\neq 0,j\in[n]\}, Mi\mathbf{M_{i}} is iith row of the coding matrix M\mathbf{M}. Suppose that the set I={i1,i2,…,ik}I=\{i_{1},i_{2},\ldots,i_{k}\} with i1<i2<…<iki_{1}<i_{2}<\ldots<i_{k} and supp(Mik1)∩supp(Mik2)≠∅\text{supp}(\mathbf{M}_{i_{k_{1}}})\cap\text{supp}(\mathbf{M}_{i_{k_{2}}})\neq\emptyset. Otherwise, we can divide the set into two parts Il={i1,i2,…,ik1}I_{l}=\{i_{1},i_{2},\ldots,i_{k_{1}}\} and Ir={ik2,i2,…,ik}I_{r}=\{i_{k_{2}},i_{2},\ldots,i_{k}\} and prove the similar result in these two sets. Based on our construction of diagonal code, we have

The above, step (a) is based on the fact that ik−i1≥ki_{k}-i_{1}\geq k. Therefore, the lemma follows.

-E Proof of Corollary 1

To prove 11-diagonal code achieves the recovery threshold of nn, we need to show that, for each subset U⊆[n+1]U\subseteq[n+1] with ∣U∣=n|U|=n, submatrix MU\mathbf{M}^{U} is full rank. Let U=[n+1]\{k}U=[n+1]\backslash\{k\}, the submatrix MU\mathbf{M}^{U} satisfies

where E\mathbf{E} is a (k−1)(k-1) dimensional square submatrix consists of first (k−1)(k-1) rows and columns, and F\mathbf{F} is a (n−k+1)(n-k+1) dimensional square submatrix consists of the last (n−k+1)(n-k+1) rows and columns. The matrix E\mathbf{E} is a lower diagonal matrix due to the fact that, for i<ji<j,

The matrix E\mathbf{E} is a upper diagonal matrix due to the fact that, for i>ji>j,

The above, (a) utilizes the fact that (i+k)−(j+k−1)≥2(i+k)-(j+k-1)\geq 2 when i>ji>j. Based on the above analysis, we have

which implies that matrix MU\mathbf{M}^{U} is full rank. Therefore, the corollary follows.

-F Proof of Theorem 4

Based on our analysis in Section III-C, the full rank probability of an n×nn\times n submatrix MU\mathbf{M}^{U} can be lower bounded by a constant multiplying the probability of the existence of perfect matching in a bipartite graph.

Therefore, to prove that the pp-Bernoulli code achieves the probabilistic recovery threshold of nn, we need to show that each subgraph contains a perfect matching with high probability. Without loss of generality, we can define the following random bipartite graph model.

(pp-Bernoulli random graph) Graph Gb(U,V2,p)G^{b}(U,V_{2},p) initially contains the isolated nodes with ∣U∣=∣V2∣=n|U|=|V_{2}|=n. Each node v1∈Uv_{1}\in U and node v2∈V2v_{2}\in V_{2} is connected with probability pp independently.

Clearly, the above model describes the support structure of each submatrix MU\mathbf{M}^{U} of pp-Bernoulli code. The rest is to show that, with specific choice of pp, the subgraph Gb(U,V2,p)G^{b}(U,V_{2},p) contains a perfect matching with high probability.

The technical idea is to use the Hall’s theorem. Assume that the bipartite graph Gb(U,V2,p)G^{b}(U,V_{2},p) does not have a perfect matching. Then by Hall’s condition, there exists a violating set S⊆US\subseteq U or S⊆V2S\subseteq V_{2} such that ∣N(S)∣<∣S∣|N(S)|<|S|. Formally, by choosing such SS of smallest cardinality, one immediate consequence is the following technical statement.

If the bipartite graph Gb(U,V2,p)G^{b}(U,V_{2},p) does not contain a perfect matching, then there exists a set S⊆US\subseteq U or S⊆V2S\subseteq V_{2} with the following properties.

For each node t∈N(S)t\in N(S), there exists at least two adjacent nodes in SS.

Case 1: We consider S⊆US\subseteq U and ∣S∣=1|S|=1. In this case, we have ∣N(S)∣=0|N(S)|=0 and need to estimate the probability that there exists one isolated node in partition UU. Let random variable XiX_{i} be the indicator function of the event that node viv_{i} is isolated. Then we have the probability that

Let XX be the total number of isolated nodes in partition UU. Then we have

The above, step (a) utilizes the assumption that p=2log⁡(n)/np=2\log(n)/n and the inequality that (1+x/n)n≤ex(1+x/n)^{n}\leq e^{x}.

Case 2: We consider S⊆US\subseteq U and 2≤∣S∣≤n/22\leq|S|\leq n/2. Let EE be the event that such an SS exists, we have

The above, step (a) is based on the inequality

The step (b) utilizes the fact that p=2log⁡(n)/np=2\log(n)/n, k≤n/2k\leq n/2 and the inequality (1+x/n)n≤ex(1+x/n)^{n}\leq e^{x}; step (c) is based on the fact that k(n−k+1)≥2(n−1)k(n-k+1)\geq 2(n-1) and k(n−k)≥2(n−2)k(n-k)\geq 2(n-2), n/(n−2)<5/3n/(n-2)<5/3 for k≥2k\geq 2 and n≥5n\geq 5. Utilizing the union bound to sum the results in case 1 and case 2, we can obtain that the probability that graph G(U,V2,p)G(U,V_{2},p) contains a perfect matching is at least

Therefore, incorporating this result into estimating (11), the theorem follows.

-G Proof of Theorem 5

To prove the (d1,d2)(d_{1},d_{2})-cross code achieves the probabilistic recovery threshold of nn, we need to show that each subgraph of the following random bipartite graph contains a perfect matching with high probability.

((d1,d2)(d_{1},d_{2})-regular random graph) Graph Gc(V1,V2,d1,d2)G^{c}(V_{1},V_{2},d_{1},d_{2}) initially contains the isolated nodes with ∣V1∣=m|V_{1}|=m and ∣V2∣=n|V_{2}|=n. Each node v1∈V1v_{1}\in V_{1} (v2∈V2v_{2}\in V_{2}) randomly and uniformly connects to d1d_{1} (d2d_{2}) nodes in V1V_{1} (V2V_{2}).

The corresponding subgraph is defined as follows.

For each U⊆V1U\subseteq V_{1} with ∣U∣=n|U|=n, the subgraph Gc(U,V2,d1,dˉ)G^{c}(U,V_{2},d_{1},\bar{d}) is obtained by deleting the nodes in V1\UV_{1}\backslash U and corresponding arcs.

Clearly, the above definitions of (d1,d2)(d_{1},d_{2})-regular graph and corresponding subgraph describe the support structure of the coding matrix M\mathbf{M} and submatrix MU\mathbf{M}^{U} of the (d1,d2)(d_{1},d_{2})-cross code. Moreover, we have the following result regarding the structure of each subgraph Gc(U,V2,d1,dˉ)G^{c}(U,V_{2},d_{1},\bar{d})

Claim. For each U⊆V1U\subseteq V_{1} with ∣U∣=n|U|=n, the subgraph Gc(U,V2,d1,dˉ)G^{c}(U,V_{2},d_{1},\bar{d}) can be constructed from the following procedure" (i) Initially, graph Gc(U,V2,d1,dˉ)G^{c}(U,V_{2},d_{1},\bar{d}) contain the isolated nodes with ∣U∣=∣V2∣=n|U|=|V_{2}|=n; (ii) Each node v1∈Uv_{1}\in U randomly and uniformly connects to d1d_{1} nodes in V2V_{2}; (iii) Each node v2∈V2v_{2}\in V_{2} randomly and uniformly connects to ll nodes in V1V_{1}, where ll is chosen according to the distribution:

Then, the rest is to show that the subgraph Gc(U,V2,d1,dˉ)G^{c}(U,V_{2},d_{1},\bar{d}) contains a perfect matching with high probability.

(Forbidden kk-pair) For a bipartite graph G(U,V2,d1,dˉ)G(U,V_{2},d_{1},\bar{d}), a pair (A,B)(A,B) is called a kk-blocking pair if A⊆UA\subseteq U with ∣A∣=k|A|=k, B⊆V2B\subseteq V_{2} with ∣B∣=n−k+1|B|=n-k+1, and there exists no arc between the nodes of sets AA and BB. A blocking kk-pair (A,B)(A,B) is called a forbidden pair if at least one of the following holds:

2≤k<(n+1)/22\leq k<(n+1)/2, and for any v1∈Av_{1}\in A and v2∈V2\Bv_{2}\in V_{2}\backslash B, (A\{v1},B∪{v2})(A\backslash\{v_{1}\},B\cup\{v_{2}\}) is not a (k−1)(k-1)-blocking pair.

(n+1)/2≤k≤n−1(n+1)/2\leq k\leq n-1, and for any v1∈U\Av_{1}\in U\backslash A and v2∈Bv_{2}\in B, (A∪{v1},B\{v2})(A\cup\{v_{1}\},B\backslash\{v_{2}\}) is not a (k+1)(k+1)-blocking pair.

The following technical lemma modified from is useful in our proof.

If the graph Gc(U,V2,d1,dˉ)G^{c}(U,V_{2},d_{1},\bar{d}) does not contain a perfect matching, then there exists a forbidden kk-pair for some kk.

One direct application of the Konig’s theorem to bipartite graph shows that Gc(U,V2,d1,dˉ)G^{c}(U,V_{2},d_{1},\bar{d}) contains a perfect matching if and only if it does not contain any blocking kk-pair. It is rest to show that the existence of a kk-blocking pair implies that there exists a forbidden l−l-pair for some ll. Suppose that there exists a kk-blocking pair (A,B)(A,B) with k<(n+1)/2k<(n+1)/2, and it is not a forbidden kk-pair. Otherwise, we already find a forbidden pair. Then, it implies that there exists v1∈Av_{1}\in A and v2∈V2\Bv_{2}\in V_{2}\backslash B such that (A\{v1},B∪{v2})(A\backslash\{v_{1}\},B\cup\{v_{2}\}) is a (k−1)(k-1)-blocking pair. Similarly, we can continue above argument on blocking pair (A\{v1},B∪{v2})(A\backslash\{v_{1}\},B\cup\{v_{2}\}) until we find a forbidden pair. Otherwise, we will find a 1−1-blocking pair (A′,B′)(A^{\prime},B^{\prime}), which is a contradiction to our assumption that each node v1∈Uv_{1}\in U connects d1d_{1} nodes in V2V_{2}. The proof for k≥(n+1)/2k\geq(n+1)/2 is same. ∎

Let EE be the event that graph G(U,V2,d1,dˉ) contains perfect matchingG(U,V_{2},d_{1},\bar{d})\text{ contains perfect matching}. Based on the the results of Lemma 5, we have

The above, AA and BB are defined as node sets such that A⊆UA\subseteq U with ∣A∣=k|A|=k and B⊆V2B\subseteq V_{2} with ∣B∣=n−k+1|B|=n-k+1. The α(k)\alpha(k) and β(k)\beta(k) are defined as follows.

From the Definition 12, it can be obtained the following estimation of probability β(k)\beta(k).

The first factor gives the probability that there exists no arc from nodes of AA to nodes of BB. The second factor gives the probability that there exists no arc from nodes of BB to nodes of AA. The summation operation in the second factor comes from conditioning such probability on the distribution (29). Based on the Chu-Vandermonde identity, one can simplify β(k)\beta(k) as

In the above, step (a) is based on fact that (1+x/n)n≤ex(1+x/n)^{n}\leq e^{x} and parameters c1c_{1} and c2c_{2} are defined as

Combining the equations (35)-(38), we can obtain that

The constant c3c_{3} is given by c3=ed12+13d1/12+d2+1/3/2πc_{3}=e^{d_{1}^{2}+13d_{1}/12+d_{2}+1/3}/2\pi. The third term satisfies

which is based on the fact that, if d2>2d_{2}>2, the function

is monotonically decreasing when k≤(d2n−2m)/(d2−2)k\leq(d_{2}n-2m)/(d_{2}-2) and increasing when k≥(d2n−2m)/(d2−2)k\geq(d_{2}n-2m)/(d_{2}-2). If d2=2d_{2}=2, it is monotonically increasing for k≥0k\geq 0.

We then estimate the conditional probability α(k)\alpha(k). Given a blocking pair A⊆UA\subseteq U with ∣A∣=k|A|=k and B⊆V2B\subseteq V_{2} with ∣B∣=n−k+1|B|=n-k+1, and a node vi∈Av_{i}\in A, let EiE_{i} be the set of nodes in V2\BV_{2}\backslash B on which d1d_{1} arcs from node viv_{i} terminate. Let E′E^{\prime} be the set of nodes vv in V2\BV_{2}\backslash B such that at least 22 arcs leaving from vv to nodes in AA. Then we have the following technical lemma.

Given a blocking pair A⊆UA\subseteq U with ∣A∣=k|A|=k and B⊆V2B\subseteq V_{2} with ∣B∣=n−k+1|B|=n-k+1, if (A,B)(A,B) is kk-forbidden pair, then

Suppose that there exists node v∈V2\(E∗∪B)v\in V_{2}\backslash(E^{*}\cup B), then there exists no arc from AA to vv and there exists at most 11 arc from vv to AA. If such an arc exists, let v′v^{\prime} be the corresponding terminating node in AA. Then we have (A\{v′},B∪{v})(A\backslash\{v^{\prime}\},B\cup\{v\}) is a blocking pair, which is contradictory to the definition of forbidden pair. If such an arc does not exist, let v′v^{\prime} be the an arbitrary node in AA. Then we have (A\{v′},B∪{v})(A\backslash\{v^{\prime}\},B\cup\{v\}) is a blocking pair, which is also contradictory to the definition of forbidden pair. ∎

The lemma 6 implies that we can upper bound the conditional probability by

where P1P_{1} and P2P_{2} is defined as: for any node v∈V2\Bv\in V_{2}\backslash B,

The above, step (a) utilizes Chu-Vandermonde identity twice; step (b) is adopts the inequality (34); step (c) is based on the fact that if nn is sufficiently large, (1−x/n)n≥c5e−x(1-x/n)^{n}\geq c_{5}e^{-x}, where c5c_{5} is a constant. Combining the above estimation of P1P_{1} and P2P_{2}, we have the following upper bound of α(k)\alpha(k).

We finally estimate the probability that the graph Gc(U,V2,d1,dˉ)G^{c}(U,V_{2},d_{1},\bar{d}) contains a perfect matching under the following two cases.

Case 1: The number of stragglers s=poly(log⁡(n))s=\text{poly}(\log(n)). Let d1=2,d2=3d_{1}=2,d_{2}=3. Based on the estimation (40), we have

Combining the above results with the estimation of β(k)\beta(k), we have

The above, step (a) utilizes the estimation of α(k)\alpha(k) in (43); step (b) is based on estimating the tail of the summation as geometric series (1−c6/e2)k−1(1-c_{6}/e^{2})^{k-1}.

Case 2: The number of stragglers s=Θ(nα)s=\Theta(n^{\alpha}), α<1\alpha<1. Let d1=2,d2=2/(1−α)d_{1}=2,d_{2}=2/(1-\alpha). For 2≤k≤n−22\leq k\leq n-2, we have

Combining the above results with the estimation of β(k)\beta(k), we have

If α≥1/2\alpha\geq 1/2, we can directly obtain that γ(k)≤c3e−2n\gamma(k)\leq\frac{c_{3}e^{-2}}{n}. If α<1/2\alpha<1/2, we have

Therefore, in both cases, incorporating the above results into estimating (11), the theorem follows.

-H Proof of Theorem 6

Based on our construction, the cardinality of the set

Since n<ik+1<⋯<in≤n+sn<i_{k+1}<\cdots<i_{n}\leq n+s, we have n−k≤sn-k\leq s. We next show, after we recover the blocks indexed by [n]\{i1,…,ik}[n]\backslash\{i_{1},\ldots,i_{k}\}, the rest blocks can be recovered by peeling decoding without rooting steps. Combining these results together, the total number of rooting steps is at most ss.