Coded Computation over Heterogeneous Clusters

Amirhossein Reisizadeh, Saurav Prakash, Ramtin Pedarsani, Amir Salman Avestimehr

I Introduction

General distributed computing frameworks, such as MapReduce and Spark , along with the availability of large-scale commodity servers, such as Amazon EC2, have made it possible to carry out large-scale data analytics at the production level. These “virtualized data centers” enjoy an abundance of storage space and computing power, and are cheaper to rent by the hour than maintaining dedicated data centers round the year. However, these systems suffer from various forms of “system noise” which reduce their efficiency: system failures, limited communication bandwidth, straggler nodes, etc.

The current state-of-the-art approaches to mitigate the impact of system noise in cloud computing environments involve creation of some form of “computation redundancy”. For example, replicating the straggling task on another available node is a common approach to deal with stragglers , while partial data replication is also used to reduce the communication load in distributed computing . However, there have been recent results demonstrating that coding can play a transformational role for creating and exploiting computation redundancy to effectively alleviate the impact of system noise. In particular, there have been two coding concepts proposed to deal with the communication and straggler bottlenecks in distributed computing.

The first coding concept introduced in enables an inverse-linear tradeoff between computation load and communication load in distributed computing. This result implies that increasing the computation load by a factor of rr (i.e. evaluating each computation at rr carefully chosen nodes) can create novel coding opportunities that reduce the required communication load for computing by the same factor rr. Hence, these codes can be utilized to pool the underutilized computing resources at network edge to slash the communication load of Fog computing . Other related works tackling the communication bottleneck in distributed computation include .

In the second coding concept introduced in , an inverse-linear tradeoff between computation load and computation latency (i.e. the overall job response time) is established for distributed matrix multiplication in homogeneous computing environments. More specifically, this approach utilizes coding to effectively inject redundant computations to alleviate the effects of stragglers and speed up the computations. Hence, by utilizing more computation resources, this can significantly speed up distributed computing applications. A number of related works have been proposed recently to mitigate stragglers in distributed computation. In , the authors propose the use of redundant short dot products to speed up distributed computation of linear transforms. The work in proposes coding schemes for mitigating stragglers in distributed batch gradient computation. Coding schemes for high-dimensional matrix-matrix multiplication have been developed in . Techniques for efficient straggler mitigation for matrix-vector computation in distributed wireless settings have been developed in . In , the potential of the multicore nature of computing machines is studied. In , the authors propose an anytime approach to distributed computing, developing an approximate matrix multiplication scheme. The authors in propose a novel encoding scheme for achieving large sparsity in the encoded matrix. Work in develops a coding strategy for mitigating straggling decoders in cloud radio access network. Speeding up the computation of linear transformations with unreliable components is studied in . Straggler mitigation through data encoding in distributed optimization is proposed in . A coded scheme based on LT codes is proposed in for multiplying a matrix by a set of vectors in a distributed computing environment. Addressing stragglers has attracted a lot of attention in the queuing-based frameworks for large-scale computation as well . These works utilize the technique of dynamically replicating the tasks in a careful manner to minimize run-time.

We extend the problem of distributed matrix multiplication in homogeneous clusters in to heterogeneous environments. As discussed in , the computing environments in virtualized data centers are heterogeneous and algorithms based on homogeneous assumptions can result in significant performance reduction. In this paper, we focus on general heterogeneous distributed computing clusters consisting of a variety of computing machines with different capabilities. Specifically, we propose a coding framework for speeding up distributed matrix multiplication in heterogeneous clusters with straggling servers, named Heterogeneous Coded Matrix Multiplication (HCMM). Matrix multiplication is a crucial computation module in many engineering and scientific disciplines. In particular, it is a fundamental component of many popular machine learning algorithms such as logistic regression, reinforcement learning and gradient descent-based algorithms. Implementations that speed up matrix multiplication would naturally speed up the execution of a wide variety of popular algorithms. Therefore, we envision HCMM to play a fundamental role in speeding up big data analytics in virtualized data centers by leveraging the wide range of computing capabilities provided by these heterogeneous environments.

We now describe the main ideas behind HCMM, which results in asymptotically optimal performance. In a coded implementation of distributed matrix-vector multiplication, each worker node is assigned the task of computing inner products of the assigned coded rows with the input vector, where the assigned coded rows are random linear combinations of the rows of the original matrix. Computation time at each worker is a random variable, which is first assumed to have shifted exponential distribution, and we later generalize it to shifted Weibull distribution. The master node receives the results from the worker nodes and aggregates them until it receives a decodable set of inner products and recovers the matrix-vector multiplication. We are interested in finding the optimal load allocation that minimizes the expected waiting time to complete this computation. However, due to heterogeneity, finding the exact solution to the optimization problem seems intractable.

As the main contribution of the paper, we propose an alternative optimization that focuses on maximizing the expected number of returned computation results from the workers. Apart from being computationally tractable, the alternative optimization asymptotically approximates the problem of finding the optimal computation load allocation. Specifically, we develop the HCMM algorithm that is derived as a solution to the alternative formulation, and prove it is asymptotically optimal. Furthermore, we prove that given a heterogeneous cluster of nn workers, HCMM is Θ(log⁡n)\Theta(\log n) times faster than uncoded schemes under the shifted exponential distribution for run-time. We further generalize the proposed HCMM algorithm to shifted Weibull model and provide similar unbounded gains over uncoded scenarios.

In addition to proving the asymptotic optimality of HCMM, we carry out numerical studies and experiments over Amazon EC2 clusters to demonstrate how HCMM can be used in practice. We compare HCMM with three benchmark schemes – Uniform Uncoded, Load-balanced Uncoded, and Uniform Coded. In our numerical analysis, HCMM results in significant speedups of up to 73%73\%, 56%56\% and 42%42\% over the three aforementioned benchmark schemes, respectively. In experiments using Amazon EC2 clusters, we use the Luby transform (LT) codes for coding and demonstrate that HCMM combined with LT codes significantly reduces the overall execution time in comparison to uncoded and coded schemes. In particular, HCMM achieves gains of up to 61%61\%, 46%46\% and 36%36\%, respectively over Uniform Uncoded, Load-balanced Uncoded and Uniform Coded. Furthermore, the overall computation load of HCMM is less than the one of Uniform Coded. Our results demonstrate that HCMM combines the benefits of both Load-balanced Uncoded and Uniform Coded schemes by achieving efficient load balancing along with minimal number of redundant computations.

Furthermore, we consider the problem of load allocation under budget constraints, considering an intuitive and convincing pricing model. In particular, we show that HCMM is the (asymptotically) optimal load allocation in feasible budget-constrained scenarios as well, and determine whether a budget-constrained computation task is feasible given a cluster of machines. We then develop a heuristic algorithm to find the (sub)optimal load allocations using the proposed HCMM scheme. The heuristic is based on the observation that given a computation task and a set of machines, decreasing the number of fastest machines participating in HCMM results in smaller average cost.

II Problem Formulation and Main Results

In this section, we describe our computation model, the network model and the precise problem formulation. We then conclude with four theorems highlighting the main contributions of the paper.

We now present the formal definition of Coded Distributed Computation.

(Coded Distributed Computation) The coded distributed implementation of a computation task fA(⋅)f_{\mathbf{A}}(\cdot) is specified by:

local data blocks <Ai>i=1n\left<\mathbf{A}_{i}\right>_{i=1}^{n} and local computation tasks <fAii(⋅)>i=1n\left<f_{\mathbf{A}_{i}}^{i}(\cdot)\right>_{i=1}^{n};

a decoding function that outputs fA(⋅)f_{\mathbf{A}}(\cdot) given the results from a decodable set of local computations.

II-B Network Model

The network model is based on a master-worker setup illustrated in Fig. 1. The master node receives an input x\mathbf{x} and broadcasts it to all the workers. Each worker computes its assigned set of computations and unicasts the result to the master node. The master node aggregates the results from the worker nodes until it receives a decodable set of computations and recovers the output Ax\mathbf{A}\mathbf{x}.

II-C Problem Formulation

For a homogeneous cluster, to achieve a coded solution, one can divide A\mathbf{A} into kk equal size submatrices, and apply an (n,k)(n,k) MDS code to these submatrices. The master node can then obtain the final result from any kk responses. In , the authors find the optimal kk for minimizing the average running time for the shifted exponential run-time model.

For heterogeneous clusters, however, assigning equal loads to servers is clearly not optimal. Moreover, directly finding the optimal solution to Pmain\mathcal{P}_{\textnormal{main}} is hard. In homogeneous clusters, the problem of finding a sufficient number of inner products can be mapped to the problem of finding the waiting time for a set of fastest responses, and thus closed form expressions for the expected computation time can be found using order statistics of i.i.d. run-times. However, this is not straight-forward in heterogeneous clusters, where the load allocation is non-uniform. In Section III, we present an alternative formulation to Pmain\mathcal{P}_{\textnormal{main}} in (3), and show that the solution to the alternative formulation – which we shall name HCMM – is tractable and provably asymptotically optimal.

Assumptions. From now onward, we consider the practically relevant regime where the size of the problem scales linearly with the size of the network, while the computing power and the storage capacity of each worker node remain constant. Specifically, we assume r=Θ(n)r=\Theta{(n)}, ai=Θ(1)a_{i}=\Theta{(1)}, μi=Θ(1)\mu_{i}=\Theta{(1)} and αi=Θ(1)\alpha_{i}=\Theta{(1)} for each worker ii.

II-D Main Results

Having set the model and formulation of the problem, we now present the main contributions of this paper. The following theorem characterizes the asymptotic optimality of HCMM for the shifted exponential run-time model.

Theorem 1 demonstrates that our proposed HCMM algorithm is asymptotically optimal as the number of workers nn approaches infinity. In other words, the optimal computation load allocation problem Pmain\mathcal{P}_{\textnormal{main}} in (3) can be optimally solved using the proposed HCMM algorithm as nn gets large.

We note that Pmain\mathcal{P}_{\textnormal{main}} in (3) is a hard combinatorial optimization problem since it will require checking all load combinations to minimize the overall expected execution time. The key idea in Theorem 1 is to consider an alternative formulation to (3) focusing on maximizing the expected number of returned computation results from the workers, i.e. maximizing the aggregate return. As we describe in Section III, the alternative optimization problem not only can be solved efficiently in a tractable way giving rise to HCMM algorithm, it also asymptotically approximates Pmain\mathcal{P}_{\textnormal{main}} and allows us to establish Theorem 1.

While Theorem 1 theoretically characterizes the optimality of our proposed scheme HCMM, we also demonstrate gains that one can get in practice. In particular, we carry out numerical studies and experiments over Amazon EC2 clusters that demonstrate that HCMM can provide significant gains in a wide variety of computing scenarios. In particular, we compare HCMM’s performance with three benchmark load allocation policies – Uniform Uncoded, Load-balanced Uncoded, and Uniform Coded. In numerical studies, HCMM achieves speedups of up to 71%71\% over Uniform Uncoded, up to 53%53\% over Load-balanced Uncoded, and up to 39%39\% over Uniform Coded. In EC2 experiments, HCMM combined with the Luby transform (LT) codes provides speedups of up to 61%61\%, 46%46\% and 36%36\% over Uniform Uncoded, Load-balanced Uncoded and Uniform Coded, respectively.

Let TUCT_{\mathsf{UC}} denote the completion time of the uncoded distributed matrix multiplication algorithm. Then, for the shifted exponential run-times with constant parameters and r=Θ(n)r=\Theta{(n)},

As Theorem 2 shows, our proposed HCMM guarantees an improvement of \Theta\big{(}\log n\big{)} in expected execution time over any uncoded scheme, including the one that optimally allocates the workers’ loads. This result illustrates that by leveraging coded computing, one achieves the same order-wise gain over heterogeneous clusters as over homogeneous clusters .

Although Theorems 1 and 2 are based on the shifted exponential model (1) for run-time random variables for the workers, our analyses are general and can be extended to other models. The following two theorems generalize the results when the execution time of each worker follows the Weibull distribution as described in (2).

Under the Weibull distribution for run-times with constant parameters and r=Θ(n)r=\Theta{(n)}, the proposed HCMM scheme unboundedly outperforms the uncoded scheme, i.e.,

As stated in Theorem 4, HCMM provides an unbounded gain over any uncoded scheme – including the optimal uncoded load allocation – under the Weibull distribution for workers’ run-times. Furthermore, our numerical simulations demonstrate speedups of up to 73%73\%, 56%56\% and 42%42\% over Uniform Uncoded, Load-balanced Uncoded and Uniform Coded, respectively.

In the following section, we describe our alternative formulation based on aggregate return and describe our proposed HCMM algorithm that solves the alternative optimization.

III The Proposed HCMM Scheme and Proofs of Theorems 1 and 2

In this section, we prove Theorems 1 and 2 for the exponential model (1). In particular, we start by describing the HCMM algorithm and show that it asymptotically achieves the optimal performance, as stated in Theorem 1, and lastly conclude the section by characterizing the gain of HCMM over uncoded scheme.

To derive HCMM, we start by reformulating Pmain\mathcal{P}_{\textnormal{main}} defined in (3) and show that the alternative formulation can be efficiently solved, as opposed to solving Pmain\mathcal{P}_{\textnormal{main}} that needs an exhaustive search over all possible load allocations. The solution to the alternative problem gives rise to HCMM. We will further prove the optimality of HCMM and compare its average run-time to uncoded schemes.

We propose the following two-step alternative formulation for Pmain\mathcal{P}_{\textnormal{main}} defined in (3). First, for a fixed feasible time tt, we maximize the aggregate return over different load allocations, i.e., we solve

where X∗(t){X^{*}}(t) is the aggregate return at time tt for load allocation obtained from Palt(1)\mathcal{P}_{\textnormal{alt}}^{(1)}, that is

III-B Solving the Alternative Formulation

Considering the exponential distribution for workers’ run-times, we first proceed to solve Palt(1)\mathcal{P}_{\textnormal{alt}}^{(1)} in (4). The expected number of equations aggregated at the master node at time tt is:

Since there is no constraint on load allocations, Palt(1)\mathcal{P}_{\textnormal{alt}}^{(1)} can be decomposed to nn decoupled optimization problems, i.e.,

for all workers i∈[n]i\in[n]. The solution to (6) satisfies the following optimality condition:

where λi=Θ(1)\lambda_{i}=\Theta{(1)} is a constant independent of tt and is the positive solution to the following equation:

for all workers ii. In the following, we formally define the HCMM algorithm which is basically the solution to Palt\mathcal{P}_{\textnormal{alt}}.

We now provide an approximation to t∗t^{*} and show it asymptotically converges to t∗t^{*}. The expected aggregate return at time tt for optimal loads obtained in (7) is

since μi=Θ(1)\mu_{i}=\Theta{(1)} and λi=Θ(1)\lambda_{i}=\Theta{(1)}. Let τ∗\tau^{*} be the solution to the following equation when solved for tt:

In other words, τ∗\tau^{*} is the time for which there are exactly rr inner products – on average – aggregated at the master node, when the workers are loaded according to the loading obtained in (7). Using (7), (III-B) and (10), we find that

We now present the following lemma, which shows that τ∗\tau^{*} converges to t∗t^{*} for large nn (see Appendix for proof).

Let t∗t^{*} be the solution to the alternative formulation Palt\mathcal{P}_{\textnormal{alt}} in (4-5) and τ∗\tau^{*} be the solution to (10). Then,

III-C Asymptotic Optimality of HCMM

In this subsection, we prove the asymptotic optimality of HCMM as claimed in Theorem 1.

Consider the HCMM load assignment in (8). Let the random variable THCMMT_{\mathsf{HCMM}} denote the finish time associated to this load allocation, i.e. the waiting time to receive at least rr inner products from the workers. Let TmaxT_{\mathsf{max}} be the random variable denoting the finish time of all the workers for the HCMM load assignment.

Let us define two events E1\mathcal{E}_{1} and E2\mathcal{E}_{2} as follows:

Conditioning on these events, we can write

We can write the second term in RHS of (III-C) as follows:

To prove (a)(a), we note that HCMM returns rr inner products by time THCMMT_{\mathsf{HCMM}}. Moreover, the aggregate return is increasing in time. Therefore,

Moreover, the third term in RHS of (III-C) can be written as

where proof of (b)(b) is similar to proof of (a)(a) in (III-C). Regarding the first term in RHS of (III-C), we have

for some k1=Θ(1)k_{1}=\Theta(1) and large enough nn. To derive inequality (c)(c), we find a stochastic upper bound on TmaxT_{\mathsf{max}} by considering nn i.i.d. copies of the worker run-times with largest shift and smallest straggling parameters that are also Θ(1)\Theta(1), and use the PDF of the maximum of nn i.i.d. exponential random variables. As we later use in the proof of Theorem 3, one can similarly write for the shifted Weibull distribution:

for some constants k1k_{1} and k2k_{2}. Therefore, using (III-C), (III-C) and (III-C) (or (III-C) for the shifted Weibull model) in (III-C) we have

To this end, we show the following two inequalities,

which implies inequality (d)(d). We proceed to prove (e)(e) by showing the following two inequalities,

where τ∗\tau^{*} is obtained in (11). Given the fact that HCMM maximizes the expected aggregate return, we have

for every feasible tt, which implies (18). Moreover, Lemma 1 proves (19). All in all, we have

III-D Comparison with Uncoded Schemes

This subsection provides the proof of Theorem 2 by comparing the performance of HCMM to uncoded scheme. In an uncoded scheme, the redundancy factor is 11; thus, the master node has to wait for the results from all the worker nodes in order to complete the computation.

We start by characterizing the expected run-time of the best uncoded scheme. Particularly, we show that

where TUCT_{\mathsf{UC}} denotes the completion time of the optimum uncoded distributed matrix multiplication algorithm. To do so, we start by showing that

for all i∈[n]i\in[n]. Since the master needs to wait for all of the machines to return their results, the total run-time is T~UC=max⁡i∈[n]T~i\widetilde{T}_{\mathsf{UC}}=\max_{i\in[n]}\widetilde{T}_{i}. Therefore,

where Hn=1+12+13+⋯+1nH_{n}=1+\frac{1}{2}+\frac{1}{3}+\cdots+\frac{1}{n} is the sum of the harmonic series. We can further bound (20) using the fact that

Now consider another set of nn machines, where each machine is replaced with a slower one with parameters (a^,μ^)(\hat{a},\hat{\mu}) for a^=max⁡iai\hat{a}=\max_{i}a_{i} and μ^=min⁡iμi\hat{\mu}=\min_{i}\mu_{i}. By an argument similar to the one employed the lower bound, we can write

for another constant CC. From (21) and (22), one can conclude that

Further, by Theorem 1 and Lemma 1, we find that

Comparing (23) to (24) demonstrates that HCMM outperforms the best uncoded scheme by a factor of Θ(log⁡n)\Theta(\log n), i.e.,

IV Generalization to the Shifted Weibull Model and Proofs of Theorems 3 and 4

In this section, we consider the shifted Weibull distribution for the workers’ execution times, which captures a broader class of run-time models than the exponential distribution. We particularly generalize our proposed HCMM algorithm to the class of shifted Weibull distributed run-times and prove Theorems 3 and 4. More specifically, we argue that asymptotic optimality of HCMM is derived similar to the shifted exponential case and further show that HCMM provides unbounded gain over uncoded schemes, asymptotically.

A random variable TT has Weibull distribution with shape parameter α>0\alpha>0 and scale parameter μ>0\mu>0, denoted by T∼W(α,μ)T\sim\mathcal{W}(\alpha,\mu), if the CDF of TT is of the following form:

As in the exponential case, we begin by maximizing the expected aggregate return at the master node (Palt(1)\mathcal{P}_{\textnormal{alt}}^{(1)}) under the shifted Weibull distribution, which is given by

The optimal load allocation that maximizes the individual expected aggregate returns at each worker (and thus the total aggregate return) can be found by solving the following equation:

Similar to Section III, we can define ss as follows,

With the aforementioned reparametrizations of λi{\lambda_{i}}, ss and τ∗\tau^{*}, the HCMM algorithm defined in Algorithm 1 is identically applicable to the Weibull model. Proof of the asymptotic optimality of HCMM under the Weibull distribution follows the similar steps as in the proof for the exponential case in Section III-C (unless specifically justified, e.g. (III-C)). We avoid rewriting these steps for the purpose of readability of the paper, but we note that the concentration inequalities used to establish the proof of Theorem 1 can be applied to a wide class of distributions including the Weibull distribution. ∎

Let {Ti}i=1∞\{T_{i}\}_{i=1}^{\infty} be a sequence of i.i.d. W(α,μ)\mathcal{W}(\alpha,\mu) random variables and Tn∗=max⁡i∈[n]TiT^{*}_{n}=\max_{i\in[n]}T_{i} denote the maximum of the first nn variables. Then,

V Numerical Studies and Experiments using Amazon EC2 Machines

In this section, we present our results both from simulations as well as from experiments over Amazon EC2 clusters. These results demonstrate how HCMM can provide significant speedups in comparison to state-of-the-art load allocation schemes.

We now present numerical results evaluating the performance of HCMM. We consider both the shifted exponential model in (1) and the shifted Weibull model in (2) for run-time distributions in our simulations, assuming the unit seconds per row (s/row\text{s}/\text{row}) for aa and 1/μ1/\mu. The underlying computation task is to compute r=10000r=10000 inner products using a heterogeneous cluster of n=100n=100 workers, where different scenarios for heterogeneity are considered. For each scenario under consideration, we implement the following load allocation schemesFor each scheme, the load number for each worker is approximated to the nearest larger integer using the ceil()\mathtt{ceil()} function. For the practical large load regime considered in simulations, this rounding step has negligible impact on load allocation and on the overall results.:

Uniform Coded: Equal number of coded rows are assigned to each worker. Redundancy is numerically optimized for minimizing the average computation time for receiving results of at least rr inner products at the master node.

For simulations under the shifted exponential model, we consider the following three scenarios:

Scenario 1\mathbf{1} (2\mathbf{2}-mode heterogeneity): (ai,μi)=(1,1)(a_{i},\mu_{i})=(1,1) for 5050 workers, and (ai,μi)=(4,0.5)(a_{i},\mu_{i})=(4,0.5) for the other 5050 workers.

Scenario 2\mathbf{2} (3\mathbf{3}-mode heterogeneity): (ai,μi)=(1,0.5)(a_{i},\mu_{i})=(1,0.5) for 2525 workers, (ai,μi)=(4,2)(a_{i},\mu_{i})=(4,2) for 2525 workers, and (ai,μi)=(12,0.25)(a_{i},\mu_{i})=(12,0.25) for the remaining 5050 workers.

Scenario 3\mathbf{3} (Random heterogeneity): For each worker ii, parameters aia_{i} and μi\mu_{i} are sampled from the sets {1,4,12}\{1,4,12\}, {0.5,2,0.25}\{0.5,2,0.25\}, respectively and all uniformly at random.

The following three scenarios are considered for simulations under the shifted Weibull distribution for run-times:

Scenario 1\mathbf{1} (2\mathbf{2}-mode heterogeneity): (ai,μi,αi)=(1,1,1.2)(a_{i},\mu_{i},\alpha_{i})=(1,1,1.2) for 5050 workers, and (ai,μi,αi)=(4,0.5,0.8)(a_{i},\mu_{i},\alpha_{i})=(4,0.5,0.8) for the other 5050 workers.

Scenario 2\mathbf{2} (3\mathbf{3}-mode heterogeneity): (ai,μi,αi)=(1,0.5,0.9)(a_{i},\mu_{i},\alpha_{i})=(1,0.5,0.9) for 2525 workers, (ai,μi,αi)=(4,2,1.2)(a_{i},\mu_{i},\alpha_{i})=(4,2,1.2) for 2525 workers, and (ai,μi,αi)=(12,0.25,1.5)(a_{i},\mu_{i},\alpha_{i})=(12,0.25,1.5) for the remaining 5050 workers.

Scenario 3\mathbf{3} (Random heterogeneity): For each worker ii, parameters aia_{i}, μi\mu_{i} and αi\alpha_{i} are sampled from the sets {1,4,12}\{1,4,12\}, {0.5,2,0.25}\{0.5,2,0.25\} and {0.9,1.2,1.5}\{0.9,1.2,1.5\}, respectively and all uniformly at random.

Fig. 2 and 3 illustrate the performance comparison of the four schemes for the two run-time models. We make the following conclusions from the results.

HCMM significantly outperforms the benchmark load allocation schemes. In particular, for the shifted exponential model, HCMM provides speedups of up to 71%71\% over Uniform Uncoded, up to 53%53\% over Load-balanced Uncoded, and up to 39%39\% over Uniform Coded, among the three scenarios. When the machine run-time is assumed to have a shifted Weibull distribution, among the three scenarios HCMM results in gains of up to 73%73\%, 56%56\% and 42%42\% over Uniform Uncoded, Load-balanced Uncoded, and Uniform Coded respectively.

Both Load-balanced Uncoded and Uniform Coded improve upon the performance of Uniform Uncoded. In Load-balanced Uncoded scheme, assigning larger loads to faster machines leads to better performance, while for Uniform Coded, repeated computations lead to better performance as the master does not need to wait for all the results. HCMM provides the best expected execution time among the four schemes as it combines the gains of Load-balanced Uncoded and Uniform Coded by employing efficient load balancing along with minimal number of redundant computations.

Next, we present the results from our experiments over Amazon EC2 clusters. These results show agreement with our numerical studies.

V-B Experiments using Amazon EC2 machines

We use Python with mpi4py package to implement our developed HCMM scheme over Amazon EC2 clusters. To emulate the straggler effects in large-scale systems , we inject artificial delays.Artificial delays are injected since stragglers are rarely observed in small clusters in Amazon EC2. Though other emerging platforms such as federated learning, computation with deadline, mobile edge computing, fog computing, etc., still suffer from stragglers where our ideas can be employed . This is achieved by selecting some workers to be stragglers at the beginning of experiments and slowing down each such worker by making it wait for 33 times the amount of time it spends in computation before it sends its results to the master. This is done using the sleep() function in time package. For each scenario, the choice of stragglers is made by drawing a sample from the Bernoulli(0.5)(0.5) distribution for each worker, i.e., each worker is chosen to be a straggler with probability 0.50.5.

For performance comparison of the four schemes, we consider the following three computing scenarios:

Scenario 1\mathbf{1}: Each row has 500000500000 elements. We use a heterogeneous cluster of 1111 machines – one master of instance type m4.xlarge , four workers of instance type r4.2xlarge , and six workers of instance type r4.xlarge .

Scenario 2\mathbf{2}: Each row has 500000500000 elements. We use a heterogeneous cluster of 1616 machines – one master of instance type m4.xlarge , six workers of instance type r4.2xlarge , and nine workers of instance type r4.xlarge .

Scenario 3\mathbf{3}: Each row has 10000001000000 elements. We use the same heterogeneous cluster as in the previous scenario.

Fig. 4 provides a performance comparison of HCMM with the benchmark load allocation schemes for the three scenarios, where the decoding time is taken into account as well. Fig. 5 presents the typical cumulative distribution functions for the instances used in the experiments. We make the following conclusions from the results:

As demonstrated in Fig. 5, the shifted exponential model is a good first order fit for the run-times of the workers.

HCMM achieves significant speedups over the benchmark load allocation policies. In particular, HCMM combined with LT codes provides gains in the overall execution time of up to 61%61\% over Uniform Uncoded, up to 46%46\% over Load-balanced Uncoded, and up to 36%36\% over Uniform Coded.

As presented in Table I, HCMM has significantly lower total computation load compared to Uniform Coded. Hence, HCMM leads to efficient utilization of the computing resources, combining the benefits of both Load-balanced Uncoded and Uniform Coded schemes.

These results demonstrate that HCMM can provide significant speedups in large-scale computing environments.

VI Generalization to Computing Scenarios under Budget Constraints

In this section, we consider the optimization problem in (3) under the shifted exponential distribution with a monetary constraint for carrying out the overall computation. Running computation tasks on a commodity server costs depending on several factors including CPU, memory, ECU, storage, bandwidth, etc. Different cloud computing platforms employ different pricing policies, and these need to be taken into account for developing efficient task allocation and execution algorithms . For example, Table II summarizes the cost per hour of using Amazon EC2 clusters with different parameters (at the time of writing this manuscript) . In this section, we take into account the monetary constraint in the optimization problem in (3) and provide a heuristic algorithm towards finding the optimal load allocation under cost budget constraint.

We now present the precise problem formulation we are interested in. For a computation task and a given set of NN machines, the goal is to minimize the expected run-time while satisfying the budget constraint CC, that is

where cic_{i} represents the cost per time unit of using machine i∈[N]i\in[N]. According to the pricing polices provided by AWS, e.g. Table II, a linear model for cost (versus performance parameters) is intuitive and convincing. Considering the last two rows of Table II for instance, doubling the parameters results in doubled cost. To be general, we model the computation cost of a single machine as c=κμγc=\kappa\mu^{\gamma} per unit of time, which captures a convex dependency of the speed parameter μ\mu for constants γ≥1\gamma\geq 1 and κ>0\kappa>0.

We assume that there are KK types of machines parameterized with {(ak,μk)}k=1K\{(a_{k},\mu_{k})\}_{k=1}^{K}, and Nk,k∈[K]N_{k},k\in[K] of each type is available to run a distributed computation task, where N=∑k=1KNkN=\sum_{k=1}^{K}N_{k} is the total number of available machines. We also assume that μ1≤⋯≤μK\mu_{1}\leq\cdots\leq\mu_{K} and a1μ1=⋯=aKμK=ξa_{1}\mu_{1}=\cdots=a_{K}\mu_{K}=\xi for a constant ξ\xi.The latter assumption can be intuitively justified as follows. If a machine is cc times more powerful than another machine, as the first order estimation, one can assume that both the shift (aka_{k}) and the straggling parameter (μk\mu_{k}) of the computation are cc times stronger. As we showed in Theorem 1, HCMM is asymptotically optimal (i.e. optimal within a vanishing deviation) regarding the average run-time. In this section, we also consider the asymptotic regime, i.e. for large enough number of machines and hence HCMM attains the optimality per Pmain\mathcal{P}_{\textnormal{main}} in (3).

The following lemma states a useful observation regarding the solutions to the constrained problem Pmain-constrained\mathcal{P}_{\textnormal{main-constrained}} and the minimum possible cost for carrying out a computation task.

HCMM is the (asymptotic) solution to the feasible Pmain-constrained\mathcal{P}_{\textnormal{main-constrained}}. Moreover, given a computation task and a set of machines, decreasing the number of fastest (slowest) machines in HCMM, results in smaller (greater) expected cost. And, the minimum (maximum) cost of HCMM is induced by running the task only on any number of the slowest (fastest) machines.

We first argue that if the budget-constrained problem defined in Pmain-constrained\mathcal{P}_{\textnormal{main-constrained}} is feasible, then HCMM determines the asymptotically optimal load allocation. Consider a set of NN machines and assume that MM of them are assigned non-zero loads in the optimal budget-constrained scheme. Now, one can run HCMM load allocation over the set of these MM machines and according to asymptotic optimality results, HCMM asymptotically attains the optimal run-time while satisfying the budget constraint.

Now assume that nkn_{k} number of type k∈[K]k\in[K] machine is used. Then, by assigning the loads obtained from HCMM and the result of Theorem 1, the induced expected cost (for large number of machines) can be written as

where xξ=1+μkλkx_{\xi}=1+\mu_{k}\lambda_{k} is the solution to the equation exξ−ξ−1=xξe^{x_{\xi}-\xi-1}=x_{\xi} for all machine type k∈[K]k\in[K]. In another scenario, assume that we remove one machine of type KK (the fastest machine type) and run HCMM accordingly, i.e. nkn_{k} of type k∈[K−1]k\in[K-1] and nK−1n_{K}-1 of type KK. The expected cost of this scenario can be written as follows:

where inequality (f)(f) can be easily verified given that μ1≤⋯≤μK\mu_{1}\leq\cdots\leq\mu_{K}. We can iteratively apply the same argument and conclude that the minimum expected cost is achieved when only the slowest machines are used, that is

for any 1≤n1≤N11\leq n_{1}\leq N_{1}. Similar to (VI), one can show that reducing the number of participating slowest machines increases the induced expected cost of HCMM, that is

Therefore, applying (VI) iteratively shows that the maximum expected cost occurs when only the fastest machines are employed, that is

Lemma 3 implies that if the available budget CC is less than CminC_{\text{min}} defined in (29), then Pmain-constrained\mathcal{P}_{\textnormal{main-constrained}} is infeasible and it is impossible to run the task on the given set of machines while satisfying the budget constraint. Moreover, reducing one machine from the available set of fastest machines along with HCMM results in a lower expected cost; and reducing the number of participating slowest machines results in a larger expected cost.

Now that HCMM asymptotically solves the feasible budget-constrained problem in (26), i.e. for C≥CminC\geq C_{\text{min}}, finding the optimal number of machines of each type to use in HCMM requires combinatorial search over all possible allocations. However, as Lemma 3 suggests, using faster machines induces a larger cost. Further, the computation time increases if we decrease the number of machines. This is the motivation behind our heuristic algorithm for an efficient search to find the number of machines of each type to include in HCMM, which we describe next.

First, Algorithm 2 runs HCMM algorithm using all machines, i.e., nk=Nkn_{k}=N_{k} for each k∈[K]k\in[K]. Then, it calculates the corresponding cost according to (VI). If the cost is >C>C, it starts to decrease the number of available fastest machines, i.e. nK←nK−1n_{K}\leftarrow n_{K}-1, and runs HCMM again. While the cost is >C>C, the algorithm keeps decreasing the number of used fast machines till nK=0n_{K}=0. Then, the algorithm sets nK=0n_{K}=0 and starts decreasing nK−1n_{K-1} and so on, until a feasible cost is achieved. Thus, the algorithm returns (N1,⋯ ,Nj,nj+1,0,⋯ ,0)(N_{1},\cdots,N_{j},n_{j+1},0,\cdots,0) which is the first tuple that satisfies the cost constraint. Therefore, the search space complexity of the heuristic is O(N1+⋯+NK)=O(N)\mathcal{O}(N_{1}+\cdots+N_{K})=\mathcal{O}(N) which is more efficient than the exhaustive search where the complexity is O(N1⋯NK)\mathcal{O}(N_{1}\cdots N_{K}). The pseudo-code in Algorithm 2 summarizes the heuristic.

In this example, we consider two different scenarios to demonstrate the application of the proposed heuristic search algorithm. For the cost model, we assume γ=2\gamma=2 and κ=1\kappa=1, i.e. c=μ2c=\mu^{2}. Further, we consider the task of computing r=100r=100 equations.

VII Conclusion

In this paper, we proposed a coding framework for distributed matrix-vector multiplication in heterogeneous cloud computing environments. In particular, we considered two distributions for machines’ run-times, i.e. shifted exponential and shifted Weibull and tackled the intractable problem of minimizing the average run-time of a computation task over all possible load allocations by proposing a tractable alternative formulation. The solution to the alternative problem established our proposed HCMM load allocation scheme which we proved to be asymptotically optimal. We also demonstrated the speedup of HCMM over three benchmark load allocation schemes and presented both the numerical and the experimental results. Experiments over Amazon EC2 clusters demonstrate that HCMM combined with LT codes and peeling decoders can provide significant gains in the average overall execution time. Moreover, we argued that HCMM is the asymptotically optimal allocation in budget-constrained scenarios as well, which led to providing a heuristic search in order to find a (sub)optimal load-machine assignment for a given set of machines while satisfying a pre-defined budget constraint.

VIII Acknowledgment

This material is based upon work supported by Defense Advanced Research Projects Agency (DARPA) under Contract No. HR001117C0053, ARO award W911NF1810400, NSF grants CCF-1703575, ONR Award No. N00014-16-1- 2189, CCF-1763673, CCF-1755808 and the UC Office of President under grant No. LFR-18-548175. The views, opinions, and/or findings expressed are those of the author(s) and should not be interpreted as representing the official views or policies of the Department of Defense or the U.S. Government.

References

Appendix A

for any x1,⋯ ,xn,xi′∈Xx_{1},\cdots,x_{n},x^{\prime}_{i}\in\mathcal{X} and i∈[n]i\in[n]. Then, for any ϵ>0\epsilon>0,

for any ϵ>0\epsilon>0. Now, we proceed to the proof of Lemma 1.

Let t=τ∗+δt=\tau^{*}+\delta for some δ=Θ(log⁡nn)\delta=\Theta\left(\frac{\log n}{\sqrt{n}}\right) and ϵ=δ2\epsilon=\delta^{2}. The claim is that \Pr\big{[}X^{*}(t)\leq r-\epsilon\big{]}=o\left(\frac{1}{n}\right). From McDiarmid’s inequality, we have

In above, equality (g)(g) follows from the fact that r=Θ(n)r=\Theta(n), s=Θ(n)s=\Theta(n), λi=Θ(1)\lambda_{i}=\Theta(1), δ=Θ(log⁡nn)\delta=\Theta\left(\frac{\log n}{\sqrt{n}}\right), and therefore ∑iλi2=Θ(n)\sum_{i}\lambda_{i}^{2}=\Theta(n) and s2=Θ(n2)s^{2}=\Theta(n^{2}). Moreover, if t∗<τ∗t^{*}<\tau^{*}, with a positive probability there are less than rr equations at the master node by time t∗t^{*} which is a contradiction. Therefore,