Information Theoretic Limits of Data Shuffling for Distributed Learning

Mohamed Attia, Ravi Tandon

I Introduction

Distributed computing systems for large data-sets have gained a lot of interest recently as they enable the processing of data-intensive tasks for machine learning, model tracing, and data analysis over a large number of commodity machines, and servers (e.g., Apache Spark , and MapReduce ). Generally speaking, a master node, which has the entire data-set, sends data blocks to be processed at distributed worker nodes. The workers subsequently respond with locally computed functions to the master node for the desired data analysis. This enables the processing of many terabytes of data over thousands of distributed servers in a time efficient manner.

One of the core elements in distributed learning algorithms is data shuffling . Before each iteration of the learning process, the entire data is randomly shuffled before being assigned to the worker nodes. This shuffling operation enables the worker nodes to process different data batches at each iteration, which presents large statistical gains . The statistical advantages provided by data shuffling come at the unavoidable cost of the communication overhead between the master and the worker nodes which must be incurred for every shuffling iteration. Thus, there exists a fundamental tradeoff between the communication overhead and the storage capacity of each worker node. To exemplify this tradeoff, consider two extreme scenarios: an ideal scenario in which the storage at the distributed workers is large enough to store the entire data-set, thus no communication has to be done from the master node for any shuffle. On the other extreme, when the storage is just enough to store the batch under processing, the communication load is expected to be maximum.

The focus of this paper is on formalizing and understanding this fundamental information-theoretic tradeoff between storage and the worst-case communication overhead for the data shuffling problem. Each iteration of data shuffling can be divided into two phases: data delivery, and storage update. In the data delivery phase, depending on the shuffled data points, the master node communicates a function of the data to the workers, so that each worker obtains its assigned data points. The second phase is termed as the storage update phase, which as shown in this paper is extremely critical in reducing the communication overhead of subsequently shuffling iterations. We next summarize the main contributions of this paper:

∙\bullet We first present an information-theoretic formulation of the data shuffling problem involving both data delivery and update phases, accounting for the respective constraints and formalizing the tradeoff between worst-case communication overhead and the storage capacity of distributed workers.

∙\bullet We also completely characterize this tradeoff for K=2K=2 and K=3K=3 workers, for any value of storage capacity. One of the most interesting aspects of the result for K=3K=3 worker problem is the design of data delivery and storage update algorithms. In particular, for data delivery phase, we show that transmitting coded data from the master node to the workers can significantly reduce the communication overhead. More interestingly, the proposed storage update algorithm maintains the structural properties of the storage at the workers over time. This structural invariance placement is extremely critical in leveraging the gains of coding for different shuffles.

Related Work: In the past few years, there has been a flurry of research acitivity in understanding the benefits of coding for caching starting from the work of Maddah-Ali and Niesen who showed that exploiting multi-casting opportunities by coding can reduce the communication for caching. In , coding for MapReduce was proposed in order to reduce the communication cost between mappers and reducers, however the underlying focus of is significantly different than the problem considered in this work, where we care about the communication between the master node and the workers. The paper most closely related to this work is , where the idea of coding for data shuffling problem is presented to reduce the communication overhead between the master node and worker nodes. provides a probabilistic scheme of leveraging coding based on a random storage placement. In contrast to , in this paper we provide a deterministic and systematic storage update scheme, which increases the coding opportunities in the delivery phase. The underlying metric used here is the worst-case communication cost over all the possible shuffles, unlike the average cost considered in . Finally, we also present the first information theoretic lower bounds on the communication overhead for the data shuffling problem.

II System Model

We assume a master node which has access to the entire data-set A=[x1T,x2T,…,xNT]TA=[x_{1}^{T},x_{2}^{T},\ldots,x_{N}^{T}]^{T} of size NdNd bits, i.e., AA is a matrix containing NN data points, denoted by x1,x2,…,xNx_{1},x_{2},\ldots,x_{N}, where dd is the dimensionality of each data point. Treating AA, and its data points xnx_{n} as random variables, we therefore have the entropies of these random variables as

At each iteration, indexed by tt, the master node divides the data-set AA into KK data batches given as A1t,A2t,…,AKtA^{t}_{1},A^{t}_{2},\ldots,A^{t}_{K}, where the batch AktA^{t}_{k} is designated to be processed by worker wkw_{k}, and these batches correspond to the random permutation of the data-set, πt:A→{A1t,…,AKt}\pi^{t}:A\rightarrow\{A^{t}_{1},\ldots,A^{t}_{K}\}. Note that these data chunks are disjoint, and span the whole data-set, i.e.,

Hence, the entropy of any batch AktA^{t}_{k} is given as

After getting the data batch, each worker locally computes a function (as an example, this function could correspond to the gradient or sub-gradients of the data points assigned to the worker). The local functions from the KK workers are processed subsequently at the master node. We assume that each worker wkw_{k} has a storage ZktZ^{t}_{k} of size SdSd bits, for a real number SS. For processing purposes, the assigned data blocks are needed to be stored by the workers, therefore, each worker wkw_{k} must at least store the data block AktA^{t}_{k} at time tt. If we consider ZktZ^{t}_{k} as a random variable then the storage constraint is given by

According to (3) and (4), we get the minimum storage per worker S≥NKS\geq\frac{N}{K}. We also have the processing constraint as

In the next epoch t+1t+1, the data-set is randomly reshuffled at the master node according to a random permutation πt+1:A→{A1t+1,A2t+1,…,AKt+1}\pi^{t+1}:A\rightarrow\{A^{t+1}_{1},A^{t+1}_{2},\ldots,A^{t+1}_{K}\}. The main communication bottleneck occurs during Data Delivery since the master node needs to communicate the new data batches to the workers. Trivially, if the storage (per worker) exceeds NdNd bits, i.e., S≥NS\geq N, then each worker can store the whole data-set, and no communication has to be done between the master node and the workers for any shuffle. Therefore from the constraint on minimum storage per worker, we can write the possible range for storage as NK≤S≤N\frac{N}{K}\leq S\leq N.

We next proceed to describe the data delivery mechanism, and the associated encoding and decoding functions. The main process can be divided into 2 phases, namely the data delivery phase and the storage update phase as described next:

At time t+1t+1, the master node sends a function of the data batches for the subsequent shuffles (πt,πt+1)(\pi_{t},\pi_{t+1}), X(πt,πt+1)=ϕ(A1t,…,AKt,A1t+1,…,AKt+1)=ϕ(πt,πt+1)(A)X_{(\pi_{t},\pi_{t+1})}=\phi(A^{t}_{1},\ldots,A^{t}_{K},A^{t+1}_{1},\ldots,A^{t+1}_{K})=\phi_{(\pi_{t},\pi_{t+1})}(A) over the shared link, where ϕ\phi is the data delivery encoding function

where R(πt,πt+1)R_{(\pi_{t},\pi_{t+1})} is the rate of the shared link based on the shuffles (πt,πt+1)(\pi_{t},\pi_{t+1}). Therefore, we have

Each worker wkw_{k} should decode the desired batch Akt+1A^{t+1}_{k} out of the transmitted function X(πt,πt+1)X_{(\pi_{t},\pi_{t+1})}, and the data stored in the previous time slot denoted as ZktZ^{t}_{k}. Therefore, the desired data is given by Akt+1=ψ(X(πt,πt+1),Zkt)A^{t+1}_{k}=\psi(X_{(\pi_{t},\pi_{t+1})},Z^{t}_{k}), where ψ\psi is the decoding function at the workers

which can be written in terms of a decodability constraint, at each worker as follows

II-B Storage Update Phase

At each iteration, every worker updates its stored content as follows: the new storage content Zkt+1Z^{t+1}_{k} is a function of the old storage content ZktZ^{t}_{k} as well as transmitted function X(πt,πt+1)X_{(\pi_{t},\pi_{t+1})}, i.e., Zkt+1=μ(X(πt,πt+1),Zkt)Z^{t+1}_{k}=\mu(X_{(\pi_{t},\pi_{t+1})},Z^{t}_{k}), where μ\mu is the update function

This implies the following storage-update constraint

The excess storage, if any, can be used to store opportunistically a function of the remaining data batches. Since the shuffling process at each time is done randomly, all the remaining batches are of equal importance. Consequently, the amount of excess storage, given by (S−NK)d(S-\frac{N}{K})d bits, is divided equally among the remaining K−1K-1 batches. For the scope of this work, we assume that the placement of the excess storage is uncoded, which means that (S−NK)K−1d\frac{(S-\frac{N}{K})}{K-1}d bits of the excess storage are dedicated to store a function of only one of the remaining K−1K-1 batches. We give the notation Ai,kt+1A^{t+1}_{i,k}, where i≠ki\neq k, as the part of data that worker wkw_{k} stores about Ait+1A^{t+1}_{i} in the excess storage at time t+1t+1. Considering Ai,kt+1A^{t+1}_{i,k} as a random variable, then

for k,i∈{1,…,K}k,i\in\{1,\ldots,K\}, and i≠ki\neq k.

We next define the worst-case communication as follows:

For any achievable scheme characterized by the functions (ϕ,ψ,μ)(\phi,\psi,\mu), the worst-case communication overhead over all possible consecutive data shuffles (πt,πt+1)(\pi_{t},\pi_{t+1}) is defined as

Our goal in this work is to characterize the optimal worst-case communication Rworst-case∗(K,S)R_{\textsf{worst-case}}^{*}(K,S) defined as

We next present a claim which shows that the optimal communication R∗R^{*} (for any shuffle including the worst-case) is a convex function of the storage SS:

R∗R^{*} is a convex function of SS, where SS is the available storage at each worker.

Claim 1 follows from a simple memory sharing argument which shows that for any two available storage values S1S_{1} and S2S_{2}, if (S1,R∗(K,S1))(S_{1},R^{*}(K,S_{1})), and (S2,R∗(K,S2))(S_{2},R^{*}(K,S_{2})) are achievable optimal schemes, then for any storage Sˉ=αS1+(1−α)S2\bar{S}=\alpha S_{1}+(1-\alpha)S_{2}, 0≤α≤10\leq\alpha\leq 1, there is a scheme which achieves a communication overhead of Rˉ=αR∗(K,S1)+(1−α)R∗(K,S2)\bar{R}=\alpha R^{*}(K,S_{1})+(1-\alpha)R^{*}(K,S_{2}).

This is done as follows: First, we divide the data-set AA across dd dimensions into 2 batches namely; A(α)A^{(\alpha)}, and A(1−α)A^{(1-\alpha)} of dimensions αd\alpha d, and (1−α)d(1-\alpha)d, for each point respectively. Then, we divide the storage for every worker wkw_{k} into 2 parts namely; Zk(α)Z_{k}^{(\alpha)}, and Zk(1−α)Z_{k}^{(1-\alpha)} of size S1αdS_{1}\alpha d, and S2(1−α)dS_{2}(1-\alpha)d, respectively. The former batch A(α)A^{(\alpha)} will be shuffled among the former part of the storage Zk(α)Z_{k}^{(\alpha)} to achieve the point (S1,R∗(K,S1))(S_{1},R^{*}(K,S_{1})), while the latter batch A(1−α)A^{(1-\alpha)} will be shuffled among the latter part of the storage Zk(1−α)Z_{k}^{(1-\alpha)} to achieve the point (S2,R∗(K,S2))(S_{2},R^{*}(K,S_{2})). Therefore, the total achievable load is given by

We next note that the optimal communication rate R∗(K,Sˉ)R^{*}(K,\bar{S}) is upper bounded by Rˉ(K,Sˉ)\bar{R}(K,\bar{S}), the rate of the memory sharing scheme, which completes the proof. ∎

III Main Results

For the distributed shuffling problem with K=2K=2 workers of storage SdSd bits each, and a data-set of size NdNd bits, the optimal communication versus storage tradeoff is given by

For a distributed shuffling problem system with K=3K=3 workers, the optimal communication versus storage tradeoff is given by

One interesting implication of Theorem 2 for K=3K=3 workers is that the corner point (2N3,N6)(\frac{2N}{3},\frac{N}{6}) as in Fig. 1 is better than memory sharing between the two points (N3,2N3)(\frac{N}{3},\frac{2N}{3}), and (N,0)(N,0), which falls on the line connecting the two points, i.e., memory sharing here is not optimal, and coding can be leveraged to reduce the communication overhead.

IV Proof of Theorem 1 (K=2𝐾2K=2 workers)

We start with the achievablity of the corner points S=NS=N and S=N2S=\frac{N}{2}. The point S=NS=N is trivial and represents the case when the workers can store the whole data-set. In this case, no communication is necessary, i.e., (S,R∗)=(N,0)(S,R^{*})=(N,0) is achievable.

The point corresponding to S=N2S=\frac{N}{2} is not trivial, where each worker can only store half of the data-set. Let us assume the data batches at time tt for the workers w1w_{1}, and w2w_{2} are A1tA^{t}_{1}, and A2tA^{t}_{2}, respectively. These batches should be stored at the corresponding workers, which are just enough to store them, i.e., Z1t=A1tZ^{t}_{1}=A^{t}_{1}, and Z2t=A2tZ^{t}_{2}=A^{t}_{2}. Recall that these batches are equally sized of N2d\frac{N}{2}d bits. For the next iteration t+1t+1, the data is randomly shuffled at the master node such that the new batches are A1t+1A^{t+1}_{1}, and A2t+1A^{t+1}_{2}, also of the same equal size N2d\frac{N}{2}d. In the delivery phase, the master node sends

worker w1w_{1} uses X(πt,πt+1)X_{(\pi_{t},\pi_{t+1})} and Z1t=A1tZ^{t}_{1}=A^{t}_{1} to decode A2tA^{t}_{2}, and get access to the whole data-set A={A1t,A2t}={A1t+1,A2t+1}A=\{A^{t}_{1},A^{t}_{2}\}=\{A^{t+1}_{1},A^{t+1}_{2}\}. The storage update is only storing the desired batch A1t+1A^{t+1}_{1}. The same procedure applies for w2w_{2}. Note that this choice of the transmitted function in (18) works for any possible shuffling, which gives a constant communication load H(X(πt,πt+1))=N2dH(X_{(\pi_{t},\pi_{t+1})})=\frac{N}{2}d, and hence the corner point (N2,N2)(\frac{N}{2},\frac{N}{2}) is achievable.

With the achievability of the corner points, any point in between for any real number SS can be simply achieved by memory sharing (see Claim 1). Therefore, the upper bound for the optimal worst-case communication rate is given by

If there is an overlap between AktA^{t}_{k} and Akt+1A^{t+1}_{k} for k∈{1,2}k\in\{1,2\}, then for S=N2S=\frac{N}{2}, the communication cost of the above scheme can be further improved by sending

where Akt∖Akt+1A^{t}_{k}\setminus A^{t+1}_{k} represents the part of the old batch AktA^{t}_{k} at worker wkw_{k}, which is not needed any more in the new batch Akt+1A^{t+1}_{k}. For an overlap ∣Akt∩Akt+1∣=b|A^{t}_{k}\cap A^{t+1}_{k}|=b data points, where bb is an integer number b∈{0,…,N2}b\in\{0,\ldots,\frac{N}{2}\}, then the we achieve the corner point (N2,N2−b)(\frac{N}{2},\frac{N}{2}-b), with the worst-case rate when b=0b=0.

IV-B Converse for K=2𝐾2K=2 workers

In this section, we present an information theoretic lower bound for the worst-case communication rate which matches the above scheme for K=2K=2 workers.

Since we do not know a priori the shuffle that gives the worst-case communication, we assume first a shuffle (πt,πt+1)(\pi_{t},\pi_{t+1}) with the optimal rate R(πt,πt+1)∗R_{(\pi_{t},\pi_{t+1})}^{*}, and then we lower bound R(πt,πt+1)∗R_{(\pi_{t},\pi_{t+1})}^{*}. Since the optimal worst-case rate is larger than the rate for any shuffle, i.e., Rworst-case∗≥R(πt,πt+1)∗R_{\textsf{worst-case}}^{*}\geq R_{(\pi_{t},\pi_{t+1})}^{*}, the lower bound found over R(πt,πt+1)∗R_{(\pi_{t},\pi_{t+1})}^{*} serves also as a lower bound for the worst-case communication Rworst-case∗R_{\textsf{worst-case}}^{*}. The novel part in our proof is choosing the right shuffle which leads to the optimal lower bound (also see Section V-B for the application of this idea to the converse proof for Theorem 22).

Now let us assume the data is shuffled such that A1t+1=A2tA^{t+1}_{1}=A^{t}_{2}, and A2t+1=A1tA^{t+1}_{2}=A^{t}_{1}, then from decodability constraint in (9):

Hence, we can find that H(A∣Z1t,X(πt,πt+1))=0H(A|Z^{t}_{1},X_{(\pi_{t},\pi_{t+1})})=0 as follows

where (a)(a) follows from (2), (b)(b) follows from the fact that H(A,B∣C)≤H(A∣C)+H(B∣C)H(A,B|C)\leq H(A|C)+H(B|C), (c)(c) is because conditioning reduces entropy, and (d)(d) is from (5) and (21). We next prove the lower bound as follows

where (a)(a) follows from (1), (b)(b), and (c)(c) follow from the fact that I(A;B)=H(A)−H(A∣B)=H(B)−H(B∣A)I(A;B)=H(A)-H(A|B)=H(B)-H(B|A) as well as (IV-B), (d)(d) is due to the fact that H(A,B)≤H(A)+H(B)H(A,B)\leq H(A)+H(B) and the fact that the Z1tZ^{t}_{1} and X(πt,πt+1)X_{(\pi_{t},\pi_{t+1})} are functions of the whole data-set AA, i.e., H(Z1t,X(πt,πt+1)∣A)=0H(Z^{t}_{1},X_{(\pi_{t},\pi_{t+1})}|A)=0, and (e)(e) follows from (4) and (7). Hence, from Remark 2, the lower bound for the worst-case rate is characterized as

Hence, the proof of Theorem 1 is complete from (19) and (24).

V Proof of Theorem 2 (K=3𝐾3K=3 workers)

From Fig. 1, the achievability involves three corner points: S=N3S=\frac{N}{3}, S=2N3S=\frac{2N}{3}, and S=NS=N. Achieving the point (N,0)(N,0) is trivial. The scheme for the point S=N3S=\frac{N}{3} is similar to the point S=N2S=\frac{N}{2} in the K=2K=2 worker case, where each worker can only store the desired data batches. Similar to (18) the transmission at time t+1t+1 is X(πt,πt+1)={A1t⊕A2t,A2t⊕A3t}X_{(\pi_{t},\pi_{t+1})}=\{A^{t}_{1}\oplus A^{t}_{2},A^{t}_{2}\oplus A^{t}_{3}\}, which is sufficient for each worker to access all the data-set and store what it needs, achieving the corner point (N3,2N3)(\frac{N}{3},\frac{2N}{3}).

In this section, we focus on the achievability of the corner point (2N3,N6)\left(\frac{2N}{3},\frac{N}{6}\right), which is perhaps the most interesting aspect of this result. For each worker wkw_{k} at time tt, k∈{1,2,3}k\in\{1,2,3\}, half of the storage (N/3N/3 points) is used to store the desired batch AktA^{t}_{k}, while the remaining half (excess storage of N/3N/3 points) is used opportunistically in order to minimize the communication overhead by storing some parts of the remaining batches given by the sub-batches Ai,kt,  i∈{1,2,3}∖kA^{t}_{i,k},\;i\in\{1,2,3\}\setminus k.

These parts will be formed as follows: if we consider the data batch AktA^{t}_{k} assigned for wkw_{k} of N3\frac{N}{3} points, each point x∈Aktx\in A^{t}_{k} is divided across dd dimensions into two equal subdivisions labeled as x(i),  i∈{1,2,3}∖kx^{(i)},\;i\in\{1,2,3\}\setminus k. Then, each one of these subdivisions x(i)x^{(i)} is placed in the corresponding sub-batch Ak,itA^{t}_{k,i} (the part of x∈Aktx\in A^{t}_{k} stored in the excess storage of wiw_{i}). For instance, A1tA^{t}_{1} (of size N3\frac{N}{3}) stored in the processing half of the storage of worker w1w_{1} is divided into two equal non-overlapping parts A1t={A1,2t,A1,3t}A^{t}_{1}=\{A^{t}_{1,2},A^{t}_{1,3}\} (of size N6\frac{N}{6} each). Worker w2w_{2} will use half of its excess storage to store A1,2tA^{t}_{1,2}, while worker w3w_{3} will use half of its excess storage to store A1,3tA^{t}_{1,3}.

The storage update procedure we present here maintains the above structural property of the stored data over time. The consequence of such structural invariance is that for any data point is required to be at a worker, at least half of this point is guaranteed to be already present at the worker, which decreases the communication overhead of the shuffling process.

Data Delivery: With this placement strategy, for the subsequent shuffles (πt,πt+1)(\pi_{t},\pi_{t+1}), the transmitted function is given as

We claim that the above transmission is sufficient for all the three workers to obtain the required new points for any shuffle. Without loss of generality, let us consider worker w1w_{1}. According to (25), for w1w_{1} to obtain the needed points not available in its storage, A1t+1∖Z1tA^{t+1}_{1}\setminus Z^{t}_{1}, it must have A2t+1∖Z2tA^{t+1}_{2}\setminus Z^{t}_{2} and A3t+1∖Z3tA^{t+1}_{3}\setminus Z^{t}_{3} in its storage Z1tZ_{1}^{t}. This is indeed the case and can be proved according to the following argument:

In the following, we prove that A2t+1∖Z2t∈Z1tA^{t+1}_{2}\setminus Z^{t}_{2}\in Z_{1}^{t}, and using the same argument we can show that A3t+1∖Z3t∈Z1tA^{t+1}_{3}\setminus Z^{t}_{3}\in Z_{1}^{t}. Consider a data point x∈A2t+1x\in A_{2}^{t+1}, which is newly assigned to worker w2w_{2} at time t+1t+1 and is not fully present in its storage Z2(t)Z_{2}^{(t)}, i.e., x∉A2tx\not\in A_{2}^{t}. Therefore, there are two possibilities: a) x∈A1tx\in A_{1}^{t} was being processed by worker w1w_{1} at time tt, which directly implies that xx is already available at w1w_{1}; or b) x∈A3tx\in A_{3}^{t} was being processed by worker w3w_{3} at time tt, which implies that the sub-divisions of xx are {x(1),x(2)}\{x^{(1)},x^{(2)}\} and the needed part by w2w_{2}, x(1)=x∩(A2t+1∖Z2t)x^{(1)}=x\cap(A^{t+1}_{2}\setminus Z^{t}_{2}), is available at w1w_{1} by definition (since x(2)x^{(2)} is the subdivision that is already present at w2w_{2}).

According to (25), the worst-case scenario for this scheme happens if there is no overlap between AktA^{t}_{k} and Akt+1A^{t+1}_{k} (completely new assignments). However, half of each data point x∈Akt+1x\in A^{t+1}_{k} is already stored in the excess storage at wkw_{k}, labeled x(k)x^{(k)}, then the worst-case rate of Akt+1∖ZktA^{t+1}_{k}\setminus Z^{t}_{k} and eventually X(πt,πt+1)X_{(\pi_{t},\pi_{t+1})} is N6\frac{N}{6}, achieving the corner point (2N3,N6)(\frac{2N}{3},\frac{N}{6}).

Storage Update: Now, we present a deterministic storage update strategy, which maintains the structural properties of the storage at time tt. Without loss of generality, let us analyze a data point x={x(2),x(3)}∈A1t+1x=\{x^{(2)},x^{(3)}\}\in A^{t+1}_{1} . The update procedure is done according to the following three cases at time tt

∙\bullet Case 1: x={x(2),x(3)}∈A1t{x=\{x^{(2)},x^{(3)}\}\in A^{t}_{1}}. In this case, since the point xx was already being processed at worker w1w_{1} at time tt, hence no storage update is necessary at time t+1t+1 for this point across workers.

∙\bullet Case 2: x={x(1),x(3)}∈A2tx=\{x^{(1)},x^{(3)}\}\in A^{t}_{2}, x(1)∈A2,1tx^{(1)}\in A^{t}_{2,1}, x(3)∈A2,3tx^{(3)}\in A^{t}_{2,3}. After receiving x(3)x^{(3)} from the delivery phase, worker w1w_{1} stores the full-point xx in A1t+1A^{t+1}_{1}. For worker w3w_{3}, x(3)x^{(3)} leaves A2,3tA^{t}_{2,3} and enters A1,3t+1A^{t+1}_{1,3}, and we can notice that x(3)x^{(3)} remains within the excess storage of Z3t+1Z_{3}^{t+1}. Worker w2w_{2} removes the data point xx from the processing batch A2t+1A^{t+1}_{2}, and stores x(1)x^{(1)} in A1,2t+1A_{1,2}^{t+1} after relabelling it as x(2)x^{(2)} (now stored at the excess storage of w2w_{2}), i.e., it simply moves one half of xx in its excess storage.

∙\bullet Case 3: x={x(1),x(2)}∈A3tx=\{x^{(1)},x^{(2)}\}\in A^{t}_{3}, x(1)∈A3,1tx^{(1)}\in A^{t}_{3,1}, x(2)∈A3,2tx^{(2)}\in A^{t}_{3,2}. This case is similar to case 2, where w1w_{1} stores xx in A1t+1A^{t+1}_{1} and removes x(1)x^{(1)} from A3,1t+1A^{t+1}_{3,1}, worker w2w_{2} moves x(2)x^{(2)} into A1,2t+1A^{t+1}_{1,2} instead of A3,2tA_{3,2}^{t}, and worker w3w_{3} removes xx from A3t+1A^{t+1}_{3} and stores x(1)x^{(1)} (now labeled as x(3)x^{(3)}) in A1,3t+1A^{t+1}_{1,3}.

Let us now take a representative example depicted in Fig. 2 to illustrate our proposed data-delivery and storage-update phases for the corner point (2N3,N6)(\frac{2N}{3},\frac{N}{6}). Consider a system with K=3K=3 workers, and N=3N=3 data points, {x1,x2,x3}\{x_{1},x_{2},x_{3}\}. We assume that storage per worker is S=2NK=2S=\frac{2N}{K}=2 points, i.e., each worker can store one extra data point in addition to the one under processing. We first clarify the color code used in this example to indicate the data point assigned to a certain worker labeled with these colors: blue for x1x_{1}, red for x2x_{2}, and yellow for x3x_{3}.

At time tt, consider the dataset is shuffled such that A1t=x1A_{1}^{t}=x_{1}, A2t=x2A_{2}^{t}=x_{2}, and A3t=x3A_{3}^{t}=x_{3}. The corresponding storage placement at time tt in Fig. 2(a) is as follows: After using half the storage to store the desired data point (which can be depicted in this example as the desired color), each data point is divided equally among the unintended workers (depicted in this example as the unintended colors). For example, if we take the batch A1t=x1={x1(2),x1(3)}A_{1}^{t}=x_{1}=\{x^{(2)}_{1},x^{(3)}_{1}\}, worker w2w_{2} stores A1,2t=x1(2)A_{1,2}^{t}=x_{1}^{(2)}, and worker w3w_{3} stores A1,3t=x1(3)A_{1,3}^{t}=x_{1}^{(3)}.

At time t+1t+1 in Fig. 2(b), the data is randomly shuffled again such that the new batches are: A1t+1=x3A_{1}^{t+1}=x_{3}, A2t+1=x1A_{2}^{t+1}=x_{1}, and A3t+1=x2A_{3}^{t+1}=x_{2}. We take this particular shuffle since it represents one of the possible worst cases, where every worker is assigned a completely different batch. According to the previous storage content shown in Fig. 2(a), worker w1w_{1} already has x3(1)x_{3}^{(1)} but still needs x3(2)x_{3}^{(2)}, which is stored at w2w_{2} and w3w_{3}. Similarly, worker w2w_{2} needs x1(3)x_{1}^{(3)} which is stored at w1w_{1} and w3w_{3}, and worker w3w_{3} needs x2(1)x_{2}^{(1)} which is stored at w1w_{1} and w2w_{2}. Following the data delivery as in (25), the master node transmits:

Each worker has two out of these three subdivisions, therefore it can decode the remaining needed one. The rate of this transmission is R=12R=\frac{1}{2}, which is N6\frac{N}{6} where N=3N=3.

For the storage update in Fig. 2(b), we only discuss the changes in A1t+1A^{t+1}_{1} and the corresponding A1,2t+1A^{t+1}_{1,2}, and A1,3t+1A^{t+1}_{1,3}. The storage update of the remaining parts can be done in a similar manner. From the delivery phase and the previous storage, worker w1w_{1} gets x3x_{3} (labeled yellow) and stores it in the processing half A1t+1A^{t+1}_{1} (above the dotted line). Worker w2w_{2} already has a part of A1t+1A^{t+1}_{1}, x3(2)x^{(2)}_{3} (labeled yellow), which was previously stored as A3,2tA^{t}_{3,2}. Therefore, it remains in the excess storage (below the dotted line) as A1,2t+1A^{t+1}_{1,2}. Worker w3w_{3} already has x3x_{3} previously labeled as A3tA^{t}_{3}, so it keeps in its excess storage the part that is not stored in A1,2t+1A^{t+1}_{1,2}, i.e., x3(1)x_{3}^{(1)}, to be stored in A1,3t+1A^{t+1}_{1,3}, after relabelling it to x3(3)x_{3}^{(3)} (now stored in w3w_{3}).

Using the achievability of the three corner points for K=3K=3, and Claim 1, we get an upper bound on Rworst-case∗(3,S)R^{*}_{\textsf{worst-case}}(3,S) as

V-B Converse for K=3𝐾3K=3 workers

We now present the information theoretic lower bounds for the three-worker case, which matches the above scheme. Following Remark 2, we first assume subsequent data shuffles at times {t,t+1,t+2}\{t,t+1,t+2\} such that A1t+1=A2tA^{t+1}_{1}=A^{t}_{2}, and A1t+2=A3tA^{t+2}_{1}=A^{t}_{3}, then from the decodability constraint in (9), we have

Hence, in a similar proof to (IV-B) we get

where (a)(a) follows from (2) and the storage update constraint in (11), (b)(b) follows from the fact that H(A,B,C∣D)≤H(A∣D)+H(B∣D)+H(C∣D)H(A,B,C|D)\leq H(A|D)+H(B|D)+H(C|D), and also the fact that conditioning reduces entropy, (c)(c) from (5) and (28), and (d)(d) from (28). Now, using (V-B) we can find the upper bound as follows

where (a)(a) follows from (1), (b)(b), and (c)(c) follow from (V-B), and due to the fact that I(A;B)=H(A)−H(A∣B)=H(B)−H(B∣A)I(A;B)=H(A)-H(A|B)=H(B)-H(B|A), (d)(d) is due to the fact that H(A,B,C)≤H(A)+H(B)+H(C)H(A,B,C)\leq H(A)+H(B)+H(C) and the fact that Z1tZ^{t}_{1}, X(πt,πt+1)X_{(\pi_{t},\pi_{t+1})}, and X(πt+1,πt+2)X_{(\pi_{t+1},\pi_{t+2})} are all functions of the whole data-set AA, and (e)(e) follows from (4) and (7). Hence, following Remark 2, we get a lower bound for the worst-case rate characterized as

The lower bound in (31) matches the upper bound obtained in (27) for the range 2N3≤S≤N\frac{2N}{3}\leq S\leq N.

We next present another lower bound on R∗(3,S)R^{*}(3,S) which proves the optimality of our scheme for the range N3≤S≤2N3\frac{N}{3}\leq S\leq\frac{2N}{3}. To this end, we now assume a data shuffle such that A1t+1=A3tA_{1}^{t+1}=A_{3}^{t}. Similar to (IV-B) and (V-B), we have

where (a)(a) follows from (2), (b)(b) from the facts that H(A,B,C∣D)≤H(A∣D)+H(B∣D)+H(C∣D)H(A,B,C|D)\leq H(A|D)+H(B|D)+H(C|D), and (c)(c) using (5), and (9). We now proceed to obtain the second lower bound on Rworst-case∗R^{*}_{\textsf{worst-case}} as follows

where (a)(a) follows from (1), (b)(b), and (c)(c) from (V-B), and due to the fact that I(A;B)=H(A)−H(A∣B)=H(B)−H(B∣A)I(A;B)=H(A)-H(A|B)=H(B)-H(B|A), (d)(d) from the chain rule of entropy and the fact that Z1tZ^{t}_{1}, Z2tZ^{t}_{2}, and X(πt,πt+1)X_{(\pi_{t},\pi_{t+1})} are all functions of the whole data-set AA, (e)(e) from (9) where A1t+1=A3tA_{1}^{t+1}=A_{3}^{t} must be decoded from Z1tZ_{1}^{t}, and X(πt,πt+1)X_{(\pi_{t},\pi_{t+1})}, and because A1tA_{1}^{t} and A2,1tA^{t}_{2,1} are stored within Z1tZ^{t}_{1}, (f)(f) because conditioning reduces entropy, (g)(g) because after obtaining A1tA_{1}^{t}, A3tA_{3}^{t}, and A2,1tA_{2,1}^{t}, the only remaining part in Z2tZ_{2}^{t} is A2,3tA_{2,3}^{t}, and finally (h)(h) follows from (4), (7), and (12). From Remark 2, and by rearranging (V-B), we get the following bound

Therefore, from (34), and (31), we get the following lower bound on Rworst-case∗(3,S)R^{*}_{\textsf{worst-case}}(3,S):

Finally, the proof of Theorem 2 follows from (27) and (35).

VI Conclusions

In this paper, we presented information theoretic formulation of the data shuffling problem, where we studied the tradeoff between the worst-case communication overhead and the storage available at the worker nodes. We completely characterized the optimal worst-case communication for K=2K=2, and K=3K=3 workers with any storage capacity, where we leveraged excess storage and coding to minimize the communication overhead in subsequent data shuffling iteration. A systematic storage update and delivery scheme was presented, which preserves the structural properties of the storage across workers. Generalizing these results for any number of workers (K>3)(K>3) is part of our ongoing work.

References