Tackling the Objective Inconsistency Problem in Heterogeneous Federated Optimization

Jianyu Wang, Qinghua Liu, Hao Liang, Gauri Joshi, H. Vincent Poor

Introduction

Federated learning is an emerging sub-area of distributed optimization where both data collection and model training is pushed to a large number of edge clients that have limited communication and computation capabilities. Unlike traditional distributed optimization where consensus (either through a central server or peer-to-peer communication) is performed after every local gradient computation, in federated learning, the subset of clients selected in each communication round perform multiple local updates before these models are aggregated in order to update a global model.

Heterogeneity in the Number of Local Updates in Federated Learning. The clients participating in federated learning are typically highly heterogeneous, both in the size of their local datasets as well as their computation speeds. The original paper on federated learning proposed that each client performs EE epochs (traversals of their local dataset) of local-update stochastic gradient descent (SGD) with a mini-batch size BB. Thus, if a client has nin_{i} local data samples, the number of local SGD iterations is τi=⌊Eni/B⌋\tau_{i}=\lfloor En_{i}/B\rfloor, which can vary widely across clients. The heterogeneity in the number of local SGD iterations is exacerbated by relative variations in the clients’ computing speeds. Within a given wall-clock time interval, faster clients can perform more local updates than slower clients. The number of local updates made by a client can also vary across communication rounds due to unpredictable straggling or slowdown caused by background processes, outages, memory limitations etc. Finally, clients may use different learning rates and local solvers (instead of vanilla SGD, they may use proximal gradient methods or adaptive learning rate schedules) which may result in heterogeneity in the model progress at each client.

Heterogeneity in Local Updates Causes Objective Inconsistency. Most recent works that analyze the convergence of federated optimization algorithms assume that number of local updates is the same across all clients (that is, τi=τ\tau_{i}=\tau for all clients ii). These works show that periodic consensus between the locally trained client models attains a stationary point of the global objective function F(x)=∑i=1mniFi(x)/nF(\bm{x})=\sum_{i=1}^{m}n_{i}F_{i}(\bm{x})/n, which is a sum of local objectives weighted by the dataset size nin_{i}. However, none of these prior works provides insight into the convergence of local-update or federated optimization algorithms in the practical setting when the number of local updates τi\tau_{i} varies across clients 1,…,m1,\dots,m. In fact, as we show in Section 3, standard averaging of client models after heterogeneous local updates results in convergence to a stationary point – not of the original objective function F(x)F(\bm{x}), but of an inconsistent objective F~(x)\widetilde{F}(\bm{x}), which can be arbitrarily different from F(x)F(\bm{x}) depending upon the relative values of τi\tau_{i}. To gain intuition into this phenomenon, observe in Figure 1 that if client 11 performs more local updates, then the updated x(t+1,0)\bm{x}^{(t+1,0)} strays towards the local minimum x1∗\bm{x}_{1}^{*}, away from the true global minimum x∗\bm{x}^{*}.

The Need for a General Analysis Framework. A naive approach to overcome heterogeneity is to fix a target number of local updates τ\tau that each client must finish within a communication round and keep fast nodes idle while the slow clients finish their updates. This method will ensure objective consistency (that is, the surrogate objective F~(x)\widetilde{F}(\bm{x}) equals to the true objective F(x)F(\bm{x})), nonetheless, waiting for the slowest one can significantly increase the total training time. More sophisticated approaches such as FedProx , VRLSGD and SCAFFOLD , designed to handle non-IID local datasets, can be used to reduce (not eliminate) objective inconsistency to some extent, but these methods either result in slower convergence or require additional communication and memory. So far, there is no rigorous understanding of the objective inconsistency and the speed of convergence for this challenging setting of federated learning with heterogeneous local updates. It is also unclear how to best combine models trained with heterogeneous levels of local progress.

Proposed Analysis Framework to Understand Bias Due to Objective Inconsistency. To the best of our knowledge, this work provides the first fundamental understanding of the bias in the solution (caused by objective inconsistency) and how the convergence rate is influenced by heterogeneity in clients’ local progress. In Section 4 we propose a general theoretical framework that allows heterogeneous number of local updates, non-IID local datasets as well as different local solvers such as GD, SGD, SGD with proximal gradients, gradient tracking, adaptive learning rates, momentum, etc. It subsumes existing methods such as FedAvg and FedProx and provides novel insights on their convergence behaviors.

Proposed Normalized Averaging Method FedNova. In Section 5 we propose FedNova, a method that correctly normalizes local model updates when averaging. The main idea of FedNova is that instead of averaging the cumulative local gradient xi(t,τi)−x(t,0)\bm{x}_{i}^{(t,\tau_{i})}-\bm{x}^{(t,0)} returned by client ii (which performs τi\tau_{i} local updates) in tt-th training round, the aggregator averages the normalized local gradients (xi(t,τi)−x(t,0))/τi(\bm{x}_{i}^{(t,\tau_{i})}-\bm{x}^{(t,0)})/\tau_{i}. FedNova ensures objective consistency while preserving fast error convergence and outperforms existing methods as shown in Section 6. It works with any local solver and server optimizer and is therefore complementary to existing approaches such as . By enabling aggregation of models with heterogeneous local progress, FedNova gives the bonus benefit of overcoming the problem of stragglers, or unpredictably slow nodes by allowing fast clients to perform more local updates than slow clients within each communication round.

System Model and Prior Work

The Federated Heterogeneous Optimization Setting. In federated learning, a total of mm clients aim to jointly solve the following optimization problem:

where pi=ni/np_{i}=n_{i}/n denotes the relative sample size, and Fi(x)=1ni∑ξ∈Difi(x;ξ)F_{i}(\bm{x})=\frac{1}{n_{i}}\sum_{\xi\in\mathcal{D}_{i}}f_{i}(\bm{x};\xi) is the local objective function at the ii-th client. Here, fif_{i} is the loss function (possibly non-convex) defined by the learning model and ξ\xi represents a data sample from local dataset Di\mathcal{D}_{i}. In the tt-th communication round, each client independently runs τi\tau_{i} iterations of local solver (e.g., SGD) starting from the current global model x(t,0)\bm{x}^{(t,0)} to optimize its own local objective.

In our theoretical framework, we treat τi\tau_{i} as an arbitrary scalar which can also vary across rounds. In practice, if clients run for the same local epochs EE, then τi=⌊Eni/B⌋\tau_{i}=\lfloor En_{i}/B\rfloor, where BB is the mini-batch size. Alternately, if each communication round has a fixed length in terms of wall-clock time, then τi\tau_{i} represents the local iterations completed by client ii within the time window and may change across clients (depending on their computation speeds and availability) and across communication rounds.

The Fedavg Baseline Algorithm. Federated Averaging (FedAvg) is the first and most common algorithm used to aggregate these locally trained models at the central server at the end of each communication round. The shared global model is updated as follows:

where xi(t,k)\bm{x}_{i}^{(t,k)} denotes client ii’s model after the kk-th local update in the tt-th communication round and Δi(t)=xi(t,τi)−xi(t,0)\smash{\Delta_{i}^{(t)}=\bm{x}_{i}^{(t,\tau_{i})}-\bm{x}_{i}^{(t,0)}} denotes the cumulative local progress made by client ii at round tt. Also, η\eta is the client learning rate and gig_{i} represents the stochastic gradient over a mini-batch of BB samples. When the number of clients mm is large, then the central server may only randomly select a subset of clients to perform computation at each round.

Convergence Analysis of FedAvg. first analyze FedAvg by assuming the local objectives are identical and show that FedAvg is guaranteed to converge to a stationary point of F(x)F(\bm{x}). This analysis was further expanded to the non-IID data partition and client sampling cases by . However, in all these works, they assume that the number of local steps and the client optimizer are the same across all clients. Besides, asynchronous federated optimization algorithms proposed in take a different approach of allowing clients make updates to stale versions of the global model, and their analyses are limited to IID local datasets and convex local functions.

FedProx: Improving FedAvg by Adding a Proximal Term. To alleviate inconsistency due to non-IID data and heterogeneous local updates, proposes adding a proximal term μ2∥x−x(t,0)∥2\frac{\mu}{2}\|\bm{x}-\bm{x}^{(t,0)}\|^{2} to each local objective, where μ≥0\mu\geq 0 is a tunable parameter. This proximal term pulls each local model backward closer to the global model x(t,0)\bm{x}^{(t,0)}. Although empirically shows that FedProx improves FedAvg, its convergence analysis is limited by assumptions that are stronger than previous FedAvg analysis and only works for sufficiently large μ\mu. Since FedProx is a special case of our general framework, our convergence analysis provides sharp insights into the effect of μ\mu. We show that a larger μ\mu mitigates (but does not eliminate) objective inconsistency, albeit at an expense of slower convergence. Our proposed FedNova method can improve FedProx by guaranteeing consistency without slowing down convergence.

Improving FedAvg via Momentum and Cross-client Variance Reduction. The performance of FedAvg has been improved in recent literature by applying momentum on the server side , or using cross-client variance reduction such as VRLSGD and SCAFFOLD . Again, these works do not consider heterogeneous local progress. Our proposed normalized averaging method FedNova is orthogonal to and can be easily combined with these acceleration or variance-reduction techniques. Moreover, FedNova is also compatible with and complementary to gradient compression/quantization and fair aggregation techniques .

A Case Study to Demonstrate the Objective Inconsistency Problem

Below, we show that the convergence point of FedAvg can be arbitrarily away from x∗\bm{x}^{*}.

For the objective function in 3, if client ii performs τi\tau_{i} local steps per round, then FedAvg (with sufficiently small learning rate η\eta, deterministic gradients and full client participation) will converge to

Convergence Problem in Other Federated Algorithms. We can generalize Lemma 1 to the case of FedProx to demonstrate its convergence gap, as given in Appendix A. From the simulations shown in Figure 2, observe that FedProx can slightly improve on the optimality gap of FedAvg, but it converges slower. Besides, previous cross-client variance reduction methods such as variance-reduced local SGD (VRLSGD) and SCAFFOLD are only designed for homogeneous local steps case. In the considered heterogeneous setting, if we replace the same local steps τ\tau in VRLSGD by different τi\tau_{i}’s, then we observe that it has drastically different convergence under different settings and even diverge when clients perform random local steps (see the right panel in Figure 2). These observations emphasize the critical need for a deeper understanding of objective inconsistency and new federated heterogeneous optimization algorithms.

New Theoretical Framework For Heterogeneous Federated Optimization

We now present a general theoretical framework that subsumes a suite of federated optimization algorithms and helps analyze the effect of objective inconsistency on their error convergence. Although the results are presented for the full client participation setting, it is fairly easy to extend them to the case where a subset of clients are randomly sampled in each round In the case of client sampling, the update rule of FedAvg 2 should hold in expectation in order to guarantee convergence . One can achieve this by either (i) sampling qq clients with replacement with respect to probability pip_{i}, and then averaging the cumulative local changes with equal weights, or (ii) sampling qq clients without replacement uniformly at random, and then weighted averaging local changes, where the weight of client ii is re-scaled to pim/qp_{i}m/q. Our convergence analysis can be easily extended to these two cases..

Recall from (2) that the update rule of federated optimization algorithms can be written as x(t+1,0)−x(t,0)=∑i=1mpiΔi(t)\smash{\bm{x}^{(t+1,0)}-\bm{x}^{(t,0)}=\sum_{i=1}^{m}p_{i}\Delta_{i}^{(t)}}, where Δi(t):=x(t,τi)−x(t,0)\smash{\Delta_{i}^{(t)}:=\bm{x}^{(t,\tau_{i})}-\bm{x}^{(t,0)}} denote the local parameter changes of client ii at round tt and pi=ni/np_{i}=n_{i}/n, the fraction of data at client ii. We re-write this update rule in a more general form as follows:

The three key elements τeff,wi\tau_{\text{eff}},w_{i} and di(t)\bm{d}_{i}^{(t)} of this update rule take different forms for different algorithms. Below, we provide detailed descriptions of these key elements.

Aggregation weights wiw_{i}: Each client’s normalized gradient di\bm{d}_{i} is multiplied with weight wiw_{i} when computing the aggregated gradient ∑i=1mwidi\sum_{i=1}^{m}w_{i}\bm{d}_{i}. By definition, these weights satisfy ∑i=1mwi=1\sum_{i=1}^{m}w_{i}=1. Observe that these weights determine the surrogate objective F~(x)=∑i=1mwiFi(x)\smash{\widetilde{F}}(\bm{x})=\sum_{i=1}^{m}w_{i}F_{i}(\bm{x}), which is optimized by the general algorithm in (4) instead of the true global objective F(x)=∑i=1mpiFi(x)F(\bm{x})=\sum_{i=1}^{m}p_{i}F_{i}(\bm{x}) – we will prove this formally in Theorem 1.

Effective number of steps τeff\tau_{\text{eff}}: Since client ii makes τi\tau_{i} local updates, the average number of local SGD steps per communication round is τˉ=∑i=1mτi/m\bar{\tau}=\sum_{i=1}^{m}\tau_{i}/m. However, the server can scale up or scale down the effect of the aggregated updates by setting the parameter τeff\tau_{\text{eff}} larger or smaller than τˉ\bar{\tau} (analogous to choosing a global learning rate ). We refer to the ratio τˉ/τeff\bar{\tau}/\tau_{\text{eff}} as the slowdown, and it features prominently in the convergence analysis presented in Section 4.2.

The general rule (4) enables us to freely choose τeff\tau_{\text{eff}} and wiw_{i} for a given local solver ai\bm{a}_{i}, which helps design fast and consistent algorithms such as FedNova, the normalized averaging method proposed in Section 5. In Figure 3, we further illustrate how the above key elements influence the algorithm and compare the novel generalized update rule and FedAvg in the model parameter space. Besides, in terms of the implementation of the generalized update rule, each client can send the normalized update −ηdi(t)-\eta\smash{\bm{d}_{i}^{(t)}} to the central server, which is just a re-scaled version of Δi(t)\smash{\Delta_{i}^{(t)}}, the accumulated local parameter update sent by clients in the vanilla update rule (2). The server is not necessary to know the specific form of local accumulation vector ai\bm{a}_{i}.

Previous Algorithms as Special Cases. Any previous algorithm whose accumulated local changes Δi(t)=−ηGi(t)ai\smash{\Delta_{i}^{(t)}=-\eta\bm{G}_{i}^{(t)}\bm{a}_{i}}, a linear combination of local gradients is subsumed by the above formulation. One can validate this as follows:

Unlike the more general form 4, in (6), which subsumes the following previous methods, τeff\tau_{\text{eff}} and wiw_{i} are implicitly fixed by the choice of the local solver (i.e., the choice of ai\bm{a}_{i}). Due to space limitations, the derivations of following examples are relegated to Appendix B.

When α=0\alpha=0, FedProx is equivalent to FedAvg. As α=ημ\alpha=\eta\mu increases, the wiw_{i} in FedProx is more similar to pip_{i}, thus making the surrogate objective F~(x)\widetilde{F}(\bm{x}) more consistent. However, a larger α\alpha corresponds to smaller τeff\tau_{\text{eff}}, which slows down convergence, as we discuss more in the next subsection.

SGD with Decayed Learning Rate as Local Solver. Suppose the clients’ local learning rates are exponentially decayed, then we have ai=[1,γi,…,γiτi−1]\bm{a}_{i}=[1,\gamma_{i},\dots,\gamma_{i}^{\tau_{i}-1}] where γi≥0\gamma_{i}\geq 0 can vary across clients. As a result, we have ∥ai∥1=(1−γiτi)/(1−γi)\|\bm{a}_{i}\|_{1}=(1-\gamma_{i}^{\tau_{i}})/(1-\gamma_{i}) and wi∝pi(1−γiτi)/(1−γi)w_{i}\propto p_{i}(1-\gamma_{i}^{\tau_{i}})/(1-\gamma_{i}). Comparing with the case of FedProx 7, changing the values of γi\gamma_{i} has a similar effect as changing (1−α)(1-\alpha).

Momentum SGD as Local Solver. If we use momentum SGD where the local momentum buffers of active clients are reset to zero at the beginning of each round due to the stateless nature of FL , then we have ai=[1−ρτi,1−ρτi−1,…,1−ρ]/(1−ρ)\bm{a}_{i}=[1-\rho^{\tau_{i}},1-\rho^{\tau_{i}-1},\dots,1-\rho]/(1-\rho), where ρ\rho is the momentum factor, and ∥ai∥1=[τi−ρ(1−ρτi)/(1−ρ)]/(1−ρ)\|\bm{a}_{i}\|_{1}=[\tau_{i}-\rho(1-\rho^{\tau_{i}})/(1-\rho)]/(1-\rho).

More generally, the new formulation 6 suggests that wi≠piw_{i}\neq p_{i} whenever clients have different ∥ai∥1\|\bm{a}_{i}\|_{1}, which may be caused by imbalanced local updates (i.e., ai\bm{a}_{i}’s have different dimensions), or various local learning rate/momentum schedules (i.e., ai\bm{a}_{i}’s have different scales).

2 Convergence Analysis for Smooth Non-Convex Functions

In Theorem 1 and Theorem 2 below we provide a convergence analysis for the general update rule 4 and quantify the solution bias due to objective inconsistency. The analysis relies on Assumptions 1 and 2 used in the standard analysis of SGD and Assumption 3 commonly used in the federated optimization literature to capture the dissimilarities of local objectives.

Each local objective function is Lipschitz smooth, that is, ∥∇Fi(x)−∇Fi(y)∥≤L∥x−y∥,∀i∈{1,2,…,m}\left\|\nabla F_{i}(\bm{x})-\nabla F_{i}(\bm{y})\right\|\leq L\left\|\bm{x}-\bm{y}\right\|,\forall i\in\{1,2,\dots,m\}.

For any sets of weights {wi≥0}i=1m,∑i=1mwi=1\{w_{i}\geq 0\}_{i=1}^{m},\sum_{i=1}^{m}w_{i}=1, there exist constants β2≥1,κ2≥0\beta^{2}\geq 1,\kappa^{2}\geq 0 such that ∑i=1mwi∥∇Fi(x)∥2≤β2∥∑i=1mwi∇Fi(x)∥2+κ2\sum_{i=1}^{m}w_{i}\left\|\nabla F_{i}(\bm{x})\right\|^{2}\leq\beta^{2}\left\|\sum_{i=1}^{m}w_{i}\nabla F_{i}(\bm{x})\right\|^{2}+\kappa^{2}. If local functions are identical to each other, then we have β2=1,κ2=0\beta^{2}=1,\kappa^{2}=0.

Under Assumptions 1, 2 and 3, any federated optimization algorithm that follows the update rule 4, will converge to a stationary point of a surrogate objective F~(x)=∑i=1mwiFi(x)\smash{\widetilde{F}}(\bm{x})=\sum_{i=1}^{m}w_{i}F_{i}(\bm{x}). More specifically, if the total communication rounds TT is pre-determined and the learning rate η\eta is small enough η=m/τ‾T\eta=\sqrt{m/\overline{\tau}T} where τ‾=1m∑i=1mτi\overline{\tau}=\frac{1}{m}\sum_{i=1}^{m}\tau_{i}, then the optimization error will be bounded as follows:

where O\mathcal{O} swallows all constants (including LL), and quantities A,B,CA,B,C are defined as follows:

where ai,−1a_{i,-1} is the last element in the vector ai\bm{a}_{i}.

In Appendix C, we also provide another version of this theorem that explicitly contains the local learning rate η\eta. Moreover, since the surrogate objective F~(x)\smash{\widetilde{F}}(\bm{x}) and the original objective F(x)F(\bm{x}) are just different linear combinations of the local functions, once the algorithm converges to a stationary point of F~(x)\smash{\widetilde{F}}(\bm{x}), one can also obtain some guarantees in terms of F(x)F(\bm{x}), as given by Theorem 2 below.

Under the same conditions as Theorem 1, the minimal gradient norm of the true global objective function F(x)=∑i=1mpiFi(x)F(\bm{x})=\sum_{i=1}^{m}p_{i}F_{i}(\bm{x}) will be bounded as follows:

where ϵopt\epsilon_{\text{opt}} denotes the vanishing optimization error given by 8 and χp∥w2=∑i=1m(pi−wi)2/wi\smash{\chi^{2}_{\bm{p}\|\bm{w}}}=\sum_{i=1}^{m}(p_{i}-w_{i})^{2}/w_{i} represents the chi-square divergence between vectors p=[p1,…,pm]\bm{p}=[p_{1},\dots,p_{m}] and w=[w1,…,wm]\bm{w}=[w_{1},\dots,w_{m}].

Discussion: Theorems 1 and 2 describe the convergence behavior of a broad class of federated heterogeneous optimization algorithms. Observe that when all clients take the same number of local steps using the same local solver, we have p=w\bm{p}=\bm{w} such that χ2=0\chi^{2}=0. Also, when all local functions are identical to each other, we have β2=1,κ2=0\beta^{2}=1,\kappa^{2}=0. Only in these two special cases, is there no objective inconsistency. For most other algorithms subsumed by the general update rule in (4), both wiw_{i} and τeff\tau_{\text{eff}} are influenced by the choice of ai\bm{a}_{i}. When clients have different local progress (i.e., different ai\bm{a}_{i} vectors), previous algorithms will end up with a non-zero error floor χ2κ2\chi^{2}\kappa^{2}, which does not vanish to even with sufficiently small learning rate. In Section D.1, we further construct a lower bound and show that lim⁡T→∞min⁡t∈[T]∥∇F(x(t,0))∥2=Ω(χ2κ2)\lim_{T\rightarrow\infty}\min_{t\in[T]}\|\nabla F(\bm{x}^{(t,0)})\|^{2}=\Omega(\chi^{2}\kappa^{2}), suggesting 10 is tight.

Novel Insights Into the Convergence of FedProx and the Effect of μ\mu. Recall that in FedProx ai=[(1−α)τi−1,…,(1−α),1]\bm{a}_{i}=[(1-\alpha)^{\tau_{i}-1},\dots,(1-\alpha),1], where α=ημ\alpha=\eta\mu. Accordingly, substituting the effective steps and aggregated weight, given by 7, into 8 and 10, we get the convergence guarantee for FedProx. Again, it has objective inconsistency because wi≠piw_{i}\neq p_{i}. As we increase α\alpha, the weights wiw_{i} come closer to pip_{i} and thus, the non-vanishing error χ2κ2\chi^{2}\kappa^{2} in 10 decreases (see blue curve in Figure 4). However increasing α\alpha worsens the slowdown τ‾/τeff\overline{\tau}/\tau_{\text{eff}}, which appears in the first error term in 8 (see the red curve in Figure 4). In the extreme case when α=1\alpha=1, although FedProx achieves objective consistency, it has a significantly slower convergence because τeff=1\tau_{\text{eff}}=1 and the first term in 8 is τ‾\overline{\tau} times larger than that with FedAvg (eq. to α=0\alpha=0).

Theorem 1 also reveals that, in FedProx, there should exist a best value of α\alpha that balances all terms in 8. In Appendix, we provide a corollary showing that α=O(\nicefracm12τ‾12T16)\alpha=\mathcal{O}(\nicefrac{{m^{\frac{1}{2}}}}{{\overline{\tau}^{\frac{1}{2}}T^{\frac{1}{6}}}}) optimizes the error bound 8 of FedProx and yields a convergence rate of O(\nicefrac1mτ‾T+\nicefrac1T23)\mathcal{O}(\nicefrac{{1}}{{\sqrt{m\overline{\tau}T}}}+\nicefrac{{1}}{{T^{\frac{2}{3}}}}) on the surrogate objective. This can serve as a guideline on setting α\alpha in practice.

Linear Speedup Analysis. Another implication of Theorem 1 is that when the communication rounds TT is sufficiently large, then the convergence of the surrogate objective will be dominated by the first two terms in 8, which is \nicefrac1mτ‾T\nicefrac{{1}}{{\sqrt{m\overline{\tau}T}}}. This suggests that the algorithm only uses T/γT/\gamma total rounds when using γ\gamma times more clients (i.e., achieving linear speedup) to reach the same error level.

FedNova: Proposed Federated Normalized Averaging Algorithm

Theorems 1 and 2 suggest an extremely simple solution to overcome the problem of objective inconsistency. When we set wi=piw_{i}=p_{i} in 4, then the second non-vanishing term χp∥w2κ2\smash{\chi^{2}_{\bm{p}\|\bm{w}}}\kappa^{2} in 10 will just become zero. This simple intuition yields the following new algorithm:

Flexibility in Choosing Hyper-parameters and Local Solvers. Besides vanilla SGD, the new formulation of FedNova naturally allows clients to choose various local solvers (i.e., client-side optimizer). As discussed in Section 4.1, the local solver can also be GD/SGD with decayed local learning rate, GD/SGD with proximal updates, GD/SGD with local momentum, etc. Furthermore, the value of τeff\tau_{\text{eff}} is not necessarily to be controlled by the local solver as previous algorithms. For example, when using SGD with proximal updates, one can simply set τeff=∑i=1mpiτi\tau_{\text{eff}}=\sum_{i=1}^{m}p_{i}\tau_{i} instead of its default value ∑i=1mpi[1−(1−α)τi]/α\sum_{i=1}^{m}p_{i}[1-(1-\alpha)^{\tau_{i}}]/\alpha. This can help alleviate the slowdown problem discussed in Section 4.2.

Combination with Acceleration Techniques. If clients have additional communication bandwidth, they can use cross-client variance reduction techniques to further accelerate the training . In this case, each local gradient step at the tt-round will be corrected by ∑i=1mpidi(t−1)−di(t−1)\sum_{i=1}^{m}p_{i}\bm{d}_{i}^{(t-1)}-\bm{d}_{i}^{(t-1)}. That is, the local gradient at the kk-th local step becomes gi(x(t,k))+∑i=1mpidi(t−1)−di(t−1)g_{i}(\bm{x}^{(t,k)})+\sum_{i=1}^{m}p_{i}\bm{d}_{i}^{(t-1)}-\bm{d}_{i}^{(t-1)}. Besides, on the server side, one can also implement server momentum or adaptive server optimizers , in which the aggregated normalized gradient −τeff∑i=1mηpidi-\tau_{\text{eff}}\sum_{i=1}^{m}\eta p_{i}\bm{d}_{i} is used to update the server momentum buffer instead of directly updating the server model.

Convergence Analysis. In FedNova, the local solvers at clients do not necessarily need to be the same or fixed across rounds. In the following theorem, we obtain strong convergence guarantee for FedNova, even with arbitrarily time-varying local updates and client optimizers.

Suppose that each client performs arbitrary number of local updates τi(t)\tau_{i}(t) using arbitrary gradient accumulation method ai(t),t∈[T]\bm{a}_{i}(t),t\in[T] per round. Under Assumptions 1, 2 and 3, and local learning rate as η=m/(τ~T)\eta=\sqrt{m/(\widetilde{\tau}T)}, where τ~=∑t=0T−1τ‾(t)/T\widetilde{\tau}=\sum_{t=0}^{T-1}\overline{\tau}(t)/T denotes the average local steps over all rounds at clients, then FedNova converges to a stationary point of F(x)F(\bm{x}) in a rate of O(1/mτ~T)\mathcal{O}(1/\sqrt{m\widetilde{\tau}T}). The detailed bound is the same as the right hand side of 8, except that τ‾,A,B,C\overline{\tau},A,B,C are replaced by their average values over all rounds.

Using the techniques developed in , Theorem 3 can be further generalized to incorporate client sampling schemes. We provide a corresponding corollary in Appendix G. When qq clients are selected per round, then the convergence rate of FedNova is O(1/qτ~T)\mathcal{O}(1/\sqrt{q\widetilde{\tau}T}), where τ~=∑t=0T−1∑i∈S(t)τi(t)/(qT)\widetilde{\tau}=\sum_{t=0}^{T-1}\sum_{i\in\mathcal{S}^{(t)}}\tau_{i}^{(t)}/(qT) and S(t)\mathcal{S}^{(t)} is the set of selected client indices at the tt-th round.

Experimental Results

Experimental Setup. We evaluate all algorithms on two setups with non-IID data partitioning: (1) Logistic Regression on a Synthetic Federated Dataset: The dataset Synthetic(1,1)(1,1) is originally constructed in . The local dataset sizes ni,i∈n_{i},i\in follows a power law. (2) DNN trained on a Non-IID partitioned CIFAR-10 dataset: We train a VGG-11 network on the CIFAR-10 dataset , which is partitioned across 1616 clients using a Dirichlet distribution Dir16(0.1)\text{Dir}_{16}(0.1), as done in . The original CIFAR-10 test set (without partitioning) is used to evaluate the generalization performance of the trained global model. The local learning rate η\eta is decayed by a constant factor after finishing 50%50\% and 75%75\% of the communication rounds. The initial value of η\eta is tuned separately for FedAvg with different local solvers. When using the same solver, FedNova uses the same η\eta as FedAvg to guarantee a fair comparison. On CIFAR-10, we run each experiment with 33 random seeds and report the average and standard deviation. More details are provided in Appendix I.

Synthetic Dataset Simulations. In Figure 5, we observe that by simply changing wiw_{i} to pip_{i}, FedNova not only converges significantly faster than FedAvg but also achieves consistently the best performance under three different settings. Note that the only difference between FedNova and FedAvg is the aggregated weights when averaging the normalized gradients.

Non-IID CIFAR-10 Experiments. In Table 1 we compare the performance of FedNova and FedAvg on non-IID CIFAR-10 with various client optimizers run for 100100 communication rounds. When the client optimizer is SGD or SGD with momentum, simply changing the weights yields a 66-9%9\% improvement on the test accuracy; When the client optimizer is proximal SGD, FedAvg is equivalent to FedProx. By setting τeff=∑i=1mpiτi\tau_{\text{eff}}=\sum_{i=1}^{m}p_{i}\tau_{i} and correcting the weights wi=piw_{i}=p_{i} while keeping ai\bm{a}_{i} same as FedProx, FedNova-Prox achieves about 10%10\% higher test accuracy than FedProx. In Figure 6, we further compare the training curves. It turns out that FedNova consistently converges faster than FedAvg. When using variance-reduction methods such as SCAFFOLD (that requires doubled communication), FedNova-based method preserves the same test accuracy. Furthermore, combining local momentum and variance-reduction can be easily achieved in FedNova. It yields the highest test accuracy among all other local solvers. This kind of combination is non-trivial and has not appeared yet in the literature. We provide its pseudocode in Appendix H.

Effectiveness of Local Momentum. From Table 1, it is worth noting that using momentum SGD as the local solver is an effective way to improve the performance. It generally achieves 33-7%7\% higher test accuracy than vanilla SGD. This local momentum scheme can be further combined with server momentum . When Ei(t)∼U(2,5)E_{i}(t)\sim\mathcal{U}(2,5), the hybrid momentum scheme achieves test accuracy 81.15±0.38%81.15\pm 0.38\% As a reference, using server momentum alone achieves 77.49±0.25%77.49\pm 0.25\%.

Concluding Remarks

In federated learning, the participated clients (e.g., IoT sensors, mobile devices) are typically highly heterogeneous, both in the size of their local datasets as well as their computation speeds. Clients can also join and leave the training at any time according to their availabilities. Therefore, it is common that clients perform different amount of works within one round of local computation. However, previous analyses on federated optimization algorithms are limited to the homogeneous case where all clients have the same local steps, hyper-parameters, and client optimizers. In this paper, we develop a novel theoretical framework to analyze the challenging heterogeneous setting. We show that original FedAvg algorithm will converge to stationary points of a mismatched objective function which can be arbitrarily different from the true objective. To the best of our knowledge, we provide the first fundamental understanding of how the convergence rate and bias in the final solution of federated optimization algorithms are influenced by the heterogeneity in clients’ local progress. The new framework naturally allows clients to have different local steps and local solvers, such as GD, SGD, SGD with momentum, proximal updates, etc. Inspired by the theoretical analysis, we propose FedNova, which can automatically adjust the aggregated weight and effective local steps according to the local progress. We validate the effectiveness of FedNova both theoretically and empirically. On a non-IID version of CIFAR-10 dataset, FedNova generally achieves 66-9%9\% higher test accuracy than FedAvg. Future directions include extending the theoretical framework to adaptive optimization methods or gossip-based training methods.

Acknowledgements

This research was generously supported in part by NSF grants CCF-1850029, the 2018 IBM Faculty Research Award, and the Qualcomm Innovation fellowship (Jianyu Wang). We thank Anit Kumar Sahu, Tian Li, Zachary Charles, Zachary Garrett, and Virginia Smith for helpful discussions.

References

Appendix A Proof of Lemma 1: Objective Inconsistency in Quadratic Model

Consider a simple setting where each local objective function is strongly convex and defined as follows:

where H‾=∑i=1mpiHi\overline{\bm{H}}=\sum_{i=1}^{m}p_{i}\bm{H}_{i} and e‾=∑i=1mpiei\overline{\bm{e}}=\sum_{i=1}^{m}p_{i}\bm{e}_{i}. As a result, the global minimum is x∗=H‾−1e‾\bm{x}^{*}=\overline{\bm{H}}^{-1}\overline{\bm{e}}. Now, let us study whether previous federated optimization algorithms can converge to this global minimum.

Local Update Rule.

The local update rule of FedProx for the ii-th device can be written as follows:

where xi(t,k)\bm{x}_{i}^{(t,k)} denotes the local model parameters at the kk-th local iteration after tt communication rounds, η\eta denotes the local learning rate and μ\mu is a tunable hyper-parameter in FedProx. When μ=0\mu=0, the algorithm will reduce to FedAvg. We omit the device index in x(t,0)\bm{x}^{(t,0)}, since it is synchronized and the same across all devices.

where ci(t)=(Hi+μI)−1(ei+μx(t,0))\bm{c}_{i}^{(t)}=\left(\bm{H}_{i}+\mu\bm{I}\right)^{-1}\left(\bm{e}_{i}+\mu\bm{x}^{(t,0)}\right). Then, after performing τi\tau_{i} steps of local updates, the local model becomes

For the ease of writing, we define Ki(η,μ)=[I−(I−ημI−ηHi)τi](Hi+μI)−1\bm{K}_{i}(\eta,\mu)=\left[\bm{I}-\left(\bm{I}-\eta\mu\bm{I}-\eta\bm{H}_{i}\right)^{\tau_{i}}\right]\left(\bm{H}_{i}+\mu\bm{I}\right)^{-1}.

Server Aggregation.

For simplicity, we only consider the case when all devices participate in the each round. In FedProx, the server averages all local models according to the sample size:

Accordingly, we get the following update rule for the central model:

After TT communication rounds, one can get

Accordingly, when ∥I−∑i=1mpiKi(η,μ)Hi∥2<1\left\|\bm{I}-\sum_{i=1}^{m}p_{i}\bm{K}_{i}(\eta,\mu)\bm{H}_{i}\right\|_{2}<1, the iterates will converge to

Recall that Ki(η,μ)=[I−(I−ημI−ηHi)τi](Hi+μI)−1\bm{K}_{i}(\eta,\mu)=\left[\bm{I}-\left(\bm{I}-\eta\mu\bm{I}-\eta\bm{H}_{i}\right)^{\tau_{i}}\right]\left(\bm{H}_{i}+\mu\bm{I}\right)^{-1}.

Concrete Example in Lemma 1.

Now let us focus on a concrete example where p1=p2=⋯=pm=1/m,H1=H2=⋯=Hm=Ip_{1}=p_{2}=\cdots=p_{m}=1/m,\bm{H}_{1}=\bm{H}_{2}=\cdots=\bm{H}_{m}=\bm{I} and μ=0\mu=0. Then, in this case, Ki=1−(1−η)τi\bm{K}_{i}=1-(1-\eta)^{\tau_{i}}. As a result, we have

Furthermore, when the learning rate is sufficiently small (e.g., can be achieved by gradually decaying the learning rate), according to L’Hospital’s rule, we obtain

Appendix B Detailed Derivations for Various Local Solvers

In this section, we will derive the specific expression of the vector ai\bm{a}_{i} when using different local solvers. Recall that the local change at client ii is Δi(t)=−ηGi(t)ai\Delta_{i}^{(t)}=-\eta\bm{G}_{i}^{(t)}\bm{a}_{i} where Gi(t)\bm{G}_{i}^{(t)} stacks all stochastic gradients in the current round and a\bm{a} is a non-negative vector.

In this case, we can write the update rule of local models as follows:

Subtracting xi(t,0)\bm{x}_{i}^{(t,0)} on both sides, we obtain

Repeating the above procedure, it follows that

According to the definition, we have ai=[(1−α)τi−1,(1−α)τi−2,…,(1−α),1]\bm{a}_{i}=[(1-\alpha)^{\tau_{i}-1},(1-\alpha)^{\tau_{i}-2},\dots,(1-\alpha),1] where α=ημ\alpha=\eta\mu.

B.2 SGD with Local Momentum

Let us firstly write down the update rule of the local models. Suppose that ρ\rho denotes the local momentum factor and ui\bm{u}_{i} is the local momentum buffer at client ii. Then, the update rule of local momentum SGD is:

One can expand the expression of local momentum buffer as follows:

where the last equation comes from the fact ui(t,0)=0\bm{u}_{i}^{(t,0)}=0. Substituting 39 into 36, we have

Repeating the above procedure, it follows that

Then, the coefficient of gi(xi(t,k))g_{i}(\bm{x}_{i}^{(t,k)}) is

Appendix C Proof of Theorem 1: Convergence of Surrogate Objective

For the ease of writing, let us define a surrogate objective function F~(x)=∑i=1mwiFi(x)\smash{\widetilde{F}}(\bm{x})=\sum_{i=1}^{m}w_{i}F_{i}(\bm{x}), where ∑i=1mwi=1\sum_{i=1}^{m}w_{i}=1, and define the following auxiliary variables

According to the Lipschitz-smooth assumption, it follows that

where the expectation is taken over mini-batches ξi(t,k),∀i∈{1,2,…,m},k∈{0,1,…,τi−1}\xi_{i}^{(t,k)},\forall i\in\{1,2,\dots,m\},k\in\{0,1,\dots,\tau_{i}-1\}. Before diving into the detailed bounds for T1T_{1} and T2T_{2}, we would like to firstly introduce several useful lemmas.

Assume i<ji<j. Then, using the law of total expectation,

C.2 Bounding First term in 49

For the first term on the right hand side (RHS) in 49, we have

where the last equation uses the fact: 2<a, b>=∥a∥2+∥b∥2−∥a−b∥22\left<{a},\,{b}\right>=\left\|a\right\|^{2}+\left\|b\right\|^{2}-\left\|a-b\right\|^{2}.

C.3 Bounding Second term in 49

For the second term on the right hand side (RHS) in 49, we have

where 60 is derived using Lemma 2, 61 follows Assumption 2.

C.4 Intermediate Result

When τeffηL≤1/2\tau_{\text{eff}}\eta L\leq 1/2, it follows that

where the last inequality uses the fact F~(x)=∑i=1mwiFi(x)\smash{\widetilde{F}}(\bm{x})=\sum_{i=1}^{m}w_{i}F_{i}(\bm{x}) and Jensen’s Inequality: ∥∑i=1mwizi∥2≤∑i=1mwi∥zi∥2\left\|\sum_{i=1}^{m}w_{i}z_{i}\right\|^{2}\leq\sum_{i=1}^{m}w_{i}\left\|z_{i}\right\|^{2}. Next, we will focus on bounding the last term in 64.

C.5 Bounding the Difference Between Server Gradient and Normalized Gradient

Recall the definition of hi(t)\bm{h}_{i}^{(t)}, one can derive that

where 67 uses Jensen’s Inequality again: ∥∑i=1mwizi∥2≤∑i=1mwi∥zi∥2\left\|\sum_{i=1}^{m}w_{i}z_{i}\right\|^{2}\leq\sum_{i=1}^{m}w_{i}\left\|z_{i}\right\|^{2}, and 68 follows Assumption 1. Now, we turn to bounding the difference between the server model x(t,0)\bm{x}^{(t,0)} and the local model xi(t,k)\bm{x}_{i}^{(t,k)}. Plugging into the local update rule and using the fact ∥a+b∥2≤2∥a∥2+2∥b∥2\left\|a+b\right\|^{2}\leq 2\left\|a\right\|^{2}+2\left\|b\right\|^{2},

where 74 follows from Jensen’s Inequality. Furthermore, note that

where ai,−1a_{i,-1} is the last element in the vector ai\bm{a}_{i}. As a result, we have

In addition, we can bound the second term using the following inequality:

Define D=4η2L2max⁡i{∥ai∥1(∥ai∥1−ai,−1)}<1D=4\eta^{2}L^{2}\max_{i}\{\left\|\bm{a}_{i}\right\|_{1}(\left\|\bm{a}_{i}\right\|_{1}-a_{i,-1})\}<1. We can simplify 84 as follows

Taking the average across all workers and applying Assumption 3, one can obtain

Now, we are ready to derive the final result.

C.6 Final Results

If D≤12β2+1D\leq\frac{1}{2\beta^{2}+1}, then it follows that 11−D≤1+12β2\frac{1}{1-D}\leq 1+\frac{1}{2\beta^{2}} and Dβ21−D≤12\frac{D\beta^{2}}{1-D}\leq\frac{1}{2}. These facts can help us further simplify inequality 89.

Taking the average across all rounds, we get

For the ease of writing, we define the following auxiliary variables:

C.7 Constraint on Local Learning Rate

Here, let us summarize the constraints on local learning rate:

For the second constraint, we can further tighten it as follows:

C.8 Further Optimizing the Bound

By setting η=mτ‾T\eta=\sqrt{\frac{m}{\overline{\tau}T}} where τ‾=1m∑i=1mτi\overline{\tau}=\frac{1}{m}\sum_{i=1}^{m}\tau_{i}, we have

Here, we complete the proof of Theorem 1.

Appendix D Proof of Theorem 2: Including Bias in the Error Bound

For any model parameter x\bm{x}, the difference between the gradients of F(x)F(\bm{x}) and F~(x)\smash{\widetilde{F}}(\bm{x}) can be bounded as follows:

where χp∥w2\chi^{2}_{\bm{p}\|\bm{w}} denotes the chi-square distance between p\bm{p} and w\bm{w}, i.e., χp∥w2=∑i=1m(pi−wi)2/wi\chi^{2}_{\bm{p}\|\bm{w}}=\sum_{i=1}^{m}(p_{i}-w_{i})^{2}/w_{i}.

According to the definition of F(x)F(x) and F~(x)\smash{\widetilde{F}}(\bm{x}), we have

Applying Cauchy–Schwarz inequality, it follows that

where the last inequality uses Assumption 3. ∎

where ϵopt\epsilon_{\text{opt}} denotes the optimization error.

One can manually construct a strongly convex objective function such that FedAvg with heterogeneous local updates cannot converge to its global optimum. In particular, the gradient norm of the objective function does not vanish as learning rate approaches to zero. We have the following lower bound:

where χp∥w2\chi^{2}_{\bm{p}\|\bm{w}} denotes the chi-square divergence between weight vectors and κ2\kappa^{2} quantifies the dissimilarities among local objective functions and is defined in Assumption 3.

Suppose that there are only two clients with local objectives F1(x)=12(x−a)2F_{1}(x)=\frac{1}{2}(x-a)^{2} and F2(x)=12(w+a)2F_{2}(x)=\frac{1}{2}(w+a)^{2}. The global objective is defined as F(x)=12F1(x)+12F2(x)F(x)=\frac{1}{2}F_{1}(x)+\frac{1}{2}F_{2}(x). For any set of weights w1,w2,w1+w2=1w_{1},w_{2},w_{1}+w_{2}=1, we define the surrogate objective function as F~(x)=w1F1(x)+w2F2(x)\widetilde{F}(\bm{x})=w_{1}F_{1}(\bm{x})+w_{2}F_{2}(\bm{x}). As a consequence, we have

Comparing with Assumption 3, we can define κ2=2w1w2a2\kappa^{2}=2w_{1}w_{2}a^{2} and β2=1\beta^{2}=1 in this case. Furthermore, according to the derivations in Appendix A, the iterate of FedAvg can be written as follows:

where χp∥w2=∑i=1m(pi−wi)2/wi=(w1−1/2)2/w1+(w2−1/2)2/w2=(τ2−τ1)2/(2τ1τ2)\chi^{2}_{\bm{p}\|\bm{w}}=\sum_{i=1}^{m}(p_{i}-w_{i})^{2}/w_{i}=(w_{1}-1/2)^{2}/w_{1}+(w_{2}-1/2)^{2}/w_{2}=(\tau_{2}-\tau_{1})^{2}/(2\tau_{1}\tau_{2}). ∎

Appendix E Special Cases of Theorem 1

Here, we provide several instantiations of Theorem 1 and check its consistency with previous results.

In the case where all clients have the same local dataset size, i.e., pi=1/m,∀ip_{i}=1/m,\forall i. It follows that

Substituting 125 into Theorem 1, we get the convergence guarantee for FedAvg. We formally state it in the following corollary.

Under the same conditions as Theorem 1, if pi=1/mp_{i}=1/m, then FedAvg algorithm (vanilla SGD with fixed local learning rate as local solver) will converge to the stationary point of a surrogate objective F~(x)=∑i=1mτiFi(x)/∑i=1mτi\widetilde{F}(\bm{x})=\sum_{i=1}^{m}\tau_{i}F_{i}(\bm{x})/\sum_{i=1}^{m}\tau_{i}. The optimization error will be bounded as follows:

where O\mathcal{O} swallows all constants (including LL), and var⁡[τ]=∑i=1mτi2/m−τ‾2\operatorname{var}[\bm{\tau}]=\sum_{i=1}^{m}\tau_{i}^{2}/m-\overline{\tau}^{2} denotes the variance of local steps.

Consistent with Previous Results. When all clients perform the same local steps, i.e.,τi=τ\textit{i.e.},\tau_{i}=\tau, then var⁡[τ]=0\operatorname{var}[\bm{\tau}]=0 and the above error bound 126 recovers previous results . When τi=1\tau_{i}=1, then FedAvg reduces to fully synchronous SGD and the error bound 126 becomes 1/mT1/\sqrt{mT}, which is the same as standard SGD convergence rate .

E.2 FedProx

As a consequence, we can derive the closed-form expression of τeff,A,B,C\tau_{\text{eff}},A,B,C as follows:

Substituting AFedProx,BFedProx,CFedProxA_{\texttt{FedProx}},B_{\texttt{FedProx}},C_{\texttt{FedProx}} back into Theorem 1, one can obtain the convergence guarantee for FedProx. Again, it will converge to the stationary points of a surrogate objective due to wi≠piw_{i}\neq p_{i}.

Consistency with FedAvg. From the update rule of FedProx, we know that when μ=0\mu=0 (or α=0\alpha=0), FedProx is equivalent to FedProx. This can also be validated from the expressions of AFedProx,BFedProx,CFedProxA_{\texttt{FedProx}},B_{\texttt{FedProx}},C_{\texttt{FedProx}}. Using L’Hospital law, it is easy to show that

Best value of α\alpha in FedProx. Given the expressions of τeff\tau_{\text{eff}} and A,B,CA,B,C, we can further select a best value of α\alpha that optimizes the error bound of FedProx, as stated in the following corollary.

Under the same conditions as Theorem 1 and suppose pi=1/mp_{i}=1/m and τi≫1\tau_{i}\gg 1, then α=O(\nicefracm12τ‾12T16))\alpha=\mathcal{O}(\nicefrac{{m^{\frac{1}{2}}}}{{\overline{\tau}^{\frac{1}{2}}T^{\frac{1}{6}}}})) minimizes the optimization error bound of FedProx in terms of converging to the stationary points of the surrogate objective. In particular, we have

Discussion: Corollary 2 shows that there exists a non-zero value of α\alpha that optimizes the error upper bound of FedProx. That is to say, FedProx (α>0\alpha>0) is better than FedAvg (α=0\alpha=0) by a constant in terms of error upper bound. However, on the other hand, it is worth noting that the minimal communication rounds of FedProx to achieve 1/mK1/\sqrt{mK} rate, given by Corollary 2, is exactly the same as FedAvg . In this sense, FedProx has the same convergence rate as FedAvg and cannot further reduce the communication overhead.

First of all, let us relax the error terms of FedProx. Under the assumption of τi≫1\tau_{i}\gg 1, the quantities A,B,CA,B,C can be bounded or approximated as follows:

Accordingly, the error upper bound of FedProx can be rewritten as follows:

In order to optimize the above bound, we can simply take the derivative with respect to α\alpha. When the derivative equals to zero, we get

Plugging the expression of best α\alpha into 138, we have

where K=τ‾TK=\overline{\tau}T denotes the average total gradient steps at clients. In order to let the first term dominates the convergence rate, it requires that

As a results, the total communication rounds T=K/τ‾T=K/\overline{\tau} should be greater than O(K34m34)\mathcal{O}(K^{\frac{3}{4}}m^{\frac{3}{4}}). ∎

Appendix F Proof of Theorem 3

In the case of FedNova, the aggregated weights wiw_{i} equals to pip_{i}. Therefore, the surrogate objective F~(x)=∑i=1mwiFi(x)\smash{\widetilde{F}}(\bm{x})=\sum_{i=1}^{m}w_{i}F_{i}(\bm{x}) is the same as the original objective function F(x)=∑i=1mpiFi(x)F(\bm{x})=\sum_{i=1}^{m}p_{i}F_{i}(\bm{x}). We can directly reuse the intermediate results in the proof of Theorem 1. According to 91, we have

where quantities A(t),B(t),C(t)A^{(t)},B^{(t)},C^{(t)} are defined as follows:

Taking the total expectation and averaging over all rounds, it follows that

where A~=∑t=0T−1A(t)/T,B~=∑t=0T−1B(t)/T\widetilde{A}=\sum_{t=0}^{T-1}A^{(t)}/T,\widetilde{B}=\sum_{t=0}^{T-1}B^{(t)}/T, and C~=∑t=0T−1C(t)/T\widetilde{C}=\sum_{t=0}^{T-1}C^{(t)}/T. After minor rearranging, we have

Bt setting η=mτ~T\eta=\sqrt{\frac{m}{\widetilde{\tau}T}} where τ~=∑t=0T−1τ‾(t)/T\widetilde{\tau}=\sum_{t=0}^{T-1}\overline{\tau}^{(t)}/T, the above upper bound can be further optimized as follows:

Here, we complete the proof of Theorem 3.

Moreover, it is worth mentioning the constraints on the local learning rate. Recall that, at the tt-th round, we have the following constraint:

In order to guarantee the convergence, the above inequality should hold in every round. That is to say,

Appendix G Extension: Incorporating Client Sampling

In this section, we extend the convergence guarantee of FedNova to the case of client sampling. Following previous works , we assume the sampling scheme guarantees that the update rule 12 hold in expectation. This can be achieved by sampling with replacement from {1,2,…,m}\{1,2,\dots,m\} with probabilities {pi}\{p_{i}\}, and averaging local updates from selected clients with equal weights. Specifically, we have

Under the same condition as Theorem 1, suppose at each round, the server randomly selects q(≤m)q(\leq m) clients with replacement to perform local computation. The probability of choosing the ii-th client is pi=ni/np_{i}=n_{i}/n. In this case, FedNova will converge to the stationary points of the global objective F(x)F(\bm{x}). If we set η=q/τ~T\eta=\sqrt{q/\widetilde{\tau}T} where τ~\widetilde{\tau} is the average local updates across all rounds, then the expected gradient norm is bounded as follows:

where O\mathcal{O} swallows all other constants (including L,σ2,κ2L,\sigma^{2},\kappa^{2}).

According to the Lipschitz-smooth assumption, it follows that

where the expectation is taken over randomly selected indices {lj}\{l_{j}\} as well as mini-batches ξi(t,k),∀i∈{1,2,…,m},k∈{0,1,…,τi−1}\xi_{i}^{(t,k)},\forall i\in\{1,2,\dots,m\},k\in\{0,1,\dots,\tau_{i}-1\}.

For the first term in 156, we can first take the expectation over indices and obtain

This term is exactly the same as the first term in 49. We can directly reuse previous results in the proof of Theorem 1. Comparing with 56, we have

where the last inequality comes from Lemma 5, stated below.

For the first term, by Cauchy-Schwarz inequality, we have

The second term can be bounded as following

Substituting (169) and (170) into (G) completes the proof. ∎

Substituting 160 and 164 into 156, we have

When ηL≤1/(2τeff)\eta L\leq 1/(2\tau_{\text{eff}}) and 6τeffηL+6τeffηLβ2/q≤126\tau_{\text{eff}}\eta L+6\tau_{\text{eff}}\eta L\beta^{2}/q\leq\frac{1}{2}, it follows that

Recall that the third term in 175 can be bounded as follows (see 87):

where D=4η2L2max⁡i{∥ai∥1(∥ai∥1−ai,−1)}<1D=4\eta^{2}L^{2}\max_{i}\{\left\|\bm{a}_{i}\right\|_{1}(\left\|\bm{a}_{i}\right\|_{1}-a_{i,-1})\}<1. If D≤112β2+1D\leq\frac{1}{12\beta^{2}+1}, then it follows that 11−D≤1+112β2≤2\frac{1}{1-D}\leq 1+\frac{1}{12\beta^{2}}\leq 2 and 3Dβ21−D≤14\frac{3D\beta^{2}}{1-D}\leq\frac{1}{4}. These facts can help us further simplify inequality 176. One can obtain

where the last inequality uses the fact that ∥a∥2≤∥a∥1\left\|\bm{a}\right\|_{2}\leq\left\|\bm{a}\right\|_{1}, for any vector a\bm{a}. Taking the total expectation and averaging all rounds, one can obtain

After minor rearranging, the above inequality is equivalent to

If we set the learning rate to be small enough, i.e., η=qτ~T\eta=\sqrt{\frac{q}{\widetilde{\tau}T}} where τ~=∑t=0T−1τ‾/T\widetilde{\tau}=\sum_{t=0}^{T-1}\overline{\tau}/T, then we get

where O\mathcal{O} swallows all other constants. ∎

Appendix H Pseudo-code of FedNova

Here we provide a pseudo-code of FedNova (see Algorithm 1) as a general algorithmic framework. Then, as an example, we show the pseudo-code of a special case of FedNova, where the local solver is specified as momentum SGD with cross-client variance reduction (see Algorithm 2). Note that when the server updates the global model, we set τeff\tau_{\text{eff}} to be the same as FedAvg, i.e., τeff=∑i∈Stpi∥ai(t)∥1\tau_{\text{eff}}=\sum_{i\in\mathcal{S}_{t}}p_{i}\|\bm{a}_{i}^{(t)}\|_{1} where St\mathcal{S}_{t} denotes the randomly selected subset of clients. Alternatively, the server can also choose other values of τeff\tau_{\text{eff}}.

Appendix I More Experiments Details

Platform. All experiments in this paper are conducted on a cluster of 1616 machines, each of which is equipped with one NVIDIA TitanX GPU. The machines communicate (i.e., transfer model parameters) with each other via Ethernet. We treat each machine as one client in the federated learning setting. The algorithms are implemented by PyTorch. We run each experiments for 33 times with different random seeds.

Hyper-parameter Choices. On non-IID CIFAR10 dataset, we fix the mini-batch size per client as 3232. When clients use momentum SGD as the local solver, the momentum factor is 0.90.9; when clients use proximal SGD, the proximal parameter μ\mu is selected from {0.0005,0.001,0.005,0.01}\{0.0005,0.001,0.005,0.01\}. It turns out that when Ei=2E_{i}=2, μ=0.005\mu=0.005 is the best and when Ei(t)∼U(2,5)E_{i}(t)\sim\mathcal{U}(2,5), μ=0.001\mu=0.001 is the best. The client learning rate η\eta is tuned from {0.005,0.01,0.02,0.05,0.08}\{0.005,0.01,0.02,0.05,0.08\} for FedAvg with each local solver separately. When using the same local solver, FedNova uses the same client learning rate as FedAvg. Specifically, if the local solver is momentum SGD, then we set η=0.02\eta=0.02. In other cases, η=0.05\eta=0.05 consistently performs the best. On the synthetic dataset, the mini-batch size per client is 2020 and the client learning rate is 0.020.02.

Training Curves on Non-IID CIFAR10. The training curves of FedAvg and FedNova are presented in Figure 7. Observe that FedNova (red curve) outperforms FedAvg (blue curve) by a large margin. FedNova only requires about half of the total rounds to achieve the same test accuracy as FedAvg. Besides, note that in , the test accuracy of FedAvg is higher than ours. This is because the authors of let clients to perform 2020 local epochs per round, which is 1010 times more than our setting. In , after 100100 communication rounds, FedAvg equivalently runs 100×20=2000100\times 20=2000 epochs.