Provably Accelerated Randomized Gossip Algorithms

Nicolas Loizou, Michael Rabbat, Peter Richtárik

Introduction

Distributed averaging is a fundamental problem in the area of distributed computing and multi-agent systems . Randomized gossip algorithms are one of the most popular class of methods for solving it. The seminal 2006 paper of Boyd et al. on randomized gossip algorithms motivated a flurry of subsequent research, and now gossip algorithms appear in many applications, including distributed data fusion in sensor networks , load balancing and clock synchronization . The development and design of efficient gossip algorithms was studied extensively in the last decade. For a survey of selected relevant work prior to 2010, we refer the reader to the survey . For more recent results on randomized gossip algorithms we suggest . See also .

In the literature of gossip algorithms, an important task is the design of fast and efficient algorithms. Surprisingly, to the best of our knowledge, there are no variants of gossip algorithms that converge to consensus with an accelerated linear rate. In this work, our focus is precisely this. We design two provably accelerated randomized gossip protocols which converge to consensus fast.

The average consensus problem. In the average consensus (AC) problem we are given an undirected connected network G=(V,E)\mathcal{G}=(\mathcal{V},\mathcal{E}) with node set V={1,2,…,n}\mathcal{V}=\{1,2,\dots,n\} and edges E\mathcal{E}. Each node i∈Vi\in\mathcal{V} holds a local value ci∈Rc_{i}\in\mathcal{R}. The goal of AC is for every node to compute the average of these private values, cˉ:=1n∑ici\bar{c}:=\tfrac{1}{n}\sum_{i}c_{i}, in a distributed fashion. That is, the exchange of information can only occur between connected nodes (neighbors).

Main contributions. In this work, building upon a recent framework for the design and analysis of randomized gossip algorithms , we present two novel and provably accelerated randomized gossip protocols where in each step all nodes of the network update their values using their own information but only a pair of them exchange messages. The accelerated convergence rates of the proposed protocols are obtained by establishing a connection with the area of accelerated randomized Kaczmarz methods for solving consistent linear systems.

To the best of our knowledge, our protocols are the first randomized gossip algorithms that converge to consensus with an accelerated linear rate. The theoretical results are validated via computational testing on typical network topologies.

Structure of the paper. Section 2 introduces important technical preliminaries and the necessary background for understanding of our methods. Two accelerated variants of the randomized Kaczmarz (RK) method for solving linear systems and their theoretical convergence results are described. In Section 3 we present the two provably accelerated gossip protocols, along with some remarks on their implementation. Numerical evaluation of the new gossip protocols is presented in Section 4. Finally, concluding remarks are given in Section 5.

Notation. The following notational conventions are used in this paper. We write [n]:={1,2,…,n}[n]:=\{1,2,\dots,n\}. Boldface upper-case letters denote matrices; I\mathbf{I} is the identity matrix. By L\mathcal{L} we denote the solution set of the linear system Ax=b\mathbf{A}x=b, where A∈Rm×n\mathbf{A}\in\mathcal{R}^{m\times n} and b∈Rmb\in\mathcal{R}^{m}. By Ai:\mathbf{A}_{i:} and A:j\mathbf{A}_{:j} we indicate the ithi_{th} row and the jthj_{th} column of matrix A\mathbf{A}, respectively. Throughout the paper, x∗x^{*} is the projection of x0x^{0} onto L\mathcal{L} (that is, x∗x^{*} is the solution of the best approximation problem; see equation (2)). With λmin⁡+(⋅)\lambda_{\min}^{+}(\cdot) we indicate the smallest nonzero eigenvalue of matrix (⋅)(\cdot). ∥⋅∥\|\cdot\| and ∥⋅∥F\|\cdot\|_{F} are used to denote the Euclidean norm and the Frobenius norm, respectively. Finally, xk=(x1k,…,xnk)∈Rnx^{k}=(x^{k}_{1},\dots,x^{k}_{n})\in\mathcal{R}^{n} represents the vector with the local values of the nn nodes of the network at the kthk^{th} iteration. Here, xikx_{i}^{k} denotes the value of node i∈[n]i\in[n] at the kthk^{th} iteration.

Technical Preliminaries

In this section we present the connections between the randomized Kaczmarz methods for solving linear systems and the gossip algorithms for solving the AC problem, as discussed in more details in . In particular, we focus on the presentation of the two recently proposed accelerated variants of Kaczmarz methods and on their theoretical convergence analysis.

Kaczmarz-type methods are popular algorithms for solving linear systems Ax=b\mathbf{A}x=b with many equations. The randomized Kaczmarz method (RK) for solving consistent linear systems was first proposed and proved to converge with linear rate in . This work triggered much research into developing and analyzing randomized linear solvers and several improved variants of RK have been proposed .

In particular, in its simplest version, RK works as follows; In each step, one row Ai:\mathbf{A}_{i:} of matrix A\mathbf{A} is sampled with probability pi>0p_{i}>0 and then is used to obtain the next iterate by following the update rule:

For the case of consistent linear systems, it was shown that RK and its variants solves the following problem (known as best approximation problem) :

where x0x_{0} is the initial vector of the method.

In it was shown how RK works as a gossip algorithm when applied to a special linear system encoding the underlying network. The following definition is used to describe the class of linear systems considered here.

A linear system Ax=b\mathbf{A}x=b is called an “average consensus (AC) system” when all solutions xx satisfy that xi=xjx_{i}=x_{j} for all (i,j)∈E(i,j)\in\mathcal{E}.

Many linear systems satisfy the above definition. In this work we focus on the case where b=0b=0 and A∈R∣E∣×n\mathbf{A}\in\mathcal{R}^{|\mathcal{E}|\times n} is the incidence matrix of G\mathcal{G} (or its normalized form where ∥Ai:∥=1\|\mathbf{A}_{i:}\|=1). In this case, the row of the system corresponding to edge (i,j)(i,j) directly encodes the constraint xi=xjx_{i}=x_{j}.

Since the right hand side of the above system is b=0b=0, the update rule of equation (1) simplifies to: xk+1=xk−Ai:xk∥Ai:∥22Ai:⊤=[I−Ai:⊤Ai:∥Ai:∥22]xk.x^{k+1}=x^{k}-\tfrac{\mathbf{A}_{i:}x^{k}}{\|\mathbf{A}_{i:}\|_{2}^{2}}\mathbf{A}_{i:}^{\top}=\left[\mathbf{I}-\tfrac{\mathbf{A}_{i:}^{\top}\mathbf{A}_{i:}}{\|\mathbf{A}_{i:}\|_{2}^{2}}\right]x^{k}. In the case that the starting point is x0=cx_{0}=c it can be shown that RK solves the average consensus probem and that the above udpate rule is equivalent with the pairwise randomized gossip algorithm of (see for more details). The convergence performance of RK for solving the best approximation problem (and as a result the average consensus problem) is described by the following theorem.

2 Accelerated Kaczmarz methods

There are two different but very similar ways to accelerate the randomized Kaczmarz method. The first paper that proves asymptotic convergence with an accelerated linear rate is . The proof technique is similar to the framework developed by Nesterov in for the acceleration of coordinate descent methods. In a modified version for the selection of the parameters was proposed and a non-asymptotic accelerated linear rate was established. In Algorithm 1, pseudocode of the Accelerated Kaczmarz method (AccRK) is presented where both variants can be cast as special cases, by choosing the parameters with the correct way.

There are two options for selecting the parameters, which we describe next.

From : Choose λ∈[0,λmin⁡+(A⊤A)]\lambda\in[0,\lambda_{\min}^{+}(\mathbf{A}^{\top}\mathbf{A})] and set γ−1=0\gamma_{-1}=0. Generate the sequence {γk:k=0,1,…,K+1}\{\gamma_{k}:k=0,1,\dots,K+1\} by choosing γk\gamma_{k} to be the largest root of

and generate the sequences {αk:k=0,1,…,K+1}\{\alpha_{k}:k=0,1,\dots,K+1\} and {βk:k=0,1,…,K+1}\{\beta_{k}:k=0,1,\dots,K+1\} by setting

Choose the three sequences to be fixed constants as follows: βk=β=1−λmin⁡+(W)ν\beta_{k}=\beta=1-\sqrt{\frac{\lambda_{\min}^{+}(\mathbf{W})}{\nu}}, γk=γ=1λmin⁡+(W)ν\gamma_{k}=\gamma=\sqrt{\frac{1}{\lambda_{\min}^{+}(\mathbf{W})\nu}}, αk=α=11+γν∈(0,1)\alpha_{k}=\alpha=\frac{1}{1+\gamma\nu}\in(0,1) where W=A⊤Am\mathbf{W}=\frac{\mathbf{A}^{\top}\mathbf{A}}{m}.

3 Theoretical guarantees of AccRK

The two variants (Option 1 and Option 2) of AccRK are closely related, however their convergence analyses are different. Below we present the theoretical guarantees of the two options as presented in and .

Let {xk}k=0∞\{x^{k}\}_{k=0}^{\infty} be the sequence of random iterates produced by Algorithm 1 with the Option 1 for the parameters. Let λ∈[0,λmin⁡+(A⊤A)]\lambda\in[0,\lambda_{\min}^{+}(\mathbf{A}^{\top}\mathbf{A})] and define σ1=1+λ2m\sigma_{1}=1+\frac{\sqrt{\lambda}}{2m} and σ2=1−λ2m\sigma_{2}=1-\frac{\sqrt{\lambda}}{2m}. Then for any k≥0k\geq 0 we have that:

Note that as k→∞k\rightarrow\infty, we have that σ2k→0\sigma^{k}_{2}\rightarrow 0. This means that the decrease of the right hand side is governed mainly by the behavior of the term σ1\sigma_{1} in the denominator and as a result the method converge asymptotically with a decrease factor per iteration: σ1−2=(1+λ2m)−2≈1−λm.\sigma_{1}^{-2}=(1+\frac{\sqrt{\lambda}}{2m})^{-2}\approx 1-\frac{\sqrt{\lambda}}{m}.

Thus, by choosing λ=λmin⁡+\lambda=\lambda_{\min}^{+} and for the case that λmin⁡+\lambda_{\min}^{+} is small, Algorithm 1 will have significantly faster convergence rate than RK. Note that the above convergence results hold for normalized matrices A∈Rm×n\mathbf{A}\in\mathcal{R}^{m\times n}, that is matrices that have ∥Ai:∥=1\|\mathbf{A}_{i:}\|=1 for any i∈mi\in m.

Let W=A⊤Am\mathbf{W}=\frac{\mathbf{A}^{\top}\mathbf{A}}{m} and assume that Null(W)=Null(A){\rm Null}({\bf W})={\rm Null}({\bf A}). Let {xk,yk,vk}\{x^{k},y^{k},v^{k}\} be the iterates of Algorithm 1 with the Option 2 for the parameters. Then

The above result implies that Algorithm 1 converges linearly with rate 1−λmin⁡+(W)/ν1-\sqrt{\lambda_{\min}^{+}(\mathbf{W})/\nu}, which translates to a total of O(ν/λmin⁡+(W)log⁡(1/ϵ))O\left(\sqrt{\nu/\lambda_{\min}^{+}(\mathbf{W})}\log(1/\epsilon)\right) iterations to bring the quantity Ψk\Psi^{k} below ϵ>0\epsilon>0. It can be shown that 1≤ν≤1/λmin⁡+(W)1\leq\nu\leq 1/\lambda_{\min}^{+}(\mathbf{W}), (Lemma 2 in ) where ν\nu is as defined in (3). Thus, 1λmin⁡+(W)≤νλmin⁡+(W)≤1λmin⁡+(W)\sqrt{\frac{1}{\lambda_{\min}^{+}(\mathbf{W})}}\leq\sqrt{\frac{\nu}{\lambda_{\min}^{+}(\mathbf{W})}}\leq\frac{1}{\lambda_{\min}^{+}(\mathbf{W})} which means that the rate of AccRK (Option 2) is always better than that of the RK which (see Theorem 2.2) is equal to O(1/λmin⁡+(W)log⁡(1/ϵ)){\cal O}(1/\lambda_{\min}^{+}(\mathbf{W})\log(1/\epsilon)) for normalized matrices (∥A∥F2=m\|\mathbf{A}\|^{2}_{F}=m).

Accelerated randomized gossip algorithms

In the previous section we presented the complexity analysis guarantees of AccRK for solving consistent linear systems with normalized matrices. Now, let us explain how the two options of AccRK behave as gossip algorithms when they are used to solve the linear system Ax=0\mathbf{A}x=0 where A∈R∣E∣×n\mathbf{A}\in\mathcal{R}^{|\mathcal{E}|\times n} is the normalized incidence matrix of the network. That is, each row e=(i,j)e=(i,j) of A\mathbf{A} can be represented as (Ae:)⊤=12(ei−ej)(\mathbf{A}_{e:})^{\top}=\frac{1}{\sqrt{2}}(e_{i}-e_{j}) where eie_{i} (resp.eje_{j}) is the ithi^{th} (resp. jthj^{th}) unit coordinate vector in Rn\mathcal{R}^{n}.

By using this particular linear system, the expression Ai:yk−bi∥Ai:∥22Ai:⊤\tfrac{\mathbf{A}_{i:}y^{k}-b_{i}}{\|\mathbf{A}_{i:}\|_{2}^{2}}\mathbf{A}_{i:}^{\top} that appears in steps 8 and 9 of AccRK takes the following form when the row e=(i,j)∈Ee=(i,j)\in\mathcal{E} is sampled: Ae:yk−bi∥Ae:∥22Ae:⊤=b=0Ae:yk∥Ae:∥22Ae:⊤=form of Ayik−yjk2(ei−ej).\tfrac{\mathbf{A}_{e:}y^{k}-b_{i}}{\|\mathbf{A}_{e:}\|_{2}^{2}}\mathbf{A}_{e:}^{\top}\overset{b=0}{=}\tfrac{\mathbf{A}_{e:}y^{k}}{\|\mathbf{A}_{e:}\|_{2}^{2}}\mathbf{A}_{e:}^{\top}\overset{\text{form of A}}{=}\frac{y_{i}^{k}-y_{j}^{k}}{2}(e_{i}-e_{j}).

Let L\mathbf{L} be the Laplacian matrix of the network. For solving the above AC system (see Definition 2.1), the simple RK requires O((2mλmin⁡+(L))log⁡(1/ϵ))O\left((\frac{2m}{\lambda_{\min}^{+}(\mathbf{L})})\log(1/\epsilon)\right) iterations to achieve expected accuracy ϵ>0\epsilon>0. To understand the acceleration in the gossip framework this should be compared to the O(m2/λmin⁡+(L)log⁡(1/ϵ))O(m\sqrt{2/\lambda_{\min}^{+}(\mathbf{L})}\log(1/\epsilon)) of AccRK (Option 1) and the O(2mν/λmin⁡+(L)log⁡(1/ϵ))O(\sqrt{2m\nu/\lambda_{\min}^{+}(\mathbf{L})}\log(1/\epsilon)) of AccRK (Option 2).

The parameter λmin⁡+(L)\lambda^{+}_{\min}(\mathbf{L}) can be estimated by all nodes in a decentralized manner using the method described in . In order to implement this algorithm, we assume that all nodes have synchronized clocks and that they know the rate at which gossip updates are performed, so that inactive nodes also update their local values. This may not be feasible in all applications, but when it is possible (e.g., if nodes are equipped with inexpensive GPS receivers, or have reliable clocks) then they can benefit from the significant speedup achieved.

Related work on accelerated gossip algorithms: The idea of having gossip updates in a network with two registers in each node is not new. It was first proposed in and its analysis under strong conditions was presented in . There local memory is exploited by installing shift registers at each agent where the first register stores the agent’s current value and the second the agent’s value before the latest update. In , the Stochastic Heavy Ball (SHB) method is used for solving the AC problem and an accelerated method is proposed which was shown to be in practice faster than the algorithm of . is the first paper that presents gossip algorithms where in each step all nodes of the network update their values but only a subset of them exchange their private values.

Numerical Evaluation

We devote this section to numerically evaluate the performance of the proposed accelerated gossip protocols. In all of our experiments we compare the simple RK (equivalent to pairwise gossip algorithm of ) the Stochastic Heavy Ball method (SHB) proposed in and the AccRK (Algorithm 2) with the two options for the selection of the parameters presented in Section 2.2. In comparing the methods we use the relative error measure ∥xk−x∗∥2/∥x0−x∗∥2\|x^{k}-x^{*}\|^{2}/\|x^{0}-x^{*}\|^{2} where the starting vector of values x0=cx^{0}=c is taken to be always Gaussian vector. For all of our experiments the horizontal axis represents the number of iterations. The networks used in the experiments are the cycle (ring graph), the 2-dimension grid and the randomized geometric graph (RGG) with radius r=log⁡(n)/nr=\sqrt{\log(n)/n}. Code was written in Julia 0.6.3.

For the implementation of SHB we use the same parameters with the ones used in . For the AccRK (Option 1) we use λ=λmin⁡+(A⊤A)\lambda=\lambda_{\min}^{+}(\mathbf{A}^{\top}\mathbf{A}). Note that for all networks under study the two proposed protocols are faster than both the pairwise gossip algorithm of and the SHB of .

Conclusion and Future Research

We proposed novel provably accelerated randomized gossip algorithms for solving the AC problem. Our approach is based on connections established between the gossip algorithms and the Kaczmarz methods for solving linear systems. We believe that many novel and efficient gossip protocols can be discovered using results from the literature of Kaczmarz methods either by using different AC linear systems or using other Kaczmarz-type algorithms than the one presented in this manuscript. We speculate that the gossip algorithms presented in this work can be extended to the more general setting of minimizing the average of convex functions (1/n)∑i=1nfi(x)(1/n)\sum_{i=1}^{n}f_{i}(x) in a decentralized way . While preparing this work we become aware of where an accelerated gossip algorithm is developed for solving the dual of the best approximation problem (2) using the accelerated coordinate descent method of . A comparison of our protocols and the algorithm of is an ongoing research work.

References