Minimizing Latency for Secure Distributed Computing

Rawad Bitar, Parimal Parag, Salim El Rouayheb

I Introduction

We consider the setting of distributed computing in which a server M, referred to as Master, possesses confidential data, such as personal information of online users, genomic and medical data etc., and wants to perform intensive computations on it. M wants to divide these computations into smaller computational tasks and distribute them to nn worker machines that can perform these smaller tasks in parallel. The workers then return their results to the master, who can process them to obtain the result of its original task. The well celebrated MapReduce framework falls under this model and is implemented in many computing clusters.

In this paper, we are interested in applications in which the worker machines do not belong to the same system or cluster as the master. Rather, the workers are online computing machines that can be hired or can volunteer to help the master in its computations. Existing applications that fall under this model include the SETI@home project for search for extraterrestrial intelligence , the folding@home project for disease research that simulates protein folding and Amazon mechanical turkAmazon mechanical turk hires humans to perform tasks. But, one can imagine a similar application where computing machines are hired. . The additional constraint that we worry about here, and which does not exist in the previous applications, is that the workers cannot be trusted with the sensitive data, which must remain hidden from them. Our privacy constraint is information theoretic, meaning that each worker must obtain zero information about the data irrespective of its computational power. We choose information theoretic privacy instead of homomorphic encryption, due to the high computation and memory overheads of the latter .

We focus on linear computations (matrix multiplication) since they form a basic building block of many iterative algorithms. The workers introduce random delays due to the difference of their workloads or network congestion. This causes the Master to wait for the slowest workers, referred to as stragglers in the distributed computing community . In addition, some workers may never respond. Our goal is to reduce the delay at the Master caused by the workers.

Privacy can be achieved by encoding the data using a linear secret sharing codes as illustrated in Example 1. However, these codes are not specifically designed to minimize latency as we will highlight later.

Let the matrix AA denote the data set owned by M and let x\mathbf{x} be a given vector. M wants to compute AxA\mathbf{x} . Suppose that M gets the help of n=3n=3 workers out of which at most n−k=1n-k=1 may be unresponsive. M generates a random matrix RR of same dimensions as AA and over the same field and encodes AA and RR into 3 shares S1=RS_{1}=R, S2=R+AS_{2}=R+A and S3=R+2AS_{3}=R+2A using a secret sharing scheme . First, M sends share SiS_{i} to worker Wi{\text{W}}_{i} (Figure 3) and then sends x\mathbf{x} to all the workers. Each worker computes SixS_{i}\mathbf{x} and sends it back to M (Figure 3). M can decode AxA\mathbf{x} after receiving any k=2k=2 responses. For instance, if the first two workers respond, M can obtain Ax=S2x−S1xA\mathbf{x}=S_{2}\mathbf{x}-S_{1}\mathbf{x}. No information about AA is revealed to the workers, because AA is one-time padded by RR.

The delay experienced by M in the previous example results from the fact it has to wait until k=2k=2 workers finish their whole tasks in order to decode AxA\mathbf{x}, even when the 33 workers are all responsive. This is due to the fact that classical secret sharing codes are designed for the worst-case scenario of one worker being unresponsive. We overcome this limitation by using Staircase codes which were introduced in and are explained in the next example.

Consider the same setting as Example 1. Instead of using a classical secret sharing code, M now encodes AA and RR using the Staircase code given in Table I.

The Staircase code requires M to divide the matrices AA and RR into A=[A1A2]TA=\begin{bmatrix}A_{1}&A_{2}\end{bmatrix}^{T} and R=[R1R2]TR=\begin{bmatrix}R_{1}&R_{2}\end{bmatrix}^{T}. In this setting, M sends two subshares to each worker, hence each task consists of 22 subtasks. The master sends x\mathbf{x} to all the workers. Each worker multiplies the subshares by x\mathbf{x} (going top to bottom) and sends each multiplication back to M independently. Now, M has two possibilities for decoding: 1) Mreceives the first subtask from all the workers, i.e., receives (A1+A2+R1)x(A_{1}+A_{2}+R_{1})\mathbf{x}, (A1+2A2+4R1)x(A_{1}+2A_{2}+4R_{1})\mathbf{x} and (A1+2A2+4R1)x(A_{1}+2A_{2}+4R_{1})\mathbf{x} and decodes AxA\mathbf{x} which is the concatenation of A1xA_{1}\mathbf{x} and A2xA_{2}\mathbf{x}. Note that M decodes only R1xR_{1}\mathbf{x} and does not need to decode R2xR_{2}\mathbf{x}. 2) Mreceives all the subtasks from any 22 workers and decodes AxA\mathbf{x}. Here M has to decode R1xR_{1}\mathbf{x} and R2xR_{2}\mathbf{x}. One can check that no information about AA is revealed to the workers.

Under an exponential delay model for each worker, we show that the Staircase code given in Example 2 can lead to a 25%25\% improvement in delay over the secret sharing code given in Example 1. Our goal is to give a general systematic study of the delay incurred by Staircase codes and compare it to classical secret sharing codes.

Related work: Straggler mitigation and privacy concerns are studied separately in the literature. In Liang et al. adaptively encoded the tasks depending on the workload at the workers’ end. Lee et al. used MDS codes to mitigate stragglers in linear distributed machine learning algorithms. Tandon et al. introduced new codes for straggler mitigation in distributed gradient descent algorithms. Li et al. studied the effect of the workers’ computation load on the communication complexity.

On the other hand, privacy concerns have been studied in the machine learning literature, see e.g., . The main model assumes that several parties owning private data sets want to train a model based on all the data sets without revealing them, e.g., . However, the techniques extensively rely on cryptographic assumptions and secure multi-party computation. Atallah and Frikken studied the problem of distributively multiplying two private matrices assuming that k−1k-1 workers can collude (with k2<n2k^{2}<n^{2}). The provided solution ensures information theoretic privacy, but does not account for straggler mitigation. Another related problem is federated learning . A large number of users own different amounts of data and a central server aims to train a high-quality model based on all the data with the smallest communication complexity. However, privacy is ensured by keeping the data local to the users.

Contributions: In this paper, we consider the model in which M owns the whole data set on which it wants to perform a distributed linear computation. We introduce a new approach for securely outsourcing the linear computations to nn workers which do not own any parts of the data. The data set is to be kept private in an information theoretic sense. We assume that at most n−k, k<nn-k,\ k<n, workers may be unresponsive, the remaining respond at random times. This is similar to the straggler problem. We study the master’s waiting time, i.e., the aggregate delays caused by the workers, under the exponential model when using Staircase codes. More specifically, we make the following contributions: (i) we derive an upper bound and a lower bound on the mean waiting time; (ii) we derive an integral expression leading to the CDF of the waiting time and use this expression to find the exact mean waiting time for the cases when k=n−1k=n-1 and k=n−2k=n-2; and (iii) we compare our approach to the approach using secret sharing and show that for high rates, k/nk/n, and small number of workers our approach saves about 40%40\% of the waiting time. Moreover, we ran simulations to check the tightness of the bounds and show that for low rates our approach saves at least 10%10\% of the waiting time for all values of nn.

II System Model

Workers model: The workers have the following properties: 1) At most n−kn-k workers may be unresponsive. The actual number of unresponsive workers is unknown a priori. 2) The responsive workers incur random delays while executing the task assigned to them by M resulting in what is known as the straggler problem . We model all the delays incurred by each worker by an independent and identical exponential random variable. 3) The workers do not collude, i.e., they do not share with each other the data they receive from M. This has implications on the privacy constraint described later.

General scheme: M encodes AA, using randomness, into nn shares SiS_{i} sent to worker Wi{\text{W}}_{i}, i=1,…,ni=1,\dots,n. Any kk or more shares can decode AA. The workers obtain zero information about AA, i.e., H(A∣Si)=H(A)H(A|S_{i})=H(A) for all i∈{1,…,n}i\in\{1,\dots,n\}.

At each iteration, the master sends x\mathbf{x} to all the workers. Then, each worker computes SixS_{i}\mathbf{x} and sends it back to the master. Since the scheme and the computations are linear, the master can decode AxA\mathbf{x} after receiving enough responsesIn some cases the attribute vectors xj\mathbf{x}^{j} contain information about AA, and therefore need to be hidden from the workers. We describe in how our scheme can be generalized to such cases.. We refer to such scheme as an (n,k)(n,k) system.

Delay model: Let TAT_{A} be the random variable representing the time spent to compute AxA\mathbf{x} at one worker. We assume a mother runtime distribution FTA(t)F_{T_{A}}(t) that is exponentialOur analysis remains true for the shifted exponential model . with rate λ\lambda. Due to the encoding, each task given to a worker is k−1k-1 times smaller than AA. Let Ti, i∈{1,…,n}T_{i},\ i\in\{1,\dots,n\} denote the time spent by worker Wi{\text{W}}_{i} to execute its task, then we assume that FTiF_{T_{i}} is a scaled distribution of FTAF_{T_{A}}, i.e.,

For an (n,k)(n,k) system using Staircase codes, we assume that TiT_{i} is evenly distributed between the subshares, i.e., the time spent by a worker Wi{\text{W}}_{i} on one subshare is equal to Ti/αT_{i}/\alpha. Let T(i)T_{(i)} be the ithi^{th} order statistic of the TiT_{i}’s and TSCT_{\text{SC}} be the time the master waits until it can decode AxA\mathbf{x}. We can write

where αi≜(k−1)/(i−1)\alpha_{i}\triangleq(k-1)/(i-1). For an (n,k)(n,k) system using classical secret sharing codes, we can write TSS=T(k).T_{\text{SS}}=T_{(k)}.

III Main Results

Our main results are summarized as follows. We provide an upper bound and a lower bound on the mean waiting time of M in Theorem 1.

where HnH_{n} is the nthn^{\text{th}} harmonic sum defined as Hn≜∑i=1n1iH_{n}\triangleq\sum_{i=1}^{n}\frac{1}{i}, and H0≜0H_{0}\triangleq 0. The mean waiting time is lower bounded by

Discussion: Our extensive simulations show that (1) is a good approximation of the mean waiting time. Moreover, by taking d=kd=k in (1), the upper bound on the mean waiting time of Staircase codes becomes the one of classical secret sharing, i.e.,

While finding the exact expression of the mean waiting time for any (n,k)(n,k) system remains open, we derive in Corollary 1 an expression for systems with 11 and 22 parities, i.e. (k+1,k)(k+1,k) and (k+2,k)(k+2,k) systems, using the result of Theorem 2. Using Corollary 1 one can compare the performance of Staircase codes an secret sharing codes. For instance, in a (4,2)(4,2) system Staircase codes reduce the mean waiting time by 40%40\%.

Let ti≜t(i−1)/(k−1)t_{i}\triangleq t(i-1)/(k-1), the CDF of the waiting time TSCT_{\text{SC}} of an (n,k)(n,k) system using Staircase codes is given by

where A(t)=∩i≥k{yi∈(ti,yi+1]}A(t)=\cap_{i\geq k}\{y_{i}\in(t_{i},y_{i+1}]\} and F(yi)=FTi(yi)F(y_{i})=F_{T_{i}}(y_{i}).

To check the tightness of the bounds we plot in Figure 4 the upper bound in (1), lower bound in (1) and the exact mean waiting time in (17) for (k+2,k)(k+2,k) systems.

IV Proof of Theorem 1

We will need the following characterization of order statistics for iid exponential random variables.

The dthd^{\text{th}} order statistic T(d)T_{(d)} of nn iid exponential random variables TiT_{i}, with distribution function F(t)=1−e−λt{F}(t)=1-e^{-\lambda t}, is equal to a random variable ZZ in the distribution, where

and ZjZ_{j} are iid random variables with distribution F(t){F}(t).

Since min⁡\min is a convex function, we can use Jensen’s inequality to write

Equations (5) and (6) conclude the proof. We give an intuitive behavior of the upper bound. The harmonic number can be approximated by Hn≈log⁡(n)+γ,H_{n}\approx\log(n)+\gamma, where γ≈0.577218\gamma\approx 0.577218 is called the Euler-Mascheroni constant. Therefore, log⁡(n)<Hn<log⁡(n+1)\log(n)<H_{n}<\log(n+1). Hence, we can write

IV-B Lower bound on the mean waiting time

where αj≜(k−1)/(j−1)\alpha_{j}\triangleq(k-1)/(j-1). For TSCT_{\text{SC}} to be greater than tt, all the jthj^{\text{th}} order statistic T(j)T_{(j)}’s must be greater than t/αjt/\alpha_{j} for j∈{k,…,n}j\in\{k,\dots,n\}. We show that if C\mathcal{C} is satisfied, then the previous condition is satisfied. If T(k)>t/αdT_{(k)}>t/\alpha_{d}, then T(i)>t/αiT_{(i)}>t/\alpha_{i} for all i∈{k,…,d}i\in\{k,\dots,d\}, because T(i)≥T(k)>t/αd>t/αiT_{(i)}\geq T_{(k)}>t/\alpha_{d}>t/\alpha_{i}. It follows that if for all j∈{d+1,…,n},j\in\{d+1,\dots,n\}, T(j)−T(j−1)>t/αj−t/αj−1T_{(j)}-T_{(j-1)}>t/\alpha_{j}-t/\alpha_{j-1}, then T(j)>t/αjT_{(j)}>t/\alpha_{j}. Therefore, Pr⁡(TSC>t)≥Pr⁡(C is satisfied)≜Pr⁡(C)\Pr\left(T_{\text{SC}}>t\right)\geq\Pr\left(\mathcal{C}\text{ is satisfied}\right)\triangleq\Pr(\mathcal{C}). Furthermore,

Next we derive an expression of ∫0∞Pr⁡(C)dt\int_{0}^{\infty}\Pr\left(\mathcal{C}\right)dt. Note that 1/αj−1/αj−1=1/(k−1)1/\alpha_{j}-1/\alpha_{j-1}=1/(k-1), using Theorem 3 we can write

where FˉZj(t)≜Pr⁡(Zj>t)\bar{F}_{Z_{j}}(t)\triangleq\Pr\left(Z_{j}>t\right). From (9) we get

Since FˉZj(t)=e−(k−1)λt\bar{F}_{Z_{j}}(t)=e^{-(k-1)\lambda t}, we can write

On the other hand, FˉT(k)(t/αd)\bar{F}_{T_{(k)}}\left(t/\alpha_{d}\right) is the probability that there are at most k−1k-1 TiT_{i}’s less than t/αdt/\alpha_{d}, therefore

Recall that FTi(t)=1−e(k−1)λt=1−FˉTi(t){F}_{T_{i}}(t)=1-e^{(k-1)\lambda t}=1-\bar{F}_{T_{i}}(t), therefore by using the binomial expansion we can write

Using (13) and the fact that FTi(t)=e−(k−1)λtF_{T_{i}}(t)=e^{-(k-1)\lambda t}, (12) becomes

Combining (11) and (14) and noting that FˉTi(t)=FˉZj(t)=e−λ(k−1)t\bar{F}_{T_{i}}(t)=\bar{F}_{Z_{j}}(t)=e^{-\lambda(k-1)t}, (10) becomes

Note that ∫0∞e−xtdt=1/x\int_{0}^{\infty}e^{-xt}dt=1/x and that the integral of a sum is equal to the sum of the integrals. Therefore, integrating (IV-B) from to ∞\infty and maximizing it over all values of d, d∈{k,…,n}d,\ d\in\{k,\dots,n\}, concludes the proof.

V Proof of Theorem 2

We derive an integral expression leading to the probability distribution of the waiting time TSCT_{\text{SC}}. Since the delays at the workers’ side TiT_{i}’s are independent and are absolutely continuous with respect to the Lebesgue measure (i.e. the probability density exists), we have

where tit_{i} denotes t/αit/\alpha_{i} and 0≤t1≤…≤tn0\leq t_{1}\leq\ldots\leq t_{n}. Therefore we can write the distribution of TSCT_{\text{SC}} as

That is, we can re-write Pr⁡{TSC>t}\Pr\{T_{\text{SC}}>t\} as

∫0yk⋯∫0y2∏i=1k−1dFTi(yi)=F(yk)k−1(k−1)!.\int_{0}^{y_{k}}\cdots\int_{0}^{y_{2}}\prod_{i=1}^{k-1}dF_{T_{i}}(y_{i})=\frac{F(y_{k})^{k-1}}{(k-1)!}.

The result of Claim 1 is straightforward, it follows from integrating k−1k-1 times the complementary CDF of an exponential random variable in respect to its derivative. This completes the proof. A more detailed proof of Claim 1 can be found in . We state the mean waiting time for the (k+2,k)(k+2,k) and (k+1,k)(k+1,k) systems in Corollary 1.

VI Simulations

We check the tightness of the bounds of Theorem 1 and measure the improvement, in terms of delays, of Staircase codes over classical secret sharing codes for systems with fixed rate R≜k/nR\triangleq k/n. In Figure 7 (a) we plot the upper bound (1), lower bound (1) and the simulated mean waiting time for R=1/4R=1/4. Our extensive simulations show that the upper bound is a good approximation of the exact mean waiting time, whereas the lower bound might be loose.

References