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 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:
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.
We provide a convergence analysis for FedBuff in the smooth non-convex setting. When clients take local SGD steps, FedBuff requires server iterations to reach accuracy (Section 4).
Empirically, we show that FedBuff is up to 3.8 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 more efficient than the closest asynchronous FL algorithm in the literature, FedAsync (Xie et al., 2019). Our extensive empirical evaluation finds that is a good setting across benchmarks and does not require tuning.
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 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: satisfies ()-DP if for all adjacent datasets and for any subset of outputs , the following holds:
The definition of adjacent datasets is domain and application dependent. In the context of cross-device FL, we consider to be two datasets of training examples, where each example is associated with a client. Then, and are adjacent if can be formed by adding or removing all of the examples associated with a single client from (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 is the total number of clients and the function measures the loss of a model with parameters on the th client’s data, and weighs the importance of the data from client . The goal is to find a model that fits all clients’ data well on (weighted) average. In FL, is only accessible by client .
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 client updates are in the buffer, where is the size of the buffer and is a tunable parameter. However, we find that 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 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 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: represents the set of all client indices, denotes the gradient with respect to the loss on client ’s data, denotes the minimum of , denotes the stochastic gradient on client , is the buffer size for aggregation before producing each server update, and denotes the number of local steps taken by each client. We make the following assumptions throughout.
(Bounded local and global variance) for all clients ,
(Bounded gradient) for all .
(Lipschitz gradient) for all client , the gradient is -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 ) For all clients and for each server step , the staleness between the model version in which FedBuff-client uses to start local training, and the model version in which is used to modify the global model is not larger than when .
Remark. More generally, the staleness upper-bound depends on the buffer size . 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 , the maximum delay is at most ; 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 .
Effect of staleness. The effect of staleness between the initialization of FedBuff-client and the server update dissipates at the rate of according to the fourth term in equation . In addition, the maximum staleness reduces as the buffer size grows. (see Appendix A)
Practical Improvements
Staleness scaling. To control the effect of staleness in client ’s contribution to the -th server update, we down-weight stale updates using the following function: , 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 , 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 . 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 , 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 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 more updates, and FedAsync requires 1.1-2.5 more updates.
Scalablility of FedBuff. Figure 3 shows that FedBuff scales much better to larger values of concurrency than FedAvgM. FedBuff with 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 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 . Table 2 presents the number of client trips to reach validation accuracy for different values of , with fixed concurrency. We find that =10 is a good setting across benchmarks. We analyze FedBuff with even larger values of 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 , and be competitive for lower values of . This result illustrates FedBuff’s flexibility to be adapted for privacy. Even with DP, we find that is good setting. We find that a small 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 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 , 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 more efficient than FedAvgM, and up to 2.5 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 is always bounded as . In this section we will show that this implies the staleness bound when running FedBuff with .
Consider an execution of FedBuff. Let denote the time when the ’th client update is received by the server, and let denote the time when the client downloaded the serve model before performing local steps that resulted in the model update received at .
When , the staleness of the ’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
When , the server waits to aggregate client updates before stepping the global model. Thus, if client updates are received between the times and , then at most server updates occur during this time. Hence , and therefore
Note that the times and 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 ). These times are not impacted by the choice of ; rather 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 reduces average delay. In particular, let
and suppose the limit exists. Clearly, if the limit exists and Assumption 5 holds, then . 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 pixels, then normalize by mean and 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 be the number of examples held by client , and let be the amount of time required by client to perform local training. Also assume that there is a constant such that
In this context, is the straggler parameter. The larger the , the longer the expected client training time. For a given round , let be number of clients required to close a round, and let be the total number of clients training concurrently with over-selection. If denote the raw times when clients complete the round and denote order statistic of the client training time, then for SyncFL is
If 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 be the number of clients required to reach a target accuracy and 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 , 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 . In Table 2, we observe that FedBuff trains fast when running with small values of , relative to the concurrency. However, large values of are useful when providing user-level differential privacy, as essentially the noise is divided among larger number of clients (larger values of ) (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 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 , FedBuff can achieve the same utility as SyncFL with amplified DP-SGD at the cost of slower convergence. At the same , 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 . 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 smaller than the server-prescribed batch size . This may occur when processing a batch at the end of one epoch, including the first batch if the client has fewer than 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 or global iteration index . Here denotes a stochastic gradient of (the client’s local objective) evaluated at using batch size .
Also assume that the stochastic gradients are unbiased and have variance satisfying a weak growth condition. Specifically, assume that with batch size ,
Note that in the proof of Theorem 2, we make the stronger assumption of bounded variance, corresponding to .
Furthermore, suppose that a mini-batch stochastic gradient with batch size is obtained by averaging the gradients evaluated at 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 . 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 and taking the total expectation yields
Now, multiplying both sides by , 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 can reach the target validation accuracy quicker than other values of . At , 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 denotes the set of clients that contribute to the ’th server update, and is the staleness of an update contributed by client to the ’th server update. Specifically, when , the update returned by client was computed by starting from and performing local gradient steps. When there is no staleness in the update, and more generally corresponds to some staleness; i.e., 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 is a uniform subset ; 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 , the server can only sample that client after the server has performed another update.
where is the client delta which is trained from using the global model after updates as initialization. We will next derive the upper bounds on and . To begin,
Using conditional expectation, the expectation operator can be written as
Now for , from the definition of ,
Further, by telescoping, can be decomposed as
The upper bound on 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 ,
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 becomes:
Inserting the upper bound on into (10), we have,
Now, plugging (18), (19) and (20) into (8),
Summing up from to and rearrange, yields