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 epochs (traversals of their local dataset) of local-update stochastic gradient descent (SGD) with a mini-batch size . Thus, if a client has local data samples, the number of local SGD iterations is , 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, for all clients ). These works show that periodic consensus between the locally trained client models attains a stationary point of the global objective function , which is a sum of local objectives weighted by the dataset size . 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 varies across clients . 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 , but of an inconsistent objective , which can be arbitrarily different from depending upon the relative values of . To gain intuition into this phenomenon, observe in Figure 1 that if client performs more local updates, then the updated strays towards the local minimum , away from the true global minimum .
The Need for a General Analysis Framework. A naive approach to overcome heterogeneity is to fix a target number of local updates 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 equals to the true objective ), 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 returned by client (which performs local updates) in -th training round, the aggregator averages the normalized local gradients . 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 clients aim to jointly solve the following optimization problem:
where denotes the relative sample size, and is the local objective function at the -th client. Here, is the loss function (possibly non-convex) defined by the learning model and represents a data sample from local dataset . In the -th communication round, each client independently runs iterations of local solver (e.g., SGD) starting from the current global model to optimize its own local objective.
In our theoretical framework, we treat as an arbitrary scalar which can also vary across rounds. In practice, if clients run for the same local epochs , then , where is the mini-batch size. Alternately, if each communication round has a fixed length in terms of wall-clock time, then represents the local iterations completed by client 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 denotes client ’s model after the -th local update in the -th communication round and denotes the cumulative local progress made by client at round . Also, is the client learning rate and represents the stochastic gradient over a mini-batch of samples. When the number of clients 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 . 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 to each local objective, where is a tunable parameter. This proximal term pulls each local model backward closer to the global model . 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 . Since FedProx is a special case of our general framework, our convergence analysis provides sharp insights into the effect of . We show that a larger 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 .
For the objective function in 3, if client performs local steps per round, then FedAvg (with sufficiently small learning rate , 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 in VRLSGD by different ’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 clients with replacement with respect to probability , and then averaging the cumulative local changes with equal weights, or (ii) sampling clients without replacement uniformly at random, and then weighted averaging local changes, where the weight of client is re-scaled to . 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 , where denote the local parameter changes of client at round and , the fraction of data at client . We re-write this update rule in a more general form as follows:
The three key elements and of this update rule take different forms for different algorithms. Below, we provide detailed descriptions of these key elements.
Aggregation weights : Each client’s normalized gradient is multiplied with weight when computing the aggregated gradient . By definition, these weights satisfy . Observe that these weights determine the surrogate objective , which is optimized by the general algorithm in (4) instead of the true global objective – we will prove this formally in Theorem 1.
Effective number of steps : Since client makes local updates, the average number of local SGD steps per communication round is . However, the server can scale up or scale down the effect of the aggregated updates by setting the parameter larger or smaller than (analogous to choosing a global learning rate ). We refer to the ratio 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 and for a given local solver , 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 to the central server, which is just a re-scaled version of , 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 .
Previous Algorithms as Special Cases. Any previous algorithm whose accumulated local changes , 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, and are implicitly fixed by the choice of the local solver (i.e., the choice of ). Due to space limitations, the derivations of following examples are relegated to Appendix B.
When , FedProx is equivalent to FedAvg. As increases, the in FedProx is more similar to , thus making the surrogate objective more consistent. However, a larger corresponds to smaller , 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 where can vary across clients. As a result, we have and . Comparing with the case of FedProx 7, changing the values of has a similar effect as changing .
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 , where is the momentum factor, and .
More generally, the new formulation 6 suggests that whenever clients have different , which may be caused by imbalanced local updates (i.e., ’s have different dimensions), or various local learning rate/momentum schedules (i.e., ’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, .
For any sets of weights , there exist constants such that . If local functions are identical to each other, then we have .
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 . More specifically, if the total communication rounds is pre-determined and the learning rate is small enough where , then the optimization error will be bounded as follows:
where swallows all constants (including ), and quantities are defined as follows:
where is the last element in the vector .
In Appendix C, we also provide another version of this theorem that explicitly contains the local learning rate . Moreover, since the surrogate objective and the original objective are just different linear combinations of the local functions, once the algorithm converges to a stationary point of , one can also obtain some guarantees in terms of , as given by Theorem 2 below.
Under the same conditions as Theorem 1, the minimal gradient norm of the true global objective function will be bounded as follows:
where denotes the vanishing optimization error given by 8 and represents the chi-square divergence between vectors and .
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 such that . Also, when all local functions are identical to each other, we have . Only in these two special cases, is there no objective inconsistency. For most other algorithms subsumed by the general update rule in (4), both and are influenced by the choice of . When clients have different local progress (i.e., different vectors), previous algorithms will end up with a non-zero error floor , which does not vanish to even with sufficiently small learning rate. In Section D.1, we further construct a lower bound and show that , suggesting 10 is tight.
Novel Insights Into the Convergence of FedProx and the Effect of . Recall that in FedProx , where . 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 . As we increase , the weights come closer to and thus, the non-vanishing error in 10 decreases (see blue curve in Figure 4). However increasing worsens the slowdown , which appears in the first error term in 8 (see the red curve in Figure 4). In the extreme case when , although FedProx achieves objective consistency, it has a significantly slower convergence because and the first term in 8 is times larger than that with FedAvg (eq. to ).
Theorem 1 also reveals that, in FedProx, there should exist a best value of that balances all terms in 8. In Appendix, we provide a corollary showing that optimizes the error bound 8 of FedProx and yields a convergence rate of on the surrogate objective. This can serve as a guideline on setting in practice.
Linear Speedup Analysis. Another implication of Theorem 1 is that when the communication rounds is sufficiently large, then the convergence of the surrogate objective will be dominated by the first two terms in 8, which is . This suggests that the algorithm only uses total rounds when using 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 in 4, then the second non-vanishing term 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 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 instead of its default value . 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 -round will be corrected by . That is, the local gradient at the -th local step becomes . Besides, on the server side, one can also implement server momentum or adaptive server optimizers , in which the aggregated normalized gradient 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 using arbitrary gradient accumulation method per round. Under Assumptions 1, 2 and 3, and local learning rate as , where denotes the average local steps over all rounds at clients, then FedNova converges to a stationary point of in a rate of . The detailed bound is the same as the right hand side of 8, except that 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 clients are selected per round, then the convergence rate of FedNova is , where and is the set of selected client indices at the -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 is originally constructed in . The local dataset sizes 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 clients using a Dirichlet distribution , 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 is decayed by a constant factor after finishing and of the communication rounds. The initial value of is tuned separately for FedAvg with different local solvers. When using the same solver, FedNova uses the same as FedAvg to guarantee a fair comparison. On CIFAR-10, we run each experiment with 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 to , 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 communication rounds. When the client optimizer is SGD or SGD with momentum, simply changing the weights yields a - improvement on the test accuracy; When the client optimizer is proximal SGD, FedAvg is equivalent to FedProx. By setting and correcting the weights while keeping same as FedProx, FedNova-Prox achieves about 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 - higher test accuracy than vanilla SGD. This local momentum scheme can be further combined with server momentum . When , the hybrid momentum scheme achieves test accuracy As a reference, using server momentum alone achieves .
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 - 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 and . As a result, the global minimum is . 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 -th device can be written as follows:
where denotes the local model parameters at the -th local iteration after communication rounds, denotes the local learning rate and is a tunable hyper-parameter in FedProx. When , the algorithm will reduce to FedAvg. We omit the device index in , since it is synchronized and the same across all devices.
where . Then, after performing steps of local updates, the local model becomes
For the ease of writing, we define .
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 communication rounds, one can get
Accordingly, when , the iterates will converge to
Recall that .
Concrete Example in Lemma 1.
Now let us focus on a concrete example where and . Then, in this case, . 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 when using different local solvers. Recall that the local change at client is where stacks all stochastic gradients in the current round and is a non-negative vector.
In this case, we can write the update rule of local models as follows:
Subtracting on both sides, we obtain
Repeating the above procedure, it follows that
According to the definition, we have where .
B.2 SGD with Local Momentum
Let us firstly write down the update rule of the local models. Suppose that denotes the local momentum factor and is the local momentum buffer at client . 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 . Substituting 39 into 36, we have
Repeating the above procedure, it follows that
Then, the coefficient of is
Appendix C Proof of Theorem 1: Convergence of Surrogate Objective
For the ease of writing, let us define a surrogate objective function , where , and define the following auxiliary variables
According to the Lipschitz-smooth assumption, it follows that
where the expectation is taken over mini-batches . Before diving into the detailed bounds for and , we would like to firstly introduce several useful lemmas.
Assume . 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: .
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 , it follows that
where the last inequality uses the fact and Jensen’s Inequality: . 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 , one can derive that
where 67 uses Jensen’s Inequality again: , and 68 follows Assumption 1. Now, we turn to bounding the difference between the server model and the local model . Plugging into the local update rule and using the fact ,
where 74 follows from Jensen’s Inequality. Furthermore, note that
where is the last element in the vector . As a result, we have
In addition, we can bound the second term using the following inequality:
Define . 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 , then it follows that and . 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 where , 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 , the difference between the gradients of and can be bounded as follows:
where denotes the chi-square distance between and , i.e., .
According to the definition of and , we have
Applying Cauchy–Schwarz inequality, it follows that
where the last inequality uses Assumption 3. ∎
where 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 denotes the chi-square divergence between weight vectors and quantifies the dissimilarities among local objective functions and is defined in Assumption 3.
Suppose that there are only two clients with local objectives and . The global objective is defined as . For any set of weights , we define the surrogate objective function as . As a consequence, we have
Comparing with Assumption 3, we can define and in this case. Furthermore, according to the derivations in Appendix A, the iterate of FedAvg can be written as follows:
where . ∎
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., . 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 , then FedAvg algorithm (vanilla SGD with fixed local learning rate as local solver) will converge to the stationary point of a surrogate objective . The optimization error will be bounded as follows:
where swallows all constants (including ), and denotes the variance of local steps.
Consistent with Previous Results. When all clients perform the same local steps, , then and the above error bound 126 recovers previous results . When , then FedAvg reduces to fully synchronous SGD and the error bound 126 becomes , which is the same as standard SGD convergence rate .
E.2 FedProx
As a consequence, we can derive the closed-form expression of as follows:
Substituting 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 .
Consistency with FedAvg. From the update rule of FedProx, we know that when (or ), FedProx is equivalent to FedProx. This can also be validated from the expressions of . Using L’Hospital law, it is easy to show that
Best value of in FedProx. Given the expressions of and , we can further select a best value of that optimizes the error bound of FedProx, as stated in the following corollary.
Under the same conditions as Theorem 1 and suppose and , then 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 that optimizes the error upper bound of FedProx. That is to say, FedProx () is better than FedAvg () 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 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 , the quantities 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 . When the derivative equals to zero, we get
Plugging the expression of best into 138, we have
where 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 should be greater than . ∎
Appendix F Proof of Theorem 3
In the case of FedNova, the aggregated weights equals to . Therefore, the surrogate objective is the same as the original objective function . We can directly reuse the intermediate results in the proof of Theorem 1. According to 91, we have
where quantities are defined as follows:
Taking the total expectation and averaging over all rounds, it follows that
where , and . After minor rearranging, we have
Bt setting where , 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 -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 with probabilities , 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 clients with replacement to perform local computation. The probability of choosing the -th client is . In this case, FedNova will converge to the stationary points of the global objective . If we set where is the average local updates across all rounds, then the expected gradient norm is bounded as follows:
where swallows all other constants (including ).
According to the Lipschitz-smooth assumption, it follows that
where the expectation is taken over randomly selected indices as well as mini-batches .
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 and , it follows that
Recall that the third term in 175 can be bounded as follows (see 87):
where . If , then it follows that and . These facts can help us further simplify inequality 176. One can obtain
where the last inequality uses the fact that , for any vector . 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., where , then we get
where 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 to be the same as FedAvg, i.e., where denotes the randomly selected subset of clients. Alternatively, the server can also choose other values of .
Appendix I More Experiments Details
Platform. All experiments in this paper are conducted on a cluster of 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 times with different random seeds.
Hyper-parameter Choices. On non-IID CIFAR10 dataset, we fix the mini-batch size per client as . When clients use momentum SGD as the local solver, the momentum factor is ; when clients use proximal SGD, the proximal parameter is selected from . It turns out that when , is the best and when , is the best. The client learning rate is tuned from 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 . In other cases, consistently performs the best. On the synthetic dataset, the mini-batch size per client is and the client learning rate is .
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 local epochs per round, which is times more than our setting. In , after communication rounds, FedAvg equivalently runs epochs.