Distributed Optimization Based on Gradient-tracking Revisited: Enhancing Convergence Rate via Surrogation
Ying Sun, Amir Daneshmand, Gesualdo Scutari
Introduction
We study distributed optimization over networks in the form:
Distributed optimization in the form (P) has found a wide range of applications in several areas, including network information processing, telecommunications, multi-agent control, and machine learning. An instance of particular interest to this work is the distributed Empirical Risk Minimization (ERM) whereby the goal is to minimize the average loss over some dataset, distributed across the nodes of the network (cf. Sec. 2.1.2). Letting the dataset of examples available at node ’s side, the local empirical loss reads , where measures the fit between the parameter and the sample . Data sets are usually large and high-dimensional, which makes routing local data to other agents (let alone to a centralized node) infeasible or highly inefficients. Given the cost of communications (especially if compared with the speed of local processing), the challenge in such a network setting is designing communication efficient distributed algorithms.
Motivated by the aforementioned applications, our focus pertains to such a design in two possible settings (one being a special case of the other) : 1) The scenario where no significant relationship can be assumed among the local functions –this is what the literature of distributed optimization has extensively studied, and will be refereed to as the unrelated setting—and 2) the case where the ’s are related, e.g., because they reflect statistical similarity in the data residing at different nodes. For instance, in the distributed ERM problem above, when data are i.i.d. among machines, one can show that quantities such as the gradients and Hessian matrices of the local functions differ only by , due to concentrations of measure effects –we will refer to this as -related setting (cf. Sec. 2.1.2). If properly exploited in the algorithmic design, such similarity can speed up the optimization/learning process over general purpose optimization algorithms.
Problem (P) in the two settings above has been extensively studied in the centralized environment, including star-networks wherein there is a master node connected to all the other workers. Our interest is in the following (non-accelerated) algorithms:
1) Unrelated setting: (P) can be solved on star-networks employing the standard proximal gradient method: to reach precision on the objective value, one needs \mathcal{O}\big{(}\kappa_{g}\log(1/\epsilon)\big{)} iterations (which is also the number of communication rounds between the master and the workers), where is the condition number of .
2) -related setting: When the agents’ functions are sufficiently similar, a linear rate proportional to may be highly suboptimal. For instance, in the extreme case where all ’s are identical (), the number of iterations/communications to an solution would remain the same as for . In fact, when , faster rates can be obtained exploiting the similarity of the ’s. Specifically, proposed DANE: a mirror-descent type algorithm over star-networks, where each worker replaces the quadratic term in its local proximal-gradient update with the Bregman divergence of the reference function ; and the master averages the solutions of the workers. DANE is applicable to (P) with : For quadratic losses, it achieves an -solution in \mathcal{O}\big{(}(\beta/\mu)^{2}\cdot\log(1/\epsilon)\big{)} iterations/communications (it is assumed while no improvement is proved over the proximal gradient if the ’s are not quadratic. More recently, proposed CEASE, which achieves DANE’s rate for (P) with and nonquadratic losses. Using recent results in , it is not difficult to check that the mirror-descent algorithm implemented at the master (thus without averaging workers’ iterates) with the Bregman divergence of ( is the local function at the master) achieves an solution in \widetilde{\mathcal{O}}\big{(}\beta/\mu\cdot\log(1/\epsilon)\big{)} iterations/communications, improving thus on DANE/CEASE’s rates.
1 Major contributions
We provide the first linear convergence rate analysis of a distributed algorithm, SONATA (Successive cONvex Approximation algorithm over Time-varying digrAphs), applicable to the composite, constrained formulation (P) over (time-varying, directed) graphs. SONATA was earlier proposed in the companion paper for nonconvex problems. It combines the use of surrogate functions in the agents’ subproblems with a perturbed (push-sum) consensus mechanism that aims at locally tracking the gradient of . Surrogate functions replace the more classical first order approximation of the local ’s, which is the omnipresent choice in current distributed algorithms, offering the potential to better suit the geometry of the problem. For instance, (approximate) Newton-type subproblems or mirror descent-type updates naturally fit our surrogate models; they are the key enabler of provably faster rates in the -related setting. We comment SONATA’s rates below (cf. Table 3).
Unrelated setting (Table 3): When the network is sufficiently connected or it has a star-topology, SONATA reaches an -solution on the objective value in \mathcal{O}\big{(}\kappa_{g}\log(1/\epsilon)\big{)} iterations/communications, which matches the rate of the centralized proximal-gradient algorithm. For arbitrary network connectivity, the same iteration complexity is achieved at the cost of rounds of communications per iteration (employing Chebishev acceleration), where is the second largest eigenvalue modulus of the mixing matrix. Our rates improve on those of existing distributed algorithms which show a much more pessimistic dependence on the optimization parameters and are proved under more restrictive assumptions–contrast Table 2 with Table 3. Linear rates over time-varying digraphs are reported in Table 4 (cf. Sec. 4.2).
-related setting (Table 3): When the agents’ functions are sufficiently similar (specifically, ), the use of a mirror descent-type surrogate over linearization of the ’s provably yields faster rates, at higher computation costs. This improves on the rate of existing distributed algorithms, which are oblivious of function similarity (cf. Table 2). Notice that this is achieved without exchanging any Hessian matrix over the network but leveraging function homogeneity via surrogation. When customized over star-topologies, SONATA’s rates improve on DANE/CEASE’s ones too.
2 Related works
Early works on distributed optimization aimed at decentralizing the (sub)gradient algorithm. The Distributed Gradient Descent (DGD) was introduced in for unconstrained instances of (P) and in for least squares, bot over undirected graphs. A refined convergence rate analysis of DGD can be found in . Subsequent variants of DGD include the projected (sub)gradient algorithm and the push-sum gradient consensus algorithm , the latter implementable over digraphs. While different, the updates of the agents’ variables in the above algorithms can be abstracted as a combination of one (or multiple) consensus step(s) (weighted average with neighbors variables) and a local (sub)gradient descent step, controlled by a step-size (in some schemes, followed by a proximal operation). A diminishing step-size is used to reach exact consensus on the solution, converging thus at a sublinear rate. With a fixed step-size , linear rate of the iterates is achievable, but it can only converge to a -neighborhood of the solution .
Several subsequent attempts have been proposed to cope with this speed-accuracy dilemma, leading to algorithms converging to the exact solution while employing a constant step-size. Based upon the mechanism put forth to cancel the steady state error in the individual gradient direction, existing proposals can be roughly organized in three groups, namely: i) primal-based distributed methods leveraging the idea of gradient tracking ; ii) distributed schemes using ad-hoc corrections of the local optimization direction ; and iii) primal-dual-based methods . We elaborate next on these works, focusing on schemes achieving linear rate– Table 1 organizes these schemes based upon the setting their convergence is established while Table 2 reports the explicit expression of the rates.
i) Gradient-tracking-based methods: In these schemes, each agent updates its own variables along a direction that tracks the global gradient . This idea was proposed independently in the NEXT algorithm for Problem (P) and in AUG-DGM for strongly convex, smooth, unconstrained optimization. The work introduced SONATA, extending NEXT over (time-varying) digraphs. A convergence rate analysis of was later developed in , with considering also (time-varying) digraphs. Other algorithms based on the idea of gradient tracking and implementable over digraphs are ADD-OPT and . Subsequent schemes, , the Push-Pull , and the algorithms, relaxed previous conditions on the mixing matrices used in the consensus and gradient tracking steps over digraphs, which neither need to be row- nor column-stochastic. All the schemes above but NEXT and SONATA are applicable only to smooth, unconstrained instances of (P), with each strongly convex. This latter assumption is restrictive in some applications, such as distributed machine learning, where not all are strongly convex but is so.
ii) Ad-hoc gradient correction-based methods: These methods developed specific corrections of the plain DGD direction. Specifically, EXTRA and its variant over digraphs, EXTRA-PUSH , introduce two different weight matrices for any two consecutive iterations as well as leverage history of gradient information. They are applicable only to it smooth, unconstrained problems; when each is strongly convex, they generate iterates that converge linearly to the minimizer of . To deal with an additive convex nonsmooth term in the objective, proposed PG-EXTRA, which is thus applicable to (P) over undirected graphs, possibly with different local nonsmooth functions. However, linear convergence is not certified. A different approach is to use a linearly increasing number of consensus steps rather than correcting directly the gradient direction; this has been studied in for unconstrained minimization of smooth, strongly convex ’s over undirected graphs.
iii) Primal-dual methods: A common theme of these schemes is employing a prima-dual reformulation of the original multiagent problem whereby dual variables associated to a properly defined (augmented) Lagrangian function serve the purpose of correcting the plain DGD local direction. Examples of such algorithms include: i) distributed ADMM methods and their inexact implementations ; ii) distributed Augmented Lagrangian-based methods with randomized primal variable updates ; and iii) a distributed dual ascent method employing tracking of the average of the primal variable . All these schemes are applicable only to smooth, unconstrained optimization over undirected graphs, with handling time-varying graphs. The extension of these methods to digraphs seems not straightforward, because it is not clear how to enforce consensus via constraints over directed networks.
To summarize, the above literature review shows that currently there exists no distributed algorithm for the general formulation (P) that provably converges at linear rate to the exact solution, in the presence of a nonsmooth function or constraints (cf. Table 1); let alone mentioning digraphs. Furthermore, when it comes to the dependence of the rate on the optimization parameters, Table 3 shows that, even restricting to unconstrained, smooth minimization, SONATA’s rates improve on existing ones–in particular, SONATA provably obtains fast convergence if the agents’ objective functions (e.g., data) are sufficiently similar.
3 Paper organization
Sec. 2 introduces the main assumptions on the optimization problem and network, along with some motivating examples from machine learning. The SONATA algorithm over undirected graphs is studied in Sec. 3; in particular, linear convergence is proved in Sec. 3.3, while a detailed discussion on the rate expression and its scalability properties is provided in Sec. 3.4. The case of time-varying, possibly directed, graphs is considered in Sec. 4. Finally, some numerical results supporting our theoretical findings are reported in Sec. 5. The study of SONATA when is nonconvex can be found in the technical report .
Problem & Network Setting
This section summarizes the assumptions on the optimization problem and network setting. We also introduce a general learning problem over networks, which will be used as case study throughout the paper.
Our algorithmic design and convergence results pertain to two problem settings, namely: i) the one where the local functions are generic and unrelated (cf. Sec. 2.1.1), and ii) the case where they are related (cf. Sec. 2.1.2). These two settings are formally introduced below.
Consider the following standard assumption.
The set is closed and convex;
Each is twice differentiable on the open set and convex;
is convex possibly nonsmooth.
for some and . Unlike existing works (cf. Table 1), we do not require each to be strongly convex but just (cf. A3). Also, twice differentiability of is not really necessary, but assumed here to simplify our derivations.
Under Assumption A, we define the global conditional number associated to (P):
Related quantities determining the (linear) convergence rate of existing distributed algorithms are (cf. Table 2):
Example 1: Consider the following instance of Problem (P):
which all grow indefinitely as or increase.
In the setting above, our goal is to design linearly convergent distributed algorithms whose iterations complexity is proportional to , instead of the larger quantities in (3).
1.2 The β𝛽\beta-related setting
This setting considers explicitly the case where the functions are similar, in the sense defined below .
The local functions ’s (satisfying Assumption A) are called -related if , for all and some .
The more similar the ’s, the smaller . For arbitrary ’s, is of the order of
The interesting case is when ; a specific example is discussed next.
Consider a stochastic learning setting whereby the ultimate goal is to minimize some population objective
To solve (6), the agents have access only to a finite number, say , of i.i.d. samples from the distribution , evenly and randomly distributed over the network. Using the notation introduced in Sec. 1, the ERM problem reads:
where is regularized empirical loss of agent , -strongly convex. Clearly (7) is an instance of (P), satisfying Assumption A.
For the ERM problems (7) we derive next the associated and contrasts with . is -strongly convex; therefore, we can set . The optimal choice of is the one minimizing the statistical error resulting in using as proxy for . We have [36, Th. 7], with high probability, F(\widehat{\mathbf{x}})-F({\mathbf{x}}^{\star})\leq\frac{\lambda}{2}\|\boldsymbol{\theta}^{\star}\|^{2}+\mathcal{O}\big{(}\frac{G_{f}^{2}}{\lambda\,N}\big{)}\leq\mathcal{O}\big{(}\lambda\,B^{2}+\frac{G_{f}^{2}}{\lambda\,N}\big{)}, where is the Lipschitz constant of on , for all . The optimal choice of and resulting minimum error rate are then
An estimate of can be obtained exploring the statistical similarity of the local empirical losses in (7). Under the additional assumption that is -Lipchitz on , for all , a minor modification of [58, Lemma 6] applied to (6)-(7), yields: with high probability,
where hides the log-factor dependence. Note that when is quadratic (i.e., ), scales favorably with the dimension .
Based on (8)-(9), an estimate of and for (7) reads:
Note that increases with the local sample size while does not (neglecting log-factors). It turns out that algorithms converging at a rate depending on exhibit a speed-accuracy dilemma: small statistical errors in (8) (larger ) are achieved at the cost of more iterations (larger ). In this setting, it is thus desirable to design distributed algorithms whose rate depends on rather than .
2 Network setting
We will consider separately two network settings: i) the case where the underlying communication graph is fixed and undirected; and ii) the more general setting of time-varying directed graphs.
When the network of the agent is modeled as a fixed, undirected graph, we write , where denotes the vertex set–the set of agents–while represents the set of edges–the communication links; iff there exists a communication link between agent and . We make the following standard assumption on the graph connectivity.
In this setting, communication network is modeled as a time-varying digraph: time is slotted, and at time-frame , the digraph reads , where the set of edges represents the agents’ communication links: there is a link going from agent to agent . We make the following standard assumption on the “long-term” connectivity property of the graphs.
Assumption B (On the network). The graph sequence , , is -strongly connected, i.e., there exists a finite integer such that the graph with edge set is strongly connected, for all .
The network setting covers, as special case, star-networks, i.e., architectures with a centralized node (a.k.a. master node) connected to all the others (a.k.a. workers). This is the typical computational architecture of several federated learning systems.
The SONATA algorithm over undirected graphs
In words, each agent , given the current iterates and , first solves a strongly convex optimization problem wherein is an approximation of the sum-cost at ; in (11a) is a strongly convex function, which plays the role of a surrogate of (cf. Assumption C below) while acts as approximation of the gradient of at , that is, (see discussion below). Then, agent updates along the local direction [cf. (11b)], using the step-size ; the resulting point is broadcast to its neighbors. The update is obtained via the consensus step (11c) while the -variables are updated via the perturbed consensus (11d), aiming at tracking .
The main assumptions underlying the convergence of SONATA are discussed next.
The surrogate functions satisfy the following conditions.
Each is and satisfies
, for all ;
is -Lipschitz continuous on , for all ;
is -strongly convex on , for all ;
where is the partial gradient of at with respect to the first argument.
The assumption states that should be regarded as a surrogate of that preserves at each iterate the first order properties of . Conditions (i)-(iii) are certainly satisfied if one uses the classical linearization of , that is,
with . Using (13), one can rewrite (11a) as:
which can be interpreted as a mirror-descent update (with step-size one) for the composite minimization of , based on the Bregman distance associated with the reference function .
We refer the reader to as good sources of examples of nonlinear surrogates satisfying Assumption C; here we only anticipate that, when the ’s are sufficiently similar, higher order models such as (13) yield indeed faster rates of SONATA than those achievable using linear surrogates (12). Further intuition is provided next.
For instance, (15) holds with . Roughly speaking, the smaller the better in (11a) approximates . To see this, compare and up to the second order: there exist such that
Noting that [Assumption C(i)] and , and anticipating as (see discussion below), it follows that approximates asymptotically, up to the first order. A better match, is achieved when is sufficiently small. One can then expect that, if the local functions are sufficiently similar ( is small), surrogates exploiting higher order information of , such as (13), may be more effective than mere linearization. Our theoretical findings confirm the above intuition–see Sec. 3.4.
In the consensus and tracking steps, the weights ’s satisfy the following standard assumption.
The weight matrix has a sparsity pattern compliant with , that is
, if ; and otherwise;
Furthermore, is doubly stochastic, that is, and .
Several rules have been proposed in the literature compliant with Assumption D, such as the Laplacian, the Metropolis-Hasting, and the maximum-degree weights rules .
Finally, we comment the anticipated gradient tracking property of the -variables, that is, as . Define the average processes
Summing (11d) over and invoking the doubly stochasticity of ; we have
Applying (18) inductively and using the initial condition , yield
That is, the average of all the ’s in the network is equal to that of the ’s, at every iteration . Assuming that consensus on ’s and ’s is asymptotically achieved, that is, and , , (19) would imply the desired gradient tracking property as , for all .
1 A special instance: SONATA on star-networks
Although the main focus of the paper is the study of SONATA over meshed-networks, it is worth discussing here its special instance over star networks. Specifically, consider a star (unidirected) graph with nodes, where one of them (the master node) connects with all the others (workers). The workers still own only one function of the sum-cost . Two common approaches developed in the literature to solve (P) in this setting are: (i) based upon receiving the gradients from the workers, the master solves (P) and broadcasts the updated vector variables to the workers; (ii) based upon receiving the full gradient and the current iterate from the master, all the workers solve locally an instance of (P) and send their outcomes to the master that averages them out, producing then the new iterate. Here we follow the latter approach; the algorithm is described in Algorithm 2, which corresponds to SONATA (up to a proper initialization), with weight matrix .
SONATA-star, employing linear surrogates [cf. (12)] and , reduces to the proximal gradient algorithm. When the surrogates (13) are used (and still ), SONATA-star coincides with the DANE algorithm if and to the CEASE (with averaging) algorithm if . Nevertheless, our convergence rates improve on those of DANE and CEASE–see Sec. 3.4.1.
2 Intermediate definitions
We conclude this section introducing some quantities that will be used in the rest of the paper. We define the optimality gap as
where is the unique solution of Problem (P).
We stack the local variables and gradients in the column vectors
The average of each of the vectors above is defined as . The consensus disagreements on ’s and ’s are
respectively, while the gradient tracking error is defined as
Finally, given the weight matrix , we define
Under Assumptions B and D, it is well known that (see, e.g., )
where denotes the largest singular value of its argument.
3 Linear convergence rate
Our proof of linear rate of SONATA passes through the following steps. Step 1: We begin showing that the optimality gap converges linearly up to an error of the order of , see Proposition 3.4. Step 2 proves that and are also linearly convergent up to an error , see Proposition 3.5. In Step 3 we close the loop establishing , see Proposition 3.6. Finally, in Step 4, we properly chain together the above inequalities (cf. Proposition 3.8), so that linear rate is proved for the sequences , , , and –see Theorems 3.9 and 3.10. We will tacitly assume that Assumptions A, B, C, and D are satisfied.
Invoking the convexity of and the doubly stochasticity of , we can bound as
We can now bound , regarding the local optimization (11a)-(11b) as a perturbed descent on the objective, whose perturbation is due to the tracking error . In fact, Lemma 3.1 below shows that, for sufficiently small , the local update (11b) will decrease the objective value up to some error, related to .
Let be the sequence generated by SONATA; there holds:
where .
Invoking the optimality of and defining , we have
where the equality follows from and the integral form of the mean value theorem. Substituting (30) in (29) and using the convexity of yield
It remains to bound . We proceed as follows:
We can now substitute (28) into (27) and get
where in (a) we used Young’s inequality, with satisfying
Next we lower bound in terms of the optimality gap.
The following lower bound holds for :
where is defined in (24).
Invoking the optimality condition of , yields
Using the -strong convexity of , we can write
Rearranging the terms and summing over , yields
Using (27) in conjunction with leads to
Combining (37) with (38) provides the desired result (35). ∎
As last step, we upper bound in (33) in terms of the consensus errors and .
The following upper bound holds for the tracking error :
where is defined in (4).
We are ready to prove the linear convergence of the optimality gap up to consensus errors. The result is summarized in Proposition 3.4 below. The proof follows readily multiplying (33) and (35) by and , respectively, adding them together to cancel out , and using (39) to bound .
The optimality gap [cf. (20)] satisfies
where and are defined as
We upper bound and in terms of . We begin rewriting the SONATA algorithm (11a)-(11d) in vector-matrix form; using (21) and (25),we have
Noting that [similarly, ] and (due to the doubly stochasticity of ), it follows from (43) that
Using (44)-(45), Proposition 3.5 below establishes linear convergence of the consensus errors and , up to a perturbation.
with and defined in (26) and (4), respectively.
We prove next (46b); (46a) follows readily from (44). Using (43a), (45), and the Lipschitz continuity of [cf. (1)], we can bound as
where in the last inequality we used . ∎
The following upper bound holds for :
where and , , are defined in (4) and (24), respectively.
By optimality of and we have
Summing the two inequalities above yields
Rearranging terms and using the reverse triangle inequality we obtain the following bound for :
3.4 Step 4: Proof of the linear rate (chaining the inequalities)
We are now ready to prove linear rate of the SONATA algorithm. We build on the following intermediate result, introduced in .
Given the sequence , define the transformations
for . If is bounded, then .
We show next how to chain the inequalities (40), (46) and (47) so that Lemma 3.7 can be applied to the sequences , , and , establishing thus their linear convergence.
Let , , and denote the transformation (49) applied to the sequences , , and , respectively. Given the constants and (defined in Proposition 3.4) and the free parameters (to be determined), the following hold
Squaring (46) and using Young’s inequality yield
for arbitrary . The proof is completed by taking the maximum of both sides of (40), (47), and (53) over and using , for any sequence and . ∎
Chaining the inequalities in Proposition 3.8 in the way shown in Fig. 1, we can bound as (see Appendix A for the proof)
where is defined as
and is a remainder, which is bounded under (51).
Therefore, as long as , (54) implies
where is a constant independent of . Therefore, and thus converges R-linearly to zero at rate at least (cf. Lemma 3.7). Applying the same argument to the other inequalities in Proposition 3.8, one can conclude that also the sequences , and converge R-linearly to zero.
The last step consists to showing that there exist a sufficiently small step-size and satisfying (51), such that . This is proved in the Theorem 3.9 below.
The proof is organized in following two steps: Step 1) We first consider the “marginal” stable case by letting , and show that there exists so that , for all ; Step 2) Then, invoking the continuity of , we argue that, for any , one can find such that \mathcal{P}\big{(}\alpha,\bar{z}(\alpha)\big{)}<1. This implies the boundedness of D^{K}\big{(}\bar{z}(\alpha)\big{)}, and thus \|\mathbf{d}^{\nu}\|^{2}=\mathcal{O}\big{(}\bar{z}(\alpha)^{\nu}\big{)} (cf. Lemma 3.7).
Step 1: We begin optimizing the free parameters , , and . Since the goal is to find the largest so that , for all , the optimal choice of , , and is the one that minimizes , that is,
We then set , and proceed to optimize , which appears in and . Recalling the definition of and (cf. Proposition 3.4) and the constraint (34), the problem boils down to minimize
Let denote the value of corresponding to the optimal choice of the above parameters. The expression of reads
The next theorem provides an explicit expression of the convergence rate in Theorem 3.9 in terms of the step-size ; the constants , , and therein are defined in (103), (101) with , and (105), respectively.
In the setting of Theorem 3.9, suppose that the step-size satisfies , with Then, , for all , where
4 Discussion
Theorem 3.10 provides a unified set of convergence conditions for different choices of surrogates and network topologies. To shed light on the expression of the rate and its dependence on the key optimization and network parameters, we customize here Theorem 3.10 to specific network topologies and surrogate functions. We begin considering star-networks (cf. Sec. 3.4.1) and then move to general graph topologies with no master node (cf. Sec. 3.4.2). We will customize the rate achieved by SONATA employing the following two surrogate functions , representing the two extreme choices in the spectrum of admissible surrogates:
Convergence of SONATA-Star (Algorithm 2) is established in Corollary 3.11 below.
In particular, when the surrogates (62) and (63) are employed along with , the rate above reduces to the following expressions:
Linearization (62): . Therefore, in at most \mathcal{O}\Big{(}\kappa_{g}\log({1}/{\epsilon})\Big{)} iterations (communications);
Therefore, in at most
The following comments are in order. When linearization is employed, SONATA-Star matches the iteration complexity of the centralized proximal-gradient algorithm. When the ’s are sufficiently similar, (65)-(66) proves that faster rates can be achieved if surrogates (63) are chosen over first-order approximations: when , (66) is significantly faster than \mathcal{O}\big{(}\kappa_{g}\log({1}/{\epsilon})\big{)}. As case study, consider Example 2 (cf. Sec. 2.1.2): plugging (10) into Corollary 3.11 shows that using the surrogates (63) yields \widetilde{\mathcal{O}}\big{(}L\,\sqrt{{d\,m}}\cdot\log(1/\epsilon)\big{)} iterations (communications); this contrasts with \widetilde{\mathcal{O}}\big{(}L\,\sqrt{d\,m\,n}\cdot\log(1/\epsilon)\big{)}, achieved by first-order methods (and SONATA-Star using linearization), which instead increases with the sample size .
Since SONATA-Star contains as special cases the DANE and CEASE algorithms, we contrast here Corollary 3.11 with their convergence rates. We recall that DANE is applicable to (P) when : For quadratic losses, it achieves an -optimal objective value in \mathcal{O}\big{(}(\beta/\mu)^{2}\cdot\log(1/\epsilon)\big{)} iterations/communications (here ). This rate is worse than (66). For nonquadratic losses, did not show any rate improvement of DANE over plain gradient algorithms, i.e., \mathcal{O}\big{(}\kappa_{g}\cdot\log(1/\epsilon)\big{)} while SONATA-star still retains \mathcal{O}\big{(}\beta/\mu\cdot\log(1/\epsilon)\big{)}. The CEASE algorithm is proved to achieve an -solution on the iterates in \mathcal{O}\big{(}(\beta/\mu)^{2}\cdot\log(1/\epsilon)\big{)} iterations/communications (with ); SONATA reaches the same error on the iterates in \mathcal{O}\big{(}\beta/\mu\cdot\log(\kappa_{g}/\epsilon)\big{)} iterations/communications, which matches the order of the mirror-decent algorithm.
In the next section we extend the study to networks with no centralized nodes, sheding lights on the role of the network in achieving the same kind of results.
4.2 The general case
The convergence rate of SONATA over general graphs is summarized in Corollary 3.12 for the linearization surrogates (62) while Corollaries 3.13 and 3.14 consider the surrogates (63) based on local , with Corollary 3.13 addressing the case and Corollary 3.14 the case . The step-size is tuned to obtain favorable rate expressions.
In the setting of Theorem 3.10, let be the sequence generated by SONATA, using the surrogates (62) and step-size , , with . The number of iterations (communications) needed for , , is
Instate assumptions of Theorem 3.10 and suppose . Consider SONATA using the surrogates (63) and step-size , , with and . The number of iterations (communications) needed for , , is
Instate assumptions of Theorem 3.10 and suppose . Consider SONATA using the surrogates (63) and step-size , , with and . The number of iterations (communications) needed for , , is
The proof of Corollaries 3.13 and 3.14 can be found in Appendix E.
Order of the rate of centralized (nonaccelerated) methods (Case I): For a fixed optimization problem, if the network is sufficiently connected ( “small”), its impact on the rate becomes negligible (the bottleneck is the optimization), and SONATA matches the network-independent rate order achieved on star-topologies (cf. Corollary 3.11) by the proximal gradient algorithm when linearization is employed [cf. (67)] and by the mirror-descent scheme when the local ’s are used in the surrogates [cf. (69) and (71)].
Network-dependent rates (Case II): As expected, the convergence rate deteriorates as increases, i.e., the network connectivity gets worse. This translates in a less favorable dependence of the complexity on and (by a square factor) and network scalability of the order of . When (e.g., the network is decently connected or ), the complexity becomes , which compares favorably with that of existing distributed schemes, determined instead by the more pessimistic local quantities (3). The scalability of the rate with the network connectivity, , can be improved leveraging multiple rounds of communications or accelerated consensus protocols, as discussed below.
Linearization (62) vs. local (63) surrogates: As already observed in the setting of star-networks, the use of the local losses as surrogates employs a form of preconditioning in the local agents subproblems. When the ’s are sufficiently similar to each other, so that , exploiting local Hessian information via (63) provably reduces the iteration/communication complexity over linear models (62)–contrast (67) with (69) and (71). Note that these faster rates are achieved without exchanging any matrices over the network, which is a key feature of SONATA. On the other hand, when the functions are heterogeneous, the local surrogates (63) are no longer informative of the average-loss and using linearization might yield better rates. Although these design recommendations are based on sufficient conditions, numerical results seem to confirm the above conclusions–see Sec. 5.
Multiple communications rounds and acceleration: The discussion above shows that rates of the order of those of centralized methods can be achieved if the network is sufficiently connected (Case I). When this is not the case, one can still achieve the same iteration complexity at the cost of multiple, finite, rounds of communications per iteration. Specifically, let be the connectivity of the given network and suppose we run steps of communications per iteration (computation) in (43a)-(43b); this yields an effective network with improved connectivity . One can then choose so that the ratio satisfies the condition triggering Case I in the Corollaries 3.12–3.14, as briefly summarized next.
1) Linearization: Invoking Corollary 3.12, one can check that the order of such a is ; therefore, SONATA using the surrogates (62) reaches an -solution in iterations and communications. The dependence on the network connectivity can be further improved leveraging Chebyshev polynomials (see, e.g., ): the final communication complexity of SONATA reads
2) Local surrogates: Considering the case (Corollary 3.14), we can show that SONATA using the surrogates (63) and employing multiple rounds of communications per iteration, reaches an -solution in iterations and \mathcal{O}\left(\beta/\mu\cdot\log\big{(}(\kappa_{g}+\beta/\mu)(1+L/\beta)\big{)}(1-\rho_{0})^{-1}\log(1/\epsilon)\right) communications. If Chebyshev polynomials are used to accelerate the communications, the communication complexity further improves to
The SONATA algorithm over directed time-varying graphs
In this section we extend SONATA and its convergence analysis to solve Problem (P) over directed, time-varying graphs (Assumption B). Note that (11a)-(11d) is not readily applicable to this setting, as constructing a doubly stochastic weight matrix compliant with a directed graph is generally infeasible or computationally costly–see e.g. . Conditions on the weight matrices can be relaxed if the consensus/tracking schemes (11c)-(11d) are properly changed to deal with the lack of doubly stochasticity.
Here, we consider the perturbed push-sum protocols as proposed in the companion paper (but in the Adapt-Then-Combine (ATC) form). The resulting distributed algorithm, still termed SONATA, is formally described in Algorithm 3.
In the perturbed push-sum protocols (73c)-(73d), satisfies the assumption below.
Moreover, is column stochastic, i.e., , for all
We conclude this section stating the counterparts of the definitions introduced in Sec. 2, adjusted here to the case of directed time-varying graphs. Using the column stochasticity of and (73d), one can see that opposed to (18), the average gradient is now preserved on the weighted average of the ’s:
where is defined in (17). This suggests to decompose into its weighted average and the consensus error, defined respectively as
Accordingly, we define the weighted average of and the consensus error as
In addition, we also generalize the definition of the optimality gap as
Furthermore, we will use the following lower and upper bounds of [35, Prop. 1]
The proof of linear convergence of SONATA (Algorithm 3) follows the same path of the one developed in Sec. 3.3 for the case of undirected graphs. Hence, we omit similar derivations and highlight only the key differences. We will tacitly assume that Assumptions A, B, C, and E are satisfied.
The optimality gap sequence satisfies:
where the constants and are defined in (4) and (24), respectively; and and are defined in (41).
The proof follows closely that of Proposition 3.4 and thus is omitted. For completeness, we report it in the supporting materials. Here, we only notice that, instead of (27), we built on: \sum_{i=1}^{m}\phi_{i}^{\nu+1}U(\mathbf{x}_{i}^{\nu+1})\leq\sum_{i=1}^{m}\phi_{i}^{\nu}U\big{(}\mathbf{x}_{i}^{\nu+\frac{1}{2}}\big{)}, where we used , for all .∎
The following bounds hold for and :
where and are defined in (78), and and are arbitrary positive constants (to be determined).
Using the result in [26, Lemma 5] and [35, Lemma 3, 11], we obtain
The rest of the proof follows similar steps as [52, Lemma 2], hence it is omitted. ∎
where , , , and are defined in (4) and (24), respectively.
The proof follows similar path of that of Proposition 3.6 and thus is omitted. ∎
2 Establishing linear rate
Let , , , and denote the transformation (49) of the sequences , , and . Given the constants and , defined in Proposition 4.1, and the free parameters , the following holds:
The proof of the first two inequalities (85) and (88) follows the same steps of those used to prove Proposition 3.8. Applying [43, Lemma 21] to (81a) and (81b) respectively gives (86) and (87).
Chaining the inequalities in Proposition 4.4 as done in for (50) (cf. Fig. 1), we can bound as
where is defined as
and is a bounded remainder term.
Comparing (95) to (55) we can see that they share the same form and only differ in coefficients. Therefore, with the same argument as in the proof of Theorem 3.9 we can easily arrive at the following conclusion.
We provide the proof in the supporting material. ∎
For sake of completeness, we provide an explicit expression of the linear rates in terms of the step-size in the supporting material–see Theorem III.1. Table 4 summarizes the expression of the rates achieved by SONATA using the surrogate functions (62) and (63)–a formal statement of these results along with the proofs can be found in the supporting material-see Corollaries IV.1, V.1 and V.2.
The rate estimates in Table 4 are almost identical to those obtained in Sec. 3.4.2, with the difference that the network dependence now is expressed throughout rather than . Therefore, similar comments–as those stated in Sec. 3.4.2–apply to the rates in Table 4. For example, if the network is sufficiently connected ( “small”), its impact on the rate becomes negligible and SONATA matches the network-independent rate achieved on star-topology (cf. Corollary 3.11) or centralized settings. Specifically, when linearization surrogate (62) is used, this rate coincides with the rates of centralized proximal gradient algorithm.
Numerical Results
In this section, we corroborate numerically the complexity results proved in Corollaries 3.12–3.14. As a test problem, we consider the distributed ridge regression:
where the loss function of agent is [agent owns data ]. Problem parameters are generated as follows. Each row of the measurement matrix is independently and identically drawn from distribution ; and is generated according to the linear model , where is the ground truth, generated according to , and is the measurement noise. The covariance matrix is constructed according to the eigenvalue decomposition , where the eigenvalues are uniformly distributed in . The eigenvectors, forming , are obtained via the QR decomposition of a random matrix with standard Gaussian i.i.d. elements. The network is generated using an Erdős-Rényi model , with nodes and each edge independently included in the graph with probability .
To investigate the impact of and on the convergence rate, we specifically consider the following two scenarios:
We run SONATA using surrogates (62) (linearization) and (63) (local )–we term it as SONATA-L and SONATA-F, respectively. The simulations parameters of the different experiments are summarized in Table 5; and the algorithmic parameters are set according to Corollaries 3.12–3.14. The expressions are not tight in terms of the absolute constants. To show convergence rate in both Cases I and II in Corollary 3.12-3.14, we enlarged the second term in the expression of by a constant factor. We measure the algorithm’s complexity using .
In Table 6, we report the corresponding iteration complexity of SONATA for each simulation setup (s.1)-(s.6) in Table 5. Each figure is generated under one particular realization of the problem setting. Further, in order to compare the complexity of SONATA across different settings, all the simulations share the same network parameters, as well as the same data set whenever the problem parameters are the same. The results of our experiments are reported in Table 6; the curve are generated using only one random realization for visualization clarity. However, the behavior of the curves (e.g., scalability with respect to the parameters) is representative and consistent across all the random experiments we conducted.
Scalability with respect to . Consider setting (S.I) wherein is fixed and is changing. Figures for (s.1)-(s.3) show that when (blue curve), the iteration complexity of SONATA-L scales linearly with respect to [as predicted by Corollary 3.12], while that of SONATA-F is invariant whenever [as stated in Corollary 3.13]. When , the iteration complexity of SONATA-F grows as increases since decreases [cf. Corollary 3.14]. However, the increasing rate is much slower than SONATA-L, due to the fact that for large . When , the iteration complexity scales quadratically with respect to , in all settings, as predicted by our theory.
Scalability with respect to . Consider now setting (S.II), where we decrease the local sample size to increase . In contrast to setting (S.I), Figures for (s.4) and (s.5) show that, with , the iteration complexity of SONATA-F scales linearly with when , while that of SONATA-L is invariant–this is consistent with Corollaries 3.12 and 3.14. When , the iteration complexity scales quadratically with respect to . Finally, the plot associated with (s.6) simply reveals that when , iteration complexity of SONATA-F remains bounded, as stated in Corollary 3.13.
Appendix A Proof of (54)
Chaining the inequalities in (50) as shown in Fig. 1, we have
Notice that, under (51), , , , and , , are all bounded, which implies that the reminder in (50) is bounded as well.
Appendix B Proof of Theorem 3.10
We find the smallest satisfying (51) such that , for , with to be determined.
Let us begin considering the condition in (51). To simplify the analysis, we impose instead the following stronger version
Observe that in the expression of , the only coefficient multiplying that depends on is the optimization gain Using (97), can be upper bounded as
where , and are constants defined as
Condition (100) shows the rate must satisfy
Notice that, under , (101) implies , which are the other two conditions on in (51). Therefore, overall, must satisfy (97) and (101). Letting in (97), the condition simplifies to
Therefore, the overall convergence rate can be upper bounded by , where
Note that as increases from , the first term in the max operator above is monotonically increasing from while the second term is monotonically decreasing from . Therefore, there must exist some so that the two terms are equal, which is
To conclude, given the step-size satisfying , the sequence converges at rate , with given in (61).
Appendix C Proof of Corollary 3.11
Since , we have ; then (33a) and (35) reduce to
respectively. Combining (106) and (107) and using , yield
We customize next (64) to the specific choices of the surrogate functions.
Finally, setting in the expression above, yields (65).
Appendix D Proof of Corollary 3.12
According to Theorem 3.10, the rate can be bounded as
where and are defined in (103) and (101), respectively.
Accordingly, the expressions of and read:
where in the last inequality we have used the fact that .
Using the above expressions, in the sequel we upperbound and .
Since must be chosen so that , we impose , implying . Since [cf. (111)], the condition on reduces to . Choose , for some given . Depending on the value of , either or .
Case I: . This corresponds to the case , which happens when the network is sufficiently connected ( is small). Note that, we also have , otherwise . In this setting, , and
where in (a) we used and (b) follows from .
Case II: . This corresponds to the case . We have ,
We claim that . Suppose this is not the case, that is, . Since [cf. (111)] and , would imply . This however is in contradiction with the assumption , as it would lead to .
Using , we can bound
Appendix E Proof of Corollaries 3.13 and 3.14
We follow similar steps as in Appendix D but customized to the surrogate (63). We begin particularizing the expressions of and .
Accordingly, the expressions of and read:
Similarly to the proof of Corollary 3.12, we bound as
where and are now given by (115) and (116), respectively. For , we require , and choose , with arbitrary . We study separately the cases and .
Since , we study next the case and separately.
Case I: . We have , , and thus
Since and , it must be . Therefore, the rate can be bounded as
Case II: . This corresponds to , , and
Using the same argument as in the proof of Corollary 3.12–Case II, one can show that . Therefore,
Case I: . Following the same reasoning as , we can prove
Case II: . We claim that , otherwise , which would lead to the following contradiction . Therefore,
where is a suitable constant, independent on , and .
References
Supporting Material
Appendix I Proof of Proposition 4.1
We begin introducing some intermediate results.
Consider Problem (P) under Assumption A; and SONATA (Algorithm 3) under Assumptions C and E. Then, there holds
with defined in (23).
where .
Invoking the optimality of , we have
where the equality follows from and the integral form of the mean value theorem; and .
Substituting (123) in (122) and using the convexity of yield
It remains to bound . We proceed as follows:
We connect now the individual decreases in (121) with that of the optimality gap , defined in (77). Notice that
due to the convexity of , column-stochasticity of and , for all . Summing (121) over , and using (126), we obtain
where in (a) we used Young’s inequality, with satisfying
Next we lower bound in terms of the optimality gap.
In the setting of Lemma 3.1, there holds:
with defined in (24).
Invoking the optimality condition of , yields
Using the -strong convexity of , we can write
Rearranging the terms and summing over , yields
Using (27) in conjunction with leads to
Combining (131) with (132) yields the desired result (129). ∎
As last step, we upper bound in (33) in terms of the consensus errors and .
The tracking error can be bounded as
where is defined in (4).
The linear convergence of the optimality gap up to consensus errors as stated in Proposition follows readily multiplying (129) by and adding with (127) to cancel out , and using (39) to bound .
Appendix II Proof of Theorem 4.5
Following the same steps as in the proof of Theorem 3.9, we derive the optimal appearing in and :
Setting and denoting the corresponding as , the expression of reads
Appendix III Explicit expression of the linear rate in the time-varying directed network setting
The following theorem provides an explicit expression of the convergence rate in Theorem 4.5, in terms of the step-size ; the constants and therein are defined in (148) and (145) with , respectively.
The proof follows similar steps as the proof of Theorem 3.10. For sake of simplicity, we used the same notation as therein. We find the smallest satisfying (89) such that , for , and to be determined[recall that is defined in (95)].
Using exactly the same argument as Theorem 3.10 we have the following two conditions on :
where , and are constants defined as
Lower bounding by we obtain
Letting in (142), the condition reduces to
Therefore, the overall convergence rate can be upper bounded by , where
Appendix IV Rate estimate using linearization surrogate (62) (time-varying directed network case)
In the setting of Theorem III.1, let be the sequence generated by SONATA (Algorithm 3), using the surrogates (62) and step-size , , where and is a constant defined in (154). The number of iterations (communications) needed for , , is
According to Theorem III.1, the rate can be bounded as
where and are defined in (148) and (145), respectively.
Accordingly, the expressions of and read:
and in the first inequality we have used the fact that and , and the last inequality holds since and . Using the above expressions, in the sequel we upperbound and .
Since must be chosen so that , we impose , implying . Since [cf. (152)], the condition on reduces to . Choose , for some given . Depending on the value of , either or .
Case I: . This corresponds to the case . Note that, we also have , otherwise . In this setting, , and
where in (a) we used and (b) follows from .
Case II: . This corresponds to . We have ,
Now we can bound . Since (by the same reasoning as in proof of Proposition 3.12),
Accordingly, the expressions of and read:
and the last inequality holds since and .
Similarly, we bound as
when ; the last inequality holds due to .