Speeding Up Distributed Machine Learning Using Codes

Kangwook Lee, Maximilian Lam, Ramtin Pedarsani, Dimitris Papailiopoulos, Kannan Ramchandran

I Introduction

In recent years, the computational paradigm for large-scale machine learning and data analytics has shifted towards massively large distributed systems, comprising individually small and unreliable computational nodes (low-end, commodity hardware). Specifically, modern distributed systems like Apache Spark and computational primitives like MapReduce have gained significant traction, as they enable the execution of production-scale tasks on data sizes of the order of petabytes. However, it is observed that the performance of a modern distributed system is significantly affected by anomalous system behavior and bottlenecks , i.e., a form of “system noise”. Given the individually unpredictable nature of the nodes in these systems, we are faced with the challenge of securing fast and high-quality algorithmic results in the face of uncertainty.

In this work, we tackle this challenge using coding theoretic techniques. The role of codes in providing resiliency against noise has been studied for decades in several other engineering contexts, and is part of our everyday infrastructure (smartphones, laptops, WiFi and cellular systems, etc.). The goal of our work is to apply coding techniques to blueprint robust distributed systems, especially for distributed machine learning algorithms. The workflow of distributed machine learning algorithms in a large-scale system can be decomposed into three functional phases: a storage, a communication, and a computation phase, as shown in Fig. 1. In order to develop and deploy sophisticated solutions and tackle large-scale problems in machine learning, science, engineering, and commerce, it is important to understand and optimize novel and complex trade-offs across the multiple dimensions of computation, communication, storage, and the accuracy of results. Recently, codes have begun to transform the storage layer of distributed systems in modern data centers under the umbrella of regenerating and locally repairable codes for distributed storage which are also having a major impact on industry .

In this paper, we explore the use of coding theory to remove bottlenecks caused during the other phases: the communication and computation phases of distributed algorithms. More specifically, we identify two core blocks relevant to the communication and computation phases that we believe are key primitives in a plethora of distributed data processing and machine learning algorithms: matrix multiplication and data shuffling.

For matrix multiplication, we use codes to leverage the plethora of nodes and alleviate the effect of stragglers, i.e., nodes that are significantly slower than average. We show analytically that if there are nn workers having identically distributed computing time statistics that are exponentially distributed, the optimal coded matrix multiplication is Θ(log⁡n)\Theta(\log n)For any two sequences f(n)f(n) and g(n)g(n): f(n)=Ω(g(n))f(n)=\Omega(g(n)) if there exists a positive constant cc such that f(n)≥cg(n)f(n)\geq cg(n); f(n)=o(g(n))f(n)=o(g(n)) if lim⁡n→∞f(n)g(n)=0\lim_{n\rightarrow\infty}\frac{f(n)}{g(n)}=0. times faster than the uncoded matrix multiplication on average.

Data shuffling is a core element of many machine learning applications, and is well-known to improve the statistical performance of learning algorithms. We show that codes can be used in a novel way to trade off excess in available storage for reduced communication cost for data shuffling done in parallel machine learning algorithms. We show that when a constant fraction of the data matrix can be cached at each worker, and nn is the number of workers, coded shuffling reduces the communication cost by a factor Θ(γ(n))\Theta(\gamma(n)) compared to uncoded shuffling, where γ(n)\gamma(n) is the ratio of the cost of unicasting nn messages to nn users to multicasting a common message (of the same size) to nn users. For instance, γ(n)≃n\gamma(n)\simeq n if multicasting a message to nn users is as cheap as unicasting a message to one user.

We would like to remark that a major innovation of our coding solutions is that they are woven into the fabric of the algorithmic design, and coding/decoding is performed over the representation field of the input data (e.g., floats or doubles). In sharp contrast to most coding applications, we do not need to “re-factor code” and modify the distributed system to accommodate for our solutions; it is all done seamlessly in the algorithmic design layer, an abstraction that we believe is much more impactful as it is located “higher up” in the system layer hierarchy compared to traditional applications of coding that need to interact with the stored and transmitted “bits” (e.g., as is the case for coding solutions for the physical or storage layer).

We now provide a brief overview of the main results of this paper.

Coded Computation designs parallel tasks for a linear operation using erasure codes such that its runtime is not affected by up to a certain number of stragglers. Matrix multiplication is one of the most basic linear operations and is the workhorse of a host of machine learning and data analytics algorithms, e.g., gradient descent based algorithm for regression problems, power-iteration like algorithms for spectral analysis and graph ranking applications, etc. Hence, we focus on the example of matrix multiplication in this paper. With coded computation, we will show that the runtime of the algorithm can be significantly reduced compared to that of other uncoded algorithms. The main result on Coded Computation is stated in the following (informal) theorem.

If the number of workers is nn, and the runtime of each subtask has an exponential tail, the optimal coded matrix multiplication is Θ(log⁡n)\Theta(\log n) times faster than the uncoded matrix multiplication.

For the formal version of the theorem and its proof, see Sec. III-D.

We now overview the main results on coded shuffling. Consider a master-worker setup where a master node holds the entire data set. The generic machine learning task that we wish to optimize is the following: 1) the data set is randomly permuted and partitioned in batches at the master; 2) the master sends the batches to the workers; 3) each worker uses its batch and locally trains a model; 4) the local models are averaged at the master and the process is repeated. To reduce communication overheads between master and workers, Coded Shuffling exploits i) the locally cached data points of previous passes and ii) the “transmission strategy” of the master node.

We illustrate the basics of Coded Shuffling with a toy example. Consider a system with two worker nodes and one master node. Assume that the data set consists of 4 batches A1,…,A4{\bf A}_{1},\ldots,{\bf A}_{4}, which are stored across two workers as shown in Fig. 3. The sole objective of the master is to transmit A3{\bf A}_{3} to the first worker and A4{\bf A}_{4} to the second. For this purpose, the master node can simply multicast a coded message A2+A3\mathbf{A}_{2}+\mathbf{A}_{3} to the worker nodes since the workers can decode the desired batches using the stored batches. Compared to the naïve (or uncoded) shuffling scheme in which the master node transmits A2\mathbf{A}_{2} and A3\mathbf{A}_{3} separately, this new shuffling scheme can save 50%50\% of the communication cost, speeding up the overall machine learning algorithm. The Coded Shuffling algorithm is a generalization of the above toy example, which we explain in detail in Sec. IV.

Note that the above example assumes that multicasting a message to all workers costs exactly the same as unicasting a message to one of the workers. In general, we capture the advantage of using multicasting over unicasting by defining γ(n)\gamma(n) as follows:

Clearly, 1≤γ(n)≤n1\leq\gamma(n)\leq n: if γ(n)=n\gamma(n)=n, the cost of multicasting is equal to that of unicasting a single message (as in the above example); if γ(n)=1\gamma(n)=1, there is essentially no advantage of using multicast over unicast.

We now state the main result on Coded Shuffling in the following (informal) theorem.

Let α\alpha be the fraction of the data matrix that can be cached at each worker, and nn be the number of workers. Assume that the advantage of multicasting over unicasting is γ(n)\gamma(n). Then, coded shuffling reduces the communication cost by a factor of (α+1n)γ(n)\left(\alpha+\frac{1}{n}\right)\gamma(n) compared to uncoded shuffling.

For the formal version of the theorem and its proofs, see Sec. IV-D.

The remainder of this paper is organized as follows. In Sec. II, we provide an extensive review of the related works in the literature. Sec. III introduces the coded matrix multiplication, and Sec. IV introduces the coded shuffling algorithm. Finally, Sec. V presents conclusions and discusses open problems.

II Related Work

The straggler problem has been widely observed in distributed computing clusters. The authors of show that running a computational task at a computing node often involves unpredictable latency due to several factors such as network latency, shared resources, maintenance activities, and power limits. Further, they argue that stragglers cannot be completely removed from a distributed computing cluster. The authors of characterize the impact and causes of stragglers that arise due to resource contention, disk failures, varying network conditions, and imbalanced workload.

One approach to mitigate the adverse effect of stragglers is based on efficient straggler detection algorithms. For instance, the default scheduler of Hadoop constantly detects stragglers while running computational tasks. Whenever it detects a straggler, it relaunches the task that was running on the detected straggler at some other available node. In , Zaharia et al. propose a modification to the existing straggler detection algorithm and show that the proposed solution can effectively reduce the completion time of MapReduce tasks. In , Ananthanarayanan et al. propose a system that efficiently detects stragglers using real-time progress and cancels those stragglers, and show that the proposed system can further reduce the runtime of MapReduce tasks.

Another line of work is based on breaking the synchronization barriers in distributed algorithms . An asynchronous parallel execution can continuously make progress without having to wait for all the responses from the workers, and hence the overall runtime is less affected by stragglers. However, these asynchronous approaches break the serial consistency of the algorithm to be parallelized, and do not guarantee “correctness” of the end result, i.e., the output of the asynchronous algorithm can differ from that of a serial execution with an identical number of iterations.

Recently, replication-based approaches have been explored to tackle the straggler problem: by replicating tasks and scheduling the replicas, the runtime of distributed algorithms can be significantly improved . By collecting outputs of the fast-responding nodes (and potentially canceling all the other slow-responding replicas), such replication-based scheduling algorithms can reduce latency. In , the authors show that even without replica cancellation, one can still reduce the average task latency by properly scheduling redundant requests. We view these policies as special instances of coded computation: such task replication schemes can be seen as repetition-coded computation. In Sec. III, we describe this connection in detail, and indicate that coded computation can significantly outperform replication (as is usually the case for coding vs. replication in other engineering applications).

Another line of work that is closely related to coded computation is about the latency analysis of coded distributed storage systems. In , the authors show that the flexibility of erasure-coded distributed storage systems allows for faster data retrieval performance than replication-based distributed storage systems. Joshi et al. show that scheduling redundant requests to an increased number of storage nodes can improve the latency performance, and characterize the resulting storage-latency tradeoff. Sun et al. study the problem of adaptive redundant requests scheduling, and characterize the optimal strategies for various scenarios. In , Kadhe and Soljanin analyze the latency performance of availability codes, a class of storage codes designed for enhanced availability. In , the authors study the cost associated with scheduling of redundant requests, and propose a general scheduling policy that achieves a delicate balance between the latency performance and the cost.

We now review some recent works on coded computation, which have been published after our conference publications . In , an anytime coding scheme for approximate matrix multiplication is proposed, and it is shown that the proposed scheme can improve the quality of approximation compared with the other existing coded schemes for exact computation. In , the authors propose a coded computation scheme called ‘Short-Dot’. Short-Dot induces additional sparsity to the encoded matrices at the cost of reduced decoding flexibility, and hence potentially speeds up the computation. The authors of consider the problem of computing gradients in a distributed system, and propose a novel coded computation scheme tailored for computing a sum of functions. In many machine learning problems, the objective function is a sum of per-data loss functions, and hence the gradient of the objective function is the sum of gradients of per-data loss functions. Based on this observation, they propose Gradient Coding, which can reliably compute the exact gradient of any function in the presence of stragglers. While Gradient coding can be applied to computing gradients of any functions, it usually incurs significant storage and computation overheads. In , the authors consider a secure coded computation problem where the input data matrices need to be secured from the workers. They propose a secure computation scheme based on Staircase codes, which can speed up the distributed computation while securing the input data from the workers. In , the authors consider the problem of large matrix-matrix multiplication, and propose a new coded computation scheme based on product codes. In , the authors consider the coded computation problem on heterogenous computing clusters while our work assumes a homogeneous computing cluster. The authors show that by delicately distributing jobs across heterogenous workers, one can improve the performance of coded computation compared with the symmetric job allocation scheme, which is designed for homogeneous workers in our work. While most of the works focus on the application of coded computation to linear operations, a recent work shows that coding can be used also in distributed computing frameworks involving nonlinear operations . The authors of show that by leveraging the multi-core architecture in the worker computing units and “coding across” the multi-core computed outputs, significant (and in some settings unbounded) gains in speed-up in computational time can be achieved between the coded and uncoded schemes.

II-B Data Shuffling and Communication Overheads

Distributed learning algorithms on large-scale networked systems have been extensively studied in the literature . Many of the distributed algorithms that are implemented in practice share a similar algorithmic “anatomy”: the data set is split among several cores or nodes, each node trains a model locally, then the local models are averaged, and the process is repeated. While training a model with parallel or distributed learning algorithms, it is common to randomly re-shuffle the data a number of times . This essentially means that after each shuffling the learning algorithm will go over the data in a different order than before. Although the effects of random shuffling are far from understood theoretically, the large statistical gains have turned it into a common practice. Intuitively, data shuffling before a new pass over the data, implies that nodes get a nearly “fresh” sample from the data set, which experimentally leads to better statistical performance. Moreover, bad orderings of the data—known to lead to slow convergence in the worst case —are “averaged out”. However, the statistical benefits of data shuffling do not come for free: each time a new shuffle is performed, the entire dataset is communicated over the network of nodes. This inevitably leads to performance bottlenecks due to heavy communication.

In this work, we propose to use coding opportunities to significantly reduce the communication cost of some distributed learning algorithms that require data shuffling. Our coded shuffling algorithm is built upon the coded caching scheme by Maddah-Ali and Niesen . Coded caching is a technique to reduce the communication rate in content delivery networks. Mainly motivated by video sharing applications, coded caching exploits the multicasting opportunities between users that request different video files to significantly reduce the communication burden of the server node that has access to the files. Coded caching has been studied in many scenarios such as decentralized coded caching , online coded caching , hierarchical coded caching for wireless communication , and device-to-device coded caching . Recently, the authors in proposed coded MapReduce that reduces the communication cost in the process of transferring the results of mappers to reducers.

Our proposed approach is significantly different from all related studies on coded caching in two ways: (i) we shuffle the data points among the computing nodes to increase the statistical efficiency of distributed computation and machine learning algorithms; and (ii) we code the data over their actual representation (i.e., over the doubles or floats) unlike the traditional coding schemes over bits. In Sec. IV, we describe how coded shuffling can remarkably speed up the communication phase of large-scale parallel machine learning algorithms, and provide extensive numerical experiments to validate our results.

The coded shuffling problem that we study is related to the index coding problem . Indeed, given a fixed “side information” reflecting the memory content of the nodes, the data delivery strategy for a particular permutation of the data rows induces an index coding problem. However, our coded shuffling framework is different from index coding in at least two significant ways. First, the coded shuffling framework involves multiple iterations of data being stored across all the nodes. Secondly, when the caches of the nodes are updated in coded shuffling, the system is unaware of the upcoming permutations. Thus, the cache update rules need to be designed to target any possible unknown permutation of data in succeeding iterations of the algorithm.

We now review some recent works on coded shuffling, which have been published after our first presentation . In , the authors study the information-theoretic limits of the coded shuffling problem. More specifically, the authors completely characterize the fundamental limits for the case of 22 workers and the case of 33 workers. In , the authors consider the worse-case formulation of the coded shuffling problem, and propose a two-stage shuffling algorithm. The authors of propose a new coded shuffling scheme based on pliable index coding. While most of the existing works focus on either coded computation or coded shuffling, one notable exception is . In this work, the authors generalize the original coded MapReduce framework by introducing stragglers to the computation phases. Observing that highly flexible codes are not favorable to coded shuffling while replication codes allow for efficient shuffling, the authors propose an efficient way of coding to mitigate straggler effects as well as reduce the shuffling overheads.

III Coded Computation

In this section, we propose a novel paradigm to mitigate the straggler problem. The core idea is simple: we introduce redundancy into subtasks of a distributed algorithm such that the original task’s result can be decoded from a subset of the subtask results, treating uncompleted subtasks as erasures. For this specific purpose, we use erasure codes to design coded subtasks.

An erasure code is a method of introducing redundancy to information for robustness to noise . It encodes a message of kk symbols into a longer message of nn coded symbols such that the original kk message symbols can be recovered by decoding a subset of coded symbols . We now show how erasure codes can be applied to distributed computation to mitigate the straggler problem.

A coded distributed algorithm is specified by local functions, local data blocks, decodable sets of indices, and a decoding function: The local functions and data blocks specify the way the original computational task and the input data are distributed across nn workers; and the decodable sets of indices and the decoding function are such that the desired computation result can be correctly recovered using the decoding function as long as the local computation results from any of the decodable sets are collected.

The formal definition of coded distributed algorithms is as follows.

Consider a computational task fA(⋅)f_{\mathbf{A}}(\cdot). A coded distributed algorithm for computing fA(⋅)f_{\mathbf{A}}(\cdot) is specified by

local functions ⟨fAii(⋅)⟩i=1n\langle f^{i}_{\mathbf{A}_{i}}(\cdot)\rangle_{i=1}^{n} and local data blocks ⟨Ai⟩i=1n\langle\mathbf{A}_{i}\rangle_{i=1}^{n};

(minimal) decodable sets of indices I⊂P([n])\mathcal{I}\subset\mathcal{P}([n]) and a decoding function dec(⋅,⋅)\texttt{dec}(\cdot,\cdot),

where [n]=\makebox[0.0pt]\mboxdef{1,2,…,n}[n]\mathrel{\overset{\makebox[0.0pt]{\mbox{\tiny def}}}{=}}\{1,2,\ldots,n\}, and P(⋅)\mathcal{P}(\cdot) is the power set of a set. The decodable sets of indices I\mathcal{I} is minimal: no element of I\mathcal{I} is a subset of other elements. The decoding function takes a sequence of indices and a sequence of subtask results, and it must correctly output fA(x)f_{\mathbf{A}}(\mathbf{x}) if any decodable set of indices and its corresponding results are given.

A coded distributed algorithm can be run in a distributed computing cluster as follows. Assume that the ithi^{\text{\tiny th}} (encoded) data block Ai\mathbf{A}_{i} is stored at the ithi^{\text{\tiny th}} worker for all ii. Upon receiving the input argument x\mathbf{x}, the master node multicasts x\mathbf{x} to all the workers, and then waits until it receives the responses from any of the decodable sets. Each worker node starts computing its local function when it receives its local input argument, and sends the task result to the master node. Once the master node receives the results from some decodable set, it decodes the received task results and obtains fA(x)f_{\mathbf{A}}(\mathbf{x}).

The algorithm described in Sec. I-A is an example of coded distributed algorithms: it is a coded distributed algorithm for matrix multiplication that uses an (n,n−1)(n,n-1) MDS code. One can generalize the described algorithm using an (n,k)(n,k) MDS code as follows. For any 1≤k≤n1\leq k\leq n, the data matrix A\mathbf{A} is first divided into kk equal-sized submatricesIf the number of rows of A\mathbf{A} is not a multiple of kk, one can append zero rows to A\mathbf{A} to make the number of rows a multiple of kk.. Then, by applying an (n,k)(n,k) MDS code to each element of the submatrices, nn encoded submatrices are obtained. We denote these nn encoded submatrices by A1′,A2′,…,An′\mathbf{A}^{\prime}_{1},\mathbf{A}^{\prime}_{2},\ldots,\mathbf{A}^{\prime}_{n}. Note that the Ai′=Ai\mathbf{A}^{\prime}_{i}=\mathbf{A}_{i} for 1≤i≤k1\leq i\leq k if a systematic MDS code is used for the encoding procedure. Upon receiving any kk task results, the master node can use the decoding algorithm to decode kk task results. Then, one can find AX\mathbf{A}\mathbf{X} simply by concatenating them.

III-B Runtime of Uncoded/Coded Distributed Algorithms

In this section, we analyze the runtime of uncoded and coded distributed algorithms. We first consider the overall runtime of an uncoded distributed algorithm, ToveralluncodedT^{\text{uncoded}}_{\text{overall}}. Assume that the runtime of each task is identically distributed and independent of others. We denote the runtime of the ithi^{\text{\tiny th}} worker under a computation scheme, say s, by TisT^{\text{s}}_{i}. Note that the distributions of TiT_{i}’s can differ across different computation schemes.

where T(i)T_{(i)} is the ithi^{\text{\tiny th}} smallest one in {Ti}i=1n\{T_{i}\}_{i=1}^{n}. From (2), it is clear that a single straggler can slow down the overall algorithm. A coded distributed algorithm is terminated whenever the master node receives results from any decodable set of workers. Thus, the overall runtime of a coded algorithm is not determined by the slowest worker, but by the first time to collect results from some decodable set in I\mathcal{I}, i.e.,

We remark that the runtime of uncoded distributed algorithms (2) is a special case of (3) with I={[n]}\mathcal{I}=\{[n]\}. In the following examples, we consider the runtime of the repetition-coded algorithms and the MDS-coded algorithms.

Consider an nk\frac{n}{k}-repetition-code where each local task is replicated nk\frac{n}{k} times. We assume that each group of nk\frac{n}{k} consecutive workers work on the replicas of one local task. Thus, the decodable sets of indices I\mathcal{I} are all the minimal sets that have kk distinct task results, i.e., I={1,2,…,nk}×{nk+1,nk+2,…,nk+k}×…×{n−nk+1,n−nk+2,…,n}\mathcal{I}=\{1,2,\ldots,\frac{n}{k}\}\times\{\frac{n}{k}+1,\frac{n}{k}+2,\ldots,\frac{n}{k}+k\}\times\ldots\times\{n-\frac{n}{k}+1,n-\frac{n}{k}+2,\ldots,n\}, where A×BA\times B denotes the Cartesian product of matrix AA and BB. Thus,

If one uses an (n,k)(n,k) MDS code, the decodable sets of indices are the sets of any kk indices, i.e., I={i∣i⊂[n], ∣i∣=k}\mathcal{I}=\{\mathbf{i}|\mathbf{i}\subset[n],~{}|\mathbf{i}|=k\}. Thus,

That is, the algorithm’s runtime will be determined by the kthk^{\text{\tiny th}} response, not by the nthn^{\text{\tiny th}} response.

III-C Probabilistic Model of Runtime

In this work, we assume homogeneous clusters: that is, all the workers have independent and identically distributed computing time statistics. While our symmetric job allocation is optimal for homogeneous cases, it can be strictly suboptimal for heterogenous cases. While our work focuses on homogeneous clusters, we refer the interested reader to a recent work for a generalization of our problem setting to that of heterogeneous clusters, for which symmetric allocation strategies are no longer optimal.

We first consider an uncoded distributed algorithm with nn (uncoded) subtasks. Due to the assumptions mentioned above, the runtime of each subtask is F(nt)F(nt). Thus, the runtime distribution of an uncoded distributed algorithm, denoted by Foveralluncoded(t)F_{\text{overall}}^{\text{uncoded}}(t), is simply [F(nt)]n\left[F(nt)\right]^{n}.

When repetition codes or MDS codes are used, an algorithm is first divided into k (<n)k~{}(<n) systematic subtasks, and then n−kn-k coded tasks are designed to provide an appropriate level of redundancy. Thus, the runtime of each task is distributed according to F(kt)F(kt). Using (4) and (5), one can easily find the runtime distribution of an nk\frac{n}{k}-repetition-coded distributed algorithm, FoverallRepetitionF_{\text{overall}}^{\text{Repetition}}, and the runtime distribution of an (n,k)(n,k)-MDS-coded distributed algorithm, FoverallMDS-codedF_{\text{overall}}^{\text{MDS-coded}}. For an nk\frac{n}{k}-repetition-coded distributed algorithm, one can first find the distribution of

and then find the distribution of the maximum of kk such terms:

The runtime distribution of an (n,k)(n,k)-MDS-coded distributed algorithm is simply the kthk^{\text{\tiny th}} order statistic:

For the same values of nn and kk, the runtime distribution of a repetition-coded distributed algorithm strictly dominates that of an MDS-coded distributed algorithm. This can be shown by observing that the decodable sets of the MDS-coded algorithm contain those of the repetition-coded algorithm.

In Fig. 4, we compare the runtime distributions of uncoded and coded distributed algorithms. We compare the runtime distributions of uncoded algorithm, repetition-coded algorithm, and MDS-coded algorithm with n=10n=10 and k=5k=5. In Fig. 4(a), we use a shifted-exponential distribution as the mother runtime distribution. That is, F(t)=1−et−1F(t)=1-e^{t-1} for t≥1t\geq 1. In Fig. 4(b), we use the empirical task runtime distribution that is measured on an Amazon EC2 clusterThe detailed description of the experiments is provided in Sec. III-F.. Observe that for both cases, the runtime distribution of the MDS-coded distribution has the lightest tail.

III-D Optimal Code Design for Coded Distributed Algorithms: The Shifted-exponential Case

When a coded distributed algorithm is used, the original task is divided into a fewer number of tasks compared to the case of uncoded algorithms. Thus, the runtime of each task of a coded algorithm, which is F(kt)F(kt), is stochastically larger than that of an uncoded algorithm, which is F(nt)F(nt). If the value that we choose for kk is too small, then the runtime of each task becomes so large that the overall runtime of the distributed coded algorithm will eventually increase. If kk is too large, the level of redundancy may not be sufficient to prevent the algorithm from being delayed by the stragglers.

Given the mother runtime distribution and the code parameters, one can compute the overall runtime distribution of the coded distributed algorithm using (6) and (III-C). Then, one can optimize the design based on various target metrics, e.g., the expected overall runtime, the 99th99^{\text{\tiny th}} percentile runtime, etc.

In this section, we show how one can design an optimal coded algorithm that minimizes the expected overall runtime for a shifted-exponential mother distribution. The shifted-exponential distribution strikes a good balance between accuracy and analytical tractability. This model is motivated by the model proposed in : the authors used this distribution to model latency of file queries from cloud storage systems. The shifted-exponential distribution is the sum of a constant and an exponential random variable, i.e.,

where the exponential rate μ\mu is called the straggling parameter.

With this shifted-exponential model, we first characterize a lower bound on the fundamental limit of the average runtime.

The average runtime of any distributed algorithm, in a distributed computing cluster with nn workers, is lower bounded by 1n\frac{1}{n}.

One can show that the average runtime of any distributed algorithm strictly decreases if the mother runtime distribution is replaced with a deterministic constant 11. Thus, the optimal average runtime with this deterministic mother distribution serves as a strict lower bound on the optimal average runtime with the shifted-exponential mother distribution. The constant mother distribution implies that stragglers do not exist, and hence the uncoded distributed algorithm achieves the optimal runtime, which is 1n\frac{1}{n}. ∎

We now analyze the average runtime of uncoded/coded distributed algorithms. We assume that nn is large, and kk is linear in nn. Accordingly, we approximate Hn=\makebox[0.0pt]\mboxdef∑i=1n1i≃log⁡nH_{n}\mathrel{\overset{\makebox[0.0pt]{\mbox{\tiny def}}}{=}}\sum_{i=1}^{n}{\frac{1}{i}}\simeq\log n and Hn−k≃log⁡(n−k)H_{n-k}\simeq\log{(n-k)}. We first note that the expected value of the maximum of nn independent exponential random variables with rate μ\mu is Hnμ\frac{H_{n}}{\mu}. Thus, the average runtime of an uncoded distributed algorithm is

For the average runtime of an nk\frac{n}{k}-Repetition-coded distributed algorithm, we first note that the minimum of nk\frac{n}{k} independent exponential random variables with rate μ\mu is distributed as an exponential random variable with rate nkμ\frac{n}{k}\mu. Thus,

Finally, we note that the expected value of the kthk^{\text{\tiny th}} statistic of nn independent exponential random variables of rate μ\mu is Hn−Hn−kμ\frac{H_{n}-H_{n-k}}{\mu}. Therefore,

Using these closed-form expressions of the average runtime, one can easily find the optimal value of kk that achieves the optimal average runtime. The following lemma characterizes the optimal repetition code for the repetition-coded algorithms and their runtime performances.

If μ≥1\mu\geq 1, the average runtime of an nk\frac{n}{k}-Repetition-coded distributed algorithm, in a distributed computing cluster with nn workers, is minimized by setting k=nk=n, i.e., not replicating tasks. If μ=1v\mu=\frac{1}{v} for some integer v>1v>1, the average runtime is minimized by setting k=μnk=\mu n, and the corresponding minimum average runtime is 1nμ(1+log⁡(nμ))\frac{1}{n\mu}\left(1+\log(n\mu)\right).

It is easy to see that (10) as a function of kk has a unique extreme point. By differentiating (10) with respect to kk and equating it to zero, we have k=μnk=\mu n. Thus, if μ≥1\mu\geq 1, one should set k=nk=n; if μ=1v<1\mu=\frac{1}{v}<1 for some integer vv, one should set k=μnk=\mu n. ∎

The above lemma reveals that the optimal repetition-coded distributed algorithm can achieve a lower average runtime than the uncoded distributed algorithm if μ<1\mu<1; however, the optimal repetition-coded distributed algorithm still suffers from the factor of Θ(log⁡n)\Theta(\log n), and cannot achieve the order-optimal performance. The following lemma, on the other hand, shows that the optimal MDS-coded distributed algorithm can achieve the order-optimal average runtime performance.

The average runtime of an (n,k)(n,k)-MDS-coded distributed algorithm, in a distributed computing cluster with nn workers, can be minimized by setting k=k⋆k=k^{\star} where

and W−1(⋅)W_{-1}(\cdot) is the lower branch of Lambert W functionW−1(x)W_{-1}(x), the lower branch of Lambert W function evaluated at xx, is the unique solution of tet=xte^{t}=x and t≤−1t\leq-1. Thus,

It is easy to see that (11) as a function of kk has a unique extreme point. By differentiating (11) with respect to kk and equating it to zero, we have 1k⋆(1+1μlog⁡(nn−k⋆))=1μ1n−k⋆\frac{1}{k^{\star}}\left(1+\frac{1}{\mu}\log\left(\frac{n}{n-k^{\star}}\right)\right)=\frac{1}{\mu}\frac{1}{n-k^{\star}}. By setting k=α⋆nk=\alpha^{\star}n, we have 1α⋆(1+1μlog⁡(11−α⋆))=1μ11−α⋆\frac{1}{\alpha^{\star}}\left(1+\frac{1}{\mu}\log\left(\frac{1}{1-\alpha^{\star}}\right)\right)=\frac{1}{\mu}\frac{1}{1-\alpha^{\star}}, which implies μ+1=11−α⋆−log⁡(11−α⋆)\mu+1=\frac{1}{1-\alpha^{\star}}-\log\left(\frac{1}{1-\alpha^{\star}}\right). By defining β=11−α⋆\beta=\frac{1}{1-\alpha^{\star}} and exponentiating both the sides, we have eμ+1=eββe^{\mu+1}=\frac{e^{\beta}}{\beta}. Note that the solution of exx=t\frac{e^{x}}{x}=t, t≥et\geq e and x≥1x\geq 1 is x=−W−1(−1t)x=-W_{-1}(-\frac{1}{t}). Thus, β=−W−1(−e−μ−1)\beta=-W_{-1}(-e^{-\mu-1}). By plugging the above equation into the definition of β\beta, the claim is proved. ∎

We plot nT⋆nT^{\star} and k⋆μ\frac{k^{\star}}{\mu} as functions of μ\mu in Fig. 5.

In addition to the order-optimality of MDS-coded distributed algorithms, the above lemma precisely characterizes the gap between the achievable runtime and the optimistic lower bound of 1n\frac{1}{n}. For instance, when μ>1\mu>1, the optimal average runtime is only 3.153.15 away from the lower bound.

So far, we have considered only the runtime performance of distributed algorithms. Another important metric to be considered is the storage cost. When coded computation is being used, the storage overhead may increase. For instance, the MDS-coded distributed algorithm for matrix multiplication, described in Sec. III-A, requires 1k\frac{1}{k} of the whole data to be stored at each worker, while the uncoded distributed algorithm requires 1n\frac{1}{n}. Thus, the storage overhead factor is 1k−1n1n=nk−1\frac{\frac{1}{k}-\frac{1}{n}}{\frac{1}{n}}=\frac{n}{k}-1. If one uses the runtime-optimal MDS-coded distributed algorithm for matrix multiplication, the storage overhead is nk⋆−1=1α⋆−1\frac{n}{k^{\star}}-1=\frac{1}{\alpha^{\star}}-1.

III-E Coded Gradient Descent: An MDS-coded Distributed Algorithm for Linear Regression

In this section, as a concrete application of coded matrix multiplication, we propose the coded gradient descent for solving large-scale linear regression problems.

We first describe the (uncoded) gradient-based distributed algorithm. Consider the following linear regression,

The above algorithm is guaranteed to converge to the optimal solution if we use a small enough step size η\eta , and can be easily distributed. We describe one simple way of parallelizing the algorithm, which is implemented in many open-source machine learning libraries including Spark mllib . As AT(Ax(t)−y)=∑i=1qai(aiTx(t)−yi)\mathbf{A}^{T}(\mathbf{A}\mathbf{x}^{(t)}-\mathbf{y})=\sum_{i=1}^{q}{\mathbf{a}_{i}(\mathbf{a}_{i}^{T}\mathbf{x}^{(t)}-\mathbf{y}_{i})}, gradients can be computed in a distributed way by computing partial sums at different worker nodes and then adding all the partial sums at the master node. This distributed algorithm is an uncoded distributed algorithm: in each round, the master node needs to wait for all the task results in order to compute the gradient.Indeed, one may apply another coded computation scheme called Gradient Coding , which was proposed after our conference publications. By applying Gradient Coding to this algorithm, one can achieve straggler tolerance but at the cost of significant computation and storage overheads. More precisely, it incurs Θ(n)\Theta(n) larger computation and storage overheads in order to protect the algorithm from Θ(n)\Theta(n) stragglers. Later in this section, we will show that our coded computation scheme, which is tailor-designed for linear regression, incurs Θ(1)\Theta(1) overheads to protect the algorithm from Θ(n)\Theta(n) stragglers. Thus, the runtime of each update iteration is determined by the slowest response among all the worker nodes.

We now propose the coded gradient descent, a coded distributed algorithm for linear regression problems. Note that in each iteration, the following two matrix-vector multiplications are computed.

In Sec. III-A, we proposed the MDS-coded distributed algorithm for matrix multiplication. Here, we apply the algorithm twice to compute these two multiplications in each iteration. More specifically, for the first matrix multiplication, we choose 1≤k1<n1\leq k_{1}<n and use an (n,k1)(n,k_{1})-MDS-coded distributed algorithm for matrix multiplication to encode the data matrix A\mathbf{A}. Similarly for the second matrix multiplication, we choose 1≤k2<n1\leq k_{2}<n and use a (n,k2)(n,k_{2})-MDS-coded distributed algorithm to encode the transpose of the data matrix. Denoting the ithi^{\text{\tiny th}} row-split (column-split) of A\mathbf{A} as Ai\mathbf{A}_{i} (A~i\widetilde{\mathbf{A}}_{i}), the ithi^{\text{\tiny th}} worker stores both Ai\mathbf{A}_{i} and A~i\widetilde{\mathbf{A}}_{i}. In the beginning of each iteration, the master node multicasts x(t)\mathbf{x}^{(t)} to the worker nodes, each of which computes the local matrix multiplication for Ax(t)\mathbf{A}\mathbf{x}^{(t)} and sends the result to the master node. Upon receiving any k1k_{1} task results, the master node can start decoding the result and obtain z(t)=Ax(t)\mathbf{z}^{(t)}=\mathbf{A}\mathbf{x}^{(t)}. The master node now multicasts z(t)\mathbf{z}^{(t)} to the workers, and the workers compute local matrix multiplication for ATz(t)\mathbf{A}^{T}\mathbf{z}^{(t)}. Finally, the master node can decode ATz(t)\mathbf{A}^{T}\mathbf{z}^{(t)} as soon as it receives any k2k_{2} task results, and can proceed to the next iteration. Fig. 6 illustrates the protocol with k1=k2=n−1k_{1}=k_{2}=n-1.

The coded gradient descent requires each node to store a (1k1+1k2−1k1k2)(\frac{1}{k_{1}}+\frac{1}{k_{2}}-\frac{1}{k_{1}k_{2}})-fraction of the data matrix. As the minimum storage overhead per node is a 1n\frac{1}{n}-fraction of the data matrix, the relative storage overhead of the coded gradient descent algorithm is at least about factor of 22, if k1≃nk_{1}\simeq n and k2≃nk_{2}\simeq n.

III-F Experimental Results

In order to see the efficacy of coded computation, we implement the proposed algorithms and test them on an Amazon EC2 cluster. We first obtain the empirical distribution of task runtime in order to observe how frequently stragglers appear in our testbed by measuring round-trip times between the master node and each of 1010 worker instances on an Amazon EC2 cluster. Each worker computes a matrix-vector multiplication and passes the computation result to the master node, and the master node measures round trip times that include both computation time and communication time. Each worker repeats this procedure 500500 times, and we obtain the empirical distribution of round trip times across all the worker nodes.

In Fig. 7, we plot the histogram and complementary CDF (CCDF) of measured computing times; the average round trip time is 0.110.11 second, and the 95th95^{\text{\tiny th}} percentile latency is 0.200.20 second, i.e., roughly five out of hundred tasks are going to be roughly two times slower than the average tasks. Assuming the probability of a worker being a straggler is 5%5\%, if one runs an uncoded distributed algorithm with 1010 workers, the probability of not seeing such a straggler is only about 60%60\%, so the algorithm is slowed down by a factor of more than 22 with probability 40%40\%. Thus, this observation strongly emphasizes the necessity of an efficient straggler mitigation algorithm. In Fig. 4(a), we plot the runtime distributions of uncoded/coded distributed algorithms using this empirical distribution as the mother runtime distribution. When an uncoded distributed algorithm is used, the overall runtime distribution entails a heavy tail, while the runtime distribution of the MDS-coded algorithm has almost no tail.

We then implement the coded matrix multiplication in C++ using OpenMPI​ and benchmark on a cluster of 2626 EC2 instances (2525 workers and a master)For the benchmark, we manage the cluster using the StarCluster toolkit . Input data is generated using a Python script, and the input matrix is row-partitioned for each of the workers (with the required encoding as described in the previous sections) in a preprocessing step. The procedure begins by having all of the worker nodes read in their respective row-partitioned matrices. Then, the master node reads the input vector and distributes it to all worker nodes in the cluster through an asynchronous send (MPI_Isend). Upon receiving the input vector, each worker node begins matrix multiplication through a BLAS routine call and once completed sends the result back to the master using MPI_Send. The master node waits for a sufficient number of results to be received by continuously polling (MPI_Test) to see if any results are obtained. The procedure ends when the master node decodes the overall result after receiving enough partial results.. Also, three uncoded matrix multiplication algorithms – block, column-partition, and row-partition – are implemented and benchmarked.

We randomly draw a square matrix of size 5750×57505750\times 5750, a fat matrix of size 5750×115005750\times 11500, and a tall matrix of size 11500×575011500\times 5750, and multiply them with a column vector. For the coded matrix multiplication, we choose an (25,23)(25,23) MDS code so that the runtime of the algorithm is not affected by any 22 stragglers. Fig. 8 shows that the coded matrix multiplication outperforms all the other parallel matrix multiplication algorithms in most cases. On a cluster of m1-small, the most unreliable instances, the coded matrix multiplication achieves about 40%40\% average runtime reduction and about 60%60\% tail reduction compared to the best of the 33 uncoded matrix multiplication algorithmss. On a cluster of c1-medium instances, the coded algorithm achieves the best performance in most of the tested cases: the average runtime is reduced by at most 39.5%39.5\%, and the 95th95^{\text{\tiny th}} percentile runtime is reduced by at most 58.3%58.3\%. Among the tested cases, we observe one case in which both the uncoded row-partition and the coded row-partition algorithms are outperformed by the uncoded column-partition algorithm. This is the case of a fat matrix multiplication with c1-medium instances. Note that when a row-partition algorithm is used, the size of messages from the master node to the workers is nn times larger compared with the case of column-partition algorithms. Thus, when the variability of computational times becomes low compared with that of communication time, the larger communication overhead of row-partition algorithms seems to arise, nullifying the benefits of coding.

We also evaluate the performance of the coded gradient descent algorithm for linear regression. The coded linear regression procedure is also implemented in C++ using OpenMPI, and benchmarked on a cluster of 1111 EC2 machines (1010 workers and a master). Similar to the previous benchmarks, we randomly draw a square matrix of size 2000×20002000\times 2000, a fat matrix of size 400×10000400\times 10000, and a tall matrix of size 10000×40010000\times 400, and use them as a data matrix. We use a (10,8)(10,8)-MDS code for the coded linear regression so that each multiplication of the gradient descent algorithm is not slowed down by up to 22 stragglers. Fig. 9 shows that the gradient algorithm with the coded matrix multiplication significantly outperforms the one with the uncoded matrix multiplication; the average runtime is reduced by 31.3%31.3\% to 35.7%35.7\%, and the tail runtime is reduced by 27.9%27.9\% to 35.6%35.6\%.

IV Coded Shuffling

We shift our focus from solving the straggler problem to solving the communication bottleneck problem. In this section, we explain the problem of data-shuffling, propose the Coded Shuffling algorithm, and analyze its performance.

We consider a master-worker distributed setup, where the master node has access to the entire data-set. Before every iteration of the distributed algorithm, the master node randomly partition the entire data set into nn subsets, say A1,A2,…,An\mathbf{A}_{1},\mathbf{A}_{2},\ldots,\mathbf{A}_{n}. The goal of the shuffling phase is to distribute each of these partitioned data sets to the corresponding worker so that each worker can perform its distributed task with its own exclusive data set after the shuffling phase.

IV-B Shuffling Schemes

We now present our coded shuffling algorithm, consisting of a transmission strategy for the master node, and caching and decoding strategies for the worker nodes. Let CitC_{i}^{t} be the cache content of node ii (set of row indices stored in cache ii) at the end of iteration tt. We design a transmission algorithm (by the master node) and a cache update algorithm to ensure that (i) Sit⊂CitS_{i}^{t}\subset C_{i}^{t}; and (ii) Cit∖SitC_{i}^{t}\setminus S_{i}^{t} is distributed uniformly at random without replacement in the set [q]∖Sit[q]\setminus S_{i}^{t}. The first condition ensures that at each iteration, the workers have access to the data set that they are supposed to work on. The second condition provides the opportunity of effective coded transmissions for shuffling in the next iteration as will be explained later.

We consider the following cache update rule: the new cache will contain the subset of the data points used in the current iteration (this is needed for the local computations), plus a random subset of the previous cached contents. More specifically, q/nq/n rows of the new cache are precisely the rows in Sit+1S_{i}^{t+1}, and s−q/ns-q/n rows of the cache are sampled points from the set Cit∖Sit+1C_{i}^{t}\setminus S_{i}^{t+1}, uniformly at random without replacement. Since the permutation πt\pi^{t} is picked uniformly at random, the marginal distribution of the cache contents at iteration t+1t+1 given Sit+1, 1≤i≤nS^{t+1}_{i},~{}1\leq i\leq n is described as follows: Sit+1⊂Cit+1S^{t+1}_{i}\subset C^{t+1}_{i} and Cit+1∖Sit+1C^{t+1}_{i}\setminus S^{t+1}_{i} is distributed uniformly at random in [q]∖Sit+1[q]\setminus S^{t+1}_{i} without replacement.

IV-B2 Encoding and Transmission Schemes

We now formally describe two transmission schemes of the master node: (1) uncoded transmission and (2) coded transmission. In the following descriptions, we drop the iteration index tt (and t+1t+1) for the ease of notation.

The uncoded transmission first finds how many data rows in SiS_{i} are already cached in CiC_{i}, i.e. ∣Ci∩Si∣|C_{i}\cap S_{i}|. Since, the new permutation (partitioning) is picked uniformly at random, s/qs/q fraction of the data row indices in SiS_{i} are cached in CiC_{i}, so as qq gets large, we have ∣Ci∩Si∣+o(q)=qn(1−s/q)|C_{i}\cap S_{i}|+o(q)=\frac{q}{n}(1-s/q). Thus, without coding, the master node needs to transmit qn(1−s/q)\frac{q}{n}(1-s/q) data points to each of the nn worker nodes. The total communication rate (in data points transmitted per iteration) of the uncoded scheme is then

We now describe the coded transmission scheme. Define the set of “exclusive” cache content as C~I=(∩i∈ICi)∩(∩i′∈[n]∖ICi′∁)\widetilde{C}_{\mathcal{I}}=\left(\cap_{i\in\mathcal{I}}C_{i}\right)\cap\left(\cap_{i^{\prime}\in[n]\setminus\mathcal{I}}C^{\complement}_{i^{\prime}}\right) that denotes the set of rows that are stored at the caches of I\mathcal{I}, and are not stored at the caches of [n]∖I[n]\setminus\mathcal{I}. For each subset I\mathcal{I} with ∣I∣≥2|\mathcal{I}|\geq 2, the master node will multicast ∑i∈IA(Si∩C~I∖{i})\sum_{i\in\mathcal{I}}\mathbf{A}(S_{i}\cap\widetilde{C}_{\mathcal{I}\setminus\{i\}}) to the worker nodes. Note that in general, the matrices A\mathbf{A}’s differ in their sizes, so one has to zero-pad the shorter matrices and sum the zero-padded matrices. Algorithm 1 provides the pseudocode of the coded encoding and transmission scheme.Note that for each encoded data row, the master node also needs to transmit tiny metadata describing which data rows are included in the summation. We omit this detail in the description of the algorithm.

IV-B3 Decoding Algorithm

The decoding algorithm for the uncoded transmission scheme is straightforward: each worker simply takes the additional data rows that are required for the new iteration, and ignores the other data rows. We now describe the decoding algorithm for the coded transmission scheme. Each worker, say worker ii, decodes each encoded data row as follows. Consider an encoded data row for some I\mathcal{I} that contains ii. (All other data rows are discarded.) Such an encoded data row must be the sum of some data row in SiS_{i} and ∣I∣−1|\mathcal{I}|-1 data rows in C~I∖{i}\widetilde{C}_{\mathcal{I}\setminus\{i\}}, which are available in worker ii by the definition of C~\widetilde{C}. Hence, the worker can always subtract the data rows corresponding to C~I∖{i}\widetilde{C}_{\mathcal{I}\setminus\{i\}} and decode the data row in SiS_{i}.

IV-C Example

The following example illustrates the coded shuffling scheme.

Let n=3n=3. Recall that worker node ii needs to obtain A(Si∩Ci∁)\mathbf{A}(S_{i}\cap C^{\complement}_{i}) for the next iteration of the algorithm. Consider i=1i=1. The data rows in S1∩C1∁S_{1}\cap C^{\complement}_{1} are stored either exclusively in C2C_{2} or C3C_{3} (i.e. C~2\widetilde{C}_{2} or C~3\widetilde{C}_{3}), or stored in both C2C_{2} and C3C_{3} (i.e. C~2,3\widetilde{C}_{2,3}). The transmitted message consists of 4 parts:

(Part 11) M{1,2}=A(S1∩C~2)+A(S2∩C~1){M}_{\{1,2\}}=\mathbf{A}(S_{1}\cap\widetilde{C}_{2})+\mathbf{A}(S_{2}\cap\widetilde{C}_{1}),

(Part 22) M{1,3}=A(S1∩C~3)+A(S3∩C~1){M}_{\{1,3\}}=\mathbf{A}(S_{1}\cap\widetilde{C}_{3})+\mathbf{A}(S_{3}\cap\widetilde{C}_{1}),

(Part 33) M{2,3}=A(S2∩C~3)+A(S3∩C~2){M}_{\{2,3\}}=\mathbf{A}(S_{2}\cap\widetilde{C}_{3})+\mathbf{A}(S_{3}\cap\widetilde{C}_{2}), and

(Part 44) M{1,2,3}=A(S1∩C~2,3)+A(S2∩C~1,3)+A(S3∩C~1,2){M}_{\{1,2,3\}}=\mathbf{A}(S_{1}\cap\widetilde{C}_{2,3})+\mathbf{A}(S_{2}\cap\widetilde{C}_{1,3})+\mathbf{A}(S_{3}\cap\widetilde{C}_{1,2}).

We show that worker node 1 can recover the data rows that it does not store or A(S1∩C1∁)\mathbf{A}(S_{1}\cap C^{\complement}_{1}). First, observe that node 11 stores S2∩C~1S_{2}\cap\widetilde{C}_{1}. Thus, it can recover A(S1∩C~2)\mathbf{A}(S_{1}\cap\widetilde{C}_{2}) using part 1 of the message since A(S1∩C~2)=M1−A(S2∩C~1)\mathbf{A}(S_{1}\cap\widetilde{C}_{2})={M}_{1}-\mathbf{A}(S_{2}\cap\widetilde{C}_{1}). Similarly, node 11 recovers A(S1∩C~3)=M2−A(S3∩C~1)\mathbf{A}(S_{1}\cap\widetilde{C}_{3})={M}_{2}-\mathbf{A}(S_{3}\cap\widetilde{C}_{1}). Finally, from part 4 of the message, node 11 recovers A(S1∩C~2,3)=M4−A(S2∩C~1,3)−A(S3∩C~1,2)\mathbf{A}(S_{1}\cap\widetilde{C}_{2,3})={M}_{4}-\mathbf{A}(S_{2}\cap\widetilde{C}_{1,3})-\mathbf{A}(S_{3}\cap\widetilde{C}_{1,2}).

IV-D Main Results

We now present the main result of this section, which characterizes the communication rate of the coded scheme. Let p=s−q/nq−q/np=\frac{s-q/n}{q-q/n}.

Coded shuffling achieves communication rate

(in number of data rows transmitted per iteration from the master node), which is significantly smaller than RuR_{u} in (17).

The reduction in communication rate is illustrated in Fig. 10 for n=50n=50 and q=1000q=1000 as a function of s/qs/q, where 1/n≤s/q≤11/n\leq s/q\leq 1.

For instance, when s/q=0.1s/q=0.1, the communication overhead for data-shuffling is reduced by more than 81%81\%. Thus, at a very low storage overhead for caching, the algorithm can be significantly accelerated.

Before we present the proof of the theorem, we briefly compare our main result with similar results shown in . Our coded shuffling algorithm is related to the coded caching problem , since one can design the right cache update rule to reduce the communication rate for an unknown demand or permutation of the data rows. A key difference though is that the coded shuffling algorithm is run over many iterations of the machine learning algorithm. Thus, the right cache update rule is required to guarantee the opportunity of coded transmission at every iteration. Furthermore, the coded shuffling problem has some connections to coded MapReduce as both algorithms mitigate the communication bottlenecks in distributed computation and machine learning. However, coded shuffling enables coded transmission of raw data by leveraging the extra memory space available at each node, while coded MapReduce enables coded transmission of processed data in the shuffling phase of the MapReduce algorithm by cleverly introducing redundancy in the computation of the mappers.

To find the transmission rate of the coded scheme we first need to find the cardinality of sets Sit+1∩C~ItS^{t+1}_{i}\cap\widetilde{C}^{t}_{\mathcal{I}} for I⊂[n]\mathcal{I}\subset[n] and i∉Ii\notin\mathcal{I}. To this end, we first find the probability that a random data row, r\mathbf{r}, belongs to C~It\widetilde{C}^{t}_{\mathcal{I}}. Denote this probability by Pr⁡(r∈C~It)\Pr(\mathbf{r}\in\widetilde{C}^{t}_{\mathcal{I}}). Recall that the cache content distribution at iteration tt: q/nq/n rows of cache jj are stored with SjtS^{t}_{j} and the other s−q/ns-q/n rows are stored uniformly at random. Thus, we can compute Pr⁡(r∈C~It)\Pr(\mathbf{r}\in\widetilde{C}^{t}_{\mathcal{I}}) as follows.

(19) is by the law of total probability. (20) is by the fact that r\mathbf{r} is chosen randomly. To see (21), note that Pr⁡(r∈C~It∣r∈Sit,i∉I)=0\Pr(\mathbf{r}\in\widetilde{C}^{t}_{\mathcal{I}}|\mathbf{r}\in S_{i}^{t},i\notin\mathcal{I})=0. Thus, the summation can be written only on the indices of I\mathcal{I}. We now explain (22). Given that r\mathbf{r} belongs to SitS_{i}^{t}, and i∈Ii\in\mathcal{I}, then r∈Ci\mathbf{r}\in C_{i} with probability 1. The other ∣I∣−1|\mathcal{I}|-1 caches with indices in I∖{i}\mathcal{I}\setminus\{i\} contain r\mathbf{r} with probability s−q/nq−q/n\frac{s-q/n}{q-q/n} independently. Further, the caches with indices in [n]∖I[n]\setminus\mathcal{I} do not contain r\mathbf{r} with probability 1−s−q/nq−q/n1-\frac{s-q/n}{q-q/n}. By defining p=\makebox[0.0pt]\mboxdefs−q/nq−q/np\mathrel{\overset{\makebox[0.0pt]{\mbox{\tiny def}}}{=}}\frac{s-q/n}{q-q/n}, we have (23).

We now find the cardinality of Sit+1∩C~ItS^{t+1}_{i}\cap\widetilde{C}^{t}_{\mathcal{I}} for I⊂[n]\mathcal{I}\subset[n] and i∉Ii\notin\mathcal{I}. Note that ∣Sit+1∣=q/n|S^{t+1}_{i}|=q/n. Thus, as qq gets large (and nn remains sub-linear in qq), by the law of large numbers,

Recall that for each subset I\mathcal{I} with ∣I∣≥2|\mathcal{I}|\geq 2, the master node will send ∑i∈IA(Si∩C~I∖{i})\sum_{i\in\mathcal{I}}\mathbf{A}(S_{i}\cap\widetilde{C}_{\mathcal{I}\setminus\{i\}}) . Thus, the total rate of coded transmission is

To complete the proof, we simplify the above expression. Let x=p1−px=\frac{p}{1-p}. Taking derivative with respect to xx from both sides of the equality ∑i=1n(ni)xi−1=1x[(1+x)n−1]\sum_{i=1}^{n}{n\choose i}x^{i-1}=\frac{1}{x}\left[(1+x)^{n}-1\right], we have

Using (26) in (25) completes the proof. ∎

Consider the case that the cache sizes are just enough to store the data required for processing; that is s=q/ns=q/n. Then, Rc=12RuR_{c}=\frac{1}{2}R_{u}. Thus, one gets a factor 2 reduction gain in communication rate by exploiting coded caching.

Note that when s=q/ns=q/n, p=0p=0. Finding the limit lim⁡p→0Rc\lim_{p\to 0}R_{c} in (18), after some manipulations, one calculates

Consider the regime of interest where nn, ss, and qq get large, and s/q→c>0s/q\to c>0 and n/q→0n/q\to 0. Then,

Thus, using coding, the communication rate is reduced by Θ(n)\Theta(n).

It is reasonable to assume that γ(n)≃n\gamma(n)\simeq n for wireless architecture that is of great interest with the emergence of wireless data centers, e.g. , and mobile computing platforms . However, still in many applications, the network topology is based on point-to-point communication, and the multicasting opportunity is not fully available, i.e., γ(n)<n\gamma(n)<n. For these general cases, we have to renormalize the communication cost of coded shuffling since we have assumed that γ(n)=n\gamma(n)=n in our results. For instance, in the regime considered in Corollary 8, the renormalized communication cost of coded shuffling RcγR^{\gamma}_{c} given γ(n)\gamma(n) is

Thus, the communication cost of coded shuffling is smaller than uncoded shuffling if γ(n)>q/s\gamma(n)>q/s. Note that s/qs/q is the fraction of the data matrix that can be stored at each worker’s cache. Thus, in the regime of interest where s/qs/q is a constant independent of nn, and γ(n)\gamma(n) scales with nn, the reduction gain of coded shuffling in communication cost is still unbounded and increasing in nn.

We emphasize that even in point-to-point communication networks, multicasting the same message to multiple nodes is significantly faster than unicasting different message (of the same size) to multiple nodes, i.e., γ(n)≫1\gamma(n)\gg 1, justifying the advantage of using coded shuffling. For instance, the MPI broadcast API (MPI_Bcast) utilizes a tree multicast algorithm, which achieves γ(n)=Θ(nlog⁡n)\gamma(n)=\Theta\left(\frac{n}{\log{n}}\right). Shown in Fig. 11 is the time taken for a data block to be transmitted to an increasing number of workers on an Amazon EC2 cluster, which consists of a point-to-point communication network. We compare the average transmission time taken with MPI scatter (unicast) and that with MPI broadcast. Observe that the average transmission time increases linearly as the number of receivers increases, but with MPI broadcast, the average transmission time increases logarithmically.

V Conclusion

In this paper, we have explored the power of coding in order to make distributed algorithms robust to a variety of sources of “system noise” such as stragglers and communication bottlenecks. We propose a novel Coded Computation framework that can significantly speed up existing distributed algorithms, by introducing redundancy through codes into the computation. Further, we propose Coded Shuffling that can significantly reduce the heavy price of data-shuffling, which is required for achieving high statistical efficiency in distributed machine learning algorithms. Our preliminary experimental results validate the power of our proposed schemes in effectively curtailing the negative effects of system bottlenecks, and attaining significant speedups of up to 40%40\%, compared to the current state-of-the-art methods.

There exists a whole host of theoretical and practical open problems related to the results of this paper. For coded computation, instead of the MDS codes, one could achieve different tradeoffs by employing another class of codes. Then, although matrix multiplication is one of the most basic computational blocks in many analytics, it would be interesting to leverage coding for a broader class of distributed algorithms.

For coded shuffling, convergence analysis of distributed machine learning algorithms under shuffling is not well understood. As we observed in the experiments, shuffling significantly reduces the number of iterations required to achieve a target reliability, but missing is a rigorous analysis that compares the convergence performances of algorithms with shuffling or without shuffling. Further, the trade-offs between bandwidth, storage, and the statistical efficiency of the distributed algorithms are not well understood. Moreover, it is not clear how far our achievable scheme, which achieves a bandwidth reduction gain of Θ(1n)\Theta(\frac{1}{n}), is from the fundamental limit of communication rate for coded shuffling. Therefore, finding an information-theoretic lower bound on the rate of coded shuffling is another interesting open problem.

References