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 . 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 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- (where is a vector of all ’s) Perron eigenvector of the underlying weight matrix, see for details. Such division causes significant conservatism and stability issues .
Recently, we introduced the 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 be RS and be CS, with and , in addition to being primitive. From Perron-Frobenius theorem, we have that and . Clearly, using or alone makes an algorithm dependent on the non- Perron eigenvector ( or ) and thus the need for the aforementioned division by the iterates learning this eigenvector. Using and simultaneously, the asymptotics of are driven by, loosely speaking, , which recovers the consensus matrix, , without any scaling. It is shown in that converges linearly to the optimal for smooth and strongly-convex functions.
In this letter, we study accelerated optimization over arbitrary graphs by extending with Nesterov’s momentum. We first propose 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 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 algorithm. Section III describes the two methods, 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 agents connected over a digraph, , where is the set of agents and is the collection of edges, , such that . We define as the collection of in-neighbors of agent , i.e., the set of agents that can send information to agent . Similarly, is the set of out-neighbors of agent . Note that both and include node . The agents solve the following unconstrained optimization problem:
The graph, , is strongly-connected.
Let be the class of functions satisfying Assumption 3 and let be the class of functions that satisfy both Assumptions 2 and 3; note that . In this letter, we propose distributed algorithms to solve Problem P1 for both function classes, i.e., and . We assume that the underlying optimization is solvable in the class .
The gradient descent algorithm is given by
where is the momentum parameter. For the function class , choosing leads to an optimal oracle complexity of , while for the function class , results into an optimal oracle complexity of .
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, 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 algorithm: Consider two distinct sets of weights, and , at each agent such that
In other words, the weight matrix, , is row-stochastic, while is column-stochastic. It is straightforward to note that the construction of row-stochastic weights, , is trivial as it each agent on its own assigns arbitrary weights to incoming information (from agents in ) such that these weights sum to . The construction of column-stochastic weights is more involved as it requires that all outgoing weights at agent must sum to and thus cannot be assigned on incoming information. The simplest way to obtain such weights is for each agent to transmit to its outgoing neighbors in . This strategy, however, requires the knowledge of the out-degree at each agent .
With the help of the row- and column-stochastic weights, we can now describe the algorithm as follows :
The 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 . Extensions of 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 ’s at each is to choose them as , which does not require knowing the outgoing nodes but only the out-degree. For the function class , is a constant; for the function class , we choose .
III-B The FROZEN algorithm
Note that is restricted to communication protocols that allow column-stochastic weights, ’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 . 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 .
To this aim, we first write in the vector-matrix form. Let , , and denote the concatenated vectors with ’s, ’s, ’s, and ’s, respectively. Then can be compactly written follows:
where and , where is the Kronecker. Since is already row-stochastic, we seek a transformation that makes a row-stochastic matrix. Since is column-stochastic, we denote its left and right Perron eigenvectors as and . Let denote a matrix with on its main diagonal. With the help of , we define a state transformation, , and rewrite as follows:
where can be easily verified to be row-stochastic. Since is the right Perron vector of , 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 .
Generalizations and extensions: The method we described to convert to FROZEN leads to another variant of 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, and are preferable due to their faster convergence. It is further straightforward to conceive a time-varying implementation of 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 and 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, and FROZEN, in this letter, and compare them with well-known solutions for distributed optimization. To this aim, we generate strongly-connected digraphs with nodes using nearest-neighbor rules. We use an uniform weighting strategy to generate the row- and column-stochastic weight matrices, i.e., and . We first compare and FROZEN with the following methods over digraphs: ADDOPT/Push-DIGing , FROST , and . For comparison, we plot the average residual: .
In our setting, the feature vectors, ’s, are generated from a Gaussian distribution with zero mean. The binary labels are generated from a Bernoulli distribution. We set and . The results are shown in Fig. 1. Although FROZEN is slower than , 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, ’s, to be smooth, convex but not strongly-convex. In particular, , where ’s are randomly generated, , and is chosen as follows:
It can be verified that is not strongly-convex as . The results are shown in Fig. 2 where the momentum parameter is chosen as 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 and randomly generate three nearest-neighbor digraphs, , and , 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 and FROZEN approach centralized Nesterov method as the graphs become dense. FROZEN, however, is much slower than 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, , uses both row- and column-stochastic weights, simultaneously, to achieve agreement and optimality. We then derive a variant from , 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.