Local Exact-Diffusion for Decentralized Optimization and Learning

Sulaiman A. Alghunaim

Introduction

This work examines the distributed consensus optimization problem, as formally presented in (2). In this setup, a network of nodes (also referred to as agents, workers, or clients) collaboratively seeks to minimize the average of the nodes’ objectives. This formulation is appealing for large scale data problems because it is more efficient to use distributed solution methods to reduce the computational burden for large data sets. In terms of communication protocol, distributed methods can be classified as either centralized or decentralizedIn this work, the term “distributed methods” refers to the class of methods that includes both centralized (server-workers) and decentralized approaches.. Centralized distributed methods require all nodes to communicate with a central server (e.g., server-workers connection) without sharing private data, as seen in parallel optimization and federated learning . In this setup, there is a central node that is responsible for aggregating local variables and updating model estimates. In contrast, decentralized distributed methods are “fully distributed” that are designed for arbitrary connected network topologies such as line, ring, grid, and random graphs. These methods require nodes to communicate only with their immediate neighbors . It’s important to note that decentralized methods can adapt to a centralized setting when the network is fully connected.

In this paper, we consider a group of NN nodes, connected via an undirected decentralized network, collaborating to solve the optimization problem:

where fi:m→ℜf_{i}:{}^{m}\rightarrow\real (i=1,…,Ni=1,\dots,N) represents a smooth function known only to node ii. This function is defined as the expected value of some loss function Fi(⋅;ξi)F_{i}(\cdot;\xi_{i}) over the random variable or data ξi\xi_{i}. We focus on the stochastic online setting, in which each node has access only to random samples of its data {ξi}\{\xi_{i}\}. Problems of the form (2) have received a lot of attention in control and engineering communities , as well as in the machine learning community .

Our main contribution is the proposal and study of a local variant of the decentralized Exact-Diffusion algorithm (see also ), where nodes employ multiple local updates between communication rounds. We establish the algorithm’s convergence in both convex and nonconvex settings. . Our bounds improve upon those of decentralized methods and match the best known results for centralized methods. Before we formally state our contributions, we will first discuss some related works.

We begin by discussing relevant centralized methods, which require a central server for implementation. One of the most popular centralized methods is FedAvg, which involves a random subset of nodes performing multiple local stochastic gradient descent (SGD) updates at each round; They then send their estimates (parameters) to the central server, which averages these estimates and sends them back to replace the local estimates . It should be noted that FedAvg is often called Local-SGDIn this paper, we refer to the case where all of the nodes participate in each round as Local-SGD.. Several works have analyzed FedAvg and Local-SGD . It has been observed that the performance of FedAvg and Local-SGD is suboptimal for heterogeneous data and that an increased number of local steps can lead to worse performance . One major reason is that node estimates drift toward their local solutions due to local updates, resulting in a biased solution . To correct this drift in FedAvg and Local-SGD, several algorithms have been proposed, including SCAFFOLD , FedDyn , FedPD , VRL-SGD , and FedGATE . These methods, however, are only applicable to centralized connections.

In this work, we focus on decentralized setups as presented in . The most extensively studied method for this setup is the decentralized stochastic gradient descent method (DSGD) DSGD has two main implementations depending on the combination step: the adapt-then-combine (ATC) implementation (aka diffusion) and the non-ATC implementation (aka consensus) . Both implementations are termed DSGD in this paper.. The study in showcased that DSGD achieves the centralized SGD rate asymptotically, a distributed trait termed as linear speedup . However, DSGD converges to a biased solution, and this bias is further negatively influenced by network sparsity . The bias arises from the heterogeneity among the local functions fi{f_{i}}, as characterized by (1/N)∑i=1N∥∇fi(x)−∇f(x)∥2(1/N)\sum_{i=1}^{N}\|{\nabla}f_{i}(x)-{\nabla}f(x)\|^{2}. This value, which is required to be bounded for the analysis of DSGD , can become quite large when the functions are heterogeneous, thereby slowing down the convergence of DSGD . Multiple studies have introduced bias-correction algorithms resistant to local function heterogeneity, such as the alternating direction method of multipliers (ADMM) based methods , EXTRA , Exact-Diffusion (ED) (also known as NIDS and D2 ), and Gradient-Tracking (GT) methods . It has been established that these bias-correction methods outperform DSGD . Yet, all these methods necessitate communication at every iteration.

Locally updated stochastic decentralized methods have received less attention than centralized methods and are more challenging to study. Federated learning can be viewed as a subset of decentralized optimization and learning under time-varying and asynchronous updates . For instance, DSGD with local steps, termed Local-DSGD, has been studied in and is analogous to Local-SGD when the network is fully connected. However, just as DSGD suffers from bias, Local-DSGD does too. Furthermore, similar to Local-SGD, it also experiences drift. Increasing the number of local steps exacerbates this drift in the solution, requiring the use of very small stepsizes, which significantly slows down convergence.

Only a few works have studied decentralized methods with local updates and bias correction. The work in studied Local Gradient-Tracking (LGT) under nonconvex costs, but it focused solely on deterministic settings. Similarly, the research presented in explored a locally updated stochastic variant of gradient-tracking, namely KK-GT, but it too addressed only nonconvex settings. In this paper, we propose and investigate a different algorithm inspired by . Our results improve upon the rates of local GT-based methods and require just half the communication cost of existing approaches. Furthermore, our rates surpass those of Local-DSGD . Importantly, when employing a constant step size, LED achieves precise convergence in the deterministic (noiseless) case. In contrast, Local-DSGD does not, due to the bias/drift introduced by the heterogeneity in local functions, as discussed earlier. We will next formally outline our contributions.

2 Contribution

We propose Local Exact-Diffusion (LED) method for distributed optimization with local updates. An advantage over previous methods is that LED is a decentralized method that only requires one single vector communication per link and is robust to the functions heterogeneity – See Table 1. Numerical results are provided to demonstrate the effectiveness of LED over other methods.

We provide insights and draw connections between our proposed method and the following state-of-the-art algorithms: Exact-Diffusion , NIDS , D2 , ProxSkip/Scaffnew , VRL-SGD , and FedGATE . For instance, we demonstrate that LED can be interpreted as Scaffnew with fixed local updates instead of random ones. We also show that LED is a decentralized variant of the centralized methods FedGATE and VRL-SGD . Furthermore, we highlight that all these methods can be traced back to the primal-dual method PDFP2O (also known as PAPC ) in the case of a single local update.

We establish the convergence of LED in both (strongly-)convex and nonconvex environments for online stochastic learning settings. Our rates improve upon existing bounds for local decentralized methods—see Table 2. It is worth noting that the analysis of LED is more challenging than that of Local-DSGD, even in the single local update scenario. Additionally, in contrast to Local-DSGD, LED is robust to heterogeneity in local functions and converges exactly in the deterministic case (no noise), as discussed earlier.

A byproduct of our result is that, when adapting our analysis to centralized networks, we achieve new and tighter analyses for the methods VRL-SGD and FedGATE .

Outline. This paper is organized as follows: Section 2 introduces our algorithm and its motivation. In Section 3, we compare our method with other leading approaches. Section 4 presents our core assumptions and convergence findings, and a discussion comparing our results with prior works. Section 5 provides simulation outcomes, and conclusions are drawn in Section 6. Detailed proofs are reserved for the appendix.

Local Exact-Diffusion

In this section, we start by describing the proposed algorithm in its decentralized implementation. We then rewrite it in network notation for reasons of analysis and interpretation.

The method under study is described in Alg. 1 and is named Local Exact-Diffusion (LED). In step 1, each node ii employs \uptau\uptau local updates, starting from the initialization xirx_{i}^{r}, which is its local estimate of the solution after the communication round rr. Step 2 is the communication round during which each node ii sends its local intermediate estimate ϕi,\uptaur\phi_{i,\uptau}^{r} to its neighbors j∈Nij\in{\mathcal{N}}_{i}, where the symbol Ni{\mathcal{N}}_{i} denotes the set of neighbors of node ii (including node ii); in this step, wijw_{ij} is a nonnegative scalar weight that node ii uses to scale the information received from node j∈Nij\in{\mathcal{N}}_{i}. The final step, step 3, is where each node ii updates its (dual) estimate yir∈my_{i}^{r}\in{}^{m}.

2 Networked description

The LED method, as listed in 1, is described at the node level. For the analysis and interpretation of the method, we will present it in a networked form. To do this, we introduce the following network weight matrix notation:

Then, Algorithm 1 can be represented in a compact networked form as follows: Given x0{\mathbf{x}}^{0}, set y0=(I−W)x0{\mathbf{y}}^{0}=({\mathbf{I}}-{\mathbf{W}}){\mathbf{x}}^{0} (or y0=0{\mathbf{y}}^{0}=\mathbf{0}) and update for r=0,1,2,…r=0,1,2,\dots

Local primal updates: set Φ0r=xr\bm{\Phi}_{0}^{r}={\mathbf{x}}^{r}, for t=0,…,\uptau−1t=0,\dots,\uptau-1:

The networked description (6) will be used for analysis purposes.

3 Motivation and relation with Exact-Diffusion

Exact-Diffusion (ED) was derived in and takes the following form:

The unified decentralized algorithm (UDA) from demonstrated that the iterates xr{\mathbf{x}}^{r} of ED in (7) can be equivalently described by:

where B=I−W{\mathbf{B}}={\mathbf{I}}-{\mathbf{W}}. The above form is convenient for analytical purposes; however, it cannot be implemented in a decentralized fashion due to B1/2=(I−W)1/2{\mathbf{B}}^{1/2}=({\mathbf{I}}-{\mathbf{W}})^{1/2}.

To derive our method, we set B=β(I−W){\mathbf{B}}=\beta({\mathbf{I}}-{\mathbf{W}}) in (8c) and introduce the change of variable yr=1βB1/2zr{\mathbf{y}}^{r}=\frac{1}{\beta}{\mathbf{B}}^{1/2}{\mathbf{z}}^{r}. This leads to the following description:

It can be observed that the update LED-1 (9c) is equivalent to LED (6) when \uptau=1\uptau=1. In other words, LED (6) is an extension of LED-1 (9c) that incorporates local updates.

where W~=(1−αc)I+αcW\widetilde{{\mathbf{W}}}=(1-\alpha c){\mathbf{I}}+\alpha c{\mathbf{W}} and cc is a stepsize parameter. When c=1/αc=1/\alpha, NIDS reduces to ED (7). We also note that ED (or NIDS with c=1/αc=1/\alpha) has been studied under the name D2 . Thus, LED can be viewed as a modification of NIDS/D2 that incorporates local updates. ■\blacksquare

It is important to note that the updates (7), (8c), and (9c) are equivalent only when there are no multiple local updates and β=1\beta=1. To understand this, observe that the local updates variants can be modeled as a time-varying graph Wr{\mathbf{W}}_{r}, where Wr=W{\mathbf{W}}_{r}={\mathbf{W}} when r=\uptau,2\uptau,3\uptau,…r=\uptau,2\uptau,3\uptau,\dots, and Wk=I{\mathbf{W}}_{k}={\mathbf{I}} otherwise. In this scenario, the updates xr{\mathbf{x}}^{r} differ for all these methods. Indeed, for this case, the updates (7) and (8c) are not guaranteed to converge and often diverge in simulations. ■\blacksquare

Connection with existing algorithms

In this section, we discuss and highlight the connections of LED to the following algorithms: Scaffnew/ProxSkip , VRL-SGD , and FedGate . We also demonstrate that all these methods can be traced back to the primal-dual method PDFP2O , which is also known as PAPC and was initially proposed in for quadratic objectives.

We now demonstrate that LED (6) with \uptau=1\uptau=1 (i.e., LED-11 (9c)) can be interpreted as the primal-dual algorithm PDFP2O applied to the following reformulation of problem (2):

where B=I−W{\mathbf{B}}={\mathbf{I}}-{\mathbf{W}} and g{\mathbf{g}} is the indicator function of zero, i.e., g(u)=0{\mathbf{g}}({\mathbf{u}})=0 if u=0{\mathbf{u}}=\mathbf{0} and g(u)=+∞g({\mathbf{u}})=+\infty otherwise. Problem (10) is equivalent to (2) because Bx=0{\mathbf{B}}{\mathbf{x}}=\mathbf{0} if and only if x1=x2=⋯=xNx_{1}=x_{2}=\dots=x_{N} – see .

The following updates are obtained when PDFP2O is applied to formulation (10):

where prox⁡αηg∗(⋅)\operatorname{prox}_{\frac{\alpha}{\eta}{\mathbf{g}}^{*}}(\cdot) denotes the proximal operator of the conjugate of gg and α,η>0\alpha,\eta>0 are stepsize parameters. The following result relates LED-1 (9c) (LED (6) with \uptau=1\uptau=1) with PDFP2O (11b).

The updates of PDFP2O (11b) can be rewritten as

It follows that PDFP2O (12) is equivalent to LED-1 (9c) when η=β=1\eta=\beta=1.

If we let Φr=xr−α∇f(xr)−ηB12vr\bm{\Phi}^{r}=\mathbf{x}^{r}-\alpha{\nabla}\mathbf{f}(\mathbf{x}^{r})-\eta\mathbf{B}^{1\over 2}\mathbf{v}^{r}, then we can rewrite equation (11b) as follows:

Since g{\mathbf{g}} is the indicator function of zero, we have prox⁡αηg∗(z)=z\operatorname{prox}_{\frac{\alpha}{\eta}{\mathbf{g}}^{*}}({\mathbf{z}})={\mathbf{z}}; thus vr+1=vr+B12Φr\mathbf{v}^{r+1}=\mathbf{v}^{r}+\mathbf{B}^{1\over 2}\bm{\Phi}^{r}. Moreover, observe that

The above result demonstrates that LED (6) can be interpreted as a locally updated variant of PDFP2O (11b). It also shows that ED/D2 and NIDS are different representations of PDFP2O applied on formulation (10). ■\blacksquare

2 Relation with Scaffnew

The work studied a proximal skipping variant of PDFP2O. The decentralized Scaffnew method is given by [44, Alg. 5]:

where α,ζ\alpha,\zeta are stepsize parameters. Observe that (15) employs local updates (15a) and communicates only with a small probability pp (15b). If we let yr=(1/ζ)zr{\mathbf{y}}^{r}=(1/\zeta){\mathbf{z}}^{r} and p=1p=1 (communicate at each iteration) then (15) reduces to

The update (16) is the same as PDFP2O (12) when η=αζ\eta=\alpha\zeta. Consequently, when p=1p=1 and ζ=1/α\zeta=1/\alpha (16) is exactly LED-1 (9c) when β=1\beta=1.

LED (6) employs a fixed number \uptau\uptau of local updates between two communication rounds, whereas Scaffnew uses random number of local updates between two communication rounds. The use of random communication skipping or fixed local steps differs in analysis, however, in terms of performance they are strikingly similar with 1/p1/p playing the role of \uptau\uptau.

The work analyzes Scaffnew and shows that local steps can save communication when the network is well connected. We point out that the analysis techniques in do not show linear speedup and are only-suited for the probabilistic implementation with strongly-convex costs. The techniques we present in this work are distinct and applicable to the locally updated variant with a deterministic number of local updates (6) for both nonconvex and (strongly-)convex settings. ■\blacksquare

3 Relation with FedGATE/VRL-SGD

The work introduced and analyzed a federated learning algorithm (centralized method) named FedCOMGATE that employs compression; without compression the method reduces to FedGATE [24, Alg. 3], which is a generalization of VRL-SGD . We will now show the relationship between FedGATE/VRL-SGD and the centralized version of LED. As a first step, we will represent FedGATE/VRL-SGD in a networked form.

FedGATE is described as follows [24, Alg. 3]: For r=0,1,2,…r=0,1,2,\dots, set ϕi,0r=xr\phi_{i,0}^{r}=x^{r}, for t=0,…,\uptau−1t=0,\dots,\uptau-1:

where γ\gamma is a global stepsize parameter. By letting yir=−α\uptauδiry_{i}^{r}=-\alpha\uptau\delta_{i}^{r} and employing the network notation defined in (5) with xr=1⊗xr{\mathbf{x}}^{r}=\mathbf{1}\otimes x^{r}, the method above can be rewritten as:

It’s now evident that when αγ=1\alpha\gamma=1, the update (18) aligns with LED (6) when W=1N11T{\mathbf{W}}=\tfrac{1}{N}\mathbf{1}\mathbf{1}^{\textit{\footnotesize{T}}} and β=1/\uptau\beta=1/\uptau. It’s worth noting that the updates (18) simplify to VRL-SGD when αγ=1\alpha\gamma=1 . In essence, FedGATE with αγ=1\alpha\gamma=1 (or VRL-SGD) corresponds to LED in the fully connected network scenario. This also suggests that FedGATE and VRL-SGD are locally updated versions of PDFP2O (11b) with B=I−1N11T{\mathbf{B}}={\mathbf{I}}-\tfrac{1}{N}\mathbf{1}\mathbf{1}^{\textit{\footnotesize{T}}}.

All of the derivations in this section require appropriate stepsizes tuning and are based on the assumption that the graph is static (i.e., WW is constant). When the stepsizes differ or the graph is dynamic (as in the local update variant), these various representation may not necessarily be equivalent. We note that the analysis techniques from are particularly suited for static graphs and are limited to the strongly-convex case. Moreover, the techniques from are tailored for centralized networks. In contrast, our analyses address the more challenging decentralized connections with local updates; thus, our techniques can be specialized for these methods. In fact, our analysis can provide tighter rates compared to those in . Remark 8 explains how to adapt our techniques to the centralized scenario.

Convergence result

In this section, we present our main convergence findings and discuss how they differ from previous results. Before proceeding, we will review the assumptions necessary for our results to hold, which are standard in the literature .

The weight matrix WW is symmetric, doubly stochastic, and primitive. Moreover, we assume that WW is positive definite. ■\blacksquare

Under Assumption 1, the eigenvalues of WW, denoted by {λi}i=1N\{\lambda_{i}\}_{i=1}^{N}, are all strictly less than one (in magnitude for nonpositive definite WW), with the exception of a single eigenvalue at one, which we denote by λ1\lambda_{1}. The network’s mixing rate is defined as:

Each stochastic gradient ∇Fi(xik;ξik){\nabla}F_{i}(x_{i}^{k};\xi_{i}^{k}) is unbiased with bounded variance:

Each function fi:m→ℜf_{i}:{}^{m}\rightarrow\real is LL-smooth:

for some L>0L>0. Additionally, the aggregate function f(x)=1N∑i=1Nfi(x)f(x)=\frac{1}{N}\sum_{i=1}^{N}f_{i}(x) is bounded below, i.e., f(x)≥f⋆>−∞f(x)\geq f^{\star}>-\infty for every x∈mx\in{}^{m}, where f⋆f^{\star} denotes the optimal value of ff. ■\blacksquare

Under the aforementioned assumption, the aggregate function f(x)=1N∑i=1Nfi(x)f(x)=\frac{1}{N}\sum_{i=1}^{N}f_{i}(x) is also LL-smooth.

Assumptions 1–3 are sufficient to establish convergence under nonconvex settings. We will also study convergence under additional convexity assumption given below.

Each function fi:m→ℜf_{i}:{}^{m}\rightarrow\real is (μ\mu-strongly)convex for some 0≤μ≤L0\leq\mu\leq L. (When μ=0\mu=0, then the functions are simply convex.) ■\blacksquare

2 Main results

We are now ready to present our main findings. The convergence results for nonconvex and convex functions are presented in Theorems 1 and 2, respectively. The final convergence rates derived from these theorems are given in Corollaries 1 and 2. All proofs can be found in the appendices.

Under Assumptions 1–3, and for sufficiently small constant stepsizes α\alpha and β=1/\uptau\beta=1/\uptau, it holds that

Under Assumptions 1–4, and for sufficiently small constant stepsizes α\alpha and β=1/\uptau\beta=1/\uptau, it holds that for μ=0\mu=0 (convex case)

where ρ≜1−λ\rho\triangleq 1-\lambda and a0a_{0} is a constant that depends on the initialization. ■\blacksquare

For the nonconvex case and when σ>0\sigma>0, Theorem 1 shows that the algorithm converges to a radius around some stationary point, which can be controlled by the stepsize α\alpha. Without any additional assumptions, a stationary point is the best guarantee possible and is a satisfactory criterion to measure the performance of distributed methods with nonconvex objectives . For the convex case, Theorem 2 shows that the algorithm converges around some optimal solution controlled by the stepsize α\alpha.

Suppose the conditions of Theorem 1 are met. Then, in the noiseless deterministic case where σ=0\sigma=0, substituting this into (22) gives the nonconvex rate:

Therefore, LED converges exactly in the deterministic case with a rate of 1/R1/R. Similar results can be obtained for the convex cases. ■\blacksquare

For the stochastic case, the stepsize α\alpha is tuned based on RR to obtain the following result.

For the nonconvex, convex, and strongly convex cases, there exists a stepsize α\alpha that yields the following rates.

Here, xˉ0≜(1/N)∑i=1Nxi0\bar{x}^{0}\triangleq(1/N)\sum_{i=1}^{N}x_{i}^{0} and ς02≜1N∑i=1N∥∇fi(xˉ0)−∇f(xˉ0)∥2\varsigma_{0}^{2}\triangleq\frac{1}{N}\sum_{i=1}^{N}\|{\nabla}f_{i}(\bar{x}^{0})-{\nabla}f(\bar{x}^{0})\|^{2}.

The stepsize yielding the results in Corollary 2 is intricate and not very practical. It’s chosen mainly for theoretical reasons, as it provides the optimal convergence rate based on our bounds. We opted for this choice to ensure fair comparisons with SCAFFOLD, Local-DSGD, and K-GT, which also tune the stepsize in a similar manner.

In practice, we typically set α=1/R\alpha=1/\sqrt{R} for both the nonconvex and convex cases, and α=1/R\alpha=1/R for the strongly-convex case. For instance, if we plug in α=1L+\uptauR/N\alpha=\frac{1}{L+\sqrt{\uptau R/N}} into (22), we obtain the rate

For large RR, the above rate is also 1/N\uptauR1/\sqrt{N\uptau R} consistent with (25). ■\blacksquare

Suppose Assumptions 2–3 hold. Then, with the appropriate parameters, the LED under the server-workers scenario, as listed in Alg. 2, converges at the rate

Discussion of our results. We will discuss our results for the nonconvex case; similar arguments apply for the convex case. For large RR, the higher-order terms from (25) can be neglected, and the dominant part becomes on the order of (σN\uptauR)12\left(\frac{\sigma}{N\uptau R}\right)^{\frac{1}{2}}. This suggests that to achieve an ϵ\epsilon accuracy, we need R≥1N\uptauϵ2R\geq\frac{1}{N\uptau\epsilon^{2}}. In this scenario, the advantages of NN and \uptau\uptau are evident. Furthermore, the number of communication rounds required to achieve ϵ\epsilon accuracy decreases linearly with NN; this characteristic is termed linear speed-up . When RR is not sufficiently large, then the higher-order terms, specifically (1(1−λ)1/3(σ\uptauR)23)\left(\frac{1}{(1-\lambda)^{1/3}}(\frac{\sigma}{\sqrt{\uptau}R})^{\frac{2}{3}}\right) and (1/(1−λ)+ς02R)\left(\frac{1/(1-\lambda)+\varsigma_{0}^{2}}{R}\right), cannot be ignored as they may slow down the convergence. For instance, when the network is sparse, the quantity 1−λ1-\lambda can be extremely small as λ≈0\lambda\approx 0; in this situation, the rate becomes slower, as will be discussed in Section 5. Table 2 lists the convergence rate of LED compared to state-of-the-art results, in terms of the number of communication rounds needed to achieve ϵ\epsilon accuracy.

Compared to our results, observe that Local-DSGD introduces an additional term ς1−λ1ϵ3/2\frac{\varsigma}{1-\lambda}\frac{1}{\epsilon^{3/2}} (or ς1−λ1ϵ1/2\frac{\varsigma}{1-\lambda}\frac{1}{\epsilon^{1/2}} for the strongly convex case) where ς\varsigma represents the local functions heterogeneity constant, such that 1N∑i=1N∥∇fi(x)−∇f(x)∥2≤ς2\frac{1}{N}\sum_{i=1}^{N}\|\nabla f_{i}(x)-\nabla f(x)\|^{2}\leq\varsigma^{2}. This additional term result in suboptimal convergence rates, even in deterministic (σ=0\sigma=0) scenarios, causing a significant slow down in convergence (refer to the Simulation section for more details). When compared to KK-GT , the second and third terms are (σ(1−λ)2\uptau)1ϵ3/2+1(1−λ)2ϵ\left(\frac{\sigma}{(1-\lambda)^{2}\sqrt{\uptau}}\right)\frac{1}{\epsilon^{3/2}}+\frac{1}{(1-\lambda)^{2}\epsilon}, whereas in our rate for LED, these are (σ(1−λ)\uptau)1ϵ3/2+1(1−λ)ϵ(\frac{\sigma}{\sqrt{(1-\lambda)\uptau}})\frac{1}{\epsilon^{3/2}}+\frac{1}{(1-\lambda)\epsilon}. The factor 1−λ1-\lambda becomes notably small for sparse networks, which indicates that the performance of KK-GT may degrade significantly relative to LED in sparsely connected networks. Note that considering a single local step, with \uptau=1\uptau=1, our rates align with the best-established decentralized rates .

The table also enumerates the rate of the centralized method SCAFFOLD as cited in . For centralized networks defined by W=(1/N)1T1W=(1/N)\mathbf{1}^{\textit{\footnotesize{T}}}\mathbf{1}, our rate matches that of SCAFFOLD . In fact, our rate is more refined than VRL-SGD, presented in . Specifically, our bound allows for setting \uptau=O(1Nϵ)\uptau={\mathcal{O}}\left(\frac{1}{N\epsilon}\right). In this scenario, the number of communication rounds necessary to achieve ϵ\epsilon precision is given by R=O(1ϵ)R={\mathcal{O}}\left(\frac{1}{\epsilon}\right). In contrast, suggests that achieving ϵ\epsilon precision requires a communication round count worse by a factor of NN, specifically R=O(Nϵ)R={\mathcal{O}}\left(\frac{N}{\epsilon}\right) (as seen in [24, Table 6]). Moreover, in the convex scenario, our rate R=O(1ϵ)R={\mathcal{O}}\left(\frac{1}{\epsilon}\right) is sharper than that of FedGate from , which is R=O(1ϵlog⁡(1ϵ))R={\mathcal{O}}\left(\frac{1}{\epsilon}\log\left(\frac{1}{\epsilon}\right)\right) (also observed in [24, Table 6]). This implies that when adapting our analysis to centralized networks, we provide new and improved analyses for the methods VRL-SGD and FedGATE .

Numerical simulations

In this section, we use numerical simulations to demonstrate and validate our findings on the logistic regression problem with a nonconvex regularizer given by

Simulation results. Figure 1 compares our LED method with the decentralized methods L-DSGD , LGT , and KK-GT for different local steps \uptau=1\uptau=1, \uptau=5\uptau=5, and \uptau=10\uptau=10. In order to fairly compare the convergence rates of these methods, we individually tune the parameters of each algorithm so that each method reaches a predetermined error of 10−410^{-4} as quickly as possible. We used a ring (cycle) network and the weight matrix was generated using the Metropolis rule with λ≈0.943\lambda\approx 0.943. We observe that LED outperforms all the other methods as we increase the number of local steps (rightmost plot). L-DSGD performs poorly because it cannot handle the local functions heterogeneity across the nodes. It’s worth noting that increasing the number of local steps reduces the communication required to achieve the same level of accuracy.

Figure 2 shows the results against the centralized methods SCAFFOLD and Local-SGD; in this case, the network is fully connected W=(1/N)11TW=(1/N)\mathbf{1}\mathbf{1}^{\textit{\footnotesize{T}}}. We observe that both LED and SCAFFOLD perform similarly. Furthermore, the performance of Local-DSGD degrades as the number of local updates increases, as expected. It’s worth noting that in these figures, the horizontal-axis refers to the number of communication rounds. For GT and SCAFFOLD, each agent must communicate two vectors to its neighbors per communication round. In contrast, LED and Local-DSGD requires the communication of only one vector.

Is local steps always beneficial? From Figure 1, it can be observed that we can save on communication when we increase the number of local steps. However, these results hold for the stochastic case, and the significant benefits primarily arise from using more data per communication (similar to batch training). To further test the benefits of local steps, we consider the deterministic case, in which σ=0\sigma=0. Figures 3(a)–3(b) depict the results of LED with different local steps and various topologies under both heterogeneous and similar data regimes. When the data is heterogeneous, Figure 3(a) shows that the benefits of local steps are only visible in the fully connected network. Conversely, when the data is similar across nodes (minimal heterogeneity), there is a significant benefit to local steps, as shown in Figure 3(b). However, for a sparse network (ring topology), we don’t observe any noticeable advantage. These findings suggest that local steps can be beneficial when the network is well-connected and/or when data is consistent across the network.This can be explained by our theoretical results as follows. Assuming the network is well-connected (ρ≈1\rho\approx 1), we can disregard the term 1/ρ1/\rho from the rate expression (25). The dominant part of the rate then simplifies to O(σN\uptauR)12+O(ς02R){\mathcal{O}}(\frac{\sigma}{N\uptau R})^{\frac{1}{2}}+{\mathcal{O}}(\frac{\varsigma_{0}^{2}}{R}). In this scenario, when data heterogeneity is relatively small (e.g., ς02=O(1/\uptau)\varsigma_{0}^{2}={\mathcal{O}}(1/\sqrt{\uptau})), the benefits of \uptau\uptau become evident. Conversely, if the network is sparse (ρ≈0\rho\approx 0), then 1/ρ1/\rho can be quite large. In such a case, the dominant terms in the rate would be O(σN\uptauR)12+O(1ρ1/3(σ\uptauR)23)+O(ς02ρR){\mathcal{O}}(\frac{\sigma}{N\uptau R})^{\frac{1}{2}}+{\mathcal{O}}(\frac{1}{\rho^{1/3}}(\frac{\sigma}{\sqrt{\uptau}R})^{\frac{2}{3}})+{\mathcal{O}}(\frac{\varsigma_{0}^{2}}{\rho R}). Notice that 1/ρ1/\rho adversely affects the higher-order terms, thus slowing down convergence.

Concluding remarks

In this work, we proposed Local Exact-Diffusion (LED), a locally updated method inspired by the framework from and the Exact-Diffusion method . We demonstrated that LED can be interpreted as the primal-dual method PDFP2O/PAPC from . We also explored its connection with the following methods: Exact-Diffusion (ED) , NIDS , D2 , Scaffnew , VRL-SGD , and FedGate . We proved the convergence of LED in both convex and nonconvex settings and established bounds that offer improvements over existing decentralized methods. Finally, we provided numerical simulations to illustrate the effectiveness of the proposed algorithm.

A promising direction for future research involves examining LED in the context of probabilistic local updates. It’s worth exploring if the current analysis can be integrated with techniques from Scaffnew/ProxSkip to understand the advantages of local steps in non-strongly-convex scenarios. Another potential avenue is broadening LED to accommodate time-varying stochastic graphs. Furthermore, it would be insightful to see if decentralized local methods might also offer advantages for other network coupling constraints, as seen in multitask problems or the distributed feature problem.

Acknowledgments

The author would like to thank Kun Yuan for his insightful discussions on parts of the manuscript.

References

Appendix A Preliminary transformation

In this section, we will convert LED updates (6) into another form more suitable for our analysis. To that end, we will first introduce some notation that complements the notation (5).

The following quantities will be used in the analysis:

A.2 Weight matrix decomposition

We will use the structure of the weight matrix WW in our analysis, which is a requirement for our proof. As a result, we will now go over some facts about the matrix WW. When Assumption 1 holds, then the weight matrix WW can be decomposed as follows :

where Λ^=diag{λi}i=2N\widehat{\Lambda}={\rm diag}\{\lambda_{i}\}_{i=2}^{N} and the matrix Q^\widehat{Q} has size N×(N−1){N\times(N-1)}, and satisfies Q^Q^T=IN−1N11T\widehat{Q}\widehat{Q}^{\textit{\footnotesize{T}}}=I_{N}-\tfrac{1}{N}\mathbf{1}\mathbf{1}^{\textit{\footnotesize{T}}} and 1TQ^=0\mathbf{1}^{\textit{\footnotesize{T}}}\widehat{Q}=0. It follows that W≜W⊗Im{\mathbf{W}}\triangleq W\otimes I_{m} can be decomposed as

Below, we provide a list of relevant facts and properties concerning the aforementioned decomposition.

The matrix B^≜I−Λ^\widehat{{\mathbf{B}}}\triangleq{\mathbf{I}}-\widehat{\mathbf{\Lambda}} satisfies

where λ‾\underline{\lambda} is the smallest nonzero eigenvalue of WW and λ\lambda is the network’s mixing rate introduced in (19).

A.3 Transformed recursion

To obtain our result, we will perform a series of transformations that will eventually lead us to the critical result specified in Lemma 1. In order to study the communication complexity of LED, we begin by representing it in terms of communication rounds. It holds true when iterating through (6a) updates:

Observe that under our initialization, y0=(I−W)x0{\mathbf{y}}^{0}=({\mathbf{I}}-{\mathbf{W}}){\mathbf{x}}^{0}, yr{\mathbf{y}}^{r} will always be in the range of I−W{\mathbf{I}}-{\mathbf{W}}. As a result, we have (1T⊗Im)yr=0(\mathbf{1}^{\textit{\footnotesize{T}}}\otimes I_{m}){\mathbf{y}}^{r}=0 for all rr. Multiplying both sides of (34a) by (1/N)(1T⊗Im)(1/N)(\mathbf{1}^{\textit{\footnotesize{T}}}\otimes I_{m}) on the left and using the definitions in (30) , namely ∇f‾(xr)=1N∑i=1N∇fi(xir),sˉtr=1N∑i=1N(∇Fi(ϕi,tr;ξi,t)−∇fi(ϕi,tr))\overline{{\nabla}f}({\mathbf{x}}^{r})=\frac{1}{N}\sum_{i=1}^{N}{\nabla}f_{i}(x_{i}^{r}),\bar{s}_{t}^{r}=\frac{1}{N}\sum_{i=1}^{N}(\nabla F_{i}(\phi^{r}_{i,t};\xi_{i,t})-\nabla f_{i}(\phi^{r}_{i,t}))], we get

Equation (35) above describes how the average (centroid) vector evolves in terms of communication rounds, which will be important in our analysis. Now, using the definitions in (30), namely zr=yr+αβ∇f(xˉr){\mathbf{z}}^{r}={\mathbf{y}}^{r}+\frac{\alpha}{\beta}{\nabla}{\mathbf{f}}(\bar{{\mathbf{x}}}^{r}) and str=∇F(Φtr;ξt)−∇f(Φtr){\mathbf{s}}_{t}^{r}=\nabla{\mathbf{F}}(\bm{\Phi}^{r}_{t};\bm{\xi}_{t})-\nabla{\mathbf{f}}(\bm{\Phi}^{r}_{t}), the update (34) can be equivalently described as

The introduction of zr{\mathbf{z}}^{r} is inspired from . The quantity zr{\mathbf{z}}^{r} can be interpreted as a variable that tracks the average gradient vector 1⊗∇f(xˉr)\mathbf{1}\otimes{\nabla}f(\bar{x}^{r}). Observe that by (1T⊗Im)yr=0(\mathbf{1}^{\textit{\footnotesize{T}}}\otimes I_{m}){\mathbf{y}}^{r}=0 and (30e), we have

We will next transform the updates (36) into quantities that measure how far xr{\mathbf{x}}^{r} and zr{\mathbf{z}}^{r} deviate from the averages xˉr=1N⊗xˉr\bar{{\mathbf{x}}}^{r}=\mathbf{1}_{N}\otimes\bar{x}^{r} and zˉr≜1N⊗zˉr\bar{{\mathbf{z}}}^{r}\triangleq\mathbf{1}_{N}\otimes\bar{z}^{r}, respectively. To do so, we will leverage the structure and properties of the weight matrix W{\mathbf{W}} (32). ■\blacksquare

Using the decomposition of W{\mathbf{W}} given in (31), it holds that Q^TW=Λ^Q^T\widehat{{\mathbf{Q}}}^{\textit{\footnotesize{T}}}{\mathbf{W}}=\widehat{\mathbf{\Lambda}}\widehat{{\mathbf{Q}}}^{\textit{\footnotesize{T}}}. Therefore, multiplying both sides of (36) by Q^T\widehat{{\mathbf{Q}}}^{\textit{\footnotesize{T}}}, it holds that

Rewriting (38) in matrix notation, we have

Note that from (32), we have ∥Q^Txr∥2=∥xr−xˉr∥2\|\widehat{{\mathbf{Q}}}^{\textit{\footnotesize{T}}}{\mathbf{x}}^{r}\|^{2}=\|{\mathbf{x}}^{r}-\bar{{\mathbf{x}}}^{r}\|^{2} and similarly ∥Q^Tzr∥2=∥zr−zˉr∥2\|\widehat{{\mathbf{Q}}}^{\textit{\footnotesize{T}}}{\mathbf{z}}^{r}\|^{2}=\|{\mathbf{z}}^{r}-\bar{{\mathbf{z}}}^{r}\|^{2}. Thus, the above describes how the nodes vectors deviates from the averages. For β=1/\uptau\beta=1/\uptau, we will let

If the norm of the matrix D{\mathbf{D}} is less than one ∥D∥<1\|{\mathbf{D}}\|<1, then the updates (A.3) can be used to directly measure the deviation from the averages. However, even though the eigenvalues of D{\mathbf{D}} are less than one, its norm is not guaranteed to satisfy ∥D∥<1\|{\mathbf{D}}\|<1; but, we can decompose D{\mathbf{D}} and transform (A.3) into a more suitable form for our analysis, as shown in the important result below.

Suppose that Assumption 1 holds and β=1/\uptau\beta=1/\uptau, then

where Δ\mathbf{\Delta} is a matrix with norm δ≜∥Δ∥=λ<1\delta\triangleq\|\mathbf{\Delta}\|=\sqrt{\lambda}<1, B^≜I−Λ^\widehat{{\mathbf{B}}}\triangleq{\mathbf{I}}-\widehat{\mathbf{\Lambda}}, and

Here, P{\mathbf{P}} is some permutation matrix and ȷ=−1\jmath=\sqrt{-1} is the imaginary number.

The proof exploits the special structure of the matrix D{\mathbf{D}} given in (40). First note that if

where bi,ci,di,eib_{i},c_{i},d_{i},e_{i} are constants, then there exists a permutation matrix PP such that

Since the blocks of D{\mathbf{D}} given in (40) are diagonal matrices, there exists a permutation matrix P{\mathbf{P}} such that

The eigenvalues of DiD_{i} (i=2,…,Ni=2,\ldots,N) are

Notice that ∣δ(1,2),i∣<1|\delta_{(1,2),i}|<1 when −13<λi<1-\frac{1}{3}<\lambda_{i}<1, which holds under Assumption 1 since W>0W>0, i.e., 0<λi<10<\lambda_{i}<1 (i=2,…,Ni=2,\dots,N). For 0<λi<10<\lambda_{i}<1, the eigenvalues of DiD_{i} are complex and distinct:

where ȷ2=−1\jmath^{2}=-1. Through algebraic multiplication it can be verified that Di=ViΔiVi−1D_{i}=V_{i}\Delta_{i}V_{i}^{-1} where

We conclude that D=PTVΔV−1P{\mathbf{D}}={\mathbf{P}}^{\textit{\footnotesize{T}}}{\mathbf{V}}\mathbf{\Delta}{\mathbf{V}}^{-1}{\mathbf{P}} where V=blkdiag{Vi}i=2N⊗Im{\mathbf{V}}={\rm blkdiag}\{V_{i}\}_{i=2}^{N}\otimes I_{m} and Δ=blkdiag{Δi}i=2N⊗Im\mathbf{\Delta}={\rm blkdiag}\{\Delta_{i}\}_{i=2}^{N}\otimes I_{m}. Therefore, left multiplying both sides of (A.3) by V^−1\widehat{{\mathbf{V}}}^{-1} where V^=PTV\widehat{{\mathbf{V}}}={\mathbf{P}}^{\textit{\footnotesize{T}}}{\mathbf{V}} gives (41). Exploiting the structure of V−1=blkdiag{Vi−1}i=2N⊗Im{\mathbf{V}}^{-1}={\rm blkdiag}\{V_{i}^{-1}\}_{i=2}^{N}\otimes I_{m} where Vi−1V_{i}^{-1} is defined in (46) and using (44), we get

Using similar arguments for V{\mathbf{V}} gives (42b).

The transformation in Lemma 1 is needed for decentralized network analysis since, as explained before, the norm of the matrix D{\mathbf{D}} (40) is not necessarily less than one. However, when the network is fully-connected (centralized case), we have Λ^=0\widehat{\mathbf{\Lambda}}=\mathbf{0}, and thus, it follows from (A.3) that Q^Txr=0\widehat{{\mathbf{Q}}}^{\textit{\footnotesize{T}}}{\mathbf{x}}^{r}=\mathbf{0} and when β=1/\uptau\beta=1/\uptau, we have

In this case, the analysis can be greatly simplified since D=0{\mathbf{D}}=\mathbf{0}. This specialization covers the method VRL-SGD . A similar approach can also be adapted to FedGate , as described in (18). Doing so eventually leads to the result in Corollary LABEL:cor_centralized. See Appendix D. ■\blacksquare

Appendix B Convergence analysis

In this section, we will prove Theorems 1 and 2. The proof utilizes equations (35) and (41) derived in the previous section. We want to emphasize that some of the bounds may appear tedious, even though they only involve basic algebra. Such complexity is common in decentralized analysis, as seen in, for example, .

In the proof, we will use the following useful results and facts. (You may overlook and refer to this subsection later in the proofs.)

Since the squared norm ∥⋅∥2\|\cdot\|^{2} is convex, applying Jensen’s inequality, it holds that

Many bounds later in the proofs uses inequality (48a) without referring to it to avoid redundancies.

Taking the squared norm on both sides of (42a) and using (48a), it holds that:

where PuT{\mathbf{P}}_{u}^{\textit{\footnotesize{T}}} and PlT{\mathbf{P}}_{l}^{\textit{\footnotesize{T}}} are the upper and lower blocks of PT=[PuT;PlT]{\mathbf{P}}^{\textit{\footnotesize{T}}}=[{\mathbf{P}}_{u}^{\textit{\footnotesize{T}}};{\mathbf{P}}_{l}^{\textit{\footnotesize{T}}}]. It follows that:

where we used ∥PuT+PlT∥2≤4\|{\mathbf{P}}_{u}^{\textit{\footnotesize{T}}}+{\mathbf{P}}_{l}^{\textit{\footnotesize{T}}}\|^{2}\leq 4 and ∥PuT−PlT∥2≤4\|{\mathbf{P}}_{u}^{\textit{\footnotesize{T}}}-{\mathbf{P}}_{l}^{\textit{\footnotesize{T}}}\|^{2}\leq 4 since P{\mathbf{P}} is a permutation matrix ∥P∥=1\|{\mathbf{P}}\|=1.

The last step holds by using Jensen’s inequality (48) and (32)–(33). Following similar arguments, it can be shown that the squared norm of

B.2 Key bounds

In this section, we derive some key bounds that will be used to establish our result for both nonconvex and convex cases.

The first bound involves the cumulative deviation of the local updates from the averaged vector at the previous communication round defined as

Let Assumptions 2–3 hold, then for α≤122L\uptau\alpha\leq\frac{1}{2\sqrt{2}L\uptau} we have

The proof extends the techniques from [10, Lemma 8]. When \uptau=1\uptau=1, then ϕi,0=xir\phi_{i,0}=x_{i}^{r} for all ii and

Now suppose that \uptau≥2\uptau\geq 2. Then, using (6a), it holds that

The second inequality uses (48b) with θ=1−1\uptau\theta=1-\frac{1}{\uptau}. The last inequality holds for 2α2\uptauL2≤14(\uptau−1)2\alpha^{2}\uptau L^{2}\leq\frac{1}{4(\uptau-1)}, which is satisfied if α≤122L\uptau\alpha\leq\frac{1}{2\sqrt{2}L\uptau}. Iterating the inequality above for t=0,…,\uptau−1t=0,\dots,\uptau-1:

where in the second and third inequalities we used (1+a\uptau−1)t≤exp⁡(at\uptau−1)≤exp⁡(a)(1+\tfrac{a}{\uptau-1})^{t}\leq\exp(\frac{at}{\uptau-1})\leq\exp(a) for t≤\uptau−1t\leq\uptau-1. Summing over ii and tt:

The next result measures the deviation from the average vector introduced in Lemma 1.

where δ\delta and B^\widehat{{\mathbf{B}}} were defined in Lemma 1.

From now on we use the notation ∑t≡∑t=0\uptau−1\sum\limits_{t}\equiv\sum\limits_{t=0}^{\uptau-1} and ∑i≡∑i=1N\sum\limits_{i}\equiv\sum\limits_{i=1}^{N}. From (41), (53), and (55), we have

where δ=∥Δ∥=λ<1\delta=\|\bm{\Delta}\|=\sqrt{\lambda}<1. The second line follows from the unbiased stochastic gradient condition (20a) and the last step uses Jensen’s inequality (48). Using 2⟨sr,hr+1⟩≤∥sr∥2+∥hr+1∥22\langle{\mathbf{s}}^{r},{\mathbf{h}}^{r+1}\rangle\leq\|{\mathbf{s}}^{r}\|^{2}+\|{\mathbf{h}}^{r+1}\|^{2} and 1≤1/(1−δ)1\leq 1/(1-\delta) gives

Substituting the previous two bounds into (60) and taking expectation gives

Substituting (58) into the above inequality yields

B.3 Nonconvex case (Theorem 1)

The nonconvex proof begins with the following bound for any LL-smooth function ff :

Recall from (35) that \bar{x}^{r+1}=\bar{x}^{r}-\frac{\alpha}{N}\sum_{t=0}^{\uptau-1}\sum_{i=1}^{N}\big{(}\nabla f_{i}(\phi^{r}_{i,t})+s^{r}_{i,t}\big{)}. Substituting y=xˉr+1y=\bar{x}^{r+1} and z=xˉrz=\bar{x}^{r} into inequality (63), and taking conditional expectation, we get

where the second bound holds from Jensen’s inequality (48). Combining the last two equations and taking expectation yields

Substituting the bound (58) into inequality (66) and taking expectation yields

When α≤142\uptauL\alpha\leq\frac{1}{4\sqrt{2}\uptau L}, we can upper bound the previous inequality by

The last step uses Jensen’s inequality and 36α2\uptau2L2∥B^−1∥(1−δ)≤1\frac{36\alpha^{2}\uptau^{2}L^{2}\|\widehat{{\mathbf{B}}}^{-1}\|}{(1-\delta)}\leq 1, i.e.,

Averaging over r=1,…,Rr=1,\dots,R and using (71), it holds that

Adding ∥d^0∥2R\frac{\|\widehat{{\mathbf{d}}}^{0}\|^{2}}{R} to both sides of the previous inequality and using ∥d^0∥2R≤∥d^0∥2(1−δ)R\frac{\|\widehat{{\mathbf{d}}}^{0}\|^{2}}{R}\leq\frac{\|\widehat{{\mathbf{d}}}^{0}\|^{2}}{(1-\delta)R}, we get

Substituting inequality (74) into (68) and rearranging, we obtain

Now from (49), we can bound ∥d^0∥2\|\widehat{{\mathbf{d}}}^{0}\|^{2} by

where the second inequality we used (30e) (z0=y0+α\uptau∇f(xˉ0){\mathbf{z}}^{0}={\mathbf{y}}^{0}+\alpha\uptau{\nabla}{\mathbf{f}}(\bar{{\mathbf{x}}}^{0})) and Jensen’s inequality. The last inequality holds under initialization y0=0{\mathbf{y}}^{0}=\mathbf{0}. We conclude that

where ς02=1N∥∇f(xˉ0)−1⊗∇f(xˉ0)∥2\varsigma_{0}^{2}=\frac{1}{N}\|{\nabla}{\mathbf{f}}(\bar{{\mathbf{x}}}^{0})-\mathbf{1}\otimes{\nabla}f(\bar{{\mathbf{x}}}^{0})\|^{2}.

B.4 Convex cases (Theorem 2)

Recall from (35) that \bar{x}^{r+1}=\bar{x}^{r}-\frac{\alpha}{N}\sum_{t=0}^{\uptau-1}\sum_{i=1}^{N}\big{(}\nabla f_{i}(\phi^{r}_{i,t})+s^{r}_{i,t}\big{)}. Thus, it holds that

The first inequality we used the zero mean condition (20a) and Jensen’s inequality. The second inequality holds from bounded noise variance condition from Assumption 2 and Jensen’s inequality. We now bound the cross term by using the bound [10, Lemma 5]:

for any LL-smooth and μ\mu-strongly convex function gg. Using (81), we can bound the cross term as follows

where ∥Φ^r∥2=∑i=1N∑t=0\uptau−1∥ϕi,tr−xˉr∥2\|\bm{\widehat{\Phi}}^{r}\|^{2}=\sum_{i=1}^{N}\sum_{t=0}^{\uptau-1}\|\phi^{r}_{i,t}-\bar{x}^{r}\|^{2}. Substituting the previous bound into (80) and taking expectation gives

where ∇f‾(Φtr)=1N∑i∇fi(ϕi,tr)\overline{\nabla f}(\bm{\Phi}^{r}_{t})=\tfrac{1}{N}\textstyle\sum\limits_{i}{\nabla}f_{i}(\phi^{r}_{i,t}). Note that

Substituting the previous bound into (B.4) gives

The last inequality holds for 1+2α\uptauL≤3/21+2\alpha\uptau L\leq 3/2 or

Substituting the bound (58) into the above inequality gives

In the last step, we used ∥∇f(xˉr)∥2≤2L(f(xˉr)−f(x⋆))\|{\nabla}f(\bar{x}^{r})\|^{2}\leq 2L(f(\bar{x}^{r})-f(x^{\star})), which is satisfied under Assumption 4 . Using (2α\uptau−8Lα2\uptau2−96α3\uptau3L2)≤α\uptau(2\alpha\uptau-8L\alpha^{2}\uptau^{2}-96\alpha^{3}\uptau^{3}L^{2})\leq\alpha\uptau, i.e.,

where the last inequality holds when 36α2\uptau2L2∥B^−1∥(1−δ)≤1\frac{36\alpha^{2}\uptau^{2}L^{2}\|\widehat{{\mathbf{B}}}^{-1}\|}{(1-\delta)}\leq 1, i.e.,

Substituting the bound (58) into the above inequality yields

Using the condition δ+527\uptau2α2L2(1−δ)≤1+δ2≜δˉ\delta+\tfrac{527\uptau^{2}\alpha^{2}L^{2}}{(1-\delta)}\leq\frac{1+\delta}{2}\triangleq\bar{\delta}, which holds when

the right hand side can be upper bounded by

For convex but not strongly-convex, we have μ=0\mu=0 and equation (87) becomes

Plugging ∥∇f(xˉr)∥2≤2L(f(xˉr)−f(x⋆))\|{\nabla}f(\bar{x}^{r})\|^{2}\leq 2L(f(\bar{x}^{r})-f(x^{\star})) into (92) gives

Iterating and averaging over r=1,…,Rr=1,\dots,R

Adding ∥d^0∥2R\frac{\|\widehat{{\mathbf{d}}}^{0}\|^{2}}{R} to both sides of the previous inequality and using ∥d^0∥2R≤∥d^0∥2(1−δ)R\frac{\|\widehat{{\mathbf{d}}}^{0}\|^{2}}{R}\leq\frac{\|\widehat{{\mathbf{d}}}^{0}\|^{2}}{(1-\delta)R}, we get

Substituting inequality (98) into (95) and rearranging, we obtain

where ς02=1N∥∇f(xˉ0)−1⊗∇f(xˉ0)∥2\varsigma_{0}^{2}=\frac{1}{N}\|{\nabla}{\mathbf{f}}(\bar{{\mathbf{x}}}^{0})-\mathbf{1}\otimes{\nabla}f(\bar{{\mathbf{x}}}^{0})\|^{2}.

B.4.2 Strongly-convex case μ>0𝜇0\mu>0

where the last inequality follows from ∥∇f(xˉr)∥2≤L2∥xˉr−x⋆∥2\|{\nabla}f(\bar{x}^{r})\|^{2}\leq L^{2}\|\bar{x}^{r}-x^{\star}\|^{2}. It follows that

where the last inequality holds under the step size condition:

Since ρ(A)<1\rho(A)<1, we can iterate inequality (105) to get

Taking the 11-induced-norm and using properties of the (induced) norms, it holds that

The last step holds for 14−768α3\uptau3L3≥18\frac{1}{4}-768\alpha^{3}\uptau^{3}L^{3}\geq\frac{1}{8} or 768α3\uptau3L3≤18768\alpha^{3}\uptau^{3}L^{3}\leq\frac{1}{8}, which holds under condition (85). Therefore,

Substituting the above into (109) and using (106), we obtain

Appendix C Proof of Corollary 2

The final rate can be obtained by tuning the stepsize in a way similar to .

If all nodes use equal initialization, then equation (B.3) ((22) from Theorem 1) reduces to

Note that the above holds under the condition:

where 1α‾\frac{1}{\underline{\alpha}} satisfies all stepsize conditions used to derive (22). Setting α=min⁡{(c0c1R)12,(c0c2R)13,1α‾}≤1α‾\alpha=\min\left\{\left(\frac{c_{0}}{c_{1}R}\right)^{\frac{1}{2}},\left(\frac{c_{0}}{c_{2}R}\right)^{\frac{1}{3}},\frac{1}{\underline{\alpha}}\right\}\leq\frac{1}{\underline{\alpha}}. Then we have three cases.

When α=1α‾\alpha=\frac{1}{\underline{\alpha}} and is smaller than both (c0c1R)12\left(\frac{c_{0}}{c_{1}R}\right)^{\frac{1}{2}} and (c0c2R)13\left(\frac{c_{0}}{c_{2}R}\right)^{\frac{1}{3}}, then

When α=(c0c1R)12≤(c0c2R)13\alpha=\left(\frac{c_{0}}{c_{1}R}\right)^{\frac{1}{2}}\leq\left(\frac{c_{0}}{c_{2}R}\right)^{\frac{1}{3}}, then

When α=(c0c2R)13≤(c0c1R)12\alpha=\left(\frac{c_{0}}{c_{2}R}\right)^{\frac{1}{3}}\leq\left(\frac{c_{0}}{c_{1}R}\right)^{\frac{1}{2}}, then

Combining the above three cases together it holds that

Substituting the above into (111), we conclude that

The rate (25) follows by plugging the parameters (112) and using (113).

If we start from equal initialization then the convex bound (B.4.1) also satisfies (111) under condition

Therefore, the rate can be obtained by following the same arguments used for the noncovex case.

Using the stepsize condition used to derive Theorem 2, namely,

and starting from equal initialization, inequality (B.4.2) ((24)) can be upper bounded by

Now we select α=min⁡{ln⁡(max⁡{1,μ\uptau(c0+b0/α‾2)R/c1})μ\uptauR,1α‾}≤1α‾\alpha=\min\left\{\frac{\ln\left(\max\left\{1,\mu\uptau(c_{0}+b_{0}/\underline{\alpha}^{2})R/c_{1}\right\}\right)}{\mu\uptau R},\frac{1}{\underline{\alpha}}\right\}\leq\frac{1}{\underline{\alpha}} to get the following cases.

If α=ln⁡(max⁡{1,μ(c0+b0/α‾2)R/c1})μ\uptauR≤1α‾\alpha=\frac{\ln\left(\max\left\{1,\mu(c_{0}+b_{0}/\underline{\alpha}^{2})R/c_{1}\right\}\right)}{\mu\uptau R}\leq\frac{1}{\underline{\alpha}} then

Otherwise α=1α‾≤ln⁡(max⁡{1,μ\uptau(c0+b0/α‾2)/c1})μ\uptauR\alpha=\frac{1}{\underline{\alpha}}\leq\frac{\ln\left(\max\left\{1,\mu\uptau(c_{0}+b_{0}/\underline{\alpha}^{2})/c_{1}\right\}\right)}{\mu\uptau R} and

Collecting these cases together into (116), we obtain

Plugging in the parameters and using (115) gives the final rate (27).

Appendix D LED analysis in the centralized server-workers setup

In this section, we will analyze LED within the server-workers setup. Specifically, we will examine the algorithm listed in 2. It can be verified that this algorithm is equivalent to Algorithm 1 for the fully connected network case, W=(1/N)11TW=(1/N)\mathbf{1}\mathbf{1}^{\textit{\footnotesize{T}}}, when γ=1\gamma=1. Here, γ\gamma is an additional parameter that allows us to derive tighter bounds. To make this section self-contained, we will revisit steps similar to those in the decentralized case but specialized for the centralized case, leading to simpler steps.

Using the above notation, Algorithm 2 can be described in compact form as follows: Set \upphi0r=xˉr\bm{\upphi}_{0}^{r}=\bar{{\mathbf{x}}}^{r} and do:

For analysis purposes, we also introduce the notation:

D.1 Centroid and gradient deviation

When y0=0{\mathbf{y}}^{0}=\mathbf{0}, the iterates {yr}\{{\mathbf{y}}^{r}\} will always be in the range of I−A{\mathbf{I}}-{\mathbf{A}}, consequently, (1T⊗Im)yr=0(\mathbf{1}^{\textit{\footnotesize{T}}}\otimes I_{m}){\mathbf{y}}^{r}=0 for all rr. Using this and the fact Axˉr=xˉr{\mathbf{A}}\bar{{\mathbf{x}}}^{r}=\bar{{\mathbf{x}}}^{r}, the updates (122) become

Using the definitions in (121) into (123b), we have

Let zˉr=Azr\bar{{\mathbf{z}}}^{r}={\mathbf{A}}{\mathbf{z}}^{r}, then we have for β=1/\uptau\beta=1/\uptau:

where in the last step we used Azˉr=zˉr{\mathbf{A}}\bar{{\mathbf{z}}}^{r}=\bar{{\mathbf{z}}}^{r} and A(I−A)=0{\mathbf{A}}({\mathbf{I}}-{\mathbf{A}})=\mathbf{0}. Therefore,

It follows from (123a) and (125) and that

D.2 Auxiliary bounds

Let Assumptions 2–3 hold, then for α≤122L\uptau\alpha\leq\frac{1}{2\sqrt{2}L\uptau} we have

The last inequality holds for 2α2\uptauL2≤14(\uptau−1)2\alpha^{2}\uptau L^{2}\leq\frac{1}{4(\uptau-1)}, which is satisfied if α≤122L\uptau\alpha\leq\frac{1}{2\sqrt{2}L\uptau}. Iterating the inequality above for t=0,…,\uptau−1t=0,\dots,\uptau-1:

where in the second and third inequalities we used (1+a\uptau−1)t≤exp⁡(at\uptau−1)≤exp⁡(a)(1+\tfrac{a}{\uptau-1})^{t}\leq\exp(\frac{at}{\uptau-1})\leq\exp(a) for t≤\uptau−1t\leq\uptau-1. Summing over ii and tt:

The result follows by using zˉr=α\uptau1N∑i=1N∇fi(xr)\bar{z}^{r}=\alpha\uptau\frac{1}{N}\sum_{i=1}^{N}{\nabla}f_{i}(x^{r}). ∎

For α≤122L\uptau\alpha\leq\frac{1}{2\sqrt{2}L\uptau}, it holds that

From now on we use the notation ∑t≡∑t=0\uptau−1\sum\limits_{t}\equiv\sum\limits_{t=0}^{\uptau-1} and ∑i≡∑i=1N\sum\limits_{i}\equiv\sum\limits_{i=1}^{N}. From (126b)

Substituting (128) into the above inequality yields

D.3 Nonconvex case

Recall from (126a) that x^{r+1}=x^{r}-\frac{\alpha\gamma}{N}\sum_{t=0}^{\uptau-1}\sum_{i=1}^{N}\big{(}\nabla f_{i}(\phi^{r}_{i,t})+s^{r}_{i,t}\big{)}. Substituting y=xr+1y=x^{r+1} and z=xrz=x^{r} into inequality (63), and taking conditional expectation, we get

where the second bound holds from Jensen’s inequality. Combining the last two equations and taking expectation yields

Substituting the bound (128) into inequality (132) and taking expectation yields

When αγ≤142\uptauL\alpha\gamma\leq\frac{1}{4\sqrt{2}\uptau L}, we can upper bound the previous inequality by

Substituting inequality (139) into (135) and rearranging, we obtain

If we set 1−64α2\uptau2L2≥1/21-64\alpha^{2}\uptau^{2}L^{2}\geq 1/2, then it holds that

The proof for the convex cases can also be specialized for the centralized case to obtain the rate given in Table 2.