FedPD: A Federated Learning Framework with Optimal Rates and Adaptivity to Non-IID Data

Xinwei Zhang, Mingyi Hong, Sairaj Dhople, Wotao Yin, Yang Liu

Introduction

Federated learning (FL)—a distributed machine learning approach proposed in —has gained popularity for applications involving learning from distributed data. In FL, a cloud server (the “server”) can communicate with distributed data sources (the “agents”). The goal is to train a global model that works well for all the distributed data, but without requiring the agents to reveal too much local information. Since its inception, the broad consensus on FL’s implementation appears to involve a generic “computation then aggregation” (CTA) protocol. This involves the following steps: S1) the server sends the model x\mathbf{x} to the agents; S2) the agents update their local models xi\mathbf{x}_{i}’s for several iterations based on their local data; S3) the server aggregates xi\mathbf{x}_{i}’s to obtain a new global model x\mathbf{x}. It is widely acknowledged that multiple local steps save communication efforts, while only transmitting local models protects data privacy .

Even though the FL paradigm has attracted significant research from both academia and industry, and many algorithms such as Federated Averaging (FedAvg) have been proposed, several attributes are not clearly established. In particular, the commonly adopted CTA protocol poses significant theoretical and practical challenges to designing effective FL algorithms. This work attempts to provide a deeper understanding of FL, by raising and resolving a few theoretical questions, as well as by developing an effective algorithmic framework with several desirable features.

Problem Formulation. Consider the vanilla FL solving the following problem:

where Pi\mathcal{P}_{i} denotes the data distribution on the ii-th agent. Throughout the paper, we will make the following blanket assumptions for problem (1):

Each fi(⋅)f_{i}(\cdot), as well as ff in (1) is LL–smooth:

In addition to these standard assumptions, state-of-the-art efforts on analysis of FL algorithms oftentimes invoke a number of more restrictive assumptions.

(Bounded Gradient (BG)) The gradients ∇fi\nabla f_{i}’s are upper bounded (by a constant G>0G>0)

(I.I.D. Data) Either one of the following holds:

First, the BG assumption typically does not hold for (1), in particular, fi(x)=∥Aix−bi∥2f_{i}(\mathbf{x})=\|\mathbf{A}_{i}\mathbf{x}-\mathbf{b}_{i}\|^{2} (where Ai\mathbf{A}_{i} and bi\mathbf{b}_{i} are related to data). However, the BG assumption is critical for analyzing FedAvg-type algorithms because it bounds the distance traveled after multiple local iterates. Second, (4) is typically used in FL to characterize homogeneity about local data . However, an assumption of this type does not hold for FL applications where the data (such as medical records, keyboard input data) are generated by the individual agents . A reasonable relaxation to this i.i.d. assumption is the following notion of δ\delta-non-i.i.d.-ness of the data distribution.

(δ\delta-Non-I.I.D. Data) The local functions are called δ\delta-non-i.i.d. if either one of the equivalent conditions below holds:

By varying δ\delta from to ∞\infty, (6) provides a characterization of data non-i.i.d.-ness. In Appendix A, we give a few examples of loss functions with different values of δ\delta. Note that the second inequality in (6) is often used in decentralized optimization to quantify the similarity of local problems . Third, (5) does not hold for many practical problems. To see this, note that this condition is parameterized by ϵ\epsilon, which is typically the desired optimization accuracy . Since ϵ\epsilon can be chosen arbitrarily small, (5) essentially requires that the problem is realizable, that is, ∥∇f(x)∥\|\nabla f(\mathbf{x})\| approaches zero only when all the local gradients approach zero at x\mathbf{x}, that is, when the local data are “similar”.

Finally, we mention that our objective is to understand FL algorithm from an optimization perspective. So we say that a solution x\mathbf{x} is an ϵ\epsilon-stationary solution if the following holds:

We are interested in finding the minimum system resources required, such as the number of local updates, the number of times local variables are transmitted to the server, and the number of times local samples F(x;ξi)F(\mathbf{x};\xi_{i})’s are accessed, before computing an ϵ\epsilon-solution (7). These quantities are referred to as local computation, communication complexity, and sample complexities, respectively.

Questions to address Despite extensive recent research, the FL framework, and in particular, the CTA protocol is not yet well understood. Below, we list four questions regarding the CTA protocol.

Q1 (local updates). What are the best local update directions for the agents to take so as to achieve the best overall system performance (stability, sample complexity, etc.)?

Q2 (global aggregation). Can we use more sophisticated processing in the aggregation step to help improve the system performance (sample or communication complexity)?

Q3 (communication efficiency). If multiple local updates are preformed between two aggregation steps, will it reduce the communication overhead?

Q4 (assumptions). What is the best performance that the CTA type algorithms can achieve while relying on a minimum set of assumptions about the problem?

Although these questions are not directly related to data privacy, another important aspect of FL, we argue that answering these fundamental questions can provide much needed understanding on algorithms following the CTA, and thus the FL approach. A few recent works (to be reviewed shortly) have touched upon those questions, but to our knowledge, none of them has conducted a thorough investigation of the questions listed above.

Related Works. We start with a popular method following the CTA protocol, the FedAvg in Algorithm 1, which covers the original FedAvg , the Local SGD , PR-SGD and the RI-SGD among others.

In FedAvg, TT is the total stage number, QQ the number of local updates, rr the index of the stage, qq the index of the inner iteration, and ηr,q\eta^{r,q}’s are the stepsizes. It has two options for local updates:

Many recent works are extensions of FedAvg. The algorithm proposed in adds momentum to the algorithm. In , the data on the local agents are separated into blocks and shared with other agents. In the local GD version (9) is studied. In , a cooperative-SGD is considered; it includes virtual agents, extra variables, and relaxes the parameter server topology.

It is pertinent to consider how these algorithms address questions Q1–Q4. For Q1, most FedAvg-type algorithms perform multiple local (stochastic) GD steps to minimize the local loss function. However, we will see shortly that in many cases, successive local GD steps lead to algorithm divergence. For Q2, most algorithms use simple averaging, and there is little discussion on whether other types of (linear) processing will be helpful. For Q3, a number of recent works such as show that, for non-convex problems, to achieve ϵ\epsilon-solution (7), a total of O(1/ϵ3/2)O(1/\epsilon^{3/2}) aggregation steps are needed. However, it is not clear if this achieves the best communication complexity. As for Q4, the FedAvg-type algorithm typically requires either some variant of the BG assumption, or some i.i.d. assumption, or both; See Table 1 rows 1–5 for detailed discussions.

A number of more recent works have improved upon FedAvg in various aspects. FedProx addresses Q1 and Q4 by perturbing the update direction. This algorithm does not need the BG, but it still requires the i.i.d. assumption (5). The VRL-SGD proposed in addresses Q1 and Q4 by using the variance reduction (VR) technique to update the directions for local agents and achieves O(1/ϵ)\mathcal{O}(1/\epsilon) communication complexity without the i.i.d. assumption. F-SVRG is another recent algorithm that uses VR. This algorithm does not follow the CTA protocol as the agents have to transmit the local gradients, but it does not require A3 and A4. The PR-SPIDER further improves upon FSVRG by reducing sample complexity (SC) from O(M/ϵ)\mathcal{O}(M/\epsilon) to O(M/ϵ)\mathcal{O}(\sqrt{M}/\epsilon) (where MM is typically larger than 1/ϵ1/\epsilon). Although FSVRG and PR-SPIDER neither require the BG or the i.i.d. assumptions, they require the agents to transmit local gradients to the server and thus do not follow the CTA protocol. This is undesirable, as it has been shown that local gradient information can leak private data . Additionally, questions Q2-Q3 are not addressed in these works.

Our Main Contributions. First, we address Q1-Q4 and provide an in-depth examination of the CTA protocol. We show that algorithms following the CTA protocol that are based on successive local gradient updates, the best possible communication efficiency is O(1/ϵ)\mathcal{O}(1/\epsilon); neither additional local processing nor general linear processing can improve this rate. We then show that the BG and/or i.i.d. data assumption is important for the popular FedAvg to work as intended.

Our investigation suggests that the existing FedAvg-based algorithms are (provably) insufficient in dealing with many practical problems, calling for a new design strategy. We then propose a meta-algorithm called Federated Primal-Dual (FedPD), which also follows the CTA protocol and can be implemented in several different forms with desirable properties. In particular, it i) can deal with the general non-convex problem, ii) achieve the best possible optimization and communication complexity when data is non-i.i.d., iii) achieve convergence under only Assumptions A1– A2.

Most importantly, the communication pattern of the proposed algorithm can be adapted to the degree of non-i.i.d.-ness of the local data. That is, we show that under the δ\delta-non-i.i.d. condition (6), communication saving and data heterogeneity interestingly exhibit a linear-logarithmic relationship; see Fig. 1 for an illustration. To our knowledge, this is the first algorithm for FL that achieves all the above properties.

Addressing Open Questions

We focus on the case where the Vt(⋅)V^{t}(\cdot)’s and Wit(⋅)W^{t}_{i}(\cdot)’s are linear operators, which implies that xit,qx^{t,q}_{i} can use all past iterates and (sample) gradients for its update. Clearly, (10) covers both the local-GD and local-SGD versions of FedAvg as special cases. In the following, we provide an informal statement of the result. The formal statement and the full proof are given in Appendix B and Theorem 2.

(Informal) Consider any algorithm AA that belongs to the class described in (10), with Vt(⋅)V^{t}(\cdot) and Wit(⋅)W^{t}_{i}(\cdot)’s being linear and possibly time-varying operators. Then, there exists a non-convex problem instance satisfying Assumptions 1–2 such that for any Q>0Q>0, algorithm AA takes at least O(1/ϵ)\mathcal{O}(1/\epsilon) communication rounds to reach an ϵ\epsilon-stationary solution satisfying (7).

Remark 1. The proof technique is related to those developed from both classical and recent works that characterize lower bounds for first-order methods, in both centralized and decentralized settings. The main technical difference is that our processing model (10) additionally allows local processing iterations, and there is a central aggregator. Our goal is not to establish lower bounds on the number of (centralized) gradient access, nor to show the optimal graph dependency, but to characterize (potential) communication savings when allowing multiple steps of local processing.

In the proof, we construct difficult problem instances in which fif_{i}’s are not i.i.d. (more precisely, δ\delta in assumption (4) grows with the total number of iterations TT). Then we show that it is necessary to aggregate (thus communicate) to make any progress. On the other hand, it is obvious that in another extreme case where the data are -non-i.i.d., only O(1)\mathcal{O}(1) communication rounds are needed. An open question is: when the local data are related to each other, i.e., δ\delta lies between and infinity, is it possible to reduce the total communication rounds? This question is addressed below in Sec. 3. ■\blacksquare

We now address Q1 and Q4. We consider the FedAvg Algorithm 1, and show that BN and/or i.i.d. assumptions are critical for them to perform well. Our result suggests that, despite its popularity, components in FedAvg, such as the pure local (stochastic) gradient directions and linear aggregation are not compatible with each other. The proof of the results below are given in Appendix C.

Fix any constant η>0\eta>0, Q>1Q>1 for Algorithm 1. There exists a problem that satisfies A1 and A2 but fails to satisfy A3 and A4, on which FedAvg diverges to infinity.

Remark 2. A recent work has shown that FedAvg with constant stepsize η>0\eta>0 can only converge to a neighborhood of the global minimizer for convex problems. Moreover, the error to the global optima is related to QQ and the degree of non-i.i.d.-ness as measured by the size of ∑i=1N∥∇fi(x⋆)∥2\sum^{N}_{i=1}\left\lVert\nabla f_{i}(\mathbf{x}^{\star})\right\rVert^{2} where x⋆\mathbf{x}^{\star} is the global optimal solution. On the other hand, our result indicates that when fif_{i}’s are non-convex, FedAvg can perform much worse without the BN and the i.i.d. assumption. Even if Q=2Q=2 and there exists a solution such that ∑i=1N∥fi(x^)∥2=0\sum_{i=1}^{N}\|f_{i}(\hat{\mathbf{x}})\|^{2}=0, FedAvg (with constant stepsize η\eta) diverges and the iteration can go to ∞\infty. ■\blacksquare

One may think that using a constant stepsize is the culprit for the divergence in Claim 2.2. In fact, we can show that having BG or not can still impact the performance of FedAvg, even when diminishing stepsize is used. In particular, we show in Appendix D, that FedAvg converges under the BG assumption for any diminishing stepsize, but without it, the choice of the stepsize can be significantly restricted.

The FedPD Framework

The previous section reveals a number of properties about FedAvg and, broadly speaking, the CTA protocol. But why does FedAvg only work under very restrictive conditions? Is it because the local update directions are not chosen correctly? Is it possible to make it work without any additional assumptions? Can we reduce communication effort when the local data becomes more i.i.d.?

In this section, we propose a meta-algorithm called Federated Primal-Dual (FedPD), which can be specialized into different sub-variants to address the above questions. These algorithms possess a few desirable features: They can achieve the best optimization and communication complexity when data is non-i.i.d. (i.e., achieving the bounds in Claim 2.1); they only require A1 –A2, while being able to utilize both full or sampled local gradients. Most importantly, the communication pattern of the proposed algorithm can be made adaptive to the degree of data non-i.i.d.-ness across the agents.

Our algorithm is based on the following global consensus reformulation of the original problem (1):

Similar to traditional primal-dual based algorithms , the idea is that, when relaxing the equality constraints, the resulting problem is separable across different nodes. However, different from ADMM, the agents can now perform either a single (or multiple) local update(s) between two communication rounds. Importantly, such flexibility makes it possible to adapt the communication frequency to the degree of δ\delta-non-i.i.d.-ness of the local data. In particular, we identify that under δ\delta-non-i.i.d. (4), the fraction ϵ/δ2\epsilon/\delta^{2} is the key quantity that determines communication saving; see Fig. 1. Intuitively, significant reduction can be achieved when δ\delta is smaller than ϵ\epsilon; otherwise, the reduction goes to zero linearly as δ\delta increases. To our knowledge, none of the existing ADMM based algorithms, nor FL based algorithms, are able to provably achieve such a reduction.

To present our algorithm, let us define the augmented Lagrangian (AL) function of (11) as

Fixing x0\mathbf{x}_{0}, the AL is separable over all local pairs {(xi,λi)}\{(\mathbf{x}_{i},\lambda_{i})\}. The key technique in the design is to specify how each local AL Li(⋅)\mathcal{L}_{i}(\cdot) should be optimized, and when to perform model aggregation.

Federated primal-dual algorithm (FedPD) captures the main idea of the classical primal-dual based algorithm while meeting the flexibility need of FL; see Algorithm 2. In particular, its update rules share a similar pattern as ADMM, but it does not specify how the local models are updated. Instead, an oracle \mboxOraclei(⋅)\mbox{\rm Oracle}_{i}(\cdot) is used as a placeholder for local processing, and we will see that careful instantiations of these oracles lead to algorithms with different properties. Importantly, we introduce a critical constant p∈[0,1)p\in[0,1), which determines the frequency at which the aggregation and communication steps are skipped. In Algorithm 3 and Algorithm 4, we provide two useful examples of the local oracles.

In Algorithm 3, QiQ_{i}’s are chosen so that the local problems are solved accurately enough to satisfy:

We provide two ways for solving this subproblem by using GD and SGD, but any other solver that achieves (12) can be used. For the SGD version, the stochastic gradient is defined as

where ∼\sim denotes uniform sampling. Despite the simplicity of the local updates, we will show that using Oracle I makes FedPD adaptive to the non-i.i.d. parameter δ\delta.

In Algorithm 4, the oracle applies the variance reduction technique to reduce the sample complexity. The detailed descriptions and the analyses are given in Appendix G due to the space limitation.

2 Convergence and Complexity Analysis

We analyze the convergence of FedPD with Oracle I. The detailed proofs are given in Appendix F. The convergence analysis of FedPD with an alternative Oracle is given in Appendix G

Suppose A1 –A2 hold. Consider FedPD with Oracle I, where QiQ_{i} are selected by (12).

Case I) Suppose A5 holds with δ=∞\delta=\infty. Set 0<η<5−14L0<\eta<\frac{\sqrt{5}-1}{4L}, p=0p=0. Then we have:

Case II) Suppose 0<η<5−14L0<\eta<\frac{\sqrt{5}-1}{4L}, 0≤p<10\leq p<1, and A5 holds with a finite δ\delta. Then we have:

Here C2,C4,C5>0C_{2},C_{4},C_{5}>0 are constants independent of T,δ,pT,\delta,p; C_{3}:={\color[rgb]{0,0,0}\frac{{p(1+L\eta)+L\eta}}{1-L\eta}}\geq 0.

Remark 3. (Communication complexity) Case I says if one does not skip communication (p=0p=0), then to achieve ϵ\epsilon-stationarity (i.e., ∥∇f(x0t)∥2≤ϵ\left\lVert\nabla f(\mathbf{x}^{t}_{0})\right\rVert^{2}\leq\epsilon for some r∈(1,T)r\in(1,T)), they need to set T=1/(2C2D0ϵ)T=1/(2C_{2}D_{0}\epsilon), ϵ1=ϵ/(2C4)\epsilon_{1}=\epsilon/(2C_{4}), and the total communication rounds is TT.

In Case II, the second term on the right hand side of (14) involves both δ\delta and pp, so we can select them appropriately to reduce communication while maintaining the same accuracy. Specifically, we choose ϵ1=min⁡{ϵ/(4C4),δ2}\epsilon_{1}=\min\{\epsilon/(4C_{4}),\delta^{2}\} and T=1/(2C2D0ϵ)T=1/(2C_{2}D_{0}\epsilon). Then we have

Then we can achieve the same ϵ\epsilon-accuracy as Case I. The communication rounds here is T(1−p)T(1-p).

The relation between the saving pp and ϵδ2\frac{\epsilon}{\delta^{2}} is showed in Table 2. Roughly speaking, pp is inversely proportional to degree of non-i.i.d-ness δ2\delta^{2} when δ2∈(O(ϵ),∞)\delta^{2}\in(\mathcal{O}(\epsilon),\infty); further, p→1p\rightarrow 1 at a log-rate when δ2→0\delta^{2}\rightarrow 0. Our result also indicates that, when using stage-wise training for neural networks, the algorithm can communicate less at the early stages since they typically have lower accuracy target to enable larger stepsizes .

Remark 4. (Computation complexity) To achieve ϵ\epsilon accuracy, we need both T=O(1/ϵ)T=\mathcal{O}(1/\epsilon) and ϵ1=O(ϵ)\epsilon_{1}=\mathcal{O}(\epsilon). As the local AL is strongly convex with respect to xi\mathbf{x}_{i}, optimizing it to ϵ\epsilon accuracy requires O(log⁡(ϵ))\mathcal{O}(\log(\epsilon)) iterations for GD and O(1/ϵ)\mathcal{O}(1/\epsilon) for SGD . So the total number of times that the local gradients (respectively, stochastic gradients) are accessed is given by O(1/ϵ×log⁡(1/ϵ))\mathcal{\mathcal{O}(\text{1}/\epsilon\times\log(\text{1}/\epsilon))} (respectively, O(1/ϵ2)\mathcal{O}(1/\epsilon^{2})).

We conclude this section by noting that the above communication and computation complexity results we have obtained are the best so far among all FL algorithms for non-convex problems satisfying A1 – A2. Please see the last three rows of Table 1 for a summary of the results.

3 Connection with Other Algorithms

Before we close this section, we discuss the relation of FedPD with a few existing algorithms.

The FedProx/FedDANE In FedProx the agents optimize the following local objective: fi(xi)+ρ2∥xi−x0r∥2f_{i}(\mathbf{x}_{i})+\frac{\rho}{2}\left\lVert\mathbf{x}_{i}-\mathbf{x}^{r}_{0}\right\rVert^{2}. This fails to converge to the global stationary solution. In contrast, FedPD introduces extra local dual variables {λi}\{\lambda_{i}\} that record the gap between the local model xi\mathbf{x}_{i} and the global model x0\mathbf{x}_{0} which help the global convergence. FedDANE also proposes a way of designing the subproblem by using the global gradient, but this violates the CTA protocol. Compared with these two algorithms, the proposed FedPD has weaker assumptions, and it achieves better sample and/or communication complexity.

Event Triggering Algorithms. A number of recent works such as Lazily Aggregated Gradient (LAG) and COLA have been proposed to occasionally skip message exchanges among the agents to save communication. In LAG, each agent receives the global model every iteration, and decides whether to send some local gradient differences by checking certain conditions. Since gradient information is transmitted, LAG does not belong to the algorithm class (10). When the local problems are unbalanced, in the sense that the discrepancy between the local Lipschitz gradients LiL_{i}’s is large, then the agents with smaller LiL_{i}’s can benefit from the lazy aggregation. Meanwhile, instead of measuring whether the local problems are balanced, the δ\delta-non-i.i.d. criteria characterizes if local problems are similar by measuring the uniform difference between arbitrary pairs of the local problems. If the data is i.i.d., then the agents benefit equally from the communication reduction.

Numerical Experiments

In the first experiment, we show the convergence of the proposed algorithms on synthetic data with FedAvg and FedProx as baselines. We use the non-convex penalized logistic regression as the loss function. We use two ways to generate the dataset, in the first case (referred to as the “weakly non-i.i.d” case), the data is generated in an i.i.d. way. In the second case (referred to as the “strongly non-i.i.d” case), we generate the data using non-i.i.d. distribution. In both cases there are 400400 samples on each agent with total 100100 agents.

We run FedPD with Oracle I (FedPD-SGD and FedPD-GD) and Oracle II (FedPD-VR). For FedPD-SGD, we set Q=600Q=600, and for FedPD-GD and FedPD-VR we set Q=8Q=8. For FedPD-GD we set p=0p=0 and p=0.5p=0.5, wherein the later case the agents skips half of the communication rounds. For FedPD-VR, we set mini-batch size B=1B=1 and gradient computation frequency I=20I=20. For comparison, we also run FedAvg with local GD/SGD and FedProx. For FedAvg with GD, Q=8Q=8, and for FedAvg with SGD, Q=600Q=600. For FedProx, we solve the local problem using variance reduction for Q=8Q=8 iterations. The total number of iterations TT is set as 600600 for all algorithms.

Fig. 2 shows the results with respect to the number of communication rounds. In Fig. 2(a), we compare the convergence of the tested algorithms on weakly non-i.i.d. data set. It is clear that FedProx and FedPD with p=0p=0 (i.e., no communication skipping) are comparable. Meanwhile, FedAvg with local GD will not converge to the stationary point with a constant stepsize when local update step Q>1Q>1. By skipping half of the communication, FedPD with local GD can still achieve a similar error as FedAvg, but using fewer communication rounds. In Fig. 2(b), we compare the convergence results of different algorithms with the strongly non-i.i.d. data set. We can see that the algorithms using stochastic solvers become less stable compared with the case when the data sets are weakly non-i.i.d. Further, FedPD-VR and FedPD-GD with p=0p=0 are able still to converge to the global stationary point while FedProx will achieve a similar error as the FedAvg with local GD.

We included more details on the experimental results and additional experiments in Appendix H.

Conclusion

We study federated learning under the CTA protocol. We explore a number of theoretical properties of this protocol, and design a meta-algorithm called FedPD, which contains various algorithms with desirable properties, such achieving the best communication/computation complexity, as being able to adapt its communication pattern with data heterogeneity.

References

Appendix A Examples of Cost Functions Satisfy A5

In this part, we provide a commonly used function that satisfies A5.

Define the scalar exp⁡(bk−akTx)(1+exp⁡(bk−akTx))2\frac{\exp(b_{k}-\mathbf{a}^{T}_{k}\mathbf{x})}{(1+\exp(b_{k}-\mathbf{a}^{T}_{k}\mathbf{x}))^{2}} as v(ak,bk,x)v(\mathbf{a}_{k},b_{k},\mathbf{x}), we have v(ak,bk,x)∈(0,1),  ∀x,ak,bkv(\mathbf{a}_{k},b_{k},\mathbf{x})\in(0,1),\;\forall x,\mathbf{a}_{k},b_{k}. Further stack v(ak,bk,x)v(\mathbf{a}_{k},b_{k},\mathbf{x}) as v(Di,x)\mathbf{v}(\mathcal{D}_{i},\mathbf{x}), that is v(Di,x)=[v(a1,b1,x);…,;v(a∣Di∣,b∣Di∣,x)]\mathbf{v}(\mathcal{D}_{i},\mathbf{x})=[v(\mathbf{a}_{1},b_{1},\mathbf{x});\dots,;v(\mathbf{a}_{|\mathcal{D}_{i}|},b_{|\mathcal{D}_{i}|},\mathbf{x})]. Further we define AiA_{i} as the stacked matrix of all ak∈Di\mathbf{a}_{k}\in\mathcal{D}_{i} (i.e., Ai=[a1,…,a∣Di∣]A_{i}=[\mathbf{a}_{1},\dots,\mathbf{a}_{|\mathcal{D}_{i}|}]), then we can express ∇fi(x)\nabla f_{i}(\mathbf{x}) as

The difference between the gradients of fif_{i} and fjf_{j} is

As v(a,b,x)∈(0,1)v(\mathbf{a},b,\mathbf{x})\in(0,1), we know ∥v(Di,x)∥≤∥[1,…,1]∥=∣Di∣\left\lVert\mathbf{v}(\mathcal{D}_{i},\mathbf{x})\right\rVert\leq\left\lVert[1,\dots,1]\right\rVert=\sqrt{|\mathcal{D}_{i}|}, which implies:

Utilizing the above inequality in (18), we obtain:

So we can define δ=max⁡i,j{1∣Di∣∥Ai∥+1∣Dj∣∥Aj∥}\delta=\max_{i,j}\left\{\frac{1}{\sqrt{|\mathcal{D}_{i}|}}\left\lVert A_{i}\right\rVert+\frac{1}{\sqrt{|\mathcal{D}_{j}|}}\left\lVert A_{j}\right\rVert\right\} which is a finite constant. Note that the above analysis holds true for any DiD_{i} and x\mathbf{x}. Note that with finer analysis we can obtain better bounds for δ\delta.

Similar to logistic regression, we can also show that A5 holds for hyperbolic tangent function which is commonly used in neural network models. First, notice that the hyperbolic tangent is a rescaled version of logistic regression:

So, δ\delta for tanh⁡\tanh is 44 times that applicable to the logistic regression problem. Note that this analysis can further cover a wide range of neural network training problems that uses cross entropy loss and sigmoidal activation functions (e.g. MLP, CNN and RNN).

Then if the feature AiA_{i}’s satisfy AiTAi=AjTAj,∀ i≠j,A_{i}^{T}A_{i}=A_{j}^{T}A_{j},\forall~{}i\neq j, we have

Appendix B Proof of Claim 2.1

The proof is related to techniques developed in classical and recent works that characterize lower bounds for first-order methods in centralized and decentralized settings. Technically, our computational / communication model is different compared to the aforementioned works, since we allow arbitrary number of local processing iterations, and we have a central aggregator. The difference here is that our goal is not to show the lower bounds on the number of total (centralized) gradient access, nor to show the optimal graph dependency. The main point we would like to make is that there exist constructions of local functions fif_{i}’s such that no matter how many times that local first-order processing is performed, without communication and aggregation, no significant progress can be made in reducing the stationarity gap of the original problem.

For notational simplicity, we will assume that the full local gradients {∇fi(xik)}\{\nabla f_{i}(x^{k}_{i})\} can be evaluated. Later we will comment on how to extend this result to enable access to the sample gradients ∇F(xik;ξi)\nabla F(x^{k}_{i};\xi_{i}). In particular, we consider the following slightly simplified model for now:

In this section, we will call each tt a “stage,” and call each local iteration qq an “iteration.” We use xx to denote the variable located at the server. We use xix_{i} (and sometimes xqx_{q}) to denote the local variable at node ii, and use xi[j]x_{i}[j] and xi[k]x_{i}[k] to denote its jjth and kkth elements, respectively. We use gi(⋅)g_{i}(\cdot) and fi(⋅)f_{i}(\cdot) to denote some functions related to node ii, and g(⋅)g(\cdot) and f(⋅)f(\cdot) to denote the average functions of gig_{i}’s and fif_{i}’s, respectively. We use NN to denote the total number of nodes.

B.2 Main Constructions.

Suppose there are NN distributed nodes in the system, and they can all communicate with the server. To begin, we construct the following two non-convex functions

where we have defined the following functions

Clearly, each Θ(x,j)\Theta(x,j) is only related to two components in xx, i.e., x[j−1]x[j-1] and x[j]x[j].

By the above definition, the average function becomes:

See Fig. 3 for an illustration of the construction discussed above.

Further, for a given error constant ϵ>0\epsilon>0 and a given the Lipschitz constant LL, let us define

B.3 Properties.

First we present some properties of the component functions hih_{i}’s.

The functions Ψ\Psi and Φ\Phi satisfy the following:

For all w≤0w\leq 0, Ψ(w)=0\Psi(w)=0, Ψ′(w)=0\Psi^{\prime}(w)=0.

The following bounds hold for the functions and their first- and second-order derivatives:

The function hh is lower bounded as follows:

To prove Property 2), note that following holds for w>0w>0:

Obviously, Ψ(w)\Psi(w) is an increasing function over w>0w>0, therefore the lower and upper bounds are Ψ(0)=0,Ψ(∞)=1\Psi(0)=0,\Psi(\infty)=1; Ψ′(w)\Psi^{\prime}(w) is increasing on [0,12][0,\frac{1}{\sqrt{2}}] and decreasing on [12,∞][\frac{1}{\sqrt{2}},\infty], where Ψ′′(12)=0\Psi^{\prime\prime}(\frac{1}{\sqrt{2}})=0, therefore the lower and upper bounds are Ψ′(0)=Ψ′(∞)=0,Ψ′(12)=2e\Psi^{\prime}(0)=\Psi^{\prime}(\infty)=0,\Psi^{\prime}(\frac{1}{\sqrt{2}})=\sqrt{\frac{2}{e}}; Ψ′′(w)\Psi^{\prime\prime}(w) is decreasing on (0,32](0,\sqrt{\frac{3}{2}}] and increasing on [32,∞)[\sqrt{\frac{3}{2}},\infty) (this can be verified by checking the signs of Ψ′′′(w)=4e−w2w(2w2−3)\Psi^{\prime\prime\prime}(w)=4e^{-w^{2}}w(2w^{2}-3) in these intervals). Therefore the lower and upper bounds are Ψ′′(32)=−4e32,Ψ′′(0+)=2\Psi^{\prime\prime}(\sqrt{\frac{3}{2}})=-\frac{4}{e^{\frac{3}{2}}},\Psi^{\prime\prime}(0^{+})=2, i.e.,

Similarly, as above, we can obtain the following bounds:

To show Property 3), note that for all w≥1w\geq 1 and ∣v∣<1|v|<1,

where the first inequality is true because Ψ(w)\Psi(w) is strictly increasing and Φ′(v)\Phi^{\prime}(v) is strictly decreasing for all w>0w>0 and v>0v>0, and that Φ′(v)=Φ′(∣v∣)\Phi^{\prime}(v)=\Phi^{\prime}(|v|).

Next we show Property 4). Note that 0≤Ψ(w)<10\leq\Psi(w)<1 and 0<Φ(w)<4π0<\Phi(w)<4\pi. Therefore we have g(0)=−Ψ(1)Φ(0)<0{g}(0)=-\Psi(1)\Phi(0)<0 and using the construction in (22)

where the first inequality follows from Ψ(w)Φ(v)>0\Psi(w)\Phi(v)>0, the second follows from Ψ(w)Φ(v)<4π\Psi(w)\Phi(v)<4\pi, and the last is true because T/N≥1T/N\geq 1.

Finally, we show Property 5), using the fact that a function is Lipschitz if it is piecewise smooth with bounded derivative. Before proceeding, let us note a few properties of the construction in (24) (also see Fig. 3). First, for a given node qq, its local function hqh_{q} is only related to the following x[j]x[j]’s

Then the first-order partial derivative of gq(y)g_{q}(y) can be expressed below.

From the above derivation, it is clear that for any j,qj,q, ∂gq∂x[j]\frac{\partial g_{q}}{\partial x[j]} is either zero or is a piecewise smooth function separated at the non-differentiable point x[j]=0x[j]=0, because the function Ψ′(⋅)\Psi^{\prime}(\cdot) is not differentiable at .

By using (32) and the above facts, the second-order partial derivative of gq(x)g_{q}(x) (∀x≠0\forall x\neq 0) is given as follows when j≠1j\neq 1:

By applying Lemma 1 – i) [i.e., Ψ(w)=Ψ′(w)=Ψ′′(w)=0\Psi(w)=\Psi^{\prime}(w)=\Psi^{\prime\prime}(w)=0 for ∀  w≤0\forall\;w\leq 0], we can obtain that at least one of the terms Ψ(−x[j−1])Φ′′(−x[j])\Psi\left(-x[j-1]\right)\Phi^{\prime\prime}\left(-x[j]\right) or −Ψ(x[j−1])Φ′′(x[j])-\Psi\left(x[j-1]\right)\Phi^{\prime\prime}\left(x[j]\right) is zero. It follows that

Taking the maximum over equations (34) to (40) and plug in the above inequalities, we obtain

where the equality comes from Lemma 1 – ii).

When j=1j=1, by using (33), we have the following:

Summarizing the above results, we obtain:

The following lemma is a simple extension of the previous result.

We have the following properties for the functions ff defined in (26) and (25):

The first-order derivatives of f{f} and that for each fi,i∈[N]f_{i},i\in[N] are Lipschitz continuous, with the same constant U>0U>0.

Proof. To show that property 1) is true, note that we have the following:

Then by applying Lemma 1 we have that for any T≥1T\geq 1, the following holds

Property 2) is true is due to the definition of fif_{i}, so that we have:

Property 3) is true because the following:

where the last inequality comes from Lemma 1 – (5). This completes the proof. ■\blacksquare

Next let us analyze the size of ∇g\nabla{g}. We have the following result.

If there exists k∈[T]k\in[T] such that ∣x[k]∣<1|x[k]|<1, then

Proof. The first inequality holds for all k∈[T]k\in[T], since 1N∑i=1N∂∂y[k]gi(x)\frac{1}{N}\sum_{i=1}^{N}\frac{\partial}{\partial y[k]}g_{i}(x) is one element of 1N∑i=1N∇gi(x)\frac{1}{N}\sum_{i=1}^{N}\nabla g_{i}(x). We divide the proof for the second inequality into two cases.

Case 1. Suppose ∣x[j−1]∣<1|x[j-1]|<1 for all 2≤j≤k2\leq j\leq k. Therefore, we have ∣x∣<1|x|<1. Using (33), we have the following inequalities:

where (i){\rm(i)} is true because Ψ′(w),Φ(w)\Psi^{\prime}(w),\Phi(w) are all non-negative from Lemma 1 -(2); (ii){\rm(ii)} is true due to Lemma 1 – (3). Therefore, we have the following

Case 2) Suppose there exists 2≤j≤k2\leq j\leq k such that ∣x[j−1]∣≥1|x[j-1]|\geq 1.

We choose jj so that ∣x[j−1]∣≥1|x[j-1]|\geq 1 and ∣x[j]∣<1|x[j]|<1. Therefore, depending on the choices of (i,j)(i,j) we have three cases:

First, note that ∂gi(x)∂x[j]≤0\frac{\partial g_{i}(x)}{\partial x[j]}\leq 0, for all i,ji,j, by checking the definitions of Ψ(⋅),Φ′(⋅),Ψ′(⋅),Φ(⋅)\Psi(\cdot),\Phi^{\prime}(\cdot),\Psi^{\prime}(\cdot),\Phi(\cdot).

Then for (i,j)(i,j) satisfying the first condition, because ∣x[j−1]∣≥1|x[j-1]|\geq 1 and ∣x[j]∣<1|x[j]|<1, using Lemma 1 – (3), and the fact that the negative part is zero for Ψ\Psi, and Φ′\Phi^{\prime} is even function, the expression further simplifies to:

If the second condition holds true, the expression is obviously non-positive because both Ψ′\Psi^{\prime} and Φ\Phi are non-negative. Overall, we have

Consider using an algorithm of the form (20) to solve the following problem:

Assume the initial solution: xi=0,  ∀ i∈[N]x_{i}=0,\;\forall~{}i\in[N]. Let xˉ=1N∑i=1Nαixi\bar{x}=\frac{1}{N}\sum_{i=1}^{N}\alpha_{i}x_{i} denote some linear combination of local variables, where {αi>0}\{\alpha_{i}>0\} are the coefficients (possibly time-varying and dependent on tt). Then no matter how many local computation steps (20b) are performed, at least TT communication steps (20a) are needed to ensure xˉ[T]≠0\bar{x}[T]\neq 0.

Proof. For a given j≥2j\geq 2, suppose that xi[j],xi[j+1],...,xi[T]=0x_{i}[j],x_{i}[j+1],...,x_{i}[T]=0, ∀i\forall i, that is, \mboxsupport{xi}⊆{1,2,3,...,j−1}\mbox{support}\{x_{i}\}\subseteq\{1,2,3,...,j-1\} for all ii. Then Ψ′(xi[j])=Ψ′(−xi[j])=0\Psi^{\prime}\left(x_{i}[j]\right)=\Psi^{\prime}\left(-x_{i}[j]\right)=0 for all ii, and gig_{i} has the following partial derivative (see (32))

Clearly, if xi[j−1]=0x_{i}[j-1]=0, then by the definition of Ψ(⋅)\Psi(\cdot), the above partial gradient is also zero. In other words, the above partial gradient is only non-zero if xi[j−1]≠0x_{i}[j-1]\neq 0.

Recall that we have assumed that the server aggregation is performed using a linear combination xˉ=1N∑i=1Nαixi\bar{x}=\frac{1}{N}\sum_{i=1}^{N}\alpha_{i}x_{i}, with the coefficients αi\alpha_{i}’s possibly depending on the stage tt (but such a dependency will be irrelevant for our purpose, as will be see shortly). Therefore, at a given stage tt, for a given node ii, when j≥3j\geq 3, its jjth element will become nonzero only if one of the following two cases hold true:

If before the aggregation step (i.e., at stage t−1t-1), some other node qq has xq[j]x_{q}[j] being nonzero.

If ∂gi(xi)∂xi[j]\frac{\partial g_{i}(x_{i})}{\partial x_{i}[j]} is nonzero at stage tt.

Now suppose that the initial solution is xi[j]=0x_{i}[j]=0 for all (i,j)(i,j). Then at the first iteration only ∂gi(xi)∂xi\frac{\partial g_{i}(x_{i})}{\partial x_{i}} is non-zero for all ii, due to the fact that ∂gi(xi)∂xi=Ψ(1)Φ′(0)=4(1−e−1)\frac{\partial g_{i}(x_{i})}{\partial x_{i}}=\Psi(1)\Phi^{\prime}(0)=4(1-e^{-1}) for all ii from (33). It is also important to observe that, if all nodes i≠1i\neq 1 were to perform subsequent local updates (20b), the local variable xjx_{j} will have the same support (i.e., only the first element is non-zero). To see this, suppose k=2k=2, then for i=2i=2, we have

since x=0x=0 implies Ψ′(−x)=0\Psi^{\prime}\left(-x\right)=0. Similarly reasoning applies when i=2i=2, k≥3k\geq 3.

If i≥3i\geq 3, then these local functions are not related to xix_{i}, so the partial derivative is also zero.

Now let us look at node i=1i=1. For this node, according to (46), we have

Since x1x_{1} can be non-zero, then this partial gradient can also be non-zero. Further, with a similar argument as above, we can also confirm that no matter how many local computation steps that node 11 performs, only the first two elements of x1x_{1} can be non-zero.

So for the first stage t=1t=1, we conclude that, no matter how many local computation that the nodes perform (in the form of the computation step given in (20b)), only x1x_{1} can have two non-zero entries, while the rest of the local variables only have one non-zero entries.

Then suppose that the communication and aggregation step is performed once. It follows that after broadcasting xˉ=1N∑i=1Nαixi\bar{x}=\frac{1}{N}\sum_{i=1}^{N}\alpha_{i}x_{i} to all the nodes, everyone can have two non-zero entries. Then the nodes proceed with local computation, and by the same argument as above, one can show that this time only x2x_{2} can have three non-zero entries. Following the above procedure, it is clear that each aggregation step can advance the non-zero entry of xˉ\bar{x} by one, while performing multiple local updates does not advance the non-zero entry. Then we conclude that we need at least TT communication steps, and local gradient computation steps, to make xi[T]x_{i}[T] possibly non-zero. ■\blacksquare

B.4 Main Result for Claim 2.1.

Below we state and prove a formal version of Claim 2.1 given in the main text.

Let ϵ\epsilon be a positive number. Let xi0[j]=0x_{i}^{0}[j]=0 for all i∈[N]i\in[N], and all j=1,⋯ ,T+1j=1,\cdots,T+1. Consider any algorithm obeying the rules given in (10), where the Vt(⋅)V^{t}(\cdot) and Wit(⋅)W^{t}_{i}(\cdot)’s are linear operators. Then regardless of the number of local updates there exists a problem satisfying Assumption 1 – 2, such that it requires at least the following number of stages tt (and equivalently, aggregation and communications rounds in (20a))

Proof of Claim 2.1. First, let us show that the algorithm obeying the rules given in (20) has the desired property. Note that the difference between two rules is whether the sampled local gradients are used for the update, or the full local gradients are used.

By Lemma 4 we have xˉ[T]=0\bar{x}[T]=0 for all t<Tt<T. Then by applying Lemma 2 – (2) and Lemma 3, we can conclude that the following holds

where the second inequality follows that there exists k∈[T]k\in[T] such that ∣xˉ[k]Uπ2ϵ∣=0<1|\frac{\bar{x}[k]U}{\pi\sqrt{2\epsilon}}|=0<1, then we can directly apply Lemma 3.

The third part of Lemma 2 ensures that fif_{i}’s are LL-Lipschitz continuous gradient, and the first part shows

Second, consider the algorithm obeying the rules give in (10), in which local sampled gradients are used. By careful inspection, the result for this case can be trivially extended from the previous case. We only need to consider the following local functions

where each sampled loss function F(x;ξi)F(x;\xi_{i}) is defined as

where δ(ξi)\delta(\xi_{i})’s satisfy δ(ξi)>0\delta(\xi_{i})>0 and ∑ξi∈Diδ(ξi)=1\sum_{\xi_{i}\in D_{i}}\delta(\xi_{i})=1. It is easy to see that, the local sampled gradients have the same dependency on xx as their averaged version (by dependency we meant the structure that is depicted in Fig. 3). Therefore, the progression of the non-zero pattern of the average xˉ=1N∑i=1Nxi\bar{x}=\frac{1}{N}\sum_{i=1}^{N}x_{i} is exactly the same as the batch gradient version. Additionally, since the local function f^(x)\hat{f}(x) is exactly the same as the previous local function f(x)f(x), so other estimates, such as the one that bounds f(0)−inf⁡f(x)f(0)-\inf f(x), also remain the same.

Appendix C Proof of Claim 2.2

First let us consider FedAvg with local-GD update (9). We consider the following problem with N=2N=2, which satisfies both Assumptions 1 and 2, with f(x)=0,  ∀ xf(\mathbf{x})=0,\;\forall~{}\mathbf{x}

Each local iteration of the FedAvg is given by

For simplicity, let us define y=[x1,x2]T\mathbf{y}=[\mathbf{x}_{1},\mathbf{x}_{2}]^{T}, and define the matrix D=[1−η,0;0,1+η]\mathbf{D}=[1-\eta,0;0,1+\eta]. Then running QQ rounds of the FedAvg algorithm starting with r=kQr=kQ for some non-negative integer k≥0k\geq 0, can be expressed as

It is easy to show that for any Q>1Q>1, the eigenvalues of the matrix 12DQ−111TD\frac{1}{2}\mathbf{D}^{Q-1}\mathbf{1}\mathbf{1}^{T}\mathbf{D} are and (1+η)Q+(1−η)Q2>1\frac{(1+\eta)^{Q}+(1-\eta)^{Q}}{2}>1.

It follows that the above iteration will diverge for any Q>1Q>1 starting from any non-zero initial point.

Moreover, when the sample on one agent are the same (e.g., agent 1 has two samples that both has loss function x2x^{2}), then using SGD as local update will be identical to the update of GD.

Appendix D Results showing the role of GB for FedAvg with diminishing stepsizes

This section is used to show the role of A3 in FedAvg. First we prove that with A3, FedAvg can use arbitrary stepsize as long as it diminishes to zero. Next, we show by a simple example that without A3, FedAvg will diverge with the same diminishing stepsize choice.

Suppose A1–A3 hold and the stepsizes satisfy: 1) ηr,0=η∈(0,1/L)\eta^{r,0}=\eta\in(0,1/L) for all rr; 2) set 0<ηr,q≤min⁡{12(Q−1)L,ηQ},  lim⁡r→∞ηr,q=0, q≠00<\eta^{r,q}\leq\min\{\frac{1}{2(Q-1)L},\frac{\eta}{Q}\},\;\lim_{r\rightarrow\infty}\eta^{r,q}=0,~{}q\neq 0. Then the average gradient converges to zero for FedAvg with local-GD update (9):

Suppose that all the assumptions made in Claim D.1 hold, except that A3 does not hold. Then FedAvg with local-GD can diverge for any Q>1Q>1.

Before we prove Claim D.1, the following lemma is needed.

Under A1 and A3, following the update steps in Algorithm 1, between each outer iterations we have:

where (a)(a) comes from the update rule in Algorithm 1, in (b)(b) we use Jensen’s inequality and notice xir,0=xr\mathbf{x}^{r,0}_{i}=\mathbf{x}^{r} so in (c)(c) we extract the terms with index (r,0)(r,0) from the inner product.

Note that for any vector a,ba,b of the same length, the equality 2⟨a,b⟩=∥a∥2+∥b∥2−∥a−b∥2,2\left\langle a,b\right\rangle=\left\lVert a\right\rVert^{2}+\left\lVert b\right\rVert^{2}-\left\lVert a-b\right\rVert^{2}, holds, we have

where we use Jensen’s inequality in (a)(a) and A1 in (b)(b).

The first equality comes from the update rule of xir,q\mathbf{x}^{r,q}_{i}, which basically performs qq steps of updates on xr\mathbf{x}^{r}; (a)(a) comes from Jensen’s inequality; in (b)(b) we use A3.

Substitute (63) to (62) and then to (61), rearrange the terms we obtain (60), which ends the proof of the lemma. ■\blacksquare

Proof: By choosing ηr,0=η1=∈(0,1/L)\eta^{r,0}=\eta_{1}=\in(0,1/L) as constant and ηr,q≤1/(2QL)  ,∀ q≠0\eta^{r,q}\leq 1/(2QL)\;,\forall~{}q\neq 0 then applying Lemma 5 we have

where C1=η1(1−Lη1)>0C_{1}=\eta_{1}(1-L\eta_{1})>0. Using telescope sum from r=0r=0 to r=T−1r=T-1 we have

Rearrange the terms and multiply both side by 2/(TC1)2/(TC_{1}), then we have

Choose ηr,q≤η1/Q\eta^{r,q}\leq\eta_{1}/Q, then (η1)2+∑q=1Q−1(ηr,q)2≤2(η1)2(\eta_{1})^{2}+\sum^{Q-1}_{q=1}(\eta^{r,q})^{2}\leq 2(\eta_{1})^{2}. Choose {ηr,q}\{\eta^{r,q}\} as a sequence that diminishes to , then for all q≠0q\neq 0, as T→∞T\rightarrow\infty, 2η1Q2G2C11QT∑r=0T−1∑q=1Q−1ηr,q→0\frac{2\eta_{1}Q^{2}G^{2}}{C_{1}}\frac{1}{QT}\sum^{T-1}_{r=0}\sum^{Q-1}_{q=1}\eta^{r,q}\rightarrow 0. Therefore the right hand side converges to 0, Claim D.1 is proved.

Appendix E Proof of Claim D.2

Proof. We consider the following problem with N=2N=2, which satisfies both Assumptions 1 and 2, with f(x)=0,  ∀ xf(\mathbf{x})=0,\;\forall~{}\mathbf{x}

Each local iteration of the FedAvg is given by

For simplicity, let us define y=[x1,x2]T\mathbf{y}=[\mathbf{x}_{1},\mathbf{x}_{2}]^{T}, and define the matrix Dr=[1−ηr,0;0,1+ηr]\mathbf{D}_{r}=[1-\eta^{r},0;0,1+\eta^{r}]. Then running QQ rounds of the FedAvg algorithm starting with r=kQr=kQ for some non-negative integer k≥0k\geq 0, can be expressed as

In particular, we pick ηr=1r\eta^{r}=\frac{1}{\sqrt{r}} when r≠kQ+1r\neq kQ+1 and ηkQ+1=1/2\eta^{kQ+1}=1/2. Then for Q>1Q>1, it is easy to compute the eigenvalues of the matrix 12∏r=kQ+1(k+1)Q−1Dr11TDkQ\frac{1}{2}\prod_{r=kQ+1}^{(k+1)Q-1}\mathbf{D}_{r}\mathbf{1}\mathbf{1}^{T}\mathbf{D}_{kQ} to be:

It is clear that λ2\lambda_{2} is strictly larger than one which indicates that the algorithm will diverge. ■\blacksquare

Appendix F Proofs for Results in Section 3

First let us prove Theorem 1 about the FedPD algorithm with Oracle I.

Towards this end, let us first introduce some notations. First recall that when Oracle I is used, the local problem is solved such that the following holds true:

Note that if SGD is applied in Oracle I to solve the local problem, then this condition (71) is replaced with the following

The difference does not significantly change the proofs and the results. So throughout the proof of Theorem 1, we use (71) as the condition.

Then we define the error between different nodes as

Here, △x0r\triangle\mathbf{x}^{r}_{0} denotes the maximum difference of estimated center model among all the nodes and △xr\triangle\mathbf{x}^{r} denotes the maximum difference of local models among all nodes.

From the termination condition that generates xir+1\mathbf{x}^{r+1}_{i} (given in (71)), we have

where the first equality holds because of the update rule of λi\lambda_{i}. Furthermore, from the update step of λir+1\lambda^{r+1}_{i}, we can explicitly write down the following expression

The main lemmas that we need are outlined below. Their proofs can be found in Sec. F.1.1– F.1.4.

The first lemma shows the sufficient descent of the local AL function.

Suppose A1 holds true. Consider FedPD with Algorithm 4 (Oracle I) as the update rule. When the local problem is solved such that (71) is satisfied, the difference of the local augmented Lagrangian is bounded by

Then we derive a key lemma about how the error propagates if the communication step is skipped.

Suppose A1 and A5 hold. Consider FedPD with Algorithm 4 (Oracle I) as the update rule. When the local problem is solved such that (71) is satisfied, the difference between the local models xir\mathbf{x}^{r}_{i}’s and the difference between local copies of the global models x0,ir\mathbf{x}^{r}_{0,i}’s are bounded by

is a rank 11 matrix with eigenvalues (0,Lη+p(1+Lη))(0,L\eta+p(1+L\eta)) and B=[p(3+Lη),2]TB=[p(3+L\eta),2]^{\scriptscriptstyle T}.

We define a virtual sequence {x‾0r}\{\overline{\mathbf{x}}^{r}_{0}\} where x‾0r≜1N∑i=1Nx0,ir\overline{\mathbf{x}}^{r}_{0}\triangleq\frac{1}{N}\sum^{N}_{i=1}\mathbf{x}^{r}_{0,i} which is the average of the local x0,ir\mathbf{x}^{r}_{0,i} and we know that x0,ir=x0r\mathbf{x}^{r}_{0,i}=\mathbf{x}^{r}_{0} when rmod  R=1r\mod R=1, that is, when the communication and aggregation step is performed. Next, we bound the error between the local AL and the global AL evaluated at the virtual sequence.

Suppose A1 holds. Consider FedPD with Algorithm 4 (Oracle I) as the update rule. When the local problem is solved such that (71) is satisfied, the difference between local AL and the global AL is bounded as below:

Lastly we bound the original objective function using the global AL.

Under A1 and A2, when the local problem is solved to ϵ1\epsilon_{1} accuracy, the difference between the original loss and the augmented Lagrangian is bounded.

Using the previous lemmas, we can then prove Theorem 1.

We divide the left hand side (LHS) of (75), i.e., Li(xir+1,x0,ir+,λir+1)−Li(xir,x0,ir,λir)\mathcal{L}_{i}(\mathbf{x}^{r+1}_{i},\mathbf{x}^{r+}_{0,i},\lambda^{r+1}_{i})-\mathcal{L}_{i}(\mathbf{x}^{r}_{i},\mathbf{x}^{r}_{0,i},\lambda^{r}_{i}), into the sum of three parts:

which correspond to the three steps in the algorithm’s update steps.

We bound the first difference by first applying A1 to −f(⋅)-f(\cdot):

and obtain the following series of inequalities:

In the above equation, in (a)(a) we use the fact that ∥a∥2−∥b∥2=⟨a+b,a−b⟩\left\lVert a\right\rVert^{2}-\left\lVert b\right\rVert^{2}=\left\langle a+b,a-b\right\rangle when vector a,ba,b has the same length to the last two terms; in (b)(b) we split the last term into 2xir+1−2x0,ir2\mathbf{x}^{r+1}_{i}-2\mathbf{x}^{r}_{0,i} and −xir+1+xir-\mathbf{x}^{r+1}_{i}+\mathbf{x}^{r}_{i}; in (c)(c) we use the fact that ⟨a,b⟩≤L2∥a∥2+12L∥b∥2)\left\langle a,b\right\rangle\leq\frac{L}{2}\left\lVert a\right\rVert^{2}+\frac{1}{2L}\left\lVert b\right\rVert^{2}); in (d)(d) we apply the fact that xir+1\mathbf{x}^{r+1}_{i} is the inexact solution; see (74).

Then we bound the second difference in (79) by the following:

where (a)(a) directly comes from the update rule of λir+1\lambda^{r+1}_{i}.

Further we bound the third difference in (79) by the following:

where, in (a)(a), we use the same reasoning as in (80) (a)(a) and (b)(b); in (b)(b) we apply the update rule of x0,ir+\mathbf{x}^{r+}_{0,i} in the FedPD algorithm, which implies that the first term becomes zero.

Finally we sum up (80), (81), (82) and Lemma 6 is proved.

F.1.2 Proof of Lemma 7

First we derive the relation between ∥xir+1−xjr+1∥\left\lVert\mathbf{x}^{r+1}_{i}-\mathbf{x}^{r+1}_{j}\right\rVert for arbitrary i≠ji\neq j and △r\triangle^{r} by using the definition of ϵ1\epsilon_{1} (74):

where in (a)(a) we plug the definition of △x0r\triangle\mathbf{x}^{r}_{0} and eir+1\mathbf{e}^{r+1}_{i}; in (b)(b) we use A1; (c)(c) comes form A5; in (d)(d) we move the second term to the left and divide both side by 1−Lη1-L\eta.

Then we bound the difference ∥λir−λjr∥\left\lVert\lambda^{r}_{i}-\lambda^{r}_{j}\right\rVert by plugging in the expression of λir\lambda^{r}_{i} in (74), and note that λir+1η(xir+1−x0,ir)=λir+1\lambda^{r}_{i}+\frac{1}{\eta}(\mathbf{x}^{r+1}_{i}-\mathbf{x}^{r}_{0,i})=\lambda^{r+1}_{i}:

where (a)(a) and (b)(b) follow the same argument in (a)(a), (b)(b) and (c)(c) of (83) ; in (c)(c) we plug in the definition of △xr\triangle\mathbf{x}^{r}.

Next we bound the difference ∥x0,ir+1−x0,jr+1∥\left\lVert\mathbf{x}^{r+1}_{0,i}-\mathbf{x}^{r+1}_{0,j}\right\rVert. With probability 1−p1-p the aggregation step has just been done at iteration rr, x0,ir+1=x0,jr+1\mathbf{x}^{r+1}_{0,i}=\mathbf{x}^{r+1}_{0,j}.With probability pp, they are not equal, then we take expectation with communication probability pp, and get

where in (a)(a) we plug in the definition of △xr+1\triangle\mathbf{x}^{r+1} and (84). As these relations hold true for arbitrary (i,j)(i,j) pairs, they are also true for the maximum of ∥xir+1−xjr+1∥\left\lVert\mathbf{x}^{r+1}_{i}-\mathbf{x}^{r+1}_{j}\right\rVert and ∥x0,ir+1−x0,jr+1∥\left\lVert\mathbf{x}^{r+1}_{0,i}-\mathbf{x}^{r+1}_{0,j}\right\rVert.

Therefore stacking (83) and (85) and plug in (84), we have

Rewrite it into matrix form then we complete the proof of Lemma 7.

F.1.3 Proof of Lemma 8

Let us first recall that the definition of local AL is given below:

where (a)(a) follows the same argument in (82); in (b)(b),we plug in the definition of x‾0r+1\overline{\mathbf{x}}^{r+1}_{0}; in (c)(c) we use Jensen’s inequality and we bound the term with △x0r+1\triangle\mathbf{x}^{r+1}_{0}. Then the lemma is proved.

F.1.4 Proof of Lemma 9

Taking an average over NN agents we are able to prove Lemma 9.

F.1.5 Proof of Theorem 1

First notice that from the optimality condition (74), the following holds:

Then we bound the gradients of L(xir,x0,ir,λir)\mathcal{L}(\mathbf{x}^{r}_{i},\mathbf{x}^{r}_{0,i},\lambda^{r}_{i}).

Further, we note that, when no aggregation has been performed at iteration rr, then x0,ir=xir+ηλir\mathbf{x}^{r}_{0,i}=\mathbf{x}^{r}_{i}+\eta\lambda^{r}_{i}, so the following holds

When there the aggregation has been performed at iteration rr, then x0,ir=1N∑j=1N(xjr+ηλjr), ∀i\mathbf{x}^{r}_{0,i}=\frac{1}{N}\sum^{N}_{j=1}(\mathbf{x}^{r}_{j}+\eta\lambda^{r}_{j}),~{}\forall i, so we have

Summing (90) and (93), denote ∥∇xiLi(xir,x0,ir,λir)∥+∥∇λiLi(xir,x0,ir,λir)∥\left\lVert\nabla_{\mathbf{x}_{i}}\mathcal{L}_{i}(\mathbf{x}^{r}_{i},\mathbf{x}^{r}_{0,i},\lambda^{r}_{i})\right\rVert+\left\lVert\nabla_{\lambda_{i}}\mathcal{L}_{i}(\mathbf{x}^{r}_{i},\mathbf{x}^{r}_{0,i},\lambda^{r}_{i})\right\rVert as ∥∇Li(xir,x0,ir,λir)∥\left\lVert\nabla\mathcal{L}_{i}(\mathbf{x}^{r}_{i},\mathbf{x}^{r}_{0,i},\lambda^{r}_{i})\right\rVert we have

Squaring both sides of the above inequality, we obtain:

where C6≥max⁡{(1+Lηη)2,(1+2η)2,L2η2}C_{6}\geq\max\{(\frac{1+L\eta}{\eta})^{2},(1+2\eta)^{2},L^{2}\eta^{2}\}.

Notice that when communication is not performed ∥x0,ir−x0,ir+1∥2≤∥x0,ir−x0,ir+∥2\left\lVert\mathbf{x}^{r}_{0,i}-\mathbf{x}^{r+1}_{0,i}\right\rVert^{2}{\leq}\left\lVert\mathbf{x}^{r}_{0,i}-\mathbf{x}^{r+}_{0,i}\right\rVert^{2}, and when communication is performed

where the last inequality holds due to the use of Jensen’s inequality, and the definition of △x0r+1\triangle\mathbf{x}^{r+1}_{0} in (73). It follows that summing both sides of (96) over ii, we have

Taking the expectation over the randomness in pp, conditioning on the information before the communication the successive difference of Li\mathcal{L}_{i},

where (a) expands the expectation on pp, and use the fact that with probability pp, x0,ir+1=x0,ir+\mathbf{x}^{r+1}_{0,i}=\mathbf{x}^{r+}_{0,i}, and with probability (1−p)(1-p) x0r+1\mathbf{x}^{r+1}_{0} will be updated; in (b) we apply Lemma 8 to the last term.

Combining (95), (98) and (100), define C7=2C6/min⁡{1−2Lη−4L2η22η,12η,1+8Lη2L}C_{7}=2C_{6}/\min\{\frac{1-2L\eta-4L^{2}\eta^{2}}{2\eta},\frac{1}{2\eta},\frac{1+8L\eta}{2L}\} and sum up the iterations, we have

Next we bound the last term. By iteratively applying Lemma 7 from τ=0\tau=0 to rr and use the fact that Δ0=0\Delta^{0}=0, we have

So by taking norm square on both side of (102), we have

Substitute (103) into (101) and divide both side by TT we have

From the initial conditions we have L(x00,xi0,λi0)=f(x00)\mathcal{L}(\mathbf{x}^{0}_{0},\mathbf{x}^{0}_{i},\lambda^{0}_{i})=f(\mathbf{x}^{0}_{0}) and apply Lemma 9 we obtain

Finally we bound ∥∇f(x0r)∥2\left\lVert\nabla f(\mathbf{x}^{r}_{0})\right\rVert^{2} by

where in (a)(a) we use the same argument in (91) and (92).

Therefore Theorem 1 is proved. During the proof, we need all C2,…,C7,C8>0C_{2},\dots,C_{7},C_{8}>0, therefore, 0<η<5−14L0<\eta<\frac{\sqrt{5}-1}{4L}.

Finally, let us note that if the local problems are solved with SGD, then the local problem needs to be solved such that the condition (72) holds true. As no other information of the local solvers except error term eir\mathbf{e}^{r}_{i} is used in the proof, the proofs and results of FedPD with SGD as local solver will not change much, except that all the results hold in expectation. Therefore we skip the proof for the SGD version.

F.1.6 Constants used in the proofs

In this subsection we list all the constants C2,…,C8C_{2},\dots,C_{8} used in the proof of Theorem 1.

we can see that when 0<η<5−14L0<\eta<\frac{\sqrt{5}-1}{4L}, all the terms are positive.

Appendix G FedPD with Variance Reduction

In this section we provided an alternative oracle for FedPD which has a lower sample complexity.

Alternatively, when instantiating the local oracle using Algorithm 4, the original local problems are not required to solve to ϵ1\epsilon_{1} accuracy. Instead, we successively optimize a linearized AL function:

In the above expression, we linearize fi(xi)f_{i}(\mathbf{x}_{i}) at inner iteration xir,q\mathbf{x}^{r,q}_{i} as (γ\gamma is a constant)

where gir,qg^{r,q}_{i} is an approximation of ∇fi(xir,q)\nabla f_{i}(\mathbf{x}^{r,q}_{i}). The optimizer has a closed-form expression:

In Oracle II, an agent ii first decides whether to compute the full gradient ∇fi(xir,0)\nabla f_{i}(\mathbf{x}^{r,0}_{i}), or to keep using the previous estimate gir−1,Qg^{r-1,Q}_{i}. Then QQ local steps are performed, each requires BB local data samples. In this scheme, QQ can be chosen as any positive integer.

It is important to note that this oracle does not simply apply the VR technique (such as F-SVRG) to solve the subproblem of optimizing Li(xi,x0,ir,λir)\mathcal{L}_{i}(\mathbf{x}_{i},\mathbf{x}^{r}_{0,i},\lambda^{r}_{i}). That is, it is not a variation of Oracle I. Instead, the VR technique is applied to the entire primal-dual iteration, and the full gradient evaluation ∇fi(xir,0)\nabla f_{i}(\mathbf{x}^{r,0}_{i}) is only needed every II iteration rr. Later we will see that if II is large enough, then there is an O(M)\mathcal{O}(\sqrt{M}) reduction of sample complexity.

G.2 Algorithm Convergence and Complexity

The convergence result of FedPD with Oracle II is given as follows:

Suppose A1–A2 hold, and consider FedPD with Oracle II. Choose p=0p=0, \eta\in\big{(}0,\frac{1}{3(Q+\sqrt{QI/B})L}\big{)}, and γ>5ηBL\gamma>\frac{5\eta}{B\sqrt{L}}. Then, the following holds (where C9>0C_{9}>0 is a constant):

Remark 5. (Communication complexity): As p=0p=0, the communication round to achieve ϵ\epsilon accuracy is T=O(1/ϵ)T=\mathcal{O}(1/\epsilon) which is independent of QQ.

Remark 6. (Computation complexity): Note that the total number full gradient evaluation is T/I+1T/I+1, each uses MM samples. Meanwhile, the total number of mini-batch stochastic gradient evaluation is TQTQ, each uses 2B2B samples per node. So the total sample complexity is O(M+MT/I+2TQBN)\mathcal{O}(M+MT/I+2TQBN). In order to keep the same convergence speed, we need stepsize η\eta to be unchanged. Therefore, we choose I=M,B=I/QN=M/QNI=\sqrt{M},B=I/QN=\sqrt{M}/QN, then the SC of Algorithm 4 is O(M+Mϵ)\mathcal{O}(M+\frac{\sqrt{M}}{\epsilon}).

G.3 Proof of Theorem 3

Following the similar proof of Theorem 1, we first analyze the descent between each outer iteration. Notice throughout the proof, we assume that p=0p=0, that is, there is no delayed communication. It follows that the following holds:

We also recall that rr is the (outer) stage index, and qq is the local update index. First we provide a series of lemmas.

Under Assumption 1, consider FedPD with Algorithm 4 (Oracle II) as the update rule. The difference of the local AL is bounded by:

Then we deal with the variance of the stochastic gradient estimations.

Suppose A1 holds true and the samples are randomly sampled according to (13), consider FedPD with Algorithm 4 (Oracle II) as the update rule. The expected norm square of the difference between gir,q+1g^{r,q+1}_{i} and ∇fi(xir,q+1)\nabla f_{i}(\mathbf{x}^{r,q+1}_{i}) is bounded by

Lastly we upper bound the original loss function.

Under A1 and A2, the difference between the original loss and the AL is bounded as below:

Let us first express the difference of the local AL as following:

where the above three differences respectively correspond to the three steps in the algorithm’s update steps.

Let us bound the above three differences one by one. First, note that we have the following decomposition (by using the fact that xir,Q+1=xir+1\mathbf{x}^{r,Q+1}_{i}=\mathbf{x}^{r+1}_{i} and xir,1=xir\mathbf{x}^{r,1}_{i}=\mathbf{x}^{r}_{i}):

Each term on the right hand side (RHS) of the above equality can be bounded by (see a similar arguments in (80)):

in (b)(b) we use the fact that 2⟨a,b⟩≤L∥a∥2+1L∥b∥22\left\langle a,b\right\rangle\leq L\left\lVert a\right\rVert^{2}+\frac{1}{L}\left\lVert b\right\rVert^{2}. Therefore, the first difference in the RHS of (108) is given by

The other two differences in (108) can be explicitly expressed as:

Next we bound ∥λir+1−λir∥2\left\lVert\lambda^{r+1}_{i}-\lambda^{r}_{i}\right\rVert^{2}. Notice that the from the update rule the following holds:

where in (a)(a) we apply Cauchy-Schwarz inequality. Next we bound ∥gir,Q−1−gir−1,Q−1∥2\left\lVert g^{r,Q-1}_{i}-g^{r-1,Q-1}_{i}\right\rVert^{2} by

where in (a)(a) and (b)(b) we both apply Cauchy-Schwarz inequality, in (a)(a) we use A1 to the last term and in (b)(b) we notice xir−1,Q=xir,0\mathbf{x}^{r-1,Q}_{i}=\mathbf{x}^{r,0}_{i}.

Substitute (G.3.1) to (116) and sum the three parts, we have

G.3.2 Proof of Lemma 11

That is, r0r_{0} is a multiple of II and there is no more than IQIQ local update steps between step {r0,0}\{r_{0},0\} and step {r,q}\{r,q\}. By the update rule of gir,qg^{r,q}_{i}, we have

Iteratively taking expectation until {r,q}={r0,0}\{r,q\}=\{r_{0},0\}, we have

G.3.3 Proof of Lemma 12

where in (a)(a) we apply Cauchy-Schwarz inequality twice, that is

in (b)(b) we apply Lemma 11 to the first term and apply A1 to the second term.

Substitute (120) to (119) and average over the agents, Lemma 12 is proved.

G.3.4 Proof of Theorem 3

By the update step of x0r\mathbf{x}^{r}_{0}, following (91) we have

where in (a)(a), the first term is obtained by plugging in (115) given below

Next we take expectation and substitute (116), (G.3.1),

where we substitute Lemma 11 and (G.3.1) in (a)(a).

Taking expectation of (10), summing over r=0r=0 to r=T−1r=T-1 and average over the agents, we have the following

where in (a)(a) we apply Lemma 11 and (91).

Finally, in the last equation of (121), we have defined the constant C10C_{10} as

Then by taking expectation and applying Lemma 12, we obtain

where by the initialization that xi0=x00\mathbf{x}^{0}_{i}=\mathbf{x}^{0}_{0} we have f(x00)=1N∑i=1NLi(xi0,x0,i0,λi0).f_{(}\mathbf{x}^{0}_{0})=\frac{1}{N}\sum^{N}_{i=1}\mathcal{L}_{i}(\mathbf{x}^{0}_{i},\mathbf{x}^{0}_{0,i},\lambda^{0}_{i}).

Combine (G.3.4) and (LABEL:eq:th_VR_2), we can find a positive constant C11C_{11} satisfying

Similar to the proof of Theorem 1, we can bound ∥∇f(x0r)∥2\left\lVert\nabla f(\mathbf{x}^{r}_{0})\right\rVert^{2} by 1N∑i=1N∥∇Li(xir,x0r,λir)∥2,\frac{1}{N}\sum^{N}_{i=1}\left\lVert\nabla\mathcal{L}_{i}(\mathbf{x}^{r}_{i},\mathbf{x}^{r}_{0},\lambda^{r}_{i})\right\rVert^{2}, therefore Theorem 3 is proved.

to be positive constant. By selecting γ>5BLη\gamma>\frac{5}{B\sqrt{L}}\eta, and 0<η<13(Q+QI/B)L0<\eta<\frac{1}{3(Q+\sqrt{QI/B})L}, this is guaranteed.

Appendix H Additional Numerical Results

In this experiment, we consider the penalized regression problem , whose loss function evaluated on a single sample (a,b)=ξ(\mathbf{a},b)=\xi is given by:

In the experiment, we use two ways to generate the data. In the first case (referred to as the “weakly non-i.i.d” case), the features and the labels on the agents are randomly generated, so the local data sets are not very non-i.i.d. In the second case (referred to as the “strong non-i.i.d.” case), we first generate the feature vector a\mathbf{a}’s following the standard Normal distribution, then we generate the local model xi\mathbf{x}_{i} on the ithi^{th} agent by using uniform distribution in the range of $foreachcomponent.Thenwecomputethelabelfor each component. Then we compute the labelb’saccordingtothelocalmodelsandthefeaturesandaddsomeuniformnoise.Inthiscase,thedatadistributionontheagentsaremorenon−i.i.d.comparedtothefirstcase.Inbothcases,thereare’s according to the local models and the features and add some uniform noise. In this case, the data distribution on the agents are more non-i.i.d. compared to the first case. In both cases, there are400samplesoneachagentwithtotalsamples on each agent with total100$ agents.

The total number of iterations TT is set as 600600 for all algorithms. We choose the stepsize to be η=4\eta=4 for FedAvg-GD with local update number Q=8Q=8 and for FedAvg-SGD we use diminishing stepsize η=4/Qr+q+1\eta=4/\sqrt{Qr+q+1} with Q=600Q=600. For FedProx we use VR algorithm as the local solver and set Q=8Q=8, ρ=1\rho=1 and stepsize η=4\eta=4. For FedPD, we also use the same stepsize η=4\eta=4 with Q=8Q=8 with local GD. For FedPD-SGD, we also set η=4\eta=4 and uses local step size η1=1Q\eta_{1}=\frac{1}{Q} with inner iteration number Q=600Q=600. Lastly for FedPD with VR, we set the parameters to be η=4\eta=4, γ=4\gamma=4, I=100I=100, Q=2Q=2 and B=1B=1. The choice of the stepsize is the same among all the algorithms. We used grid search on stepsizes η∈{5,2,1,0.1,0.01}\eta\in\{5,2,1,0.1,0.01\} and the relative performance of the algorithms are similar to what we will show shortly.

Fig. 4 shows the convergence results of the penalized logistic regression problem with the first data set. In Fig. 4(a), we compare the convergence of the tested algorithms w.r.t the communication rounds. It is clear that FedProx and FedPD with R=1R=1 (i.e., no communication skipping) are comparable. Meanwhile, FedAvg with local GD will not converge to the stationary point with a constant stepsize when local update step Q>1Q>1. By skipping half of the communication, FedPD with local GD can still achieve a similar error as FedAvg, but using fewer communication rounds. In Fig. 4(b), we compare the sample complexity of different algorithms. It can be shown that when using the same number of samples for computation, FedPD with Oracle II (FedPD-VR) converges the fastest among all the algorithms. FedProx uses VR to solve the inner problem and converges the second fastest. Fig 5 shows the convergence results with the strongly non-i.i.d. data set. We can see that the algorithms using stochastic solvers become less stable compared with the case when the data sets are weakly non-i.i.d. Further, FedPD-VR and FedPD-GD with R=1R=1 are able still to converge to the global stationary point while FedProx will achieve a similar error as the FedAvg with local GD.

H.2 Handwritten Character Classification

In the second experiment, we compare FedPD with FedAvg and FedProx on the FEMNIST data set . The FEMNIST data set collects the handwritten characters, including numbers 1–10 and the upper- and lower-case letters A–Z and a–z, from different writers and is separated by the writers, therefore the data set naturally preserves non-i.i.d-ness.

The entire data set contains 805,000 samples collected from 3,550 writers. In our experiments, we use the data collected from 100 writers with an average of 300 samples per writer and the size of the whole data set is 29,214. We set the number of agent N=90N=90, the first ten agents are assigned with data from two writers, and the rest of the agents are assigned with data form one writer. Therefore, the data distribution is neither i.i.d. nor balanced. We use the neural network given in as the training model, which consists of 2 convolutional layers and two fully connected layers. The output layer has 62 neurons that matches the number of classes in the FEMNIST data set.

The numerical results shown in Fig. 6 in the main text were generated by running MATLAB codes on Amazon Web Services (AWS), with Intel Xeon E5-2686 v4 CPUs. In the training phase, we train the CNN model with FedAvg, FedProx and FedPD. In Fig. 6(a), for FedAvg, we use gradient descent for Q=8Q=8 local update steps between each communication rounds; to solve the local problem for FedProx, we use SARAH with Q=20Q=20 local steps; we use FedPD with Oracle II, computing full gradient every I=20I=20 communication rounds and perform Q=2Q=2 local steps between two communication rounds. The hyper-parameters we use for FedAvg is η=0.005\eta=0.005; for FedProx we use ρ=1\rho=1 and stepsize η=0.01\eta=0.01; for FedPD we use η=100\eta=100 and γ=400\gamma=400. In Fig. 6(b), we use FedPD with Oracle I, with Q=20Q=20, η=100\eta=100 and γ=400\gamma=400 and the mini-batch size 22. We set the communication saving to p=0p=0 and p=0.5p=0.5.

The results shown in Fig. 7 were generated by running Python codes (using the the PyTorch package PyTorch: An Imperative Style, High-Performance Deep Learning Library, https://pytorch.org/) with AMD EPYC 7702 CPUs and an NVIDIA V100 GPU.

In the training phase, we train with FedProx, FedAvg and FedPD with a total T=1000T=1000 outer iterations. The local problems are solved with SGD for Q=300Q=300 local iterations and the mini-batch size in evaluating the stochastic gradient is 22. The stepsize choice for FedAvg, FedProx and FedPD are 0.0010.001, 0.010.01 and 0.010.01, the hyper-parameter of FedProx is ρ=1\rho=1 and for FedPD η=1\eta=1. In the experiment, we set the communication saving for FedPD to be p=0p=0, p=0.5p=0.5 and p=0.25p=0.25. Note that we also tested FedAvg with larger stepsize 0.010.01, but the algorithm becomes unstable, and its performance degrages significantly. As shown in Fig. 7 and 8, FedAvg is slower than FedPD and FedProx, while FedProx has similar performance as FedPD when R=1R=1. Further, we can see that as the frequency of communication of FedPD decreases, the final accuracy decreases and the final loss increases. However, the drop of accuracy is not significant, so FedPD is able to achieve a better performance with the same number of communication rounds.