Block-Diagonal and LT Codes for Distributed Computing With Straggling Servers

Albin Severinson, Alexandre Graell i Amat, Eirik Rosnes

I Introduction

Distributed computing systems have emerged as one of the most effective ways of solving increasingly complex computational problems, such as those in large-scale machine learning and data analytics . These systems, referred to as “warehouse-scale computers” (WSCs) , may be composed of thousands of relatively homogeneous hardware and software components. Achieving high availability and efficiency for applications running on WSCs is a major challenge. One of the main reasons is the large number of components that may experience transient or permanent failures . As a result, several distributed computing frameworks have been proposed . In particular, MapReduce has gained significant attention as a means of effectively utilizing large computing clusters. For example, Google routinely performs computations over several thousands of servers using MapReduce . Among the challenges brought on by distributed computing systems, the problems of straggling servers and bandwidth scarcity have recently received significant attention. The straggler problem is a synchronization problem characterized by the fact that a distributed computing task must wait for the slowest server to complete its computation, which may cause large delays . On the other hand, distributed computing tasks typically require that data is moved between servers during the computation, the so-called data shuffling, which is a challenge in bandwidth-constrained networks.

Coding for distributed computing to reduce the computational delay and the communication load between servers has recently been considered in . In , a structure of repeated computation tasks across servers was proposed, enabling coded multicast opportunities that significantly reduce the required bandwidth to shuffle the results. In , the authors showed that maximum distance separable (MDS) codes can be applied to a linear computation task (e.g., multiplying a vector with a matrix) to alleviate the effects of straggling servers and reduce the computational delay. In , a unified coding framework was presented and a fundamental tradeoff between computational delay and communication load was identified. The ideas of can be seen as particular instances of the framework in , corresponding to the minimization of the communication load and the computational delay, respectively. The code proposed in is an MDS code of code length proportional to the number of rows of the matrix to be multiplied, which may be very large in practice. For example, Google performs matrix-vector multiplications with matrices of dimension of the order of 1010×101010^{10}\times 10^{10} when ranking the importance of websites . In , the computational delay incurred by the encoding and decoding is not considered. However, the encoding and decoding may incur a high computational delay for large matrices.

Coding has previously been applied to several related problems in distributed computing. For example, the scheme in has been extended to distributed matrix-matrix multiplication where both matrices are too large to be stored at one server . Whereas the schemes in are based on MDS codes, the scheme in is based on a novel coding scheme that exploits the algebraic properties of matrix-matrix multiplication over a finite field to reduce the computational delay. In , it was shown that introducing sparsity in a structured manner during encoding can speed up computing dot products between long vectors. Distributed computing over heterogeneous clusters has been considered in .

In this paper, we propose two coding schemes for the problem of multiplying a matrix by a set of vectors. The first is a block-diagonal coding (BDC) scheme equivalent to partitioning the matrix and applying smaller MDS codes to each submatrix separately (we originally introduced the BDC scheme in ). The storage design for the BDC scheme can be cast as an integer optimization problem, whose computation scales exponentially with the problem size. We propose a heuristic solver for efficiently solving the optimization problem, and a branch-and-bound approach for improving on the resulting solution iteratively. Furthermore, we prove that up to a certain level of partitioning the BDC scheme has identical computational delay (as defined in ) and communication load to those of the scheme in . Interestingly, when the delay incurred by encoding and decoding is taken into account, the proposed scheme achieves an overall computational delay significantly lower than that of the scheme in . We further propose a second coding scheme based on Luby Transform (LT) codes under inactivation decoding , which in some scenarios achieves a lower computational delay than that of the BDC scheme at the expense of a higher communication load. We show that for the LT code-based scheme it is possible to trade an increase in communication load for a lower computational delay. We finally consider distributed computing under a deadline, where we are interested in completing a computation within some computational delay, and show numerically that both the BDC and the LT code-based schemes significantly increase the probability of meeting a deadline over the scheme in . In particular, the LT code-based scheme achieves the highest probability of meeting a deadline for the scenarios considered.

II System Model and Preliminaries

where we assume that rr divides KmKm and hence qq is an integer. The rr coded rows of C\bm{C}, c1,…,cr\bm{c}_{1},\ldots,\bm{c}_{r}, are divided into (Kηq)\binom{K}{\eta q} disjoint batches, each containing r/(Kηq)r/\binom{K}{\eta q} coded rows. Each batch is assigned to ηq\eta q servers. Correspondingly, a batch BB is labeled by a unique set T⊂{S1,…,SK}\mathcal{T}\subset\{S_{1},\ldots,S_{K}\}, of size ∣T∣=ηq|\mathcal{T}|=\eta q, denoting the subset of servers that store that batch. We write BTB_{\mathcal{T}} to denote the batch stored at the unique set of servers T\mathcal{T}. Server SkS_{k}, k=1,…,Kk=1,\ldots,K, stores the coded rows of BTB_{\mathcal{T}} if and only if Sk∈TS_{k}\in\mathcal{T}.

We assume that running a computation on a single server takes a random amount of time, which is denoted by the random variable HH, according to the shifted-exponential cumulative probability distribution function (CDF)

When an algorithm is split into KK parallel subtasks that are run across KK servers, we denote the runtime of the subtask running on server SkS_{k} by HkH_{k}. As in , we assume that H1,…,HKH_{1},\ldots,H_{K} are independent and identically distributed random variables with CDF FH(Kh;σ)F_{H}(Kh;\sigma). For i=1,…,Ki=1,\dots,K, we denote the ii-th order statistic by H(i)H_{(i)}, i.e., the ii-th smallest random variable of H1,…,HKH_{1},\ldots,H_{K}. The runtime of the ii-th fastest server to complete its subtask is thus given by H(i)H_{(i)}, which is a Gamma distributed random variable with expectation and variance given by

We parameterize the Gamma distribution by its inverse scale factor aa and its shape parameter bb. We give these in terms of the distribution mean and variance as

Denote by FH(i)(h(i);σ,K)F_{H_{(i)}}(h_{(i)};\sigma,K) the CDF of H(i)H_{(i)}. It is given by

where Γ\Gamma denotes the Gamma function and γ\gamma the lower incomplete Gamma function,

We remark that FH(i)(h(i);σ,K)F_{H_{(i)}}(h_{(i)};\sigma,K) is the probability of a computation finishing prior to some deadline t=h(i)t=h_{(i)}.

II-B Distributed Computing Model

We consider the coded computing framework introduced in , which extends the MapReduce framework . The overall computation proceeds in three phases, the map, shuffle, and reduce phases, which are augmented to make use of the coded multicasting strategy proposed in to address the bandwidth scarcity problem and the coded scheme proposed in to alleviate the straggler problem. Furthermore, we consider the delay incurred by the encoding of A\bm{A} that takes place before the start of the map phase. We refer to this as the encoding phase. Also, we assume that the matrices A\bm{A} and Ψ\bm{\Psi} as well as the input vectors x1,…,xN\bm{x}_{1},\ldots,\bm{x}_{N} are known to all servers at the start of the computation. The overall computation proceeds in the following manner.

In the encoding phase, the coded matrix C\bm{C} is computed from A\bm{A} and Ψ\bm{\Psi} in a distributed fashion. Specifically, denote by R(S)\mathcal{R}^{(S)} the set of indices of rows of C\bm{C} that are assigned to server SS and denote by Ψ(S)\bm{\Psi}^{(S)} the matrix consisting of the rows of Ψ\bm{\Psi} with indices from R(S)\mathcal{R}^{(S)}. Then, server SS computes the coded rows it needs by multiplying Ψ(S)\bm{\Psi}^{(S)} by A\bm{A}. Note that since we assign each coded row to ηq\eta q servers, each row of C\bm{C} is computed separately by ηq\eta q servers. We define the computational delay of the encoding phase as its average runtime per source row and vector y\bm{y}, i.e.,

where σencode\sigma_{\mathsf{encode}} is the complexity of the encoding. During the encoding process, the rows of Ψ\bm{\Psi} are multiplied by the columns of A\bm{A}. Therefore, the complexity scales with the product of the number of nonzero elements of Ψ\bm{\Psi} and the number of columns of A\bm{A}. Specifically,

Alternatively, we compute C\bm{C} by performing a decoding operation on A\bm{A}. In this case σencode\sigma_{\mathsf{encode}} is the decoding complexity (see Section IV-B). Furthermore, since the decoding algorithms are designed to decode the entire codeword, each server has to compute all rows of C\bm{C}. Using this strategy the encoding delay is

For each case we choose the strategy that minimizes the delay.

II-B2 Map Phase

In the map phase, we compute in a distributed fashion coded intermediate values, which will be later used to obtain vectors y1,…,yN\bm{y}_{1},\ldots,\bm{y}_{N}. Server SS multiplies the input vectors xj\bm{x}_{j}, j=1,…,Nj=1,\ldots,N, by all the coded rows of matrix C\bm{C} it stores, i.e., it computes

The map phase terminates when a set of servers G⊆{S1,…,SK}\mathcal{G}\subseteq\{S_{1},\ldots,S_{K}\} that collectively store enough values to decode the output vectors have finished their map computations. We denote the cardinality of G\mathcal{G} by gg. The (r,m)(r,m) linear code proposed in is an MDS code for which y1,…,yN\bm{y}_{1},\ldots,\bm{y}_{N} can be obtained from any subset of qq servers, i.e., g=qg=q. We illustrate the completion of subtasks in Fig. 1.

We define the computational delay of the map phase as its average runtime per source row and vector y\bm{y}, i.e.,

where σmap=KηmN((n−1)σA+nσM)\sigma_{\mathsf{map}}=K\eta mN\left((n-1)\sigma_{\mathsf{A}}+n\sigma_{\mathsf{M}}\right), as all KK servers compute ηm\eta m inner products, each requiring n−1n-1 additions and nn multiplications, for each of the NN input vectors. In , DmapD_{\mathsf{map}} is referred to simply as the computational delay.

After the map phase, the computation of y1,…,yN\bm{y}_{1},\ldots,\bm{y}_{N} proceeds using only the servers in G\mathcal{G}. We denote by Q⊆G\mathcal{Q}\subseteq\mathcal{G} the set of the first qq servers to complete the map phase. Each of the qq servers in Q\mathcal{Q} is responsible to compute N/qN/q of the vectors y1,…,yN\bm{y}_{1},\ldots,\bm{y}_{N}. Let WS\mathcal{W}_{S} be the set containing the indices of the vectors y1,…,yN\bm{y}_{1},\ldots,\bm{y}_{N} that server S∈QS\in\mathcal{Q} is responsible for. The remaining servers in G\mathcal{G} assist the servers in Q\mathcal{Q} in the shuffle phase.

II-B3 Shuffle Phase

In the shuffle phase, intermediate values calculated in the map phase are exchanged between servers in G\mathcal{G} until all servers in Q\mathcal{Q} hold enough values to compute the vectors they are responsible for. As in , we allow creating and multicasting coded messages that are simultaneously useful for multiple servers. Furthermore, as in , we denote by ϕ(j)\phi(j) the ratio between the communication load of unicasting the same message to each of jj recipients and multicasting that message to jj recipients. For example, if the communication load of multicasting a message to jj recipients and unicasting a message to a single recipient is the same, we have ϕ(j)=j\phi(j)=j. On the other hand, if the communication load of multicasting a message to jj recipients is equal to that of unicasting that same message to each recipient, ϕ(j)=1\phi(j)=1. The total communication load of a multicast message is then given by jϕ(j)\frac{j}{\phi(j)}. The shuffle phase proceeds in three steps as follows.

Coded messages composed of several intermediate values are multicasted among the servers in Q\mathcal{Q}.

Intermediate values are unicasted among the servers in Q\mathcal{Q}.

Any intermediate values still missing from servers in Q\mathcal{Q} are unicasted from the remaining servers in G\mathcal{G}, i.e., from the servers in G∖Q\mathcal{G}\setminus\mathcal{Q}.

For a subset of servers S⊂Q\mathcal{S}\subset\mathcal{Q} and S∈Q∖SS\in\mathcal{Q}\setminus\mathcal{S}, we denote the set of intermediate values needed by server SS and known exclusively by the servers in S\mathcal{S} by VS(S)\mathcal{V}_{\mathcal{S}}^{(S)}. More formally,

We transmit coded multicasts only between the servers in Q\mathcal{Q}, and each coded message is simultaneously sent to multiple servers. We denote by

the smallest number of recipients of a coded message . We remark that mαjm\alpha_{j} is the total number of coded values delivered to each server via the coded multicast messages with exactly jj recipients. More specifically, for each j∈{ηq,ηq−1,…,sq}j\in\{\eta q,\eta q-1,\ldots,s_{q}\}, and every subset S⊆Q\mathcal{S}\subseteq\mathcal{Q} of size j+1j+1, the shuffle phase proceeds as follows.

The communication load, denoted by LL, is the number of unicasts and multicasts (weighted by their cost relative to a unicast) per source row and vector y\bm{y} exchanged during the shuffle phase. Specifically, each unicasted message increases LL by 1mN\frac{1}{mN}, and each message multicasted to jj recipients increases LL by jmNϕ(j)\frac{j}{mN\phi(j)}.

The communication load after completing the shuffle phase is given in . If the shuffle phase finishes by unicasting the remaining needed values (strategy 11), the communication load after completing the multicast phase is

If instead steps 1)1) and 2)2) are repeated for j=sq−1j=s_{q}-1 (strategy 22), the communication load is

For the scheme in , the total communication load is

where 1−η−∑j=sqηqαj1-\eta-\sum_{j=s_{q}}^{\eta q}\alpha_{j} is the communication load due to unicasting the remaining needed values.

II-B4 Reduce Phase

Finally, in the reduce phase, the vectors y1,…,yN\bm{y}_{1},\ldots,\bm{y}_{N} are computed. More specifically, server S∈QS\in\mathcal{Q} uses the locally computed sets Z1(S),…,ZN(S)\mathcal{Z}_{1}^{(S)},\ldots,\mathcal{Z}_{N}^{(S)} and the received messages to compute the vectors yj\bm{y}_{j}, ∀j∈WS\forall j\in\mathcal{W}_{S}. The computational delay of the reduce phase is its average runtime per source row and output vector y\bm{y}, i.e.,

where σreduce\sigma_{\mathsf{reduce}} is the computational complexity (see Section II-A) of the reduce phase.

The overall computational delay, DD, is the sum of the encoding, map, and reduce phase delays, i.e., D=Dencode+Dmap+DreduceD=D_{\mathsf{encode}}+D_{\mathsf{map}}+D_{\mathsf{reduce}}.

II-C Previously Proposed Coded Computing Schemes

The uncoded scheme uses no erasure coding and no coded multicasting and has parameters ηUC=1K\eta_{\mathsf{UC}}=\frac{1}{K} and qUC=Kq_{\mathsf{UC}}=K, implying ηUCqUC=1\eta_{\mathsf{UC}}q_{\mathsf{UC}}=1. Furthermore, the encoding matrix ΨUC\bm{\Psi}_{\mathsf{UC}} is the m×mm\times m identity matrix and the coded matrix is CUC=A\bm{C}_{\mathsf{UC}}=\bm{A}.

The CMR scheme uses only coded multicasting, i.e., CCMR=A\bm{C}_{\mathsf{CMR}}=\bm{A} and qCMR=Kq_{\mathsf{CMR}}=K. Furthermore, the fraction of rows stored at each server is ηCMR=ηqK\eta_{\mathsf{CMR}}=\frac{\eta q}{K}. We remark that there is no reduce delay for this scheme, i.e., Dreduce=0D_{\mathsf{reduce}}=0.

The SC scheme uses an erasure code but no coded multicasting. For the corresponding SC scheme, the code rate is unchanged, i.e., qSC=qq_{\mathsf{SC}}=q, and the fraction of rows stored at each server is ηSC=1qSC\eta_{\mathsf{SC}}=\frac{1}{q_{\mathsf{SC}}}. The encoding matrix ΨSC\bm{\Psi}_{\mathsf{SC}} of the SC scheme is obtained by splitting the rows of A\bm{A} into qSCq_{\rm{SC}} equally tall submatrices A1,…,AqSC\bm{A}_{1},\ldots,\bm{A}_{q_{\mathsf{SC}}} and applying a (K,qSC)(K,q_{\mathsf{SC}}) MDS code to the elements of each submatrix, thereby creating KK coded submatrices C1,…,CK\bm{C}_{1},\ldots,\bm{C}_{K}. The coded matrix CSC\bm{C}_{\mathsf{SC}} is the concatenation of C1,…,CK\bm{C}_{1},\ldots,\bm{C}_{K}, i.e.,

The unified scheme uses both an erasure code and coded multicasting and has parameters ηunified=η\eta_{\mathsf{unified}}=\eta and qunified=qq_{\mathsf{unified}}=q. Furthermore, the encoding matrix of the unified scheme, Ψunified\bm{\Psi}_{\mathsf{unified}}, is an (r,m)(r,m) MDS code encoding matrix.

III Block-Diagonal Coding

In this section, we introduce a BDC scheme for the problem of multiplying a matrix by a set of vectors. For large matrices, the encoding and decoding complexity of the proposed scheme is significantly lower than that of the scheme in , leading to a lower overall computational delay, as will be shown in Section VII. Specifically, the scheme is based on a block-diagonal encoding matrix of the form

where ψ1,…,ψT\bm{\psi}_{1},\ldots,\bm{\psi}_{T} are rT×mT\frac{r}{T}\times\frac{m}{T} encoding matrices of an (rT,mT)(\frac{r}{T},\frac{m}{T}) MDS code, for some integer TT that divides mm and rr. Note that the encoding given by ΨBDC\bm{\Psi}_{\mathsf{BDC}} amounts to partitioning the rows of A\bm{A} into TT disjoint submatrices A1,…,AT\bm{A}_{1},\ldots,\bm{A}_{T} and encoding each submatrix separately. We refer to an encoding ΨBDC\bm{\Psi}_{\mathsf{BDC}} with TT disjoint submatrices as a TT-partitioned scheme, and to the submatrix of C=ΨBDCA\bm{C}=\bm{\Psi}_{\mathsf{BDC}}\bm{A} corresponding to ψi\bm{\psi}_{i} as the ii-th partition. We remark that all submatrices can be encoded using the same encoding matrix, i.e., ψi=ψ\bm{\psi}_{i}=\bm{\psi}, i=1,…,Ti=1,\ldots,T, reducing the storage requirements, and encoding/decoding can be performed in parallel if many servers are available. Notably, by keeping the ratio mT\frac{m}{T} constant, the decoding complexity scales linearly with mm. We further remark that the case ΨBDC=ψ\bm{\Psi}_{\mathsf{BDC}}=\bm{\psi} (i.e., the number of partitions is T=1T=1) corresponds to the scheme in , which we will sometimes refer to as the unpartitioned scheme. We illustrate the BDC scheme with T=3T=3 partitions in Fig. 2.

For a block-diagonal encoding matrix ΨBDC\bm{\Psi}_{\mathsf{BDC}}, we denote by ci(t)\bm{c}_{i}^{(t)}, t=1,…,Tt=1,\ldots,T and i=1,…,r/Ti=1,\ldots,r/T, the ii-th coded row of C\bm{C} within partition tt. For example, c1(2)\bm{c}_{1}^{(2)} denotes the first coded row of the second partition. As described in Section II, the coded rows are divided into (Kηq)\binom{K}{\eta q} disjoint batches. To formally describe the assignment of coded rows to batches we use a (Kηq)×T\binom{K}{\eta q}\times T integer matrix P=[pi,j]\bm{P}=[p_{i,j}], where pi,jp_{i,j} is the number of rows from partition jj that are stored in batch ii. In the sequel, P\bm{P} will be referred to as the assignment matrix. Note that, due to the MDS property, any set of m/Tm/T rows of a partition is sufficient to decode the partition. Thus, without loss of generality, we consider a sequential assignment of rows of a partition into batches. More precisely, when first assigning a row of partition tt to a batch, we pick c1(t)\bm{c}_{1}^{(t)}. Next time a row of partition tt is assigned to a batch we pick c2(t)\bm{c}_{2}^{(t)}, and so on. In this manner, each coded row is assigned to a unique batch exactly once. The rows of P\bm{P} are labeled by the subset of servers the corresponding batch is stored at, and the columns are labeled by their partition indices. For convenience, we refer to the pair (ΨBDC,P)(\bm{\Psi}_{\mathsf{BDC}},\bm{P}) as the storage design. The assignment matrix P\bm{P} must satisfy the following conditions.

The entries of each row of P\bm{P} must sum up to the batch size, i.e.,

The entries of each column of P\bm{P} must sum up to the number of rows per partition, i.e.,

We clarify the assignment of coded rows to batches and the coded computing scheme in the following example.

For these parameters, there are r/T=6r/T=6 coded rows per partition, of which m/T=4m/T=4 are sufficient for decoding, and (Kηq)=15\binom{K}{\eta q}=15 batches, each containing r/(Kηq)=2r/\binom{K}{\eta q}=2 coded rows. We construct the storage design shown in Fig. 3 with (Kηq)×T=15×5\binom{K}{\eta q}\times T=15\times 5 assignment matrix

where rows are labeled by the subset of servers the batch is stored at, and columns are labeled by the partition index. In this case rows c1(1)\bm{c}_{1}^{(1)} and c2(1)\bm{c}_{2}^{(1)} are assigned to batch 11, c3(1)\bm{c}_{3}^{(1)} and c4(1)\bm{c}_{4}^{(1)} are assigned to batch 22, and so on. For this storage design, any g=4g=4 servers collectively store at least 44 coded rows from each partition. However, some servers store more rows than needed to decode some partitions, suggesting that this storage design is suboptimal.

Assume that G={S1,S2,S3,S4}\mathcal{G}=\{S_{1},S_{2},S_{3},S_{4}\} is the set of g=4g=4 servers that finish their map computations first. Also, assign vector yi\bm{y}_{i} to server SiS_{i}, i=1,2,3,4i=1,2,3,4. We illustrate the coded shuffling scheme for S={S1,S2,S3}\mathcal{S}=\{S_{1},S_{2},S_{3}\} in Fig. 4. Server S1S_{1} multicasts c1(1)x3⊕c3(1)x2\bm{c}_{1}^{(1)}\bm{x}_{3}\mathop{\oplus}\bm{c}_{3}^{(1)}\bm{x}_{2} to S2S_{2} and S3S_{3}. Since S2S_{2} and S3S_{3} can cancel c1(1)x3\bm{c}_{1}^{(1)}\bm{x}_{3} and c3(1)x2\bm{c}_{3}^{(1)}\bm{x}_{2}, respectively, both servers receive one needed intermediate value. Similarly, S2S_{2} multicasts c2(1)x3⊕c5(2)x1\bm{c}_{2}^{(1)}\bm{x}_{3}\mathop{\oplus}\bm{c}_{5}^{(2)}\bm{x}_{1}, while S3S_{3} multicasts c4(1)x2⊕c6(2)x1\bm{c}_{4}^{(1)}\bm{x}_{2}\mathop{\oplus}\bm{c}_{6}^{(2)}\bm{x}_{1}. This process is repeated for S={S2,S3,S4}\mathcal{S}=\{S_{2},S_{3},S_{4}\}, S={S1,S3,S4}\mathcal{S}=\{S_{1},S_{3},S_{4}\}, and S={S1,S2,S4}\mathcal{S}=\{S_{1},S_{2},S_{4}\}. After the shuffle phase, we have sent 1212 multicast messages and 3030 unicast messages, resulting in a communication load of (12+30)/20/4=0.525(12+30)/20/4=0.525, a 50%50\% increase from the load of the unpartitioned scheme (0.350.35, given by (2)). In this case, S1S_{1} received additional intermediate values from partition 22, despite already storing enough, further indicating that the assignment in (3) is suboptimal.

IV Performance of the Block-Diagonal Coding

In this section, we analyze the impact of partitioning on the performance. We also prove that we can partition up to the batch size, i.e., T=r/(Kηq)T=r/\binom{K}{\eta q}, without increasing the communication load and the computational delay of the map phase with respect to the original scheme in .

For the unpartitioned scheme of , G=Q\mathcal{G}=\mathcal{Q}, and the number of remaining values that need to be unicasted after the multicast phase is constant regardless which subset Q\mathcal{Q} of servers finish first their map computations. However, for the BDC (partitioned) scheme, both gg and the number of remaining unicasts may vary.

For a given assignment matrix P\bm{P} and a specific Q\mathcal{Q}, we denote by UQ(S)(P)U_{\mathcal{Q}}^{(S)}(\bm{P}) the number of remaining values needed after the multicast phase by server S∈QS\in\mathcal{Q}, and by

We first explain how UQ(S)U_{\mathcal{Q}}^{(S)} is evaluated. Let uQ(S)\bm{u}_{\mathcal{Q}}^{(S)} be a vector of length TT, where the tt-th element is the number of intermediate values from partition tt stored by server SS at the end of the multicast phase. Furthermore, each row of P\bm{P} corresponds to a batch, and coded multicasting is made possible by storing each batch at multiple servers. The intermediate values transmitted during the multicast phase thus correspond to rows of P\bm{P}. The vector uQ(S)\bm{u}_{\mathcal{Q}}^{(S)} is then computed by adding some set of rows of P\bm{P}. The indices of the rows to add depend on Q\mathcal{Q} and SS (see Section II-B3).

We denote by (uQ(S))t\left(\bm{u}_{\mathcal{Q}}^{(S)}\right)_{t} the tt-th element of the vector uQ(S)\bm{u}_{\mathcal{Q}}^{(S)}. The number of values UQ(S)U_{\mathcal{Q}}^{(S)} is given by adding the number of intermediate values still needed for each partition, i.e.,

We consider the same system as in Example 1. We again assume that G=Q={S1,S2,S3,S4}\mathcal{G}=\mathcal{Q}=\{S_{1},S_{2},S_{3},S_{4}\} is the set of g=q=4g=q=4 servers that finish their map computations first. During the multicast phase server S1S_{1} receives the intermediate values in VS∖S1(S1)\mathcal{V}_{\mathcal{S}\setminus S_{1}}^{(S_{1})} for all sets S\mathcal{S} of cardinality j+1=3j+1=3 (see Section II-B3). In this case, we perform coded multicasting within the sets

S={S1,S2,S3}\mathcal{S}=\{S_{1},S_{2},S_{3}\}, VS∖S1(S1)={c5(2)x1,c6(2)x1}\mathcal{V}_{\mathcal{S}\setminus S_{1}}^{(S_{1})}=\{\bm{c}_{5}^{(2)}\bm{x}_{1},\bm{c}_{6}^{(2)}\bm{x}_{1}\},

S={S1,S2,S4}\mathcal{S}=\{S_{1},S_{2},S_{4}\}, VS∖S1(S1)={c1(3)x1,c2(3)x1}\mathcal{V}_{\mathcal{S}\setminus S_{1}}^{(S_{1})}=\{\bm{c}_{1}^{(3)}\bm{x}_{1},\bm{c}_{2}^{(3)}\bm{x}_{1}\},

S={S1,S3,S4}\mathcal{S}=\{S_{1},S_{3},S_{4}\}, VS∖S1(S1)={c1(4)x1,c2(4)x1}\mathcal{V}_{\mathcal{S}\setminus S_{1}}^{(S_{1})}=\{\bm{c}_{1}^{(4)}\bm{x}_{1},\bm{c}_{2}^{(4)}\bm{x}_{1}\}.

Note that V{S2,S3}(S1)\mathcal{V}_{\{S_{2},S_{3}\}}^{(S_{1})} contains the intermediate values computed from the coded rows stored in the batch that labels the 66-th row of the assignment matrix P\bm{P}. In the same manner, V{S2,S4}(S1)\mathcal{V}_{\{S_{2},S_{4}\}}^{(S_{1})} and V{S3,S4}(S1)\mathcal{V}_{\{S_{3},S_{4}\}}^{(S_{1})} correspond to rows 77 and 1010 of P\bm{P}, respectively. Furthermore, prior to the shuffle phase server S1S_{1} stores the batches corresponding to rows 11 to 55 of P\bm{P}. Thus, u{S1,S2,S3,S4}(S1)\bm{u}_{\{S_{1},S_{2},S_{3},S_{4}\}}^{(S_{1})} is equal to the sum of rows 1,2,3,4,5,6,71,2,3,4,5,6,7, and 1010 of P\bm{P}. In this case, u{S1,S2,S3,S4}(S1)=(6,6,2,2,0)\bm{u}_{\{S_{1},S_{2},S_{3},S_{4}\}}^{(S_{1})}=(6,6,2,2,0), and S1S_{1} needs 88 more intermediate values, i.e., U{S1,S2,S3,S4}(S1)=8U_{\{S_{1},S_{2},S_{3},S_{4}\}}^{(S_{1})}=8. Computing uQ(S)\bm{u}_{\mathcal{Q}}^{(S)} for arbitrary Q\mathcal{Q} and SS then corresponds to summing the rows of P\bm{P} corresponding to batches either stored by server SS prior to the shuffle phase or received by SS in the multicast phase. The row indices are computed as explained in Section II-B3.

For a given ΨBDC\bm{\Psi}_{\mathsf{BDC}}, the assignment of rows into batches can be formulated as an optimization problem, where one would like to minimize LBDC(ΨBDC,P)L_{\text{BDC}}(\bm{\Psi}_{\mathsf{BDC}},\bm{P}) over all assignments P\bm{P}. More precisely, the optimization problem is

IV-B Computational Delay

We consider the delay incurred by the encoding, map, and reduce phases (see Definition 2). As in , we do not consider the delay incurred by the shuffle phase as the computations it requires are simple in comparison. Note that in only DmapD_{\mathsf{map}} is considered, i.e., D=DmapD=D_{\mathsf{map}}. However, one should not neglect the computational delay incurred by the encoding and reduce phases. Thus, we consider the overall computational delay

The encoding delay DencodeD_{\mathsf{encode}} is a function of the number of nonzero elements of ΨBDC\bm{\Psi}_{\mathsf{BDC}}. As there are at most mT\frac{m}{T} nonzero elements in each row of a block-diagonal encoding matrix, for an encoding scheme with TT partitions we have

The reduce phase consists of decoding the NN output vectors and hence the delay it incurs depends on the underlying code and decoding algorithm. We assume that each partition is encoded using a Reed-Solomon (RS) code and is decoded using either the Berlekamp-Massey (BM) algorithm or the FFT-based algorithm proposed in , whichever yields the lowest complexity. To the best of our knowledge the algorithm proposed in is the lowest complexity algorithm for decoding long RS codes. We measure the decoding complexity by its associated shifted-exponential parameter σ\sigma (see Section II-A).

The number of field additions and multiplications required to decode an (r/T,m/T(r/T,m/T) RS code using the BM algorithm is (r/T)(ξ(r/T)−1)(r/T)\left(\xi(r/T)-1\right) and (r/T)2ξ(r/T)^{2}\xi, respectively, where ξ\xi is the fraction of erased symbols . With ξ\xi upper bounded by 1−qK1-\frac{q}{K} (the map phase terminates when a fraction of at least qK\frac{q}{K} symbols from each partition is available), the complexity of decoding the TT partitions for all NN output vectors is upper bounded as

On the other hand, the FFT-based algorithm has complexity O(rlog⁡r)\mathcal{O}(r\log r) . We estimate the number of additions and multiplications required for a given code length rr by fitting a curve of the form a+brlog⁡2(cr)a+br\log_{2}(cr), where (a,b,c)(a,b,c) are coefficients, to empiric results derived from the authors’ implementation of the algorithm. For additions the resulting parameters are (2,8.5,0.867)(2,8.5,0.867) and for multiplications they are (2,1,4)(2,1,4). The resulting curves diverge negligibly at the measured points. The total decoding complexity for the FFT-based algorithm is

The encoding and decoding complexity of the unified scheme in is given by evaluating (8) and either (9) or (10) (whichever gives the lowest complexity), respectively, for T=1T=1. For the BDC scheme, by choosing TT close to rr we can thus significantly lower the delay of the encoding and reduce phases. On the other hand, the scheme in uses codes of length proportional to the number of servers KK. The encoding and decoding complexity of the SC scheme in is thus given by evaluating (8) and either (9) or (10) for T=mqT=\frac{m}{q}.

IV-C Lossless Partitioning

For T≤r/(Kηq)T\leq r/\binom{K}{\eta q}, there exists an assignment matrix P\bm{P} such that the communication load and the computational delay of the map phase are equal to those of the unpartitioned scheme.

The computational delay of the map phase is equal to that of the unpartitioned scheme if any qq servers hold enough coded rows to decode all partitions. For T=r/(Kηq)T=r/\binom{K}{\eta q} we let P\bm{P} be a (Kηq)×T\binom{K}{\eta q}\times T all-ones matrix and show that it has this property by repeating the argument from [9, Sec. IV.B] for each partition. In this case, any set of qq servers collectively store ηqmT\frac{\eta qm}{T} rows from each partition, and since each coded row is stored by at most ηq\eta q servers, any qq servers collectively store at least ηqmηqT=mT\frac{\eta qm}{\eta qT}=\frac{m}{T} unique coded rows from each partition. The computational delay of the map phase is thus unchanged from the unpartitioned scheme. The communication load is unchanged if UQ(S)U_{\mathcal{Q}}^{(S)} is equal to that of the unpartitioned scheme for all Q\mathcal{Q} and SS. The number of values needed UQ(S)U_{\mathcal{Q}}^{(S)} is computed from uQ(S)\bm{u}_{\mathcal{Q}}^{(S)} (see (7)), which is the sum of ll rows of P\bm{P}, for some integer ll. For the all-ones assignment matrix, because all rows of P\bm{P} are identical, we have

which is the number of remaining values for the unpartitioned scheme.

Next, we consider the case where T<r/(Kηq)T<r/\binom{K}{\eta q}. First, consider the case T=r/(Kηq)−jT=r/\binom{K}{\eta q}-j, for some integer jj, 0≤j<r2(Kηq)0\leq j<\frac{r}{2\binom{K}{\eta q}}. We first set all entries of P\bm{P} equal to 11. At this point, the total number of unique rows of C\bm{C} per partition stored by any set of qq servers is at least

The number of coded rows per partition that are not yet assigned is given by r/Tr/T multiplied by the fraction of partitions removed jr/(Kηq)\frac{j}{r/\binom{K}{\eta q}}, i.e.,

We assign these rows to batches such that an equal number of coded rows is assigned to each of the KK servers, which is always possible due to the limitations imposed by the system model. Any set of qq servers will thus store a fraction q/Kq/K of these rows. The total number of unique coded rows per partition stored among any set of qq servers is then lower bounded by the sum of (12) weighted by q/Kq/K and (11), i.e.,

showing that it is possible to decode all partitions using the coded rows stored over any set of qq servers.

The communication load is unchanged with respect to the case where the number of partitions is r/(Kηq)r/\binom{K}{\eta q} if and only if no server receives rows it does not need in the multicast phase. Due to decreasing the number of partitions from r/(Kηq)r/\binom{K}{\eta q} to T=r/(Kηq)−jT=r/\binom{K}{\eta q}-j, we increase the number of coded rows needed to decode each partition by

Furthermore, reducing the number of partitions increases the number of coded rows per partition stored among any set of qq servers (see (12) and the following text) by

Note that the number of additional rows needed to decode each partition (see (13)) is greater than or equal to the number of additional rows stored among the qq servers (see (14)). It is thus impossible that too many coded rows are delivered for any partition.

Second, we consider the case T=r/(Kηq)−jiT=\frac{r/\binom{K}{\eta q}-j}{i}, where jj is chosen as for the first case above and where ii is a positive integer. Now, we first set all elements of P\bm{P} to ii. At this point the number of unique rows of C\bm{C} per partition stored by any set of qq servers is given by (11) multiplied by a factor ii (since we set each element of P\bm{P} to ii instead of one). Furthermore, the number of coded rows per partition that are not yet assigned is given by (12). Therefore, by using the same strategy as for i=1i=1 and assigning the remaining rows to batches such that an equal number of rows is assigned to each of the KK servers, we are guaranteed that the communication load and the computational delay are unchanged also in this case. ∎

V Assignment Solvers

For all solvers, we first label the batches lexiographically and then optimize LBDCL_{\text{BDC}} in (6). For example, for ηq=2\eta q=2, we label the first batch by S1,S2S_{1},S_{2}, the second by S1,S3S_{1},S_{3}, and so on. The solvers are available under the Apache 2.0 license . We remark that choosing P\bm{P} is similar to the problem of designing the coded matrices stored by each server in .

The heuristic solver is inspired by the assignment matrices created by the branch-and-bound solver for small instances. It creates an assignment matrix P\bm{P} in two steps. We first set each entry of P\bm{P} to

thus assigning the first (Kηq)Y\binom{K}{\eta q}Y rows of each partition to batches such that each batch is assigned YTYT rows. Let d=r/(Kηq)−YTd=r/\binom{K}{\eta q}-YT be the number of rows that still need to be assigned to each batch. The r/T−(Kηq)Yr/T-\binom{K}{\eta q}Y rows per partition not assigned yet are assigned in the second step as shown in Algorithm 1.

Interestingly, for T≤r/(Kηq)T\leq r/\binom{K}{\eta q} the heuristic solver creates an assignment matrix satisfying the requirements outlined in the proof of Theorem 1. In the special case of T=r/(Kηq)T=r/\binom{K}{\eta q}, the all-ones matrix is produced.

V-B Branch-and-Bound Solver

The branch-and-bound solver finds an optimal solution by recursively branching at each batch for which there is more than one possible assignment and considering all options. The solver is initially given an empty assignment matrix, i.e., an all-zeros (Kηq)×T\binom{K}{\eta q}\times T matrix. For each branch, we lower bound the value of the objective function of any assignment in that branch and only investigate branches with possibly better assignments. The branch-and-bound operations given below are repeated until there are no more potentially better solutions to consider.

For the first row of P\bm{P} with remaining assignments, branch on every available assignment for that row. More precisely, find the smallest index ii of a row of the assignment matrix P\bm{P} whose entries do not sum up to the batch size, i.e.,

For row ii, branch on incrementing the element pi,jp_{i,j} by 11 for all columns (with index jj) such that their entries do not sum up to the number of coded rows per partition, i.e.,

V-B2 Bound

We use a dynamic programming approach to lower bound LBDCL_{\text{BDC}} for a subtree. Specifically, for each row ii and column jj of P\bm{P}, we store the number of vectors uQ(S)\bm{u}_{\mathcal{Q}}^{(S)} that are indexed by row ii and where the jj-th element satisfies

V-C Hybrid Solver

The branch-and-bound solver can only be used by itself for small instances. However, it can be used to complete a partial assignment matrix, i.e., a matrix P\bm{P} for which not all rows have entries that sum up to the batch size. The branch-and-bound solver then completes the assignment optimally. We first find a candidate solution using the heuristic solver and then iteratively improve it using the branch-and-bound solver. In particular, we decrement by 11 a random set of entries of P\bm{P} and then use the branch-and-bound solver to reassign the corresponding rows optimally. We repeat this process until the average improvement between iterations drops below some threshold.

VI Luby Transform Codes

In this section, we consider LT codes for use in distributed computing. Specifically, we consider a distributed computing system where Ψ\bm{\Psi} is an LT code encoding matrix, denoted by ΨLT\bm{\Psi}_{\mathsf{LT}}, of fixed rate mr\frac{m}{r}. As explained in Section II, we divide the rr coded rows of C=ΨLTA\bm{C}=\bm{\Psi}_{\mathsf{LT}}\bm{A} into (Kηq)\binom{K}{\eta q} disjoint batches, each of which is stored at a unique subset of size ηq\eta q of the KK servers. For this scheme, due to the random nature of LT codes, we can assign coded rows to batches randomly. The distributed computation is carried out as explained in Section II-B, i.e., we wait for the fastest g≥qg\geq q servers to complete their respective computations in the map phase, perform coded multicasting during the shuffle phase, and carry out the decoding of the NN output vectors in the reduce phase.

We assume that decoding is performed using inactivation decoding . Inactivation decoding is an efficient maximum likelihood decoding algorithm that combines iterative decoding with optimal decoding in a two-step fashion and is widely used in practice. As suggested in , we assume that the optimal decoding phase is performed by Gaussian elimination. In particular, iterative decoding is used until the ripple is empty, i.e., until there are no coded symbols of degree 11, at which point an input symbol is inactivated. The iterative decoder is then restarted to produce a solution in terms of the inactivated symbol. This procedure is repeated until all input symbols are either decoded or inactivated. Note that the value of some input symbols may be expressed in terms of the values of the inactivated symbols at this point. Finally, optimal decoding of the inactivated symbols is performed via Gaussian elimination, and the decoded values are back-substituted into the decoded input symbols that depend on them. The decoding schedule has a large performance impact. Our implementation follows the recommendations in . It is important to tune the parameters MM and δ\delta to minimize the number of inactivations.

Due to the nature of LT codes, we need to collect m(1+ϵ)m(1+\epsilon) intermediate values for each vector y\bm{y} before decoding. We refer to ϵ\epsilon as the overhead. Under inactivation decoding, and for a given overhead ϵ\epsilon, the probability of decoding failure with an overhead of at most ϵ\epsilon, denoted by Pf(ϵ)P_{\mathsf{f}}(\epsilon), is lower bounded by

Note that Pf(ϵ)P_{\mathsf{f}}(\epsilon) is the CDF for the random variable “decoding is not possible at a given overhead ϵ\epsilon.” Furthermore, the lower bound (15) well approximates the failure probability for an overhead slightly larger than ϵ=0\epsilon=0. Denote by FDS(ϵ)F_{\mathsf{DS}}(\epsilon) the probability of decoding being possible at an overhead of at most ϵ\epsilon. It follows that

We find the decoding success probability density function (PDF) by numerically differentiating FDS(ϵ)F_{\mathsf{DS}}(\epsilon).

VI-B Code Design

We design LT codes for a minimum overhead ϵmin\epsilon_{\mathsf{min}}, i.e., we collect at least m(1+ϵmin)m(1+\epsilon_{\mathsf{min}}) coded symbols from the servers before attempting to decode, and a target failure probability Pf,target=Pf(ϵmin)P_{\mathsf{f,target}}=P_{\mathsf{f}}(\epsilon_{\mathsf{min}}). We remark that increasing ϵmin\epsilon_{\mathsf{min}} and Pf,targetP_{\mathsf{f,target}} leads to a lower average degree Ωˉ\bar{\Omega}, and thus to less complex encoding and decoding and subsequently to a lower computational delay for encoding and decoding. The tradeoff is that the communication load increases as more intermediate values need to be transferred over the network on average. Furthermore, increasing ϵmin\epsilon_{\mathsf{min}} and Pf,targetP_{\mathsf{f,target}} may increase the average number of servers gg required to decode. We thus need to balance the computational delay of the encoding and reduce phases against that of the map phase to achieve a low overall computational delay. Furthermore, waiting for more than g=qg=q servers typically increases the overall computational delay by more than what is saved by the less complex encoding and decoding given by the larger ϵmin\epsilon_{\mathsf{min}} and Pf,targetP_{\mathsf{f,target}}. We thus choose ϵmin\epsilon_{\mathsf{min}} and Pf,targetP_{\mathsf{f,target}} such that decoding is possible with high probability using the number of coded rows stored at any set of qq servers. Note that the overhead ϵ\epsilon required for decoding may be larger than ϵmin\epsilon_{\mathsf{min}}. We take this into account by numerically integrating the decoding success PDF multiplied by the performance of the scheme as a function of the overhead ϵ\epsilon.

For a given ϵmin\epsilon_{\mathsf{min}} and Pf,targetP_{\mathsf{f,target}}, we find a pair (M,δ)(M,\delta) that minimizes the decoding complexity (see Section VI-C) under the constraint that Pf(ϵmin)≈Pf,targetP_{\mathsf{f}}(\epsilon_{\mathsf{min}})\approx P_{\mathsf{f,target}}. Essentially, we minimize the computational delay of the reduce phase for a fixed delay of the map phase. We remark that LT codes with low decoding complexity have a low average degree Ωˉ\bar{\Omega}, and thus also low encoding complexity. Note that for a given MM, decreasing δ\delta lowers the failure probability, but also increases the decoding complexity. We find good pairs (M,δ)(M,\delta) by selecting through binary search the largest MM such that there exists a δ\delta for which the lower bound on Pf(ϵmin)P_{\mathsf{f}}(\epsilon_{\mathsf{min}}) in (15) is approximately equal to Pf,targetP_{\mathsf{f,target}}. This heuristic produces codes with complexity very close to those found using basin-hopping combined with the Powell optimization method .

VI-C Computational Delay

There are on average Ωˉ\bar{\Omega} nonzero entries in each row of the LT code encoding matrix. The LT code encoding complexity is thus given by

We simulate the complexity of the decoding σreduce,LT\sigma_{\mathsf{reduce,LT}}. Furthermore, we assume that the decoding complexity depends only on ϵmin\epsilon_{\mathsf{min}}, i.e., we evaluate the decoding complexity only at ϵ=ϵmin\epsilon=\epsilon_{\mathsf{min}}, and simulate the number of servers gg required for a given overhead ϵ\epsilon.

VI-D Communication Load

The coded multicasting scheme (see Section II-B3) is designed for the case where we need mm intermediate values per vector y\bm{y}. Here, we tune it for the case where we instead need at least m(1+ϵmin)m(1+\epsilon_{\mathsf{min}}) intermediate values by increasing the number of coded multicast messages sent. Note that the coded multicasting scheme is greedy in the sense that it starts by multicasting coded messages to the largest possible number of recipients and then gradually lowers the number of recipients. Specifically, we perform the shuffle phase with (see (1))

The communication load of the LT code-based scheme for a given ϵ≥ϵmin\epsilon\geq\epsilon_{\mathsf{min}} is then given by

VI-E Partitioning of the LT Code-Based Scheme

We can apply partitioning to the LT code-based scheme in the same manner as for the BDC scheme. Specifically, we consider a block-diagonal encoding matrix ΨBDC−LT\bm{\Psi}_{\mathsf{BDC-LT}}, where the blocks ψ1,…,ψT\bm{\psi}_{1},\ldots,\bm{\psi}_{T} are LT code encoding matrices. In particular, we consider the case where the number of partitions TT is equal to the partitioning limit of Theorem 1, i.e., T=r/(Kηq)T=r/\binom{K}{\eta q}. In this case the all-ones assignment matrix P\bm{P} introduced in the proof of Theorem 1 is a valid matrix. By using this assignment matrix and identical encoding matrices for each of the partitions, i.e., ψi=ψ\bm{\psi}_{i}=\bm{\psi}, i=1,…,Ti=1,\ldots,T, the encoding and decoding complexity of each partition is identical regardless of which set of servers G\mathcal{G} first completes the map phase. Furthermore, by the same argument as in the proof of Theorem 1, we are guaranteed that if any partition can be decoded using the coded rows stored at the set of servers G\mathcal{G}, all other partitions can also be decoded.

VII Numerical Results

We present numerical results for the proposed BDC and LT code-based schemes and compare them with the schemes in . Furthermore, we compare the performance of the BDC scheme with assignment P\bm{P} produced by the heuristic and hybrid solvers. We also evaluate the performance of the LT code-based scheme for different Pf,targetP_{\mathsf{f,target}} and ϵmin\epsilon_{\mathsf{min}}. For each plot, the field size is equal to one more than the largest number of coded rows considered for that plot, r+1r+1, rounded up to the closest power of 22. The results, except those in Fig. 12, are normalized by the performance of the uncoded scheme. Unless stated otherwise, the assignment P\bm{P} is given by the heuristic solver. As in , we assume that ϕ(j)=j\phi(j)=j.

In Fig. 5, we depict the communication load LL (see Definition 1) and the computational delay DD (see Definition 2) as a function of the number of partitions, TT. The system parameters are m=6000m=6000, n=6000n=6000, K=9K=9, q=6q=6, N=6000N=6000, and η=1/3\eta=1/3. The parameters of the CMR and SC schemes are qCMR=9q_{\mathsf{CMR}}=9, ηCMR=29\eta_{\mathsf{CMR}}=\frac{2}{9}, and ηSC=16\eta_{\mathsf{SC}}=\frac{1}{6}. The minimum overhead for the LT code-based scheme is ϵmin=0.3\epsilon_{\mathsf{min}}=0.3 and its target failure probability is Pf,target=0.1P_{\mathsf{f,target}}=0.1. For up to r/(Kηq)=250r/\binom{K}{\eta q}=250 partitions (marked by the vertical dotted line), the BDC scheme does not incur any loss in DmapD_{\mathsf{map}} and communication load with respect to the unified scheme (see Theorem 1). Furthermore, the BDC scheme yields about a 22% lower delay compared to the unified scheme for T=1000T=1000. The delay of the LT code-based scheme is slightly worse than that of the BDC scheme, and the load is about 65%65\% higher (for T=250T=250). Partitioning the LT code-based scheme increases the communication load and reduces the computational delay by about 0.50.5%. We remark that the number of partitions for the LT code-based scheme is fixed at r/(Kηq)r/\binom{K}{\eta q}. For heavy partitioning of the BDC scheme, a tradeoff between partitioning level, communication load, and map phase delay is observed. For example, with 30003000 partitions (the maximum possible), there is about a 10%10\% increase in communication load over the unified scheme. Note that the gain in computational delay saturates, thus there is no reason to partition beyond a given level. The load of the SC scheme is about twice that of our proposed schemes and the delay is about half. Finally, the delay of the BDC and the LT code-based scheme is about 25%25\% lower compared to the CMR scheme for T>100T>100.

In Fig. 6, we plot the performance for a constant ηq=2\eta q=2, n=m/100n=m/100, ηm=2000\eta m=2000, code rate m/r=2/3m/r=2/3, and N=500qN=500q vectors as a function of the number of servers, KK. The ratio m/nm/n is motivated by machine learning applications, where the number of rows and columns often represent the number of samples and features, respectively. Note that the number of arithmetic operations performed by each server in the map phase increases with KK. We choose the number of partitions TT that minimizes the delay under the constraint that the communication load is at most 11% higher compared to the unified scheme. The parameters of the LT code-based scheme are ϵmin=0.335\epsilon_{\mathsf{min}}=0.335 and Pf,target=0.1P_{\mathsf{f,target}}=0.1. The results shown are averages over 10001000 randomly generated realizations of G\mathcal{G}. Our proposed BDC scheme outperforms the unified scheme in terms of computational delay by between about 2525% (for K=6K=6) and 1010% (for K=201K=201). Furthermore, the delay of both the BDC and LT code-based schemes are about 5050% lower than that of the CMR scheme for K=201K=201. For K=6K=6 the computational delay of the unpartitioned and partitioned LT code-based schemes is about 55% higher and 88% lower compared to the BDC scheme, respectively. For K=201K=201 the delay of the LT code-based scheme is about 11% lower than that of the BDC scheme. However, the communication load is about 45%45\% higher. Finally, the communication load of the BDC scheme is between about 42%42\% (for K=6K=6) and 66%66\% (for K=201K=201) of that of the SC scheme.

In Fig. 7, we show the performance for code rate m/r=2/3m/r=2/3, ηq=2\eta q=2, and a fixed workload per server as a function of KK. Specifically, we fix the number of additions and multiplications computed by each server in the map phase to 10810^{8} (±5\pm 5% to find valid parameters) and scale m,n,Nm,n,N with KK. The number of rows mm of A\bm{A} takes values between 1260012600 and 5980059800, and we let n=m/100n=m/100 and N=nN=n. The number of partitions TT is selected in the same way as for Fig. 6. The results shown are averages over 10001000 randomly generated realizations of G\mathcal{G}. The computational delay of the unified scheme is about a factor 2020 higher than that of the BDC scheme for K=300K=300. The computational delay of the partitioned LT code-based scheme is similar to that of the BDC scheme, while the delay of the unpartitioned LT code-based scheme is about 6060% higher. Furthermore, the communication load of the LT code-based scheme is about 4545% higher compared to those of the unified and BDC schemes.

In Fig. 8, we plot the performance of the BDC and LT code-based schemes as a function of the number of columns nn. The system parameters are m=2400m=2400, K=9K=9, q=6q=6, N=60N=60, T=240T=240, and η=1/3\eta=1/3. The communication load of the LT code-based scheme depends primarily on the minimum overhead ϵmin\epsilon_{\mathsf{min}} and the computational delay primarily on the target failure probability Pf,targetP_{\mathsf{f,target}}. We remark that a higher Pf,targetP_{\mathsf{f,target}} allows for using codes with lower average degree and thus less complex encoding and decoding. For n=20000n=20000, the computational delay of the LT code-based scheme with Pf,target=0.1P_{\mathsf{f,target}}=0.1 is about 1.51.5% lower than that of the BDC scheme. For Pf,target=0.001P_{\mathsf{f,target}}=0.001, the computational delay is about 33% and 1.51.5% higher than that of the BDC scheme when ϵmin=0.3\epsilon_{\mathsf{min}}=0.3 and ϵmin=0.37\epsilon_{\mathsf{min}}=0.37, respectively. On the other hand, the communication load of the LT code-based scheme with ϵmin=0.3\epsilon_{\mathsf{min}}=0.3 and ϵmin=0.37\epsilon_{\mathsf{min}}=0.37 is about 41%41\% and 44%44\% higher than that of the BDC scheme, respectively.

VII-B Assignment Solver Comparison

In Figs. 9 and 10, we plot the performance of the BDC scheme with assignment P\bm{P} given by the heuristic and the hybrid solver. We also give the average performance over 100100 random assignments. The vertical dotted line marks the partitioning limit of Theorem 1. The parameters in Figs. 9 and 10 are identical to those in Figs. 5 and 6, respectively.

In Fig. 9, we plot the performance as a function of the number of partitions, TT. For TT less than about 200200, the performance for all solvers is identical. On the other hand, for T>200T>200 both the computational delay and the communication load are reduced with P\bm{P} from the heuristic solver over the random assignments (about 55% for load and 4747% for delay at T=3000T=3000). A further improvement in communication load can be achieved using the hybrid solver, but at the expense of a possibly larger computational delay.

In Fig. 10, we plot the performance as a function of the number of servers, KK. The results shown are averages over 10001000 randomly generated realizations of G\mathcal{G}. For K=6K=6, the communication load of the heuristic solver is about 5%5\% lower than that of the random assignments, but for K=201K=201 the difference is negligible. In terms of computational delay, the heuristic solver outperforms the random assignments by about 1818% and 33% for K=9K=9 and K=201K=201, respectively. The hybrid solver is too computationally complex for use with the largest systems considered.

VII-C Tradeoff Between Communication Load and Computational Delay

In Fig. 11, we show the tradeoff between communication load and computational delay. The parameters are K=14K=14, m=50000m=50000 (±3\pm 3% to find valid parameters), n=500n=500, N=840N=840, and η=1/2\eta=1/2. Note that the code rate is decreasing toward the bottom of the plot. We select the number of partitions TT that minimizes the delay while the load is at most 11% or 1010% higher compared to the unified scheme. Allowing a 1010% increased load gives up to about 77% lower delay compared to allowing a 11% increase. For the topmost data point of the BDC and unified schemes the encoding complexity dominates, and there is no reason to operate at this point since both the delay and load can be reduced. The parameters of the partitioned LT code-based scheme are ϵmin=0.3\epsilon_{\mathsf{min}}=0.3 and Pf,target=10−1P_{\mathsf{f,target}}=10^{-1}. For the data point with minimum computational delay, the LT code-based scheme yields about 1515% lower delay at the expense of about a 3030% higher load compared to the BDC scheme. Finally, the computational delay of the BDC scheme is between about 4747% and 44% lower compared to the unified scheme for the topmost and bottommost data points, respectively.

VII-D Computational Delay Deadlines

In Fig. 12, we plot the probability of a computation not finishing before a deadline tt, i.e., the probability of the computational delay being larger than tt. As in , we plot the complement of the CDF of the delay in logarithmic scale. On the horizontal axis, we show the deadline tt. The system parameters are K=201K=201, q=134q=134, m=134000m=134000, n=1340n=1340, N=67000N=67000 vectors, T=6700T=6700 partitions, and code rate m/r=2/3m/r=2/3. The parameters for the LT code-based scheme are ϵmin=0.335\epsilon_{\mathsf{min}}=0.335 and Pf,target=10−9P_{\mathsf{f,target}}=10^{-9}. The results are due to simulations. In particular, we simulate the decoding failure probability of LT codes for various tt and extrapolate from these points under the assumption that the decoding failure probability is Gamma distributed. The fitted values deviate negligibly from the simulated values.

When the deadline is t=3500t=3500, the probability of exceeding the deadline is about 0.40.4 for the unified and uncoded schemes. For the BDC scheme the probability is only about 7⋅10−37\cdot 10^{-3}. The probability is slightly lower for the LT code-based scheme, about 4⋅10−34\cdot 10^{-3}. If we instead consider a deadline t=4000t=4000, the probability of exceeding the deadline is about 10−310^{-3} and 0.150.15 for the unified and uncoded schemes, respectively. For the BDC scheme the probability of exceeding the deadline is about 9⋅10−89\cdot 10^{-8}, i.e., 44 orders of magnitude lower compared to the unified scheme. The LT code-based scheme further improves the performance with a probability of exceeding the deadline of about 3⋅10−83\cdot 10^{-8}. We remark that for the data point with minimum delay in Fig. 11, the LT code-based scheme has a significant advantage over the BDC scheme in terms of meeting a short deadline.

VII-E Alternative Runtime Distribution

Here, we consider a runtime distribution with CDF

where σ\sigma is the shift and β\beta is a parameter that scales the tail of the distribution, i.e., it differs from the one considered previously by that the scale of the tail may be different from the shift. It is equal to the previously considered distribution if β=σ\beta=\sigma. This model has been used to model distributed computing in, e.g., . Under this model we assume that the reduce delay of the uncoded scheme follows the distribution above with parameters β\beta and σUC,reduce=0\sigma_{\mathsf{UC,reduce}}=0 since each server has to assemble the final output from the intermediate results regardless coding is used or not. We assume that the encoding delay of the uncoded scheme is zero. Denote by σc\sigma_{\mathsf{c}} the computational complexity of matrix-vector multiplication for the BDC and unified schemes. We let β=ωσc\beta=\omega\sigma_{\mathsf{c}} for ω=0,1,10,100\omega=0,1,10,100. In Fig. 13, we plot the computational delay normalized by that of the uncoded scheme. The system parameters (and thus also the communication load) are identical to those in Fig. 7.

We observe the greatest gain of the BDC scheme over the unified scheme for small ω\omega since the benefits of straggler coding are small compared to the added delay due to encoding and decoding, which is significant for the unified scheme. For larger ω\omega the benefits of straggler coding are larger while the delay due to encoding and decoding remains constant. Hence, the performance of both schemes converge. However, even for ω=100\omega=100 the delay of the unified scheme is about 3333% higher than that of the BDC scheme for the largest system considered (K=300K=300). We remark that for the example considered in the parameters β=σ=1\beta=\sigma=1, i.e., ω=1\omega=1, are used.

VIII Conclusion

We introduced two coding schemes for distributed matrix multiplication. One is based on partitioning the matrix into submatrices and encoding each submatrix separately using MDS codes. The other is based on LT codes. Compared to the earlier scheme in and to the CMR scheme in , both proposed schemes yield a significantly lower overall computational delay. For instance, for a matrix of size 59800×59859800\times 598, the BDC scheme reduces the computational delay by about a factor 2020 over the scheme in with about a 11% increase in communication load. The LT code-based scheme may reduce the computational delay further at the expense of a higher communication load. For example, for a matrix with about 5000050000 rows, the computational delay of the LT code-based scheme is about 1515% lower than that of the BDC scheme with a communication load that is about 30%30\% higher. Finally, we have shown that the proposed coding schemes significantly increase the probability of a computation finishing within a deadline. The LT code-based scheme may be the best choice in situations where high reliability is needed due to its ability to decrease the computational delay at the expense of the communication load.

Acknowledgment

The authors would like to thank Dr. Francisco Lázaro and Dr. Gianluigi Liva for fruitful discussions and insightful comments on LT codes.

References