Enquire Now
Medical Computer Vision · Clinical Diagnostics · PyTorch / TensorFlow · GPU Optimized · 2026

Speaker Verification Xvector

Tensor Pipeline · Custom Loss Formulations · Model Quantization · Accelerated Inference — A rigorous deep learning engineering project focused on automated pathological lesion segmentation and radiological disease classification. Architected for thesis defense viva presentations, IEEE reproduction, and high-throughput production deployment.

PyTorch
Core Framework
AMP FP16
Mixed Precision
TensorRT
Quantized Serving

Ai Reproducible Research

This project focuses on ai reproducible research using modern AI and machine learning techniques. The content below is adapted from research literature and practical implementation notes.

Advances and Open Problems in Federated Learning

Peter Kairouz7 * H. Brendan McMahan7∗ Brendan Avent21 Aurélien Bellet9

19 13 7

Mehdi Bennis Arjun Nitin Bhagoji Kallista Bonawitz Zachary Charles7

Graham Cormode23 Rachel Cummings6 Rafael G.L. D’Oliveira14

Hubert Eichner7 Salim El Rouayheb14 David Evans22 Josh Gardner24

Zachary Garrett7 Adrià Gascón7 Badih Ghazi7 Phillip B. Gibbons2

Marco Gruteser7,14 Zaid Harchaoui24 Chaoyang He21 Lie He 4

Zhouyuan Huo 20 Ben Hutchinson7 Justin Hsu25 Martin Jaggi4 Tara Javidi17

2 2 7

Gauri Joshi Mikhail Khodak Jakub Konečný Aleksandra Korolova21

Farinaz Koushanfar17 Sanmi Koyejo7,18 Tancrède Lepoint7 Yang Liu12

Prateek Mittal13 Mehryar Mohri7 Richard Nock1 Ayfer Özgür15

Rasmus Pagh7,10 Hang Qi7 Daniel Ramage7 Ramesh Raskar11

Mariana Raykova7 Dawn Song16 Weikang Song7 Sebastian U. Stich4

Ziteng Sun3 Ananda Theertha Suresh7 Florian Tramèr15 Praneeth Vepakomma11

2 5 7 8

Jianyu Wang Li Xiong Zheng Xu Qiang Yang Felix X. Yu7 Han Yu12

Sen Zhao7

Abstract

Federated learning (FL) is a machine learning setting where many clients (e.g. mobile devices or whole organizations) collaboratively train a model under the orchestration of a central server (e.g. service provider), while keeping the training data decentralized. FL embodies the principles of focused data collection and minimization, and can mitigate many of the systemic privacy risks and costs resulting from traditional, centralized machine learning and data science approaches. Motivated by the explosive

growth in FL research, this paper discusses recent advances and presents an extensive collection of open problems and challenges.

* Peter Kairouz and H. Brendan McMahan conceived, coordinated, and edited this work. Correspondence to kairouz@

1 Introduction 4

1.1 The Cross-Device Federated Learning Setting . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 5

1.1.1 The Lifecycle of a Model in Federated Learning . . . . . . . . . . . . . . . . . . . . . . . . 7

1.1.2 A Typical Federated Training Process . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 8

1.2 Federated Learning Research . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 9

1.3 Organization . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 10

2 Relaxing the Core FL Assumptions: Applications to Emerging Settings and Scenarios 11

2.1 Fully Decentralized / Peer-to-Peer Distributed Learning . . . . . . . . . . . . . . . . . . . . . . . . . 11

2.1.1 Algorithmic Challenges . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 12

2.1.2 Practical Challenges . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 14

2.2 Cross-Silo Federated Learning . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 14

2.3 Split Learning . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 16

2.4 Executive summary . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 17

3 Improving Efficiency and Effectiveness 18

3.1 Non-IID Data in Federated Learning . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 18

3.1.1 Strategies for Dealing with Non-IID Data . . . . . . . . . . . . . . . . . . . . . . . . . . . . 19

3.2 Optimization Algorithms for Federated Learning . . . . . . . . . . . . . . . . . . . . . . . . . . . . 20

3.2.1 Optimization Algorithms and Convergence Rates for IID Datasets . . . . . . . . . . . . . . . 21

3.2.2 Optimization Algorithms and Convergence Rates for Non-IID Datasets . . . . . . . . . . . . 25

3.3 Multi-Task Learning, Personalization, and Meta-Learning . . . . . . . . . . . . . . . . . . . . . . . . 28

3.3.1 Personalization via Featurization . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 28

3.3.2 Multi-Task Learning . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 28

3.3.3 Local Fine Tuning and Meta-Learning . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 29

3.3.4 When is a Global FL-trained Model Better? . . . . . . . . . . . . . . . . . . . . . . . . . . . 30

3.4 Adapting ML Workflows for Federated Learning . . . . . . . . . . . . . . . . . . . . . . . . . . . . 30

3.4.1 Hyperparameter Tuning . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 31

3.4.2 Neural Architecture Design . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 31

3.4.3 Debugging and Interpretability for FL . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 32

3.5 Communication and Compression . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 32

3.6 Application To More Types of Machine Learning Problems and Models . . . . . . . . . . . . . . . . 34

3.7 Executive summary . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 34

4 Preserving the Privacy of User Data 36

4.1 Actors, Threat Models, and Privacy in Depth . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 37

4.2 Tools and Technologies . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 38

4.2.1 Secure Computations . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 40

4.2.2 Privacy-Preserving Disclosures . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 44

4.2.3 Verifiability . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 46

4.3 Protections Against External Malicious Actors . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 48

4.3.1 Auditing the Iterates and Final Model . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 49

4.3.2 Training with Central Differential Privacy . . . . . . . . . . . . . . . . . . . . . . . . . . . . 49

4.3.3 Concealing the Iterates . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 51

4.3.4 Repeated Analyses over Evolving Data . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 52

4.3.5 Preventing Model Theft and Misuse . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 52

4.4 Protections Against an Adversarial Server . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 53

4.4.1 Challenges: Communication Channels, Sybil Attacks, and Selection . . . . . . . . . . . . . . 53

4.4.2 Limitations of Existing Solutions . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 54

4.4.3 Training with Distributed Differential Privacy . . . . . . . . . . . . . . . . . . . . . . . . . . 55

4.4.4 Preserving Privacy While Training Sub-Models . . . . . . . . . . . . . . . . . . . . . . . . . 58

4.5 User Perception . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 59

4.5.1 Understanding Privacy Needs for Particular Analysis Tasks . . . . . . . . . . . . . . . . . . . 59

4.5.2 Behavioral Research to Elicit Privacy Preferences . . . . . . . . . . . . . . . . . . . . . . . . 60

4.6 Executive Summary . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 60

5 Defending Against Attacks and Failures 62

5.1 Adversarial Attacks on Model Performance . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 62

5.1.1 Goals and Capabilities of an Adversary . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 63

5.1.2 Model Update Poisoning . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 66

5.1.3 Data Poisoning Attacks . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 67

5.1.4 Inference-Time Evasion Attacks . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 69

5.1.5 Defensive Capabilities from Privacy Guarantees . . . . . . . . . . . . . . . . . . . . . . . . . 70

5.2 Non-Malicious Failure Modes . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 71

5.3 Exploring the Tension between Privacy and Robustness . . . . . . . . . . . . . . . . . . . . . . . . . 73

5.4 Executive Summary . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 73

6 Ensuring Fairness and Addressing Sources of Bias 75

6.1 Bias in Training Data . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 75

6.2 Fairness Without Access to Sensitive Attributes . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 76

6.3 Fairness, Privacy, and Robustness . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 77

6.4 Leveraging Federation to Improve Model Diversity . . . . . . . . . . . . . . . . . . . . . . . . . . . 78

6.5 Federated Fairness: New Opportunities and Challenges . . . . . . . . . . . . . . . . . . . . . . . . . 79

6.6 Executive Summary . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 79

7 Addressing System Challenges 81

7.1 Platform Development and Deployment Challenges . . . . . . . . . . . . . . . . . . . . . . . . . . . 81

7.2 System Induced Bias . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 82

7.2.1 Device Availability Profiles . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 82

7.2.2 Examples of System Induced Bias . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 83

7.2.3 Open Challenges in Quantifying and Mitigating System Induced Bias . . . . . . . . . . . . . 84

7.3 System Parameter Tuning . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 85

7.4 On-Device Runtime . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 86

7.5 The Cross-Silo Setting . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 87

7.6 Executive Summary . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 88

8 Concluding Remarks 89

A Software and Datasets for Federated Learning 119

1 Introduction

Federated learning (FL) is a machine learning setting where many clients (e.g. mobile devices or whole or- ganizations) collaboratively train a model under the orchestration of a central server (e.g. service provider), while keeping the training data decentralized. It embodies the principles of focused collection and data minimization, and can mitigate many of the systemic privacy risks and costs resulting from traditional, cen- tralized machine learning. This area has received significant interest recently, both from research and applied

perspectives. This paper describes the defining characteristics and challenges of the federated learning set- ting, highlights important practical constraints and considerations, and then enumerates a range of valuable research directions. The goals of this work are to highlight research problems that are of significant theo- retical and practical interest, and to encourage research on problems that could have significant real-world impact. Federated Learning, since the learning task is solved by a loose federation of participating devices (which

we refer to as clients) which are coordinated by a central server.” An unbalanced and non-IID (identically and independently distributed) data partitioning across a massive number of unreliable devices with limited communication bandwidth was introduced as the defining set of challenges.

Significant related work predates the introduction of the term federated learning. A longstanding goal

pursued by many research communities (including cryptography, databases, and machine learning) is to ana- lyze and learn from data distributed among many owners without exposing that data. Cryptographic methods for computing on encrypted data were developed starting in the early 1980s [396, 492], and Agrawal and a centralized server while preserving privacy. Conversely, even since the introduction of the term federated learning, we are aware of no single work that directly addresses the full set of FL challenges. Thus, the term

federated learning provides a convenient shorthand for a set of characteristics, constraints, and challenges that often co-occur in applied ML problems on decentralized data where privacy is paramount. This paper originated at the Workshop on Federated Learning and Analytics held June 17–18th, 2019, hosted at Google’s Seattle office. During the course of this two-day event, the need for a broad paper surveying the many open challenges in the area of federated learning became clear.1

A key property of many of the problems discussed is that they are inherently interdisciplinary — solving them likely requires not just machine learning, but techniques from distributed optimization, cryptography, security, differential privacy, fairness, compressed sensing, systems, information theory, statistics, and more. Many of the hardest problems are at the intersections of these areas, and so we believe collaboration will be essential to ongoing progress. One of the goals of this work is to highlight the ways in which techniques from

these fields can potentially be combined, raising both interesting possibilities as well as new challenges. Since the term federated learning was initially introduced with an emphasis on mobile and edge device applications [337, 334], interest in applying FL to other applications has greatly increased, including some which might involve only a small number of relatively reliable clients, for example multiple organizations collaborating to train a model. We term these two federated learning settings “cross-device” and “cross-silo”

respectively. Given these variations, we propose a somewhat broader definition of federated learning:

Federated learning is a machine learning setting where multiple entities (clients) collaborate

in solving a machine learning problem, under the coordination of a central server or service provider. Each client’s raw data is stored locally and not exchanged or transferred; instead,

focused updates intended for immediate aggregation are used to achieve the learning objective.

Focused updates are updates narrowly scoped to contain the minimum information necessary for the specific learning task at hand; aggregation is performed as early as possible in the service of data minimization. We note that this definition distinguishes federated learning from fully decentralized (peer-to-peer) learning techniques as discussed in Section 2.1. Although privacy-preserving data analysis has been studied for more than 50 years, only in the past decade have solutions been widely deployed at scale (e.g. [177, 154]). Cross-device federated learning and

federated data analysis are now being applied in consumer digital products. Google makes extensive use of federated learning in the Gboard mobile keyboard [376, 222, 491, 112, 383], as well as in features on Pixel phones and in Android Messages . While Google has pioneered cross-device FL, interest in this setting is now much broader, for example: Apple is using cross-device FL in iOS 13 , for applications like the QuickType keyboard and the vocal classifier for “Hey Siri” ; doc.ai is developing cross-device

FL solutions for medical research , and Snips has explored cross-device FL for hotword detection .

Cross-silo applications have also been proposed or described in myriad domains including finance risk

prediction for reinsurance , pharmaceuticals discovery , electronic health records mining , medical data segmentation [15, 139], and smart manufacturing . The growing demand for federated learning technology has resulted in a number of tools and frameworks becoming available. These include TensorFlow Federated , Federated AI Technology Enabler , PySyft , Leaf , PaddleFL and Clara Training Framework ; more details in Appendix A.

Commercial data platforms incorporating federated learning are in development from established technology

companies as well as smaller start-ups. distributed learning across a range of axes. These characteristics establish many of the constraints that practical federated learning systems must typically satisfy, and hence serve to both motivate and inform the open challenges in federated learning. They will be discussed at length in the sections that follow. These two FL variants are called out as representative and important examples, but different FL settings may have different combinations of these characteristics. For the remainder of this paper, we consider the

cross-device FL setting unless otherwise noted, though many of the problems apply to other FL settings as well. Section 2 specifically addresses some of the many other variations and applications.

Next, we consider cross-device federated learning in more detail, focusing on practical aspects common

particular production system, including a discussion of specific architectural choices and considerations.

1.1 The Cross-Device Federated Learning Setting

This section takes an applied perspective, and unlike the previous section, does not attempt to be definitional. Rather, the goal is to describe some of the practical issues in cross-device FL and how they might fit into a broader machine learning development and deployment ecosystem. The hope is to provide useful context and motivation for the open problems that follow, as well as to aid researchers in estimating how straight- forward it would be to deploy a particular new approach in a real-world system. We begin by sketching the

lifecycle of a model before considering a FL training process.

Datacenter Cross-silo Cross-device

distributed learning federated learning federated learning

Setting Training a model on a large Training a model on siloed data. The clients are a very large number of but “flat” dataset. Clients Clients are different organiza- mobile or IoT devices. are compute nodes in a sin- tions (e.g. medical or financial) gle cluster or datacenter. or geo-distributed datacenters. Data Data is centrally stored and Data is generated locally and remains decentralized. Each client stores distribution can be shuffled and balanced its own data and cannot read the data of other clients. Data is not indepen-

across clients. Any client can dently or identically distributed. read any part of the dataset.

Orchestration Centrally orchestrated. A central orchestration server/service organizes the training, but never

Wide-area None (fully connected Typically a hub-and-spoke topology, with the hub representing a coordi-

communication clients in one datacen- nating service provider (typically without data) and the spokes connecting ter/cluster). to clients. Data All clients are almost always available. Only a fraction of clients are available at availability any one time, often with diurnal or other variations. Distribution Typically 1 - 1000 clients. Typically 2 - 100 clients. Massively parallel, up to 1010 clients.

scale Primary Computation is more often Might be computation or com- Communication is often the primary bottleneck the bottleneck in the datacen- munication. bottleneck, though it depends on the ter, where very fast networks task. Generally, cross-device federated can be assumed. computations use wi-fi or slower con- nections. Addressability Each client has an identity or name that allows the system to Clients cannot be indexed directly (i.e.,

access it specifically. no use of client identifiers). Client Stateful — each client may participate in each round of the com- Stateless — each client will likely par- statefulness putation, carrying state from round to round. ticipate only once in a task, so gener- ally a fresh sample of never-before-seen clients in each round of computation is assumed.

Client Relatively few failures. Highly unreliable — 5% or more of the

reliability clients participating in a round of com- putation are expected to fail or drop out (e.g. because the device becomes ineli- gible when battery, network, or idleness requirements are violated). Data partition Data can be partitioned / re- Partition is fixed. Could be Fixed partitioning by example (horizon- axis partitioned arbitrarily across example-partitioned (horizontal) tal).

clients. or feature-partitioned (vertical).

Cross-device and cross-silo federated learning are two examples of FL domains, but are not intended to be exhaustive. The primary defining characteristics of FL are highlighted in bold, but the other characteristics are also critical in determining which techniques are applicable.

engineers & analysts federated rest of the world learning model deployment

1.1.1 The Lifecycle of a Model in Federated Learning

The FL process is typically driven by a model engineer developing a model for a particular application. For example, a domain expert in natural language processing may develop a next word prediction model for use in a virtual keyboard. Figure 1 shows the primary components and actors. At a high level, a typical workflow is:

1. Problem identification: The model engineer identifies a problem to be solved with FL.

2. Client instrumentation: If needed, the clients (e.g. an app running on mobile phones) are instru-

mented to store locally (with limits on time and quantity) the necessary training data. In many cases, the app already will have stored this data (e.g. a text messaging app must store text messages, a photo management app already stores photos). However, in some cases additional data or metadata might need to be maintained, e.g. user interaction data to provide labels for a supervised learning task.

3. Simulation prototyping (optional): The model engineer may prototype model architectures and test

learning hyperparameters in an FL simulation using a proxy dataset.

4. Federated model training: Multiple federated training tasks are started to train different variations

of the model, or use different optimization hyperparameters.

5. (Federated) model evaluation: After the tasks have trained sufficiently (typically a few days, see below), the models are analyzed and good candidates selected. Analysis may include metrics com- puted on standard datasets in the datacenter, or federated evaluation wherein the models are pushed to held-out clients for evaluation on local client data.

6. Deployment: Finally, once a good model is selected, it goes through a standard model launch process,

including manual quality assurance, live A/B testing (usually by using the new model on some devices and the previous generation model on other devices to compare their in-vivo performance), and a

Total population size 106 –1010 devices

Devices selected for one round of training 50 – 5000

Total devices that participate in training one model 105 –107

Number of rounds for model convergence 500 – 10000

Wall-clock training time 1 – 10 days

staged rollout (so that poor behavior can be discovered and rolled back before affecting too many users). The specific launch process for a model is set by the owner of the application and is usually independent of how the model is trained. In other words, this step would apply equally to a model trained with federated learning or with a traditional datacenter approach.

One of the primary practical challenges an FL system faces is making the above workflow as straight- forward as possible, ideally approaching the ease-of-use achieved by ML systems for centralized training. While much of this paper concerns federated training specifically, there are many other components in- cluding federated analytics tasks like model evaluation and debugging. Improving these is the focus of Section 3.4. For now, we consider in more detail the training of a single FL model (Step 4 above).

1.1.2 A Typical Federated Training Process

We now consider a template for FL training that encompasses the Federated Averaging algorithm of McMa- A server (service provider) orchestrates the training process, by repeating the following steps until train- ing is stopped (at the discretion of the model engineer who is monitoring the training process):

1. Client selection: The server samples from a set of clients meeting eligibility requirements. For

example, mobile phones might only check in to the server if they are plugged in, on an unmetered wi-fi connection, and idle, in order to avoid impacting the user of the device.

2. Broadcast: The selected clients download the current model weights and a training program (e.g. a

TensorFlow graph ) from the server.

3. Client computation: Each selected device locally computes an update to the model by executing the

training program, which might for example run SGD on the local data (as in Federated Averaging).

4. Aggregation: The server collects an aggregate of the device updates. For efficiency, stragglers might

be dropped at this point once a sufficient number of devices have reported results. This stage is also the integration point for many other techniques which will be discussed later, possibly including: secure aggregation for added privacy, lossy compression of aggregates for communication efficiency, and noise addition and update clipping for differential privacy.

5. Model update: The server locally updates the shared model based on the aggregated update computed

from the clients that participated in the current round.

The separation of the client computation, aggregation, and model update phases is not a strict require- ment of federated learning, and it indeed excludes certain classes of algorithms, for example asynchronous SGD where each client’s update is immediately applied to the model, before any aggregation with updates from other clients. Such asynchronous approaches may simplify some aspects of system design, and also be beneficial from an optimization perspective (though this point can be debated). However, the approach

presented above has a substantial advantage in affording a separation of concerns between different lines of research: advances in compression, differential privacy, and secure multi-party computation can be devel- oped for standard primitives like computing sums or means over decentralized updates, and then composed with arbitrary optimization or analytics algorithms, so long as those algorithms are expressed in terms of aggregation primitives. It is also worth emphasizing that in two respects, the FL training process should not impact the user

experience. First, as outlined above, even though model parameters are typically sent to some devices during the broadcast phase of each round of federated training, these models are an ephemeral part of the training process, and not used to make “live” predictions shown to the user. This is crucial, because training ML models is challenging, and a misconfiguration of hyperparameters can produce a model that makes bad predictions. Instead, user-visible use of the model is deferred to a rollout process as detailed above in Step 6

of the model lifecycle. Second, the training itself is intended to be invisible to the user — as described under client selection, training does not slow the device or drain the battery because it only executes when the device is idle and connected to power. However, the limited availability these constraints introduce leads directly to open research challenges which will be discussed subsequently, such as semi-cyclic data availability and the potential for bias in client selection.

1.2 Federated Learning Research

The remainder of this paper surveys many open problems that are motivated by the constraints and chal- lenges of real-world federated learning settings, from training models on medical data from a hospital sys- tem to training using hundreds of millions of mobile devices. Needless to say, most researchers working on federated learning problems will likely not be deploying production FL systems, nor have access to fleets of millions of real-world devices. This leads to a key distinction between the practical settings that motivate the

work and experiments conducted in simulation which provide evidence of the suitability of a given approach to the motivating problem. This makes FL research somewhat different than other ML fields from an experimental perspective, lead- ing to additional considerations in conducting FL research. In particular, when highlighting open problems, we have attempted, when possible, to also indicate relevant performance metrics which can be measured in simulation, the characteristics of datasets which will make them more representative of real-world per-

formance, etc. The need for simulation also has ramifications for the presentation of FL research. While not intended to be authoritative or absolute, we make the following modest suggestions for presenting FL research that addresses the open problems we describe:

• As shown in Table 1, the FL setting can encompass a wide range of problems. Compared to fields where the setting and goals are well-established, it is important to precisely describe the details of the particular FL setting of interest, particularly when the proposed approach makes assumptions that may not be appropriate in all settings (e.g. stateful clients that participate in all rounds).

• Of course, details of any simulations should be presented in order to make the research reproducible. But it is also important to explain which aspects of the real-world setting the simulation is designed to capture (and which it is not), in order to effectively make the case that success on the simulated

problem implies useful progress on the real-world objective. We hope that the guidance in this paper will help with this.

• Privacy and communication efficiency are always first-order concerns in FL, even if the experiments are simulations running on a single machine using public data. More so than with other types of ML, for any proposed approach it is important to be unambiguous about where computation happens as well as what is communicated.

Software libraries for federated learning simulation as well as standard datasets can help ease the chal- lenges of conducting effective FL research; Appendix A summarizes some of the currently available options.

Developing standard evaluation metrics and establishing standard benchmark datasets for different federated

learning settings (cross-device and cross-silo) remain highly important directions for ongoing work.

1.3 Organization

Section 2 builds on the ideas in Table 1, exploring other FL settings and problems beyond the original focus on cross-device settings. Section 3 then turns to core questions around improving the efficiency and effectiveness of federated learning. Section 4 undertakes a careful consideration of threat models and considers a range of technologies toward the goal of achieving rigorous privacy protections. As with all machine learning systems, in federated learning applications there may be incentives to manipulate the

models being trained, and failures of various kinds are inevitable; these challenges are discussed in Section 5. Finally, we address the important challenges of providing fair and unbiased models in Section 6.

2 Relaxing the Core FL Assumptions: Applications to Emerging Settings

In this section, we will discuss areas of research related to the topics discussed in the previous section. Even though not being the main focus of the remainder of the paper, progress in these areas could motivate design of the next generation of production systems.

2.1 Fully Decentralized / Peer-to-Peer Distributed Learning

In federated learning, a central server orchestrates the training process and receives the contributions of all clients. The server is thus a central player which also potentially represents a single point of failure. While large companies or organizations can play this role in some application scenarios, a reliable and powerful central server may not always be available or desirable in more collaborative learning scenarios . Furthermore, the server may even become a bottleneck when the number of clients is very large, as

The key idea of fully decentralized learning is to replace communication with the server by peer-to- peer communication between individual clients. The communication topology is represented as a connected graph in which nodes are the clients and an edge indicates a communication channel between two clients. The network graph is typically chosen to be sparse with small maximum degree so that each node only needs to send/receive messages to/from a small number of peers; this is in contrast to the star graph of the

server-client architecture. In fully decentralized algorithms, a round corresponds to each client performing a local update and exchanging information with their neighbors in the graph2 . In the context of machine learning, the local update is typically a local (stochastic) gradient step and the communication consists in averaging one’s local model parameters with the neighbors. Note that there is no longer a global state of the model as in standard federated learning, but the process can be designed such that all local models converge

to the desired global solution, i.e., the individual models gradually reach consensus. While multi-agent optimization has a long history in the control community, fully decentralized variants of SGD and other optimization algorithms have recently been considered in machine learning both for improved scalability in datacenters as well as for decentralized networks of devices [127, 459, 443, 59, 278, 291, 173].

They consider undirected network graphs, although the case of directed networks (encoding unidirectional

channels which may arise in real-world scenarios such as social networks or data markets) has also been studied in [29, 226]. It is worth noting that even in the decentralized setting outlined above, a central authority may still be in charge of setting up the learning task. Consider for instance the following questions: Who decides what is the model to be trained in the decentralized setting? What algorithm to use? What hyperparameters? Who is responsible for debugging when something does not work as expected? A certain degree of trust of the

participating clients in a central authority would still be needed to answer these questions. Alternatively, the decisions could be taken by the client who proposes the learning task, or collaboratively through a consensus scheme (see Section 2.1.2). assumptions of decentralized learning are distinct from those of federated learning, it can often be applied to similar problem domains, many of the same challenges arise, and there is significant overlap in the research communities. Thus, we consider decentralized learning in this paper as well; in this section challenges

Note, however, that the notion of a round does not need to even make sense in this setting. See for instance the discussion on clock models in .

Federated learning Fully decentralized

Orchestration A central orchestration server or ser- No centralized orchestration. vice organizes the training, but never sees raw data.

Wide-area communication Typically a hub-and-spoke topology, Peer-to-peer topology, with a possi-

with the hub representing a coor- bly dynamic connectivity graph. dinating service provider (typically without data) and the spokes con- necting to clients.

Note that as with FL, decentralized learning can be further divided into different use-cases, with distinctions similar to those made in Table 1 comparing cross-silo and cross-device FL.

specific to the decentralized approach are explicitly considered, but many of the open problems in other sections also arise in the decentralized case.

2.1.1 Algorithmic Challenges

A large number of important algorithmic questions remain open on the topic of real-world usability of de- centralized schemes for machine learning. Some questions are analogous to the special case of federated learning with a central server, and other challenges come as an additional side-effect of being fully decen- tralized or trust-less. We outline some particular areas in the following.

Effect of network topology and asynchrony on decentralized SGD Fully decentralized algorithms for

learning should be robust to the limited availability of the clients (with clients temporarily unavailable, dropping out or joining during the execution) and limited reliability of the network (with possible message drops). While for the special case of generalized linear models, schemes using the duality structure could enable some of these desired robustness properties , for the case of deep learning and SGD this remains an open question. When the network graph is complete but messages have a fixed probability to be dropped,

work. Additional open research questions concern non-IID data distributions, update frequencies, efficient communication patterns and practical convergence time , as we outline in more detail below.

Well-connected or denser networks encourage faster consensus and give better theoretical convergence

rates, which depend on the spectral gap of the network graph. However, when data is IID, sparser topologies do not necessarily hurt the convergence in practice: this was analyzed theoretically in . Denser net- works typically incur communication delays which increase with the node degrees. Most of optimization- theory works do not explicitly consider how the topology affects the runtime, that is, wall-clock time re- based on matching decomposition sampling, that reduces the communication delay per iteration for any

given node topology while maintaining the same error convergence speed. The key idea is to decompose the graph topology into matchings consisting of disjoint communication links that can operate in parallel, and carefully choose a subset of these matchings in each iteration. This sequence of subgraphs results in more

frequent communication over connectivity-critical links (ensuring fast error convergence) and less frequent communication over other links (saving communication delays). The setting of decentralized SGD also naturally lends itself to asynchronous algorithms in which each client becomes active independently at random times, removing the need for global synchronization and potentially improving scalability [127, 459, 59, 29, 306].

Local-update decentralized SGD The theoretical analysis of schemes which perform several local update

steps before a communication round is significantly more challenging than those using a single SGD step, as in mini-batch SGD. While this will also be discussed later in Section 3.2, the same also holds more generally in the fully decentralized setting of interest here. Schemes relying on a single local update step are typically proven to converge in the case of non-IID local datasets [278, 279]. For the case with several local update steps, [467, 280] recently provided convergence analysis. Further, provides a convergence analysis

for the non-IID data case, but for the specific scheme based on matching decomposition sampling described above. In general, however, understanding the convergence under non-IID data distributions and how to design a model averaging policy that achieves the fastest convergence remains an open problem.

Personalization, and trust mechanisms Similarly to the cross-device FL setting, an important task for

the fully decentralized scenario under the non-IID data distributions available to individual clients is to design algorithms for learning collections of personalized models. The work of [459, 59] introduces fully decentralized algorithms to collaboratively learn a personalized model for each client by smoothing model further learn the similarity graph together with the personalized models. One of the key unique challenges in the decentralized setting remains the robustness of such schemes to malicious actors or contribution of

unreliable data or labels. The use of incentives or mechanism design in combination with decentralized learning is an emerging and important goal, which may be harder to achieve in the setting without a trusted central server.

Gradient compression and quantization methods In potential applications, the clients would often

be limited in terms of communication bandwidth available and energy usage permitted. Translating and generalizing some of the existing compressed communication schemes from the centralized orchestrator- facilitated setting (see Section 3.5) to the fully decentralized setting, without negatively impacting the con- vergence is an active research direction [278, 391, 444, 279]. A complementary idea is to design decentral- ized optimization algorithms which naturally give rise to sparse updates .

Privacy An important challenge in fully decentralized learning is to prevent any client from reconstructing the private data of another client from its shared updates while maintaining a good level of utility for the learned models. Differential privacy (see Section 4) is the standard approach to mitigate such privacy risks. In decentralized federated learning, this can be achieved by having each client add noise locally, as done in [239, 59]. Unfortunately, such local privacy approaches often come at a large cost in utility. Furthermore,

distributed methods based on secure aggregation or secure shuffling that are designed to improve the privacy- utility trade-off in the standard FL setting (see Section 4.4.3) do not easily integrate with fully decentralized algorithms. A possible direction to achieve better trade-offs between privacy and utility in fully decentralized algorithms is to rely on decentralization itself to amplify differential privacy guarantees, for instance by considering appropriate relaxations of local differential privacy .

2.1.2 Practical Challenges

An orthogonal question for fully decentralized learning is how it can be practically realized. This section outlines a family of related ideas based on the idea of a distributed ledger, but other approaches remain unexplored. A blockchain is a distributed ledger shared among disparate users, making possible digital transactions, including transactions of cryptocurrency, without a central authority. In particular, smart contracts allow execution of arbitrary code on top of the blockchain, essentially a massively replicated eventually-consistent

state machine. In terms of federated learning, use of the technology could enable decentralization of the global server by using smart contracts to do model aggregation, where the participating clients executing the smart contracts could be different companies or cloud services. However, on today’s blockchain platforms such as Ethereum , data on the blockchains is publicly available by default, this could discourage users from participating in the decentralized federated learning

protocol, as the protection of the data is typically the primary motivating factor for FL. To address such concerns, it might be possible to modify the existing privacy-preserving techniques to fit into the scenario of decentralized federated learning. First of all, to prevent the participating nodes from exploiting individually submitted model updates, existing secure aggregation protocols could be used. A practical secure aggre- dropping out participants at the cost of complexity of the protocol. An alternative system would be to have

each client stake a deposit of cryptocurrency on blockchain, and get penalized if they drop out during the execution. Without the need of handling dropouts, the secure aggregation protocol could be significantly simplified. Another way of achieving secure aggregation is to use confidential smart contract such as what is enabled by the Oasis Protocol which runs inside secure enclaves. With this, each client could simply submit an encrypted local model update, knowing that the model will be decrypted and aggregated inside

the secure hardware through remote attestation (though see discussion of privacy-in-depth in Section 4.1). In order to prevent any client from trying to reconstruct the private data of another client by exploiting the global model, client-level differential privacy has been proposed for FL. Client-level differential privacy is achieved by adding random Gaussian noise on the aggregated global model that is enough to hide any single client’s update. In the context of blockchain, each client could locally add a certain amount of

Gaussian noise after local gradient descent steps and submit the model to blockchain. The local noise scale should be calculated such that the aggregated noise on blockchain is able to achieve the same client-level differential privacy as in . Finally, the aggregated global model on blockchain could be encrypted and only the participating clients hold the decryption key, which protects the model from the public.

2.2 Cross-Silo Federated Learning

In contrast with the characteristics of cross-device federated learning, see Table 1, cross-silo federated learn- ing admits more flexibility in certain aspects of the overall design, but at the same time presents a setting where achieving other properties can be harder. This section discusses some of these differences. The cross-silo setting can be relevant where a number of companies or organizations share incentive to train a model based on all of their data, but cannot share their data directly. This could be due to constraints

imposed by confidentiality or due to legal constraints, or even within a single company when they cannot centralize their data between different geographical regions. These cross-silo applications have attracted substantial attention.

Data partitioning In the cross-device setting the data is assumed to be partitioned by examples. In the cross-silo setting, in addition to partitioning by examples, partitioning by features is of practical relevance. An example could be when two companies in different businesses have the same or overlapping set of customers, such as a local bank and a local retail company in the same city. This difference has been also

Cross-silo FL with data partitioned by features, employs a very different training architecture compared

to the setting with data partitioned by example. It may or may not involve a central server as a neutral party, and based on specifics of the training algorithm, clients exchange specific intermediate results rather than model parameters, to assist other parties’ gradient calculations; see for instance [490, Section 2.4.2]. In this setting, application of techniques such as secure multi-party computation or homomorphic encryption have been proposed in order to limit the amount of information other participants can infer from observing the

training process. The downside of this approach is that the training algorithm is typically dependent on the type of machine learning objective being pursued. Currently proposed algorithms include trees , linear and logistic regression [490, 224, 316], and neural networks . Local updates similar to Federated Av- eraging (see Section 3.2) has been proposed to address the communication challenges of feature-partitioned systems , and [238, 318] study the security and privacy related challenges inherent in such systems.

Federated transfer learning is another concept that considers challenging scenarios in which data

parties share only a partial overlap in the user space or the feature space, and leverage existing transfer learning techniques to build models collaboratively. The existing formulation is limited to the case of

2 clients.

Partitioning by examples is usually relevant in cross-silo FL when a single company cannot centralize their data due to legal constraints, or when organizations with similar objectives want to collaboratively im- prove their models. For instance, different banks can collaboratively train classification or anomaly detection models for fraud detection , hospitals can build better diagnostic models , and so on.

An open-source platform supporting the above outlined applications is currently available as Federated

AI Technology Enabler (FATE) . At the same time, the IEEE P3652.1 Federated Machine Learning

Working Group is focusing on standard-setting for the Federated AI Technology Framework. Other plat-

forms include focused on a range of medical applications and for enterprise use cases. See

Appendix A for more details.

Incentive mechanisms In addition to developing new algorithmic techniques for FL, incentive mechanism

design for honest participation is an important practical research question. This need may arise in cross- device settings (e.g. [261, 260]), but is particularly relevant in the cross-silo setting, where participants may at the same time also be business competitors. The incentive can be in the form of monetary payout or final models with different levels of performance . The option to deliver models with performance commensurate to the contributions of each client is especially relevant in collaborative learning situations

in which competitions exist among FL participants. Clients might worry that contributing their data to training federated learning models will benefit their competitors, who do not contribute as much but receive the same final model nonetheless (i.e. the free-rider problem). Related objectives include how to divide earnings generated by the federated learning model among contributing data owners in order to sustain long-term participation, and also how to link the incentives with decisions on defending against adversarial

data owners to enhance system security, optimizing the participation of data owners to enhance system efficiency.

Differential privacy The discussion of actors and threat models in Section 4.1 is largely relevant also for the cross-silo FL. However, protecting against different actors might have different priorities. For example, in many practical scenarios, the final trained model would be released only to those who participate in the training, which makes the concerns about “the rest of the world” less important. On the other hand, for a practically persuasive claim, we would usually need a notion of local differential

privacy, as the potential threat from other clients is likely to be more important. In cases when the clients are not considered a significant threat, each client could control the data from a number of their respective users, and a formal privacy guarantee might be needed on such user-level basis. Depending on application, other objectives could be worth pursuing. This area has not been systematically explored.

Tensor factorization Several works have also studied cross-silo federated tensor factorization where mul-

tiple sites (each having a set of data with the same feature, i.e. horizontally partitioned) jointly perform tensor factorization by only sharing intermediate factors with the coordination server while keeping data private at each site. Among the existing works, used an alternating direction method of multipli- ers (ADMM) based approach and improved the efficiency with the elastic averaging SGD (EASGD) algorithm and further ensures differential privacy for the intermediate factors.

2.3 Split Learning

In contrast with the previous settings which focus on data partitioning and communication patterns, the key idea behind split learning [215, 460]3 is to split the execution of a model on a per-layer basis between the clients and the server. This can be done for both training and inference. In the simplest configuration of split learning, each client computes the forward pass through a deep network up to a specific layer referred to as the cut layer. The outputs at the cut layer, referred to as

smashed data, are sent to another entity (either the server or another client), which completes the rest of the computation. This completes a round of forward propagation without sharing the raw data. The gradients can then be back propagated from its last layer until the cut layer in a similar fashion. The gradients at the cut layer – and only these gradients – are sent back to the clients, where the rest of back propagation is completed. This process is continued until convergence, without having clients directly access each others

raw data. This setup is shown in Figure 2(a) and a variant of this setup where labels are also not shared along with raw data is shown in Figure 2(b). Split learning approaches for data partitioned by features have been studied in .

In several settings, the overall communication requirements of split learning and federated learning

were compared in . Split learning brings in another dimension of parallelism in the training, paral- lelization among parts of a model, e.g. client and server. The ideas in [245, 240], where the authors break the dependencies between partial networks and reduced total centralized training time by parallelizing the computations in different parts, can be relevant here as well. However, it is still an open question to explore such parallelization of split learning on edge devices. Split learning also enables matching client-side model

components with the best server-side model components for automating model selection as shown in the ExpertMatcher . The values communicated can nevertheless, in general, reveal information about the underlying data. How much, and whether this is acceptable, is likely going to be application and configuration specific. A variation of split learning called NoPeek SplitNN reduces the potential leakage via communicated ac- tivations, by reducing their distance correlation [461, 442] with the raw data, while maintaining good model

See also split learning project website - https://splitlearning.github.io/.

(a) Vanilla split learning (b) U-shaped split learning

data as well as labels are not transferred between the client and server entities in the U-shaped split learning setting.

performance via categorical cross-entropy. The key idea is to minimize the distance correlation between the raw data points and communicated smashed data. The objects communicated could otherwise contain infor- mation highly correlated with the input data if used without NoPeek SplitNN, the use of which also enables the split to be made relatively early-on given the decorrelation it provides. One other engineering driven approach to minimize the amount of information communicated in split learning has been via a specifically

learnt pruning of channels present in the client side activations . Overall, much of the discussion in Section 4 is relevant here as well, and analysis providing formal privacy guarantees specifically for split learning is still an open problem.

2.4 Executive summary

The motivation for federated learning is relevant for a number of related areas of research.

• Fully decentralized learning (Section 2.1) removes the need for a central server coordinating the over- all computation. Apart from algorithmic challenges, open problems are in practical realization of the idea and in understanding of what form of trusted central authority is needed to set up the task.

• Cross-silo federated learning (Section 2.2) admits problems with different kinds of modelling con- straints, such as data partitioned by examples and/or features, and faces different set of concerns when formulating formal privacy guarantees or incentive mechanisms for clients to participate.

• Split learning (Section 2.3) is an approach to partition the execution of a model between the clients and the server. It can deliver different options for overall communication constraints, but detailed analysis of when the communicated values reveal sensitive information is still missing.

3 Improving Efficiency and Effectiveness

In this section we explore a variety of techniques and open questions that address the challenge of making federated learning more efficient and effective. This encompasses a myriad of possible approaches, includ- ing: developing better optimization algorithms; providing different models to different clients; making ML tasks like hyperparameter search, architecture search, and debugging easier in the FL context; improving communication efficiency; and more. One of the fundamental challenges in addressing these goals is the presence of non-IID data, so we begin

by surveying this issue and highlighting potential mitigations.

3.1 Non-IID Data in Federated Learning

While the meaning of IID is generally clear, data can be non-IID in many ways. In this section, we provide a taxonomy of non-IID data regimes that may arise for any client-partitioned dataset. The most common sources of dependence and non-identicalness are due to each client corresponding to a particular user, a particular geographic location, and/or a particular time window. This taxonomy has a close mapping to notions of dataset shift [353, 380], which studies differences between the training distribution and testing

distribution; here, we consider differences in the data distribution on each client. For the following, consider a supervised task with features x and labels y. A statistical model of feder- ated learning involves two levels of sampling: accessing a datapoint requires first sampling a client i ∼ Q, the distribution over available clients, and then drawing an example (x, y) ∼ Pi (x, y) from that client’s local data distribution. When non-IID data in federated learning is referenced, this typically refers to differences between Pi

and Pj for different clients i and j. However, it is also important to note that the distribution Q and Pi may change over time, introducing another dimension of “non-IIDness”. For completeness, we note that even considering the dataset on a single device, if the data is in an insufficiently-random order, e.g. ordered by time, then independence is violated locally as well. For exam- ple, consecutive frames in a video are highly correlated. Sources of intra-client correlation can generally be

Non-identical client distributions We first survey some common ways in which data tend to deviate from being identically distributed, that is Pi 6= Pj for different clients i and j. Rewriting Pi (x, y) as Pi (y | x)Pi (x) and Pi (x | y)Pi (y) allows us to characterize the differences more precisely.

• Feature distribution skew (covariate shift): The marginal distributions Pi (x) may vary across clients, even if P(y | x) is shared.4 For example, in a handwriting recognition domain, users who write the same words might still have different stroke width, slant, etc.

• Label distribution skew (prior probability shift): The marginal distributions Pi (y) may vary across clients, even if P(x | y) is the same. For example, when clients are tied to particular geo-regions, the distribution of labels varies across clients — kangaroos are only in Australia or zoos; a person’s face is only in a few locations worldwide; for mobile device keyboards, certain emoji are used by one demographic but not others. We write “P(y | x) is shared” as shorthand for Pi (y | x) = Pj (y | x) for all clients i and j.

• Same label, different features (concept drift): The conditional distributions Pi (x | y) may vary across clients even if P(y) is shared. The same label y can have very different features x for different clients, e.g. due to cultural differences, weather effects, standards of living, etc. For example, images of homes can vary dramatically around the world and items of clothing vary widely. Even within the U.S., images of parked cars in the winter will be snow-covered only in certain parts of the country. The

same label can also look very different at different times, and at different time scales: day vs. night, seasonal effects, natural disasters, fashion and design trends, etc.

• Same features, different label (concept shift): The conditional distribution Pi (y | x) may vary across clients, even if P(x) is the same. Because of personal preferences, the same feature vectors in a training data item can have different labels. For example, labels that reflect sentiment or next word predictors have personal and regional variation.

• Quantity skew or unbalancedness: Different clients can hold vastly different amounts of data.

Real-world federated learning datasets likely contain a mixture of these effects, and the characterization

of cross-client differences in real-world partitioned datasets is an important open question. Most empirical work on synthetic non-IID datasets (e.g. [337, 236]) have focused on label distribution skew, where a non- IID dataset is formed by partitioning a “flat” existing dataset based on the labels. A better understanding of the nature of real-world non-IID datasets will allow for the construction of controlled but realistic non-IID datasets for testing algorithms and assessing their resilience to different degrees of client heterogeneity.

Further, different non-IID regimes may require the development of different mitigation strategies. For

example, under feature-distribution skew, because P(y | x) is assumed to be common, the problem is at least in principle well specified, and training a single global model that learns P(y | x) may be appropriate. When the same features map to different labels on different clients, some form of personalization (Section 3.3) may be essential to learning the true labeling functions.

Violations of independence Violations of independence are introduced any time the distribution Q changes

over the course of training; a prominent example is in cross-device FL, where devices typically need to meet eligibility requirements in order to participate in training (see Section 1.1.2). Devices typically meet those requirements at night local time (when they are more likely to be charging, on free wi-fi, and idle), and so there may be significant diurnal patterns in device availability. Further, because local time of day corre- described this issue and some mitigation strategies, but many open questions remain.

Dataset shift Finally, we note that the temporal dependence of the distributions Q and P may introduce dataset shift in the classic sense (differences between the train and test distributions). Furthermore, other criteria may make the set of clients eligible to train a federated model different from the set of clients where that model will be deployed. For example, training may require devices with more memory than is needed for inference. These issues are explored in more depth in Section 6. Adapting techniques for handling

dataset shift to federated learning is another interesting open question.

3.1.1 Strategies for Dealing with Non-IID Data

The original goal of federated learning, training a single global model on the union of client datasets, be- comes harder with non-IID data. One natural approach is to modify existing algorithms (e.g. through

different hyperparameter choices) or develop new ones in order to more effectively achieve this objective. This approach is considered in Section 3.2.2. For some applications, it may be possible to augment data in order to make the data across clients more similar. One approach is to create a small dataset which can be shared globally. This dataset may originate from a publicly available proxy data source, a separate dataset from the clients’ data which is not privacy The heterogeneity of client objective functions gives additional importance to the question of how to

craft the objective function — it is no-longer clear that treating all examples equally makes sense. Alterna- tives include limiting the contributions of the data from any one user (which is also important for privacy, see Section 4) and introducing other notions of fairness among the clients; see discussion in Section 6. But if we have the capability to run training on the local data on each device (which is necessary for federated learning of a global model), is training a single global model even the right goal? There are

many cases where having a single model is to be preferred, e.g. in order to provide a model to clients with no data, or to allow manual validation and quality assurance before deployment. Nevertheless, since local training is possible, it becomes feasible for each client to have a customized model. This approach can turn the non-IID problem from a bug to a feature, almost literally — since each client has its own model, the client’s identity effectively parameterizes the model, rendering some pathological but degenerate non-IID

distributions trivial. For example, if for each i, Pi (y) has support on only a single label, finding a high- accuracy global model may be very challenging (especially if x is relatively uninformative), but training a high-accuracy local model is trivial (only a constant prediction is needed). Such multi-model approaches are considered in depth in Section 3.3. In addition to addressing non-identical client distributions, using a plurality of models can also address violations of independence stemming from changes in client availability.

order to provide different models for inference based on the timezone / longitude of clients.

3.2 Optimization Algorithms for Federated Learning

In prototypical federated learning tasks, the goal is to learn a single global model that minimizes the em- pirical risk function over the entire training dataset, that is, the union of the data across all the clients. The main difference between federated optimization algorithms and standard distributed training methods is the need to address the characteristics of Table 1 — for optimization, non-IID and unbalanced data, limited communication bandwidth, and unreliable and limited device availability are particularly salient.

FL settings where the total number of devices is huge (e.g. across mobile devices) necessitate algorithms that only require a handful of clients to participate per round (client sampling). Further, each device is likely to participate no more than once in the training of a given model, so stateless algorithms are necessary. This rules out the direct application of a variety of approaches that are quite effective in the datacenter context, for example stateful optimization algorithms like ADMM, and stateful compression strategies that modify

updates based on residual compression errors from previous rounds.

Another important practical consideration for federated learning algorithms is composability with other

techniques. Optimization algorithms do not run in isolation in a production deployment, but need to be combined with other techniques like cryptographic secure aggregation protocols (Section 4.2.1), differential privacy (DP) (Section 4.2.2), and model and update compression (Section 3.5). As noted in Section 1.1.2, many of these techniques can be applied to primitives like “sum over selected clients” and “broadcast to selected clients”, and so expressing optimization algorithms in terms of these

primitives provides a valuable separation of concerns, but may also exclude certain techniques such as ap-

N Total number of clients Server executes:

M Clients per round initialize x0

T Total communication rounds for each round t = 1, 2, . . . , T do

K Local steps per round. St ← (random set of M clients)

for each client i ∈ St in parallel do algorithms including Federated Averaging. xt+1 ← M 1 i P k=1 M xt+1

ClientUpdate(i, x): for local step j = 1, . . . , K do x ← x − ηOf (x; z) for z ∼ Pi return x to server

Algorithm 1: Federated Averaging (local SGD), when all

clients have the same amount of data.

plying updates asynchronously. One of the most common approaches to optimization for federated learning is the Federated Averaging algorithm , an adaption of local-update or parallel SGD.5 Here, each client runs some number of SGD steps locally, and then the updated local models are averaged to form the updated global model on the coordinating server. Pseudocode is given in Algorithm 1.

Performing local updates and communicating less frequently with the central server addresses the core

challenges of respecting data locality constraints and of the limited communication capabilities of mobile device clients. However, this family of algorithms also poses several new algorithmic challenges from an optimization theory point of view. In Section 3.2, we discuss recent advances and open challenges in federated optimization algorithms for the cases of IID and non-IID data distribution across the clients respectively. The development of new algorithms that specifically target the characteristics of the federated

learning setting remains an important open problem.

3.2.1 Optimization Algorithms and Convergence Rates for IID Datasets

While a variety of different assumptions can be made on the per-client functions being optimized, the most basic split is between assuming IID and non-IID data. Formally, having IID data at the clients means that each mini-batch of data used for a client’s local update is statistically identical to a uniformly drawn sample (with replacement) from the entire training dataset (the union of all local datasets at the clients). Since the clients independently collect their own training data which vary in both size and distribution, and

these data are not shared with other clients or the central node, the IID assumption clearly almost never holds in practice. However, this assumption greatly simplifies theoretical convergence analysis of federated optimization algorithms, as well as establishes a baseline that can be used to understand the impact of non- IID data on optimization rates. Thus, a natural first step is to obtain an understanding of the landscape of optimization algorithms for the IID data case. Federated Averaging applies local SGD to a randomly sampled subset of clients on each round, and proposes a specific update

Formally, for the IID setting let us standardize the stochastic optimization problem

min F (x) := E [f (x; z)] . x∈Rm z∼P

stateless clients participate in each of T rounds, and during each round, each client can compute gradients for K samples (e.g. minibatches) z1 , . . . , zK sampled IID from P (possibly using these to take sequential steps). In the IID-data setting clients are interchangeable, and we can without loss of generality assume M = N . Table 4 summarizes the notation used in this section. Different assumptions on f will produce different guarantees. We will first discuss the convex setting

and later review results for non-convex problems.

Baselines and state-of-the-art for convex problems In this section we review convergence results for

H-smooth, convex (but not necessarily strongly convex) functions under the assumption that the variance

of the stochastic gradients is bounded by σ 2 . More formally, by H-smooth we mean that for all z, f (·; z) is differentiable and has a H-Lipschitz gradient, that is, for all choices of x, y

k∇f (x, z) − ∇f (y, z)k ≤ Hkx − yk.

We also assume that for all x, the stochastic gradient ∇x f (x; z) satisfies

E k∇x f (x; z) − ∇F (x)k ≤ σ 2 . z∼P

When analyzing the convergence rate of an algorithm with output xT after T iterations, we consider the term

E[F (xT )] − F (x∗ ) (1)

where x∗ = arg minx F (x). All convergence rates discussed herein are upper bounds on this term. A

summary of convergence results for such functions is given in Table 5.

Federated averaging (a.k.a. parallel SGD/local SGD) competes with two natural baselines: First, we

may keep x fixed in local updates during each round, and compute a total of KM gradients at the current x, in order to run accelerated minibatch SGD. Let x̄ denote the average of T iterations of this algorithm. We then have the upper bound  

H σ

O +√

T2 T KM

for convex objectives [294, 137, 151]. Note that the first expectation is taken with respect to the randomness of z in the training procedure as well. A second natural baseline is to ignore all but 1 of the M active clients, which allows (accelerated) sequential SGD to execute for KT steps. Applying the same general bounds cited above, this approach offers an upper bound of  

H σ

O +√ .

(T K)2 TK

√ Comparing these two results, we see that minibatch SGD attains the optimal ‘statistical’ term (σ/ T KM ), whilst SGD on a single device (ignoring the updates of the other devices) achieves the optimal ‘optimization’ term (H/(T K)2 ). The convergence analysis of local-update SGD methods is an active current area of research [434, 310, 500, 467, 390, 371, 269, 481]. The first convergence results for local-update SGD methods were derived

Method Comments Convergence

Baselines 

H mini-batch SGD batch size KM O T + √T σKM  H

SGD (on 1 worker, no communication) O TK + √Tσ K

Baselines with accelerationa 

A-mini-batch SGD [294, 137] batch size KM O 2 + √T σKM

T 

A-SGD (on 1 worker, no communication) O (T K)2

Parallel SGD / Fed-Avg / Local SGD  

HKM G2 √ σ

+ T KM

HM √ σ

Wang and Joshi b , Stich and Karimireddy O T + T KM

Other algorithms 

SCAFFOLD control variates and two stepsizes O T + √T σKM

There are no accelerated fed-avg/local SGD variants so far

b This paper considers the smooth non-convex setting, we adapt here the results for our setting. c This paper considers the smooth strongly convex setting, we adapt here the results for our setting.

data setting. We assume M devices participate in each iterations, and the loss functions are H-smooth, convex, and we have access to stochastic gradients with variance at most σ 2 . All rates are upper bounds on (1) after T iterations (potentially with some iterate averaging scheme).

under the bounded gradient norm assumption in Stich for strongly-convex non-convex objective functions. These analyses could attain the desired σ/ T KM statistical term with suboptimal optimization term (in Table 5 we summarize these results for the middle ground of convex functions).

By removing the bounded gradient assumption, Wang and Joshi and Stich and Karimireddy

could further improve the optimization term to HM/T . These result show that if the number of local steps K is smaller than T /M 3 then the (optimal) statistical term is dominating the rate. However, for typical cross-device applications we might have T = 106 and M = 100 (Table 2), implying K = 1. Often in the literature the convergence bounds are accompanied by a discussion on how large K may be chosen in order to reach asymptotically the same statistical term as the convergence rate of mini-batch

For non-convex

bound 1/ T KM if the number of local updates K are smaller than T 1/3 /M . This convergence guarantee was further improved by Wang and Joshi who removed the bounded gradient norm assumption and showed that the number of local updates can be as large as T /M 3 . The analysis in can also be applied to other algorithms with local updates, and thus yields the first convergence guarantee for decentralized

SGD with local updates (or periodic decentralized SGD) and elastic averaging SGD . Haddadpour

for PL functions, T 2 /M local updates per round leads to a O(1/T KM ) convergence. While the above works focus on convergence as a function of the number of iterations performed, prac- titioners often care about wall-clock convergence speed. Assessing this must take into account the effect of the design parameters on the time spent per iteration based on the relative cost of communication and local computation. Viewed in this light, the focus on seeing how large K can be while maintaining the

statistical rate may not be the primary concern in federated learning, where one may assume almost infinite datasets (very large N ). The costs (at least in wall-clock time) are small for increasing M , and so it may be more natural to increase M sufficiently to match the optimization term, and then tune K to maximize wall-clock optimization performance. How then to choose K? Performing more local updates at the clients will increase the divergence between the resulting local models at the clients, before they are averaged. As a

result, the error convergence in terms of training loss versus the total number of sequential SGD steps T K is slower. However, performing more local updates saves significant communication cost and reduces the time spent per iteration. The optimal number of local updates strikes a balance between these two phenomena and achieves the fastest error versus wallclock time convergence. Wang and Joshi propose an adaptive communication strategy that adapts K according to the training loss at regular intervals during the training.

Another important design parameter in federated learning is the model aggregation method used to

update the global model using the updates made by the selected clients. In the original federated learning size of local datasets. For IID data, where each client is assumed to have a infinitely large dataset, this reduces to taking a simple average of the local models. However, it is unclear whether this aggregation method will result in the fastest error convergence. highlights several gaps between upper and lower bounds for optimization relevant to the federated learning setting, particularly for “intermittent communication graphs”, which captures local SGD approaches, but

convergence rates for such approaches are not known to match the corresponding lower bounds. In Table 5 we highlight convergence results for the convex setting. Whilst most schemes are able to reach the asymp- totically dominant statistical term, none are able to match the convergence rate of accelerated mini-batch SGD. It is an open problem if federated averaging algorithms can close this gap. Local-update SGD methods where all M clients perform the same number of local updates may suffer

from a common scalability issue—they can be bottlenecked if any one client unpredictably slows down or fails. Several approaches for dealing with this are possible, but it is far from clear which are optimal, provisioning clients (e.g., request updates from 1.3M clients), and then accepting the first M updates re- ceived and rejecting updates from stragglers. A slightly more sophisticated solution is to fix a time window and allow clients to perform as many local updates Ki as possible within this time, after which their models

this approach in theory. An alternative method to overcome the problem of straggling clients is to fix the number of local updates at τ , but allow clients to update the global model in an asynchronous or lock-free fashion. Although some previous works [505, 306, 163] have proposed similar methods, the error conver- gence analysis is an open and challenging problem. A larger challenge in the FL setting, however, is that as discussed at the beginning of Section 3.2, asynchronous approaches may be difficult to combine with

complimentary techniques like differential privacy or secure aggregation. Besides the number of local updates, the choice of the size of the set of clients selected per training round presents a similar trade-off as the number of local updates. Updating and averaging a larger number of client models per training round yields better convergence, but it makes the training vulnerable to slowdown due

to unpredictable tail delays in computation/communication at/with the clients. The analysis of local SGD / Federated Averaging in the non-IID setting is even more challenging; results and open questions related to this are considered in the next section, along with specialized algorithms which directly address the non-IID problem.

3.2.2 Optimization Algorithms and Convergence Rates for Non-IID Datasets

In contrast to well-shuffled mini-batches consisting of independent and identically distributed (IID) ex-

amples in centralized learning, federated learning uses local data from end user devices, leading to many varieties of non-IID data (Section 3.1). In this setting, each of N clients has a local data distribution Pi and a local objective function

where we recall that f (x; z) is the loss of a model x at an example z. We typically wish to minimize N

1 X

F (x) = fi (x) . (2)

Note that we recover the IID setting when each Pi is identical. We will let F ∗ denote the minimum value of F , obtained the point x∗ . Analogously, we will let fi∗ denote the minimum value of fi . Sec. 4.4]), where M stateless clients participate in each of T rounds, and during each round, each client can compute gradients for K samples (e.g. minibatches). The difference here is that the samples zi,1 , . . . , zi,K sampled at client i are drawn from the client’s local distribution Pi . Unlike the IID setting, we cannot

necessarily assume M = N , as the client distributions are not all equal. In the following, if an algorithm relies on M = N , we will omit M and simply write N . We note that while such an assumption may be compatible with the cross-silo federated setting in Table 1, it is generally infeasible in the cross-device setting. While [434, 500, 467, 435] mainly focused on the IID case, the analysis technique can be extended to the non-IID case by adding an assumption on data dissimilarities, for example by constraining the difference

between client gradients and the global gradient [305, 300, 304, 469, 471] or the difference between client of local SGD in the non-IID case becomes worse. In order to achieve the rate of 1/ T KN (under non- convex objectives), the number of local updates K should be smaller than T 1/3 /N , instead of T /N 3 as in to make the algorithm be more robust to the heterogeneity across local objectives. The proposed FedProx clients participate, and uses batch gradient descent on clients, which can potentially converge faster than

stochastic gradients on clients. Recently, a number of works have made progress in relaxing the assumptions necessary for analysis so as of Federated Averaging in a more realistic setting where only a subset of clients are involved in each round. In order to guarantee the convergence, they assumed that the clients are selected either uniformly at random or with probabilities that are in proportion to the sizes of local datasets. Nonetheless, in practice the server may not be able to sample clients in these idealized ways — in particular, in cross-device settings only

Non-IID assumptions

Symbol Full name Explanation

BCGV bounded inter-client gradient variance Ei k∇fi (x) − ∇F (x)k2 ≤ η 2

BOBD bounded optimal objective difference F ∗ − Ei [fi∗ ] ≤ η 2

BOGV bounded optimal gradient variance Ei k∇fi (x∗ )k2 ≤ η 2

BGV bounded gradient dissimilarity Ei k∇fi (x)k2/k∇F (x)k2 ≤ η 2

Other assumptions and variants

Symbol Explanation

CVX Each client function fi (x) is convex. SCVX Each client function fi (x) is µ-strongly convex. BNCVX Each client function has bounded nonconvexity with ∇2 fi (x)  −µI. BLGV The variance of stochastic gradients on local clients is bounded. BLGN The norm of any local gradient is bounded. LBG Clients use the full batch of local samples to compute updates. Dec Decentralized setting, assumes the the connectivity of network is good. AC All clients participate in each round.

1step One local update is performed on clients in each round. Prox Use proximal gradient steps on clients. VR Variance reduction which needs to track the state.

Convergence rates

Method Non-IID Other assumptions Variant Rate

PD-SGD BCGV BLGV Dec; AC O(N/T ) + O(1/ N T )

MATCHA BCGV BLGV Dec O(1/ T KM ) + O(M/KT )

FedProx BGV BNCVX Prox O(1/ T )

SCAFFOLD - SCVX; BLGV VR O(1/T KM ) + O(e−T )

settings. We summarize the key assumptions for non-IID data, local functions on each client, and other assumptions. We also present the variant of the algorithm comparing to Federated Averaging and the con- vergence rates that eliminate constant.

devices that meet strict eligibility requirements (e.g. charging, idle, free WiFi) will be selected to participate in the computation. At different times within a day, the clients characteristics can vary significantly. Eichner of clients with different characteristics are sampled from following a regular cyclic pattern (e.g. diurnal).

Clients can perform different local steps because of heterogeneity in their computing capacities. Wang

points of a mismatched objective function in the presence of heterogeneous local steps. They refer to this problem as objective inconsistency and propose a simple technique to eliminate the inconsistency problem from federated learning algorithms. We summarize recent theoretical results in Table 6. All the methods in Table 6 assume smoothness or Lipschitz gradients for the local functions on clients. The error bound is measured by optimal objective (1) for convex functions and norm of gradient for nonconvex functions. For each method, we present the

key non-IID assumption, assumptions on each client function fi (x), and other auxiliary assumptions. We also briefly describe each method as a variant of the federated averaging algorithm, and show the simplified convergence rate eliminating constants. Assuming the client functions are strongly convex could help the convergence rate [303, 265]. Bounded gradient variance, which is a widely used assumption to analyze stochastic gradient methods, is often used when clients use stochastic local updates [305, 303, 304, 469,

updates on randomly sampled M clients in each round, and presents a rate that suggests local updates (K > 1) could slow down the convergence. Clarifying the regimes where K > 1 may hurt or help convergence is an important open problem.

Connections to decentralized optimization The objective function of federated optimization has been

studied for many years in the decentralized optimization community. As first shown in Wang and Joshi , the convergence analysis of decentralized SGD can be applied to or combined with local SGD with a proper setting of the network topology matrix (mixing matrix). In order to reduce the communication overhead,

Wang and Joshi proposed periodic decentralized SGD (PD-SGD) which allows decentralized SGD to

non-IID case. MATCHA further improves the performance of PD-SGD by randomly sampling clients for computation and communication, and provides a convergence analysis showing that local updates can accelerate convergence.

Acceleration, variance reduction and adaptivity Momentum, variance-reduction, and adaptive learn-

ing rates are all promising techniques to improve convergence and generalization of first-order methods. However, there is no single manner in which to incorporate these techniques into FedAvg. SCAFFOLD models the difference in client updates using control variates to perform variance reduction. Notably, this allows convergence results not relying on bounding the amount of heterogeneity among clients. As for age the local buffers and the local model parameters at each communication round. Although this method

empirically improves the final accuracy of local SGD, this doubles the per-round communication cost. A developing federated versions of adaptive optimization methods with the same communication cost as Fe- showed that the momentum variants of local SGD can converge to stationary points of non-convex objective

functions at the same rate as synchronous mini-batch SGD, it is challenging to prove momentum accelerates approach for adapting centralized optimization algorithms to the heterogeneous federated setting (MIME framework and algorithms).

3.3 Multi-Task Learning, Personalization, and Meta-Learning

In this section we consider a variety of “multi-model” approaches — techniques that result in effectively using different models for different clients at inference time. These techniques are particularly relevant when faced with non-IID data (Section 3.1), since they may outperform even the best possible shared global model. We note that personalization has also been studied in the fully decentralized setting [459, 59, 504, 19], where training individual models is particularly natural.

3.3.1 Personalization via Featurization

The remainder of this section specifically considers techniques that result in different users running inference with different model parameters (weights). However, in some applications similar benefits can be achieved by simply adding user and context features to the model. For example, consider a language model for next- differently, and in fact on-device personalization of model parameters has yielded significant improvements for this problem . However, a complimentary approach may be to train a federated model that takes

as input not only the words the user has typed so far, but a variety of other user and context features— What words does this user frequently use? What app are they currently using? If they are chatting, what messages have they sent to this person before? Suitably featurized, such inputs can allow a shared global model to produce highly personalized predictions. However, largely because few public datasets contain such auxiliary features, developing model architectures that can effectively incorporate context information

for different tasks remains an important open problem with the potential to greatly increase the utility of FL-trained models.

3.3.2 Multi-Task Learning

If one considers each client’s local problem (the learning problem on the local dataset) as a separate task (rather than as a shard of a single partitioned dataset), then techniques from multi-task learning im- federated learning, directly tackling challenges of communication efficiency, stragglers, and fault tolerance. In multi-task learning, the result of the training process is one model per task. Thus, most multi-task learn- ing algorithms assume all clients (tasks) participate in each training round, and also require stateful clients

since each client is training an individual model. This makes such techniques relevant for cross-silo FL applications, but harder to apply in cross-device scenarios.

Another approach is to reconsider the relationship between clients (local datasets) and learning tasks

(models to be trained), observing that there are points on a spectrum between a single global model and different models for every client. For example, it may be possible to apply techniques from multi-task learning (as well as other approaches like personalization, discussed next), where we take the “task” to be a subset of the clients, perhaps chosen explicitly (e.g. based on geographic region, or characteristics of the device or user), or perhaps based on clustering or the connected components of a learned graph over

the clients . The development of such algorithms is an important open problem. See Section 4.4.4

for a discussion of how sparse federated learning problems, such as those arising naturally in this type of multi-task problem, might be approached without revealing to which client subset (task) each client belongs.

3.3.3 Local Fine Tuning and Meta-Learning

By local fine tuning, we refer to techniques which begin with the federated training of a single model, and then deploy that model to all clients, where it is personalized by additional training on the local dataset before use in inference. This approach integrates naturally into the typical lifecycle of a model in federated learning (Section 1.1.1). Training of the global model can still proceed using only small samples of clients on each round (e.g. 100s); the broadcast of the global model to all clients (e.g. many millions) only happens once,

when the model is deployed. The only difference is that before the model is used to make live predictions on the client, a final training process occurs, personalizing the model to the local dataset.

Given a global model that

Frequently Asked Questions

What is this project about?

This project covers practical implementation and research aspects of the topic using AI/ML techniques.