TiFL: A Tier-based Federated Learning System

Zheng Chai, Ahsan Ali, Syed Zawad, Stacey Truex, Ali Anwar, Nathalie Baracaldo, Yi Zhou, Heiko Ludwig, Feng Yan, Yue Cheng

Introduction

Modern mobile and IoT devices (such as smart phones, smart wearable devices, smart home devices) are generating massive amount of data every day, which provides opportunities for crafting sophisticated machine learning (ML) models to solve challenging AI tasks (He et al., 2016). In conventional high-performance computing (HPC), all the data is collected and centralized in one location and proceed by supercomputers with hundreds to thousands of computing nodes. However, security and privacy concerns have led to new legislation such as the General Data Protection Regulation (GDPR) (Tankard, 2016) and the Health Insurance Portability and Accountability Act (HIPAA) (O’herrin et al., 2004) that prevent transmitting data to a centralized location, thus making conventional high performance computing difficult to be applied for collecting and processing the decentralized data. Federated Learning (FL) (McMahan et al., 2017) shines light on a new emerging high performance computing paradigm by addressing the security and privacy challenges through utilizing decentralized data that is training local models on the local data of each client (data parties) and using a central aggregator to accumulate the learned gradients of local models to train a global model. Though the computing resource of individual client may be far less powerful than the computing nodes in conventional supercomputers, the computing power from the massive number of clients can accumulate to form a very powerful “decentralized virtual supercomputer”. Federated learning has demonstrated its success in a range of applications. From consumer-end devices such as GBoard (Hard et al., 2018; Yang et al., 2018) and keyword spotting (Leroy et al., 2019) to pharmaceuticals (https://cordis.europa.eu/project/id/831472?WT.mc_id=RSS-Feed&WT.rss_f=project&WT.rss_a=223634&WT.rss_ev=a, 2020), medical research (Courtiol et al., 2019), finance (https://finance.yahoo.com/news/webank-swiss-signed-cooperation-mou-112300218.html, 2020) , and manufacturing (Hao et al., 2019). There has also been a rise of FL tools and framework development, such as Tensorflow Federated (https://www.tensorflow.org/federated, 2020), LEAF (Caldas et al., 2018b), PaddleFL (https://github.com/PaddlePaddle/PaddleFL, 2020) and PySyft (Ryffel et al., 2018) to facilitate these demands. Depending on the usage scenarios, FL is usually categorized into cross-silo FL and cross-device FL (Kairouz et al., 2019). In cross-device FL, the clients are usually a massive number (e.g., up to 101010^{10}) of mobile or IoT devices with various computing and communication capacities (McMahan et al., 2016; Kairouz et al., 2019; Konečnỳ et al., 2016) while in cross-silo FL, the clients are a small number of organizations with ample computing power and reliable communications (Yang et al., 2019; Kairouz et al., 2019). In this paper, we focus on the cross-device FL (for simplicity, we call it FL in the following), which intrinsically pushes the heterogeneity of computing and communication resources to a level that is rarely found in datacenter distributed learning and cross-silo FL. More importantly, the data in FL is also owned by clients where the quantity and content can be quite different from each other, causing severe heterogeneity in data that usually does not appear in datacenter distributed learning, where data distribution is well controlled.

We first conduct a case study to quantify how data and resource heterogeneity in clients impacts the performance of FL with FedAvg in terms of training performance and model accuracy, and we summarize the key findings below: (1) training throughput is usually bounded by slow clients (a.k.a. stragglers) with less computational capacity and/or slower communication, which we name as the resource heterogeneity. Asynchronous training is often employed to mitigate this problem in datacenter distributed learning, but literature has shown that synchronization is a better approach for secure aggregation (Bonawitz et al., 2017) and differential privacy (McMahan et al., 2017). Moreover, FedAvg (McMahan et al., 2016) has become a common algorithm for FL that uses synchronous approach for training. (2) Different clients may train on different quantity of samples per training round and results in different round time that is similar to the straggler effect, which impacts the training time and potentially also the accuracy. We name this observation the data quantity heterogeneity. (3) In datacenter distributed learning, the classes and features of the training data are uniformly distributed among all clients, namely Independent Identical Distribution (IID). However, in FL, the distribution of data classes and features depends on the data owners, thus resulting in a non-uniform data distribution, known as non-Identical Independent Distribution (non-IID data heterogeneity). Our experiments show that such heterogeneity can significantly impact the training time and accuracy.

Driven by the above observations, we propose TiFL, a Tier-based Federated Learning System. The key idea here is adaptively selecting clients with similar per round training time so that the heterogeneity problem can be mitigated without impacting the model accuracy. Specifically, we first employ a lightweight profiler to measure the training time of each client and group them into different logical data pools based on the measured latency, called tiers. During each training round, clients are selected uniform randomly from the same tier based on the adaptive client selection algorithm of TiFL. In this way, the heterogeneity problem is mitigated as clients belonging to the same tier have similar training time. In addition to heterogeneity mitigation, such tiered design and adaptive client selection algorithm also allows controlling the training throughput and accuracy by adjusting the tier selection intelligently, e.g., selecting tiers such that the model accuracy is maintained while prioritizing selection of faster tiers. We further prove that the tiering method is compatible with privacy-preserving FL.

While resource heterogeneity and data quantity heterogeneity information can be reflected in the measured training time, the non-IID data heterogeneity information is difficult to capture. This is because any attempt to measure the class and feature distribution violates the privacy-preserving requirements. To solve this challenge, TiFL offers an adaptive client selection algorithm that uses the accuracy as indirect measure to infer the non-IID data heterogeneity information and adjust the tiering algorithm on-the-fly to minimize the training time and accuracy impact. Such approach also serves as an online version to be used in an environment where the characteristics of heterogeneity change over time.

We prototype TiFL in a FL testbed that follows the architecture design of Google’s FL system (Bonawitz et al., 2019) and perform extensive experimental evaluation to verify its effectiveness and robustness using both the popular ML benchmarks and state-of-the-art FL benchmark LEAF (Caldas et al., 2018b). The experimental results show that in the resource heterogeneity case, TiFL can improve the training time by a magnitude of 6×\times without affecting the accuracy. In the data quantity heterogeneity case, a 3×\times speedup is observed in training time with comparable accuracy to the conventional FL. Overall, TiFL outperforms the conventional FL with legal parameter number in definition of 3×\times improvement in training time and 8% improvement in accuracy in CIFAR10 (Krizhevsky et al., 2014) and 3×\times improvement in training time using FEMINIST(Caldas et al., 2018b) under LEAF.

Related Work

Federated Learning. The most recent research efforts in FL have been focusing on the functionality (Konečnỳ et al., 2016), scalability (Bonawitz et al., 2019), privacy (McMahan et al., 2017; Bonawitz et al., 2017), and tackling heterogeneity (Ghosh et al., 2019; Li et al., 2018). Existing FL approaches (Konečnỳ et al., 2016; McMahan et al., 2016; Caldas et al., 2018a) do not account for the resource and data heterogeneity, mainly focusing on weight and model compression to reduce communication overhead. They are not straggler-aware even though there can be significant latency issues with stragglers. In synchronous FL, a fixed number of clients are queried in each learning epoch to ensure performance and data privacy. Recent synchronous FL algorithms focus on reducing the total training time without considering the straggler clients. For example, (McMahan et al., 2016) proposes to reduce network communication costs by performing multiple SGD (stochastic gradient descent) updates locally and batching clients. (Konečnỳ et al., 2016) reduces communication bandwidth consumption by structured and sketched updates. FedCS (Nishio and Yonetani, 2019) proposes to solve client selection issue via a deadline-based approach that filters out slowly-responding clients. However, FedCS does not consider how this approach effects the contributing factors of straggler clients in model training. Similarly, (Wang et al., 2019) proposes an algorithm for running FL on resource-constrained devices. However, they do not aim to handle straggler clients and treat all clients as resource-constrained. In contrast, we focus on scenarios where resource-constrained devices are paired with more powerful devices to perform FL.

Most asynchronous FL algorithms do not take into consideration the effects of client drop-outs. For instance, (Smith et al., 2017) provides performance guarantee only for convex loss functions with bounded delay assumption. Furthermore, the comparison of synchronous and asynchronous methods of distributed gradient descent (Bonawitz et al., 2019) suggest that FL should use the synchronous approach (McMahan et al., 2016; Bonawitz et al., 2017), as it is more secure than the asynchronous approaches.

In the interest of space, for further background information, we recommend readers to read these papers (Li et al., 2019; Kairouz et al., 2019), where an in-depth discussion is offered for the current challenges and state-of-the-art systems in Federated Learning.

Stragglers in Datacenter Distributed Learning. Similar to datacenter distributed learning systems, the straggler problem also exists in FL and it is greatly pronounced in the synchronous learning setting. The significantly higher heterogeneity levels in FL makes it more challenging. (Bonawitz et al., 2019) proposes a simple approach to handle stragglers problem in FL, where the aggregator selects 130% of the target number of devices to initially participate, and discards stragglers during training process. With this method, the aggregator can get 30% tolerance for the stragglers by ignoring the updates from the slower edge devices. However, the 30% is set arbitrarily which requires further tuning. Furthermore simply dropping the slower clients might exclude certain data distributions available on the slower clients from contributing towards training the global model. FedProx (Li et al., 2018) also tackles resource and data heterogeneity by making improvements on the FedAvg algorithm. However, they also discard training data to make up for the systems heterogeneity.

(Li et al., 2018) takes into account the resource heterogeneity. However, the proposed approach is mainly focused on only two types of clients - stragglers and non-stragglers. In a real FL environment there is a wide range of heterogeneity levels. The proposed approach performs well in case of high ration of stragglers vs non-stragglers (80-90%). Moreover, their proposed solution involves partial training on stragglers which can further lead to biasness in trained model and sub-optimal model accuracy as explained in Section 5.2.4. (Ho et al., 2013) proposes that adding local cache is an efficient and reliable technique to deal with stragglers in datacenter distributed Learning. However, FL is powerless in governing the resources on client side, so it’s impractical to implement similar mechanisms in FL.

(Ghosh et al., 2019) proposes a general statistical model for Byzantine machines and clients with data heterogeneity that clusters based on data distribution. While data distribution may cause stragglers, (Ghosh et al., 2019) focuses on grouping edge devices such that their datasets are similar. The authors do not consider the impact of clustering on training time or accuracy. (Harlap et al., 2016) propose a novel design named RapidReassignment to handle straggles by specializing work shedding. It uses P2P communication among workers to detect slowed workers, performs work re-assignment, and exploits iteration knowledge to further reduce how much data needs to be preloaded on helpers. However, as stated in Section 1, migrating a user’s private data to other unknown users’ devices is strictly restricted in FL. An analogical approach named SpecSync is proposed in (Zhang et al., 2018), where each worker speculates about the parameter updates from others, and if necessary, it aborts the ongoing computation, pulls fresher parameters to start over, so as to opportunistically improve the training quality. However, information sharing between clients is not allowed in FL.

Heterogeneity Impact Study

Compared with datacenter distributed learning and cross-silo FL, one of the key features of cross-device FL is the significant resource and data heterogeneity among clients, which can potentially impact both the training throughput and the model accuracy. Resource heterogeneity arises as a result of vast number of computational devices with varying computational and communication capabilities involved in the training process. The data heterogeneity arises as a result of two main reasons - (1) the varying number of training data samples available at each client and (2) the non-uniform distribution of classes and features among the clients.

Cross-device FL is performed as an iterative process whereby the model is trained over a series of global training rounds, and the trained model is shared by all the involved clients. We define KK as the total pool of clients available to select from for each global training round, and CC as the set of clients selected per round. In every global training round, the aggregator selects a random fraction of clients CrC_{r} from KK. The vanilla cross-device FL algorithm is briefly summarized in Alg. 1. The aggregator first randomly initializes weights of the global model denoted by ω0\omega_{0}. At the beginning of each round, the aggregator sends the current model weights to a subset of randomly selected clients. Each selected client then trains its local model with its local data and sends back the updated weights to the aggregator after local training. At each round, the aggregator waits until all selected clients respond with their corresponding trained weights. This iterative process keeps on updating the global model until a certain number of rounds are completed or a desired accuracy is reached.

The state-of-the-art cross-device FL system proposed in (Bonawitz et al., 2017) adopts a client selection policy where clients are selected randomly. A coordinator is responsible for creating and deploying a master aggregator and multiple child aggregators for achieving scalability as the real world cross-device FL system can involve up to 101010^{10} clients (Bonawitz et al., 2017; Li et al., 2019; Kairouz et al., 2019).

At each round, the master aggregator collects the weights from all the child aggreagtors to update the global model.

2. Heterogeneity Impact Analysis

The resource and data heterogeneity among involved clients may lead to varying response latencies (i.e., the time between a client receives the training task and returns the results) in the cross-device FL process, which is usually referred as the straggler problem.

We denote the response latency of a client cic_{i} as LiL_{i}, and the latency of a global training round is defined as

where LrL_{r} is the latency of round rr. From Equation (1), we can see the latency of a global training round is bounded by the maximum training latency of clients in CC, i.e., the slowest client.

We define τ\tau levels of clients, i.e., within the same level, the clients have similar response latencies. Assume that the total number of levels is mm and τm\tau_{m} is the slowest level with ∣τm∣|\tau_{m}| clients inside. In the baseline case (Alg. 1), the aggregator selects the clients randomly, resulting in a group of selected clients with composition spanning multiple client levels.

We formulate the probability of selecting ∣C∣|C| clients from all client levels except the slowest level τm\tau_{m} as follows:

Accordingly, the probability of at least one client in CC comes from τm\tau_{m} can be formulated as:

a−1b−1<ab,  while  1<a<b\frac{a-1}{b-1}<\frac{a}{b},\;while\;1<a<b

Since 1<a<b1<a<b, we could get ab−b<ab−aab-b<ab-a, that is (a−1)b<(b−1)a  and  a−1b−1<ab(a-1)b<(b-1)a\;and\;\frac{a-1}{b-1}<\frac{a}{b}. ∎

In real-world scenarios, large number of clients can be selected at each round, which makes ∣K∣|K| extremely large. As a subset of KK, the size of CC can also be sufficiently large. Since ∣K∣−∣τm∣∣K∣<1\frac{|K|-|\tau_{m}|}{|K|}<1, we get (∣K∣−∣τm∣∣K∣)∣C∣≈0(\frac{|K|-|\tau_{m}|}{|K|})^{|C|}\approx 0, which makes Prs≈1Pr_{s}\approx 1, meaning in a vanilla cross-device FL training process, the probability of selecting at least one client from the slowest level is reasonably high for each round. According to Equation (1), the random selection strategy adopted by state-of-the-art cross-device FL system may suffer from a slow training performance.

3. Experimental Study

To experimentally verify the above analysis and demonstrate the impact of resource heterogeneity and data quantity heterogeneity, we conduct a study with a setup similar to the paper (Chai et al., 2019). The testbed is briefly summarized as follows -

We use a total of 20 clients and each client is further divided into 5 groups with 4 client per group.

We allocate 4 CPUs, 2 CPUs, 1 CPU, 1/3 CPU, 1/5 CPU for every client from group 1 through 5 respectively to emulate the resource heterogeneity.

The model is trained on the image classification dataset CIFAR10 (Krizhevsky et al., 2014) using the vanilla cross-device FL process 3.1 (model and learning parameters are detailed in Section 5).

Experiments with different data size for every client are conducted to produce data heterogeneity results.

As shown in Fig. 1 (a), with the same amount of CPU resource, increasing the data size from 500 to 5000 results in a near-linear increase in training time per round. As the amount of CPU resources allocated to each client increases, the training time gets shorter. Additionally, the training time increases as the number of data points increase with the same number of CPUs. These preliminary results imply that the straggler issues can be severe under a complicated and heterogeneous FL environment.

To evaluate the impact of data distribution heterogeneity, we keep the same CPU resources for every client (i.e., 2 CPUs) and generate a biased class and feature distribution following (Zhao et al., 2018). Specifically, we distribute the dataset in such a way that every client has equal number of images from 2 (non-IID(2)), 5 (non-IID(5)) and 10 (non-IID(10)) classes, respectively. We train the model on Cifar10 dataset using the vanilla FL system as described in Section 3.1 with the model and training parameters detailed in Section 5. As seen in Fig. 1 (b), there is a clear difference in the accuracy with different non-IID distributions. The best accuracy is given by the IID since it represents a uniform class and feature distribution. As the number of classes per client is reduced, we observe a corresponding decrease in accuracy. Using 10 classes per client reduces the final accuracy by around 6% compared to IID (it is worth noting that non-IID(10) is not the same as IID as the feature distribution in non-IID(10) is skewed compare to IID). In the case of 5 classes per client, the accuracy is further reduced by 8%. The lowest accuracy is observed in the 2 classes per client case, which has a significant 18% drop in accuracy.

These studies demonstrate that the data and resource heterogeneity can cause significant impact on training time and training accuracy in cross-device FL. To tackle these problems, we propose TiFL— a tier-based FL system which introduces a heterogeneity-aware client selection methodology that selects the most profitable clients during each round of the training to minimize the heterogeneity impact while preserving the FL privacy proprieties, thus improving the overall training performance of cross-device FL (in the following, we use FL to denote cross-device FL for simplicity).

TiFL: A Tier-based Federated Learning System

In this section, we present the design of the proposed tier-based federated learning system TiFL. The key idea of a tier-based system is that given the global training time of a round is bounded by the slowest client selected in that round (see Equation 1), selecting clients with similar response latency in each round can significantly reduce the training time. We first give an overview of the architecture and the main flow of TiFL system. Then we introduce the profiling and tiering approach. Based on the profiling and tiering results, we explain how a tier selection algorithm can potentially mitigate the heterogeneity impact through a straw-man proposal as well as the limitations of such static selection approach. To this end, we propose an adaptive tier selection algorithm to address the limitations of the straw-man proposal. Finally, we propose an analytical model through which one can estimate the expected training time using selection probabilities of tiers and the total number of training rounds.

The overall system architecture of TiFL is present in Fig. 2. TiFL follows the system design to the state-of-the-art FL system (Bonawitz et al., 2017) and adds two new components: a tiering module (a profiler & tiering algorithms) and a tier scheduler. These newly added components can be incorporated into the coordinator of the existing FL system (Bonawitz et al., 2019). It is worth to note that in Fig. 2, we only show a single aggregator rather than the hierarchical master-child aggregator design for a clean presentation purpose. TiFL supports master-child aggregator design for scalability and fault tolerance.

In TiFL, the first step is to collect the latency metrics of all the available clients through a lightweight profiling as detailed in Section 4.2. The profiled data is further utilized by our tiering algorithm. This groups the clients into separate logical pools called tiers. Once the scheduler has the tiering information (i.e., tiers that the clients belong to and the tiers’ average response latencies), the training process begins. Different from vanilla FL that employs a random client selection policy, in TiFL the scheduler selects a tier and then randomly selects targeted number of clients from that tier. After the selection of clients, the training proceeds as state-of-the-art FL system does. By design, TiFL is non-intrusive and can be easily plugged into any existing FL system in that the tiering and scheduler module simply regulate client selection without intervening the underlying training process.

2. Profiling and Tiering

Given the global training time of a round is bounded by the slowest client selected in that round (see Equation 1), if we can select clients with similar response latency in each round, the training time can be improved. However, in FL, the response latency is unknown a priori, which makes it challenging to carry out the above idea. To solve this challenge, we introduce a process through which the clients are tiered (grouped) by the Profiling and Tiering module as shown in Fig. 2. As the first step, all available clients are initialized with a response latency LiL_{i} of 0. The profiling and tiering module then assigns all the available clients the profiling tasks. The profiling tasks execute for sync_roundssync\_rounds rounds and in each profiling round, the aggregator asks every client to train on the local data and waits for their acknowledgement for TmaxT_{max} seconds. All clients that respond within TmaxT_{max} have their response latency value RTiRT_{i} incremented with the actual training time, while the ones that have timed out are incremented by TmaxT_{max}. After sync_roundssync\_rounds rounds are completed, the clients with Li>=sync_rounds∗TmaxL_{i}>=sync\_rounds*T_{max} are considered dropouts and excluded from the rest of the calculation. The collected training latencies from clients creates a histogram, which is split into mm groups and the clients that fall into the same group forms a tier. The average response latency is then calculated for each group and recorded persistently which is used later for scheduling and selecting tiers. The profiling and tiering can be conducted periodically for systems with changing computation and communication performance over the time so that clients can be adaptively grouped into the right tiers.

3. Straw-man Proposal: Static Tier Selection Algorithm

In this section, we present a naive static tier-based client selection policy and discuss its limitations, which motivates us to develop an advanced adaptive tier selection algorithm in the next section. While the profiling and tiering module introduced in Section 4.2 groups clients into mm tiers based on response latencies, the tier selection algorithm focuses on how to select clients from the proper tiers in the FL process to improve the training performance. The natural way to improve training time is to prioritize towards faster tiers, rather than selecting clients randomly from all tiers (i.e., the full KK pool). However, such selection approach reduces the training time without taking into consideration of the model accuracy and privacy properties. To make the selection more general, one can specify each tier njn_{j} is selected based on a predefined probability, which sums to 1 across all tiers. Within each tier, ∣C∣\mid C\mid clients are uniform randomly selected.

In a real-world FL scenarios, there can be a large number of clients involved in the FL process (e.g., up to 101010^{10}) (Bonawitz et al., 2017; Li et al., 2019; Kairouz et al., 2019). Thus in our tiering-based approach, the number of tiers is set such that mm <<<< ∣K∣|K| and number of clients per tier njn_{j} is always greater than ∣C∣|C|. The selection probability of a tier is controllable, which results in different trade-offs. If the users’ objective is to reduce the overall training time, they may increase the chances of selecting the faster tiers. However, drawing clients only from the fastest tier may inevitably introduce training bias due to the fact that different clients may own a diverse set of heterogeneous training data spread across different tiers; as a result, such bias may end up affecting the accuracy of the global model. To avoid such undesired behavior, it is preferable to involve clients from different tiers so as to cover a diverse set of training datasets. We perform an empirical analysis on the latency-accuracy trade-off in Section 5.

4. Adaptive Tier Selection Algorithm

While the above naive static selection method is intuitive, it does not provide a method to automatically tune the trade-off to optimize the training performance nor adjust the selection based on changes in the system. In this section, we propose an adaptive tier selection algorithm that can automatically strike a balance between training time and accuracy, and adapt the selection probabilities adaptively over training rounds based on the changing system conditions.

The observation here is that heavily selecting certain tiers (e.g., faster tiers) may eventually lead to a biased model, TiFL needs to balance the client selection from other tiers (e.g., slower tiers). The question being which metric should be used to balance the selection. Given the goal here is to minimize the bias of the trained model, we can monitor the accuracy of each tier throughout the training process. A lower accuracy value of a tier tt typically indicates that the model has been trained with less involvement of this tier, therefore tier tt should contribute more in the next training rounds. To achieve this, we can increase the selection probabilities for tiers with lower accuracy. To achieve good training time, we also need to limit the selection of slower tiers across training rounds. Therefore, we introduce CreditstCredits_{t}, a constraint that defines how many times a certain tier can be selected.

Specifically, a tier is initialized randomly with equal selection probability. After the weights are received and the global model is updated, the global model is evaluated on every client for every tier on their respective TestDataTestData and their resulting accuracies are stored as the corresponding tier tt’s accuracy for that round rr. This is stored in AtrA^{r}_{t}, which is the mean accuracy for all the clients in tier tt in training round rr. In the subsequent training rounds, the adaptive algorithm updates the probability of each tier based on that tier’s test accuracy at every II rounds. This is done in the function ChangeProbsChangeProbs, which adjusts the probabilities such that the lower accuracy tiers get higher probabilities to be selected for training; then with the new tier-wise selection probabilities (NewProbsNewProbs), a tier which has remaining CreditstCredits_{t} is selected from all available tiers τ\tau. The selected tier will have its CreditstCredits_{t} decremented. As clients from a particular tier gets selected over and over throughout the training rounds, the CreditstCredits_{t} for that tier ultimately reduces down to zero, meaning that it will not be selected again in the future. This serves as a control knob for the number of times a tier is selected and by setting this upper-bound, we can limit the amount of times a slower tier contributes to the training, thereby effectively gaining some control over setting a soft upper-bound on the total training time. For the straw-man implementation, we used a skewed probability of selection to manipulate training time. Since we now wish to adaptively change the probabilities, we add the CreditstCredits_{t} to gain control over limiting training time.

On one hand, the tier-wise accuracy ArtA^{t}_{r} essentially makes TiFL’s adaptive tier selection algorithm data heterogeneity aware; as such, TiFL makes the tier selection decision by taking into account the underlying dataset selection biasness, and automatically adapt the tier selection probabilities over time. On the other hand, CreditstCredits_{t} is introduced to intervene the training time by enforcing a constraint over the selection of the relatively slower tiers. While CreditstCredits_{t} and AtrA^{r}_{t} mechanisms optimize towards two different and sometimes contradictory objectives — training time and accuracy, TiFL cohesively synergizes the two mechanisms to strike a balance for the training time-accuracy trade-off. More importantly, with TiFL, the decision making process is automated, thus relieving the users from intensive manual effort. The adaptive algorithm is summarized in Algo. 2.

5. Training Time Estimation Model

In real-life scenarios, the training time and resource budget is typically finite. As a result, FL users may need to compromise between training time and accuracy. A training time estimation model would facilitate users to navigate the training time-accuracy trade-off curve to effectively achieve desired training goals. Therefore, we build a training time estimation model that can estimate the overall training time based on the given latency values and the selection probability of each tier:

where LallL_{all} is the total training time, Ltier_iL_{tier\_i} is the response latency of tier ii, PiP_{i} is the probability of tier ii, and RR is the total number of training rounds. The model is a sum of products of the tier and latencies, which gives the latency expectation per round. This is multiplied by the total number of training rounds to get the total training time.

6. Discussion: Compatibility with Privacy-Preserving Federated Learning

FL has been used together with privacy preserving approaches such as differential privacy to prevent attacks that aim to extract private information (Shokri et al., 2017; Nasr et al., 2018). Privacy-preserving FL is based on client-level differential privacy, where the privacy guarantee is defined at each individual client. This can be accomplished by each client implementing a centralized private learning algorithm as their local training approach. For example, with neural networks this would be one or more epochs using the approach proposed in (Abadi et al., 2016b). This requires each client to add the appropriate noise into their local learning to protect the privacy of their individual datasets. Here we demonstrate that TiFL is compatible with such privacy preserving approaches.

Assume that for client cic_{i}, one round of local training using a differentially private algorithm is (ϵ\epsilon, δ\delta)-differentially private, where ϵ\epsilon bounds the impact any individual may have on the algorithm’s output and δ\delta defines the probability that this bound is violated. Smaller ϵ\epsilon values therefore signify tighter bounds and a stronger privacy guarantee. Enforcing smaller values of ϵ\epsilon requires more noise to be added to the model updates sent by clients to the FL server which leads to less accurate models. Selecting clients at each round of FL has distinct privacy and accuracy implications for client-level privacy-preserving FL approaches. For simplicity we assume that all clients are adhering to the same privacy budget and therefore same (ϵ\epsilon, δ\delta) values. Let us first consider the scenario wherein CC is chosen uniformly at random each round. Compared with each client participating in each round, the overall privacy guarantee, using random sampling amplification (Beimel et al., 2010), improves from (ϵ\epsilon, δ\delta) to (O(qϵq\epsilon), qδq\delta) where q=∣C∣∣K∣q=\frac{|C|}{|K|}. This means that there is a stronger privacy guarantee with the same noise scale. Clients may therefore add less noise per round or more rounds may be conducted without sacrificing privacy. For the tiered approach the guarantee also improves. Compared to (ϵ\epsilon, δ\delta) in the all client scenario, the tiered approach improves to an (O(qmaxϵ),qmaxδO(q_{max}\epsilon),q_{max}\delta) privacy guarantee where the probability of selecting tier with weight θj\theta_{j} is given by 1ntiers∗θj\frac{1}{n_{tiers}}*\theta_{j}, qmax=max⁡j=1...∣ntiers∣qjq_{max}=\max_{j=1...|n_{tiers}|}q_{j} and qj=(1ntiers∗θj)∣C∣∣nj∣q_{j}=(\frac{1}{n_{tiers}}*\theta_{j})\frac{|C|}{|n_{j}|}.

Experimental Evaluation

We prototype TiFL with both the naive and our adaptive selection approach and perform extensive testbed experiments under three scenarios: resource heterogeneity, data heterogeneity, and resource plus data heterogeneity.

Testbed. As a proof of concept case study, we build a FL testbed for the syntehtic datasets by deploying 50 clients on a CPU cluster where each client has its own exclusive CPU(s) using Tensorflow (Abadi et al., 2016a). In each training round, 5 clients are selected to train on their own data and send the trained weights to the server which aggregates them and updates the global model similar to (McMahan et al., 2016; Bonawitz et al., 2019). (Bonawitz et al., 2019) introduces multiple levels of server aggregators in order to achieve scalability and fault tolerance in extreme scale situations, i.e., with millions of clients. In our prototype, we simplify the system to use a powerful single aggragator as it is sufficient for our purpose here, i.e., our system does not suffer from scalabiltiy and fault tolerance issues, though multiple layers of aggregator can be easily integrated into TiFL.

We also extend the widely adopted large scale distributed FL framework LEAF (Caldas et al., 2018b) in the same way. LEAF provides inherently non-IID with data quantity and class distributions heterogeneity. LEAF framework does not provide the resource heterogeneity among the clients, which is one of the key properties of any real-world FL system. The current implementation of the LEAF framework is a simulation of a FL system where the clients and server are running on the same machine. To incorporate the resource heterogeneity we first extend LEAF to support the distributed FL where every client and the aggregator can run on separate machines, making it a real distributed system. Next, we deploy the aggregator and clients on their own dedicated hardware. This resource assignment for every client is done through uniform random distribution resulting in equal number of clients per hardware type. By adding the resource heterogeneity and deploying them to separate hardware, each client mimics a real-world edge-device. Given that LEAF already provides non-IIDness, with the newly added resource heterogeneity feature the new framework provides a real world FL system which supports data quantity, quality and resource heterogeneity. For our setup, we use exactly the same sampling size used by the LEAF (Caldas et al., 2018b) paper (0.05) resulting in a total of 182 clients, each with a variety of image quantities.

2. Experimental Results

Models and Datasets. We use four image classification applications for evaluating TiFL. We use MNIST and Fashion-MNIST (Xiao et al., 2017), where each contains 60,000 training images and 10,000 test images, where each image is 28x28 pixels. We use a CNN model for both datasets, which starts with a 3x3 convolution layer with 32 channels and ReLu activation, followed by a 3x3 convolution layer with 64 channels and ReLu activation, a MaxPooling layer of size 2x2, a fully connected layer with 128 units and ReLu activation, and a fully connected layer with 10 units and ReLu activation. Dropout 0.25 is added after the MaxPooling layer, dropout 0.5 is added before the last fully connected layer. We use Cifar10 (Krizhevsky et al., 2014), which contains richer features compared to MNIST and Fashion-MNIST. There is a total of 60,000 colour images, where each image has 32x32 pixels. The full dataset is split evenly between 10 classes, and partitioned into 50,000 training and 10,000 test images. The model is a four-layer convolution network ending with two fully-connected layers before the softmax layer. It was trained with a dropout of 0.25. Lastly we also use the FEMNIST data set from LEAF framework (Caldas et al., 2018b). This is an image classification dataset which consists of 62 classes and the dataset is inherently non-IID with data quantity and class distributions heterogeneity. We use the standard model architecture as provided in LEAF (Caldas et al., 2018a).

Training Hyperparameters. We use RMSprop as the optimizer in local training and set the initial learning rate (η\eta) as 0.01 and decay as 0.995. Local batch size of each client is 10, and local epochs is 1. The total number of clients (∣K∣|K|) is 5050 and the number of participated clients (∣C∣|C|) at each round is 55. For FEMNIST we use the default training parameters provided by the LEAF Framework (SGD with lr 0.004, batch size 10). We train for a total of 2000 rounds for FEMNIST and 500 rounds for the synthetic datasets. Every experiment is run 5 times and we use the average values.

Heterogeneous Resource Setup. Among all the clients, we split them into 5 groups with equal clients per group. For MNIST and Fashion-MNIST, each group is assigned with 22 CPUs, 11 CPU, 0.750.75 CPU, 0.50.5 CPU, and 0.250.25 CPU per part respectively. For the larger Cifar10 and FEMINIST model, each group is assigned with 44 CPUs, 22 CPUs, 11 CPU, 0.50.5 CPU, and 0.10.1 CPU per part respectively. This leads to varying training time for clients belong to different groups. By using the tiering algorithm of TiFL, there are 5 tiers

Heterogeneous Data Distribution. FL differs from the datacenter distributed learning in that the clients involved in the training process may have non-uniform data distribution in terms of amount of data per client and the non-IID data distribution. ∙\bullet For data quantity heterogeneity, the training data sample distribution is 10%, 15%, 20%, 25%, 30% of total dataset for difference groups, respectively, unless otherwise specifically defined. ∙\bullet For non-IID heterogeneity, we use different non-IID strategies for different datasets. For MNIST and Fashion-MNIST, we adopt the setting in (McMahan et al., 2016), where we sort the labels by value first, divide into 100 shards evenly, and then assign each client two shards so that each client holds data samples from at most two classes. For Cifar10, we shard the dataset unevenly in a similar way and limit the number of classes to 5 per client (non-IID(5)) following (Zhao et al., 2018), (Liu et al., 2019) unless explicitly mentioned otherwise. In the case of FEMINIST we use its default non-IID-ness.

Scheduling Policies. We evaluate several different naive scheduling policies of the proposed tier-based selection approach, defined by the selection probability from each tier, and compare it with the state-of-the-practice policy (or no policy) that existing FL works adopt, i.e., randomly select 5 clients from all clients in each round (McMahan et al., 2016; Bonawitz et al., 2019), agnostic to any heterogeneity in the system. We name it as vanilla. fast is a policy that TiFL only selects the fastest clients in each round. random demonstrates the case where the selection of the fastest tier is prioritized over slower ones. uniform is a base case for our tier-based naive selection policy where every tier has an equal probability of being selected. slow is the worst policy that TiFL only selects clients from the slowest tiers and we only include it here for reference purpose so that we can see a performance range between the best case and the worst case scenarios for static tier-based selection approach. We use the above policies for CIFAR-10 and FEMINIST training. For MNIST and Fashion-MNIST, given it is a much more lightweight workload, we focus on demonstrating the sensitivity analysis when the policy prioritizes more aggressively towards the fast tier, i.e., from fast1 to fast3, the slowest tier’s selection probability has reduced from 0.1 to 0 while all other tiers got equal probability. We also include the uniform policy for comparison, which is the same as in CIFAR-10. Table 1 summarizes all these scheduling policies by showing their selection probabilities.

In this section, we evaluate the accuracy of our training time estimation model on different naive tier selection policies by comparing the estimation results of the model with the measurements obtained from test-bed experiments. The estimation model takes as input of the profiled average latency of each tier, the selection probabilities, and total number of training rounds to estimate the training time. We use mean average prediction error (MAPE) as the evaluation metric, which is defined as follows:

where LallestL_{all}^{est} is the estimated training time calculated by the estimation model and LallactL_{all}^{act} is the actual training time measured during the training process. Table 2 demonstrates the comparison results. The results suggest the analytical model is very accurate as the estimation error never exceeds more than 6 %.

2.2. Resource Heterogeneity

In this sections, we evaluate the performance of TiFL in terms of training time and model accuracy in a resource heterogeneous environment as depicted in 5.1 and we assume there is no data heterogeneity. In practice, data heterogeneity is a norm in FL, we evaluate this scenario to demonstrate how TiFL tame resource heterogeneity alone and we evaluate the scenario with both resource and data heterogeneity in Section 5.2.4.

In the interest of space, we only present the Cifar10 results here as MNIST and Fashion-MNIST share the similar observations. The results are organized in Fig. 3 (column 1), which clearly indicate that when we prioritize towards the fast tiers, the training time reduces significantly. Compared with vanilla, fast achieves almost 11 times improvement in training time, see Fig. 3 (a). One interesting observation is that even uniform has an improvement of over 6 times over the vanilla. This is because the training time is always bounded by the slowest client selected in each training round. In TiFL, selecting clients from the same tier minimizes the straggler issue in each round, and thus greatly improves the training time. For accuracy comparison, Fig. 3 (c) shows that the difference between polices are very small, i.e., less than 3.71% after 500 rounds. However, if we look at the accuracy over wall-clock time, TiFL achieves much better accuracy compared to vanilla, i.e., up to 6.19% better if training time is constraint, thanks to the much faster per round training time brought by TiFL, see Fig. 3 (e). Note here that different policies may take very different amount of wall-clock time to finish 500 rounds.

2.3. Data Heterogeneity

In this section, we evaluate data heterogeneity due to both data quantity heterogeneity and non-IID heterogeneity as depicted in Section 5.1. To demonstrate only the impact from data heterogeneity, we allocate homogeneous resource to each client, i.e., 2 CPUs per client.

∙\bullet Data quantity heterogeneity. The training time and accuracy results are show in Fig. 3 (column 2). In the interest of space, we only show Cifar10 results here. From the training time comparison in Fig. 3 (b), it is interesting that TiFL also helps in data heterogeneity only case and achieves up to 3 times speedup. The reason is that data quantity heterogeneity may also result in different round time, which shares the similar effect as resource heterogeneity. Fig. 3 (d) and (f) show the accuracy comparison, where we can see fast has relatively obvious drop compared to others because Tier 1 only contains 10% of the data, which is a significant reduction in volume of the training data. slow is also a heavily biased policy towards only one tier, but Tier 5 contains 30% of the data thus slow maintains good accuracy while worst training time. These results imply that like resource heterogeneity only, data heterogeneity only can also benefit from TiFL. However, policies that are too aggressive toward faster tier needs to be used very carefully as clients in fast tier achieve faster round time due to using less samples. It is also worth to point out that in our experiments the total amount of data is relatively limited. In a practical case where data is significantly more, the accuracy drop of fast is expected to be less pronounced.

∙\bullet non-IID heterogeneity. We observe that non-IID heterogeneity does not impact the training time. Hence, we omit the results here. However, non-IID heterogeneity effects the accuracy. Fig. 4 shows the accuracy over rounds given 2, 5, and 10 classes per client in a non-IID setting. We also show the IID results in plot for comparison. These results show that as the heterogeneity level in non-IID heterogeneity increases, the accuracy impact also increases for all policies due to the strongly biased training data. Another important observation is that vanilla case and uniform have a better resilience than other policies, thanks to the unbiased selection behavior, which helps minimize further bias introduced during the client selection process.

2.4. Resource and Data Heterogeneity

This section presents the most practical case study, since here we evaluate with both resource and data heterogeneity combined.

MNIST and Fashion-MNIST (FMNIST) results are shown in Fig. 5 columns 1 and 2 respectively. Overall, policies that are more aggressive towards the fast tiers bring more speedup in training time. For accuracy, all polices of TiFL are close to vanilla, except fast3 falls short as it completely ignores the data in Tier 5.

Cifar10 results are shown in Fig. 6 column 1. It presents the case of resource heterogeneity plus non-IID data heterogeneity with equal data quantities per client and the results are similar to resource heterogeneity only since non-IID data with the same amount of data quantity per client results in a similar effect of resource heterogeneity in terms of training time. However, the accuracy degrades slightly more here as because of the non-IID-ness the features are skewed, which results in more training bias among different classes.

Fig. 6 column 2 shows the case of resource heterogeneity plus both the data quantity heterogeneity and non-IID heterogeneity. As expected, the training time shown in Fig. 6 (b) is similar to Fig. 6 (a) since the training time impact from different data amounts can be corrected by TiFL. However, the behaviors of round accuracy are quite different here as shown in Fig. 6 (d). The accuracy of fast has degraded a lot more due to the data quantity heterogeneity as it further amplifies the training class bias (i.e., the data of some classes become very little to none) in the already very biased data distribution caused by the non-IID heterogeneity. Similar reasons can explain for other policies The best performing policy in accuracy here is the uniform case and is almost the same as vanilla, thanks to the even selection nature which results in little increase in training class bias. Fig. 6 (f) shows the wall-clock time accuracy. As expected, the significantly improved per round time in TiFL shows its advantage here as within the same time budget, more iterations can be done with shorter round time and thus remedies the accuracy disadvantage per round. fast still falls short than vanilla in the long run as the limited and biased data limits the benefits of more iterations. fast also perform worse than vanilla as it has no training advantage.

2.5. Adaptive Selection Policy

The above evaluation demonstrate the naive selection approach in TiFL can significantly improve the training time, but sometimes can fall short in accuracy, especially when strong data heterogeneity presents as such approach is data-heterogeneity agnostic. In this section, we evaluate the proposed adaptive tier selection approach of TiFL, which takes into consideration of both resource and data heterogeneity when making scheduling decisions without privacy violation. We compare adaptive with vanilla and uniform, and the later is the best accuracy performing static policy.

Fig. 7 shows adaptive outperforms vanilla and uniform in both training time and accuracy for resource heterogeneity with data quantity heterogeneity (Amount) and non-IID heterogeneity (Class), thanks to the data heterogeneity-aware schemes. In the combined resource and data heterogeneity case (Combine), adaptive achieves comparable accuracy with vanilla with almost half of the training time, and performs similar as uniform in training time while improves significantly in accuracy. The above robust performance of adaptive is credited to both the resource and data heterogeneity-aware schemes. To demonstrate the robustness of adaptive, we compare the accuracy over rounds for different policies under different non-IID heterogeneity in Fig. 8. It is clear that adaptive consistently outperforms vanilla and uniform in different level of non-IID heterogeneity.

2.6. Adaptive Selection Policy(LEAF)

This section provides the evaluation of TiFL using a widely adopted large scale distributed FL dataset FEMINIST from the LEAF framework (Caldas et al., 2018b) . We use exactly the same configurations (data distribution, total number of clients, model and training hyperparameters) as mentioned in (Caldas et al., 2018b) resulting in total number of 182 clients, i.e. deploy-able edge devices. Since LEAF provides it’s own data distribution among devices the addition of resource heterogeneity results in a range of training times thus generating a scenario where every edge device has a different training latency. We further incorporated TiFL’s tiering module and selection policy to the extended LEAF framework. The profiling modules collects the training latency of each clients and creates a logical pool of tiers which is further utilized by the scheduler. The scheduler selects a tier and then the edge clients within the tier in each training round. For our experiments with LEAF we limit the total number of tiers to 5 and during each round we select 10 clients, with 1 local epoch per round.

Figure 9 shows the training time and accuracy over rounds for LEAF with different client selection policies. Figure 9a shows the training time for different selection policies. The least training time is achieved by using the fast selection policy however, it impact the final model accuracy by almost 10% compared to vanilla selection policy. The reason for the least accuracy for fast is the result of less training point among the clients in tier 1. One interesting observation is slow out performs the selection policy fast in terms of accuracy even though each of these selection policies rely on data from only one tier. It must be noted that the slow tier is not only the reason of less computing resources but also the higher quantity of training data points. These results are consistent with our observations from the results presented in Section 5.2.3.

Figure 9b shows the accuracy over-rounds for different selection policies. Our proposed adaptive selection policy achieves 82.1% accuracy and outperforms the slow and fast selection policies by 7% and 10% respectively. The adaptive policy is on par with the vanilla and uniform ( 82.4% and 82.6% respectively). when comparing the total training time for 2000 rounds adaptive achieves 7 ×\times and 2 ×\times improvement compare to vani and uniform respectively. fast and random both outperformed the adaptive in terms of training time however, even after convergence the accuracy for both of these selection policies show a noticeable impact on the final model accuracy. The results for FEMINIST using the extended LEAF framework for both accuracy as well as training time are also consistent with the results reported in Section 5.2.5.

Conclusion

In this paper, we investigate and quantify the heterogeneity impact on “decentralized virtual supercomputer” - FL systems. Based on the observations of our case study, we propose and prototype a Tier-based Federated Learning System called TiFL. Tackling the resource and data heterogeneity, TiFL employs a tier-based approach that groups clients in tiers by their training response latencies and selects clients from the same tier in each training round. To address the challenge that data heterogeneity information cannot be directly measured due to the privacy constraints, we further design an adaptive tier selection approach that enables TiFL be data heterogeneity aware and outperform conventional FL in various heterogeneous scenarios: resource heterogeneity, data quantity heterogeneity, non-IID data heterogeneity, and their combinations. Specifically, TiFL achieves an improvement over conventional FL by up to 3×\times speedup in overall training time and by 6% in accuracy.

References