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:
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.
We also completely characterize this tradeoff for and workers, for any value of storage capacity. One of the most interesting aspects of the result for 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 of size bits, i.e., is a matrix containing data points, denoted by , where is the dimensionality of each data point. Treating , and its data points as random variables, we therefore have the entropies of these random variables as
At each iteration, indexed by , the master node divides the data-set into data batches given as , where the batch is designated to be processed by worker , and these batches correspond to the random permutation of the data-set, . Note that these data chunks are disjoint, and span the whole data-set, i.e.,
Hence, the entropy of any batch 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 workers are processed subsequently at the master node. We assume that each worker has a storage of size bits, for a real number . For processing purposes, the assigned data blocks are needed to be stored by the workers, therefore, each worker must at least store the data block at time . If we consider as a random variable then the storage constraint is given by
According to (3) and (4), we get the minimum storage per worker . We also have the processing constraint as
In the next epoch , the data-set is randomly reshuffled at the master node according to a random permutation . 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 bits, i.e., , 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 .
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 , the master node sends a function of the data batches for the subsequent shuffles , over the shared link, where is the data delivery encoding function
where is the rate of the shared link based on the shuffles . Therefore, we have
Each worker should decode the desired batch out of the transmitted function , and the data stored in the previous time slot denoted as . Therefore, the desired data is given by , where 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 is a function of the old storage content as well as transmitted function , i.e., , where 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 bits, is divided equally among the remaining batches. For the scope of this work, we assume that the placement of the excess storage is uncoded, which means that bits of the excess storage are dedicated to store a function of only one of the remaining batches. We give the notation , where , as the part of data that worker stores about in the excess storage at time . Considering as a random variable, then
for , and .
We next define the worst-case communication as follows:
For any achievable scheme characterized by the functions , the worst-case communication overhead over all possible consecutive data shuffles is defined as
Our goal in this work is to characterize the optimal worst-case communication defined as
We next present a claim which shows that the optimal communication (for any shuffle including the worst-case) is a convex function of the storage :
is a convex function of , where 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 and , if , and are achievable optimal schemes, then for any storage , , there is a scheme which achieves a communication overhead of .
This is done as follows: First, we divide the data-set across dimensions into 2 batches namely; , and of dimensions , and , for each point respectively. Then, we divide the storage for every worker into 2 parts namely; , and of size , and , respectively. The former batch will be shuffled among the former part of the storage to achieve the point , while the latter batch will be shuffled among the latter part of the storage to achieve the point . Therefore, the total achievable load is given by
We next note that the optimal communication rate is upper bounded by , the rate of the memory sharing scheme, which completes the proof. ∎
III Main Results
For the distributed shuffling problem with workers of storage bits each, and a data-set of size bits, the optimal communication versus storage tradeoff is given by
For a distributed shuffling problem system with workers, the optimal communication versus storage tradeoff is given by
One interesting implication of Theorem 2 for workers is that the corner point as in Fig. 1 is better than memory sharing between the two points , and , 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 and . The point is trivial and represents the case when the workers can store the whole data-set. In this case, no communication is necessary, i.e., is achievable.
The point corresponding to is not trivial, where each worker can only store half of the data-set. Let us assume the data batches at time for the workers , and are , and , respectively. These batches should be stored at the corresponding workers, which are just enough to store them, i.e., , and . Recall that these batches are equally sized of bits. For the next iteration , the data is randomly shuffled at the master node such that the new batches are , and , also of the same equal size . In the delivery phase, the master node sends
worker uses and to decode , and get access to the whole data-set . The storage update is only storing the desired batch . The same procedure applies for . Note that this choice of the transmitted function in (18) works for any possible shuffling, which gives a constant communication load , and hence the corner point is achievable.
With the achievability of the corner points, any point in between for any real number 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 and for , then for , the communication cost of the above scheme can be further improved by sending
where represents the part of the old batch at worker , which is not needed any more in the new batch . For an overlap data points, where is an integer number , then the we achieve the corner point , with the worst-case rate when .
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 workers.
Since we do not know a priori the shuffle that gives the worst-case communication, we assume first a shuffle with the optimal rate , and then we lower bound . Since the optimal worst-case rate is larger than the rate for any shuffle, i.e., , the lower bound found over serves also as a lower bound for the worst-case communication . 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 ).
Now let us assume the data is shuffled such that , and , then from decodability constraint in (9):
Hence, we can find that as follows
where follows from (2), follows from the fact that , is because conditioning reduces entropy, and is from (5) and (21). We next prove the lower bound as follows
where follows from (1), , and follow from the fact that as well as (IV-B), is due to the fact that and the fact that the and are functions of the whole data-set , i.e., , and 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: , , and . Achieving the point is trivial. The scheme for the point is similar to the point in the worker case, where each worker can only store the desired data batches. Similar to (18) the transmission at time is , which is sufficient for each worker to access all the data-set and store what it needs, achieving the corner point .
In this section, we focus on the achievability of the corner point , which is perhaps the most interesting aspect of this result. For each worker at time , , half of the storage ( points) is used to store the desired batch , while the remaining half (excess storage of points) is used opportunistically in order to minimize the communication overhead by storing some parts of the remaining batches given by the sub-batches .
These parts will be formed as follows: if we consider the data batch assigned for of points, each point is divided across dimensions into two equal subdivisions labeled as . Then, each one of these subdivisions is placed in the corresponding sub-batch (the part of stored in the excess storage of ). For instance, (of size ) stored in the processing half of the storage of worker is divided into two equal non-overlapping parts (of size each). Worker will use half of its excess storage to store , while worker will use half of its excess storage to store .
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 , 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 . According to (25), for to obtain the needed points not available in its storage, , it must have and in its storage . This is indeed the case and can be proved according to the following argument:
In the following, we prove that , and using the same argument we can show that . Consider a data point , which is newly assigned to worker at time and is not fully present in its storage , i.e., . Therefore, there are two possibilities: a) was being processed by worker at time , which directly implies that is already available at ; or b) was being processed by worker at time , which implies that the sub-divisions of are and the needed part by , , is available at by definition (since is the subdivision that is already present at ).
According to (25), the worst-case scenario for this scheme happens if there is no overlap between and (completely new assignments). However, half of each data point is already stored in the excess storage at , labeled , then the worst-case rate of and eventually is , achieving the corner point .
Storage Update: Now, we present a deterministic storage update strategy, which maintains the structural properties of the storage at time . Without loss of generality, let us analyze a data point . The update procedure is done according to the following three cases at time
Case 1: . In this case, since the point was already being processed at worker at time , hence no storage update is necessary at time for this point across workers.
Case 2: , , . After receiving from the delivery phase, worker stores the full-point in . For worker , leaves and enters , and we can notice that remains within the excess storage of . Worker removes the data point from the processing batch , and stores in after relabelling it as (now stored at the excess storage of ), i.e., it simply moves one half of in its excess storage.
Case 3: , , . This case is similar to case 2, where stores in and removes from , worker moves into instead of , and worker removes from and stores (now labeled as ) in .
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 . Consider a system with workers, and data points, . We assume that storage per worker is 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 , red for , and yellow for .
At time , consider the dataset is shuffled such that , , and . The corresponding storage placement at time 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 , worker stores , and worker stores .
At time in Fig. 2(b), the data is randomly shuffled again such that the new batches are: , , and . 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 already has but still needs , which is stored at and . Similarly, worker needs which is stored at and , and worker needs which is stored at and . 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 , which is where .
For the storage update in Fig. 2(b), we only discuss the changes in and the corresponding , and . The storage update of the remaining parts can be done in a similar manner. From the delivery phase and the previous storage, worker gets (labeled yellow) and stores it in the processing half (above the dotted line). Worker already has a part of , (labeled yellow), which was previously stored as . Therefore, it remains in the excess storage (below the dotted line) as . Worker already has previously labeled as , so it keeps in its excess storage the part that is not stored in , i.e., , to be stored in , after relabelling it to (now stored in ).
Using the achievability of the three corner points for , and Claim 1, we get an upper bound on 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 such that , and , then from the decodability constraint in (9), we have
Hence, in a similar proof to (IV-B) we get
where follows from (2) and the storage update constraint in (11), follows from the fact that , and also the fact that conditioning reduces entropy, from (5) and (28), and from (28). Now, using (V-B) we can find the upper bound as follows
where follows from (1), , and follow from (V-B), and due to the fact that , is due to the fact that and the fact that , , and are all functions of the whole data-set , and 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 .
We next present another lower bound on which proves the optimality of our scheme for the range . To this end, we now assume a data shuffle such that . Similar to (IV-B) and (V-B), we have
where follows from (2), from the facts that , and using (5), and (9). We now proceed to obtain the second lower bound on as follows
where follows from (1), , and from (V-B), and due to the fact that , from the chain rule of entropy and the fact that , , and are all functions of the whole data-set , from (9) where must be decoded from , and , and because and are stored within , because conditioning reduces entropy, because after obtaining , , and , the only remaining part in is , and finally 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 :
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 , and 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 is part of our ongoing work.