Distributed Nesterov gradient methods over arbitrary graphs

Ran Xin, Dusan Jakovetic, Usman A. Khan

I Introduction

Distributed optimization has recently seen a surge of interest particularly with the emergence of modern signal processing and machine learning applications. A well-studied problem in this domain is finite sum minimization that also has some relevance to empirical risk formulations, i.e.,

Since the focus is on distributed implementation, the information exchange mechanism among the agents becomes a key ingredient of the solutions. Such inter-agent information exchange is modeled by a graph and significant work has focused on algorithm design under various graph topologies. The associated algorithms require two key steps: (i) consensus, i.e., reaching agreement among the agents; and, (ii) optimality, i.e., showing that the agreement is on the optimal solution. Naturally, consensus algorithms have been predominantly used as the basic building block of distributed optimization on top of which a gradient correction is added to steer the agreement to the optimal solution. Initial work thus follows closely the progress achieved in the consensus algorithms and extensions to various graph topologies, see e.g., .

Early work on consensus assumes doubly-stochastic (DS) weights , which require the underlying graphs to be undirected (or balanced) since both incoming and outgoing weights must sum to 11. The subsequent work on optimization over undirected graphs includes where the convergence is sublinear and with linear convergence. For directed (and unbalanced) graphs, it is not possible to construct DS weights, i.e., the weights can be chosen such that they sum to 11 either only on incoming edges or only on outgoing edges. Optimization over digraphs thus has been built on consensus with non-DS weights . Required now is a division with additional iterates that learn the non-1\mathbf{1} (where 1\mathbf{1} is a vector of all 11’s) Perron eigenvector of the underlying weight matrix, see for details. Such division causes significant conservatism and stability issues .

Recently, we introduced the AB\mathcal{AB} algorithm that removes the need of eigenvector learning by utilizing both row-stochastic (RS) and column-stochastic (CS) weights, simultaneously, . The algorithm thus is applicable to arbitrary strongly-connected graphs. The intuition behind using both sets of weights is as follows: Let AA be RS and BB be CS, with w⊤A=w⊤\mathbf{w}^{\top}A=\mathbf{w}^{\top} and Bv=vB\mathbf{v}=\mathbf{v}, in addition to being primitive. From Perron-Frobenius theorem, we have that A∞=1w⊤A^{\infty}=\mathbf{1}\mathbf{w}^{\top} and B∞=v1⊤B^{\infty}=\mathbf{v}\mathbf{1}^{\top}. Clearly, using AA or BB alone makes an algorithm dependent on the non-1\mathbf{1} Perron eigenvector (w\mathbf{w} or v\mathbf{v}) and thus the need for the aforementioned division by the iterates learning this eigenvector. Using AA and BB simultaneously, the asymptotics of AB\mathcal{AB} are driven by, loosely speaking, A∞B∞=(w⊤v)⋅11⊤A^{\infty}B^{\infty}=(\mathbf{w}^{\top}\mathbf{v})\cdot\mathbf{1}\mathbf{1}^{\top}, which recovers the consensus matrix, 11⊤\mathbf{1}\mathbf{1}^{\top}, without any scaling. It is shown in that AB\mathcal{AB} converges linearly to the optimal for smooth and strongly-convex functions.

In this letter, we study accelerated optimization over arbitrary graphs by extending AB\mathcal{AB} with Nesterov’s momentum. We first propose ABN\mathcal{ABN} that uses both RS and CS weights. Construct CS weights requires each agent to know at least its out-degree, which may not be possible in broadcast-type communication scenarios. To address this challenge, we provide an alternate algorithm, termed as FROZEN, that only uses RS weights. We show that FROZEN can be derived from ABN\mathcal{ABN} with the help of a simple state transformation. Finally, we note that a rigorous theoretical analysis is beyond the scope of this letter and we present extensive simulations to highlight and verify different aspects of the proposed methods.

We now describe the rest of this paper. Section II formulates the problem and recaps the AB\mathcal{AB} algorithm. Section III describes the two methods, ABN\mathcal{ABN} and FROZEN, and Section IV provides simulations comparing the proposed methods with the state-of-the-art in distributed optimization over both convex and strongly-convex functions, and over various digraphs.

II Problem Formulation and Preliminaries

Consider nn agents connected over a digraph, G=(V,E)\mathcal{G}=(\mathcal{V},\mathcal{E}), where V={1,⋯ ,n}\mathcal{V}=\{1,\cdots,n\} is the set of agents and E\mathcal{E} is the collection of edges, (i,j),i,j∈V(i,j),i,j\in\mathcal{V}, such that j→ij\rightarrow i. We define Ni\mboxin\mathcal{N}_{i}^{{\scriptsize\mbox{in}}} as the collection of in-neighbors of agent ii, i.e., the set of agents that can send information to agent ii. Similarly, Ni\mboxout\mathcal{N}_{i}^{{\scriptsize\mbox{out}}} is the set of out-neighbors of agent ii. Note that both Ni\mboxin\mathcal{N}_{i}^{{\scriptsize\mbox{in}}} and Ni\mboxout\mathcal{N}_{i}^{{\scriptsize\mbox{out}}} include node ii. The agents solve the following unconstrained optimization problem:

The graph, G\mathcal{G}, is strongly-connected.

Let FL1,1\mathcal{F}_{L}^{1,1} be the class of functions satisfying Assumption 3 and let Fμ,L1,1\mathcal{F}_{\mu,L}^{1,1} be the class of functions that satisfy both Assumptions 2 and 3; note that μ≤L\mu\leq L. In this letter, we propose distributed algorithms to solve Problem P1 for both function classes, i.e., F∈FL1,1F\in\mathcal{F}_{L}^{1,1} and F∈Fμ,L1,1F\in\mathcal{F}_{\mu,L}^{1,1}. We assume that the underlying optimization is solvable in the class FL1,1\mathcal{F}_{L}^{1,1}.

The gradient descent algorithm is given by

where βk\beta_{k} is the momentum parameter. For the function class FL1,1\mathcal{F}_{L}^{1,1}, choosing βk=kk+3\beta_{k}=\frac{k}{k+3} leads to an optimal oracle complexity of O(1ϵ)\mathcal{O}(\frac{1}{\sqrt{\epsilon}}), while for the function class Fμ,L1,1\mathcal{F}_{\mu,L}^{1,1}, βk=L−μL+μ\beta_{k}=\frac{\sqrt{L}-\sqrt{\mu}}{\sqrt{L}+\sqrt{\mu}} results into an optimal oracle complexity of O(Qlog⁡1ϵ)\mathcal{O}(\sqrt{\mathcal{Q}}\log\frac{1}{\epsilon}).

II-B Distributed Optimization: The 𝒜​ℬ𝒜ℬ\mathcal{AB} algorithm

When the objective functions are not available at a central location, distributed solutions are required to solve Problem P1. Most existing work is restricted to undirected graphs, since the weights assigned to neighboring agents must be doubly-stochastic. The work on directed graphs is largely based on push-sum consensus that requires eigenvector learning. Recently, AB\mathcal{AB} algorithm was introduced in that does not require eigenvector learning by utilizing a novel approach to deal with the non-doubly-stochasticity in digraphs.

We now describe the AB\mathcal{AB} algorithm: Consider two distinct sets of weights, {aij}\{a_{ij}\} and {bij}\{b_{ij}\}, at each agent such that

In other words, the weight matrix, A={aij}A=\{a_{ij}\}, is row-stochastic, while B={bij}B=\{b_{ij}\} is column-stochastic. It is straightforward to note that the construction of row-stochastic weights, AA, is trivial as it each agent ii on its own assigns arbitrary weights to incoming information (from agents in Ni\mboxin\mathcal{N}_{i}^{{\scriptsize\mbox{in}}}) such that these weights sum to 11. The construction of column-stochastic weights is more involved as it requires that all outgoing weights at agent ii must sum to 11 and thus cannot be assigned on incoming information. The simplest way to obtain such weights is for each agent ii to transmit ski/∣Ni\mboxout∣{\mathbf{s}_{k}^{i}}/{|\mathcal{N}_{i}^{{\scriptsize\mbox{out}}}|} to its outgoing neighbors in Ni\mboxout\mathcal{N}_{i}^{{\scriptsize\mbox{out}}}. This strategy, however, requires the knowledge of the out-degree at each agent ii.

With the help of the row- and column-stochastic weights, we can now describe the AB\mathcal{AB} algorithm as follows :

The AB\mathcal{AB} algorithm for undirected graphs where both weights are doubly-stochastic was studied earlier in . It is shown in that the oracle complexity with doubly-stochastic weights is O(Q2log⁡1ϵ)\mathcal{O}(Q^{2}\log\frac{1}{\epsilon}). Extensions of AB\mathcal{AB} include: non-coordinated step-sizes and heavy-ball momentum ; time-varying graphs ; analysis for non-convex functions . Related work on distributed Nesterov-type methods can be found in , which is restricted to undirected graphs. There is no prior work on Nesterov’s method that is applicable to arbitrary strongly-connected graphs.

III Distributed Nesterov Gradient Methods

In this section, ww introduce two distributed Nesterov gradient methods, both of which are applicable to arbitrary, strongly-connected, graphs.

A valid choice for bijb_{ij}’s at each ii is to choose them as 1/∣Ni\mboxout∣1/{|\mathcal{N}_{i}^{{\scriptsize\mbox{out}}}|}, which does not require knowing the outgoing nodes but only the out-degree. For the function class Fμ,L1,1\mathcal{F}^{1,1}_{\mu,L}, β\beta is a constant; for the function class FL1,1\mathcal{F}^{1,1}_{L}, we choose βk=kk+3\beta_{k}=\frac{k}{k+3}.

III-B The FROZEN algorithm

Note that ABN\mathcal{ABN} is restricted to communication protocols that allow column-stochastic weights, {bij}\{b_{ij}\}’s. When this is not possible, it is desirable to have algorithms that only use row-stochastic weights. Row-stochasticity is trivially established at the receiving agent by assigning a weight to each incoming information such that the sum of weights is 11. To avoid CS weights altogether, we now develop a distributed Nesterov gradient method that only row-stochastic weights and show the procedure of constructing this new algorithm from ABN\mathcal{ABN}.

To this aim, we first write ABN\mathcal{ABN} in the vector-matrix form. Let xk,yk\mathbf{x}_{k},\mathbf{y}_{k}, sk\mathbf{s}_{k}, and ∇f(xk)\nabla\mathbf{f}(\mathbf{x}_{k}) denote the concatenated vectors with xki\mathbf{x}_{k}^{i}’s, yki\mathbf{y}_{k}^{i}’s, ski\mathbf{s}_{k}^{i}’s, and ∇fi(xki)\nabla f_{i}(\mathbf{x}_{k}^{i})’s, respectively. Then ABN\mathcal{ABN} can be compactly written follows:

where A=A⊗Ip\mathcal{A}=A\otimes I_{p} and B=B⊗Ip\mathcal{B}=B\otimes I_{p}, where ⊗\otimes is the Kronecker. Since A\mathcal{A} is already row-stochastic, we seek a transformation that makes B\mathcal{B} a row-stochastic matrix. Since BB is column-stochastic, we denote its left and right Perron eigenvectors as 1n⊤B=1n⊤\mathbf{1}_{n}^{\top}{B}=\mathbf{1}_{n}^{\top} and Bv=vB\mathbf{v}=\mathbf{v}. Let \mboxdiag(v)\mbox{diag}(\mathbf{v}) denote a matrix with v\mathbf{v} on its main diagonal. With the help of V=\mboxdiag(v)⊗IpV=\mbox{diag}(\mathbf{v})\otimes I_{p}, we define a state transformation, s~k=V−1sk\widetilde{\mathbf{s}}_{k}=V^{-1}\mathbf{s}_{k}, and rewrite ABN\mathcal{ABN} as follows:

where A~=V−1BV\widetilde{\mathcal{A}}=V^{-1}\mathcal{B}V can be easily verified to be row-stochastic. Since v\mathbf{v} is the right Perron vector of A~\widetilde{\mathcal{A}}, it is not locally known to any agent and thus the above equations are not practically possible to implement. We thus add an independent eigenvector learning algorithm to the above set equations and obtain FROZEN (Fast Row-stochastic OptimiZation with Nesterov’s momentum) described in Algorithm 2. The momentum parameter is chosen the same way as in ABN\mathcal{ABN}.

Generalizations and extensions: The method we described to convert ABN\mathcal{ABN} to FROZEN leads to another variant of ABN\mathcal{ABN} with only CS weights, see for details. The resulting methods add Nesterov’s momentum to ADDOPT and Push-DIGing . Since these variants only require CS weights, AB\mathcal{AB} and ABN\mathcal{ABN} are preferable due to their faster convergence. It is further straightforward to conceive a time-varying implementation of ABN\mathcal{ABN} and FROZEN over gossip based protocols or random graphs, see e.g., the related work in on non-accelerated methods. Asynchronous schemes may also be derived following the methodologies studied in . Finally, we note that a rigorous theoretical analysis of AB\mathcal{AB} and ABN\mathcal{ABN} is beyond the scope of this letter. We thus rely on simulations to highlight and verify different aspects of the proposed methods.

IV Numerical Results

In this section, we numerically verify the convergence of the proposed algorithms, ABN\mathcal{ABN} and FROZEN, in this letter, and compare them with well-known solutions for distributed optimization. To this aim, we generate strongly-connected digraphs with n=30n=30 nodes using nearest-neighbor rules. We use an uniform weighting strategy to generate the row- and column-stochastic weight matrices, i.e., aij=1/∣Ni\mboxin∣,∀i,a_{ij}=1/|\mathcal{N}_{i}^{{\scriptsize\mbox{in}}}|,\forall i, and bij=1/∣Nj\mboxout∣,∀jb_{ij}=1/|\mathcal{N}_{j}^{{\scriptsize\mbox{out}}}|,\forall j. We first compare ABN\mathcal{ABN} and FROZEN with the following methods over digraphs: ADDOPT/Push-DIGing , FROST , and AB\mathcal{AB} . For comparison, we plot the average residual: 1n∑i=1n∥xi(k)−x∗∥2\frac{1}{n}\sum_{i=1}^{n}\|\mathbf{x}_{i}(k)-\mathbf{x}^{*}\|_{2}.

In our setting, the feature vectors, cij\mathbf{c}_{ij}’s, are generated from a Gaussian distribution with zero mean. The binary labels are generated from a Bernoulli distribution. We set p=10p=10 and mi=5,∀im_{i}=5,\forall i. The results are shown in Fig. 1. Although FROZEN is slower than ABN\mathcal{ABN}, it is applicable broadcast-based protocols as it only requires row-stochastic weights. The step-size and momentum parameters are manually chosen to obtain the best performance for each algorithm.

IV-B Non strongly-convex case

We next choose the objective functions, fif_{i}’s, to be smooth, convex but not strongly-convex. In particular, fi(x)=u(x)+bixf_{i}(x)=u(x)+b_{i}x, where bib_{i}’s are randomly generated, bn=−∑i=1n−1bib_{n}=-\sum_{i=1}^{n-1}b_{i}, and u(x)u(x) is chosen as follows:

It can be verified that f=∑ifif=\sum_{i}f_{i} is not strongly-convex as f′′(x∗)=0f^{{}^{\prime\prime}}(x^{*})=0. The results are shown in Fig. 2 where the momentum parameter is chosen as βk=kk+3\beta_{k}=\frac{k}{k+3} and other parameters are manually optimized.

IV-C Influence of graph sparsity

Finally, we study the influence of graph sparsity with the help of the logistic regression problem discussed earlier. We fix the number of nodes to n=30n=30 and randomly generate three nearest-neighbor digraphs, G1\mathcal{G}_{1}, G2\mathcal{G}_{2} and G3\mathcal{G}_{3}, with decreasing sparsity, see Fig. 3 (Top). In Fig. 3 (Bottom), we compare the performance of the proposed methods with centralized Nesterov over the three graphs. It can be verified that ABN\mathcal{ABN} and FROZEN approach centralized Nesterov method as the graphs become dense. FROZEN, however, is much slower than ABN\mathcal{ABN} because it additionally requires eigenvector learning.

V Conclusions

In this letter, we present accelerated methods for optimization based on Nesterov’s momentum over arbitrary, strongly-connected, graphs. The fundamental algorithm, ABN\mathcal{ABN}, uses both row- and column-stochastic weights, simultaneously, to achieve agreement and optimality. We then derive a variant from ABN\mathcal{ABN}, termed as FROZEN, that only uses row-stochastic weights and thus is applicable to a larger set of communication protocols, however, at the expense of eigenvector learning, thus resulting into slower convergence. Although a theoretical analysis is beyond the scope of this letter, we provide an extensive set of numerical results to study the behavior of the proposed methods for both convex and strongly-convex cases.

References