On the Convergence of FedAvg on Non-IID Data

Xiang Li, Kaixuan Huang, Wenhao Yang, Shusen Wang, Zhihua Zhang

Introduction

Federated Learning (FL), also known as federated optimization, allows multiple parties to collaboratively train a model without data sharing (Konevcnỳ et al., 2015; Shokri and Shmatikov, 2015; McMahan et al., 2017; Konevcnỳ, 2017; Sahu et al., 2018; Zhuo et al., 2019). Similar to the centralized parallel optimization (Jakovetic, 2013; Li et al., 2014a; b; Shamir et al., 2014; Zhang and Lin, 2015; Meng et al., 2016; Reddi et al., 2016; Richtárik and Takác, 2016; Smith et al., 2016; Zheng et al., 2016; Shusen Wang et al., 2018), FL let the user devices (aka worker nodes) perform most of the computation and a central parameter server update the model parameters using the descending directions returned by the user devices. Nevertheless, FL has three unique characters that distinguish it from the standard parallel optimization Li et al. (2019).

First, the training data are massively distributed over an incredibly large number of devices, and the connection between the central server and a device is slow. A direct consequence is the slow communication, which motivated communication-efficient FL algorithms (McMahan et al., 2017; Smith et al., 2017; Sahu et al., 2018; Sattler et al., 2019). Federated averaging (FedAvg) is the first and perhaps the most widely used FL algorithm. It runs EE steps of SGD in parallel on a small sampled subset of devices and then averages the resulting model updates via a central server once in a while.In original paper (McMahan et al., 2017), EE epochs of SGD are performed in parallel. For theoretical analyses, we denote by EE the times of updates rather than epochs. In comparison with SGD and its variants, FedAvg performs more local computation and less communication.

Second, unlike the traditional distributed learning systems, the FL system does not have control over users’ devices. For example, when a mobile phone is turned off or WiFi access is unavailable, the central server will lose connection to this device. When this happens during training, such a non-responding/inactive device, which is called a straggler, appears tremendously slower than the other devices. Unfortunately, since it has no control over the devices, the system can do nothing but waiting or ignoring the stragglers. Waiting for all the devices’ response is obviously infeasible; it is thus impractical to require all the devices be active.

Third, the training data are non-iidThroughout this paper, “non-iid” means data are not identically distributed. More precisely, the data distributions in the kk-th and ll-th devices, denote DkD_{k} and DlD_{l}, can be different., that is, a device’s local data cannot be regarded as samples drawn from the overall distribution. The data available locally fail to represent the overall distribution. This does not only bring challenges to algorithm design but also make theoretical analysis much harder. While FedAvg actually works when the data are non-iid McMahan et al. (2017), FedAvg on non-iid data lacks theoretical guarantee even in convex optimization setting.

There have been much efforts developing convergence guarantees for FL algorithm based on the assumptions that (1) the data are iid and (2) all the devices are active. Khaled et al. (2019); Yu et al. (2019); Wang et al. (2019) made the latter assumption, while Zhou and Cong (2017); Stich (2018); Wang and Joshi (2018); Woodworth et al. (2018) made both assumptions. The two assumptions violates the second and third characters of FL. Previous algorithm Fedprox Sahu et al. (2018) doesn’t require the two mentioned assumptions and incorporates FedAvg as a special case when the added proximal term vanishes. However, their theory fails to cover FedAvg.

Let NN be the total number of user devices and KK (≤N\leq N) be the maximal number of devices that participate in every round’s communication. Let TT be the total number of every device’s SGDs, EE be the number of local iterations performed in a device between two communications, and thus TE\frac{T}{E} is the number of communications.

Contributions.

For strongly convex and smooth problems, we establish a convergence guarantee for FedAvg without making the two impractical assumptions: (1) the data are iid, and (2) all the devices are active. To the best of our knowledge, this work is the first to show the convergence rate of FedAvg without making the two assumptions.

We show in Theorem 1, 2, and 3 that FedAvg has O(1T){\mathcal{O}}(\frac{1}{T}) convergence rate. In particular, Theorem 3 shows that to attain a fixed precision ϵ\epsilon, the number of communications is

Here, GG, Γ\Gamma, pkp_{k}, and σk\sigma_{k} are problem-related constants defined in Section 3.1. The most interesting insight is that EE is a knob controlling the convergence rate: neither setting EE over-small (E=1E=1 makes FedAvg equivalent to SGD) nor setting EE over-large is good for the convergence.

This work also makes algorithmic contributions. We summarize the existing samplingThroughout this paper, “sampling” refers to how the server chooses KK user devices and use their outputs for updating the model parameters. “Sampling” does not mean how a device randomly selects training samples. and averaging schemes for FedAvg (which do not have convergence bounds before this work) and propose a new scheme (see Table 1). We point out that a suitable sampling and averaging scheme is crucial for the convergence of FedAvg. To the best of our knowledge, we are the first to theoretically demonstrate that FedAvg with certain schemes (see Table 1) can achieve O(1T){\mathcal{O}}(\frac{1}{T}) convergence rate in non-iid federated setting. We show that heterogeneity of training data and partial device participation slow down the convergence. We empirically verify our results through numerical experiments.

Paper organization.

In Section 2, we elaborate on FedAvg. In Section 3, we present our main convergence bounds for FedAvg. In Section 4, we construct a special example to show the necessity of learning rate decay. In Section 5, we discuss and compare with prior work. In Section 6, we conduct empirical study to verify our theories. All the proofs are left to the appendix.

Federated Averaging (FedAvg)

In this work, we consider the following distributed optimization model:

where NN is the number of devices, and pkp_{k} is the weight of the kk-th device such that pk≥0p_{k}\geq 0 and ∑k=1Npk=1\sum_{k=1}^{N}p_{k}=1. Suppose the kk-th device holds the nkn_{k} training data: xk,1,xk,2,⋯ ,xk,nkx_{k,1},x_{k,2},\cdots,x_{k,n_{k}}. The local objective Fk(⋅)F_{k}(\cdot) is defined by

Algorithm description.

Here, we describe one around (say the tt-th) of the standard FedAvg algorithm. First, the central server broadcasts the latest model, wt{\bf w}_{t}, to all the devices. Second, every device (say the kk-th) lets wtk=wt{\bf w}_{t}^{k}={\bf w}_{t} and then performs EE (≥1\geq 1) local updates:

where ηt+i\eta_{t+i} is the learning rate (a.k.a. step size) and ξt+ik\xi_{t+i}^{k} is a sample uniformly chosen from the local data. Last, the server aggregates the local models, wt+E1,⋯ ,wt+EN{\bf w}_{t+E}^{1},\cdots,{\bf w}_{t+E}^{N}, to produce the new global model, wt+E{\bf w}_{t+E}. Because of the non-iid and partial device participation issues, the aggregation step can vary.

IID versus non-iid.

Suppose the data in the kk-th device are i.i.d. sampled from the distribution Dk{\mathcal{D}}_{k}. Then the overall distribution is a mixture of all local data distributions: D=∑k=1NpkDk{\mathcal{D}}=\sum_{k=1}^{N}p_{k}{\mathcal{D}}_{k}. The prior work Zhang et al. (2015a); Zhou and Cong (2017); Stich (2018); Wang and Joshi (2018); Woodworth et al. (2018) assumes the data are iid generated by or partitioned among the NN devices, that is, Dk=D{\mathcal{D}}_{k}={\mathcal{D}} for all k∈[N]k\in[N]. However, real-world applications do not typically satisfy the iid assumption. One of our theoretical contributions is avoiding making the iid assumption.

Full device participation.

The prior work Coppola (2015); Zhou and Cong (2017); Stich (2018); Yu et al. (2019); Wang and Joshi (2018); Wang et al. (2019) requires the full device participation in the aggregation step of FedAvg. In this case, the aggregation step performs

Unfortunately, the full device participation requirement suffers from serious “straggler’s effect” (which means everyone waits for the slowest) in real-world applications. For example, if there are thousands of users’ devices in the FL system, there are always a small portion of devices offline. Full device participation means the central server must wait for these “stragglers”, which is obviously unrealistic.

Partial device participation.

This strategy is much more realistic because it does not require all the devices’ output. We can set a threshold KK (1≤K<N1\leq K<N) and let the central server collect the outputs of the first KK responded devices. After collecting KK outputs, the server stops waiting for the rest; the K+1K+1-th to NN-th devices are regarded stragglers in this iteration. Let St{\mathcal{S}}_{t} (∣St∣=K|{\mathcal{S}}_{t}|=K) be the set of the indices of the first KK responded devices in the tt-th iteration. The aggregation step performs

It can be proved that NK∑k∈Stpk\frac{N}{K}\sum_{k\in{\mathcal{S}}_{t}}p_{k} equals one in expectation.

Communication cost.

The FedAvg requires two rounds communications— one broadcast and one aggregation— per EE iterations. If TT iterations are performed totally, then the number of communications is ⌊2TE⌋\lfloor\frac{2T}{E}\rfloor. During the broadcast, the central server sends wt{\bf w}_{t} to all the devices. During the aggregation, all or part of the NN devices sends its output, say wt+Ek{\bf w}_{t+E}^{k}, to the server.

Convergence Analysis of FedAvg in Non-iid Setting

In this section, we show that FedAvg converges to the global optimum at a rate of O(1/T){\mathcal{O}}(1/T) for strongly convex and smooth functions and non-iid data. The main observation is that when the learning rate is sufficiently small, the effect of EE steps of local updates is similar to one step update with a larger learning rate. This coupled with appropriate sampling and averaging schemes would make each global update behave like an SGD update. Partial device participation (K<NK<N) only makes the averaged sequence {wt}\{{\bf w}_{t}\} have a larger variance, which, however, can be controlled by learning rates. These imply the convergence property of FedAvg should not differ too much from SGD. Next, we will first give the convergence result with full device participation (i.e., K=NK=N) and then extend this result to partial device participation (i.e., K<NK<N).

F1,⋯ ,FNF_{1},\cdots,F_{N} are all LL-smooth: for all v{\bf v} and w{\bf w}, Fk(v)≤Fk(w)+(v−w)T∇Fk(w)+L2∥v−w∥22F_{k}({\bf v})\leq F_{k}({\bf w})+({\bf v}-{\bf w})^{T}\nabla F_{k}({\bf w})+\frac{L}{2}\|{\bf v}-{\bf w}\|_{2}^{2}.

F1,⋯ ,FNF_{1},\cdots,F_{N} are all μ\mu-strongly convex: for all v{\bf v} and w{\bf w}, Fk(v)≥Fk(w)+(v−w)T∇Fk(w)+μ2∥v−w∥22F_{k}({\bf v})\geq F_{k}({\bf w})+({\bf v}-{\bf w})^{T}\nabla F_{k}({\bf w})+\frac{\mu}{2}\|{\bf v}-{\bf w}\|_{2}^{2}.

Assumptions 3 and 4 have been made by the works Zhang et al. (2013); Stich (2018); Stich et al. (2018); Yu et al. (2019).

Let F∗F^{*} and Fk∗F_{k}^{*} be the minimum values of FF and FkF_{k}, respectively. We use the term Γ=F∗−∑k=1NpkFk∗\Gamma=F^{*}-\sum_{k=1}^{N}p_{k}F_{k}^{*} for quantifying the degree of non-iid. If the data are iid, then Γ\Gamma obviously goes to zero as the number of samples grows. If the data are non-iid, then Γ\Gamma is nonzero, and its magnitude reflects the heterogeneity of the data distribution.

2 Convergence Result: Full Device Participation

Here we analyze the case that all the devices participate in the aggregation step; see Section 2 for the algorithm description. Let the FedAvg algorithm terminate after TT iterations and return wT{\bf w}_{T} as the solution. We always require TT is evenly divisible by EE so that FedAvg can output wT{\bf w}_{T} as expected.

Let Assumptions 1 to 4 hold and L,μ,σk,GL,\mu,\sigma_{k},G be defined therein. Choose κ=Lμ\kappa=\frac{L}{\mu}, γ=max⁡{8κ,E}\gamma=\max\{8\kappa,E\} and the learning rate ηt=2μ(γ+t)\eta_{t}=\frac{2}{\mu(\gamma+t)}. Then FedAvg with full device participation satisfies

3 Convergence Result: Partial Device Participation

As discussed in Section 2, partial device participation has more practical interest than full device participation. Let the set St{\mathcal{S}}_{t} (⊂[N]\subset[N]) index the active devices in the tt-th iteration. To establish the convergence bound, we need to make assumptions on St{\mathcal{S}}_{t}.

Assumption 5 assumes the KK indices are selected from the distribution pkp_{k} independently and with replacement. The aggregation step is simply averaging. This is first proposed in (Sahu et al., 2018), but they did not provide theoretical analysis.

Assume St{\mathcal{S}}_{t} contains a subset of KK indices randomly selected with replacement according to the sampling probabilities p1,⋯ ,pNp_{1},\cdots,p_{N}. The aggregation step of FedAvg performs wt⟵1K∑k∈Stwtk{\bf w}_{t}\longleftarrow\frac{1}{K}\sum_{k\in{\mathcal{S}}_{t}}{\bf w}_{t}^{k}.

Let Assumptions 1 to 4 hold and L,μ,σk,GL,\mu,\sigma_{k},G be defined therein. Let κ,γ\kappa,\gamma, ηt\eta_{t}, and BB be defined in Theorem 1. Let Assumption 5 hold and define C=4KE2G2C=\frac{4}{K}E^{2}G^{2}. Then

Alternatively, we can select KK indices from [N][N] uniformly at random without replacement. As a consequence, we need a different aggregation strategy. Assumption 6 assumes the KK indices are selected uniformly without replacement and the aggregation step is the same as in Section 2. However, to guarantee convergence, we require an additional assumption of balanced data.

Assume St{\mathcal{S}}_{t} contains a subset of KK indices uniformly sampled from [N][N] without replacement. Assume the data is balanced in the sense that p1=⋯=pN=1Np_{1}=\cdots=p_{N}=\frac{1}{N}. The aggregation step of FedAvg performs wt⟵NK∑k∈Stpk wtk{\bf w}_{t}\longleftarrow\frac{N}{K}\sum_{k\in{\mathcal{S}}_{t}}p_{k}\,{\bf w}_{t}^{k}.

Replace Assumption 5 by Assumption 6 and CC by C=N−KN−14KE2G2C=\frac{N-K}{N-1}\frac{4}{K}E^{2}G^{2}. Then the same bound in Theorem 2 holds.

Scheme II requires p1=⋯=pN=1Np_{1}=\cdots=p_{N}=\frac{1}{N} which obviously violates the unbalance nature of FL. Fortunately, this can be addressed by the following transformation. Let F~k(w)=pkNFk(w)\widetilde{F}_{k}({\bf w})=p_{k}NF_{k}({\bf w}) be a scaled local objective FkF_{k}. Then the global objective becomes a simple average of all scaled local objectives:

Theorem 3 still holds if L,μ,σk,GL,\mu,\sigma_{k},G are replaced by L~≜νL\widetilde{L}\triangleq\nu L, μ~≜ςμ\widetilde{\mu}\triangleq\varsigma\mu, σ~k=νσ\widetilde{\sigma}_{k}=\sqrt{\nu}\sigma, and G~=νG\widetilde{G}=\sqrt{\nu}G, respectively. Here, ν=N⋅max⁡kpk\nu=N\cdot\max_{k}p_{k} and ς=N⋅min⁡kpk\varsigma=N\cdot\min_{k}p_{k}.

4 Discussions

Since ∥w0−w∗∥2≤4μ2G2\|{\bf w}_{0}-{\bf w}^{*}\|^{2}\leq\frac{4}{\mu^{2}}G^{2} for μ\mu-strongly convex FF, the dominating term in eqn. (6) is

Let TϵT_{\epsilon} denote the number of required steps for FedAvg to achieve an ϵ\epsilon accuracy. It follows from eqn. (7) that the number of required communication rounds is roughlyHere we use γ=O(κ+E)\gamma={\mathcal{O}}(\kappa+E).

Thus, TϵE\frac{T_{\epsilon}}{E} is a function of EE that first decreases and then increases, which implies that over-small or over-large EE may lead to high communication cost and that the optimal EE exists.

Stich (2018) showed that if the data are iid, then EE can be set to O(T){\mathcal{O}}(\sqrt{T}). However, this setting does not work if the data are non-iid. Theorem 1 implies that EE must not exceed Ω(T)\Omega(\sqrt{T}); otherwise, convergence is not guaranteed. Here we give an intuitive explanation. If EE is set big, then wtk{\bf w}_{t}^{k} can converge to the minimizer of FkF_{k}, and thus FedAvg becomes the one-shot average Zhang et al. (2013) of the local solutions. If the data are non-iid, the one-shot averaging does not work because weighted average of the minimizers of F1,⋯ ,FNF_{1},\cdots,F_{N} can be very different from the minimizer of FF.

Choice of K𝐾K.

Stich (2018) showed that if the data are iid, the convergence rate improves substantially as KK increases. However, under the non-iid setting, the convergence rate has a weak dependence on KK, as we show in Theorems 2 and 3. This implies FedAvg is unable to achieve linear speedup. We have empirically observed this phenomenon (see Section 6). Thus, in practice, the participation ratio KN\frac{K}{N} can be set small to alleviate the straggler’s effect without affecting the convergence rate.

Choice of sampling schemes.

We considered two sampling and averaging schemes in Theorems 2 and 3. Scheme I selects KK devices according to the probabilities p1,⋯ ,pNp_{1},\cdots,p_{N} with replacement. The non-uniform sampling results in faster convergence than uniform sampling, especially when p1,⋯ ,pNp_{1},\cdots,p_{N} are highly non-uniform. If the system can choose to activate any of the NN devices at any time, then Scheme I should be used.

However, oftentimes the system has no control over the sampling; instead, the server simply uses the first KK returned results for the update. In this case, we can assume the KK devices are uniformly sampled from all the NN devices and use Theorem 3 to guarantee the convergence. If p1,⋯ ,pNp_{1},\cdots,p_{N} are highly non-uniform, then ν=N⋅max⁡kpk\nu=N\cdot\max_{k}p_{k} is big and ς=N⋅min⁡kpk\varsigma=N\cdot\min_{k}p_{k} is small, which makes the convergence of FedAvg slow. This point of view is empirically verified in our experiments.

Necessity of Learning Rate Decay

In this section, we point out that diminishing learning rates are crucial for the convergence of FedAvg in the non-iid setting. Specifically, we establish the following theorem by constructing a ridge regression model (which is strongly convex and smooth).

where we hide some problem dependent constants.

Theorem 4 and its proof provide several implications. First, the decay of learning rate is necessary of FedAvg. On the one hand, Theorem 1 shows with E>1E>1 and a decaying learning rate, FedAvg converges to the optimum. On the other hand, Theorem 4 shows that with E>1E>1 and any fixed learning rate, FedAvg does not converges to the optimum.

Third, Theorem 4 shows the requirement of learning rate decay is not an artifact of our analysis; instead, it is inherently required by FedAvg. An explanation is that constant learning rates, combined with EE steps of possibly-biased local updates, form a sub-optimal update scheme, but a diminishing learning rate can gradually eliminate such bias.

The efficiency of FedAvg principally results from the fact that it performs several update steps on a local model before communicating with other workers, which saves communication. Diminishing step sizes often hinders fast convergence, which may counteract the benefit of performing multiple local updates. Theorem 4 motivates more efficient alternatives to FedAvg.

Related Work

Federated learning (FL) was first proposed by McMahan et al. (2017) for collaboratively learning a model without collecting users’ data. The research work on FL is focused on the communication-efficiency Konevcnỳ et al. (2016); McMahan et al. (2017); Sahu et al. (2018); Smith et al. (2017) and data privacy Bagdasaryan et al. (2018); Bonawitz et al. (2017); Geyer et al. (2017); Hitaj et al. (2017); Melis et al. (2019). This work is focused on the communication-efficiency issue.

FedAvg, a synchronous distributed optimization algorithm, was proposed by McMahan et al. (2017) as an effective heuristic. Sattler et al. (2019); Zhao et al. (2018) studied the non-iid setting, however, they do not have convergence rate. A contemporaneous and independent work Xie et al. (2019) analyzed asynchronous FedAvg; while they did not require iid data, their bound do not guarantee convergence to saddle point or local minimum. Sahu et al. (2018) proposed a federated optimization framework called FedProx to deal with statistical heterogeneity and provided the convergence guarantees in non-iid setting. FedProx adds a proximal term to each local objective. When these proximal terms vanish, FedProx is reduced to FedAvg. However, their convergence theory requires the proximal terms always exist and hence fails to cover FedAvg.

When data are iid distributed and all devices are active, FedAvg is referred to as LocalSGD. Due to the two assumptions, theoretical analysis of LocalSGD is easier than FedAvg. Stich (2018) demonstrated LocalSGD provably achieves the same linear speedup with strictly less communication for strongly-convex stochastic optimization. Coppola (2015); Zhou and Cong (2017); Wang and Joshi (2018) studied LocalSGD in the non-convex setting and established convergence results. Yu et al. (2019); Wang et al. (2019) recently analyzed LocalSGD for non-convex functions in heterogeneous settings. In particular, Yu et al. (2019) demonstrated LocalSGD also achieves O(1/NT)\mathcal{O}(1/\sqrt{NT}) convergence (i.e., linear speedup) for non-convex optimization. Lin et al. (2018) empirically shows variants of LocalSGD increase training efficiency and improve the generalization performance of large batch sizes while reducing communication. For LocalGD on non-iid data (as opposed to LocalSGD), the best result is by the contemporaneous work (but slightly later than our first version) (Khaled et al., 2019). Khaled et al. (2019) used fixed learning rate η\eta and showed O(1T){\mathcal{O}}(\frac{1}{T}) convergence to a point O(η2E2){\mathcal{O}}(\eta^{2}E^{2}) away from the optimal. In fact, the suboptimality is due to their fixed learning rate. As we show in Theorem 4, using a fixed learning rate η\eta throughout, the solution by LocalGD is at least Ω((E−1)η)\Omega((E-1)\eta) away from the optimal.

If the data are iid, distributed optimization can be efficiently solved by the second-order algorithms Mahajan et al. (2018); Reddi et al. (2016); Shamir et al. (2014); Shusen Wang et al. (2018); Zhang and Lin (2015) and the one-shot methods Lee et al. (2017); Lin et al. (2017); Wang (2019); Zhang et al. (2013; 2015b). The primal-dual algorithms Hong et al. (2018); Smith et al. (2016; 2017) are more generally applicable and more relevant to FL.

Numerical Experiments

We examine our theoretical results on a logistic regression with weight decay λ=1e−4\lambda=1e-4. This is a stochastic convex optimization problem. We distribute MNIST dataset (LeCun et al., 1998) among N=100N=100 workers in a non-iid fashion such that each device contains samples of only two digits. We further obtain two datasets: mnist balanced and mnist unbalanced. The former is balanced such that the number of samples in each device is the same, while the latter is highly unbalanced with the number of samples among devices following a power law. To manipulate heterogeneity more precisly, we synthesize unbalanced datasets following the setup in Sahu et al. (2018) and denote it as synthetic(α\alpha, β\beta) where α\alpha controls how much local models differ from each other and β\beta controls how much the local data at each device differs from that of other devices. We obtain two datasets: synthetic(0,0) and synthetic(1,1). Details can be found in Appendix D.

Experiment settings

For all experiments, we initialize all runnings with w0=0{\bf w}_{0}=0. In each round, all selected devices run EE steps of SGD in parallel. We decay the learning rate at the end of each round by the following scheme ηt=η01+t\eta_{t}=\frac{\eta_{0}}{1+t}, where η0\eta_{0} is chosen from the set {1,0.1,0.01}\{1,0.1,0.01\}. We evaluate the averaged model after each global synchronization on the corresponding global objective. For fair comparison, we control all randomness in experiments so that the set of activated devices is the same across all different algorithms on one configuration.

Impact of E𝐸E

We expect that Tϵ/ET_{\epsilon}/E, the required communication round to achieve curtain accuracy, is a hyperbolic finction of EE as equ (8) indicates. Intuitively, a small EE means a heavy communication burden, while a large EE means a low convergence rate. One needs to trade off between communication efficiency and fast convergence. We empirically observe this phenomenon on unbalanced datasets in Figure 1a. The reason why the phenomenon does not appear in mnist balanced dataset requires future investigations.

Impact of K𝐾K

Our theory suggests that a larger KK may slightly accelerate convergence since Tϵ/ET_{\epsilon}/E contains a term O(EG2K)\mathcal{O}\left(\frac{EG^{2}}{K}\right). Figure 1b shows that KK has limited influence on the convergence of FedAvg in synthetic(0,0) dataset. It reveals that the curve of a large enough KK is slightly better. We observe similar phenomenon among the other three datasets and attach additional results in Appendix D. This justifies that when the variance resulting sampling is not too large (i.e., B≫CB\gg C), one can use a small number of devices without severely harming the training process, which also removes the need to sample as many devices as possible in convex federated optimization.

Effect of sampling and averaging schemes.

We compare four schemes among four federated datasets. Since the original scheme involves a history term and may be conservative, we carefully set the initial learning rate for it. Figure 1c indicates that when data are balanced, Schemes I and II achieve nearly the same performance, both better than the original scheme. Figure 1d shows that when the data are unbalanced, i.e., pkp_{k}’s are uneven, Scheme I performs the best. Scheme II suffers from some instability in this case. This is not contradictory with our theory since we don’t guarantee the convergence of Scheme II when data is unbalanced. As expected, transformed Scheme II performs stably at the price of a lower convergence rate. Compared to Scheme I, the original scheme converges at a slower speed even if its learning rate is fine tuned. All the results show the crucial position of appropriate sampling and averaging schemes for FedAvg.

Conclusion

Federated learning becomes increasingly popular in machine learning and optimization communities. In this paper we have studied the convergence of FedAvg, a heuristic algorithm suitable for federated setting. We have investigated the influence of sampling and averaging schemes. We have provided theoretical guarantees for two schemes and empirically demonstrated their performances. Our work sheds light on theoretical understanding of FedAvg and provides insights for algorithm design in realistic applications. Though our analyses are constrained in convex problems, we hope our insights and proof techniques can inspire future work.

Acknowledgements

Li, Yang and Zhang have been supported by the National Natural Science Foundation of China (No. 11771002 and 61572017), Beijing Natural Science Foundation (Z190001), the Key Project of MOST of China (No. 2018AAA0101000), and Beijing Academy of Artificial Intelligence (BAAI).

References

Appendix A Proof of Theorem 1

We analyze FedAvg in the setting of full device participation in this section.

Let wtk{\bf w}_{t}^{k} be the model parameter maintained in the kk-th device at the tt-th step. Let IE{\mathcal{I}}_{E} be the set of global synchronization steps, i.e., IE={nE ∣ n=1,2,⋯ }{\mathcal{I}}_{E}=\{nE\ |\ n=1,2,\cdots\}. If t+1∈IEt+1\in{\mathcal{I}}_{E}, i.e., the time step to communication, FedAvg activates all devices. Then the update of FedAvg with partial devices active can be described as

Here, an additional variable vt+1k{\bf v}_{t+1}^{k} is introduced to represent the immediate result of one step SGD update from wtk{\bf w}_{t}^{k}. We interpret wt+1k{\bf w}_{t+1}^{k} as the parameter obtained after communication steps (if possible).

A.2 Key Lemmas

To convey our proof clearly, it would be necessary to prove certain useful lemmas. We defer the proof of these lemmas to latter section and focus on proving the main theorem.

Assume Assumption 1 and 2. If ηt≤14L\eta_{t}\leq\frac{1}{4L}, we have

where Γ=F∗−∑k=1NpkFk⋆≥0\Gamma=F^{*}-\sum_{k=1}^{N}p_{k}F_{k}^{\star}\geq 0.

Assume Assumption 3 holds. It follows that

Assume Assumption 4, that ηt\eta_{t} is non-increasing and ηt≤2ηt+E\eta_{t}\leq 2\eta_{t+E} for all t≥0t\geq 0. It follows that

A.3 Completing the Proof of Theorem 1

For a diminishing stepsize, ηt=βt+γ\eta_{t}=\frac{\beta}{t+\gamma} for some β>1μ\beta>\frac{1}{\mu} and γ>0\gamma>0 such that η1≤min⁡{1μ,14L}=14L\eta_{1}\leq\min\{\frac{1}{\mu},\frac{1}{4L}\}=\frac{1}{4L} and ηt≤2ηt+E\eta_{t}\leq 2\eta_{t+E}. We will prove Δt≤vγ+t\Delta_{t}\leq\frac{v}{\gamma+t} where v=max⁡{β2Bβμ−1,(γ+1)Δ1}v=\max\left\{\frac{\beta^{2}B}{\beta\mu-1},(\gamma+1)\Delta_{1}\right\}.

We prove it by induction. Firstly, the definition of vv ensures that it holds for t=1t=1. Assume the conclusion holds for some tt, it follows that

Then by the LL-smoothness of F(⋅)F(\cdot),

Specifically, if we choose β=2μ,γ=max⁡{8Lμ,E}−1\beta=\frac{2}{\mu},\gamma=\max\{8\frac{L}{\mu},E\}-1 and denote κ=Lμ\kappa=\frac{L}{\mu}, then ηt=2μ1γ+t\eta_{t}=\frac{2}{\mu}\frac{1}{\gamma+t}. One can verify that the choice of ηt\eta_{t} satisfies ηt≤2ηt+E\eta_{t}\leq 2\eta_{t+E} for t≥1t\geq 1. Then, we have

A.4 Deferred proofs of key lemmas

Notice that v‾t+1=w‾t−ηtgt\overline{{\bf v}}_{t+1}=\overline{{\bf w}}_{t}-\eta_{t}{\bm{g}}_{t}, then

From the the LL-smoothness of Fk(⋅)F_{k}(\cdot), it follows that

By the convexity of ∥⋅∥2\left\|\cdot\right\|^{2} and eqn. (16), we have

By Cauchy-Schwarz inequality and AM-GM inequality, we have

By the μ\mu-strong convexity of Fk(⋅)F_{k}(\cdot), we have

By combining eqn. (15), eqn. (A.4), eqn. (18) and eqn. (19), it follows that

We next aim to bound CC. We define γt=2ηt(1−2Lηt)\gamma_{t}=2\eta_{t}(1-2L\eta_{t}). Since ηt≤14L\eta_{t}\leq\frac{1}{4L}, ηt≤γt≤2ηt\eta_{t}\leq\gamma_{t}\leq 2\eta_{t}. Then we split CC into two terms:

where in the last equation, we use the notation Γ=∑k=1Npk(F∗−Fk∗)=F∗−∑k=1NpkFk∗\Gamma=\sum_{k=1}^{N}p_{k}\left(F^{*}-F_{k}^{*}\right)=F^{*}-\sum_{k=1}^{N}p_{k}F_{k}^{*}.

where the first inequality results from the convexity of Fk(⋅)F_{k}(\cdot), the second inequality from AM-GM inequality and the third inequality from eqn. (16).

where in the last inequality, we use the following facts: (1) ηtL−1≤−34≤0\eta_{t}L-1\leq-\frac{3}{4}\leq 0 and ∑k=1Npk(Fk(w‾t)−F∗)=F(w‾t)−F∗≥0\sum_{k=1}^{N}p_{k}\left(F_{k}(\overline{{\bf w}}_{t})-F^{*}\right)=F(\overline{{\bf w}}_{t})-F^{*}\geq 0 (2) Γ≥0\Gamma\geq 0 and 4Lηt2+γtηtL≤6ηt2L4L\eta_{t}^{2}+\gamma_{t}\eta_{t}L\leq 6\eta_{t}^{2}L and (3) γt2ηt≤1\frac{\gamma_{t}}{2\eta_{t}}\leq 1.

Recalling the expression of A1A_{1} and plugging CC into it, we have

Using eqn. (A.4) and taking expectation on both sides of eqn. (A.4), we erase the randomness from stochastic gradients, we complete the proof. ∎

From Assumption 3, the variance of the stochastic gradients in device kk is bounded by σk2\sigma_{k}^{2}, then

Since FedAvg requires a communication each EE steps. Therefore, for any t≥0t\geq 0, there exists a t0≤tt_{0}\leq t, such that t−t0≤E−1t-t_{0}\leq E-1 and wt0k=w‾t0{\bf w}_{t_{0}}^{k}=\overline{{\bf w}}_{t_{0}} for all k=1,2,⋯ ,Nk=1,2,\cdots,N. Also, we use the fact that ηt\eta_{t} is non-increasing and ηt0≤2ηt\eta_{t_{0}}\leq 2\eta_{t} for all t−t0≤E−1t-t_{0}\leq E-1, then

Appendix B Proofs of Theorems 2 and 3

We analyze FedAvg in the setting of partial device participation in this section.

All sampling schemes can be divided into two groups, one with replacement and the other without replacement. For those with replacement, it is possible for a device to be activated several times in a round of communication, even though each activation is independent with the rest. We denote by Ht{\mathcal{H}}_{t} the multiset selected which allows any element to appear more than once. Note that Ht{\mathcal{H}}_{t} is only well defined for t∈IEt\in{\mathcal{I}}_{E}. For convenience, we denote by St=HN(t,E){\mathcal{S}}_{t}={\mathcal{H}}_{N(t,E)} the most recent set of chosen devices where N(t,E)=max⁡{n∣n≤t,n∈IE}N(t,E)=\max\{n|n\leq t,n\in{\mathcal{I}}_{E}\}.

Updating scheme.

Limited to realistic scenarios (for communication efficiency and low straggler effect), FedAvg first samples a random multiset St{\mathcal{S}}_{t} of devices and then only perform updates on them. This make the analysis a little bit intricate, since St{\mathcal{S}}_{t} varies each EE steps. However, we can use a thought trick to circumvent this difficulty. We assume that FedAvg always activates all devices at the beginning of each round and then uses the parameters maintained in only a few sampled devices to produce the next-round parameter. It is clear that this updating scheme is equivalent to the original. Then the update of FedAvg with partial devices active can be described as: for all k∈[N]k\in[N],

Sources of randomness.

B.2 Key Lemmas

For full device participation, we always have w‾t+1=v‾t+1\overline{{\bf w}}_{t+1}=\overline{{\bf v}}_{t+1}. This is true when t+1∉IEt+1\notin{\mathcal{I}}_{E} for partial device participation. When t+1∈IEt{+}1\in{\mathcal{I}}_{E}, we hope this relation establish in the sense of expectation. To that end, we require the sampling and averaging scheme to be unbiased in the sense that

We find two sampling and averaging schemes satisfying the requirement and provide convergence guarantees.

The server establishes St+1{\mathcal{S}}_{t+1} by i.i.d. with replacement sampling an index k∈{1,⋯ ,N}k\in\{1,\cdots,N\} with probabilities p1,⋯ ,pNp_{1},\cdots,p_{N} for KK times. Hence St+1{\mathcal{S}}_{t+1} is a multiset which allows a element to occur more than once. Then the server averages the parameters by wt+1k=1K∑k∈St+1vt+1k{\bf w}_{t+1}^{k}=\frac{1}{K}\sum_{k\in{\mathcal{S}}_{t+1}}{\bf v}_{t+1}^{k}. This is first proposed in [Sahu et al., 2018] but lacks theoretical analysis.

The server samples St+1{\mathcal{S}}_{t+1} uniformly in a without replacement fashion. Hence each element in St+1{\mathcal{S}}_{t+1} only occurs once.Then server averages the parameters by wt+1k=∑k∈St+1pkNKvt+1k{\bf w}_{t+1}^{k}=\sum_{k\in{\mathcal{S}}_{t+1}}p_{k}\frac{N}{K}{\bf v}_{t+1}^{k}. Note that when the pkp_{k}’s are not all the same, one cannot ensure ∑k∈St+1pkNK=1\sum_{k\in{\mathcal{S}}_{t+1}}p_{k}\frac{N}{K}=1.

Unbiasedness and bounded variance.

Lemma 4 shows the mentioned two sampling and averaging schemes are unbiased. In expectation, the next-round parameter (i.e., w‾t+1\overline{{\bf w}}_{t+1}) is equal to the weighted average of parameters in all devices after SGD updates (i.e., v‾t+1\overline{{\bf v}}_{t+1}). However, the original scheme in [McMahan et al., 2017] (see Table 1) does not enjoy this property. But it is very similar to Scheme II except the averaging scheme. Hence our analysis cannot cover the original scheme.

If t+1∈IEt+1\in{\mathcal{I}}_{E}, for Scheme I and Scheme II, we have

For t+1∈It+1\in{\mathcal{I}}, assume that ηt\eta_{t} is non-increasing and ηt≤2ηt+E\eta_{t}\leq 2\eta_{t+E} for all t≥0t\geq 0. We have the following results.

For Scheme I, the expected difference between v‾t+1\overline{{\bf v}}_{t+1} and w‾t+1\overline{{\bf w}}_{t+1} is bounded by

For Scheme II, assuming p1=p2=⋯=pN=1Np_{1}=p_{2}=\cdots=p_{N}=\frac{1}{N}, the expected difference between v‾t+1\overline{{\bf v}}_{t+1} and w‾t+1\overline{{\bf w}}_{t+1} is bounded by

B.3 Completing the Proof of Theorem 2 and 3

When expectation is taken over St+1{\mathcal{S}}_{t+1}, the last term (A3A_{3}) vanishes due to the unbiasedness of w‾t+1\overline{{\bf w}}_{t+1}.

If t+1∉IEt+1\notin{\mathcal{I}}_{E}, A1A_{1} vanishes since w‾t+1=v‾t+1\overline{{\bf w}}_{t+1}=\overline{{\bf v}}_{t+1}. We use Lemma 5 to bound A2A_{2}. Then it follows that

If t+1∈IEt+1\in{\mathcal{I}}_{E}, we additionally use Lemma 5 to bound A1A_{1}. Then

Then by the strong convexity of F(⋅)F(\cdot),

Specifically, if we choose β=2μ,γ=max⁡{8Lμ,E}−1\beta=\frac{2}{\mu},\gamma=\max\{8\frac{L}{\mu},E\}-1 and denote κ=Lμ\kappa=\frac{L}{\mu}, then ηt=2μ1γ+t\eta_{t}=\frac{2}{\mu}\frac{1}{\gamma+t} and

B.4 Deferred proofs of key lemmas

We first give a key observation which is useful to prove the followings. Let {xi}i=1N\{x_{i}\}_{i=1}^{N} denote any fixed deterministic sequence. We sample a multiset St{\mathcal{S}}_{t} (with size KK) by the procedure where for each sampling time, we sample xkx_{k} with probability qkq_{k} for each time. Pay attention that two samples are not necessarily independent. We only require each sampling distribution is identically. Let St={i1,⋯ ,iK}⊂[N]{\mathcal{S}}_{t}=\{i_{1},\cdots,i_{K}\}\subset[N] (some iki_{k}’s may have the same value). Then

For Scheme I, qk=pkq_{k}=p_{k} and for Scheme II, qk=1Nq_{k}=\frac{1}{N}. It is easy to prove this lemma when equipped with this observation. ∎

We separately prove the bounded variance for two schemes. Let St+1={i1,⋯ ,iK}{\mathcal{S}}_{t+1}=\{i_{1},\cdots,i_{K}\} denote the multiset of chosen indexes.

(1) For Scheme I, w‾t+1=1K∑l=1Kvt+1il\overline{{\bf w}}_{t+1}=\frac{1}{K}\sum_{l=1}^{K}{\bf v}_{t+1}^{i_{l}}. Taking expectation over St+1{\mathcal{S}}_{t+1}, we have

where the first equality follows from vt+1il{\bf v}_{t+1}^{i_{l}} are independent and unbiased.

To bound eqn. (26), we use the same argument in Lemma 5. Since t+1∈IEt+1\in{\mathcal{I}}_{E}, we know that the time t0=t−E+1∈IEt_{0}=t-E+1\in{\mathcal{I}}_{E} is the communication time, which implies {wt0k}k=1N\{{\bf w}_{t_{0}}^{k}\}_{k=1}^{N} is identical. Then

where in the last inequality we use the fact that ηt\eta_{t} is non-increasing and ηt0≤2ηt\eta_{t_{0}}\leq 2\eta_{t}.

(2) For Scheme II, when assuming p1=p2=⋯=pN=1Np_{1}=p_{2}=\cdots=p_{N}=\frac{1}{N}, we again have w‾t+1=1K∑l=1Kvt+1il\overline{{\bf w}}_{t+1}=\frac{1}{K}\sum_{l=1}^{K}{\bf v}_{t+1}^{i_{l}}.

where in the last inequality we use the same argument in (1). ∎

Appendix C The empirical risk minimization example in Section 4

Let p>1p>1 be a positive integer. To avoid the trivial case, we assume N>1N>1. Consider the following quadratic optimization

where i,ji,j are row and column indices, respectively. We partition A{\bf A} into a sum of NN symmetric matrices (A=∑k=1NAk{\bf A}=\sum_{k=1}^{N}{\bf A}_{k}) and b{\bf b} into b=∑k=1Nbk{\bf b}=\sum_{k=1}^{N}{\bf b}_{k}. Specifically, we choose b1=b=e1{\bf b}_{1}={\bf b}={\bf e}_{1} and b2=⋯=bN=0{\bf b}_{2}=\cdots={\bf b}_{N}=0. To give the formulation of Ak{\bf A}_{k}’s, we first introduce a series of sparse and symmetric matrices Bk (1≤k≤N){\bf B}_{k}\ (1\leq k\leq N):

Now Ak{\bf A}_{k}’s are given by A1=B1+E1,1,Ak=Bk (2≤k≤N−1){\bf A}_{1}={\bf B}_{1}+{\bf E}_{1,1},{\bf A}_{k}={\bf B}_{k}\ (2\leq k\leq N-1) and AN=BN+ENp+1,Np+1{\bf A}_{N}={\bf B}_{N}+{\bf E}_{Np+1,Np+1}, where Ei,j{\bf E}_{i,j} is the matrix where only the (i,j)(i,j)th entry is one and the rest are zero.

Back to the federated setting, we distribute the kk-th partition (Ak,bk)({\bf A}_{k},{\bf b}_{k}) to the kk-th device and construct its corresponding local objective by

In the next subsection (Appendix C.3), we show that the quadratic minimization with the global objective (27) and the local objectives (30) is actually a distributed linear regression. In this example, training data are not identically but balanced distributed. Moreover, data in each device are sparse in the sense that non-zero features only occur in one block. The following theorem (Theorem 5) shows that FedAvg might converge to sub-optimal points even if the learning rate is small enough. We provide a numerical illustration in Appendix C.2 and a mathematical proof in Appendix C.4.

In the above problem of the distributed linear regression, assume that each device computes exact gradients (which are not stochastic). With a constant and small enough learning rate η\eta and E>1E>1, FedAvg converges to a sub-optimal solution, whereas FedAvg with E=1E=1 (i.e., gradient descent) converges to the optimum. Specifically, in a quantitative way, we have

where w~∗\widetilde{{\bf w}}^{*} is the solution produced by FedAvg and w∗{\bf w}^{*} is the optimal solution.

C.2 Numerical illustration on the example

We conduct a few numerical experiments to illustrate the poor performance of FedAvg on the example introduced in Section 4. Here we set N=5,p=4,μ=2×10−4N=5,p=4,\mu=2\times 10^{-4}. The annealing scheme of learning rates is given by ηt=1/55+t⋅a\eta_{t}=\frac{1/5}{5+t\cdot a} where aa is the best parameter chosen from the set {10−2,10−4,10−6}\{10^{-2},10^{-4},10^{-6}\}.

C.3 Some properties of the example

which implies that 0≺A⪯4I0\prec{\bf A}\preceq 4{\bf I}.

The sparse and symmetric matrices Bk (1≤k≤N){\bf B}_{k}\ (1\leq k\leq N) defined in eqn. (29) can be rewritten as

From theory of linear algebra, it is easy to follow this proposition.

By the way of construction, Ak{\bf A}_{k}’s have following properties:

Ak{\bf A}_{k} is positive semidefinite with ∥Ak∥2≤4\|{\bf A}_{k}\|_{2}\leq 4;

From Proposition 1, we can rewrite these local quadratic objectives in form of a ridge linear regression. Specifically, for k=1k=1,

where CC is some constant irrelevant with w{\bf w}). For 2≤k≤N2\leq k\leq N,

Similarly, the global quadratic objective eqn. (27) can be written as F(w)=12N∥X(w−w∗)∥22+12μ∥w∥2F({\bf w})=\frac{1}{2N}\|{\bf X}({\bf w}-{\bf w}^{*})\|_{2}^{2}+\frac{1}{2}\mu\|{\bf w}\|^{2} .

Data in each device are sparse in the sense that non-zero features only occur in the block Ik{\mathcal{I}}_{k} of coordinates. Blocks on neighboring devices only overlap one coordinate, i.e., ∣Ik∩Ik+1∣=1|{\mathcal{I}}_{k}\cap{\mathcal{I}}_{k+1}|=1. These observations imply that the training data in this example is not identically distributed.

The kk-th device has rk (=p or p+1)r_{k}\ (=p\ \text{or}\ p+1) non-zero feature vectors which are vertically concatenated into the feature matrix Xk{\bf X}_{k}. Without loss of generality, we can assume all devices hold p+1p+1 data points since we can always add additional zero vectors to expand the local dataset. Therefore n1=⋯=nN=p+1n_{1}=\cdots=n_{N}=p+1 in this case, which implies that the training data in this example is balanced distributed.

C.4 Proof of Theorem 5.

To prove the theorem, we assume that (i) all devices hold the same amount of data points, (ii) all devices perform local updates in parallel, (iii) all workers use the same learning rate η\eta and (iv) all gradients computed by each device make use of its full local dataset (hence this case is a deterministic optimization problem). We first provide the result when μ=0\mu=0.

For convenience, we slightly abuse the notation such that wt{\bf w}_{t} is the global parameter at round tt rather than step tt. Let wt(k){\bf w}_{t}^{(k)} the updated local parameter at kk-th worker at round tt. Once the first worker that holds data (A1,b1)({\bf A}_{1},{\bf b}_{1}) runs EE step of SGD on F1(w)F_{1}({\bf w}) from wt{\bf w}_{t}, it follows that

For the rest of workers, we have wt(k)=(I−ηAi)Ewt (2≤k≤N){\bf w}_{t}^{(k)}=({\bf I}-\eta{\bf A}_{i})^{E}{\bf w}_{t}\ (2\leq k\leq N).

since 0≺A⪯4I0\prec{\bf A}\preceq 4{\bf I} means 0⪯(I−ηNA)≺I0\preceq({\bf I}-\frac{\eta}{N}{\bf A})\prec{\bf I}.

Then ∥wt+1−wt∥2≤ρ∥wt−wt−1∥2≤ρt∥w1−w0∥2\|{\bf w}_{t+1}-{\bf w}_{t}\|_{2}\leq\rho\|{\bf w}_{t}-{\bf w}_{t-1}\|_{2}\leq\rho^{t}\|{\bf w}_{1}-{\bf w}_{0}\|_{2}. By the triangle inequality,

which implies that {wt}t≥1\{{\bf w}_{t}\}_{t\geq 1} is a Cauchy sequence and thus has a limit denoted by w~∗\widetilde{{\bf w}}^{*}. We have

When E=1E=1, it follows from eqn. (32) that w~∗=A−1b=w∗\widetilde{{\bf w}}^{*}={\bf A}^{-1}{\bf b}={\bf w}^{*}, i.e., FedAvg converges to the global minimizer.

The right hand side of the last equation cannot be zero. Quantificationally speaking, we have the following lemma. We defer the proof for the next subsection.

If the step size η\eta is sufficiently small, then in this example, we have

Since A1A2≠0{\bf A}_{1}{\bf A}_{2}\neq{\bf 0} and w∗{\bf w}^{*} is dense, the lower bound in eqn. (34) is not vacuous.

Now we have proved the result when μ=0\mu=0. For the case where μ>0\mu>0, we replace Ai{\bf A}_{i} with Ai+μI{\bf A}_{i}+\mu{\bf I} and assume μ<14+μ\mu<\frac{1}{4+\mu} instead of the original. The discussion on different choice of EE is unaffected. ∎

C.5 Proof of Lemma 6

We will derive the conclusion mainly from the expression eqn. (33). Let f(η)f(\eta) be a function of η\eta. We say a matrix T{\bf T} is Θ(f(η))\Theta(f(\eta)) if and only if there exist some positive constants namely C1C_{1} and C2C_{2} such that C1f(η)≤∥T∥≤C2f(η)C_{1}f(\eta)\leq\|{\bf T}\|\leq C_{2}f(\eta) for all η>0\eta>0. In the following analysis, we all consider the regime where η\eta is sufficiently small.

Denote by V=∑i=1NAi2{\bf V}=\sum_{i=1}^{N}{\bf A}_{i}^{2}. First we have

Then by plugging this equation into the right hand part of eqn. (33), we have

Plugging the last two equations into eqn. (33), we have

where the last inequality holds because (i) we require η\eta to be sufficiently small and (ii) ∥A−1x∥≥14∥x∥\|{\bf A}^{-1}{\bf x}\|\geq\frac{1}{4}\|{\bf x}\| for any vector x{\bf x} as a result of 0<∥A∥≤40<\|{\bf A}\|\leq 4. The last equality uses the fact (i) V−A1A=A1∑i=2nAi{\bf V}-{\bf A}_{1}{\bf A}={\bf A}_{1}\sum_{i=2}^{n}{\bf A}_{i} and (ii) A1Ai=0{\bf A}_{1}{\bf A}_{i}={\bf 0} for any i≥3i\geq 3. ∎

Appendix D Experimental Details

We examine our theoretical results on a multinomial logistic regression. Specifically, let f(w;xi)f({\bf w};x_{i}) denote the prediction model with the parameter w=(W,b){\bf w}=({\bf W},{\bf b}) and the form f(w;xi)=softmax(Wxi+b)f({\bf w};{\bf x}_{i})=\text{softmax}({\bf W}{\bf x}_{i}+{\bf b}). The loss function is given by

This is a convex optimization problem. The regularization parameter is set to λ=10−4\lambda=10^{-4}.

Datasets.

We evaluate our theoretical results on both real data and synthetic data. For real data, we choose MNIST dataset [LeCun et al., 1998] because of its wide academic use. To impose statistical heterogeneity, we distribute the data among N=100N=100 devices such that each device contains samples of only two digits. To explore the effect of data unbalance, we further vary the number of samples among devices. Specifically, for unbalanced cases, the number of samples among devices follows a power law, while for balanced cases, we force all devices to have the same amount of samples.

We summarize the information of federated datasets in Table 2.

Experiments.

For all experiments, we initialize all runnings with w0=0{\bf w}_{0}=0. In each round, all selected devices run EE steps of SGD in parallel. We decay the learning rate at the end of each round by the following scheme ηt=η01+t\eta_{t}=\frac{\eta_{0}}{1+t}, where η0\eta_{0} is chosen from the set {1,0.1,0.01}\{1,0.1,0.01\}. We evaluate the averaged model after each global synchronization on the corresponding global objective. For fair comparison, we control all randomness in experiments so that the set of activated devices is the same across all different algorithms on one configuration.

D.2 Theoretical verification

From our theory, when the total steps TT is sufficiently large, the required number of communication rounds to achieve a certain precision is

which is s a function of EE that first decreases and then increases. This implies that the optimal local step E∗E^{*} exists. What’s more, the Tϵ/ET_{\epsilon}/E evaluated at E∗E^{*} is

which implies that FedAvg needs more communication rounds to tackle with severer heterogeneity.

To validate these observations, we test FedAvg with Scheme I on our four datasets as listed in Table 2. In each round, we activate K=30K=30 devices and set η0=0.1\eta_{0}=0.1 for all experiments in this part. For unbalanced MNIST, we use batch size b=64b=64. The target loss value is 0.290.29 and the minimum loss value found is 0.25910.2591. For balanced MNIST, we also use batch size b=64b=64. The target loss value is 0.500.50 and the minimum loss value found is 0.34290.3429. For two synthetic datasets, we choose b=24b=24. The target loss value for synthetic(0,0) is 0.950.95 and the minimum loss value is 0.79990.7999. Those for synthetic(1,1) are 1.151.15 and 1.0751.075.

The impact of K𝐾K.

Our theory suggests that a larger KK may accelerate convergence since Tϵ/ET_{\epsilon}/E contains a term O(EG2K)\mathcal{O}\left(\frac{EG^{2}}{K}\right). We fix E=5E=5 and η0=0.1\eta_{0}=0.1 for all experiments in this part. We set the batch size to 64 for two MNIST datasets and 24 for two synthetic datasets. We test Scheme I for illustration. Our results show that, no matter what value KK is, FedAvg converges. From Figure 3, all the curves in each subfigure overlap a lot. To show more clearly the differences between the curves, we zoom in the last few rounds in the upper left corner of the figure. It reveals that the curve of a large enough KK is slightly better. This result also shows that there is no need to sample as many devices as possible in convex federated optimization.

Sampling and averaging schemes.

We analyze the influence of sampling and averaging schemes. As stated in Section 3.3, Scheme I iid samples (with replacement) KK indices with weights pkp_{k} and simply averages the models, which is proposed by Sahu et al. . Scheme II uniformly samples (without replacement) KK devices and weightedly averages the models with scaling factor N/KN/K. Transformed Scheme II scales each local objective and uses uniform sampling and simple averaging. We compare Scheme I, Scheme II and transformed Scheme II, as well as the original scheme [McMahan et al., 2017] on four datasets. We carefully tuned the learning rate for the original scheme. In particular, we choose the best step size from the set {0.1,0.5,0.9,1.1}\{0.1,0.5,0.9,1.1\}. We did not fine tune the rest schemes and set η0=0.1\eta_{0}=0.1 by default. The hyperparameters are the same for all schemes: E=20,K=10E=20,K=10 and b=64b=64. The results are shown in Figure 1c and 1d.

Our theory does not guarantee FedAvg with Scheme II could converge when the training data are unbalanced distributed. Actually, if the number of training samples varies too much among devices, Scheme II may even diverge. To illustrate this point, we have shown the terrible performance on mnist unbalanced dataset in Figure 1b. In Figure 4, we show additional results of Scheme II on the two synthetic datasets, which are the most unbalanced. We choose b=24,K=10,E=10b=24,K=10,E=10 and η0=0.1\eta_{0}=0.1 for these experiments. However, transformed Scheme II performs well except that it has a lower convergence rate than Scheme I.