Papaya: Practical, Private, and Scalable Federated Learning

Dzmitry Huba, John Nguyen, Kshitiz Malik, Ruiyu Zhu, Mike Rabbat, Ashkan Yousefpour, Carole-Jean Wu, Hongyuan Zhan, Pavel Ustinov, Harish Srinivas, Kaikai Wang, Anthony Shoumikhin, Jesik Min, Mani Malek

Introduction

Cross-device federated learning (FL) is a distributed learning paradigm where a large collection of clients collaborate to train a machine learning model while the raw training data stays on client devices. FL promises to train high-quality models by leveraging data from massive client populations, while ensuring security and privacy of client data.

In traditional parallel systems, concurrency refers to the number of processors running a parallel application, and utilization refers to the fraction of processors actively computing at any time. In this paper we focus on the scalability of FL systems: “a measure of [their] capacity to effectively utilize an increasing number of processors” Kumar & Gupta (1994). In the context of FL, concurrency refers to the number of clients training simultaneously, and our aim is to develop systems that can accelerate training by training concurrently on more clients. Companies like Apple, Meta, Google, and others have the potential to scale FL training to hundreds of millions or billions of clients.

Prior work describing FL systems has focused on synchronous training Bonawitz et al. (2019); Paulik et al. (2021); Ludwig et al. (2020); NVIDIA ; WeBank . Synchronous FL (SyncFL) training proceeds in rounds, as illustrated in Figure 1. The number of clients participating in each round corresponds to the concurrency.In the SyncFL literature, concurrency is also referred to as clients per round McMahan et al. (2016). In each round, clients download the current server model, train this model locally on their respective data, and report a model update back to the server. Once all client updates are ready, they are aggregated, and then the server computes the new model using the aggregated updates.

SyncFL faces two main challenges when scaling. First, there are many sources of heterogeneity in cross-device FL Kairouz et al. (2019): clients have different hardware capabilities (processor speeds, memory sizes), and data can be highly imbalanced across clients, with some clients having multiple orders of magnitude more data than others. In synchronous systems, heterogeneity results in stragglers — clients in the tail take much longer to complete local training and prolong the time to complete each round of training, hampering utilization. Over-selection is commonly used to reduce the impact of stragglers on the runtime of SyncFL methods Bonawitz et al. (2019). Over-selection results in discarding updates from the slowest-responding clients selected in each round, and it has been noted that this may bias the trained model against slow-responding clients.

The second challenge SyncFL faces is that increasing concurrency in synchronous training corresponds to using larger cohorts (group of clients participating in a round); i.e., more user updates averaged before performing a server update. This leads to similar effects as using large batches in traditional data-parallel training Keskar et al. (2017). Large-cohort training has been found to make inefficient use of client updates Bonawitz et al. (2019); Charles et al. (2021). Consequently, increasing cohort size does not reduce wall-clock training time proportionally.

Asynchronous FL (AsyncFL, see Section 3) can potentially alleviate these challenges. In AsyncFL, clients return updates to be aggregated as soon as the updates are ready, and a new client may then begin computing updates immediately. Client training is decoupled from server model updates. Consequently, AsyncFL is not impacted by stragglers and utilization can be kept high (essentially at 100%) throughout training. However, as with all asynchronous systems, AsyncFL must handle staleness — updates from clients, especially slow-responding clients, based on a server model that has been updated many times in the interim, and hence may not provide useful information for training Bertsekas & Tsitsiklis (1989). AsyncFL methods have been previously explored Xie et al. (2019); Nguyen et al. (2021); Xu et al. (2021), but none has yet been demonstrated and evaluated at scale.

Contributions. This paper presents Papaya,Why Papaya? Say “privacy-preserving AI” five times in a row, fast. the first production FL system to support asynchronous and synchronous training at scale. We introduce a novel asynchronous secure aggregation protocol, allowing clients to communicate updates to the server in a cryptographically secure manner without needing to wait until other clients are ready to perform secure aggregation. This enables the implementation of FL with buffered asynchronous aggregation that has been recently introduced in Nguyen et al. (2021).

We evaluate Papaya in Section 7 by training a language model for next-word prediction on a population of millions of devices in the field. We demonstrate that AsyncFL is substantially more scalable than SyncFL. Although asynchronous execution results in some stale client responses, staleness in AsyncFL can be controlled by choosing an appropriate aggregation goal in buffered asynchronous aggregation Nguyen et al. (2021). The aggregation goal is the number of client updates that need to be received before the server performs a model update. Consequently, AsyncFL can compute many more server updates than SyncFL in a fixed amount of time, leading to much better scaling than SyncFL. Moreover, with AsyncFL, the number of server updates per unit time increases nearly linearly with concurrency. When comparing both approaches in terms of wall-clock time to reach a target test loss, we show that AsyncFL is almost 5×\times faster and 8×\times more communication-efficient than SyncFL.

Finally, we demonstrate that AsyncFL achieves more fair models than SyncFL with over-selection. We observe very high correlation between slow devices and devices with many training samples. Discarding the updates from slow devices results in biasing the model trained using SyncFL with over-selection: the test perplexity for clients in the 99th percentile increases by 53% when enabling over-selection. This bias is not introduced when training with AsyncFL.

Understanding the landscape of federated learning at-scale

Building a robust federated learning system faces key design challenges:

System and data heterogeneity, where client devices participating in FL exhibit different system characteristics and possess different amounts of training data, leading to large differences in training time, and

Scalability, where the training time speedup with higher degree of concurrency experiences diminishing return and plateaus quickly.

To demonstrate the impact of the aforementioned challenges faced by FL, we begin by examining the degree of data and system heterogeneity observed in production when hundreds of millions of client devices jointly train a global model. To understand the limit of SyncFL approaches, we take a data-driven approach to demonstrate the impact of scale on the state-of-the-art synchronous model aggregation protocol.

System and Data Heterogeneity. Compute capabilities of mobile devices in the field differ by an order of magnitude Wu et al. (2019). Moreover, the number of training examples also varies widely across users Caldas et al. (2018). In combination, system and data heterogeneity can result in large differences in training time. Variance in training time results in stragglers that slow down the overall training time in SyncFL.

Figure 2 shows the distribution of training times across millions of clients for a common FL application (language model training, Section 7). The per-client training time distribution spans more than two orders of magnitude. When running SyncFL with concurrency and aggregation goal set to 1000, the average round completion time is 21×21\times larger than the mean client training time.

To mitigate the impact of stragglers in SyncFL, some systems use over-selection Bonawitz et al. (2019). In Section 7.4, we show that over-selection causes sampling bias, thus producing models that are unfair to stragglers.

Scalability. To further minimize the training time to convergence, a straightforward approach is to scale up the overall training throughput of the FL system by increasing the degree of training concurrency. Figure 3 illustrates the training time to convergence and the communication overhead for the SyncFL method FedAdam Reddi et al. (2020) as the number of concurrently training users increases from 130 to 2600. As concurrency increases, training time decreases slowly, while communication resource consumption increases much faster. For example, doubling the concurrency from 1300 to 2600 decreases the overall training time by only 17% while increasing communication costs by 73%.

We need resilient solutions that handle data and system heterogeneity at scale. At the same time, as shown in Figure 3, we are at the scaling limit of synchronous model aggregation. To build FL suitable for billions of clients, we need a fundamentally different model aggregation protocol that is resilient to heterogeneity (client independence), scalable to large cohort sizes (beyond the order of hundreds), and secure (asynchronous secure aggregation). Next we describe the proposed design of Papaya and demonstrate how AsyncFL can improve large-scale FL by improving scalability and straggler resilience.

Proposed Design

In this section, we first describe the AsyncFL algorithm Papaya uses. Next, we discuss the challenges in implementing AsyncFL in a large-scale production system.

Papaya implements a recently proposed AsyncFL algorithm, FedBuff Nguyen et al. (2021). In FedBuff, there is no notion of rounds: clients download, train, and upload updates asynchronously (Figure 4). After a client finishes local training, it uploads the model update (difference between the trained local model and initial model it received from the server before training). The aggregator tracks progress towards an aggregation goal, the number of client updates that need to be received before the server performs a model update. As soon as the aggregation goal has been achieved, the aggregated update is released and the server model update is performed. Each client update is weighted by the number of examples the client trained on and a factor depending on the staleness of the update. Staleness is defined as the difference between the model version that a client uses to start local training and the server model version at the time when a client uploads its model update. For example, Figure 4 shows FedBuff with 4 concurrent users and an aggregation goal of 2. Device A’s update has a staleness of 1 since the server model was updated once while Device A was training. In the rest of the paper, AsyncFL refers to our implementation of the FedBuff algorithm in Papaya.

We show in Section 7 that in a large-scale production setting with system and data heterogeneity, AsyncFL is faster and more resource efficient than SyncFL. However, AsyncFL brings a unique set of challenges that require careful system design.

2 System Design challenges in AsyncFL

Existing large-scale FL systems are designed to run SyncFL Bonawitz et al. (2019); Paulik et al. (2021). Hence, their architectures are not compatible with asynchronous training. There are four main reasons for this incompatibility, which we discuss next.

Client Selection. Client selection in SyncFL is based on forming synchronous cohorts. For example, in Bonawitz et al. (2019) a client cannot begin training until the entire cohort of clients has been selected. To support AsyncFL, we design a client selection mechanism that avoids any inter-client dependencies (Section 6.1).

Secure Aggregation. Secure Aggregation (SecAgg) improves the privacy of FL algorithms by hiding individual client model updates ensuring that the server can only view the final aggregation of all model updates. Most FL systems implement SecAgg based on secure multi-party computation (SMPC) Bonawitz et al. (2016); So et al. (2021b). SMPC-based SecAgg requires clients participating in a round to form a cohort and run a multi-leg protocol through the duration of the round. These requirements are not compatible with asynchronous training.

Motivated by these challenges, we propose a novel incremental Asynchronous Secure Aggregation algorithm that uses a Trusted Execution Environment Karl et al. (2020b) in Section 5.

Client Replacement for High Utilization. Cohort-based SyncFL systems do not replace clients in the middle of a round (Bonawitz et al., 2019). However, AsyncFL requires continuous replacement of clients that have finished training or have failed. We describe a fast client replacement mechanism that enables our AsyncFL implementation to achieve close to 100% client utilization, significantly higher than SyncFL (Section 6.2).

Support for Fast Model Aggregation. AsyncFL generates up to 30×\times more server model updates per unit time than SyncFL, as shown below in Figure 8. We design our system for fast model aggregation that can support much higher throughput of server model updates (Section 6.3) than what typical SyncFL systems can achieve.

In the next sections, we describe the design of our production system and explain how it supports the four requirements above.

System Components

The Papaya high-level design involves two applications: a server application that runs on a server in the data center, and the client application that runs on end-user devices. The server has three main components: Coordinator, Selector, and Aggregator. While the number of Selectors and Aggregators can scale elastically based on the workload demand, there is only one Coordinator; see Figure 5.

Papaya’s system architecture is influenced by the Google FL stack (GFL) described in Bonawitz et al. (2019). We use the same names for the main components, and their functions are similar to those of GFL. However, their implementation and interactions are substantially different. GFL supports only SyncFL, whereas Papaya supports both SyncFL and AsyncFL. As a result, our design has fundamental differences in the protocol, execution, and scalability which enable it to achieve faster model convergence and straggler resilience; these differences are discussed further in Section 8. First, we briefly describe the responsibilities of the main components and their interactions.

Coordinator. The Coordinator performs three main functions. First, it assigns FL tasks to Aggregators, as discussed in Section 6.3. Second, the Coordinator assigns clients to FL tasks, as described in Section 6.1. Finally, it provides centralized coordination and ensures that tasks progress in the face of Aggregator failures.

Selector. The Selector is the only component that directly communicates with clients. When necessary, it forwards client requests to other components. The Selector has two main responsibilities. For client selection, it advertises available tasks to clients, and summarizes current client availability for the Coordinator, as described in Section 6.2. For client participation, the Selector routes client requests to the corresponding Aggregator, as described in Section 6.3.

Aggregator. Every task is assigned to a single Aggregator for the duration of the task (apart from failures and network partitions), as described in Section 6.3. The Aggregator has three main responsibilities. First, it aggregates client model updates to produce new versions of the server model. Second, it drives participating clients to run the client execution protocol, as described in Section 6.1. Finally, it tracks whether or not a task needs more clients and reports this to the Coordinator, as discussed in Section 6.2.

Client Runtime. The client runs on end-user devices and monitors training eligibility criteria such as whether or not the device is idle. It also tracks prior participation history to enable fair and unbiased client selection. If a client is eligible for training, the client checks in with the server to execute the FL client protocol as described in Section 6.1.

Secure Aggregation

In this section, we summarize our SecAgg mechanism to enable AsyncFL. In an honest-but-curious threat model, SecAgg allows the server to compute aggregated client updates without observing individual client updates. There are two main approaches for implementing SecAgg: using Secure Multiparty Computation (SMPC) or a Trusted Execution Environment (TEE).

Existing SMPC-based SecAgg approaches Bonawitz et al. (2016); Bell et al. (2020); So et al. (2021b) hinder asynchronous training, as they require cohort formulation and inter-client communication in each round.A concurrent work So et al. (2021a) describes an SMPC method that may overcome some of these issues. This approach could be an alternative to the TEE-based approach described here. Meanwhile, AsyncFL does not have a discrete notion of rounds; clients join and finish training asynchronously.

On the other hand, naive TEE aggregation is unscalable. Asymptotically, this approach transmits O(K⋅m)O(K\cdot m) data across the host-TEE boundary, where KK is the aggregation goal and mm is the model size. Transferring data across the TEE boundary is time-consuming (Figure 6): taking nearly 650 milliseconds for 100 clients (K=100K=100), each with a 20MB model. This data transfer time increases with aggregation goal. Trusted hardware trades performance for security guarantees.

Motivated by these challenges, we propose an Asynchronous SecAgg mechanism, relying on a TEE and an attestation mechanism; ensuring the Trusted Secure Aggregator (TSA) has not been tampered with. In this approach, random masking relies on an additive one-time-pad to protect client updates and utilizes the TSA to generate aggregated random masks, unmasking aggregated client updates. The overall mechanism depends on a secure virtual channel established between each client and the TSA using the Diffie-Hellman key exchange protocol (Merkle, 1978). Then, the mechanism leverages the TSA’s ability to regenerate a random unmask based on the clients’ secret received by the TSA over the secure channel.

Asynchronous SecAgg empowers client independence and fast incremental aggregation. The protocol consists of the following steps: (1) A participating client establishes a secure virtual channel with the TSA and validates the secure aggregation configuration and integrity of the TSA; (2) The client shares the masked model update with a corresponding Aggregator and the random seed used to generate the mask with the TSA; (3) The aggregator incrementally aggregates masked model updates; (4) The aggregator requests the TSA to generate the unmasking vector once the configured aggregation goal is reached; (5) The aggregator unmasks the aggregated model updates using the unmasking vector and creates a new server model.

The random seed, usually 16 bytes shared between each client and the TSA, allows the two parties to share an as-large-as-the-model mask at a constant cost. Asymptotically, this approach only transmits O(K+m)O(K+m) data across the boundary of the TSA. Appendix B presents more details about our secure aggregation protocol, including a security proof.

System Design

In this section, we describe the system requirements to run AsyncFL at scale and the design choices we made to fulfill these requirements. We focus on the three most important requirements for AsyncFL outside of Asynchronous SecAgg (discussed in Section 5). For completeness, other requirements for running AsyncFL are described in Appendix E.

There are three main requirements. First, AsyncFL relies on clients training asynchronously. Hence, the client protocol must not introduce any dependence between clients. Second, AsyncFL can support higher client utilization than SyncFL. To enable this, our system must perform fast client replacement. Third, AsyncFL takes many more server model steps than SyncFL per unit time. Hence, our system must support fast model aggregation.

To enable asynchronous training, Papaya’s client protocol deliberately avoids any inter-client dependency. Moreover, transient client failures do not cause clients to dropout because the client protocol is based on virtual sessions instead of persistent connections. At a high level, the protocol can be split into two phases: selection and participation. To explain the selection process, we first define client demand for a task as the difference between the target concurrency and the number of users already participating in the training of the task.

Selection. For a client, the goal of the selection phase is to find a task with positive client demand. Thus, a client can complete the selection phase with either acceptance (client is accepted for participation) or rejection (client will try to participate at another time).

Participation. Once a client is accepted, the goal of the participation phase is for a client to share a trained model with the server. Participation consists of four stages. 1. A client first downloads model parameters, model code and configuration from a content delivery network. 2. Next, the client trains the downloaded model on its local data. 3. Once the training finishes, the client reports its status to the server. The server shares an upload configuration with the client and, if enabled, the SecAgg configuration. 4. In the final stage, the client uploads the model in chunks, potentially after masking the model if SecAgg is enabled. All stages happen within a virtual session established during selection.

2 High Client Utilization

AsyncFL is capable of higher client utilization compared to SyncFL. This is mainly because in SyncFL the number of active clients increases at the beginning of a round as clients join the cohort, and it falls gradually towards the end of the round as the server waits for all clients to finish training (Figure 7). On the contrary, in AsyncFL there is no cohort formation; as soon as one client completes training or fails, a new one is selected. Thus AsyncFL achieves high utilization throughout training. We show in Figure 7 that utilization in our AsyncFL implementation is close to 100% throughout training. To realize high utilization, an AsyncFL system needs to replace completed and failed clients quickly. Achieving high utilization is especially challenging in a multi-tenant FL system, where multiple FL tasks are running in parallel, and a single client may be compatible with many tasks.

We now describe the client assignment process which is responsible for maintaining high utilization. There are three important steps to assigning clients to tasks: tracking client demand for each task, tracking task eligibility for each client, and performing the actual assignment.

Tracking client demand for each task. First, each Aggregator tracks client demand for the tasks that are assigned to it. When a client finishes training or fails, the Aggregator increases client demand for the associated task. Next, the Coordinator pools together information from all Aggregators into a consolidated view of client demand for every task in the system. Note that the Coordinator must explicitly account for clients that have been assigned to a task, but have not yet confirmed the assignment.

Tracking task eligibility for each client. For each available client, the Coordinator constructs a list of eligible tasks. A task is eligible if the client is compatible with its requirements (e.g., can train the model of the task), and if the task has positive client demand.

Task assignment. Once an eligible task list is constructed for a client, the Coordinator randomly assigns the client to an eligible task. Concretely, the Coordinator instructs Selectors to forward the client to the Aggregator responsible for the task.

3 Fast Model Aggregation

As shown in Figure 8, AsyncFL generates server model updates up to 30×\times more frequently than SyncFL. Thus, fast model aggregation in a scalable AsyncFL system is critical. In this section, we describe how Papaya efficiently aggregates client updates.

Persistent Aggregator. In our system, Aggregators are persistent and stateful because creating a new Aggregator for each task incurs a substantial overhead. Therefore, the Coordinator moves tasks between Aggregators only when it detects failed or overloaded Aggregators. The Coordinator evenly distributes tasks among available Aggregators using the estimated workload of a task. The Coordinator estimates this workload using the task concurrency and model size.

Parallel Model Aggregation. Once a client completes training, it uploads the trained serialized model update to the server. This update is then pushed into an in-memory queue on the Aggregator. A different thread drains the queue by de-serializing the updates into trainable parameters and aggregating them. To speed up this aggregation, we parallelize the aggregation process across available cores. To reduce lock contention, the ID of the thread performing intermediate aggregation is hashed to choose one of the intermediate aggregates. Once the cumulative number of aggregated model updates reaches the aggregation goal, the final aggregation is performed and a new server model is generated. Note that the aggregation goal in SyncFL is typically 1.3×1.3\times concurrency (30% over-selection), while in AsyncFL it is independent of concurrency.

Evaluation

In this section, we present evaluation results for AsyncFL. We first compare the convergence speed, scalability, and communication efficiency of AsyncFL with SyncFL. Next, we analyze the source of AsyncFL’s speed up. We then show that SyncFL can be either straggler resilient (with over-selection) or be unbiased (without over-selection), but not both simultaneously. In contrast, AsyncFL is straggler resilient without introducing bias.

To study the performance of AsyncFL, we train an LSTM-based language model Kim et al. (2015), a common FL application Hard et al. (2019), on a population of nearly 100 million Android phones. We repeat each experiment 3 times, each at the same time of the day, and report the average. AsyncFL and SyncFL are run at the same time, so they have access to the same client population.

Following the requirements in Hard et al. (2019), a client device can participate in FL training only when idle, charging, and on an unmetered network. Similar to Bonawitz et al. (2019), a timeout is imposed to limit the client training time; we set the timeout to 4 minutes. The distribution of client execution times is analyzed in Section 7.4.

For both SyncFL and AsyncFL, we use SGD on the client and FedAdam Reddi et al. (2020) on the server. For the server optimizer, we use Adam’s default learning rate and tune the first-moment parameter in simulation. We run hyperparameter sweeps in simulation, using a representative dataset, for the client optimizer to find the best client learning rate. Each client runs one local epoch of training with batch size B=32B=32. We partition each client’s data into train, test, and validation sets randomly.

Our system has two configuration parameters. First, for both SyncFL and AsyncFL tasks, concurrency specifies the maximum number of concurrently participating devices. Second, for AsyncFL tasks, KK is the aggregation goal, controlling the size and frequency of server update. The server produces a new model every KK client model updates. In our experience, we find that choosing KK to be 10-30% of concurrency works well in practice. Finally, unless otherwise stated, we use 30% over-selection with SyncFL to alleviate the impact stragglers, as proposed in Bonawitz et al. (2019).

2 Results on Convergence Time and Scalability

To begin, we evaluate the training time performance and scalability of AsyncFL and SyncFL. We measure the convergence time as the wall-clock training time to reach a target loss. To measure scalability, we compare AsyncFL and SyncFL in terms of their speedup and the number of communication trips with increasing concurrency. Figure 9 shows (left) the time taken by the two algorithms to reach a target loss for varying levels of concurrency, (middle) the speedup of AsyncFL over SyncFL, and (right) the number of communication trips to reach a target loss. As Figure 9 (left and middle) illustrates, the speedup gap widens as concurrency increases, from 2×2\times to 5×5\times. Furthermore, the SyncFL communication efficiency worsens. The overall communication efficiency gain of AsyncFL increases from 2×2\times to 8×8\times as concurrency increases. The evaluation results demonstrate that AsyncFL handles the system heterogeneity and scalability challenges more effectively than SyncFL. In the following sections, we unravel why our AsyncFL system is more suitable for scaling FL training to hundreds of millions of clients.

3 Analysis of Server-Model Step Frequency

To understand the performance of our AsyncFL implementation in detail, we study how the aggregation goal impacts convergence. The aggregation goal determines how many client updates contribute to each server model update, and for a fixed concurrency, it also affects the frequency of server model updates. We fix concurrency to be 1300 and vary aggregation goal (KK) from 100 to 1300. Figure 10 (top) depicts the time for each configuration to reach a target loss, while Figure 10 (bottom) describes the server-model step frequency per hour. As KK increases, the batch size increases, and the server takes less frequent model steps. Thus, the larger the KK is, the slower the convergence time. It is natural to ask if convergence time could be further reduced for KK smaller than 100. However, Nguyen et al. (2021) found that moderate values of KK can lead to more stable convergence. Moreover, the frequency of server updates is limited by the system’s write bandwidth. Thus, we cannot create a new server model too often. We leave improvements to overcome write bandwidth limitations as future work.

4 Analysis of Sampling Bias from Over-Selection

To compare the effectiveness of over-selection and asynchronous training in combating stragglers, we examine the distribution of participating clients, their execution time, and the number of training examples. Figure 11 illustrates the discrepancy between the client execution time distribution of SyncFL with and without over-selection. Since over-selection discards updates from some clients, the distribution of SyncFL without over-selection should be considered representative of the entire client population. Figure 11 (top-left) shows that over-selection drops slow clients, as desired. However, as illustrated in Figure 11 (top-right), the slowest clients often have more training examples.

To rigorously assess the difference, we perform a two-sample Kolmogorov-Smirnov test Chakravarti et al. (1967) to measure the goodness of fit between AsyncFL, SyncFL with over-selection, and the ground truth distribution (SyncFL without over-selection). We find that the D-statistic, representing the absolute max distance between the cumulative distribution functions of the two samples, for AsyncFL and the ground truth is 8.8×10−48.8\times 10^{-4} (pp-value = 0.98). In comparison, the D-statistic for SyncFL with over-selection and the ground truth is 6.6×10−26.6\times 10^{-2} (pp-value = 0.0). This result shows that AsyncFL and the ground truth have similar distributions while SyncFL with over-selection does not. Thus, over-selection introduces sampling bias while AsyncFL does not. The sampling probability is conditioned on the client’s device speed or the number of training examples. Next, we show that sampling bias hurts model performance, especially for clients with more training examples.

Table 1 reports the model quality in test perplexity for all clients and those with data volume in the 75% and 99% percentile. Perplexity is a measure of language model accuracy (lower is better). The sampling bias from over-selection in SyncFL causes a 6% drop in model quality overall and a 50% drop in model quality for clients with more examples. Although SyncFL without over-selection is unbiased, it is also 10×\times slower. On the other hand, AsyncFL combines fast training with high model quality and no sampling bias. Meanwhile, SyncFL with over-selection has to choose between sampling bias or straggler resilence. In summary, AsyncFL is a more desirable method to address the impact of stragglers.

5 Understanding AsyncFL Advantages

The previous sections showed that AsyncFL has two main advantages over SyncFL: better scalability because of more frequent model steps and straggler resilience without adding sampling bias. To quantify the benefits from these two properties, we present the training curves for AsyncFL alongside the current state-of-art SyncFL. We remove the frequent update advantage of AsyncFL by increasing the aggregation goal for AsyncFL to be the same as SyncFL.

Figure 12 depicts production training curves for the best synchronous setup, SyncFL with over-selection (orange), and two AsyncFL configurations: aggregation goal KK = 100 (red) and KK = 1000 (blue). All three use concurrency 1,300. Note that overall, AsyncFL with KK = 100 is 4.3×\times faster than SyncFL with over-selection, as shown in Figure 13. We find that about half of this speedup comes from using smaller KK and the rest from avoiding sampling bias (i.e., using AsyncFL rather than SyncFL with over-selection).

To read Figure 12, start with AsyncFL with KK = 100 (red), which is the best configuration since it takes more frequent server model steps and is resilient to stragglers. Next, see AsyncFL with KK = 1000 (blue), which is straggler resilient but takes less frequent model steps. Finally, move to SyncFL with over-selection (orange), which adds sampling bias.

The figure also shows SyncFL without over-selection (green) for reference, using concurrency 1000. The large gap between this configuration and AsyncFL with KK = 1000 is attributable to stragglers.

It is instructive to compare the training loss at a fixed point, e.g., at the 10-hour mark. By minimizing sampling bias, AsyncFL with KK = 1000 reduces training loss by 3.4% compared to SyncFL with over-selection. Taking more frequent server-model steps (KK = 100) in AsyncFL decreases training loss by an additional 3.5%.

Related Work

The Papaya system described in this paper is inspired by the GFL system Bonawitz et al. (2019). Another FL system is described by Apple (AFL) in Paulik et al. (2021). We focus on comparison with GFL and AFL given the similarity of production scale. At a high level, both GFL and AFL only implement SyncFL, while Papaya implements both SyncFL and AsyncFL. Diving deeper we find similarities and differences in how clients are selected for participation, client availability and participation outcome impact on model training progress, model update aggregation and privacy mechanisms. Papaya actively (through the Coordinator) selects available clients for participation at any point in time based on demand by active tasks (driven by desired task concurrency), unlike GFL which actively selects clients before the rounds starts, and AFL which uses passive probabilistic selection. Papaya enables incremental progress by making participating clients independent and replaceable whenever they complete or drop out, unlike GFL where no client can join after a round has started, potentially leading to failed rounds, and similar to AFL where clients can contribute as long as the model version is the same. Papaya moves tasks between long living Aggregators only when failure or load imbalance is detected to minimize client progress loss and reduce placement overhead, unlike GFL where tasks are dynamically placed to ephemeral Aggregators every round and AFL where aggregation is performed by an offline service. Papaya implements Asynchronous SecAgg based on TEEs, whereas GFL uses SMPC-based Synchronous SecAgg and AFL does not report using SecAgg.

Another line of related work is the FL software tool kits, offered by other technology companies. Notable among these are Clara NVIDIA , IBM-FL IBM , OpenFL Reina et al. (2021) and FATE WeBank . While related to the Papaya system described in this paper, these software tools are distinct from production FL systems training across hundreds of millions of devices, which is the focus of this paper.

Conclusions

We presented our design of a production asynchronous FL system for training at scale. Designing for a production FL system, Papaya, we find that AsyncFL is faster, more straggler resilient, and provides better model quality than SyncFL. Papaya is flexible and supports both synchronous and asynchronous FL. Empirically, we demonstrated that in high concurrency settings, asynchronous FL achieves 5×\times faster speed up and conserves nearly 8×\times more resources than synchronous FL. Finally, Papaya can be extended with features to enable differential privacy, which we leave as future work.

Acknowledgements

We thank Ilya Mironov and Rachad Alao for meaningful discussions and their valuable support which significantly improved the quality of this paper.

References

Supplementary Material

Appendix A Cryptographic Primitives

Diffie–Hellman key exchange protocol allows two parties to securely agree on a randomly-generated shared secret via an untrusted communication channel. Viewed in the server-client setting, the protocol consists of an initial message from one party (server) and a completing message as a response from the other one (client). The server can prepare the initial messages in advance, without knowing the identities of the clients. The client can solely determine the shared secret once it receives the initial message. The client needs to interact with the server only once to finish the protocol by sending the completing message.

A.2 Additive One-time pad (OTP)

There are many existing additive homomorphic encryption schemes such as the works in ElGamal (1985); Goldwasser & Micali (1982); Paillier (1999); Cohen & Fischer (1985). The complexity of the decryption algorithms in these protocols is usually linear in the ciphertext size and is independent of the number of additions. However, these schemes often operate on a large finite group whose elements can be as large as 1024–3072 bits. Such requirement inflates the ciphertext size even if the plaintext space is much smaller (e.g., 32 bit integers). Such blow-up makes these schemes less desirable when ciphertexts are transmitted via network and traffic is at a premium, for example, on mobile devices.

A PRNG-generated additive one-time pad (OTP) is a good alternative to avoid the expansion of the ciphertext. The protocol is summarized in Figure 14.

Mobile devices are often restricted in both computation power and communication bandwidth. An additive OTP is more efficient in both computation and bandwidth cost, compared to the group operations needed in other encryption schemes.

The server usually has much more computation resources to perform the relatively more expensive decryption. Furthermore, hardware acceleration optimizations are often available server-side, reducing the costs of the decryption algorithm.

Appendix B Protocol Design and Security Proof

In this section we will go over the detailed design of our protocol and formally prove its security. We adopt the same strategy from Cryptonite Karl et al. (2020a). A trusted party realized by trusted hardware (e.g., Intel SGX) will assist with the procedure and help make up for the dropped clients. With the assistance of the trusted hardware, clients no longer rely on each other to protect their own private inputs or mitigate the dropout of their peers. Without client interdependence, clients no longer need to communicate with each other via the server and no longer need to know about each others’ identities. The absence of interdependency requirement among clients allows them to participate asynchronously, making our protocol compatible with FedBuff Nguyen et al. (2021).

The thread model is composed of a server, a trusted third party and nn clients. The trusted third party and the clients can only communicate directly with the server. The clients can choose to participate in the protocol at any time, not necessarily in the beginning. Instead, they will check-in with the server when they become available. Clients may have limited availability. The availability of any two clients may have no overlap on the timeline.

The ideal functionality is summarized in Figure 15.

B.2 Overview of Our Solution

To avoid sending big chunks of data across the boundary of the secure enclave, we will aggregate random masks, instead of the actual data, inside the secure enclave. The high-level idea is to mask clients’ private inputs with some additive masks, while the server will be responsible for aggregating the masked inputs and trusted party will be responsible for aggregating the masks. Note that a 128-bit seed is sufficient to represent a random mask. The amount of data transferred into the secure enclave for each client will be a constant, despite of amount of data to aggregate. Our protocol can be divided into three steps:

New client checks in and validates the identity of the trusted party.

Client sends masked input to the untrusted server and demasking information to the trusted party.

The trusted party instructs the untrusted server how to demask the sum of all masked inputs.

B.3 Protocol Detail

Our protocol is detailed in Figure 16. We use Diffie–Hellman key exchange protocol to establish private communication channels between the trusted party and the clients.

B.4 Security Proof

We adopt the simulation-based proof technique to show that the ideal functionality and the real world protocol are computationally indistinguishable. Let Cc⊂C\mathcal{C}_{c}\subset\mathcal{C} denotes the indices of all the clients corrupted by the adversary and Ccˉ⊂C\bar{\mathcal{C}_{c}}\subset\mathcal{C} denotes the indices of the honest clients. Our strategy is to construct a simulator with the following properties (Figure 17):

The simulator runs the adversary as a subroutine.

The simulator executes the real world protocol with the adversary. The simulator plays the role of the trusted party and client ii for all i∈Ccˉi\in\bar{\mathcal{C}_{c}}. The adversary plays the role of the server and corrupted clients ii for all i∈Cci\in\mathcal{C}_{c}.

The simulator executes the ideal functionality with the ideal functionality and honest clients. The simulator will play the role of the server and all the clients in Cc\mathcal{C}_{c}.

We prove that the joint view of the adversary as the simulator’s subroutine is computationally indistinguishable from that of a real world execution. The detailed description of the simulator is in Figure 18.

We now argue that the adversary’s views are the same in either the simulation or a real world execution with real honest clients. We will show that by a series of computationally indistinguishable hybrid experiments.

Hybrid0\mathbf{Hybrid}_{0}: The simulator executes the real world protocol with the adversary. The simulator plays the role of honest clients with their private inputs viv_{i}. The adversary plays the role of the server and corrupted clients. This is exactly the real world protocol execution.

Hybrid1\mathbf{Hybrid}_{1}: The same as Hybrid0\mathbf{Hybrid}_{0}, except:

In Item 6, upon receiving a DH key exchange response with successfully decrypting the encrypted seed from client ii:

if i∈Cci\in\mathcal{C}_{c}, the simulator follows the protocol;

if i∈Ccˉi\in\bar{\mathcal{C}_{c}}, the simulator adds ii to Ca\mathcal{C}_{a};

If decrypting the encrypted seed fails, ignore the update.

Hybrid0≈Hybrid1\mathbf{Hybrid}_{0}\approx\mathbf{Hybrid}_{1}:

Hybrid2\mathbf{Hybrid}_{2}: The same as Hybrid1\mathbf{Hybrid}_{1}, except:

The simulator no longer has the inputs of honest clients, but runs the ideal functionality with real honest clients.

In Item 7, if the trusted party is expected to generate an unmasking vector, the simulator interacts with the real honest clients via the ideal functionality:

For each i∈Cci\in\mathcal{C}_{c}, the simulator sends out a -vector to the ideal functionality on behalf of the corrupted client ii.

The simulator sends Ca\mathcal{C}_{a} to the ideal functionality.

The simulator receives VV, which is the sum of all inputs of honest clients, as the server from the ideal functionality.

In Item 8, no matter what the adversary outputs in the real world protocol, the simulator outputs the same content as the same role in the ideal functionality.

This is exactly the ideal functionality execution with the simulator.

Hybrid1≈Hybrid2\mathbf{Hybrid}_{1}\approx\mathbf{Hybrid}_{2}: The indistinguishability comes from the fact that

The ideal functionality will correctly sum up the honest clients’ inputs.

The simulator’s output in the ideal functionality for each role it plays is identical to the adversary’s output in the real world protocol for the same role.

Appendix C Deployment with Intel SGX

The secure aggregation protocol (Figure 16) we present in Section B.2 involves a trusted party. When deploying this protocol, we use a Intel SGX enclave to play the role of this trusted party. To enforce the honest behavior of this trusted party, we have to ensure the following two security guarantees.

Confidentiality: the trusted party realized by the enclave shares no information with any other party except what is specified in the protocol.

In this section, we see how we employ remote attestations and verifiable logs to ensure these properties.

Remote attestation technique was originally designed to allow an enclave owner to verify the identity of a trusted binary executed in the cloud. In our use case, there is nothing secret about the code or the initial parameters inside the enclave. Therefore these data can be provided to the clients to allow them play the role of enclave owner and to verify the identity of the trusted binary that plays the role of the trusted party. To be more specific, the extra steps specified in Figure 19 are taken on top of the secure aggregation protocol in Figure 16 to ensure the honest behavior of the trusted party.

We follow the standard assumptions of SGX:

It is infeasible to forge an attestation quote that does not match the running trusted binary and/or the hash of public parameters as the custom payload, but can be verified against Intel’s collateral.

It is infeasible to tamper with the trusted binary executed inside the enclave.

It is infeasible to access the data stored inside the enclave except via the predefined APIs.

Under these assumptions and other standard assumptionsIncluding: 1. the hash algorithm we use is collusion resistant; and 2. AES is a secure block cipher. , clients accept an attestation quote only if:

the quote is generated by a legitimate enclave;

the enclave is running the predefined code;

the enclave is running with server-claimed parameters;

These arguments jointly assert the enclave is faithfully playing the role of the trusted party. The clients will proceed in the protocol with their private inputs only if they can validate the faithful trusted party. With that said, the server will not hear back from clients unless attestation quotes from a legitimate enclave with correct trusted binary and parameters are forwarded to the clients.

In addition, the server cannot successfully tamper with the data that is meant to be sent into the enclave, i.e. the DH key exchange response and the encrypted seed. This is because the decryption fails if any of them is modified by the server. Furthermore, the encrypted seed and the response is not accepted by another enclave instance either since it will not have the necessary private randomness to recover the shared key correctly. In summary, the server must use exactly the same enclave during the whole protocol, otherwise it is effectively dropping clients.

C.2 Updating the Trusted Binary with Verifiable Logs

Remote attestations allow clients to validate the trusted binary’s identity against a hardcoded hash. Such design makes it impossible to update the trusted binary in the future without updating the clients at the same time. To ease the updating process, verifiable logsver ; tri can be used to note down any changes made to the code that will run inside an enclave.

A verifiable log is implemented by a Merkle tree and append-only. Each new record appended to the end of the log is added as a new leaf in the underlying Merkle tree. The hash of the root of the Merkle tree serves as the snapshot of the whole log. An inclusion proof can be generated to demonstrate a record is indeed included in the log. A consistency check can be performed between two snapshots to decide if the corresponding append-only logs are consistent with each other.

There are several steps to integrate this technique, as detailed in Figure 20.

Note that both clients and auditors use the same API to request the log’s latest snapshot. Therefore the auditors and clients share the same snapshots. Due to the unforgeability of the underlying secure hashes, any logged trusted binary cannot avoid audition without being noticed. On the other hand, clients will only proceed in the protocol only if the trusted binary is logged. In summary, no trusted binary that interacts with clients can avoid audition without getting caught.

With this auditing mechanism in place and sufficient public auditors watching the latest snapshots, the trusted binary can be updated on a regular basis without updating on the client side.

Appendix D Fixed Point Conversion

Appendix E Additional System Design Details

To prevent unbounded client participation, the system enforces an upper bound of concurrently participating clients (C) for every task based on task configuration. A client can be selected for a task only if number of active clients is below the configured threshold. An active client may become inactive for various reasons. The client may have completed execution, or it may be considered dead due to missed heartbeats or execution error. Finally, clients may also be aborted by the server if staleness (measured as the difference between current and initial model versions) is higher than a configurable value.

E.2 Handling Staleness

The cost of asynchronous training is staleness of model updates. In this section, we describe how our AsyncFL system tracks and handles staleness. The server model is identified by a model version—a non-negative natural number that is incremented every time a new server model is generated. A new server model is generated when KK client updates have been aggregated. In an asynchronous FL systems, clients can download a model with an initial version, but upload results when the server model has moved to a different final version. Recall that we define staleness as the difference between the model version in which a client uses to start local training, and the server model version at the time instance when a client uploads its update. For each client, the aggregator records initial model version to track staleness. Let sis_{i} be the staleness of client ii with initial version VinitialV_{\text{initial}} and final version VfinalV_{\text{final}}; thus si=Vfinal−Vinitials_{i}=V_{\text{final}}-V_{\text{initial}}. We down-weight client ii’s update using the same scheme as Nguyen et al. (2021). Formally, let wiw_{i} be the weight of client ii whose staleness is sis_{i}, then wi≔1/1+siw_{i}\coloneqq 1/\sqrt{1+s_{i}}. Finally, to bound staleness, the aggregator abort clients whose staleness is larger than a configurable parameter, maximum staleness.

After every server model update, the aggregator aborts clients whose staleness is larger than a configurable parameter, maximum staleness.

E.3 Switching between SyncFL and AsyncFL

Papaya highlights that an FL system can support both synchronous and asynchronous training by using client independence, fast model aggregation, high client utilization and asynchronous secure aggregation. These properties improve the performance of both training regime.

Switching from SyncFL to AsyncFL in our system requires three small changes in behavior: client demand computation, handling of stale clients, and model aggregation.

Client demand computation. In AsyncFL, client demand is computed as concurrency−active_clients\mathit{concurrency}-\mathit{active\_clients}. However, in a typical SyncFL round, client demand is high in the beginning of a round, but decreases as clients report results (see Figure 7). In SyncFL, client demand is computed as concurrency⋅(1+o)−completed_clients\mathit{concurrency}\cdot(1+o)-\mathit{completed\_clients}, where oo is the over-selection factor.

Aborting stale clients. When a server model update is performed in SyncFL, users that are still training are aborted (users may still be training because of over-selection). In AsyncFL, users that are still training continue normally, unless their staleness would exceed maximum staleness.

Model Aggregation. AsyncFL and SyncFL use different model aggregation algorithms.

These three behavior changes are relatively minor. Thus, switching between SyncFL and AsyncFL can be done via a configuration change.

E.4 Failure Recovery

Fast recovery and isolated impact from failures help the system minimize model training progress impact. Below we outline mechanisms employed:

Client Routing. Client requests are routed by selectors using assignment maps (model training task to corresponding aggregator identity) refreshed from coordinator on every report. Upon selector failure or selector having stale assignment map clients retry through a different selector. Failed or stale selector refreshes assignment map on next report to coordinator.

Client Participation. Coordinator assigns clients to tasks. Upon coordinator failure participating clients are not affected, only for the duration of the recovery no new clients are assigned. Selectors and aggregators wait until a new leader coordinator is elected meanwhile continuing to operate based on last known assignments. After the leader election coordinator enters the recovery period (typically 30s) to rebuild the current assignment map from aggregator reports and then resumes assignments.

Task Execution. Aggregator executes assigned tasks. Upon aggregator failure or unresponsiveness, coordinator detects failures after several missed heartbeats and reassigns all tasks to other aggregators, updates and distributes the new assignment map to selectors. Coordinator detects stale assignments in aggregator reports via sequence numbers and requests to stop executing stale assignments.

E.5 Edge Training Engine

The Papaya client is built to be both a hosting platform and an ML framework. An Example Store collects training data in persistent storage and enforces the data use and retention policy. An Executor abstracts model training logic in a general way that supports easily swapping in different ML tasks (data source, model, loss, etc.). The implementation is based on PyTorch Mobile and relies on two features: selective build and the mobile interpreter. Selective build only compiles in ops used by the application to reduce the binary size. The mobile interpreter facilitates efficient cross-platform execution (Android, iOS, Linux) by providing common functionality to save and load model code and parameters, execute forward and backward passes, and optimizer steps.