Communication-Optimal Distributed Clustering
Jiecao Chen, He Sun, David P. Woodruff, Qin Zhang
Introduction
Both the spectral clustering and the geometric clustering algorithms mentioned above have been widely used in practice, and have been the subject of extensive theoretical and experimental studies over the decades. However, these algorithms are designed for the centralized setting, and are not applicable in the setting of large-scale datasets that are maintained remotely by different sites. In particular, collecting the information from all the remote sites and performing a centralized clustering algorithm is infeasible due to high communication costs, and new distributed clustering algorithms with low communication cost need to be developed.
There are several natural communication models, and we focus on two of them: (1) a point-to-point model, and (2) a model with a broadcast channel. In the former, sometimes referred to as the message-passing model, there is a communication channel between each pair of users. This may be impractical, and the so-called coordinator model can often be used in place; in the coordinator model there is a centralized site called the coordinator, and all communication goes through the coordinator. This affects the total communication by a factor of two, since the coordinator can forward a message from one server to another and therefore simulate a point-to-point protocol. There is also an additional additive bits per message, where is the number of sites, since a server must specify to the coordinator where to forward its message. In the model with a broadcast channel, sometimes referred to as the blackboard model, the coordinator has the power to send a single message which is received by all sites at once. This can be viewed as a model for single-hop wireless networks.
In both models we study the total number of bits communicated among all sites. Although the blackboard model is at least as powerful as the message-passing model, it is often unclear how to exploit its power to obtain better bounds for specific problems. Also, for a number of problems the communication complexity is the same in both models, such as computing the sum of length- bit vectors modulo two, where each site holds one bit vector , or estimating large moments . Still, for other problems like set disjointness it can save a factor of in the communication .
We present algorithms for graph clustering: for any -vertex graph whose edges are arbitrarily partitioned across sites, our algorithms have communication cost in the message passing model, and have communication cost in the blackboard model, where the notation suppresses polylogarithmic factors. The algorithm in the message passing model has each site send a spectral sparsifier of its local data to the coordinator, who then merges them in order to obtain a spectral sparsifier of the union of the datasets, which is sufficient for solving the graph clustering problem. Our algorithm in the blackboard model is technically more involved, as we show a particular recursive sampling procedure for building a spectral sparsifier can be efficiently implemented using a broadcast channel. It is unclear if other natural ways of building spectral sparsifiers can be implemented with low communication in the blackboard model. Our algorithms demonstrate the surprising power of the blackboard model for clustering problems. Since our algorithms compute spectral sparsifiers, they also have applications to solving symmetric diagonally dominant linear systems in a distributed model. Any such system can be converted into a system involving a Laplacian (see, e.g., ), from which a spectral sparsifier serves as a good preconditioner.
Next we show that bits of communication is necessary in the message passing model to even recover a constant fraction of a cluster, and bits of communication is necessary in the blackboard model. This shows the optimality of our algorithms up to poly-logarithmic factors.
We then study clustering problems in constant-dimensional Euclidean space. We show for any , computing a -approximation for -median, -means, or -center correctly with constant probability in the message passing model requires bits of communication. We then strengthen this lower bound, and show even for bicriteria clustering algorithms, which may output a constant factor more clusters and a constant factor approximation, our bit lower bound still holds. Our proofs are based on communication and information complexity. Our results imply that existing algorithms for -median and -means with bits of communication, as well as the folklore parallel guessing algorithm for -center with bits of communication, are optimal up to poly-logarithmic factors. For the blackboard model, we present an algorithm for -median and -means that achieves an -approximation using bits of communication. This again separates the models.
We give empirical results which show that using spectral sparsifiers preserves the quality of spectral clustering surprisingly well in real-world datasets. For example, when we partition a graph with over million edges (the Sculpture dataset) into sites, only of the input edges are communicated in the blackboard model and are communicated in the message passing model, while the values of the normalized cut (the objective function of spectral clustering) given in those two models are at most larger than the ones given by the centralized algorithm, and the visualized results are almost identical. This is strong evidence that spectral sparsifiers can be a powerful tool in practical, distributed computation. When the number of sites is large, the blackboard model incurs significantly less communication than the message passing model, e.g., in the Twomoons dataset when there are sites, the message passing model communicates times as many edges as communicated in the blackboard model, illustrating the strong separation between these models that our theory predicts.
2 Related work
There is a rich literature on spectral and geometric clustering algorithms from various aspects (see, e.g., ). Balcan et al. and Feldman et al. study distributed -means ( also studies -median), and present provable guarantees on the clustering quality. Very recently Guha et al. studied distributed -median/center/means with outliers. Cohen et al. study dimensionality reduction techniques for the input data matrices that can be used for distributed -means. The main takeaway is that there is no previous work which develops protocols for spectral clustering in the common message passing and blackboard models, and lower bounds are lacking as well. For geometric clustering, while upper bounds exist (e.g., ), no provable lower bounds in either model existed, and our main contribution is to show that previous algorithms are optimal. We also develop a new protocol in the blackboard model.
Preliminaries
For two sets and , the symmetric difference of and is defined as .
2 Spectral sparsification
For any undirected and weighted graph , we say a subgraph of with proper reweighting of the edges is a -spectral sparsifier if
The following lemma shows that a spectral sparsifier preserves the clustering structure of a graph.
Let be a -spectral sparsifier of for some . Then, it holds for any set that .
which implies that for any subset .
where the last inequality holds by assuming . Similarly, we have that
Hence, and differ by at most a factor of 2 for any vertex set . ∎
3 Models of computation
We study distributed clustering in two models for distributed data: the message passing model and the blackboard model. The message passing model represents those distributed computation systems with point-to-point communication, and the blackboard model represents those where messages can be broadcast to all parties.
More precisely, in the message passing model there are sites , and one coordinator. These sites can talk to the coordinator through a two-way private channel. In fact, this is referred to as the coordinator model in Section 1, where it is shown to be equivalent to the point-to-point model up to small factors. The input is initially distributed at the sites. The computation is in terms of rounds: at the beginning of each round, the coordinator sends a message to some of the sites, and then each of those sites that have been contacted by the coordinator sends a message back to the coordinator. At the end, the coordinator outputs the answer. In the alternative blackboard model, the coordinator is simply a blackboard where these sites can share information; in other words, if one site sends a message to the coordinator/blackboard then all the other sites can see this information without further communication. The order for the sites to speak is decided by the contents of the blackboard.
For both models we measure the communication cost as the total number of bits sent through the channels. The two models are now standard in multiparty communication complexity (see, e.g., ). They are similar to the congested clique model studied in the distributed computing community; the main difference is that in our models we do not post any bandwidth limitations at each channel but instead consider the total number of bits communicated.
4 Communication complexity
For any problem and any protocol solving , the communication complexity of a protocol is the maximum communication cost of over all possible inputs . When the protocol is randomised, we define the error of by
where the is over all inputs and the probability is over all random strings of the coordinator and sites. The -error randomised communication complexity of a problem in the message passing model is the minimum communication complexity of any randomised protocol that solves with error at most .
Let be an input distribution on . We call a deterministic protocol if it gives the correct answer for on at least a fraction of all input pairs, weighted by the distribution . We denote as the cost of the minimum-communication protocol. A standard lemma in communication complexity called Yao’s minimax lemma shows that
5 Information complexity
We abuse notation by using for both the protocol and its transcript (its concatenation of messages). In the message passing model, let be the transcript (set of messages exchanged) between the -th site and the coordinator. Then can be seen as a concatenation ordered by the timestamps of the messages. We define the information complexity of a problem in the message passing model by
where is the mutual information function. It has been shown in that for any input distribution .
Distributed graph clustering
In this section we study distributed graph clustering. We assume that the vertex set of the input graph can be partitioned into clusters, where vertices in each cluster are highly connected to each other, and there are fewer edges between and . To formalize this notion, we define the -way expansion constant of graph by
Notice that a graph has clusters if the value of is small. It was shown in that closely relates to by the following higher-order Cheeger inequality:
Hence, a large gap between and implies (i) the existence of a -way partition such that every has small conductance , and (ii) any -way partition of contains a subset with high conductance . Therefore, a large gap between and ensures that has exactly clusters. In the following, we assume that
to ensure that the input graph has exactly clusters. The same assumption has been used in the literature for studying graph clustering in the centralized setting .
(iii) run -means on the embedded points , and group the vertices of into clusters according to the output of -means.
We assume the edges of the input graph are arbitrarily allocated among sites , and we use to denote the edge set maintained by site . Our proposed algorithm consists of two steps: (i) every computes a linear-sized -spectral sparsifier of , for a small constant , and sends the edge set of , denoted by , to the coordinator; (ii) the coordinator runs a spectral clustering algorithm on the union of received graphs . The theorem below summarizes the performance of this algorithm, and shows the approximation guarantee of this algorithm is as good as the provable guarantee of spectral clustering known in the centralized setting, which is shown in the lemma below.
Let be an -vertex graph with , and suppose the edges of are arbitrarily allocated among sites. Assume is an optimal partition that achieves . Then, the algorithm above computes a partition satisfying for any . The total communication cost of this algorithm is bits.
By the definition of the Laplacian matrix, we have that . Since every is a -spectral sparsifier of graph , we have that . This implies that , by the definition of and graph Laplacians. Now we show that our assumption on is preserved in . By Lemma 2.1, we have for any that , which implies that has low conductance in , and . To show that is a constant approximation of , notice that
Since and , we have that , and the assumption on in is preserved from up to a constant factor. By Lemma 3.1, the output of a spectral clustering algorithm on satisfies the claimed properties. The total communication cost of bits follows from the fact that every has edges. ∎
Next we show that the communication cost of our proposed algorithm is optimal up to a logarithmic factor. Our analysis is based on a reduction from graph clustering to the Multiparty Set-Disjointness problem (): for any sites , where each has a set , let be the characteristic vector of , and let be the input matrix with being the -th row. Let be the -th column of the input matrix . We define a function on an -bit vector as \text{{\mathsf{ALLONE}_{s}}}(Y)=\bigwedge_{i\in[s]}Y_{i}, and \text{{\mathsf{DISJ}_{s,n}}}(X)=\bigvee_{j\in[n]}\text{{\mathsf{ALLONE}_{s}}}(X^{j}). Then the problem asks the value of \text{{\mathsf{DISJ}_{s,n}}}(X). We introduce two hard input distributions for and respectively.
Hard input distribution on for : with probability , we choose each to be or with equal probability; with probability we choose to be an all- vector; and with the remaining probability we choose to be a random vector with coordinates being ’s and a random coordinate being .
Hard input distribution on for : For each , we choose .
It holds that \mathsf{IC}_{0.49,\nu}(\text{{\mathsf{ALLONE}_{s}}})=\Omega(s), and \mathsf{IC}_{0.49,\nu}(\text{{\mathsf{DISJ}_{s,n}}})=\Omega(sn).
In the message passing model, any randomized algorithm that computes correctly with probability needs bits of communication.
The lemma follows from Theorem 3.3 and Yao’s minimax lemma. ∎
Let be an undirected graph with vertices, and suppose the edges of are distributed among sites. Then, any algorithm that correctly outputs a constant fraction of a cluster in requires bits of communication. This holds even if each cluster has constant expansion.
As a remark, it is easy to see that this lower bound also holds for constructing spectral sparsifiers: for any matrix whose entries are arbitrarily distributed among sites, any distributed algorithm that constructs a -spectral sparsifier of requires bits of communication. This follows since such a spectral sparsifier can be used to solve the spectral clustering problem. Spectral sparsification has played an important role in designing fast algorithms from different areas, e.g., machine learning, and numerical linear algebra. Hence our lower bound result for constructing spectral sparsifiers may have applications to studying other distributed learning algorithms.
2 The blackboard model
Next we present a graph clustering algorithm with bits of communication cost in the blackboard model. Our result is based on the observation that a spectral sparsifier preserves the structure of clusters, which was used for proving Theorem 3.2. So it suffices to design a distributed algorithm for constructing a spectral sparsifier in the blackboard model.
where and . Notice that in the chain above every is obtained by adding weights to the diagonal entries of , and approximates as long as the weights added to the diagonal entries are small. We will construct this chain recursively, so that has heavy diagonal entries and can be approximated by a diagonal matrix. Moreover, since is the Laplacian matrix of a graph , it is easy to see that as long as the edge weights of are polynomially upper-bounded in .
Based on Lemma 3.6, we will construct a chain of matrices
Let be an undirected graph on vertices, where the edges of are allocated among sites, and the edge weights are polynomially upper bounded in . Then, a spectral sparsifier of can be constructed with bits of communication in the blackboard model. That is, the chain (3) can be constructed with bits of communication in the blackboard model.
First of all, notice that , and the value of can be obtained with communication cost (different sites sequentially write the new IDs of the vertices on the blackboard). In the following we assume that is the upper bound of that we actually obtained in the blackboard.
Combining Theorem 3.7 and the fact that a spectral sparsifier preserves the structure of clusters, we obtain a distributed algorithm in the blackboard model with total communication cost bits, and the performance of our algorithm is the same as in the statement of Theorem 3.2. Notice that bits of communication are needed for graph clustering in the blackboard model, since the output of a clustering algorithm contains bits of information and each site needs to communicate at least one bit. Hence the communication cost of our proposed algorithm is optimal up to a poly-logarithmic factor.
Distributed geometric clustering
We now consider geometric clustering, including -median, -means and -center. Let be a set of points of size in a metric space with distance function , and let be an integer. In the -center problem we want to find a set such that is minimized, where . In -median and -means we replace the objective function with and , respectively.
As mentioned, for constant dimensional Euclidean space and a constant , there are algorithms that -approximate -median and -means using bits of communication . For -center, the folklore parallel guessing algorithms (see, e.g., ) achieve a -approximation using bits of communication.
The following theorem states that the above upper bounds are tight up to logarithmic factors. The proof uses tools from multiparty communication complexity. We in fact can prove a stronger statement that any algorithm that can differentiate whether we have points or points in total in the message passing model needs bits of communication.
For any , computing -approximation for -median, -means or -center correctly with probability in the message passing model needs bits of communication.
A number of works on clustering consider bicriteria solutions (e.g., ). An algorithm is a -approximation if the optimal solution costs when using centers, then the output of the algorithm costs at most when using at most centers. We can show that for -median and -means, the lower bound holds even for algorithms with bicriteria approximations.
For any , computing -bicriteria-approximation for -median or -means correctly with probability in the message passing model needs bits of communication.
Before proving Theorem 4.2, we first show the following technical lemma.
By a Markov inequality, there must exist coordinates such that the algorithm computes with error probability at most . Call each of these coordinates good. Let be the protocol transcript. We have
By Lemma 4.3, we have that for any , computing -bicriteria-approximation for -median or -means in the message passing model correctly with probability under distribution needs bits of communication. The theorem follows by Yao’s minimax principle. ∎
2 The blackboard model
We can show that there is an algorithm that achieves an -approximation using bits of communication for -median and -means. For -center, it is straightforward to implement the parallel guessing algorithm in the blackboard model using bits of communication.
Our algorithm for -median/means is an easy adaptation of the successive sampling algorithm proposed by Mettu and Plaxton in the (centralized) RAM model. We first summarize their algorithm and then describe how to port it to the blackboard model.
Let be the point sets at sites respectively. The successive sampling algorithm proceeds in rounds. At each round it does the following:
sites jointly sample point centers, denoted by ;
sites grow balls from each of the point centers in synchronously until a time step when a fraction of points in are covered;
each site updates by removing those points that are covered by any of the balls centered at points in ;
sites remove all the points covered by balls centered at points in , and proceed to the next round .
It is easy to see that the computation will finish in rounds since at each round we remove a constant fraction of points. At the end we compute an -approximation of -median or -means on the points . In it has been shown that this algorithm gives an -approximation to -median or -means with high probability.
We now describe how to implement this centralized algorithm in the blackboard model. We first consider each round. Step 1 can be done by the distributed sampling algorithm in using bits of communication; note that at the end of this step the sampled points in are written on the blackboard. Step can be done by a binary search for the minimum ball radius such that covers at least a fraction of points in , where denotes the ball centered at with radius ; this binary search can be done using bits of communication. Step and can be done locally without any communication. After rounds, the final clustering step can be done by any of the sites since all points in have already been written on the blackboard. Therefore the total communication cost can be bounded by .
Finally, we would like to mention that is an obvious lower bound, and thus our upper bound is tight up to logarithmic factors. To see this, notice that is the size of the output, and the coordinator has to communication with each of the sites for at least bit.
Experiments
In this section we present experimental results for graph clustering in the message passing and blackboard models. We will compare the following three algorithms. (1) Baseline: each site sends all the data to the coordinator directly; (2) MsgPassing: our algorithm in the message passing model (Section 3.1); (3) Blackboard: our algorithm in the blackboard model (Section 3.2).
Besides giving the visualized results of these algorithms on various datasets, we also measure the qualities of the results via the normalized cut, defined as
which is a standard objective function to be minimized for spectral clustering algorithms.
We implemented the algorithms using multiple languages, including Matlab, Python and C++. Our experiments were conducted on an IBM NeXtScale nx360 M4 server, which is equipped with 2 Intel Xeon E5-2652 v2 8-core processors, 32GB RAM and 250GB local storage.
We test the algorithms in the following real and synthetic datasets, which is visualized in Figure 1.
In the distributed model edges are randomly partitioned across sites.
2 Results on clustering quality
We visualize the clustered results for the Twomoons, Gauss and Sculpture in Figure 2. It can be seen that Baseline, MsgPassing and Blackboard give results of very similar qualities. For simplicity, here we only present the visualization for . Similar results were observed when we varied the values of .
We also compare the normalized cut (ncut) values of the clustering results of different algorithms. The results are presented in Figure 3. In all datasets, the ncut values of different algorithms are very close. The ncut value of MsgPassing slightly decreases when we increase the value of , while the ncut value of Blackboard is independent of .
3 Results on communication costs
We compare the communication costs of different algorithms in Figure 4. We observe that while achieving similar clustering qualities as Baseline, both MsgPassing and Blackboard are significantly more communication-efficient (by one or two orders of magnitudes in our experiments). We also notice that the value of does not affect the communication cost of Blackboard, while the communication cost of MsgPassing grows almost linearly with ; when is large, MsgPassing uses significantly more communication than Blackboard. These confirm our theory.
4 Parameters in MsgPassing and Blackboard
Figure 5 shows in MsgPassinghow the value of ncut is affected by the number of sites and the number of edges sampled in each site. Here, each site samples edges. When and , the ncut value diverges in all datasets. This is because with such a small , the algorithm does not generate a valid sparsifier. In general, increasing or will slightly decrease the ncut value. But once they are above some thresholds, the ncut values of MsgPassing and Baseline become very close.
Figure 6 shows in Blackboardhow the ncut value is affected by the number of iterations and the number of edges sampled. When the number of iterations is set to be , ncut values diverge in all datasets. This is because we cannot expect to generate a valid sparsifier by using such few iterations. It can be seen from 6(b) that for a fixed , performing more iterations will help to reduce ncut values. From the same figure, one can also conclude that for fixed iterations, increasing also helps to reduce the ncut values.