A Unified Coding Framework for Distributed Computing with Straggling Servers

Songze Li, Mohammad Ali Maddah-Ali, A. Salman Avestimehr

I Introduction

Recently, there have been two novel ideas proposed to exploit coding in order to speed up distributed computing applications. Specifically, a repetitive structure of computation tasks across distributed computing servers was proposed in , enabling coded multicast opportunities that significantly reduce the time to shuffle intermediate results. On the other hand, applying Maximum Distance Separable (MDS) codes to some linear computation tasks (e.g., matrix multiplication) was proposed in , in order to alleviate the effects of straggling servers and shorten the computation phase of distributed computing.

In this paper, we propose a unified coded framework for distributed computing with straggling servers, by introducing a tradeoff between “latency of computation” and “load of communication” for linear computation tasks. We show that the coding schemes of and can then be viewed as special instances of the proposed coding framework by considering two extremes of this tradeoff: minimizing either the load of communication or the latency of computation individually. Furthermore, the proposed coding framework provides a natural tradeoff between computation latency and communication load in distributed computing, and allows to systematically operate at any point on that tradeoff.

More specifically, we focus on a distributed matrix multiplication problem in which for a matrix A{\bf A} and NN input vectors x1,…,xN{\bf x}_{1},\ldots,{\bf x}_{N}, we want to compute NN output vectors y1=Ax1,…,yN=AxN{\bf y}_{1}={\bf A}{\bf x}_{1},\ldots,{\bf y}_{N}={\bf A}{\bf x}_{N}. The computation cannot be performed on a single server node since its local memory is too small to hold the entire matrix A{\bf A}. Instead, we carry out this computation using KK distributed computing servers collaboratively. Each server has a local memory, with the size enough to store up to equivalent of μ\mu fraction of the entries of the matrix A, and it can only perform computations based on the contents stored in its local memory. Matrix multiplication is one of the building blocks to solve data analytics and machine learning problems (e.g., regression and classification). Many such applications of big data analytics require massive computation and storage power over large-scale datasets, which are nowadays provided collaboratively by clusters of computing servers, using efficient distributed computing frameworks such as Hadoop MapReduce and Spark . Therefore, optimizing the performance of distributed matrix multiplication is of vital importance to improve the performance of the distributed computing applications.

A distributed implementation of matrix multiplication proceeds in three phases: Map, Shuffle and Reduce. In the Map phase, every server multiplies the input vectors with the locally stored matrix that partially represents the target matrix A{\bf A}. When a subset of servers finish their local computations such that their Map results are sufficient to recover the output vectors, we halt the Map computation and start to Shuffle the Map results across the servers in which the final output vectors are calculated by specific Reduce functions.

Within the above three-phase implementation, the coding approach of targets at minimizing the shuffling load of intermediate Map results. It introduces a particular repetitive structure of Map computations across the servers, and utilizes this redundancy to enable a specific type of network coding in the Shuffle phase (named coded multicasting) to minimize the communication load. We term this coding approach as “Minimum Bandwidth Code”. In , the Minimum Bandwidth Code was employed in a fully decentralized wireless distributed computing framework, achieving a scalable architecture with a constant load of communication. The other coding approach of , however, aims at minimizing the latency of Map computations by encoding the Map tasks using MDS codes, so that the run-time of the Map phase is not affected by up to a certain number of straggling servers. This coding scheme, which we term as “Minimum Latency Code”, results in a significant reduction of Map computation latency.

In this paper, we formalize a tradeoff between the computation latency in the Map phase (denoted by DD) and the communication (shuffling) load in the Shuffle phase (denoted by LL) for distributed matrix multiplication (in short, the Latency-Load Tradeoff), in which as illustrated in Fig. 1, the above two coded schemes correspond to the two extreme points that minimize LL and DD respectively. Furthermore, we propose a unified coded scheme that organically integrates both of the coding techniques, and allows to systematically operate at any point on the introduced tradeoff.

For a given computation latency, we also prove an information-theoretic lower bound on the minimum required communication load to accomplish the distributed matrix multiplication. This lower bound is proved by first concatenating multiple instances of the problem with different reduction assignments of the output vectors, and then applying the cut-set bound on subsets of servers. At the two end points of the tradeoff, the proposed scheme achieves the minimum communication load to within a constant factor.

We finally note that there has been another tradeoff between the computation load in the Map phase and the communication load in the Shuffle phase for distributed computing, which is introduced and characterized in . In this paper, we are fixing the amount of computation load (determined by the storage size) at each server, and focus on characterizing the tradeoff between the computation latency (determined by the number of servers that finish the Map computations) and the communication load. Hence, the considered tradeoff can be viewed as an extension of the tradeoff in by introducing a third axis, namely the computation latency of the Map phase.

II Problem Formulation

We perform the computations using KK distributed servers. Each server has a local memory of size μmnT\mu mnT bits (i.e., it can store equivalent of μ\mu fraction of the entries of the matrix A{\bf A}), for some 1K≤μ≤1\frac{1}{K}\leq\mu\leq 1.Thus enough information to recover the entire matrix A{\bf A} can be stored collectively on the KK servers.

The encoding matrices E1,…,EK{\bf E}_{1},\ldots,{\bf E}_{K} are design parameters and is denoted as storage design. The storage design is performed in prior to the computation.

For the Minimum Bandwidth Code in , each server stores μm\mu m rows of the matrix A{\bf A}. Thus, the rows of the encoding matrix Ek{\bf E}_{k} was chosen as a size-μm\mu m subset of the rows of the identity matrix Im{\bf I}_{m}, according to a specific repetition pattern. While for the Minimum Latency Code in , Ek{\bf E}_{k} was generated randomly such that every server stores μm\mu m random linear combinations of the rows of A{\bf A}, achieving a (μmK,m)(\mu mK,m) MDS code. \hfill□\hfill\square

II-B Distributed Computing Model

We assume that the input vectors x1,…,xN{\bf x}_{1},\ldots,{\bf x}_{N} are known to all the servers. The overall computation proceeds in three phases: Map, Shuffle, and Reduce.

Map Phase: The role of the Map phase is to compute some coded intermediate values according to the locally stored matrices in (1), which can be used later to re-construct the output vectors. More specifically, for all j=1,…,Nj=1,\ldots,N, Server kk, k=1,…,Kk=1,\ldots,K, computes the intermediate vectors

We denote the latency for Server kk to compute z1,k,…,zN,k{\bf z}_{1,k},\ldots,{\bf z}_{N,k} as SkS_{k}. We assume that S1,…,SKS_{1},\ldots,S_{K} are i.i.d. random variables, and denote the qqth order statistic, i.e., the qqth smallest variable of S1,…,SKS_{1},\ldots,S_{K} as S(q)S_{(q)}, for all q∈{1,…,K}q\in\{1,\ldots,K\}. We focus on a class of distributions of SkS_{k} such that

The Map phase terminates when a subset of servers, denoted by Q⊆{1,…,K}{\cal Q}\subseteq\{1,\ldots,K\}, have finished their Map computations in (2). A necessary condition for selecting Q{\cal Q} is that the output vectors y1…,yN{\bf y}_{1}\ldots,{\bf y}_{N} can be re-constructed by jointly utilizing the intermediate vectors calculated by the servers in Q{\cal Q}, i.e., {zj,k:j=1,…,N,k∈Q}\{{\bf z}_{j,k}:j=1,\ldots,N,k\in{\cal Q}\}. However, one can allow redundant computations in Q{\cal Q}, since if designed properly, they can be used to reduce the load of communicating intermediate results, for servers in Q{\cal Q} to recover the output vectors in the following stages of the computation.

The Minimum Bandwidth Code in waits for all servers to finish their computations, i.e., Q={1,…,K}{\cal Q}=\{1,\ldots,K\}. For the Minimum Latency Code in , Q{\cal Q} is the subset of the fastest ⌈1μ⌉\lceil\frac{1}{\mu}\rceil servers in performing the Map computations. \hfill□\hfill\square

We define the computation latency, denoted by DD, as the average amount of time spent in the Map phase. \hfill◊\hfill\Diamond

Shuffle Phase: The goal of the Shuffle phase is to exchange the intermediate values calculated in the Map phase, to help each server recover the output vectors it is responsible for. To do this, every server kk in Q{\cal Q} generates a message XkX_{k} from the locally computed intermediate vectors z1,k,…,zN,k{\bf z}_{1,k},\ldots,{\bf z}_{N,k} through an encoding function ϕk\phi_{k}, i.e., Xk=ϕk(z1,k,…,zN,k)X_{k}=\phi_{k}\left({\bf z}_{1,k},\ldots,{\bf z}_{N,k}\right), such that upon receiving all messages {Xk:k∈Q}\{X_{k}:k\in{\cal Q}\}, every server k∈Qk\in{\cal Q} can recover the output vectors in Wk{\cal W}_{k}. We assume that the servers are connected by a shared bus link. After generating XkX_{k}, Server kk multicasts XkX_{k} to all the other servers in Q{\cal Q}.

We define the communication load, denoted by LL, as the average total number of bits in all messages {Xk:k∈Q}\{X_{k}:k\in{\cal Q}\}, normalized by mTmT (i.e., the total number of bits in an output vector). \hfill◊\hfill\Diamond

Reduce Phase: The output vectors are re-constructed distributedly in the Reduce phase. Specifically, User kk, k∈Qk\in{\cal Q}, uses the locally computed vectors z1,k,…,zN,k{\bf z}_{1,k},\ldots,{\bf z}_{N,k} and the received multicast messages {Xk:k∈Q}\{X_{k}:k\in{\cal Q}\} to recover the output vectors with indices in Wk{\cal W}_{k} via a decoding function ψk\psi_{k}, i.e.,

We define the latency-load region, as the closure of the set of all achievable (D,L)(D,L) pairs. \hfill◊\hfill\Diamond

II-C Illustrating Example

In order to clarify the formulation, we use the following simple example to illustrate the latency-load pairs achieved by the two coded approaches discussed in Section I.

We consider a matrix A{\bf A} consisting of m=12m=12 rows a1,…,a12{\bf a}_{1},\ldots,{\bf a}_{12}. We have N=4N=4 input vectors x1,…,x4{\bf x}_{1},\ldots,{\bf x}_{4}, and the computation is performed on K=4K=4 servers each has a storage size μ=12\mu=\frac{1}{2}. We assume that the Map latency SkS_{k}, k=1,…,4k=1,\ldots,4, has a shifted-exponential distribution function

and by e.g., , the average latency for the fastest qq, 1≤q≤41\leq q\leq 4, servers to finish the Map computations is

Minimum Bandwidth Code . The Minimum Bandwidth Code in repeatedly stores each row of A{\bf A} at μK\mu K servers with a particular pattern, such that in the Shuffle phase, μK\mu K required intermediate values can be delivered with a single coded multicast message, which results in a coding gain of μK\mu K. We illustrate such coding technique in Fig. 2(a).

As shown in Fig. 2(a), a Minimum Bandwidth Code repeats the multiplication of each row of A{\bf A} with all input vectors x1,…,x4{\bf x}_{1},\ldots,{\bf x}_{4}, μK=2\mu K=2 times across the 44 servers, e.g., a1{\bf a}_{1} is multiplied at Server 1 and 2. The Map phase continues until all servers have finished their Map computations, achieving a computation latency D(4)=2×(1+∑j=141j)=376D(4)=2\times(1+\sum_{j=1}^{4}\frac{1}{j})=\frac{37}{6}. For k=1,2,3,4k=1,2,3,4, Server kk will be reducing output vector yk{\bf y}_{k}. In the Shuffle phase, as shown in Fig. 2(a), due to the specific repetition of Map computations, every server multicasts 33 bit-wise XORs, each of which is simultaneously useful for two other servers. For example, upon receiving a1x3⊕a3x2{\bf a}_{1}{\bf x}_{3}\oplus{\bf a}_{3}{\bf x}_{2} from Server 1, Server 2 can recover a3x2{\bf a}_{3}{\bf x}_{2} by canceling a1x3{\bf a}_{1}{\bf x}_{3} and Server 3 can recover a1x3{\bf a}_{1}{\bf x}_{3} by canceling a3x2{\bf a}_{3}{\bf x}_{2}. Similarly, every server decodes the needed values by canceling the interfering values using its local Map results. The Minimum Bandwidth Code achieves a communication load L=3×4/12=1L=3\times 4/12=1.

The Minimum Bandwidth Code can be viewed as a specific type of network coding , or more precisely index coding , in which the key idea is to design “side information” at the servers (provided by the Map results), enabling multicasting opportunities in the Shuffle phase to minimize the communication load.

Minimum Latency Code . The Minimum Latency Code in uses MDS codes to generate some redundant Map computations, and assigns the coded computations across many servers. Such type of coding takes advantage of the abundance of servers so that one can terminate the Map phase as soon as enough coded computations are performed across the network, without needing to wait for the remaining straggling servers. We illustrate such coding technique in Fig. 2(b).

For this example, a Minimum Latency Code first has each server kk, k=1,…,4k=1,\ldots,4, independently and randomly generate 66 random linear combinations of the rows of A{\bf A}, denoted by c6(k−1)+1,…,c6(k−1)+6{\bf c}_{6(k-1)+1},\ldots,{\bf c}_{6(k-1)+6} (see Fig. 2(b)). We note that {c1,…,c24}\{{\bf c}_{1},\ldots,{\bf c}_{24}\} is a (24,12)(24,12) MDS code of the rows of A{\bf A}. Therefore, for any subset D⊆{1,…,24}{\cal D}\subseteq\{1,\ldots,24\} of size ∣D∣=12|{\cal D}|=12, using the intermediate values {cixj:i∈D}\{{\bf c}_{i}{\bf x}_{j}:i\in{\cal D}\} can recover the output vector yj{\bf y}_{j}. The Map phase terminates once the fastest 22 servers have finished their computations (e.g., Server 1 and 3), achieving a computation latency D(2) ⁣= ⁣2 ⁣× ⁣(1+13+14) ⁣= ⁣196D(2)\!=\!2\!\times\!(1+\frac{1}{3}+\frac{1}{4})\!=\!\frac{19}{6}. Then Server 1 continues to reduce y1{\bf y}_{1} and y2{\bf y}_{2}, and Server 3 continues to reduce y3{\bf y}_{3} and y4{\bf y}_{4}. As illustrated in Fig. 2(b), Server 1 and 3 respectively unicasts the intermediate values it has calculated and needed by the other server to complete the computation, achieving a communication load L ⁣= ⁣6 ⁣× ⁣4/12 ⁣= ⁣2L\!=\!6\!\times\!4/12\!=\!2.

From the above descriptions, we note that the Minimum Bandwidth Code uses about twice of the time in the Map phase compared with the Minimum Latency Code, and achieves half of the communication load in the Shuffle phase. They represent the two end points of a general latency-load tradeoff characterized in the next section.

III Main Results

The main results of the paper are, 1) a characterization of a set of achievable latency-load pairs by developing a unified coded framework, 2) an outer bound of the latency-load region, which are stated in the following two theorems.

For a distributed matrix multiplication problem of computing NN output vectors using KK servers, each with a storage size μ≥1K\mu\geq\frac{1}{K}, the latency-load region contains the lower convex envelop of the points

where S(q)S_{(q)} is the qqth smallest latency of the KK i.i.d. latencies S1,…,SKS_{1},\ldots,S_{K} with some distribution FF to compute the Map functions in (2), g(K,q)g(K,q) is a function of KK and qq computed from FF, μˉ≜⌊μq⌋q\bar{\mu}\triangleq\frac{\lfloor\mu q\rfloor}{q}, Bj≜(q−1j)(K−q⌊μq⌋−j)qK(K⌊μq⌋)B_{j}\triangleq\frac{{q-1\choose j}{K-q\choose\lfloor\mu q\rfloor-j}}{\frac{q}{K}{K\choose\lfloor\mu q\rfloor}}, and sq≜inf⁡{s:∑j=s⌊μq⌋Bj≤1−μˉ}s_{q}\triangleq\inf\{s:\sum_{j=s}^{\lfloor\mu q\rfloor}B_{j}\leq 1-\bar{\mu}\}.

We prove Theorem 1 In Section IV, in which we present a unified coded scheme that jointly designs the storage and the data shuffling, which achieves the latency in (8) and the communication load in (9).

We numerically evaluate in Fig. 3 the latency-load pairs achieved by the proposed coded framework, for computing N ⁣= ⁣180N\!=\!180 output vectors using K ⁣= ⁣18K\!=\!18 servers each with a storage size μ ⁣= ⁣1/3\mu\!=\!1/3. The achieved tradeoff approximately exhibits an inverse-linearly proportional relationship between the latency and the load. For instance, doubling the latency from 120 to 240 results in a drop of the communication load from 43 to 23 by a factor of 1.87.\hfill□\hfill\square

The key idea to achieve D(q)D(q) and L(q)L(q) in Theorem 1 is to design the concatenation of the MDS code and the repetitive executions of the Map computations, in order to take advantage of both the Minimum Latency Code and the Minimum Bandwidth Code. More specifically, we first generate Kqm\frac{K}{q}m MDS-coded rows of A{\bf A}, and then store each of them ⌊μq⌋\lfloor\mu q\rfloor times across the KK servers in a specific pattern. As a result, any subset of qq servers would have sufficient amount of intermediate results to reduce the output vectors, and we end the Map phase as soon as the fastest qq servers finish their Map computations, achieving the latency in (8).

We also exploit coded multicasting in the Shuffle phase to reduce the communication load. In the load expression (9), BjB_{j}, j≤⌊μq⌋j\leq\lfloor\mu q\rfloor, represents the (normalized) number of coded rows of A{\bf A} repeatedly stored/computed at jj servers. By multicasting coded packets simultaneously useful for jj servers, BjB_{j} intermediate values can be delivered to a server with a communication load of Bjj\frac{B_{j}}{j}, achieving a coding gain of jj. We greedily utilize the coding opportunities with a larger coding gain until we get close to satisfying the demand of each server, which accounts for the first term in (9). Then the second term results from two follow-up strategies 1) communicate the rest of the demands uncodedly 2) continue coded multicasting with a smaller coding gain (i.e., j=sq−1j=s_{q}-1), which may however deliver more than what is needed for reduction. \hfill□\hfill\square

The latency-load region is contained in the lower convex envelop of the points

We prove Theorem 2 in Section V, by deriving an information-theoretic lower bound on the minimum required communication load for a given computation latency, using any storage design and data shuffling scheme.

We numerically compare the outer bound in Theorem 2 and the achieved inner bound in Theorem 1 in Fig. 3, from which we make the following observations.

For the intermediate latency from 70 to 270, the communication load achieved by the proposed scheme is within a multiplicative gap of at most 4.2×4.2\times from the lower bound. In general, a complete characterization of the latency-load region (or an approximation to within a constant gap for all system parameters) remains open.\hfill□\hfill\square

IV Proposed Coded Framework

In this section, we prove Theorem 1 by proposing and analyzing a general coded framework that achieves the latency-load pairs in (7). We first demonstrate the key ideas of the proposed scheme through the following example, and then give the general description of the scheme.

We assume that we can afford to wait for q=4q=4 servers to finish their computations in the Map phase, and we describe the proposed storage design and shuffling scheme.

WLOG, due to the symmetry of the storage design, we assume that Servers 11, 22, 33 and 44 are the first 44 servers that finish their Map computations. Then we assign the Reduce tasks such that Server kk reduces the output vectors y3(k−1)+1{\bf y}_{3(k-1)+1}, y3(k−1)+2{\bf y}_{3(k-1)+2} and y3(k−1)+3{\bf y}_{3(k-1)+3}, for all k∈{1,…,4}k\in\{1,\ldots,4\}.

After the Map phase, Server 1 has computed the intermediate values {c1xj,…,c10xj:j=1,…,12}\{{\bf c}_{1}{\bf x}_{j},\ldots,{\bf c}_{10}{\bf x}_{j}:j=1,\ldots,12\}. For Server 1 to recover y1=Ax1{\bf y}_{1}={\bf A}{\bf x}_{1}, it needs any subset of 10 intermediate values cix1{\bf c}_{i}{\bf x}_{1} with i∈{11,…,30}i\in\{11,\ldots,30\} from Server 22, 33 and 44 in the Shuffle phase. Similar data demands hold for all 4 servers and the output vectors they are reducing. Therefore, the goal of the Shuffle phase is to exchange these needed intermediate values to accomplish successful reductions.

Coded Shuffle. We first group the 4 servers into 4 subsets of size 3 and perform coded shuffling within each subset. We illustrate the coded shuffling scheme for Servers 11, 22 and 33 in Fig. 5. Each server multicasts 33 bit-wise XORs, denoted by ⊕\oplus, of the locally computed intermediate values to the other two. The intermediate values used to create the multicast messages are the ones known exclusively at two servers and needed by another one. After receiving 22 multicast messages, each server recovers 66 needed intermediate values. For instance, Server 1 recovers c11x1{\bf c}_{11}{\bf x}_{1}, c11x2{\bf c}_{11}{\bf x}_{2} and c11x3{\bf c}_{11}{\bf x}_{3} by canceling c2x7{\bf c}_{2}{\bf x}_{7}, c2x8{\bf c}_{2}{\bf x}_{8} and c2x9{\bf c}_{2}{\bf x}_{9} respectively, and then recovers c12x1{\bf c}_{12}{\bf x}_{1}, c12x2{\bf c}_{12}{\bf x}_{2} and c12x3{\bf c}_{12}{\bf x}_{3} by canceling c4x4{\bf c}_{4}{\bf x}_{4}, c4x5{\bf c}_{4}{\bf x}_{5} and c4x6{\bf c}_{4}{\bf x}_{6} respectively.

Similarly, we perform the above coded shuffling in Fig. 5 for another 33 subsets of 33 servers. After coded multicasting within the 44 subsets of 33 servers, each server recovers 1818 needed intermediate values (6 for each of the output vector it is reducing). As mentioned before, since each server needs a total of 3×(20−10)=303\times(20-10)=30 intermediate values to reduce the 3 assigned output vectors, it needs another 30−18=1230-18=12 after decoding all multicast messages. We satisfy the residual data demands by simply having the servers unicast enough (i.e., 12×4=4812\times 4=48) intermediate values for reduction. Overall, 9×4+48=849\times 4+48=84 (possibly coded) intermediate values are communicated, achieving a communication load of L=4.2L=4.2.

IV-B General Scheme

We first describe the storage design, Map phase computation and the data shuffling scheme that achieves the latency-load pairs (D(q),L(q))(D(q),L(q)) in (7), for all q∈{⌈1μ⌉,…,K}q\in\{\lceil\frac{1}{\mu}\rceil,\ldots,K\}. Given these achieved pairs, we can “memory share” across them to achieve their lower convex envelop as stated in Theorem 1.

Storage Design. We first use a (Kqm,m)(\frac{K}{q}m,m) MDS code to encode the mm rows of matrix A{\bf A} into Kqm\frac{K}{q}m coded rows c1…,cKqm{\bf c}_{1}\ldots,{\bf c}_{\frac{K}{q}m} (e.g., Kqm\frac{K}{q}m random linear combinations of the rows of A{\bf A}). Then as shown in Fig. 6, we evenly partitioned the Kqm\frac{K}{q}m coded rows into (Kμq){K\choose\mu q} disjoint batches, each containing a subset of mqK(Kμq)\frac{m}{\frac{q}{K}{K\choose\mu q}} coded rows. We focus on matrix multiplication problems for large matrices, and assume that m≫qK(Kμq)m\gg\frac{q}{K}{K\choose\mu q}, for all q∈{1μ,…,K}q\in\{\frac{1}{\mu},\ldots,K\}. Each batch, denoted by BT{\cal B}_{\cal T}, is labelled by a unique subset T⊂{1,…,K}\mathcal{T}\subset\{1,\ldots,K\} of size ∣T∣=μq|{\cal T}|=\mu q. That is

Server kk, k∈{1,…,K}k\in\{1,\ldots,K\} stores the coded rows in BT\mathcal{B}_{\cal T} as the rows of Uk{\bf U}_{k} if k∈Tk\in\mathcal{T}.

In the above example, q=4q=4, and Kqm=64×20=30\frac{K}{q}m=\frac{6}{4}\times 20=30 coded rows of A{\bf A} are partitioned into (Kμq)=(62)=15{K\choose\mu q}={6\choose 2}=15 batches each containing 3015=2\frac{30}{15}=2 coded rows. Every node is in 55 subsets of size two, thus storing 5×2=105\times 2=10 coded rows of A{\bf A}.

Map Phase Execution. Each server computes the inner products between each of the locally stored coded rows of A{\bf A} and each of the input vectors, i.e., Server kk computes cixj{\bf c}_{i}{\bf x}_{j} for all j=1,…,Nj=1,\ldots,N, and all i∈{BT:k∈T}i\in\{{\cal B}_{\cal T}:k\in{\cal T}\}. We wait for the fastest qq servers to finish their Map computations before halting the Map phase, achieving a computation latency D(q)D(q) in (8). We denote the set of indices of these servers as Q{\cal Q}.

The computation then moves on exclusively over the qq servers in Q{\cal Q}, each of which is assigned to reduce Nq\frac{N}{q} out of the NN output vectors y1=Ax1,…,yN=AxN{\bf y}_{1}={\bf A}{\bf x}_{1},\ldots,{\bf y}_{N}={\bf A}{\bf x}_{N}.

For a feasible shuffling scheme to exist such that the Reduce phase can be successfully carried out, every subset of qq servers (since we cannot predict which qq servers will finish first) should have collectively stored at least mm distinct coded rows ci{\bf c}_{i} for i∈{1,…,Kqm}i\in\{1,\ldots,\frac{K}{q}m\}. Next, we explain how our proposed storage design meets this requirement. First, the qq servers in Q{\cal Q} collectively provide a storage size equivalent to μqm\mu qm rows. Then since each coded row is stored by μq\mu q out of all KK servers, it can be stored by at most μq\mu q servers in Q{\cal Q}, and thus servers in Q{\cal Q} collectively store at least μqmμq=m\frac{\mu qm}{\mu q}=m distinct coded rows.

Coded Shuffle. For S⊂Q{\cal S}\subset{\cal Q} and k∈Q\Sk\in{\cal Q}\backslash{\cal S}, we denote the set of intermediate values needed by Server kk and known exclusively by the servers in S\mathcal{S} as VSk\mathcal{V}_{\mathcal{S}}^{k}. More formally:

Due to the proposed storage design, for a particular S{\cal S} of size jj, VSk\mathcal{V}_{\mathcal{S}}^{k} contains Nq⋅(K−qμq−j)mqK(Kμq)\frac{N}{q}\cdot\frac{{K-q\choose\mu q-j}m}{\frac{q}{K}{K\choose\mu q}} intermediate values.

In the above example, we have V{2,3}1={c11xj,c12xj:j=1,2,3}\mathcal{V}_{\{2,3\}}^{1}=\{{\bf c}_{11}{\bf x}_{j},{\bf c}_{12}{\bf x}_{j}:j=1,2,3\}, V{1,3}2={c3xj,c4xj:j=4,5,6}\mathcal{V}_{\{1,3\}}^{2}=\{{\bf c}_{3}{\bf x}_{j},{\bf c}_{4}{\bf x}_{j}:j=4,5,6\}, and V{1,2}3={c1xj,c2xj:j=7,8,9}\mathcal{V}_{\{1,2\}}^{3}=\{{\bf c}_{1}{\bf x}_{j},{\bf c}_{2}{\bf x}_{j}:j=7,8,9\}.

In the Shuffle phase, servers in Q{\cal Q} create and multicast coded packets that are simultaneously useful for multiple other servers, until every server in Q{\cal Q} recovers at least mm intermediate values for each of the output vectors it is reducing. The proposed shuffling scheme is greedy in the sense that every server in Q{\cal Q} will always try to multicast coded packets simultaneously useful for the largest number of servers.

The proposed shuffle scheme proceeds as follows. For each j ⁣= ⁣μq,μq−1,…,sqj\!=\!\mu q,\mu q-1,\ldots,s_{q}, where sq ⁣≜ ⁣inf⁡{s: ⁣∑j=sμq ⁣(q−1j)(K−qμq−j)qK(Kμq) ⁣≤ ⁣1 ⁣− ⁣μ}s_{q}\!\triangleq\!\inf\{s:\!\sum_{j=s}^{\mu q}\!\frac{{q-1\choose j}{K-q\choose\mu q-j}}{\frac{q}{K}{K\choose\mu q}}\!\leq\!1\!-\!\mu\}, and every subset S ⁣⊆ ⁣Q\mathcal{S}\!\subseteq\!{\cal Q} of size j ⁣+ ⁣1j\!+\!1:

For each k∈Sk\in\mathcal{S}, we evenly and arbitrarily split VS\{k}k\mathcal{V}_{\mathcal{S}\backslash\{k\}}^{k} into jj disjoint segments VS\{k}k ⁣= ⁣{VS\{k},ik ⁣: ⁣i∈S\{k}}\mathcal{V}^{k}_{\mathcal{S}\backslash\{k\}}\!=\!\{\mathcal{V}_{\mathcal{S}\backslash\{k\},i}^{k}\!:\!i\in{\cal S}\backslash\{k\}\}, and associate the segment VS\{k},ik\mathcal{V}_{\mathcal{S}\backslash\{k\},i}^{k} with the server i∈S\{k}i\in{\cal S}\backslash\{k\}.

Server ii, i∈Si\in\mathcal{S}, multicasts the bit-wise XOR, denoted by ⊕\oplus, of all the segments associated with it in S{\cal S}, i.e., Server ii multicasts ⊕k∈S\{i}VS\{k},ik\underset{k\in\mathcal{S}\backslash\{i\}}{\oplus}\mathcal{V}^{k}_{\mathcal{S}\backslash\{k\},i} to the other servers in S\{i}{\cal S}\backslash\{i\}.

For every pair of servers kk and ii in S{\cal S}, since Server kk has computed locally the segments VS\{k′},ik′\mathcal{V}^{k^{\prime}}_{\mathcal{S}\backslash\{k^{\prime}\},i} for all k′∈S\{i,k}k^{\prime}\in\mathcal{S}\backslash\{i,k\}, it can cancel them from the message ⊕k∈S\{i}VS\{k},ik\underset{k\in\mathcal{S}\backslash\{i\}}{\oplus}\mathcal{V}^{k}_{\mathcal{S}\backslash\{k\},i} sent by Server ii, and recover the intended segment VS\{k},ik\mathcal{V}^{k}_{\mathcal{S}\backslash\{k\},i}.

For each jj in the above coded shuffling scheme, each server in Q{\cal Q} recovers (q−1j)(K−qμq−j)mqK(Kμq){q-1\choose j}\frac{{K-q\choose\mu q-j}m}{\frac{q}{K}{K\choose\mu q}}intermediate values for each of the output vectors it is reducing. Therefore, j=sq+1j=s_{q}+1 is the smallest size of the subsets in which the above coded multicasting needs to be performed, before enough number of intermediate values for reduction are delivered.

In each subset S{\cal S} of size jj, since each server i∈Si\in{\cal S} multicasts a coded segment of size ∣VS\{k}k∣j\frac{|{\cal V}^{k}_{{\cal S}\backslash\{k\}}|}{j} for some k≠ik\neq i, the total communication load so far, for Bj=(q−1j)(K−qμq−j)qK(Kμq)B_{j}=\frac{{q-1\choose j}{K-q\choose\mu q-j}}{\frac{q}{K}{K\choose\mu q}}, is

Next, we can continue to finish the data shuffling in two different ways. The first approach is to have the servers in Q{\cal Q} communicate with each other uncoded intermediate values, until every server has exactly mm intermediate values for each of the output vector it is responsible for. Using this approach, we will have a total communication load of

The second approach is to continue the above 2 steps for j=sq−1j=s_{q}-1. Using this approach, we will have a total communication load of L2=∑j=sq−1μqNBjjL_{2}=\sum_{j=s_{q}-1}^{\mu q}N\frac{B_{j}}{j}.

Then we take the approach with less communication load, and achieve L(q)=min⁡{L1,L2}L(q)=\min\{L_{1},L_{2}\}.

The ideas of efficiently creating and exploiting coded multicasting opportunities have been introduced in caching problems . In this section, we illustrated how to create and utilize such coding opportunities in distributed computing to slash the communication load, when facing with straggling servers. \hfill□\hfill\square

V Converse

In this section, we prove the outer bound on the latency-load region in Theorem 2.

We start by considering a distributed matrix multiplication scheme that stops the Map phase when qq servers have finished their computations. For such scheme, as given by (8), the computation latency D(q)D(q) is the expected value of the qqth order statistic of the Map computation times at the KK servers. WLOG, we can assume that Servers 1,…,q1,\ldots,q first finish their Map computations, and they will be responsible for reducing the NN output vectors y1,…,yN{\bf y}_{1},\ldots,{\bf y}_{N}.

To proceed, we first partition the y1,…,yN{\bf y}_{1},\ldots,{\bf y}_{N} into qq groups G1,…,Gq{\cal G}_{1},\ldots,{\cal G}_{q} each of size N/qN/q, and define the output assignment

where WkA{\cal W}_{k}^{\cal A} denotes the group of output vectors reduced by Server kk in the output assignment A{\cal A}.

Next we choose an integer t∈{1,…,q−1}t\in\{1,\ldots,q-1\}, and consider the following ⌈qt⌉\lceil\frac{q}{t}\rceil output assignments which are circular shifts of (G1,…,Gq)\left({\cal G}_{1},\ldots,{\cal G}_{q}\right) with step size tt,

We note that by the Map computation in (2), at each server all the input vectors x1,…,xN{\bf x}_{1},\ldots,{\bf x}_{N} are multiplied by the same matrix (i.e., Uk{\bf U}_{k} at Server kk). Therefore, for the same set of qq servers and their storage contents, a feasible data shuffling scheme for one of the above output assignments is also feasible for all other ⌈qt⌉−1\lceil\frac{q}{t}\rceil-1 assignments by relabelling the output vectors. As a result, the minimum communication loads for all of the above output assignments are identical. \hfill□\hfill\square

For a shuffling scheme admitting an output assignment A{\cal A}, we denote the message sent by Server k∈{1,…,q}k\in\{1,\ldots,q\} as XkAX_{k}^{\mathcal{A}}, with a size of RkAmTR_{k}^{\mathcal{A}}mT bits.

Now we focus on the Servers 1,…,t1,\ldots,t and consider the compound setting that includes all ⌈qt⌉\lceil\frac{q}{t}\rceil output assignments in (17). We observe that as shown in Fig. 7, in this compound setting, the first tt servers should be able to recover all output vectors (y1…,yN)=(G1,…,Gq)({\bf y}_{1}\ldots,{\bf y}_{N})=({\cal G}_{1},\ldots,{\cal G}_{q}) using their local computation results {Ukx1,…,UkxN:k=1,…,t}\{{\bf U}_{k}{\bf x}_{1},\ldots,{\bf U}_{k}{\bf x}_{N}:k=1,\ldots,t\} and the received messages in all the output assignments {XkA1,…,XkA⌈qt⌉:k=t+1,…,q}\{X_{k}^{{\cal A}_{1}},\ldots,X_{k}^{{\cal A}_{\lceil\frac{q}{t}\rceil}}:k=t+1,\ldots,q\}. Thus we have the following cut-set bound for the first tt servers.

Next we consider qq subsets of servers each with size tt: Ni≜{i,(i+1),…,(i+t−1)}\mathcal{N}_{i}\triangleq\{i,(i+1),\ldots,(i+t-1)\}, i=1,…,qi=1,\ldots,q, where the addition is modular qq. Similarly, we have the following cut-set bound for Ni{\cal N}_{i}:

Summing up these qq cut-set bounds, we have

where (a) results from the fact mentioned in Remark 8 that the communication load is independent of the output assignment.

Since (22) holds for all t=1,…,q−1t=1,\ldots,q-1, we have

In this appendix, we prove that when all KK servers finish their Map computations, i.e., Q={1,…,K}{\cal Q}=\{1,\ldots,K\} and we operate at the point with the maximum latency, the communication load achieved by the proposed coded scheme (or the Minimum Bandwidth Code) is within a constant multiplicative factor of the lower bound on the communication load in Theorem 2. More specifically,

when μK\mu K is an integer,This always holds true for large KK. where L(K)L(K) and Lˉ(K)\bar{L}(K) are respectively given by (9) and (11).

We proceed to bound the RHS of (25) in the following two cases:

Since μK≥1\mu K\geq 1, we have K−1≥⌈K2⌉≥⌈12μ⌉K-1\geq\lceil\frac{K}{2}\rceil\geq\lceil\frac{1}{2\mu}\rceil.

In this case, we set t=⌈12μ⌉t=\lceil\frac{1}{2\mu}\rceil in (25) to have

Comparing (26) and (30) completes the proof. \hfill■\hfill\blacksquare

References