Federated Learning with Buffered Asynchronous Aggregation

John Nguyen, Kshitiz Malik, Hongyuan Zhan, Ashkan Yousefpour, Michael Rabbat, Mani Malek, Dzmitry Huba

Introduction

Federated Learning (FL) is a distributed learning paradigm that aims to train a shared model across participants while training data stays on the participant devices. In this work, we focus on cross-device FL where participants are edge devices (Kairouz et al. (2019)), and in particular, aim to address the following two challenges:

Challenge 1: Scalability. In large-scale cross-device FL settings, the number of clients can be in the millions, and only a small fraction of the client population may be available at any given time for training (Wang et al. (2021)). Additionally, client devices may have limited communication bandwidth and compute power. In these settings, an important parameter is concurrency: the number of clients training concurrently (i.e., clients-per-round or cohort size). There is a fundamental limitation when increasing concurrency in synchronous FL training: a diminishing return in the speed and quality of training. In this paper, we propose a novel buffered asynchronous aggregation optimization that makes it possible to train using significantly higher concurrency, improving the performance and efficiency of FL.

Challenge 2: Privacy. Inference attacks, methods trying to recover information from gradients, can expose sensitive information about the participating clients (Melis et al., 2019; Geiping et al., 2020). Given this privacy concern, secure aggregation (SecAgg) (Karl et al. (2020); Bonawitz et al. (2016)) and differential privacy (DP) (Kairouz et al. (2021); McMahan et al. (2018)) provide protection against inference attacks (Watson et al., 2021; Carlini et al., 2020). Using SecAgg, an honest-but-curious server cannot see the individual client updates, while DP can protect clients’ data from observations based on the inputs and the output of the computation. With SecAgg, DP clipping and noise addition can be performed on server, providing a better privacy-utility trade-off. For many real-world cross-device FL applications, compatibility with such privacy enhancing technologies is vital.

Our proposal: FedBuff. Motivated by these challenges, we propose and analyze FedBuff, a novel asynchronous federated optimization framework using buffered asynchronous aggregation. In FedBuff, clients train and communicate asynchronously with the server. Unlike other asynchronous methods, the server aggregates KK client updates in a secure buffer before performing a server update. This secure buffer can be implemented by using Trusted Execution Environments (TEEs) (Karl et al., 2020; Mo et al., 2021).

Contributions. We highlight the main contributions:

∙\bullet We propose FedBuff, a novel asynchronous federated optimization framework with buffered asynchronous aggregation to achieve scalability and privacy against the honest-but-curious threat model through secure aggregation and differential privacy.

∙\bullet We provide a convergence analysis for FedBuff in the smooth non-convex setting. When clients take QQ local SGD steps, FedBuff requires O(1/(ϵ2Q))\mathcal{O}\left(1/(\epsilon^{2}Q)\right) server iterations to reach ϵ\epsilon accuracy (Section 4).

∙\bullet Empirically, we show that FedBuff is up to 3.8×\times more efficient than competing synchronous FL algorithms, even without penalizing synchronous FL algorithms for stragglers. We also demonstrate that FedBuff is up to 2.5×\times more efficient than the closest asynchronous FL algorithm in the literature, FedAsync (Xie et al., 2019). Our extensive empirical evaluation finds that K=10K=10 is a good setting across benchmarks and does not require tuning.

∙\bullet To the best of our knowledge, we are the first to propose an asynchronous federated optimization framework that is compatible with SecAgg and global user-level DP. Under differentially private training, FedBuff can outperform both synchronous FL with amplified DP-SGD and DP-FTRL (differentially private Follow-the-Regularized-Leader) at low privacy settings, and be competitive for high privacy settings.

Background

Synchronous FL. Significant attention has been paid towards synchronous FL methods (SyncFL), as they are perhaps easier to analyze and implement. SyncFL methods are also better suited for privacy – training and aggregating updates over a large number of clients render most inference attacks ineffectual (Melis et al., 2019; Zhu and Han, 2020; Geiping et al., 2020; Lam et al., 2021). However, synchronous FL methods are prone to stragglers, proceeding at the pace of the slowest client. Bonawitz et al. (2019) proposed using over-selection to tap 30% more clients than the target cohort size and wait for the fastest replies to overcome this issue. However, over-selection comes at the cost of wasting clients’ resources and introduces selection bias. We study these problems in Appendix C.2 and C.3.

In SyncFL optimization, FedAvg, a generalization of local SGD, has been shown to work well empirically (McMahan et al., 2016). FedProx (Li et al., 2018) improves upon FedAvg by adding a proximal term μ\mu to the local SGD optimizer. FedAvgM (Hsu et al., 2019) further improves convergence by adding server-side momentum. Adaptive methods such as FedAdam (Reddi et al., 2020) are effective in cross-device FL settings and have comparable performance to FedAvgM. These optimizers often focus on heterogeneity, asymptotic convergence, and communication efficiency in low concurrency settings. In this paper, we focus on high concurrency settings. In these settings, SyncFL is not scalable and is inefficient. Typically, the optimal server learning rate increases with concurrency; aggregating over more users has a variance-reducing effect, enabling the server to take larger steps. Consequently, higher concurrency reduces the number of rounds needed to reach a target accuracy because of a larger server learning rate. However, to have stable, convergent training dynamics, the server learning rate cannot be increased indefinitely; eventually it saturates, resulting in a sub-linear speed-up similar to in large-batch training (Goyal et al. (2017); Ott et al. (2018); You et al. (2019, 2018, 2017); Shallue et al. (2018)). As a result, SyncFL systems cannot accelerate training through parallelism beyond a few hundred clients and exhibit decreasing efficiency with increasing concurrency (Figure 1).

Asynchronous FL. Asynchronous FL methods are a good match for cross-device FL settings, where clients have different compute power and intermittent availability (Wang et al., 2021). Most asynchronous FL (AsyncFL) works, such as the works (Xie et al., 2019; van Dijk et al., 2020; Chai et al., 2020; Chen et al., 2019; Wu et al., 2020; Li et al., 2021) have been focused on solving the straggler problem by designing asynchronous FL algorithms. However, these proposals include aspects that make them impractical for real-world FL deployment at scale. For instance, Chai et al. (2020); Li et al. (2021) profile client speed, Chen et al. (2019) broadcast the model updates to all clients, Xie et al. (2019) update the server model on every client, placing a significant burden on the clients and server, and van Dijk et al. (2020) assume all clients have the same speed.

In “fully” AsyncFL methods (e.g., Xie et al. (2019)), every client update results in a server model update. This has implications for privacy and scalability. Considering privacy, when every client update forces a server update, SecAgg cannot be used; secure aggregation’s benefit is in hiding individual updates by combining them in an aggregate. Additionally, providing user-level DP in AsyncFL is only feasible with local differential privacy (LDP), where the client clips the model update and adds noise locally to it before sending it to the server. LDP for high dimensional data has been criticized for poor privacy-utility trade-off (Erlingsson et al. (2020); Bittau et al. (2017)).

Secure Aggregation. SecAgg is a privacy enhancing technology based on cryptographic primitives (Bonawitz et al., 2016; Bell et al., 2020; So et al., 2021b) or hardware-based Trusted Execution Environment (TEE) (Karl et al., 2020). SecAgg enhances privacy by obfuscating a client’s update with many other clients’ updates, protecting against the honest-but-curious server threat model (Bonawitz et al., 2016). FedBuff is compatible with SecAgg.

Differential Privacy. DP (Dwork et al., 2014) provides a rigorous formulation of the release of information derived from private data. In the context of machine learning, differentially private training (Abadi et al., 2016) limits what can be learned about the original training data.

A randomized mechanism M: U↦RU\mapsto R satisfies (ϵ,δ\epsilon,\delta)-DP if for all adjacent datasets D,D′∈UD,D^{\prime}\in U and for any subset of outputs S⊆RS\subseteq R, the following holds:

The definition of adjacent datasets D,D′D,D^{\prime} is domain and application dependent. In the context of cross-device FL, we consider D,D′D,D^{\prime} to be two datasets of training examples, where each example is associated with a client. Then, DD and D′D^{\prime} are adjacent if D′D^{\prime} can be formed by adding or removing all of the examples associated with a single client from DD (i.e., user-level privacy (McMahan et al., 2018)). In this paper, we consider the global DP (GDP) setting, where a trusted server collects, clips, and aggregates the client updates. The server then adds noise to the aggregated updates. Compared to LDP, GDP provides a better privacy-utility trade-off for high dimensional data. This setting of DP relies on using SecAgg, as the server is responsible for implementing DP.

FedBuff: Federated Learning with Buffered Asynchronous Aggregation

We consider the following optimization problem:

where mm is the total number of clients and the function FiF_{i} measures the loss of a model with parameters ww on the iith client’s data, and pi>0p_{i}>0 weighs the importance of the data from client ii. The goal is to find a model that fits all clients’ data well on (weighted) average. In FL, FiF_{i} is only accessible by client ii.

SyncFL methods need to aggregate and synchronize clients after each round. Hence, concurrency in SyncFL is equal to the number of clients that participate in a given round. In asynchronous methods, concurrency is the number of clients training at a given point in time (Figure 2). In FedBuff (Algorithm 1), clients enter and finish local training asynchronously. However, the server model is not updated immediately upon receiving every client update. Instead, client updates are stored in a buffer. A server update only takes place once KK client updates are in the buffer, where KK is the size of the buffer and is a tunable parameter. However, we find that K=10K=10 is a good choice and does not require tuning. The buffer can be implemented by using a Trusted Execution Environment (TEE) (Karl et al., 2020; Mo et al., 2021) or through a cryptographic algorithm (So et al., 2021a). Note that KK is independent of concurrency — the extra degree of freedom introduced by the buffer allows the server to choose the model update frequency instead of coupling concurrency with the server model update as in SyncFL. The extra degree of freedom allows FedBuff to achieve data efficiency at high concurrency while being compatible with secure aggregation and DP.

FedBuff is compatible with SecAgg, because with K>1K>1 updates in the buffer, SecAgg provides its promise by hiding individual updates in the aggregate. Since FedBuff supports SecAgg, it can be easily extended to provide global DP. In asynchronous FL settings, the server has no control over which clients participate in a particular model update and client availability is dynamic. For such settings, privacy amplification by sampling is not feasible. DP-FTRL (Kairouz et al., 2021) has emerged as a suitable solution to address this issue. FedBuff with DP-FTRL is straightforward and we show in Algorithm 1 how one can extend FedBuff to provide global DP. The three functions, InitializeTree, AddToTree, and GetSum in Algorithm 1 correspond to those of the DP-FTRL algorithm. We defer to Section B.1 in (Kairouz et al., 2021) for more in-depth descriptions of the functions.

Convergence Analysis

In this section, we provide a convergence guarantee for FedBuff in the smooth, non-convex setting. Most previous works analyze synchronous federated learning methods, such as the works in (Lin et al., 2018; Li et al., 2018; Reddi et al., 2020; Li et al., 2020; Stich, 2019; Yu et al., 2019b; Li et al., 2019; Haddadpour and Mahdavi, 2019; Karimireddy et al., 2020). In contrast, in FedBuff, clients train asynchronously, and the client updates are first aggregated in a buffer before producing a server model update. Hence, it is essential to understand the relationship between client computation and server communication under asynchrony with buffered aggregation.

Notation. We use the following notation throughout: [m][m] represents the set of all client indices, ∇Fi(w)\nabla F_{i}(w) denotes the gradient with respect to the loss on client ii’s data, f∗f^{*} denotes the minimum of f(w)f(w), gi(w;ζi)g_{i}(w;\zeta_{i}) denotes the stochastic gradient on client ii, KK is the buffer size for aggregation before producing each server update, and QQ denotes the number of local steps taken by each client. We make the following assumptions throughout.

(Bounded local and global variance) for all clients i∈[m]i\in[m],

(Bounded gradient) ∥∇Fi∥2≤G\left\lVert\nabla F_{i}\right\rVert^{2}\leq G for all i∈[m]i\in[m].

(Lipschitz gradient) for all client i∈[m]i\in[m], the gradient is LL-smooth,

Assumptions 1–4 are commonly made in analyzing federated learning algorithms (Reddi et al. (2020); Li et al. (2020); Stich (2019); Yu et al. (2019b)). We make an additional assumption on the staleness under asynchrony.

(Bounded Staleness when K=1K=1) For all clients i∈[m]i\in[m] and for each server step tt, the staleness τi(t)\tau_{i}(t) between the model version in which FedBuff-client uses to start local training, and the model version in which Δi\Delta^{i} is used to modify the global model is not larger than τmax⁡,1\tau_{\max,1} when K=1K=1.

Remark. More generally, the staleness upper-bound depends on the buffer size KK. When the buffer size increases, the server iterates are updated less frequently, hence reducing the number of server steps in between the initialization of client training and when the client updates are used for modifying the server model. Specifically, if Assumption 5 holds, then for any execution of FedBuff with K>1K>1, the maximum delay τmax⁡,K\tau_{\max,K} is at most ⌈τmax⁡,1/K⌉\lceil\tau_{\max,1}/K\rceil; see Appendix A.

The proof of Theorem 2 is provided in Appendix D, and leverages ideas from the perturbed iterate framework (Mania et al., 2017).

where we use the relation τmax⁡,K≤⌈τmax⁡,1/K⌉\tau_{\max,K}\leq\lceil\tau_{\max,1}/K\rceil.

Effect of staleness. The effect of staleness between the initialization of FedBuff-client and the server update dissipates at the rate of O(1/T)\mathcal{O}\left(1/T\right) according to the fourth term in equation (\refeq:constantLRbound)(\ref{eq:constantLR_bound}). In addition, the maximum staleness τmax⁡,K\tau_{\max,K} reduces as the buffer size KK grows. (see Appendix A)

Practical Improvements

Staleness scaling. To control the effect of staleness τi(t)\tau_{i}(t) in client ii’s contribution to the tt-th server update, we down-weight stale updates using the following function: s(τi(t)):=1/(1+τi(t))0.5s(\tau_{i}(t)):=1/(1+\tau_{i}(t))^{0.5}, similar to (Xie et al. (2019)).

Experiments

In this section, we compare the efficiency and scalability of FedBuff with other synchronous and asynchronous FL methods from the literature via simulation. We wish to understand how FedBuff behaves under different values of KK, its scalability, and data efficiency.

Evaluation metrics. The standard evaluation metric for FL is the number of communication rounds to reach a target accuracy. However, asynchronous and synchronous methods do not have the same notion of rounds. For this reason, we compare different synchronous and asynchronous methods by the number of client trips needed to reach a target accuracy. One client trip corresponds to one client round-trip communication. A client trip involves a client pulling the latest model from the server (download communication), performing one epoch of training on the local dataset (computation), then communicating the model update to the server (upload communication). Since the number of client trips measures both communication and computation costs, we use this as a proxy for wall-clock training time. We show wall-clock time simulation with stragglers in Appendix C.2.

Datasets, models, and tasks. In order to provide a comparison with other work in the literature, we run experiments on three datasets: CelebA (Liu et al. (2015)), Sent140 (Go et al. (2009)), and CIFAR-10 (Krizhevsky et al. (2009)). Sent140 is a text classification dataset (binary sentiment analysis), whereas CelebA and CIFAR-10 are image classification datasets (multi-class classification). For Sent140 and CelebA, we use the natural non-iid client partitions and models from the LEAF benchmark (Caldas et al. (2018)). For Sent140, we train an LSTM classifier over 660,120 clients, where each Twitter account corresponds to one client. For CelebA, we train the same convolutional neural network classifier as LEAF over 9,343 clients, but with batch normalization layers replaced by group normalization layers (Wu and He (2018)) as suggested in (Hsieh et al. (2020)). For CIFAR-10, we generate 5000 non-iid clients using a Dirichlet distribution with parameter 0.1, the same approach as in (Hsu et al. (2019)). More details about datasets, models, and tasks are provided in Appendix B.1.

Experimental setup. We implement all algorithms in PyTorch (Paszke et al. (2017)). We repeat each experiment with three different seeds and report the average. For asynchronous FL methods, we assume that clients arrive at a constant rate. We sample the delay distribution, the time delay between a client’s download and upload operation, from a half-normal distribution. We choose this distribution because it best matches the delay distribution observed in our production FL system (See Appendix C.1). We also report results with two other delay distributions (uniform and exponential) in Appendix C.1. We find that FedBuff’s performance improvements are consistent across different delay distributions.

Baselines. We compare FedBuff with three SyncFL baselines, namely FedAvg (McMahan et al. (2016)), FedProx (Li et al. (2018)), FedAvgM (Hsu et al. (2019)), and one AsyncFL baseline, FedAsync (Xie et al. (2019)). For more details about the algorithms used and the experimental setup, see Appendix B.2.

Concurrency and KK. In at-scale cross-device FL, only a small fraction of all clients participate in training at any point in time. As discussed earlier, concurrency — the maximum number of clients that train in parallel — significantly impacts the performance of FL algorithms. For a fair comparison between synchronous and asynchronous algorithms, we keep concurrency the same across all configurations. Recall the example in Figure 2 where concurrency=100. For synchronous algorithms, this implies that 100 clients are training and contributing in each round. For asynchronous algorithms, this implies that 100 clients can train concurrently, and we can still vary the buffer size KK, which will control how frequently updates occur.

Comparison of Methods. Table 1 shows the number of client trips needed to converge to the target accuracy on Sent140, CelebA and CIFAR-10 for each method considered. In Table 2, we show results with other values of KK and present the learning curves in Appendix C.7. Compared to FedBuff, the best synchronous method in the experiments (FedAvgM) requires 1.7-3.3×\times more updates, and FedAsync requires 1.1-2.5×\times more updates.

Scalablility of FedBuff. Figure 3 shows that FedBuff scales much better to larger values of concurrency than FedAvgM. FedBuff with K=10K=10 scales better because it updates the server model more frequently than FedAvgM in high concurrency. When concurrency is 10, both FedAvgM and FedBuff update the server model after every 10 client updates. However, when concurrency is 1000, FedBuff with K=10K=10 updates the server model after every 10 client updates, while FedAvgM updates the server model after 1000 client updates. One might argue that FedAvgM should run at lower concurrency, e.g. 10. However, that leads to longer wall-clock training time because less parallelism is exploited. We discuss this problem in Appendix C.2. For synchronous FL methods, larger concurrency reduces training time but is also less efficient. On the other hand, taking server model steps more frequently is not free; FedBuff has to deal with staleness as a consequence. Our empirical results show that the benefits from frequent updating of the server model outweigh the cost of staleness in client model updates.

Choice of KK. Table 2 presents the number of client trips to reach validation accuracy for different values of KK, with fixed concurrency. We find that KK=10 is a good setting across benchmarks. We analyze FedBuff with even larger values of KK in Table 6, and show the training curves of FedBuff and other algorithms in Appendix C.7.

FedBuff with Differential Privacy. To evaluate the privacy-utility trade-off of FedBuff, we compare the final test accuracy of FedBuff with synchronous baselines after 600 thousands client trips, one pass over the dataset. Figure 4 illustrates that FedBuff can outperform both FedAvgM with amplified DP-SGD and FedAvgM with DP-FTRL at high values of ϵ\epsilon, and be competitive for lower values of ϵ\epsilon. This result illustrates FedBuff’s flexibility to be adapted for privacy. Even with DP, we find that K=10K=10 is good setting. We find that a small LL can counteract the additional noise from taking more steps. For more details see Appendix C.5.

Related Work

In addition to the discussion in Section 2 on related works, we discuss the other efforts in related domains.

Asynchronous stochastic optimization. Asynchronous stochastic optimization in shared-memory and distributed-memory systems has been extensively studied (Bertsekas and Tsitsiklis (1989); Chaturapruek et al. (2015); Niu et al. (2011); Lian et al. (2015, 2018); Chen et al. (2016); Zheng et al. (2017); Mania et al. (2017); Leblond et al. (2017); Reddi et al. (2015); Assran et al. (2020)). Asynchronous training is resilient to stragglers in both centralized and federated settings. The idea of aggregating KK asynchronous updates for convex objectives has been studied in (Dutta et al., 2018). Although Dutta et al. (2021) provide a guarantee for non-convex objectives, their assumption on the relationship between staleness and gradient moments is difficult enforce, in contrast to the bounded staleness assumption we consider which can be easily enforced. In this work, we consider heterogeneous objectives and show that in a federated environment with a large number of clients, the source of speed-up is not only due to avoiding stragglers but also achieving better efficiency at high concurrency.

Large-batch training. Many proposals aim to understand and characterize conditions under which linear speed-up for distributed SGD and local SGD is achievable (Lin et al. (2018); Yu et al. (2019a); Woodworth et al. (2020); Haddadpour et al. (2019)). It is well accepted that increasing concurrency eventually saturates beyond a certain batch size in synchronous methods (Yin et al. (2017); Ott et al. (2018); Goyal et al. (2017); Ott et al. (2018); You et al. (2019, 2018, 2017); Shallue et al. (2018)). However, most existing research focuses on scalability across tens of server workers, each having iid-data - very different from the FL setting.

Conclusions

In this paper, we propose FedBuff, an asynchronous FL training scheme with buffered aggregation. Compared to SyncFL proposals, FedBuff scales to large values of concurrency. Compared to AsyncFL proposals, FedBuff is more private as it is compatible with SecAgg and differential privacy. At high levels of ϵ\epsilon, we demonstrate that FedBuff can outperform major SyncFL proposals. We analyze the convergence behavior of FedBuff in the non-convex setting. Empirical evaluation shows that FedBuff is up to 3.3×\times more efficient than FedAvgM, and up to 2.5×\times more efficient than FedAsync. As for future work, we are aware that our analyses is on standard SGD. We leave extending the analysis to include momentum or adaptive learning rates as future work.

We would like to thank Ilya Mironov, Maziar Sanjabi, Graham Cormode, Samuel Horvath and Luca Melis for the meaningful discussions and their valuable suggestions which significantly improved the quality of this paper. We would like to also thank the anonymous reviewers for their insightful feedback.

References

Appendix

Appendix A Relationship Between Maximum Staleness and K𝐾K

Recall Assumption 5, that the staleness when executing FedBuff with K=1K=1 is always bounded as τi(t)≤τmax⁡,1\tau_{i}(t)\leq\tau_{\max,1}. In this section we will show that this implies the staleness bound τi(t)≤⌈τmax⁡,K⌉≤⌈τmax⁡,1/K⌉\tau_{i}(t)\leq\lceil\tau_{\max,K}\rceil\leq\lceil\tau_{\max,1}/K\rceil when running FedBuff with K>1K>1.

Consider an execution of FedBuff. Let rir_{i} denote the time when the ii’th client update is received by the server, and let si<ris_{i}<r_{i} denote the time when the client downloaded the serve model before performing local steps that resulted in the model update received at rir_{i}.

When K=1K=1, the staleness τi(1)\tau_{i}^{(1)} of the ii’th update corresponds to the number of updates that occurs between when the client downloaded the model and when it completed local training and uploaded the model to the server,

If Assumption 5 holds, then max⁡iτi(1)≤τmax⁡,1.\max_{i}\tau_{i}^{(1)}\leq\tau_{\max,1}.

When K>1K>1, the server waits to aggregate KK client updates before stepping the global model. Thus, if τi(K=1)\tau_{i}^{(K=1)} client updates are received between the times sis_{i} and rir_{i}, then at most τi(1)/K\tau_{i}^{(1)}/K server updates occur during this time. Hence τi(K)≤⌈τi(1)/K⌉\tau_{i}^{(K)}\leq\lceil\tau_{i}^{(1)}/K\rceil, and therefore

Note that the times sis_{i} and rir_{i} only depend on the number of clients training concurrently, and the distribution of client execution times (the time it takes a client to complete one round of local updates; i.e., the distribution of ri−sir_{i}-s_{i}). These times are not impacted by the choice of KK; rather KK only affects how frequently the server performs an update. Thus, the arguments above hold regardless of the distribution of client execution times, and only depend on Assumption 5.

We also remark that the same relationship holds for the average delay; i.e., increasing KK reduces average delay. In particular, let

and suppose the limit exists. Clearly, if the limit exists and Assumption 5 holds, then τ‾1≤τmax⁡,1\overline{\tau}_{1}\leq\tau_{\max,1}. Furthermore, then

Appendix B Experiment Details

Sent140. We train a sentiment classifier on tweets from the Sent140 dataset (Caldas et al., 2018; Go et al., 2009) with a two-layer LSTM binary classifier. The dataset has 660,120 clients where each client is a Twitter account. The LSTM binary classifier contains 100 hidden units with a top 10,000 pretrained word embedding from 300D GloVe (Pennington et al., ). The model has a max sequence length of 25 characters. The model first embeds each of the characters into a 300-dimensional space by looking up GloVe, passes through 2 LSTM layers and a 128 hidden unit linear layer to output labels 0 or 1. We set the dropout rate to 0.1. We split the data into 80% training set, 10% validation set, and 10% test set using script provided by Caldas et al. (2018). Due to memory constraint, we use 15% of the entire dataset using the script provided by Caldas et al. (2018), with split seed = 1549775860.

CelebA. We study an image classification problem on the CelebA dataset (Liu et al., 2015; Caldas et al., 2018) using a four layer CNN binary classifier with dropout rate of 0.1, stride of 1, and padding of 2. As it is standard with image datasets, we preprocess train, validation, and test images; we resize and center crop each image to 32×3232\times 32 pixels, then normalize by 0.50.5 mean and 0.50.5 standard deviation. The dataset has 9,343 clients where each client is a unique celebrity.

CIFAR-10. We evaluate a multi-class image classification problem on CIFAR-10 (Krizhevsky et al., 2009) using a four layer CNN binary classifier with dropout rate of 0.1, stride of 1, and padding of 2. We normalize the images by the dataset mean and standard deviation. Following Hsu et al. (2019), we partition the dataset into 5,000 clients using a Dirichlet distribution with parameter 0.1 and split seed = 0.

B.2 Implementation Details

We implemented all algorithms in Pytorch (Paszke et al., 2017) and evaluated them on a cluster of machines, each with eight NVidia V100 GPUs. Independently, we built a simulator to simulate large-scale federated learning environments. The simulator can realistically simulate clients, server, communication channels between clients and server, model aggregation schemes, and local training of clients. We intend to open-source the simulator, making it available for the research community.

For our experiments, we assume clients arrive to the FL system at a constant rate. To simulate device heterogeneity, we sample each client training duration from a half-normal, uniform, or exponential distribution. Moreover, our implementation has two other important distinctions. First, each client does one epoch of training over its local data; this distinction stems from two observations in our production stack: that our FL production stack has plenty of users to train on, and that we train small capacity models in FL (e.g., less than 10 million parameters) because of bandwidth and client compute. Second, we use the weighted sum of the client updates instead of the weighted average. This is because each client update has different levels of staleness; taking the average cannot capture the true contribution for each client.

B.3 Hyperparameters

For all experiments, we tune hyperparameters using Bayesian optimization (Snoek et al., 2012). For optimizer on clients, we use minibatch SGD for all tasks. We select the best hyperparameters based on the number of rounds to reach target validation accuracy for each dataset.

B.3.2 Best Performing Hyperparameters

Appendix C Additional Experiments

In this section, we analyze the sensitivity of FedBuff to different staleness distributions. We compare FedBuff against other competing algorithms with different staleness distributions. Table 5 demonstrates that FedBuff is robust and FedBuff’s speed up is consistent. To have an accurate view of real-world delays, we observe the delays and their resulting staleness distribution in our production stack when training over millions of clients with concurrency = 1000 and K = 100. Figure 5 demonstrates that a half-normal is a suitable delay distribution.

C.2 Wall-Clock Time Simulation

In this section, we study the speed up of FedBuff over FedAvgM in terms of wall-clock time for various concurrency levels. The results for Sent140 are in Figure 6.

In cross-device FL, each client is a mobile phone with limited compute power and communication bandwidth (Kairouz et al., 2019). Moreover, clients can have vastly different number of examples. Recall, SyncFL methods wait for all the participating clients in a round to finish before updating the server model – a round proceeds at the pace of the slowest client, the straggler effect. To mitigate the straggler problem, over-selection proposed in Bonawitz et al. (2019), which selects 30% more clients than the target number of clients to participate and waits for the fastest replies.

To confirm the speedup gain by FedBuff in terms of wall-clock time, we simulate training time of FedAvgM and FedBuff using a random exponential time model from Lee et al. (2017). The random exponential time model has been widely used to simulate the straggler effect in federated learning, e.g. in Reisizadeh et al. (2020); Tandon et al. (2017); Yu et al. (2020, 2019c); Charles et al. (2021).

We assume the time a client requires to perform local training is proportional to the number of examples the client has, same as the assumption in Charles et al. (2021). Formally, let nin_{i} be the number of examples held by client ii, and let TiT_{i} be the amount of time required by client ii to perform local training. Also assume that there is a constant λ>0\lambda>0 such that

In this context, λ\lambda is the straggler parameter. The larger the λ\lambda, the longer the expected client training time. For a given round tt, let CtC_{t} be number of clients required to close a round, and let MM be the total number of clients training concurrently with over-selection. If T1,…,TMT_{1},\dots,T_{M} denote the raw times when clients complete the round and T(1)≤⋯≤T(M)T_{(1)}\leq\dots\leq T_{(M)} denote order statistic of the client training time, then RtR_{t} for SyncFL is

If TT is the number of rounds to reach a target accuracy, the expected total training time for SyncFL is

In the case of FedBuff, the total wall-clock time is when the last client required to reach a target accuracy finishes training. Formally, let NN be the number of clients required to reach a target accuracy and T(1)≤⋯≤T(N)T_{(1)}\leq\dots\leq T_{(N)} denotes the order statistic of the client training time. Then the expected total training time for FedBuff is

C.3 Bias

In this section, we show the diversity of participating clients in FedBuff and FedAvgM. We should that SyncFL methods can introduce bias in their selection process while FedBuff does not. We use the same random exponential time model in Appendix C.2. The result is given in Figure 7. We find that in all levels of λ\lambda, FedBuff can incorporate clients with large local datasets. On the other hand, SyncFL (e.g., FedAvgM) with over-selection drops these clients leading to bias in the selection process. FedBuff does not drop the slowest clients, and can incorporate these clients with many examples.

C.4 Large values of K𝐾K

In this section, we study the performance of FedBuff with large values of KK. In Table 2, we observe that FedBuff trains fast when running with small values of KK, relative to the concurrency. However, large values of KK are useful when providing user-level differential privacy, as essentially the noise is divided among larger number of clients (larger values of KK) (Kairouz et al., 2021; McMahan et al., 2018).

We compare the training speed of FedBuff and FedAvgM in a setting where both algorithms produce a server update from the same number of aggregated client updates. We fix concurrency at 1000, and have both FedAvgM and FedBuff perform updates after aggregating responses from K=1000K=1000 clients. In this setting, FedBuff’s main advantage is robustness to stragglers. It cannot take advantage of frequent server updates, yet still needs to deal with staleness.

The synchronous FL system described in Bonawitz et al. (2019) uses over-selection, typically by 30%, to address stragglers. For example, if 1000 users are needed to produce a server model update, 1300 users are selected. The round will finish when the fastest 1000 users finish training. Results from the slowest 300 users will be thrown away. Over-selection makes synchronous FL more robust to stragglers, but at the cost of wasting some clients’ compute and bandwidth.

Table 6 reports the wall-clock training time and number of client trips to reach target accuracy for FedBuff and FedAvgM with and without over-selection. We assume a half-normal training duration distribution since that matches the behavior observed in our production system (see Figure 5). We find that over-selection reduces the impact of stragglers significantly. However, even with over-selection, FedBuff is 25%-41% faster than FedAvgM, despite using 30% lower concurrency.

C.5 FedBuff with Differential Privacy

Figure 8 shows the training curves of FedBuff with DP-FTRL, SyncFL with DP-SGD and DP-FTRL. At low values of ϵ\epsilon, FedBuff can achieve the same utility as SyncFL with amplified DP-SGD at the cost of slower convergence. At the same ϵ\epsilon, FedBuff with DP-FTRL achieves better utility and faster convergence compared to SyncFL with DP-FTRL. We find the source for this speed-up is from FedBuff’s ability to tolerate much lower clipping norm value LL. We repeat each experiment for 3 different seeds and take the average. The seeds are 0, 1, and 2.

C.6 Learning Rate Normalization (LR-Norm)

Recall that LR-Norm described in Section 5 aims to address the situation where a client performing local updates may need to perform an update using a batch size bb smaller than the server-prescribed batch size BB. This may occur when processing a batch at the end of one epoch, including the first batch if the client has fewer than BB samples in total. Since this only pertains to the local updates performed at clients, let us simply write such an update as

without referring to any specific client index ii or global iteration index tt. Here gq(bq)g_{q}^{(b_{q})} denotes a stochastic gradient of FF (the client’s local objective) evaluated at yqy_{q} using batch size bqb_{q}.

Also assume that the stochastic gradients are unbiased and have variance satisfying a weak growth condition. Specifically, assume that with batch size bq=1b_{q}=1,

Note that in the proof of Theorem 2, we make the stronger assumption of bounded variance, corresponding to M=0M=0.

Furthermore, suppose that a mini-batch stochastic gradient gq(n)g_{q}^{(n)} with batch size bq>1b_{q}>1 is obtained by averaging the gradients evaluated at bqb_{q} independent and identically distributed samples. Thus,

then it is well-known that the SGD iterates satisfy

see, for example, Theorem 4.8 in L. Bottou, F. Curtis, and J. Nocedal, “Optimization methods for large-scale machine learning,” SIAM Review, 2019.

Non-uniform batch sizes. Now suppose that some steps will use batch size 1<bq≤B1<b_{q}\leq B. In this case one can show the following result.

Consider updates as in equation 5 with per-iteration batch size

From the weak growth assumption, it follows that

Summing both sides over q=1,…,Qq=1,\dots,Q and taking the total expectation yields

Now, multiplying both sides by 2/AQ2/A_{Q}, we obtain

C.6.2 Empirical Evaluation

In Table 7, we compare LR-Norm against two other weighting schemes: Example Weight where the weight is the number of training examples for each client, and Uniform Weight where all clients have weight of 1. We see that LR-Norm performs competitively on CelebA. For CelebA, all weighting schemes, Uniform, Example, and LR-Norm perform similarly. This is because all clients in CelebA have one batch of data and number of examples per client is fairly centered around the mean. On the other hand, LR-Norm significantly outperforms Example Weight and Uniform Weight on Sent140. LR-Norm is beneficial when there is a high degree of data imbalance across clients, as in Sent140. Sent140 is more representative of real world FL applications where there is a long tail in the number of examples and number of batches per client.

C.7 Learning Curves

In this section, we show the learning curves for each algorithm in Figures 9, 10, and 11. These figures demonstrate FedBuff’s robustness to different staleness distributions. Synchronous FL algorithms, FedAvgM, FedAvg and FedProx, are unaffected by the change in staleness distribution because they simply wait for all clients in the round.

For both CelebA and Sent140, FedBuff with K=10K=10 can reach the target validation accuracy quicker than other values of KK. At K=10K=10, FedBuff appears to have the optimal balance between speed and variance reduction.

Appendix D Proof of Convergence Rate

In this appendix, we prove the main convergence result for FedBuff. A summary of the notation used is provided in Table 8.

Observe that FedBuff updates can be described succinctly as

where St\mathcal{S}^{t} denotes the set of clients that contribute to the tt’th server update, and τk(t)≥1\tau_{k}(t)\geq 1 is the staleness of an update contributed by client kk to the tt’th server update. Specifically, when k∈Stk\in\mathcal{S}^{t}, the update returned by client kk was computed by starting from wt−τk(t)w^{t-\tau_{k}(t)} and performing QQ local gradient steps. When τk(t)=1\tau_{k}(t)=1 there is no staleness in the update, and more generally τk(t)>1\tau_{k}(t)>1 corresponds to some staleness; i.e., t−τk(t)t-\tau_{k}(t) server updates have taken place between when the client last pulled a model from the server and when the client’s update is being incorporated at the server.

In addition to the assumptions stated in Section 4, in the proof below we assume that St\mathcal{S}^{t} is a uniform subset [n][n]; i.e., in any given round any client is equally likely to contribute. This can be justified in practice as follows. To avoid having any client contribute more than once to any update, after the client returns an update contributing to Δ‾t\overline{\Delta}^{t}, the server can only sample that client after the server has performed another update.

where Δkt−τk\Delta_{k}^{t-\tau_{k}} is the client delta which is trained from using the global model after t−τkt-\tau_{k} updates as initialization. We will next derive the upper bounds on T1T_{1} and T2T_{2}. To begin,

Using conditional expectation, the expectation operator can be written as

Now for T3T_{3}, from the definition of f(wt)f(w^{t}),

Further, by telescoping, T3T_{3} can be decomposed as

The upper bound on T3T_{3} can be understood as sums of bounds on the effect of staleness and local drift during client training, and local variance induced by client-side SGD. Further, we need to produce an upper bound on the staleness of initial model from which the client models are trained.

Taking the expectation in terms of H\mathcal{H},

where the last inequality follows from the assumption on maximal delay and applying Lemma 1. Similarly, the local drift term can be upper-bounded by

Thus, the upper bound on T3T_{3} becomes:

Inserting the upper bound on T3T_{3} into (10), we have,

Now, plugging (18), (19) and (20) into (8),

Summing up tt from 11 to TT and rearrange, yields