A Fundamental Tradeoff between Computation and Communication in Distributed Computing

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

I Introduction

We consider a general distributed computing framework, motivated by prevalent structures like MapReduce and Spark , in which the overall computation is decomposed into two stages: “Map” and “Reduce”. Firstly in the Map stage, distributed computing nodes process parts of the input data locally, generating some intermediate values according to their designed Map functions. Next, they exchange the calculated intermediate values among each other (a.k.a. data shuffling), in order to calculate the final output results distributedly using their designed Reduce functions.

Within this framework, data shuffling often appears to limit the performance of distributed computing applications, including self-join , tera-sort , and machine learning algorithms . For example, in a Facebook’s Hadoop cluster, it is observed that 33% of the overall job execution time is spent on data shuffling . Also as is observed in , 70% of the overall job execution time is spent on data shuffling when running a self-join application on an Amazon EC2 cluster . As such motivated, we ask this fundamental question that if coding can help distributed computing in reducing the load of communication and speeding up the overall computation? Coding is known to be helpful in coping with the channel uncertainty in telecommunication and also in reducing the storage cost in distributed storage systems and cache networks. In this work, we extend the application of coding to distributed computing and propose a framework to substantially reduce the load of data shuffling via coding and some extra computing in the Map phase.

More specifically, we formulate and characterize a fundamental tradeoff relationship between “computation load” in the Map phase and “communication load” in the data shuffling phase, and demonstrate that the two are inversely proportional to each other. We propose an optimal coded scheme, named “Coded Distributed Computing” (CDC), which demonstrates that increasing the computation load of the Map phase by a factor of rr (i.e., evaluating each Map function at rr carefully chosen nodes) can create novel coding opportunities in the data shuffling phase that reduce the communication load by the same factor.

To illustrate our main result, consider a distributed computing framework to compute QQ arbitrary output functions from NN input files, using KK distributed computing nodes. As mentioned earlier, the overall computation is performed by computing a set of Map and Reduce functions distributedly across the KK nodes. In the Map phase, each input file is processed locally, in one of the nodes, to generate QQ intermediate values, each corresponding to one of the QQ output functions. Thus, at the end of this phase, QNQN intermediate values are calculated, which can be split into QQ subsets of NN intermediate values and each subset is needed to calculate one of the output functions. In the Shuffle phase, for every output function to be calculated, all NN intermediate values corresponding to that function are transferred to one of the nodes for reduction. Of course, depending on the node that has been chosen to reduce an output function, a part of the intermediate values are already available locally, and do not need to be transferred in the Shuffle phase. This is because that the Map phase has been carried out on the same set of nodes, and the results of mapping done at a node can remain in that node to be used for the Reduce phase. This offers some saving in the load of communication. To reduce the communication load even more, we may map each input file in more than one nodes. Apparently, this increases the fraction of intermediate values that are locally available. However, as we will show, there is a better way to exploit this redundancy in computation to reduce the communication load. The main message of this paper is to show that following a particular patten in repeating Map computations along with some coding techniques, we can significantly reduce the load of communication. Perhaps surprisingly, we show that the gain of coding in reducing communication load scales with the size of the network.

To be more precise, we define the computation load rr, 1≤r≤K1\leq r\leq K, as the total number of computed Map functions at the nodes, normalized by NN. For example, r=1r=1 means that none of the Map functions has been re-computed, and r=2r=2 means that on average each Map function can be computed on two nodes. We also define communication load LL, 0≤L≤10\leq L\leq 1, as the total amount of information exchanged across nodes in the shuffling phase, normalized by the size of QNQN intermediate values, in order to compute the QQ output functions disjointly and uniformly across the KK nodes. Based on this formulation, we now ask the following fundamental question:

Given a computation load rr in the Map phase, what is the minimum communication load L∗(r)L^{*}(r), using any data shuffling scheme, needed to compute the final output functions?

We propose Coded Distributed Computing (CDC) that achieves a communication load of Lcoded(r)=1r⋅(1−rK)L_{\textup{coded}}(r)=\frac{1}{r}\cdot(1-\frac{r}{K}) for r=1,…,Kr=1,\ldots,K, and the lower convex envelop of these points. CDC employs a specific strategy to assign the computations of the Map and Reduce functions across the computing nodes, in order to enable novel coding opportunities for data shuffling. In particular, for a computation load r∈{1,…,K}r\in\{1,\ldots,K\}, CDC utilizes a carefully designed repetitive mapping of data blocks at rr distinct nodes to create coded multicast messages that deliver data simultaneously to a subset of r≥1r\geq 1 nodes. Hence, compared with an uncoded data shuffling scheme, which as we show later achieves a communication load Luncoded(r)=1−rKL_{\textup{uncoded}}(r)=1-\frac{r}{K}, CDC is able to reduce the communication load by exactly a factor of the computation load rr. Furthermore, the proposed CDC scheme applies to a more general distributed computing framework where every output function is computed by more than one, or particularly s∈{1,…,K}s\in\{1,\ldots,K\} nodes, which provides better fault-tolerance in distributed computing.

We numerically compare the computation-communication tradeoffs of CDC and uncoded data shuffling schemes (i.e., Lcoded(r)L_{\textup{coded}}(r) and Luncoded(r)L_{\textup{uncoded}}(r)) in Fig. 1. As it is illustrated, in the uncoded scheme that achieves a communication load Luncoded(r)=1−rKL_{\textup{uncoded}}(r)=1-\frac{r}{K}, increasing the computation load rr offers only a modest reduction in communication load. In fact for any rr, this gain vanishes for large number of nodes KK. Consequently, it is not justified to trade computation for communication using uncoded schemes. However, for the coded scheme that achieves a communication load of Lcoded(r)=1r⋅(1−rK)L_{\textup{coded}}(r)=\frac{1}{r}\cdot(1-\frac{r}{K}), increasing the computation load rr will significantly reduce the communication load, and this gain does not vanish for large KK. For example as illustrated in Fig. 1, when mapping each file at one extra node (r=2r=2), CDC reduces the communication load by 55.6%, while the uncoded scheme only reduces it by 11.1%.

We also prove an information-theoretic lower bound on the minimum communication load L∗(r)L^{*}(r). To prove the lower bound, we derive a lower bound on the total number of bits communicated by any subset of nodes, using induction on the size of the subset. To derive the lower bound for a particular subset of nodes, we first establish a lower bound on the number of bits needed by one of the nodes to recover the intermediate values it needs to calculate its assigned output functions, and then utilize the bound on the number of bits communicated by the rest of the nodes in that subset, which is given by the inductive argument. The derived lower bound on L∗(r)L^{*}(r) matches the communication load achieved by the CDC scheme for any computation load 1≤r≤K1\leq r\leq K. As a result, we exactly characterize the optimal tradeoff between computation load and communication load in the following:

For general 1≤r≤K1\leq r\leq K, L∗(r)L^{*}(r) is the lower convex envelop of the above points {(r,Lcoded(r)):r∈{1,…,K}}\{(r,L_{\textup{coded}}(r)):r\in\{1,\ldots,K\}\}. Note that for large KK, 1r⋅(1−rK)≈1r\frac{1}{r}\cdot(1-\frac{r}{K})\approx\frac{1}{r}, hence L∗(r)≈1rL^{*}(r)\approx\frac{1}{r}. This result reveals a fundamental inversely proportional relationship between computation load and communication load in distributed computing. This also illustrates that the gain of 1r\frac{1}{r} achieved by CDC is optimal and it cannot be improved by any other scheme (since Lcoded(r)L_{\textup{coded}}(r) is an information-theoretic lower bound on L∗(r)L^{*}(r) that applies to any data shuffling scheme).

Having theoretically characterized the optimal computation-communication tradeoff achieved by the proposed CDC scheme, we also empirically demonstrate the practical impact of this tradeoff. In particular, we apply the coding techniques of CDC to a widely used Hadoop sorting benchmark TeraSort , developing a novel coded distributed sorting algorithm CodedTeraSort . We perform extensive experiments on Amazon EC2 clusters, and observe that for typical settings of interest, CodedTeraSort speeds up the overall execution of the conventional TeraSort by a factor of 1.97×1.97\times - 3.39×3.39\times.

Finally, we discuss some future directions to extend the results of this work. In particular, we consider topics including heterogeneous networks with asymmetric tasks, straggling/failing computing nodes, multi-stage computation tasks, multi-layer networks and structured topology, joint storage and computation optimization, and coded edge/fog computing.

Related Works. The problem of characterizing the minimum communication for distributed computing has been previously considered in several settings in both computer science and information theory communities. In , a basic computing model is proposed, where two parities have xx and yy and aim to compute a boolean function f(x,y)f(x,y) by exchanging the minimum number of bits between them. Also, the problem of minimizing the required communication for computing the modulo-two sum of distributed binary sources with symmetric joint distribution was introduced in . Following these two seminal works, a wide range of communication problems in the scope of distributed computing have been studied (see, e.g., ). The key differences distinguishing the setting in this paper from most of the prior ones are 1) We focus on the flow of communication in a general distributed computing framework, motivated by MapReduce, rather than the structures of the functions or the input distributions. 2) We do not impose any constraint on the numbers of output results, input data files and computing nodes (they can be arbitrarily large), 3) We do not assume any special property (e.g. linearity) of the computed functions.

The idea of efficiently creating and exploiting coded multicasting was initially proposed in the context of cache networks in , and extended in , where caches pre-fetch part of the content in a way to enable coding during the content delivery, minimizing the network traffic. In this paper, we propose a framework to study the tradeoff between computation and communication in distributed computing. We demonstrate that the coded multicasting opportunities exploited in the above caching problems also exist in the data shuffling of distributed computing frameworks, which can be created by a strategy of repeating the computations of the Map functions specified by the Coded Distributed Computing (CDC) scheme.

Finally, in a recent work , the authors have proposed methods for utilizing codes to speed up some specific distributed machine learning algorithms. The considered problem in this paper differs from in the following aspects. We propose a general methodology for utilizing coding in data shuffling that can be applied to any distributed computing framework with a MapReduce structure, regardless of the underlying application. In other words, any distributed computing algorithm that fits in the MapReduce framework can benefit from the proposed CDC solution. We also characterize the information-theoretic computation-communication tradeoff in such frameworks. Furthermore, the coding used in is at the application layer (i.e., applying computation on coded data), while in this paper we focus on applying codes directly on the shuffled data.

II Problem Formulation

In this section, we formulate a general distributed computing framework motivated by MapReduce, and define the function characterizing the tradeoff between computation and communication.

Motivated by MapReduce, we assume that as illustrated in Fig. 2 the computation of the output function ϕq\phi_{q}, q∈{1,…,Q}q\in\{1,\ldots,Q\} can be decomposed as follows:

Note that for every set of output functions ϕ1,…,ϕQ\phi_{1},\ldots,\phi_{Q} such a Map-Reduce decomposition exists (e.g., setting gq,ng_{q,n}′s{}^{\prime}s to identity functions such that gq,n(wn)=wng_{q,n}(w_{n})=w_{n} for all n=1,…,Nn=1,\ldots,N, and hqh_{q} to ϕq\phi_{q} in (1)). However, such a decomposition is not unique, and in the distributed computing literature, there has been quite some work on developing appropriate decompositions of computations like join, sorting and matrix multiplication (see, e.g., ), for them to be performed efficiently in a distributed manner. Here we do not impose any constraint on how the Map and Reduce functions are chosen (for example, they can be arbitrary linear or non-linear functions). \hfill□\hfill\square

The above computation is carried out by KK distributed computing nodes, labelled as Node 11,…, ,\ldots,\,Node KK. They are interconnected through a multicast network. Following the above decomposition, the computation proceeds in three phases: Map, Shuffle and Reduce.

Map Phase: Node kk, k∈{1,…,K}k\in\{1,\ldots,K\} computes the Map functions of a set of files Mk\mathcal{M}_{k}, which are stored on Node kk, for some design parameter Mk⊆{w1,…,wN}\mathcal{M}_{k}\subseteq\{w_{1},\ldots,w_{N}\}. For each file wnw_{n} in Mk\mathcal{M}_{k}, Node kk computes g⃗n(wn) ⁣= ⁣(v1,n,…,vQ,n)\vec{g}_{n}(w_{n})\!=\!(v_{1,n},\ldots,v_{Q,n}). We assume that each file is mapped by at least one node, i.e., ∪k=1,…,KMk={w1,…,wN}\underset{k=1,\ldots,K}{\cup}\mathcal{M}_{k}=\{w_{1},\ldots,w_{N}\}.

We define the computation load, denoted by rr, 1≤r≤K1\leq r\leq K, as the total number of Map functions computed across the KK nodes, normalized by the number of files NN, i.e., r≜∑k=1K∣Mk∣Nr\triangleq\frac{\sum_{k=1}^{K}|\mathcal{M}_{k}|}{N}. The computation load rr can be interpreted as the average number of nodes that map each file. \hfill◊\hfill\Diamond

Beyond the symmetric task assignment considered in this paper, characterizing the optimal computation-communication tradeoff allowing general asymmetric task assignments is a challenging open problem. As the first step to study this problem, in our follow-up work in which the number of output functions QQ is fixed and the computing resources are abundant (e.g., number of computing nodes K≫QK\gg Q), we have shown that asymmetric task assignments can do better than the symmetric ones, and achieve the optimum run-time performance. \hfill□\hfill\square

Having generated the message XkX_{k}, Node kk multicasts it to all other nodes.

By the end of the Shuffle phase, each of the KK nodes receives X1,…,XKX_{1},\ldots,X_{K} free of error.

Finally, Node kk, k∈{1,…,K}k\in\{1,\ldots,K\}, computes the Reduce function uq=hq(vq,1…vq,N)u_{q}=h_{q}(v_{q,1}\ldots v_{q,N}) for all q∈Wkq\in\mathcal{W}_{k}.

We define the computation-communication function of the distributed computing framework

L∗(r)L^{*}(r) characterizes the optimal tradeoff between computation and communication in this framework. \hfill◊\hfill\Diamond

Example (Uncoded Scheme). In the Shuffle phase of a simple “uncoded” scheme, each node receives the needed intermediate values sent uncodedly by some other nodes. Since a total of QNQN intermediate values are needed across the KK nodes and rN⋅QK=rQNKrN\cdot\frac{Q}{K}=\frac{rQN}{K} of them are already available after the Map phase, the communication load achieved by the uncoded scheme

After the Map phase, each node knows the intermediate values of all QQ output functions in the files it has mapped. Therefore, for a fixed file assignment and any symmetric assignment of the Reduce functions, specified by W1,…,WK{\cal W}_{1},\ldots,{\cal W}_{K}, we can satisfy the data requirements using the same data shuffling scheme up to relabelling the Reduce functions. In other words, the communication load is independent of the assignment of the Reduce functions. \hfill□\hfill\square

III Main Results

The computation-communication function of the distributed computing framework, L∗(r)L^{*}(r) is given by

for sufficiently large TT. For general 1≤r≤K1\leq r\leq K, L∗(r)L^{*}(r) is the lower convex envelop of the above points {(r,1r⋅(1−rK)):r∈{1,…,K}}\{(r,\frac{1}{r}\cdot(1-\frac{r}{K})):r\in\{1,\ldots,K\}\}.

We prove the achievability of Theorem 1 by proposing a coded scheme, named Coded Distributed Computing, in Section V. We demonstrate that no other scheme can achieve a communication load smaller than the lower convex envelop of the points {(r,1r⋅(1−rK)):r∈{1,…,K}}\{(r,\frac{1}{r}\cdot(1-\frac{r}{K})):r\in\{1,\ldots,K\}\} by proving the converse in Section VI.

Theorem 1 exactly characterizes the optimal tradeoff between the computation load and the communication load in the considered distributed computing framework. \hfill□\hfill\square

For r∈{1,…,K}r\in\{1,\ldots,K\}, the communication load achieved in Theorem 1 is less than that of the uncoded scheme in (5) by a multiplicative factor of rr, which equals the computation load and can grow unboundedly as the number of nodes KK increases if e.g. r=Θ(K)r=\Theta(K). As illustrated in Fig. 1 in Section I, while the communication load of the uncoded scheme decreases linearly as the computation load increases, Lcoded(r)L_{\textup{coded}}(r) achieved in Theorem 1 is inversely proportional to the computation load.\hfill□\hfill\square

While increasing the computation load rr causes a longer Map phase, the coded achievable scheme of Theorem 1 maximizes the reduction of the communication load using the extra computations. Therefore, Theorem 1 provides an analytical framework to optimally trading the computation power in the Map phase for more bandwidth in the Shuffle phase, which helps to minimize the overall execution time of applications whose performances are limited by data shuffling. \hfill□\hfill\square

The computation-communication function of the cascaded distributed computing framework, L∗(r,s)L^{*}(r,s), for r∈{1,…,K}r\in\{1,\ldots,K\}, is characterized by

for some s∈{1,…,K}s\in\{1,\ldots,K\} and sufficiently large TT. For general 1≤r≤K1\leq r\leq K, L∗(r,s)L^{*}(r,s) is the lower convex envelop of the above points {(r,Lcoded(r,s)):r∈{1,…,K}}\{(r,L_{\textup{coded}}(r,s)):r\in\{1,\ldots,K\}\}.

We present the Coded Distributed Computing scheme that achieves the computation-communication function in Theorem 2 in Section V, and the converse of Theorem 2 in Section VII.

A preliminary part of this result, in particular the achievability for the special case of s ⁣= ⁣1s\!=\!1, or the achievable scheme of Theorem 1 was presented in . We note that when s=1s=1, Theorem 2 provides the same result as in Theorem 1, i.e., L∗(r,1)=1r⋅(1−rK)L^{*}(r,1)=\frac{1}{r}\cdot(1-\frac{r}{K}), for r∈{1,…,K}r\in\{1,\ldots,K\}. \hfill□\hfill\square

For any fixed s∈{1,…,K}s\in\{1,\ldots,K\} (number of nodes that compute each Reduce function), as illustrated in Fig. 3, the communication load achieved in Theorem 2 outperforms the linear relationship between computation and communication, i.e., it is superlinear with respect to the computation load rr. \hfill□\hfill\square

Before we proceed to describe the general achievability scheme for the cascaded distributed computing framework (also the distributed computing framework as a special case of s=1s=1), we first illustrate the key ideas of the proposed Coded Distributed Computing scheme by presenting two examples in the next section, for the cases of s=1s=1 and s>1s>1 respectively.

IV Illustrative Examples: Coded Distributed Computing

In this section, we present two illustrative examples of the proposed achievable scheme for Theorem 1 and Theorem 2, which we call Coded Distributed Computing (CDC), for the cases of s=1s=1 (Theorem 1) and s>1s>1 (Theorem 2) respectively.

We consider a MapReduce-type problem in Fig. 4 for distributed computing of Q=3Q=3 output functions, represented by red/circle, green/square, and blue/triangle respectively, from N=6N=6 input files, using K=3K=3 computing nodes. Nodes 11, 22, and 33 are respectively responsible for final reduction of red/circle, green/square, and blue/triangle output functions. Let us first consider the case where no redundancy is imposed on the computations, i.e., each file is mapped once and computation load r=1r=1. As shown in Fig. 4(a), Node kk maps File 2k−12k-1 and File 2k2k for k=1,2,3k=1,2,3. In this case, each node maps 22 input files locally, computing all three intermediate values needed for the three output functions from each mapped file. In Fig. 4, we represent, for example, the intermediate value of the red/circle function in File nn using a red circle labelled by nn, for all n=1,…,6n=1,\ldots,6. Similar representations follow for the green/square and the blue/triangle functions. After the Map phase, each node obtains 22 out of 66 required intermediate values to reduce the output function it is responsible for (e.g., Node 1 knows the red circles in File 1 and File 2). Hence, each node needs 44 intermediate values from the other nodes, yielding a communication load of 4×33×6=23\frac{4\times 3}{3\times 6}=\frac{2}{3}.

Now, we demonstrate how the proposed CDC scheme trades the computation load to slash the communication load via in-network coding. As shown in Fig. 4(b), we double the computation load such that each file is now mapped on two nodes (r=2r=2). It is apparent that since more local computations are performed, each node now only requires 22 other intermediate values, and an uncoded shuffling scheme would achieve a communication load of 2×33×6=13\frac{2\times 3}{3\times 6}=\frac{1}{3}. However, we can do much better with coding. As shown in Fig. 4(b), instead of unicasting individual intermediate values, every node multicasts a bit-wise XOR, denoted by ⊕\oplus, of 22 locally computed intermediate values to the other two nodes, simultaneously satisfying their data demands. For example, knowing the blue/triangle in File 33, Node 22 can cancel it from the coded packet sent by Node 11, recovering the needed green/square in File 11. Therefore, this coding incurs a communication load of 33×6=16\frac{3}{3\times 6}=\frac{1}{6}, achieving a 2×2\times gain from the uncoded shuffling. □\square

From the above example, we see that for the case of s=1s=1, i.e., each of the QQ output functions is computed on one node and the computations of the Reduce functions are symmetrically distributed across nodes, the proposed CDC scheme only requires performing bit-wise XOR as the encoding and decoding operations. However, for the case of s>1s>1, as we will show in the following example, the proposed CDC scheme requires computing linear combinations of the intermediate values during the encoding process.

In this example, we consider a job of computing Q=6Q=6 output functions from N=6N=6 input files, using K=4K=4 nodes. We focus on the case where the computation load r=2r=2, and each Reduce function is computed by s=2s=2 nodes. In the Map phase, each file is mapped by r=2r=2 nodes. As shown in Fig. 5, the sets of the files mapped by the 44 nodes are M1={w1,w2,w3}\mathcal{M}_{1}=\{w_{1},w_{2},w_{3}\}, M2={w1,w4,w5}\mathcal{M}_{2}=\{w_{1},w_{4},w_{5}\}, M3={w2,w4,w6}\mathcal{M}_{3}=\{w_{2},w_{4},w_{6}\}, and M4={w3,w5,w6}\mathcal{M}_{4}=\{w_{3},w_{5},w_{6}\}. After the Map phase, Node kk, k∈{1,2,3,4}k\in\{1,2,3,4\}, knows the intermediate values of all Q=6Q=6 output functions in the files in Mk\mathcal{M}_{k}, i.e., {vq,n:q∈{1,…,6},wn∈Mk}\{v_{q,n}:q\in\{1,\ldots,6\},w_{n}\in\mathcal{M}_{k}\}. In the Reduce phase, we assign the computations of the Reduce functions in a symmetric manner such that every subset of s ⁣= ⁣2s\!=\!2 nodes compute a common Reduce function. More specifically as shown in Fig. 5, the sets of indices of the Reduce functions computed by the 44 nodes are W1 ⁣= ⁣{1,2,3}\mathcal{W}_{1}\!=\!\{1,2,3\}, W2 ⁣= ⁣{1,4,5}\mathcal{W}_{2}\!=\!\{1,4,5\}, W3 ⁣= ⁣{2,4,6}\mathcal{W}_{3}\!=\!\{2,4,6\}, and W4 ⁣= ⁣{3,5,6}\mathcal{W}_{4}\!=\!\{3,5,6\}. Therefore, for example, Node 1 still needs the intermediate values {vq,n ⁣: ⁣q∈{1,2,3},n∈{4,5,6}}\{v_{q,n}\!:\!q\in\{1,2,3\},n\in\{4,5,6\}\} through data shuffling to compute its assigned Reduce functions h1h_{1}, h2h_{2}, h3h_{3}.

The data shuffling process consists of two rounds of communication over the multicast network. In the first round, intermediate values are communicated within each subset of 33 nodes. In the second round, intermediate values are communicated within the set of all 44 nodes. In what follows, we describe these two rounds of communication respectively.

Round 1: Subsets of 33 nodes. We first consider the subset {1,2,3}\{1,2,3\}. During the data shuffling, each node whose index is in {1,2,3}\{1,2,3\} multicasts a bit-wise XOR of two locally computed intermediate values to the other two nodes:

Node 1 multicasts v1,2⊕v2,1v_{1,2}\oplus v_{2,1} to Node 22 and Node 33,

Node 2 multicasts v4,1⊕v1,4v_{4,1}\oplus v_{1,4} to Node 11 and Node 33,

Node 3 mulicasts v4,2⊕v2,4v_{4,2}\oplus v_{2,4} to Node 11 and Node 22,

Since Node 2 knows v2,1v_{2,1} and Node 3 knows v1,2v_{1,2} locally, they can respectively decode v1,2v_{1,2} and v2,1v_{2,1} from the coded message v1,2⊕v2,1v_{1,2}\oplus v_{2,1}.

We employ the similar coded shuffling scheme on the other 3 subsets of 3 nodes. After the first round of shuffling,

Node 1 recovers (v1,4,v1,5)(v_{1,4},v_{1,5}), (v2,4,v2,6)(v_{2,4},v_{2,6}) and (v3,5,v3,6)(v_{3,5},v_{3,6}),

Node 2 recovers (v1,2,v1,3)(v_{1,2},v_{1,3}), (v4,2,v4,6)(v_{4,2},v_{4,6}) and (v5,3,v5,6)(v_{5,3},v_{5,6}),

Node 3 recovers (v2,1,v2,3)(v_{2,1},v_{2,3}), (v4,1,v4,5)(v_{4,1},v_{4,5}) and (v6,3,v6,5)(v_{6,3},v_{6,5}),

Node 4 recovers (v3,1,v3,2)(v_{3,1},v_{3,2}), (v5,1,v5,4)(v_{5,1},v_{5,4}) and (v6,2,v6,4)(v_{6,2},v_{6,4}).

Similarly, as shown in Fig. 5, each of Node 22, Node 33, and Node 44 multicasts two linear combinations of three locally computed segments to the other three nodes, using the same coefficients α1\alpha_{1}, α2\alpha_{2}, and α3\alpha_{3}.

Having received the above two linear combinations, each of Node 22, Node 33, and Node 44 first subtracts out one segment available locally from the combinations, or more specifically, v6,1(1)v_{6,1}^{(1)} for Node 22, v5,2(1)v_{5,2}^{(1)} for Node 33, and v4,3(1)v_{4,3}^{(1)} for Node 44. After the subtraction, each of these three nodes recovers the required segments from the two linear combinations. More specifically, Node 2 recovers v4,3(1)v_{4,3}^{(1)} and v5,2(1)v_{5,2}^{(1)}, Node 3 recovers v4,3(1)v_{4,3}^{(1)} and v6,1(1)v_{6,1}^{(1)}, and Node 4 recovers v5,2(1)v_{5,2}^{(1)} and v6,1(1)v_{6,1}^{(1)}. It is not difficult to see that the above decoding process is guaranteed to be successful if α1\alpha_{1}, α2\alpha_{2}, and α3\alpha_{3} are all distinct from each other, which requires the field size 2T2≥32^{\frac{T}{2}}\geq 3 (e.g., T=4T=4). Following the similar procedure, each node recovers the required segments from the linear combinations multicast by the other three nodes. More specifically, after the second round of data shuffling,

Node 1 recovers v1,6v_{1,6}, v2,5v_{2,5} and v3,4v_{3,4},

Node 2 recovers v1,6v_{1,6}, v4,3v_{4,3} and v5,2v_{5,2},

Node 3 recovers v2,5v_{2,5}, v4,3v_{4,3} and v6,1v_{6,1},

Node 4 recovers v3,4v_{3,4}, v5,2v_{5,2} and v6,1v_{6,1}.

We finally note that in the second round of data shuffling, each linear combination multicast by a node is simultaneously useful for the rest of the three nodes. \hfill□\hfill\square

V General Achievable Scheme: Coded Distributed Computing

In this section, we formally prove the upper bounds in Theorem 1 and 2 by presenting and analyzing the Coded Distributed Computing (CDC) scheme. We focus on the more general case considered in Theorem 2 with s≥1s\geq 1, and the scheme for Theorem 1 simply follows by setting s=1s=1.

We first consider the integer-valued computation load r∈{1,…,K}r\in\{1,\ldots,K\}, and then generalize the CDC scheme for any 1≤r≤K1\leq r\leq K. When r=Kr=K, every node can map all the input files and compute all the output functions locally, thus no communication is needed and L∗(K,s)=0L^{*}(K,s)=0 for all s∈{1,…,K}s\in\{1,\ldots,K\}. In what follows, we focus on the case where r<Kr<K.

In the Map phase the Nˉ\bar{N} input files are evenly partitioned into (Kr){K\choose r} disjoint batches of size η1\eta_{1}, each corresponding to a subset T⊂{1,…,K}\mathcal{T}\subset\{1,\ldots,K\} of size rr, i.e.,

where BT\mathcal{B}_{\cal T} denotes the batch of η1\eta_{1} files corresponding to the subset T\mathcal{T}.

Given this partition, Node kk, k∈{1,…,K}k\in\{1,\ldots,K\}, computes the Map functions of the files in BT\mathcal{B}_{\cal T} if k∈Tk\in\mathcal{T}. Or equivalently, BT⊆Mk\mathcal{B}_{\cal T}\subseteq\mathcal{M}_{k} if k∈Tk\in\mathcal{T}. Since each node is in (K−1r−1){K-1\choose r-1} subsets of size rr, each node computes (K−1r−1)η1=rNˉK{K-1\choose r-1}\eta_{1}=\frac{r\bar{N}}{K} Map functions, i.e., ∣Mk∣=rNˉK|\mathcal{M}_{k}|=\frac{r\bar{N}}{K} for all k∈{1,…,K}k\in\{1,\ldots,K\}. After the Map phase, Node kk, k∈{1,…,K}k\in\{1,\ldots,K\}, knows the intermediate values of all QQ output functions in the files in Mk\mathcal{M}_{k}, i.e., {vq,n:q∈{1,…,Q},wn∈Mk}\{v_{q,n}:q\in\{1,\ldots,Q\},w_{n}\in\mathcal{M}_{k}\}.

V-B Coded Data Shuffling

where DP\mathcal{D}_{\cal P} denotes the indices of the batch of η2\eta_{2} Reduce functions corresponding to the subset P\mathcal{P}.

Given this partition, Node kk, k∈{1,…,K}k\in\{1,\ldots,K\}, computes the Reduce functions whose indices are in DP\mathcal{D}_{\cal P} if k∈Pk\in\mathcal{P}. Or equivalently, DP⊆Wk\mathcal{D}_{\cal P}\subseteq\mathcal{W}_{k} if k∈Pk\in\mathcal{P}. As a result, each node computes (K−1s−1)η2=sQK{K-1\choose s-1}\eta_{2}=\frac{sQ}{K} Reduce functions, i.e., ∣Wk∣=sQK|\mathcal{W}_{k}|=\frac{sQ}{K} for all k∈{1,…,K}k\in\{1,\ldots,K\}.

For a subset S{\cal S} of {1,…,K}\{1,\ldots,K\} and S1⊂S{\cal S}_{1}\subset\mathcal{S} with ∣S1∣=r|{\cal S}_{1}|=r, we denote the set of intermediate values needed by all nodes in S\S1{\cal S}\backslash\mathcal{S}_{1}, no node outside S\mathcal{S}, and known exclusively by nodes in S1\mathcal{S}_{1} as VS1S\S1\mathcal{V}_{\mathcal{S}_{1}}^{\mathcal{S}\backslash\mathcal{S}_{1}}. More formally:

We observe that the set VS1S\S1\mathcal{V}_{\mathcal{S}_{1}}^{{\cal S}\backslash\mathcal{S}_{1}} defined above contains intermediate values of (r∣S∣−s)η2{r\choose|{\cal S}|-s}\eta_{2} output functions. This is because that the output functions whose intermediate values are included in VS1S\S1\mathcal{V}_{\mathcal{S}_{1}}^{{\cal S}\backslash\mathcal{S}_{1}} should be computed exclusively by the nodes in S\S1{\cal S}\backslash\mathcal{S}_{1} and a subset of s−(∣S∣−r)s-(|{\cal S}|-r) nodes in S1{\cal S}_{1}. Therefore, VS1S\S1\mathcal{V}_{\mathcal{S}_{1}}^{{\cal S}\backslash\mathcal{S}_{1}} contains the intermediate values of a total of (rs−(∣S∣−r))η2=(r∣S∣−s)η2{r\choose s-(|{\cal S}|-r)}\eta_{2}={r\choose|{\cal S}|-s}\eta_{2} output functions. Since every subset of rr nodes map a unique batch of η1\eta_{1} files, VS1S\S1\mathcal{V}_{\mathcal{S}_{1}}^{{\cal S}\backslash\mathcal{S}_{1}} contains ∣VS1S\S1∣=(r∣S∣−s)η1η2|\mathcal{V}_{\mathcal{S}_{1}}^{{\cal S}\backslash\mathcal{S}_{1}}|={r\choose|{\cal S}|-s}\eta_{1}\eta_{2} intermediate values.

For each k∈Sk\in{\cal S}, there are a total of (∣S∣−1r−1){|{\cal S}|-1\choose r-1} subsets of S{\cal S} with size rr that contain the element kk. We index these subsets as S(k),S(k)…,S(k)[(∣S∣−1r−1)]{\cal S}_{(k)},{\cal S}_{(k)}\ldots,{\cal S}_{(k)}[{|{\cal S}|-1\choose r-1}]. Within a subset S(k)[i]{\cal S}_{(k)}[i], the segment associated with Node kk is US(k)[i],kS\S(k)[i]U_{\mathcal{S}_{(k)}[i],k}^{{\cal S}\backslash\mathcal{S}_{(k)}[i]}, for all i=1,…,(∣S∣−1r−1)i=1,\ldots,{|{\cal S}|-1\choose r-1}. We note that each segment US(k)[i],kS\S(k)[i]U_{\mathcal{S}_{(k)}[i],k}^{{\cal S}\backslash\mathcal{S}_{(k)}[i]}, i=1,…,(∣S∣−1r−1)i=1,\ldots,{|{\cal S}|-1\choose r-1}, is known by all nodes whose indices are in S(k)[i]\mathcal{S}_{(k)}[i], and needed by all nodes whose indices are in S\S(k)[i]{\cal S}\backslash\mathcal{S}_{(k)}[i].

We note that the above encoding process is the same at all nodes whose indices are in S{\cal S}, i.e., each of them multiplies the same matrix AS{\bf A}^{\cal S} in (16) with the segments associated with it.

Having generated the above message symbols, Node kk multicasts them to the other nodes whose indices are in S{\cal S}.

When s=1s=1, i.e., every output function is computed by one node, the above shuffling scheme only takes one round for all subsets S{\cal S} of size ∣S∣=r+1|{\cal S}|=r+1. Instead of multicasting linear combinations, every node in S{\cal S} can simply multicast the bit-wise XOR of its associated segments to the other rr nodes in S{\cal S}. \hfill□\hfill\square

V-B2 Decoding

For j∈Sj\in{\cal S} and j≠kj\neq k, there are a total of (∣S∣−2r−2){|{\cal S}|-2\choose r-2} subsets of S{\cal S} that have size rr and simultaneously contain jj and kk. Hence, among all n1n_{1} segments US(k),kS\S(k),US(k),kS\S(k),…,US(k)[n1],kS\S(k)[n1]U_{\mathcal{S}_{(k)},k}^{{\cal S}\backslash\mathcal{S}_{(k)}},U_{\mathcal{S}_{(k)},k}^{{\cal S}\backslash\mathcal{S}_{(k)}},\ldots,U_{\mathcal{S}_{(k)}[n_{1}],k}^{{\cal S}\backslash\mathcal{S}_{(k)}[n_{1}]} associated with Node kk, (∣S∣−2r−2){|{\cal S}|-2\choose r-2} of them are already known at Node jj, and the rest of n1−(∣S∣−2r−2)=(∣S∣−1r−1)−(∣S∣−2r−2)=(∣S∣−2r−1)=n2n_{1}-{|{\cal S}|-2\choose r-2}={|{\cal S}|-1\choose r-1}-{|{\cal S}|-2\choose r-2}={|{\cal S}|-2\choose r-1}=n_{2} segments are needed by Node jj. We denote the indices of the subsets that contain the element kk but not the element jj as bjk1,bjk2,…,bjkn2b_{jk}^{1},b_{jk}^{2},\ldots,b_{jk}^{n_{2}}, such that 1≤bjk1<bjk2<⋯<bjkn2≤n11\leq b_{jk}^{1}<b_{jk}^{2}<\cdots<b_{jk}^{n_{2}}\leq n_{1}, and j∉S(k)[bjki]j\notin{\cal S}_{(k)}[b_{jk}^{i}] for all i=1,2,…,n2i=1,2,\ldots,n_{2}.

After receiving the symbols XkS,XkS,…,XkS[n2]X_{k}^{\cal S},X_{k}^{\cal S},\ldots,X_{k}^{\cal S}[n_{2}] from Node kk, Node jj first removes the locally known segments from the linear combinations to generate n2n_{2} symbols YjkS,YjkS,…,YjkS[n2]Y_{jk}^{\cal S},Y_{jk}^{\cal S},\ldots,Y_{jk}^{\cal S}[n_{2}], such that

V-C Correctness of CDC

We demonstrate the correctness of the above shuffling scheme by showing that after the Shuffle phase, each node can decode all of the required intermediate values to compute its assigned Reduce functions. We use Node 1 as an example, and similar arguments apply to all other nodes. WLOG we assume that the Reduce function h1h_{1} is to be computed by Node 1. Node 1 will need a total of (K−1r)η1{K-1\choose r}\eta_{1} distinct intermediate values of h1h_{1} from other nodes (it already knows rNˉK=Nˉ−(K−1r)η1\frac{r\bar{N}}{K}=\bar{N}-{K-1\choose r}\eta_{1} intermediate values of h1h_{1} by mapping the files in M1{\cal M}_{1}). By the assignment of the Reduce functions, there exits a subset S2\mathcal{S}_{2} of size ss containing Node 1 such that all nodes in S2\mathcal{S}_{2} need to compute h1h_{1}. Then, during the data shuffling process within each subset S\mathcal{S} containing S2\mathcal{S}_{2} (note that by the definition of VS1S\S1{\cal V}_{{\cal S}_{1}}^{{\cal S}\backslash{\cal S}_{1}} in (13), the intermediate values of h1h_{1} will not be communicated to Node 1 if S2⊈S{\cal S}_{2}\nsubseteq{\cal S}, and this is because that some node outside S{\cal S} also wants to compute h1h_{1}), there are (s−1∣S∣−r−1){s-1\choose|{\cal S}|-r-1} subsets S1{\cal S}_{1} of S{\cal S} with size ∣S1∣=r|{\cal S}_{1}|=r such that 1∉S11\notin{\cal S}_{1} and S\S1⊆S2{\cal S}\backslash{\cal S}_{1}\subseteq{\cal S}_{2}, and thus Node 1 decodes (s−1∣S∣−r−1)η1{s-1\choose|\mathcal{S}|-r-1}\eta_{1} distinct intermediate values of h1h_{1}. Therefore, the total number of distinct intermediate values of h1h_{1} Node 1 decodes over the entire Shuffle phase is

which matches the required number of intermediate values for h1h_{1}. This is also true for all the other Reduce functions assigned to Node 1.

V-D Communication Load

In the above shuffling scheme, for each subset S⊆{1,…,K}\mathcal{S}\subseteq\{1,\ldots,K\} of size max⁡{r+1,s}≤∣S∣≤min⁡{r+s,K}\max\{r+1,s\}\leq|\mathcal{S}|\leq\min\{r+s,K\}, each Node k∈Sk\in{\cal S} communicates n2=(∣S∣−2r−1)n_{2}={|{\cal S}|-2\choose r-1} message symbols. Each of these symbols contains (r∣S∣−s)η1η2Tr{r\choose|{\cal S}|-s}\frac{\eta_{1}\eta_{2}T}{r} bits. Hence, all nodes whose indices are in S{\cal S} communicate a total of ∣S∣(∣S∣−2r−1)(r∣S∣−s)η1η2Tr|{\cal S}|{|{\cal S}|-2\choose r-1}{r\choose|{\cal S}|-s}\frac{\eta_{1}\eta_{2}T}{r} bits. The overall communication load achieved by the proposed CDC scheme is

V-E Non-Integer Valued Computation Load

For non-integer valued computation load r≥1r\geq 1, we generalize the CDC scheme as follows. We first expand the computation load r=αr1+(1−α)r2r=\alpha r_{1}+(1-\alpha)r_{2} as a convex combination of r1≜⌊r⌋r_{1}\triangleq\lfloor r\rfloor and r2≜⌈r⌉r_{2}\triangleq\lceil r\rceil, for some 0≤α≤10\leq\alpha\leq 1. Then we partition the set of Nˉ\bar{N} input files {w1,…,wNˉ}\{w_{1},\ldots,w_{\bar{N}}\} into two disjoint subsets I1\mathcal{I}_{1} and I2\mathcal{I}_{2} of sizes ∣I1∣=αNˉ|\mathcal{I}_{1}|=\alpha\bar{N} and ∣I2∣=(1−α)Nˉ|\mathcal{I}_{2}|=(1-\alpha)\bar{N}. We next apply the CDC scheme described above respectively to the files in I1\mathcal{I}_{1} with a computation load r1r_{1} and the files in I2\mathcal{I}_{2} with a computation load r2r_{2}, to compute each of the QQ output functions at the same set of ss nodes. This results in a communication load of

where Lcoded(r,s)L_{\textup{coded}}(r,s) is the communication load achieved by CDC in (20) for integer-valued r,s∈{1,…,K}r,s\in\{1,\ldots,K\}.

Using this generalized CDC scheme, for any two integer-valued computation loads r1r_{1} and r2r_{2}, the points on the line segment connecting (r1,Lcoded(r1,s))(r_{1},L_{\textup{coded}}(r_{1},s)) and (r2,Lcoded(r2,s))(r_{2},L_{\textup{coded}}(r_{2},s)) are achievable. Therefore, for general 1≤r≤K1\leq r\leq K, the lower convex envelop of the achievable points {(r,Lcoded(r,s)):r∈{1,…,K}}\{(r,L_{\textup{coded}}(r,s)):r\in\{1,\ldots,K\}\} is achievable. This proves the upper bound on the computation-communication function in Theorem 2 (also the achievability part of Theorem 1 by setting s=1s=1).

The ideas of efficiently creating and exploiting coded multicasting opportunities have been introduced in caching problems . In this section, we illustrated how coding opportunities can be utilized in distributed computing to slash the load of communicating intermediate values, by designing a particular assignment of extra computations across distributed computing nodes. We note that the calculated intermediate values in the Map phase mimics the locally stored cache contents in caching problems, providing the “side information” to enable coding in the following Shuffle phase (or content delivery).

For the case of s=1s=1 where no two nodes are interested in computing a common Reduce function, the coded data shuffling of CDC is similar to a coded transmission strategy in wireless D2D networks proposed in , where the side information enabling coded multicasting are pre-fetched in a specific repetitive manner in the caches of wireless nodes (in CDC such information is obtained by computing the Map functions locally). When ss is larger than 11, i.e., every Reduce function needs to be computed at multiple nodes, our CDC scheme creates novel coding opportunities that exploit both the redundancy of the Map computations and the commonality of the data requests for Reduce functions across nodes, further reducing the communication load. \hfill□\hfill\square

Generally speaking, we can view the Shuffle phase of the considered distributed computing framework as an instance of the index coding problem , in which a central server aims to design a broadcast message (code) with minimum length to simultaneously satisfy the requests of all the clients, given the clients’ side information stored in their local caches. Note that while a randomized linear network coding approach (see e.g., ) is sufficient to implement any multicast communication where messages are intended by all receivers, it is generally sub-optimal for index coding problems where every client requests different messages. Although the index coding problem is still open in general, for the considered distributed computing scenario where we are given the flexibility of designing Map computation (thus the flexibility of designing side information), we prove in the next two sections tight lower bounds on the minimum communication loads for the cases s=1s=1 and s>1s>1 respectively, demonstrating the optimality of the proposed CDC scheme. \hfill□\hfill\square

VI Converse of Theorem 1

In this section, we prove the lower bound on L∗(r)L^{*}(r) in Theorem 1.

For k∈{1,…,K}k\in\{1,\ldots,K\}, we denote the set of indices of the files mapped by Node kk as Mk\mathcal{M}_{k}, and the set of indices of the Reduce functions computed by Node kk as Wk\mathcal{W}_{k}. As the first step, we consider the communication load for a given file assignment M≜(M1,M2…,MK)\mathcal{M}\triangleq({\cal M}_{1},{\cal M}_{2}\ldots,{\cal M}_{K}) in the Map phase. We denote the minimum communication load under the file assignment M\mathcal{M} by LM∗L^{*}_{\mathcal{M}}.

We denote the number of files that are mapped at jj nodes under a file assignment M{\cal M}, as aMja^{j}_{{\cal M}}, for all j∈{1,…,K}j\in\{1,\ldots,K\}:

For example, for the particular file assignment in Fig. 6, i.e., M=({1,3,5,6},{4,5,6},{2,3,4,6}){\cal M}=(\{1,3,5,6\},\{4,5,6\},\{2,3,4,6\}), aM1=2a^{1}_{\cal M}=2 since File 1 and File 2 are mapped on a single node (i.e., Node 1 and Node 3 respectively). Similarly, we have aM2=3a^{2}_{\cal M}=3 (Files 3, 4, and 5), and aM3=1a^{3}_{\cal M}=1 (File 6).

For a particular file assignment M{\cal M}, we present a lower bound on LM∗L^{*}_{\mathcal{M}} in the following lemma.

LM∗≥∑j=1KaMjN⋅K−jKjL^{*}_{\mathcal{M}}\geq\sum\limits_{j=1}^{K}\frac{a^{j}_{\cal M}}{N}\cdot\frac{K-j}{Kj}.

Next, we first demonstrate the converse of Theorem 1 using Lemma 1, and then give the proof of Lemma 1.

Converse Proof of Theorem 1. It is clear that the minimum communication load L∗(r)L^{*}(r) is lower bounded by the minimum value of LM∗L^{*}_{\mathcal{M}} over all possible file assignments which admit a computation load of rr:

For every file assignment M{\cal M} such that ∣M1∣+⋯+∣MK∣=rN|{\cal M}_{1}|+\cdots+|{\cal M}_{K}|=rN, {aMj}j=1K\{a^{j}_{\cal M}\}_{j=1}^{K} satisfy

Then since the function K−jKj\frac{K-j}{Kj} in (24) is convex in jj, and by (26) ∑j=1KaMjN=1\sum\limits_{j=1}^{K}\frac{a^{j}_{\cal M}}{N}=1, (24) becomes

where (a) is due to the requirement imposed by the computation load in (27).

Then by the convexity of the function K−jKj\frac{K-j}{Kj} in jj, we have for integer-valued j=1,…,Kj=1,\ldots,K,

where (b) is due to the constraints on {aMj}j=1K\{a^{j}_{\cal M}\}_{j=1}^{K} in (26) and (27).

Therefore, L∗(r)L^{*}(r) is lower bounded by the lower convex envelop of the points {(r,K−rKr):r∈{1,...,K}}\{(r,\frac{K-r}{Kr}):r\in\{1,...,K\}\}. This completes the proof of the converse part of Theorem 1. \hfill■\hfill\blacksquare

Although the model proposed in this paper only allows each node sending messages independently, we can show that even if the data shuffling process can be carried out in multiple rounds and dependency between messages are allowed, the lower bound on L∗(r)L^{*}(r) remains the same. \hfill□\hfill\square

We devote the rest of this section to the proof of Lemma 1. To prove Lemma 1, we develop a lower bound on the number of bits communicated by any subset of nodes, by induction on the size of the subset. In particular, for a subset of computing nodes, we first characterize a lower bound on the minimum number of bits required by a particular node in the subset, which is given by a cut-set bound separating this node and all the other nodes in the subset. Then, we combine this bound with the lower bound on the number of bits communicated by the rest of the nodes in the subset, which is given by the inductive argument.

Since each message XkX_{k} is generated as a function of the intermediate values that are computed at Node kk, the following equation holds for all k∈{1,...,K}k\in\{1,...,K\}.

where we use “:” to denote the set of all possible indices.

The validity of the shuffling scheme requires that for all k∈{1,...,K}k\in\{1,...,K\}, the following equation holds :

For a subset S⊆{1,...,K}\mathcal{S}\subseteq\{1,...,K\}, we define

which contains all the intermediate values required by the nodes in S{\cal S} and all the intermediate values known locally by the nodes in S{\cal S} after the Map phase.

For any subset S⊆{1,…,K}{\cal S}\subseteq\{1,\ldots,K\} and a file assignment M{\cal M}, we denote the number of files that are exclusively mapped by jj nodes in S\mathcal{S} as aMj,Sa^{j,\mathcal{S}}_{\cal M}:

and the message symbols communicated by the nodes whose indices are in S{\cal S} as

For any subset S⊆{1,...,K}\mathcal{S}\subseteq\{1,...,K\}, we have

where Sc≜{1,…,K}\S{\cal S}^{c}\triangleq\{1,\ldots,K\}\backslash{\cal S} denotes the complement of S{\cal S}. \hfill□\hfill\square

a. If S={k}\mathcal{S}=\{k\} for any k∈{1,…,K}k\in\{1,\ldots,K\}, obviously

b. Suppose the statement is true for all subsets of size S0S_{0}.

For any S⊆{1,...,K}\mathcal{S}\subseteq\{1,...,K\} of size ∣S∣=S0+1|\mathcal{S}|=S_{0}+1 and any k∈Sk\in\mathcal{S}, we have

For each k∈Sk\in{\cal S}, we have the following subset version of (36) and (37).

The first term on the RHS of (52) can be lower bounded as follows.

where (a) is due to the independence of intermediate values and the fact that Wk∩WSc=∅{\cal W}_{k}\cap{\cal W}_{{\cal S}^{c}}=\varnothing (different nodes calculate different output functions), (b) and (c) are due to the independence of intermediate values, and (d) is due to the independence of the intermediate values and the fact that ∣Wk∣=QK|{\cal W}_{k}|=\frac{Q}{K}.

The second term on the RHS of (52) can be lower bounded by the induction assumption:

Thus by (48), (52), (57) and (59), we have

By the definition of aMj,Sa^{j,\mathcal{S}}_{\cal M}, we have the following equations.

c. Thus for all subsets S⊆{1,...,K}\mathcal{S}\subseteq\{1,...,K\}, the following equation holds:

Then by Claim 1, let S={1,...,K}\mathcal{S}=\{1,...,K\} be the set of all KK nodes,

This completes the proof of Lemma 1. \hfill■\hfill\blacksquare

VII Converse of Theorem 2

In this section, we prove the lower bound on L∗(r,s)L^{*}(r,s) in Theorem 2, which generalizes the converse result of Theorem 1 for the case s>1s>1. Since the lower bound on L∗(r,1)L^{*}(r,1) in Theorem 2 exactly matches the lower bound on L∗(r)L^{*}(r) in Theorem 1, we focus on the case s>1s>1 (i.e., each Reduce function is calculated by 2 or more nodes) throughout this section.

We denote the minimum communication load under a particular file assignment M\mathcal{M} as LM∗(s)L^{*}_{\mathcal{M}}(s), and we present a lower bound on LM∗(s)L^{*}_{\mathcal{M}}(s) in the following lemma.

Converse Proof of Theorem 2. The minimum communication load L∗(r,s)L^{*}(r,s) is lower bounded by the minimum value of LM∗(s)L^{*}_{\mathcal{M}}(s) over all possible file assignments having a computation load of rr:

For every file assignment M{\cal M} such that ∣M1∣+⋯+∣MK∣=rN|{\cal M}_{1}|+\cdots+|{\cal M}_{K}|=rN, {aMj}j=1K\{a^{j}_{\cal M}\}_{j=1}^{K} satisfy the same conditions as the case of s=1s=1 in (25), (26) and (27).

Then by the convexity of the function Lcoded(j,s)L_{\textup{coded}}(j,s) in jj, we have for integer-valued j=1,…,Kj=1,\ldots,K,

Next, we first apply Lemma 2 to (73), then by (76), we have

where (a) is due to the constraints on {aMj}j=1K\{a^{j}_{\cal M}\}_{j=1}^{K} in (26) and (27).

Therefore, L∗(r,s)L^{*}(r,s) is lower bounded by the lower convex envelop of the points {(r,Lcoded(r,s)):r∈{1,...,K}}\{(r,L_{\textup{coded}}(r,s)):r\in\{1,...,K\}\}. This completes the proof of the converse part of Theorem 2. \hfill■\hfill\blacksquare

The proof of lemma 2 follows the same steps of the proof of Lemma 1, where a lower bound on the number of bits communicated by any subset of nodes, for the case of s>1s>1, is established by induction.

Proof of Lemma 2. We first prove the following claim.

For any subset S⊆{1,...,K}\mathcal{S}\subseteq\{1,...,K\}, we have

where aMj,Sa^{j,{\cal S}}_{\cal M} is defined in (39). \hfill□\hfill\square

a. If S={k}\mathcal{S}=\{k\} for any k∈{1,…,K}k\in\{1,\ldots,K\}, obviously

b. Suppose the statement is true for all subsets of size S0S_{0}.

For any S⊆{1,...,K}\mathcal{S}\subseteq\{1,...,K\} of size ∣S∣=S0+1|\mathcal{S}|=S_{0}+1, and all k∈Sk\in\mathcal{S}, we have as derived in (61):

where YSc=(VWSc,:,V:,MSc)Y_{\mathcal{S}^{c}}=(V_{\mathcal{W}_{\mathcal{S}^{c}},:},V_{:,\mathcal{M}_{\mathcal{S}^{c}}}).

The first term on the RHS of (82) is lower bounded by the induction assumption:

The second term on the RHS of (82) can be calculated based on the independence of intermediate values:

where (a) and (b) are due to the independence of the intermediate values, and (c) is due to the uniform distribution of the output functions such that each node in S{\cal S} calculates Q(Ks)⋅(∣S∣−1s−1)\frac{Q}{{K\choose s}}\cdot{|{\cal S}|-1\choose s-1} output functions computed exclusively by ss nodes in S{\cal S}.

For each j∈{1,…,S0+1}j\in\{1,\ldots,S_{0}+1\} in (93), we have

Since (99) holds for all subsets S{\cal S} of size ∣S∣=S0+1|{\cal S}|=S_{0}+1, we have proven Claim 2.

Then by Claim 2, let S={1,...,K}\mathcal{S}=\{1,...,K\} be the set of all KK nodes,

This completes the proof of Lemma 2. \hfill■\hfill\blacksquare

VIII Implementation and Empirical Evaluation of Coded Distributed Computing

In this section, we demonstrate the impact of the proposed Coded Distributed Computing (CDC) scheme on balancing the time spent on task execution and the time spent on data movement, in order to speed up practical distributed computing applications. In particular, let us consider a MapReduce-type application for which the total execution time is roughly composed of the time spent executing the Map tasks, denoted by TmapT_{\textup{map}}, the time spent shuffling intermediate values, denoted by TshuffleT_{\textup{shuffle}}, and the time spent executing the Reduce tasks, denoted by TreduceT_{\textup{reduce}}, i.e.,

Using CDC, we can leverage r×r\times more computations in the Map phase, in order to reduce the communication load by the same multiplicative factor. Hence, ignoring the coding overheads, CDC promises an approximate total execution time of

To minimize the above execution time, one would choose r∗=⌊TshuffleTmap⌋ or ⌈TshuffleTmap⌉r^{*}=\left\lfloor\sqrt{\tfrac{T_{\textup{shuffle}}}{T_{\textup{map}}}}\right\rfloor\textup{ or }\left\lceil\sqrt{\tfrac{T_{\textup{shuffle}}}{T_{\textup{map}}}}\right\rceil, resulting in the minimum execution time of

For example, in an application that TshuffleT_{\textup{shuffle}} is 10×10\times - 100×100\times larger than Tmap+TreduceT_{\textup{map}}+T_{\textup{reduce}}, by comparing from (101) and (103), we note that CDC can reduce the execution time by approximately 1.5×1.5\times - 5×5\times.

In the rest of this section, we empirically demonstrate the performance gain of applying CDC to TeraSort , which is a commonly used Hadoop benchmark for distributed sorting terabytes of data . In particular, we first incorporate the coding ideas in CDC into TeraSort to develop a novel coded distributed sorting algorithm, named CodedTeraSort, which imposes structured redundancy in the input data, in order to enable in-network coding opportunities that overcome the data shuffling bottleneck of TeraSort. Then, we evaluate the performance of CodedTeraSort on Amazon EC2 clusters, and observe a 1.97×1.97\times - 3.39×3.39\times speedup, compared with TeraSort, for typical settings of interest.

TeraSort is a conventional algorithm for distributed sorting of a large amount of data. The input data that is to be sorted is in the format of key-value (KV) pairs, meaning that each input KV pair consists of a key and a value. For example, the domain of the keys can be 10-byte integers, and the domain of the values can be arbitrary strings. TeraSort sorts the input data according to their keys, e.g., sorting integers.

Let us consider implementing TeraSort over KK distributed computing nodes, which consists of 5 stages: File Placement, Key Domain Partitioning, Map Phase, Shuffle Phase, and Reduce Phase. In File Placement, all input KV pairs are split into KK disjoint files, and each file is placed on one of the KK nodes. In Key Domain Partitioning, the domain of the keys is split into KK partitions, and each node will be responsible for sorting the KV pairs whose keys fall into one of the partitions. In Map Phase, each node hashes each KV pair in its locally stored file into one of the KK partitions, according to its key. In Shuffle Phase, the KV pairs in the same partition are transferred to the node that is responsible for sorting that partition. In Reduce Stage, each node locally sorts KV pairs belonging to its assigned partition. We illustrate the TeraSort algorithm using a simple example shown in Fig. 7.

VIII-A2 Performance Evaluation

To understand the performance of TeraSort, we performed an experiment on Amazon EC2 to sort 12GB of data by running TeraSort on 16 instances.We note that EC2 uses virtual machines, and each instance may not be hosted by a dedicated physical machine. The breakdown of the total execution time is shown in Table I.

We observe from Table I that for a conventional TeraSort execution, 98.4% of the total execution time was spent in data shuffling, which is 508.5×508.5\times of the time spent in the Map phase. Given the fact that data shuffling dominates the job execution time, the principle of optimally trading computation for communication of the proposed CDC scheme can be applied to significantly improve the performance of TeraSort. For example, when executing the same sorting job using a coded version of TeraSort with a computation load of r=10r=10, according to (102), we could theoretically save the total execution time by approximately 8×8\times. This motivates us to develop a novel coded distributed sorting algorithm, named CodedTeraSort, which is briefly described in the next sub-section.

VIII-B Coded TeraSort

We develop the CodedTeraSort algorithm by applying the proposed CDC scheme for the case of s=1s=1 (see Example 1 in Section IV for an illustration) to the above described TeraSort algorithm. CodedTeraSort exploits redundant computations on the input files in the Map phase, creating in-network coding opportunities to significantly slash the load of data shuffling. In particular, the execution of CodedTeraSort consists of following 66 stages of operations. Here we give high-lever descriptions of these operations, and we refer the interested readers to for more detailed descriptions.

Structured Redundant File Placement. The entire input KV pairs are split into many small files, each of which is repeatedly placed on 1≤r≤K1\leq r\leq K nodes (i.e., a computation load of rr), according to the particular pattern specified by the CDC scheme.

Map. Each node applies the hashing operation as in TeraSort on each of its assigned files.

Encoding to Create Coded Packets. Each node generates coded multicast packets from local results computed in Map phase, according to the encoding process of the CDC scheme.

Multicast Shuffling. Each node multicasts each of its generated coded packet to a specific set of rr other nodes.

Decoding. Each node locally decodes the required KV pairs from the received coded packets.

Reduce. Each node locally sorts the KV pairs within its assigned partition as in the Reduce phase of TeraSort.

VIII-C Empirical Evaluations

We imperially demonstrate the performance gain of CodedTeraSort through experiments on Amazon EC2 clusters. In this sub-section, we first present some choices we have made for the implementation. Then, we discuss the experiment results.

We first describe the following common implementation choices that we have made for both TeraSort and CodedTeraSort algorithms.

Data Format: All input KV pairs are generated from TeraGen in the standard Hadoop package. Each input KV pair consists of a 1010-byte key and a 9090-byte value. A key is a 1010-byte unsigned integer, and the value is an arbitrary string of 9090 bytes. The KV pairs are sorted based on their keys, using the standard integer ordering.

Library: We implement both TeraSort and CodedTeraSort algorithms in C++, and use Open MPI library for communications between EC2 instances.

In the TeraSort implementation, each node sequentially steps through Map, Pack, Shuffle, Unpack, and Reduce stages. The Pack stage serializes each intermediate value to a continuous memory array to ensure that a single TCP flow is created for each intermediate value (which may contain multiple KV pairs) when MPI_Send is calledCreating a TCP flow per KV pair leads to inefficiency from overhead and convergence issue.. The Unpack stage deserializes the received data to a list of KV pairs. In the Shuffle stage, intermediate values are unicast serially, meaning that there is only one sender node and one receiver node at any time instance. Specifically, as illustrated in Fig. 8(a), Node 11 starts to unicast to Nodes 2, 3, and 4 back-to-back. After Node 11 finishes, Node 22 unicasts back-to-back to Nodes 1, 3, and 4. This continues until Node 4 finishes.

In the CodedTeraSort implementation, each node sequentially steps through CodeGen, Map, Encode, Multicast Shuffling, Decode, and Reduce stages. In the CodeGen (or code generation) stage, firstly, each node generates all file indices, as subsets of rr nodes. Then each node uses MPI_Comm_split to initialize (Kr+1)\binom{K}{r+1} multicast groups each containing r+1r+1 nodes on Open MPI, such that multicast communications will be performed within each of these groups. The serialization and deserialization are implemented respectively in the Encode and the Decode stages. In Multicast Shuffling, MPI_Bcast is called to multicast a coded packet in a serial manner, so only one node multicasts one of its encoded packets at any time instance. Specifically, as illustrated in Fig. 8(b), Node 1 multicasts to the other 2 nodes in each multicast group Node 1 is in. For example, Node 1 first multicasts to Node 2 and 3 in the multicast group {1,2,3}\{1,2,3\}. After Node 11 finishes, Node 22 starts multicasting in the same manner. This process continues until Node 44 finishes.

VIII-C2 Experiment Results

We evaluate the run-time performance of TeraSort and CodedTeraSort, for different combinations of the number of workers KK and the computation load 1≤r≤K1\leq r\leq K. All experiments are repeated 55 times, and the average values are recorded.

In Table II and Table III, we list the breakdowns of the average execution times to sort 12 GB of input data using K=16K=16 workers and K=20K=20 workers respectively. Here we limit the incoming and outgoing traffic rates of each instance to 100100 Mbps. This is to alleviate the effects of the bursty behaviors of the transmission rates in the beginning of some TCP sessions, given the particular size of the data to be sorted. We observe an overall 1.97×1.97\times - 3.39×3.39\times speedup of CodedTeraSort as compared with TeraSort. From the experiment results we make the following observations:

For CodedTeraSort, the time spent in the CodeGen stage is proportional to (Kr+1)\binom{K}{r+1}, which is the number of multicast groups.

The Map time of CodedTeraSort is approximately rr times higher than that of TeraSort. This is because that each node hashes rr times more KV pairs than that in TeraSort. Specifically, the ratios of the CodedTeraSort’s Map time to the TeraSort’s Map time from Table II are 6.03/1.86≈3.26.03/1.86\approx 3.2 and 10.84/1.86≈5.810.84/1.86\approx 5.8, and from Table III are 4.68/1.47≈3.24.68/1.47\approx 3.2 and 8.59/1.47≈5.88.59/1.47\approx 5.8.

While CodedTeraSort theoretically promises a factor of more than r×r\times reduction in shuffling time, the actual gains observed in the experiments are slightly less than rr. For example, for the experiment with K=16K=16 nodes and r=3r=3, as shown in Table II, the speedup of the Shuffle stage is 945.72/412.22≈2.3<3945.72/412.22\approx 2.3<3. This phenomenon is caused by the following two factors. 1) Open MPI’s multicast API (MPI_Bcast) has an inherent overhead per a multicast group, for instance, a multicast tree is constructed before multicasting to a set of nodes. 2) Using the MPI_Bcast API, the time of multicasting a packet to rr nodes is higher than that of unicasting the same packet to a single node. In fact, as measured in , the multicasting time increases logarithmically with rr.

Further, we observe the following trends from both tables:

The impact of computation load rr: As rr increases, the shuffling time reduces by approximately rr times. However, the Map execution time increases linearly with rr, and more importantly the CodeGen time increases exponentially with rr as (Kr+1)\binom{K}{r+1}. Hence, for small values of rr (r<6r<6) we observe overall reduction in execution time, and the speedup increases. However, as we further increase rr, the CodeGen time will dominate the execution time, and the speedup decreases. Hence, in our evaluations, we have limited rr to be at most 55.The redundancy parameter rr is also limited by the total storage available at the nodes. Since for a choice of redundancy parameter rr, each piece of input KV pairs should be stored at rr nodes, we can not increase rr beyond total available storage at the worker nodesinput size.\frac{\text{total available storage at the worker nodes}}{\text{input size}}.

The impact of worker number KK: As KK increases, the speedup decreases. This is due to the following two reasons. 1) The number of multicast groups, i.e., (Kr+1)\binom{K}{r+1}, grows exponentially with KK, resulting in a longer execution time of the CodeGen process. 2) When more nodes participate in the computation, for a fixed rr, less amount of KV pairs are hashed at each node locally in the Map phase, resulting in less locally available intermediate values and a higher communication load. Hence, given more worker nodes, one would preferably use larger computation load to achieve a better run-time performance.

IX Concluding Remarks and Future Directions

We introduced a scalable distributed computing framework motivated by MapReduce, which is suited for arbitrary types of output functions. We formulated and exactly characterized an information-theoretic tradeoff between computation load and communication load within this framework. In particular, we proposed Coded Distributed Computing (CDC), a coded scheme that reduces the communication load by a factor that can grow with the network size, illustrating the role of coding in speeding up distributed computing jobs. We also proved a tight information-theoretic lower bound on the minimum communication load, using any data shuffling scheme, which exactly matches the communication load achieved by CDC. This result reveals a fundamental relationship between computation and communication in distributed computing–the two are inversely proportional to each other. Moreover, we applied the proposed CDC scheme to the conventional TeraSort algorithm to develop a novel distributed sorting algorithm, named CodedTeraSort, and empirically demonstrated the performance gain of CodedTeraSort through extensive experiments on Amazon EC2 clusters.

Finally, we discuss some follow-up research directions of this work.

Heterogeneous Networks with Asymmetric Tasks. It is common to have computing nodes with heterogeneous storage, processing and communication capacities within computer clusters (e.g., Amazon EC2 clusters composed of heterogeneous computing instances). In addition, processing different parts of the dataset can generate intermediate results with different sizes (e.g., performing data analytics on highly-clustered graphs). For computing over heterogeneous nodes, one solution is to break the more powerful nodes into multiple smaller virtual nodes that have homogeneous capability, and then apply the proposed CDC scheme for the homogeneous setting. When intermediate results have different sizes, the proposed coding scheme still applies, but the coding operations are not symmetric as in the case of homogeneous intermediate results (e.g., one may now need to compute the XOR of two data segments with different sizes). Alternatively, we can employ a low-complexity greedy approach, in which we assign the Map tasks to maximize the number of multicasting opportunities that simultaneously deliver useful information to the largest possible number of nodes. Some preliminary studies along this direction have been conducted to obtain the solutions for some special cases (see, e.g., ). Nevertheless, systematically characterizing the optimal resource allocation strategies and coding schemes for general heterogeneous networks with asymmetric tasks remains an interesting open problem.

Straggling/Failing Computing Nodes. Other than the communication bottleneck, the effect of straggling servers also severely degrades the run-time performance of distributed computing applications (see e.g., ). Recently in , Maximum-Distance-Separable (MDS) codes were utilized to encode linear computation tasks, providing robustness to a certain number of stragglers. Following the results in , coded computing strategies have been proposed to efficiently deal with the stragglers for various computation tasks and network settings (see, e.g., ). In , we have superimposed the proposed CDC scheme on top of the MDS codes, developing a unified coding framework for distributed computing with straggling servers. This framework achieves a flexible tradeoff between computation latency in the Map phase and communication load in the Shuffle phase, which has the CDC scheme (or minimum bandwidth code) and the MDS code (or minimum latency code) as the two end points. Nevertheless, designing resource allocation strategies and coding techniques to optimize the run-time performance over distributed computing clusters with stragglers is a challenging open problem.

Multi-Stage Computation Tasks. Unlike simple computation tasks like Grep, Join and Sort, many distributed computing applications contain multiples stages of MapReduce computations. Examples of these applications include machine learning algorithms , SQL queries for databases , and scientific analytics . One can express the computation logic of a multi-stage application as a directed acyclic graph (DAG) , in which each vertex represents a logical step of data transformation, and each edge represents the dataflow across processing vertices. In order to speed up multi-stage computation tasks using codes, while one straightforward approach is to apply the proposed CDC scheme for the cascaded distributed computing framework (see Theorem 2) to compute each stage locally, we expect to achieve a higher reduction in bandwidth consumption and response time by globally designing codes for the entire task graph and accounting for interactions between consecutive stages. A preliminary exploration along this direction was recently presented in .

Multi-Layer Networks and Structured Topology. So far we have only considered a single-layer topology of the distributed computing nodes, in which each node can multicast to an arbitrary number of other nodes at the same cost as unicasting to a single node. However, in practical data center networks, nodes can be connected through multiple switches at different layers with different capacities, forming a hierarchical multi-root tree topology (e.g., fat-tree topology ). In this case, we need to generalize our communication model to include more structured topologies, and develop coded shuffling strategies that account for (1) path lengths of shuffled data (2) congestion at links higher up in the topology; and (3) different link capacities and multicast-costs at different layers of network topology. We have made preliminary progress in for a star topology (motivated by wireless edge computing), where nodes are connected via only one access point (or switch layer).

Joint Storage and Computation Optimization. We have so far assumed that we can design the placement of the input files to create coding opportunities during the computation process. However, in practical file storage systems, data blocks are often stored without prior knowledge about the computations that will be performed on them, and moving the data across the nodes before the computation is often too costly. In this case, even without the capability of designing the data placement as exactly specified by the CDC scheme, one can still take advantage of the inherent data redundancy (e.g., GFS and HDFS by default place replicas of each data block on 3 distributed nodes) to create coded multicast opportunities, significantly reducing the communication load.

We plot in Fig. 9 the average communication load achieved by a coded shuffling scheme similar to the one presented in Section V-B (with the modification that each node zero-pads its associated data segments to the length of the longest one before coding), when each input file is placed and mapped at rr out of KK nodes chosen uniformly at random, and compare it with the communication load achieved by CDC where the input files are placed based on the Map phase design in Section V-A. As demonstrated in Fig. 9, without requiring the files to be placed as exactly described by the CDC scheme, one can still exploit the data redundancy to achieve a communication load that is superlinear with respect to the computation load. Therefore, the coded data shuffling scheme of CDC can effectively reduce the communication loads of computation jobs on general data storage systems. This behavior that a random data placement achieves close-to-optimum performance has also been reported in for a decentralized wireless distributed computing platform, and in for a decentralized caching system.

Coded Edge/Fog Computing. In the emerging mobile Edge/Fog computing paradigm (see, e.g., ), abundant computation resources scattered across the network edge (e.g., smartphones, tablets and smart cars) are harvested to perform data-intensive computations collaboratively. In this scenario, coding opportunities are widely available by injecting redundant storage and computations into the edge network. We envision codes to play a transformational role in Edge/Fog computing for leveraging such redundancy to substantially reduce the bandwidth consumption and the latency of computing. For an edge computing scenario where the mobile users upload the tasks to the edge nodes, and retrieve the computed results from the edge nodes, we have designed coded computing architectures in , in which coded computations that are aware of the underlying physical-layer communication are performed at the edge nodes, achieving the minimum load of computation and the maximum spectral efficiency simultaneously. In , we have formulated a wireless distributed computing framework, in which a cluster of mobile users collaborate via an access point to simultaneously meet their computational needs. For this wireless computing platform, we exploited the coding techniques of CDC to achieve a scalable design such that the platform can accommodate an unlimited number of mobile users with a constant amount of bandwidth consumption. Also in a recent magazine paper , we have demonstrated the opportunities of utilizing coding to improve the performance of Edge/Fog computing applications (e.g., navigation services and recommendation systems).

References