Coded MapReduce

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

I Introduction

MapReduce is a programming model that enables the distributed processing of large-scale datasets on a cluster of commodity servers. The desirable features of MapRuduce, such as scalability, simplicity and fault-tolerance , have made this framework popular to perform data-intensive tasks in text/graph processing, machine learning and bioinformatics. Simply speaking, in MapReduce framework, each of the input data blocks, stored across the distributed servers (e.g. HDFS in Hadoop ), is assigned to a worker node to process. The processing maps the input blocks to some intermediate (key,value) pairs. In the next step, referred as data shuffling, the intermediate (key,value) pairs are transferred via inter-server communication links to a set of processors to be reduced to the final results.

The data shuffling phase is one of the key bottlenecks for improving the runtime performance of MapReduce. In fact, as observed in , using a Hadoop cluster, 33% of the job execution time is spent on data shuffling. As a result, many optimization methods including combining intermediate (key,value) pairs before shuffling , optimal flow scheduling across network paths and employment of distributed cache memories , have been proposed to accelerate the data shuffling phase of MapReduce programs.

In this paper, we introduce “Coded MapReduce”, a new framework that enables and exploits a particular form of coding to significantly reduce the communication load of the shuffling phase. Compared to the conventional MapReduce approach, we demonstrate that Coded MapReduce can substantially cut down the inter-server communication load by a multiplicative factor that grows linearly with the number of servers in the system. We also show that Coded MapReduce achieves the minimum shuffling load within a constant multiplicative factor regardless of the system parameters.

Coded MapReduce exploits the repetitive mappings of the same data block at different servers, in order to enable coding. More specifically, for a server cluster with KK servers, interconnected through a multicast LAN network, and repetitive mapping of each data block on rr fraction of the servers, we propose a particular strategy to assign the Map tasks, such that in the shuffling phase the data demand of each group of approximately rK>1rK>1 servers can be satisfied by a single coded transmission, which we call a coded multicast opportunity. This would provide a multiplicative coding gain of rKrK, which scales linearly with the number of servers. Furthermore, even if repetition of data blocks is very small (for example, if each data block is mapped over only one other server, i.e. rK=2rK=2), we demonstrate that Coded MapReduce can still substantially reduce the communication load (for example by more than 50%). As a result, a minor increase in the replication of mapping data blocks at different servers can largely impact the shuffling load of MapReduce.

The idea of Coded MapReduce is inspired by recent results on cache networks, arguing that cache memories can be used not only to deliver part of the contents locally, but also to create some coding opportunities and reduce the traffic significantly . Surprisingly the gain of coding in reducing the traffic is significantly larger than the gain of local delivery, and indeed it can grow with the size of the network. This result has been further generalized to various cache networks in . Interestingly, we demonstrate that similar coding opportunities can also be created in the MapReduce framework via a careful assignment of repetitive Map tasks across servers.

The efficiency of Coded MapReduce in data shuffling builds upon the repetitive executions of the Map tasks, which however requires more processing time for the server cluster to compute the intermediate results. We analytically and numerically evaluate the tradeoff between the Map tasks processing time and the inter-server communication load using Coded MapReduce, based on which clients can choose appropriate operating points to minimize the overall execution time.

The rest of the paper is organized as follows. In Section II, we introduce a MapReduce framework and define the inter-server communication load for the shuffling phase. We then motivate Coded MapReduce through a word-counting example in Section III. Lower bound on the minimum shuffling load (proved in Section VI) and the shuffling load achieved by Coded MapReduce are presented in Section IV, where Coded MapReduce are shown to be approximately optimal. We formally describe the Coded MapReduce framework in Section V, and then analyze its impact on reducing the inter-server communication load. We explore the tradeoff between the computation load of Map tasks and the communication load using Coded MapReduce in Section VII. We finally conclude the paper in Section VIII.

II Problem Statement

For an input file, the job is to evaluate QQ pairs of (key,value)’s, denoted by {(Wq,uq)}q=1Q\{(W_{q},u_{q})\}_{q=1}^{Q}. The keys WqW_{q}, q∈{1,…,Q}q\in\{1,\ldots,Q\} are given. The corresponding output values uqu_{q}, q∈{1,…,Q}q\in\{1,\ldots,Q\}, are evaluated from the input file. The input file consists of NN disjoint subfiles (data blocks). The available processing resources include a cluster of KK servers.

Example (World Counting). Consider a canonical word-counting job of finding the numbers of times 4 given words, denoted by AA, BB, CC and DD, appeared in a book with N=12N=12 chapters using K=4K=4 servers. Each word itself represents a unique key. Thus we have Q=4Q=4 keys each corresponding to a word to be counted and K=4K=4 servers. The final output results are 4 (key,value) pairs (A,a)(A,a), (B,b)(B,b), (C,c)(C,c) and (D,d)(D,d) indicating that there are aa occurrences of AA, bb occurrences of BB, cc occurrences of CC and dd occurrences of DD in the book. \hfill□\hfill\square

A MapReduce approach carries out the job distributedly across the available KK servers. As illustrated in Fig. 1, the MapReduce implementation consists of 3 steps: Map Tasks Assignment, Map Tasks Execution and Reduce Tasks, run over the set of KK servers. In what follows we describe each of these 3 steps in detail.

To start a MapReduce execution, a master controller (JobTracker) assigns the tasks of mapping the NN subfiles among the available KK servers. Each server will focus on processing its assigned subfiles (see details in Step 2). For a given design parameter p∈{1K,2K,…,1}p\in\left\{\frac{1}{K},\frac{2}{K},\ldots,1\right\}, the assignment is done such that each subfile is assigned to be mapped at pKpK servers and each server is assigned pNpN subfiles (assuming that pNpN is an integer). While conventional MapReduce implementations require each subfile to be mapped at exactly one server (p=1Kp=\frac{1}{K}), the master controller can repetitively schedule two or more servers to map the same subfile.

We denote the set of indices of the subfiles assigned to Server kk as Mk\mathcal{M}_{k} and the set of indices of the servers Subfile nn is assigned to as An\mathcal{A}_{n}. Either {Mk}k=1K\{\mathcal{M}_{k}\}_{k=1}^{K} or {An}n=1N\{\mathcal{A}_{n}\}_{n=1}^{N} completely specifies the Map tasks assignment.

Example (Word-Counting: Map Tasks Assignment). In our word-counting example, Q=K=4Q=K=4, N=12N=12. For instance pp can be set to 2K=12\frac{2}{K}=\frac{1}{2}, i.e., each subfile is assigned to be mapped at pK=2pK=2 servers. We can employ a naive Map tasks assignment that assigns each of the first pNpN (6 in this example) chapters to each of the first pKpK servers, each of the second pNpN chapters to each of the second pKpK servers and so forth. For such an assignment, M1=M2={1,2,3,4,5,6}{\mathcal{M}_{1}=\mathcal{M}_{2}=\{1,2,3,4,5,6\}}, M3=M4={7,8,9,10,11,12}{\mathcal{M}_{3}=\mathcal{M}_{4}=\{7,8,9,10,11,12\}}. That is, Server 1 and Server 2 are assigned to map Chapters 1, 2, 3, 4, 5, 6 and Server 3 and Server 4 are assigned to map Chapters 7, 8, 9, 10, 11, 12. \hfill□\hfill\square

Step 2: Map Tasks Execution

Conventionally, the Map tasks execution continues until that each of the NN subfiles has been mapped to (key,value) pairs at one server. Here we introduce a key design parameter rr, r∈{1K,2K,…,p}r\in\left\{\frac{1}{K},\frac{2}{K},\ldots,p\right\} and enforce that for all n∈{1,…,N}n\in\{1,\ldots,N\}, as soon as rKrK servers finish mapping Subfile nn, the rest of the servers in An\mathcal{A}_{n} abort the task of mapping Subfile nn. Thus each subfile is mapped at rKrK servers by the end of Map tasks execution. We denote the set of indices of the servers that finish mapping Subfile nn as An′⊆An\mathcal{A}_{n}^{\prime}\subseteq\mathcal{A}_{n}, where ∣An′∣=rK|\mathcal{A}_{n}^{\prime}|=rK for all nn. rr reflects the redundancy of mapping the same subfile across different servers, and large rr can cause additional delay to the overall Map tasks execution. However, as we demonstrate later, having repetitive mapping outcomes at distinct servers can help to significantly reduce the communication load of the shuffling phase between Map and Reduce tasks.

Example (Word-Counting: Map Tasks Execution). To continue our word-counting example, having adopted the naive Map tasks assignment, the master controller sets r=p=12r=p=\frac{1}{2} such that each server has to finish mapping all of the assigned pN=6pN=6 chapters. As illustrated in Fig. 2, Server kk, k∈{1,2,3,4}k\in\{1,2,3,4\}, maps Chapter nn in Mk\mathcal{M}_{k} into 4 intermediate (key,value) pairs (A,an)[n](A,a_{n})[n], (B,bn)[n](B,b_{n})[n], (C,cn)[n](C,c_{n})[n] and (D,dn)[n](D,d_{n})[n]. For the pair (A,an)[n](A,a_{n})[n], AA is the key, ana_{n} is the value such that ana_{n} is the number of occurrences of AA in Chapter nn While a conventional implementation of the Map function in word-counting jobs maps each encountered word say AA into an intermediate pair (A,1)(A,1), we employ combiners to generate the total number of occurrences of interested words in a chapter.. The definitions of pairs with keys BB, CC and DD follow similarly. \hfill□\hfill\square

Step 3: Reduce Tasks

The goal of this step is to evaluate the final output (key, value) pairs, in a distributed fashion, from the intermediate (key,value) pairs, evaluated in Step 2, i.e. Map tasks execution. One reducer is responsible for evaluating one of the QQ keys, namely W1,…,WQW_{1},\ldots,W_{Q}, and thus a total of QQ reducers are required to be executed on KK servers. The executions of these reducers are distributed uniformly across KK servers. For simplicity, we assume that QK\frac{Q}{K} is an integer and each server executes a disjoint set of QK\frac{Q}{K} reducers. We denote the reducer distribution as D≜(W1,…,WK)\mathcal{D}\triangleq\left(\mathcal{W}_{1},\ldots,\mathcal{W}_{K}\right), where Wk\mathcal{W}_{k}, k∈{1,…,K}k\in\{1,\ldots,K\}, is the set of the indices of the keys evaluated at Server kk. Thus any valid reducer distribution must satisfy:

∣Wk∣=QK|\mathcal{W}_{k}|=\frac{Q}{K} for all k∈{1,…,K}k\in\{1,\ldots,K\},

∪k=1,…,KWk={1,…,Q}\underset{k=1,\ldots,K}{\cup}\mathcal{W}_{k}=\{1,\ldots,Q\},

Wk∩Wk′=∅\mathcal{W}_{k}\cap\mathcal{W}_{k^{\prime}}=\emptyset for k≠k′k\neq k^{\prime}.

To evaluate the key WqW_{q} for some q∈Wkq\in\mathcal{W}_{k}, the corresponding reducer at Server kk needs the values of WqW_{q} in all NN subfiles {vqn}n=1N\{v_{qn}\}_{n=1}^{N}. By the end of the Map tasks execution, Server kk has mapped a subset Mk′\mathcal{M}_{k}^{\prime} of its assigned subfiles Mk\mathcal{M}_{k} and knows the values {vqn:n∈Mk′}{\{v_{qn}:n\in\mathcal{M}_{k}^{\prime}\}}, thus the remaining values {vqn:n∉Mk′}\{v_{qn}:n\notin\mathcal{M}_{k}^{\prime}\} are needed from other servers to execute the reducer for WqW_{q}. We consider a multicast LAN network While current main-stream cloud computing platforms like Amazon EC2 do not support L2 multicast/broadcast, we will demonstrate the significant benefit of enabling multicast in the efficiency of data movement using our proposed Coded MapReduce scheme. where the KK servers are interconnected through a shared link. Next, we formally define the data shuffling scheme to exchange the required values for reduction and the resulting communication load.

A data shuffling scheme for a MapReduce job with NN subfiles, QQ keys, KK servers and parameters p,r,{Mk′}k=1Kp,r,\{\mathcal{M}_{k}^{\prime}\}_{k=1}^{K}, is defined as follows:

The communication process continues for TT times slots, until Server kk, for all k∈{1,…,K}k\in\{1,\ldots,K\}, is able to successfully construct the required values to execute the reducers for the keys in Wk\mathcal{W}_{k}, based on the messages it receives from other servers and its own Map outcomes for the keys in Wk\mathcal{W}_{k} (i.e., {vqn:q∈Wk,n∈Mk′}\{v_{qn}:q\in\mathcal{W}_{k},n\in\mathcal{M}_{k}^{\prime}\}).

Given an instance of the Map tasks execution {Mk′}k=1K\{\mathcal{M}_{k}^{\prime}\}_{k=1}^{K}, each server knows the intermediate results in mapped subfiles for all keys. Therefore, for all valid reducer distributions (i.e., KK mutually disjoint subsets each with QK\frac{Q}{K} keys are reduced respectively at the KK servers), their communication loads are identical. In other words, for a particular Map tasks execution outcome, the communication load is independent of the reducer distribution. \hfill□\hfill\square

Given a MapReduce job of using a cluster of KK servers to evaluate QQ keys in an input file consisting of NN subfiles where each subfile is assigned to be mapped at pKpK servers for some p∈{1K,…,1}p\in\{\frac{1}{K},\ldots,1\}, we are interested in the problem of minimizing the communication load over the Map tasks assignment M1,…,MK\mathcal{M}_{1},\ldots,\mathcal{M}_{K} and the data shuffling scheme. We define L∗(r)L^{*}(r) as the minimum communication load when each subfile is mapped at rKrK servers out of the pKpK servers it is assigned to.

As a baseline, we can consider the conventional MapReduce approach where each subfile is assigned to and mapped at only one server (pK=rK=1pK=rK=1). In this setting, each server maps NK\frac{N}{K} subfiles, obtaining NK\frac{N}{K} intermediate values for each of its assigned keys. Since each server reduces QK\frac{Q}{K} keys, the communication load for the conventional approach is

By increasing pp and rr beyond 1K\frac{1}{K}, each subfile is now repeatitively mapped at rK>1rK>1 servers and the total number of executed Map tasks increases by rKrK times: ∑k=1K∣Mk′∣=rKN\sum\limits_{k=1}^{K}|\mathcal{M}^{\prime}_{k}|=rKN. Server kk, k∈{1,…,K}k\in\{1,\ldots,K\}, needs another QK(N−∣Mk′∣)\frac{Q}{K}\left(N-|\mathcal{M}_{k}^{\prime}|\right) intermediate values to execute its QK\frac{Q}{K} reducers. These data requests can for example be satisfied by a simple uncoded data shuffling scheme such that each of the required intermediate values is sent over the shared link at a time, achieving the following communication load

Let us illustrate the uncoded scheme by describing the shuffling phase of our running word-counting example.

Example (Word-Counting: Data Shuffling via Uncoded Scheme). Since Q=K=4{Q=K=4}, each server executes one reducer, i.e., Server 1 evaluates AA, Server 2 evaluates BB, Server 3 evaluates CC, and Server 4 evaluates DD.

Based on the results of the Map tasks execution (see Fig. 2), Server 1 and Server 2 need the values of AA and BB respectively in Chapters 7, 8, 9, 10, 11, 12. Server 3 and Server 4 need the values of CC and DD respectively in Chapters 1, 2, 3, 4, 5, 6. An uncoded data shuffling is carried out as follows:

Server 3 sends pairs (A,11)(A,11), (A,14)(A,14), (A,10)(A,10), (A,15)(A,15), (A,15)(A,15), (A,10)(A,10) and then (B,21)(B,21), (B,21)(B,21), (B,21)(B,21), (B,21)(B,21), (B,21)(B,21), (B,21)(B,21).

Server 2 sends pairs (C,15)(C,15), (C,15)(C,15), (C,15)(C,15), (C,15)(C,15), (C,15)(C,15), (C,15)(C,15) and then (D,6)(D,6), (D,6)(D,6), (D,6)(D,6), (D,6)(D,6), (D,6)(D,6), (D,6)(D,6).

After the communication every server knows the values of its interested word in all 12 chapters, and these values are passed into a reducer to compute the final result. The shuffling phase lasts for 24 time slots and thus the communication load of the uncoded scheme is 24, which is consistent with equation (2) for N=12N=12, Q=4Q=4 and r=12r=\frac{1}{2}. Notice that, if we had employed the conventional MapReduce approach where each subfile is mapped at only one server, then the communication load of this particular job would have been 36 (this is obtained by setting N=12N=12 and Q=K=4Q=K=4 in equation (1)). \hfill□\hfill\square

Comparing equations (1) and (2), we notice that by repeatedly mapping the same subfile at more than one server (rK≥2rK\geq 2), the communication load of the shuffling phase in MapReduce can be improved by a factor of 1−1K1−r\frac{1-\frac{1}{K}}{1-r}, when using a simple uncoded scheme. This improvement in the communication load results from the fact that by mapping each subfile repeatedly at multiple servers, the servers know values from rKrK times more subfiles than the conventional approach, thus requiring less values communicated during data shuffling to execute their reducers. We denote this gain as the repetition gain, which is due to knowing values from more subfiles locally at each server.

As we will show next, in addition to the repetition gain, repeatedly mapping each subfile at multiple servers can have a much more significant impact on reducing the communication load of the data shuffling, which can be achieved by a more careful assignment of Map tasks to servers and exploiting coding in the data shuffling phase. We will next illustrate this through a motivating example, which forms the basis of the general Coded MapReduce framework that we will later present in Section V.

III Coded MapReduce: A Motivating Example

In this section we motivate Coded MapReduce via a simple example. In particular, we demonstrate through this example that, by carefully assigning Map tasks to the servers, there will be novel coding opportunities in the shuffling phase that can be utilized to significantly reduce the inter-server communication load of MapReduce.

We consider the same word-counting job of counting 4 words AA, BB, CC and DD in a book with 12 chapters using 4 servers. While maintaining the same number subfiles to map at each server (6 in this case), we consider a new Map tasks assignment as follows.

Instead of using the naive assignment, the master controller assigns the Map tasks as follows: M1={1,2,3,4,5,6}\mathcal{M}_{1}=\{1,2,3,4,5,6\}, M2={1,2,7,8,9,10}\mathcal{M}_{2}=\{1,2,7,8,9,10\}, M3={3,4,7,8,11,12}\mathcal{M}_{3}=\{3,4,7,8,11,12\}, M4={5,6,9,10,11,12}\mathcal{M}_{4}=\{5,6,9,10,11,12\}. Notice that in this assignment each chapter is assigned to exactly two servers and every two servers share exactly two chapters.

The master controller sets r=pr=p such that each server has to finish mapping all assigned chapters. The execution of the Map tasks is different from that of the naive assignment such that after generating 4 intermediate (key,value) pairs for each of the assigned chapters, each server generates 3 additional coded (key,value) pairs as follows (see Fig. 3):

Server 1 adds up the values of (B,25)(B,25) and (C,15)(C,15) to generate a pair (BC,40)(BC,40), adds up the values of (B,25)(B,25) and (D,6)(D,6) to generate another pair (BD,31)(BD,31), and adds up the values of (C,15)(C,15) and (D,6)(D,6) to generate a third pair (CD,21)(CD,21),

Server 2 adds up the values of (A,11)(A,11) and (C,15)(C,15) to generate a pair (AC,26)(AC,26), adds up the values of (A,10)(A,10) and (D,6)(D,6) to generate another pair (AD,16)(AD,16), and adds up the values of (C,16)(C,16) and (D,5)(D,5) to generate a third pair (CD,21)(CD,21),

Server 3 adds up the values of (A,14)(A,14) and (B,25)(B,25) to generate a pair (AB,39)(AB,39), adds up the values of (A,15)(A,15) and (D,6)(D,6) to generate another pair (AD,21)(AD,21), and adds up the values of (B,21)(B,21) and (D,5)(D,5) to generate a third pair (BD,26)(BD,26),

Server 4 adds up the values of (A,15)(A,15) and (B,25)(B,25) to generate a pair (AB,40)(AB,40), adds up the values of (A,10)(A,10) and (C,15)(C,15) to generate another pair (AC,25)(AC,25), and adds up the values of (B,21)(B,21) and (C,17)(C,17) to generate a third pair (BC,38)(BC,38).

A coded pair (W1W2,x)[n1,n2](W_{1}W_{2},x)[n_{1},n_{2}] has key W1W2W_{1}W_{2} and value xx, and it indicates that there are xx occurrences in total of Word W1W_{1} in Chapter n1n_{1} and Word W2W_{2} in Chapter n2n_{2}.

The Reduce tasks are distributed the same as before: Server 1 evaluates AA, Server 2 evaluates BB, Server 3 evaluates CC and Server 4 evaluates DD.

After executing the Map tasks, values from 6 chapters are missing at each server to execute the reducer. To fulfill the data requests for reduction, the data shuffling is carried out such that each server sends the 3 coded pairs generated during Map tasks execution:

Server 1 sends pairs (BC,40)(BC,40), (BD,31)(BD,31) and (CD,21)(CD,21),

Server 2 sends pairs (AC,26)(AC,26), (AD,16)(AD,16) and (CD,21)(CD,21),

Server 3 sends pairs (AB,39)(AB,39), (AD,21)(AD,21) and (BD,26)(BD,26),

Server 4 sends pairs (AB,40)(AB,40), (AC,25)(AC,25) and (BC,38)(BC,38).

Having received all coded pairs, each server performs an additional decoding operation before executing the final Reduce function:

Server 1 subtracts the values of (C,15)(C,15), (D,6)(D,6), (B,25)(B,25), (D,6)(D,6), (B,25)(B,25) and (C,15)(C,15) from the values of (AC,26)(AC,26), (AD,16)(AD,16), (AB,39)(AB,39), (AD,21)(AD,21), (AB,40)(AB,40) and (AC,25)(AC,25) respectively to decode 6 pairs it needs to count AA: (A,11)(A,11), (A,10)(A,10), (A,14)(A,14), (A,15)(A,15), (A,15)(A,15) and (A,10)(A,10),

Server 2 subtracts the values of (C,15)(C,15), (D,6)(D,6), (A,14)(A,14), (D,5)(D,5), (A,15)(A,15) and (C,17)(C,17) from the values of (BC,40)(BC,40), (BD,31)(BD,31), (AB,39)(AB,39), (BD,26)(BD,26), (AB,40)(AB,40) and (BC,38)(BC,38) respectively to decode 6 pairs it needs to count BB: (B,25)(B,25), (B,25)(B,25), (B,25)(B,25), (B,21)(B,21), (B,25)(B,25) and (B,21)(B,21),

Server 3 subtracts the values of (B,25)(B,25), (D,6)(D,6), (A,11)(A,11), (D,5)(D,5), (A,10)(A,10) and (B,21)(B,21) from the values of (BC,40)(BC,40), (CD,21)(CD,21), (AC,26)(AC,26), (CD,21)(CD,21), (AC,25)(AC,25) and (BC,38)(BC,38) respectively to decode 6 pairs it needs to count CC: (C,15)(C,15), (C,15)(C,15), (C,15)(C,15), (C,16)(C,16), (C,15)(C,15) and (C,17)(C,17),

Server 4 subtracts the values of (B,25)(B,25), (C,15)(C,15), (A,10)(A,10), (C,16)(C,16), (A,15)(A,15) and (B,21)(B,21) from the values of (BD,31)(BD,31), (CD,21)(CD,21), (AD,16)(AD,16), (CD,21)(CD,21), (AD,21)(AD,21) and (BD,26)(BD,26) respectively to decode 6 pairs it needs to count DD: (D,6)(D,6), (D,6)(D,6), (D,6)(D,6), (D,5)(D,5), (D,6)(D,6) and (D,5)(D,5).

Now each server knows the value of the interested word in each of the 12 chapters, and passes these values into a Reduce function to generate the final result. Each server accesses the shared link 3 times during the shuffling phase and the total communication load is 1212.

We notice that having successfully communicated the required values for the Reduce functions, the communication load of the proposed Coded MapReduce is 66% less than that of the conventional MapReduce approach, and 50% less than that of the naive assignment and uncoded shuffling scheme. The savings in the communication load result from the fact that each of the coded pairs (W1W2,v1n1+v2n2)[n1,n2](W_{1}W_{2},v_{1n_{1}}+v_{2n_{2}})[n_{1},n_{2}] simultaneously delivers v1n1v_{1n_{1}} to the server counting W1W_{1} and v2n2v_{2n_{2}} to the server counting W2W_{2}, given that v1n1v_{1n_{1}} is known at the server counting W2W_{2} and v2n2v_{2n_{2}} is known at the server counting W1W_{1} prior to data shuffling. For example because b3b_{3} is needed by Server 2 and known at Server 3 and c1c_{1} is needed by Server 3 and known at Server 2, having Server 1 send (BC,b3+c1)(BC,b_{3}+c_{1}) simultaneously delivers b3b_{3} to Server 2 and c1c_{1} to Server 3. We call such opportunities of effectively communicating multiple values through a single use of the shared link as Coded Multicast. Similar ideas of optimally creating and exploiting the coding opportunities have been first used in caching problems . This type of coding is closely related to the network coding problem and the index coding problem . More detailed discussion about the connections between these problems can be found in .

For this particular word-counting example, one can save communication load by adding the word counts in different chapters right after Map tasks execution and sending the sum instead of the counts in individual chapters. However, we emphasize that for a general MapReduce job where such pre-reduction might not be applicable (Reduce function is not associative or commutative), the proposed Coded MapReduce scheme, which is presented in detail in Section V, can still provide a significant gain in reducing the shuffling load. With sufficiently large number of subfiles, this gain scales linearly with the number of servers in the system. \hfill□\hfill\square

IV Main Results

In this section we present and discuss upper and lower bounds on the minimum communication load L∗(r)L^{*}(r) of a general MapReduce job. We also analytically and numerically compare the communication loads of different Map tasks assignment and data shuffling schemes to demonstrate the substantial reduction in communication load of our proposed Coded MapReduce scheme compared with the conventional MapReduce approach and the uncoded shuffling scheme.

Consider a job of using KK servers to evaluate QQ keys in an input file consisting of NN subfiles, where each subfile is repetitively assigned to pKpK servers and is randomly and uniformly mapped at rKrK of those servers for some r∈{1K,…,p}r\in\{\frac{1}{K},\ldots,p\}. The minimum communication load L∗(r)L^{*}(r) is bounded as

The upper bound LCMR(r)L_{\textup{CMR}}(r) of the minimum communication load is achieved by our proposed Coded MapReduce scheme that is described and analyzed in detail in the next section. The two lower bounds on the left hand side of (3) are derived by applying the cut-set bounds on the compound extension of the MapReduce system with multiple valid reducer distributions (i.e., which keys are reduced at which servers), and the detailed proofs are provided in Section VI.

Next, we demonstrate through the following corollary of Theorem 1 that the proposed Coded MapReduce significantly reduce the inter-server communication load compared to the conventional MapReduce approach.

Consider a job of using KK servers to evaluate QQ keys in an input file consisting of NN subfiles, where each subfile is repetitively assigned to pKpK servers according to the Coded MapReduce scheme, which will be described in Section V-A, and is randomly and uniformly mapped at rKrK of those servers for some r∈{1K,…,p}r\in\{\frac{1}{K},\ldots,p\}. Then, the communication load of the proposed Coded MapReduce scheme LCMR(r)L_{\textup{CMR}}(r) satisfies

where LconvL_{\textup{conv}} is the communication load of the conventional MapReduce approach, as defined in (1).

Corollary 1 demonstrates that compared with the conventional MapReduce approach, for sufficiently large NN, Coded MapReduce can cut down the communication load substantially by a multiplicative factor of 1−1K1−r(rK)≥rK\frac{1-\frac{1}{K}}{1-r}(rK)\geq rK, which grows linearly with the number of servers in the system. \hfill□\hfill\square

The first term on the right hand side of (4) (i.e., 1−r1−1K\frac{1-r}{1-\frac{1}{K}}) can be viewed as the repetition gain, which is due to knowing values from more subfiles locally at each server, and the second term (i.e., 1rK\frac{1}{rK}) can be viewed as the coding gain. Note that 1−r1−1K≥1−r\frac{1-r}{1-\frac{1}{K}}\geq 1-r, hence the repetition gain does not scale with KK, so the overall gain of Coded MapReduce is mostly due to coding. \hfill□\hfill\square

As an example, we have numerically compared the communication loads required by the conventional MapReduce approach, the uncoded shuffling scheme and the Coded MapReduce for a specific set of parameters in Fig. 4. One can note that for an input file consisting of 12001200 subfiles (N=1200N=1200), when each subfile is repeatedly assigned to 7 server and actually mapped at 2 servers (rK=2rK=2), we observe a repetition gain of 1.125×1.125\times, a coding gain of 1.81×1.81\times by Coded MapReduce over the uncoded scheme, and an overall 2.03×2.03\times reduction in communication load from the conventional approach by Coded MapReduce. When rKrK is increased to 77, i.e., every server has to finish mapping all its assigned subfiles, the repetition gain increases to 3×3\times, the coding gain over the uncoded scheme increases to 7×7\times and Coded MapReduce achieves an overall 21×21\times reduction of the communication load compared with the conventional MapReduce approach. Lastly, we notice that the coding gain for this particular job is proportional to rKrK, which is consistent with (4). \hfill□\hfill\square

Next we demonstrate via the following theorem that the proposed Coded MapReduce scheme is approximately optimal in the sense that it achieved the minimum communication load within a constant multiplicative factor regardless of the system parameters.

The proposed Coded MapReduce scheme achieves the minimum communication load up to a constant multiplicative factor for any MapReduce job. More precisely:

where LCMR(r)L_{\textup{CMR}}(r) is defined in (3).

Taking the limit of (7) as NN goes to infinity, we have

We recall that r∈{1K,…,1}r\in\{\frac{1}{K},\ldots,1\} and proceed to bound (8) in the following two regions:

In this case, we set s=1s=1 in the lower bound of L∗(r)L^{*}(r) and (8) becomes

In this case, we set s=⌊12r⌋s=\left\lfloor\frac{1}{2r}\right\rfloor in the lower bound of L∗(r)L^{*}(r) to obtain

where (a) is because that rKrK is a positive integer.

Comparing the two bounds in (10) and (17) completes the proof. ∎

V Coded MapReduce: General Description and Performance Analysis

In this section, we present our proposed Coded MapReduce scheme on a general MapReduce job, and analyze the corresponding inter-server communication load.

We consider a MapReduce job of evaluating QQ keys on an input file with NN subfiles using KK servers, where each subfile is assigned to pKpK servers for Map tasks and is actually mapped at rKrK of them. Before proceeding to detailed descriptions of the Coded MapReduce scheme, we recall that after the Map tasks execution, Mk′\mathcal{M}_{k}^{\prime} indicates the set of subfiles mapped at Server kk, and An′\mathcal{A}_{n}^{\prime} indicates the set of servers that have finished mapping Subfile nn.

Map Tasks Assignment

We assume that NN is sufficiently large and N=g(KpK)N=g\begin{pmatrix}K\\ pK\end{pmatrix} for some integer gg Otherwise we can introduce enough empty subfiles to have an integer-valued gg.. The master controller partitions the subfiles into (KpK)\begin{pmatrix}K\\ pK\end{pmatrix} equal-sized subsets, where each subset contains gg unique subfiles. We call each of these subsets a batch of subfiles. For each batch of gg subfiles, the master controller chooses a distinct subset of pKpK servers, and assign all gg subfiles in the batch to each of these pKpK servers.

Because each server belongs to (K−1pK−1)\begin{pmatrix}K-1\\ pK-1\end{pmatrix} subsets of size pKpK, each server is assigned g(K−1pK−1)g\begin{pmatrix}K-1\\ pK-1\end{pmatrix} subfiles by the end of Map tasks assignment, i.e., ∣Mk∣=g(K−1pK−1)|\mathcal{M}_{k}|=g\begin{pmatrix}K-1\\ pK-1\end{pmatrix} for all k∈{1,…,K}k\in\{1,\ldots,K\}.

For the motivating example in Section III, N=12N=12, Q=K=4{Q=K=4}, pK=2pK=2, thus we have g=2g=2. Every 22 servers are assigned 22 unique chapters.

Map Tasks Execution

After the Map tasks assignment, each server starts to map all of its assigned subfiles simultaneously. We assume that the times the KK servers spend mapping their assigned subfiles are i.i.d. across servers and subfiles. Thus the probability that any subset of rKrK servers in An\mathcal{A}_{n} finish mapping Subfile nn by the end of Map tasks execution is 1(pKrK)\dfrac{1}{\begin{pmatrix}pK\\ rK\end{pmatrix}}, for all n∈{1,…,N}{n\in\{1,\ldots,N\}}. Because each group of rKrK servers are assigned g(K−rKpK−rK)g\begin{pmatrix}K-rK\\ pK-rK\end{pmatrix} subfiles, by the end of Map tasks execution, the expected number of subfiles mapped at any group of rKrK servers is

By law of large numbers and the fact that N=g(KpK)>g(K−rKpK−rK){N=g\begin{pmatrix}K\\ pK\end{pmatrix}>g\begin{pmatrix}K-rK\\ pK-rK\end{pmatrix}}, for NN large enough the actual number of subfiles mapped at any group of rKrK servers is

Before proceeding to describe the Reduce tasks, we define an important quantity VSk\mathcal{V}^{k}_{\mathcal{S}} for k∈{1,…,K}k\in\{1,\ldots,K\} and S⊆{1,…,K}{{\cal S}\subseteq\{1,\ldots,K\}} such that VSk\mathcal{V}^{k}_{\mathcal{S}} is a set of information bits required by Server kk and known exclusively at all servers in S{\cal S}. More precisely, VSk≜{vqn:q∈Wk,k∉An′=S}\mathcal{V}^{k}_{{\cal S}}\triangleq\left\{v_{qn}:q\in\mathcal{W}_{k},k\notin\mathcal{A}_{n}^{\prime}={\cal S}\right\}.

Data Shuffling

The inter-server communication is carried out as follows to exchange the missing values for Reduce functions:

For each subset S{\cal S} of {1,…,K}\{1,\ldots,K\} such that ∣S∣ ⁣= ⁣rK ⁣+ ⁣1{|{\cal S}|\!=\!rK\!+\!1} and each k∈Sk\in\mathcal{S}, arbitrarily partition VS\{k}k\mathcal{V}^{k}_{\mathcal{S}\backslash\{k\}} into rKrK disjoint segments each containing ∣VS\{k}k∣rK\frac{\left|\mathcal{V}^{k}_{\mathcal{S}\backslash\{k\}}\right|}{rK} bits:

where VS\{k},ik\mathcal{V}^{k}_{\mathcal{S}\backslash\{k\},i} is the segment associated with Server ii in S\{k}\mathcal{S}\backslash\{k\}.

For each i∈Si\in\mathcal{S}, Server ii zero-pads all associated segments {VS\{k},ik:k∈S\{i}}\left\{\mathcal{V}^{k}_{\mathcal{S}\backslash\{k\},i}:k\in\mathcal{S}\backslash\{i\}\right\} to the length of the longest one and sends the following coded segment:

where ⊕\oplus denotes bitwise XOR on the zero-padded segments.

After Server kk, k∈S\{i}k\in\mathcal{S}\backslash\{i\} receives the coded segment ⊕k∈S\{i}VS\{k},ik\underset{k\in\mathcal{S}\backslash\{i\}}{\oplus}\mathcal{V}^{k}_{\mathcal{S}\backslash\{k\},i} from Server ii, because Server kk knows all the information bits in 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 coded segment and recover the intended bits in VS\{k},ik\mathcal{V}^{k}_{\mathcal{S}\backslash\{k\},i}.

Consider the the motivating example in Section III. For the subset of servers S={1,2,3}\mathcal{S}=\{1,2,3\}, V{2,3}1={a7,a8}{\mathcal{V}^{1}_{\{2,3\}}=\{a_{7},a_{8}\}}, V{1,3}2={b3,b4}\mathcal{V}^{2}_{\{1,3\}}=\{b_{3},b_{4}\}, and V{1,2}3={c1,c2}\mathcal{V}^{3}_{\{1,2\}}=\{c_{1},c_{2}\}. These sets of values are segmented such that

V{2,3},21={a7}\mathcal{V}^{1}_{\{2,3\},2}=\{a_{7}\}, V{2,3},31={a8}\mathcal{V}^{1}_{\{2,3\},3}=\{a_{8}\},

V{1,3},12={b3}\mathcal{V}^{2}_{\{1,3\},1}=\{b_{3}\}, V{1,3},32={b4}\mathcal{V}^{2}_{\{1,3\},3}=\{b_{4}\},

V{1,2},13={c1}\mathcal{V}^{3}_{\{1,2\},1}=\{c_{1}\}, V{1,2},23={c2}\mathcal{V}^{3}_{\{1,2\},2}=\{c_{2}\}.

Server 1 sends a coded pair (BC,b3+c1)(BC,b_{3}+c_{1}),

Server 2 sends a coded pair (AC,a7+c2)(AC,a_{7}+c_{2}),

Server 3 sends a coded pair (AB,a8+b4)(AB,a_{8}+b_{4}).

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 that coding opportunities can be optimally exploited to reduce network traffic during the content delivery. We notice that interestingly, such opportunities exist in a general MapReduce environment if each subfile is repeatedly mapped at different servers, and we propose Coded MapReduce to utilize them to significantly reduce the shuffling load. \hfill□\hfill\square

We summarize the proposed Coded MapReduce scheme in Algorithm 1.

V-B Performance of Coded MapReduce

We first demonstrate that the inter-server communication of Coded MapReduce successfully delivers all required bits for reduction. To start, we pick an arbitrary server, say Server kk, and consider a required bit to execute its reducers. If this bit is from subfiles that are mapped at Server kk (i.e., Mk′\mathcal{M}_{k}^{\prime}), it is readily available for reduction and no communication is needed. However if it is from some Subfile nn that is not mapped at Server kk, i.e., outside Mk′\mathcal{M}_{k}^{\prime}, we recall that An′\mathcal{A}_{n}^{\prime} (k∉An′k\notin\mathcal{A}_{n}^{\prime}) denotes the set of servers that have mapped Subfile nn and there are rKrK of them. Now we consider Line 11 of Algorithm 1 for the iteration of S={k}∪An′\mathcal{S}=\{k\}\cup\mathcal{A}_{n}^{\prime}. Suppose that this bit is in the segment VS\{k},ik\mathcal{V}^{k}_{\mathcal{S}\backslash\{k\},i} of some Server i∈An′i\in\mathcal{A}_{n}^{\prime}, because for all k′∈S\{i,k}k^{\prime}\in\mathcal{S}\backslash\{i,k\}, {i,k}⊆S\{k′}\{i,k\}\subseteq\mathcal{S}\backslash\{k^{\prime}\} and Server kk knows all the bits in VS\{k′},ik′\mathcal{V}^{k^{\prime}}_{\mathcal{S}\backslash\{k^{\prime}\},i}, Server kk can decode its required bit from the coded segment ⊕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 by cancelling 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\}. Similarly, by the end of the shuffling phase, Server kk decodes all required bits for its reducers from coded segments transmitted by other servers. The same arguments directly apply to all servers and we have demonstrated that Coded MapReduce successfully provides all values needed for reduction.

Next, we analytically characterise the inter-server communication load of Coded MapReduce.

In each iteration of the inter-server communication in Algorithm 1, by (18) we know that after Map tasks execution, every rKrK servers share the values of all keys in g(K−rKpK−rK)(pKrK)+o(N){\dfrac{g\begin{pmatrix}K-rK\\ pK-rK\end{pmatrix}}{\begin{pmatrix}pK\\ rK\end{pmatrix}}+o(N)} subfiles, thus for any subset S\mathcal{S} of rK+1rK+1 servers and sufficiently large NN,

for all k∈Sk\in\mathcal{S}, given that Server kk needs QK\frac{Q}{K} values for its reducers from each subfile.

For each S⊆{1,…,K}\mathcal{S}\subseteq\{1,\ldots,K\} with ∣S∣=rK+1|\mathcal{S}|=rK+1, by Lines 17 and 18 of Algorithm 1, Server ii sends a coded segment of length

for all i∈Si\in\mathcal{S}. Therefore, the total number of bits sent over the shared link in S\mathcal{S} is

Because Algorithm 1 has a total of (KrK+1)\begin{pmatrix}K\\ rK+1\end{pmatrix} iterations, the communication load achieved by the proposed Coded MapReduce scheme, normalized by FF, is

where (a) is because that N=g(KpK)N=g\begin{pmatrix}K\\ pK\end{pmatrix}.

In this section, two lower bounds on the inter-server communication load L∗(r)L^{*}(r) are derived by applying the cut-set bounds on the compound extensions of multiple valid reducer distributions (i.e., which servers reduce which keys). Before proceeding to the formal proofs, we re-emphasize that as highlighted in Remark 1, the communication loads of all valid reducer distributions are identical. As a first step, we partition the indices of the keys {1,…,Q}\{1,\ldots,Q\} evenly into KK groups G1,…,GK\mathcal{G}_{1},\ldots,\mathcal{G}_{K} such that Gi={(i−1)QK+1,…,iQK}\mathcal{G}_{i}=\left\{\frac{(i-1)Q}{K}+1,\ldots,\frac{iQ}{K}\right\}, for all i∈{1,…,K}i\in\{1,\ldots,K\}. Given a reducer distribution D=(W1,…,WK)\mathcal{D}=\left(\mathcal{W}_{1},\ldots,\mathcal{W}_{K}\right) that specifies which keys are reduced at which servers, we denote the messages sent by Server kk during the data shuffling as XkDX_{k}^{\mathcal{D}}, for all k∈{1,…,K}k\in\{1,\ldots,K\} with RkDFR_{k}^{\mathcal{D}}F total number of bits.

We consider the following KK valid reducer distributions on the KK servers:

For a particular reducer distribution Di\mathcal{D}_{i}, i∈{1,…,K}i\in\{1,\ldots,K\}, we denote its kkth entry (i.e., the set of indices of the keys reduced at Server kk) as Di,k\mathcal{D}_{i,k}.

Now we consider the compound setting of all KK reducer distributions at Server kk, k∈{1,…,K}k\in\{1,\ldots,K\}. Server kk needs to evaluate all keys in {Di,k:i∈{1,…,K}}\left\{\mathcal{D}_{i,k}:i\in\{1,\ldots,K\}\right\}, which is equal to the set of all keys {1,…,Q}\{1,\ldots,Q\} by the construction of the reducer distributions in (20). During the shuffling phase for a particular reducer distribution Di\mathcal{D}_{i}, Server kk receives ∑k′≠kRk′DiF\sum\limits_{k^{\prime}\neq k}R_{k^{\prime}}^{\mathcal{D}_{i}}F information bits from other servers, and overall ∑i=1K∑k′≠kRk′DiF\sum\limits_{i=1}^{K}\sum\limits_{k^{\prime}\neq k}R_{k^{\prime}}^{\mathcal{D}_{i}}F bits of information across all KK reducer distributions. Together with its mapping outcomes {vqn:q∈{1,…,Q},n∈Mk′}\{v_{qn}:q\in\{1,\ldots,Q\},n\in\mathcal{M}_{k}^{\prime}\}, Server kk should be able to evaluate all QQ keys after receiving all messages from all other servers in all reducer distributions. Thus we are essentially considering a cut at Server kk separating ({vqn:q∈{1,…,Q},n∈Mk′},{Xk′D1,…,Xk′DK:k′≠k})\left(\{v_{qn}:q\in\{1,\ldots,Q\},n\in\mathcal{M}_{k}^{\prime}\},\left\{X_{k^{\prime}}^{\mathcal{D}_{1}},\ldots,X_{k^{\prime}}^{\mathcal{D}_{K}}:k^{\prime}\neq k\right\}\right) and {vqn:q∈{1,…,Q},n∈{1,…,N}}\left\{v_{qn}:q\in\{1,\ldots,Q\},n\in\{1,\ldots,N\}\right\}, and we have

Summing up the cut-set bounds of all KK servers, we have

where (b) is due to the Map tasks execution that each subfile is completely mapped at rKrK servers, and (c) is because that by Remark 1 the communication load of the data shuffling scheme is independent of the reducer distribution.

Taking expectations of both sides of (23) over the Map tasks executions {Mk′}k=1K\{\mathcal{M}_{k}^{\prime}\}_{k=1}^{K}, we obtain the following lower bound of the minimum inter-server communication load:

Second Bound

Let s∈{1,…,K}s\in\{1,\ldots,K\} and we focus on the keys evaluated at the first ss servers. We consider a set of ⌊Ks⌋\left\lfloor\frac{K}{s}\right\rfloor valid reducer distributions such that the indices of the keys reduced by the first ss servers are respectively specified as:

where Di,(1,…,s)\mathcal{D}_{i,(1,\ldots,s)}, i∈{1,…,⌊Ks⌋}i\in\{1,\ldots,\left\lfloor\frac{K}{s}\right\rfloor\} denotes indices of the keys evaluated by the first ss servers in iith reducer distribution.

Now for the compound setting of the ⌊Ks⌋\left\lfloor\frac{K}{s}\right\rfloor reducer distributions in (25), Server kk, k∈{1,…,s}k\in\{1,\ldots,s\}, needs to reduce the keys in {G(i−1)s+k:i∈{1,…,⌊Ks⌋}}\left\{\mathcal{G}_{(i-1)s+k}:i\in\{1,\ldots,\left\lfloor\frac{K}{s}\right\rfloor\}\right\}, given the transmitted messages for all reducer distributions {XkDi:k∈{1,…,K},i∈{1,…,⌊Ks⌋}}\left\{X_{k}^{\mathcal{D}_{i}}:k\in\{1,\ldots,K\},i\in\{1,\ldots,\left\lfloor\frac{K}{s}\right\rfloor\}\right\} and its mapping outcomes {vqn:q∈{1,…,Q},n∈Mk′}\{v_{qn}:q\in\{1,\ldots,Q\},n\in\mathcal{M}_{k}^{\prime}\}.

We consider the cut containing all ss servers that separates all values needed for reduction and the transmitted messages in all reducer distributions in addition to their local mapping outcomes, i.e., {vqn:q∈{1,…,QK⌊Ks⌋s},n∈{1,…,N}}\left\{v_{qn}:q\in\left\{1,\ldots,\frac{Q}{K}\left\lfloor\frac{K}{s}\right\rfloor s\right\},n\in\{1,\ldots,N\}\right\} and ({vqn:q∈{1,…,Q},n∈Mk′}k=1s,{XkDi:k∈{1,…,K},i∈{1,…,⌊Ks⌋}})\left(\left\{v_{qn}:q\in\{1,\ldots,Q\},n\in\mathcal{M}_{k}^{\prime}\right\}_{k=1}^{s},\left\{X_{k}^{\mathcal{D}_{i}}:k\in\{1,\ldots,K\},i\in\{1,\ldots,\left\lfloor\frac{K}{s}\right\rfloor\}\right\}\right), we have

Because (27) holds for all s∈{1,…,K}s\in\{1,\ldots,K\}, we obtain another lower bound of L∗(r)L^{*}(r):

For the example in Section III where Q=4Q=4, N=12N=12, K=4K=4 and r=12r=\frac{1}{2}, the first lower bound is tighter than the second, yielding that L∗(12)≥8L^{*}(\frac{1}{2})\geq 8.

VII Map Processing Time vs. Communication Load of Coded MapReduce

As Coded MapReduce slashes the inter-server communication load through repetitive Map operations across servers, processing more Map tasks incurs extra delay to the job execution. In this section, we analytically and numerically investigate this effect and the tradeoff between the processing time of the Map tasks and the inter-server communication load of Coded MapReduce.

We first recall that after the Map tasks assignment, each server starts to process the assigned pNpN subfiles until each of the NN subfiles has been mapped at rKrK distinct servers. We assume that all subfiles are of the same size and are simultaneously available for processing. We employ the processor-sharing (PS) model in as our service discipline where the processing rate of each server, denoted as μ\mu, is evenly distributed among all assigned pNpN Map tasks. We assume that the KK servers start their respective Map tasks at the same time and the Map processing times of the subfiles within and across servers are i.i.d. exponential random variables with rate μpN\frac{\mu}{pN} For simplicity, we do not consider the rate change due to the completions of Map tasks and we assume that all subfiles are being mapped at the constant rate μpN\frac{\mu}{pN} throughout the Map tasks execution.. Suppose that the Map tasks are carried out error-free, then the Map processing time of a single subfile say Subfile nn, SnS_{n} is the rKrKth-order statistic of pKpK i.i.d. exponential random variables with rate μpN\frac{\mu}{pN}, and the probability density function of SnS_{n} by is

where f(s)=μpNe−μpNsf(s)=\frac{\mu}{pN}e^{-\frac{\mu}{pN}s} and F(s)=1−e−μpNsF(s)=1-e^{-\frac{\mu}{pN}s} are the PDF and CDF of an exponential random variable with rate μpN\frac{\mu}{pN} respectively.

We can further derive the CDF of the processing time of Subfile nn (i.e., SnS_{n}) as

The average processing time to map Subfile nn, n∈{1,…,N}n\in\{1,\ldots,N\}, at rKrK servers out of the pKpK servers it is assigned to, i.e., the expected value of SnS_{n}, is (see e.g. )

Because the Map tasks execution continues until each of the NN subfiles is mapped at rKrK servers, the overall Map processing time of NN subfiles, SS is

where S1,…,SNS_{1},\ldots,S_{N} are i.i.d. random variables with the probability density function in (29) and the cumulative distribution function in (30).

The distribution function of the overall processing time SS can be calculated as FS(s)=∏n=1NPr(Sn≤s)=(FSn(s))NF_{S}(s)=\prod\limits_{n=1}^{N}\textup{Pr}(S_{n}\leq s)=\left(F_{S_{n}}(s)\right)^{N}, and thus the average overall processing time for the Map tasks is

where FSn(s)F_{S_{n}}(s) is calculated in (30).

VII-B Numerical Evaluations

In Fig. 5 and Fig. 6, we have respectively plotted the average processing time to map one subfile and to map all subfiles and its corresponding shuffling load, using Coded MapReduce. As expected, when increasing the number of repetitive mapping (rKrK), the Map processing time becomes longer while less communication is required for shuffling. Coded MapReduce provides a flexible framework for system designers to balance between the time spent in the Map phase and the shuffling phase of MapReduce jobs, based on the underlying infrastructures (e.g., server processing speed and network speed), to optimize the run-time performance.

VIII Conclusion

In this paper, we propose Coded MapReduce, a joint framework that aims to improve the performance of data shuffling in a MapReduce environment. We demonstrate through analytical and numerical analysis that, by carefully assigning repetitive Map tasks onto different servers and smartly coding the transmitted message bits across keys and data blocks, Coded MapReduce can slash the inter-server communication load by a factor that grows linearly with number of servers in the system. We also demonstrate the optimality of Coded MapReduce by showing that it achieves the minimum communication load within a constant multiplicative factor. Further, we explore the inherent tradeoff of Coded MapReduce between the computation load of the Map tasks and the communication load of the data shuffling, shedding light on optimizing Coded MapReduce to minimize the overall job completion time. An interesting future direction is to consider a software implementation of Coded MapReduce in Hadoop clusters to demonstrate the coding gain for reducing the inter-server communication load of the shuffling phase in such clusters.

References