FedML: A Research Library and Benchmark for Federated Machine Learning

Chaoyang He, Songze Li, Jinhyun So, Xiao Zeng, Mi Zhang, Hongyi Wang, Xiaoyang Wang, Praneeth Vepakomma, Abhishek Singh, Hang Qiu, Xinghua Zhu, Jianzong Wang, Li Shen, Peilin Zhao, Yan Kang, Yang Liu, Ramesh Raskar, Qiang Yang, Murali Annavaram, Salman Avestimehr

Introduction

Federated learning (FL) is a distributed learning paradigm that aims to train machine learning models from scattered and isolated data . FL differs from data center-based distributed training in three major aspects: 1) statistical heterogeneity, 2) system constraints, and 3) trustworthiness. Solving these unique challenges calls for efforts from a variety of fields, including machine learning, wireless communication, mobile computing, distributed systems, and information security, making federated learning a truly interdisciplinary research field.

In the past few years, more and more efforts have been made to address these unique challenges. To tackle the challenge of statistical heterogeneity, distributed optimization methods such as Adaptive Federated Optimizer , FedNova , FedProx , and FedMA have been proposed. To tackle the challenge of system constraints, researchers apply sparsification and quantization techniques to reduce the communication overheads and computation costs during the training process . To tackle the challenge of trustworthiness, existing research focuses on developing new defense techniques for adversarial attacks to make FL robust , and proposing methods such as differential privacy (DP) and secure multiparty computation (SMPC) to protect privacy .

Although a lot of progress has been made, existing efforts are confronted with a number of limitations that we argue are critical to FL research:

Distributed training libraries in PyTorch , TensorFlow , MXNet , and distributed training-specialized libraries such as Horovod and BytePS are designed for distributed training in data centers. Although simulation-oriented FL libraries such as TensorFlow-Federated (TFF) , PySyft , and LEAF are developed, they only support centralized topology-based FL algorithms like FedAvg or FedProx with simulation in a single machine, making them unsuitable for FL algorithms which require the exchange of complex auxiliary information and customized training procedure. Production-oriented libraries such as FATE and PaddleFL are released by industry. However, they are not designed as flexible frameworks that aim to support algorithmic innovation for open FL problems.

Lack of support of diverse FL configurations.

FL is diverse in network topology, exchanged information, and training procedures. In terms of network topology, a variety of network topologies such as vertical FL , split learning , decentralized FL , hierarchical FL , and meta FL have been proposed. In terms of exchanged information, besides exchanging gradients and models, recent FL algorithms propose to exchange information such as pseudo labels in semi-supervised FL and architecture parameters in neural architecture search-based FL . In terms of training procedures, the training procedures in federated GAN and transfer learning-based FL are very different from the vanilla FedAvg algorithm . Unfortunately, such diversity in network topology, exchanged information, and training procedures is not supported in existing FL libraries.

Lack of standardized FL algorithm implementations and benchmarks.

The diversity of libraries used for algorithm implementation in existing work makes it difficult to fairly compare their performance. The diversity of benchmarks used in existing work also makes it difficult to fairly compare their performance. The non-I.I.D. characteristic of FL makes such comparison even more challenging : training the same DNN on the same dataset with different non-I.I.D. distributions produces varying model accuracies; one algorithm that achieves higher accuracy on a specific non-I.I.D. distribution than the other algorithms may perform worse on another non-I.I.D. distribution. In Table 8, we summarize the datasets and models used in existing work published at the top tier machine learning conferences such as NeurIPS, ICLR, and ICML in the past two years. We observe that the experimental settings of these work differ in terms of datasets, non-I.I.D. distributions, models, and the number of clients involved in each round. Any difference in these settings could affect the results.

In this work, we present FedML, an open research library and benchmark to address the aforementioned limitations and facilitate FL research. FedML provides an end-to-end toolkit to facilitate FL algorithm development and fair performance comparison under diverse computing paradigms and configurations. Table 1 summarizes the key differences between FedML and existing FL libraries and benchmarks. The highlights of FedML are summarized below:

(i) Support of diverse FL computing paradigms. FedML supports three diverse computing paradigms: 1) on-device training for edge devices including smartphones and Internet of Things (IoT), 2) distributed computing, and 3) single-machine simulation to meet algorithmic and system-level research requirements under different system deployment scenarios.

(ii) Support of diverse FL configurations. FedML introduces a worker/client-oriented programming interface to enable diverse network topologies, flexible information exchange among workers/clients, and various training procedures.

(iii) Standardized FL algorithm implementations. FedML includes standardized implementations of status quo FL algorithms. These implementations not only help users to familiarize FedML APIs but also can be used as baselines for comparisons with newly developed FL algorithms.

(iv) Standardized FL benchmarks. FedML provides standardized benchmarks with well-defined evaluation metrics, multiple synthetic and real-world non-I.I.D. datasets, as well as verified baseline results to facilitate fair performance comparison.

(v) Fully open and evolving. FL is a research field that evolves at a considerably fast pace. This requires FedML to adapt at the same pace. We will continuously expand FedML to optimize three computing paradigms and support more algorithms (distributed optimizer) and benchmarks (models and datasets) for newly explored usage scenarios. FedML is fully open and welcomes contributions from the FL research community as well. We hope researchers in diverse FL applications could contribute more valuable models and realistic datasets to our community. Promising application domains include, but are not limited to, computer vision , natural language processing , finance , transportation , digital health , recommendation , robotics , and smart cities .

FedML Library: Architecture Design

Figure 1 provides an overview of FedML library. The FedML library has two key components: FedML-API and FedML-core, which represents high-level API and low-level API, respectively.

FedML-core separates distributed communication and model training into two separate modules. The distributed communication module is responsible for low-level communication among different workers/clients. The communication backend is based on MPI (message passing interface)https://pypi.org/project/mpi4py/. Inside the distributed communication module, a TopologyManager supports a variety of network topologies that can be used in many existing FL algorithms . In addition, security/privacy-related functions are also supported. The model training module is built upon PyTorch. Users can implement workers (trainers) and coordinators according to their needs.

FedML-API is built upon FedML-core. With the help of FedML-core, new algorithms in distributed version can be easily implemented by adopting the client-oriented programming interface, which is a novel design pattern for flexible distributed computing (Section 3). Such a distributed computing paradigm is essential for scenarios in which large DNN training cannot be handled by standalone simulation due to GPU memory and training time constraints. This distributed computing design is not only used for FL, but it can also be used for conventional in-cluster large-scale distributed training. FedML-API also suggests a machine learning system practice that separates the implementations of models, datasets, and algorithms. This practice enables code reuse and fair comparison, avoiding statistical or system-level gaps among algorithms led by non-trivial implementation differences. Another benefit is that FL applications can develop more models and submit more realistic datasets without the need to understand the details of different distributed optimization algorithms.

One key feature of FedML is its support of FL on real-world hardware platforms. Specifically, FedML includes FedML-Mobile and FedML-IoT, which are two on-device FL testbeds built upon real-world hardware platforms. Currently, FedML-Mobile supports Android smartphones and FedML-IoT supports Raspberry PI 4 and NVIDIA Jetson Nano (see Appendix C for details). With such testbeds built upon real-world hardware platforms, researchers can evaluate realistic system performance, such as training time, communication, and computation cost. To support conducting experiments on those real-world hardware platforms, our FedML architecture design can smoothly transplant the distributed computing code to the FedML-Mobile and FedML-IoT platforms, reusing nearly all algorithmic implementations in the distributed computing paradigm. Moreover, for FedML-IoT, researchers only need to program with Python to customize their research experiments without the need to learn new system frameworks or programming languages (e.g., Java, C/C++)Please check here for details: https://github.com/FedML-AI/FedML/tree/master/fedml_iot.

FedML Library: Programming Interface

The goal of the FedML programming interface is to provide simple user experience to allow users to build distributed training applications (e.g. design customized message flow and topology definitions) by only focusing on algorithmic implementations while ignoring the low-level communication backend details.

As shown in Figure 2(b), FedML provides the worker-oriented programming design pattern, which can be used to program the worker behavior when participating in training or coordination in the FL algorithm. We describe it as worker-oriented because its counterpart, the standard distributed training library (as the torch.distributed exampleMore details can be found at https://pytorch.org/tutorials/intermediate/dist_tuto.html shown in Figure 2(a)), normally completes distributed training programming by describing the entire training procedure rather than focusing on the behavior of each worker.

With the worker-oriented programming design pattern, the user can customize its own worker in FL network by inheriting the WorkerManager class and utilizing its predefined APIs register_message_receive_handler and send_message to define the receiving and sending messages without considering the underlying communication mechanism (as shown in the highlighted blue box in Figure 2(b)). Conversely, existing distributed training frameworks do not have such flexibility. In order to make the comparison clearer, we use the most popular machine learning framework PyTorch as an example. Figure 2(a) illustrates a complete training procedure (distributed synchronous SGD) and aggregates gradients from all other workers with the all_reduce messaging passing interface. Although it supports multiprocessing training, it cannot flexibly customize different messaging flows in any network topology. In PyTorch, another distributed training API, torch.nn.parallel.paraDistributedDataParallelIt is recommended to use torch.nn.parallel.paraDistributedDataParallel instead of torch.nn.DataParallel. For more details, please refer to https://pytorch.org/tutorials/intermediate/ddp_tutorial.html and https://pytorch.org/docs/master/notes/cuda.html#cuda-nn-ddp-instead, also has such inflexibly.

Message definition beyond gradient and model.

FedML also supports message exchange beyond the gradient or model from the perspective of message flow. This type of auxiliary information may be due to either the need for algorithm design or the need for system-wide configuration delivery. Each worker defines the message type from the perspective of sending. Thus, in the above introduced worker-oriented programming, the WorkerManager should handle messages defined by other trainers and also send messages defined by itself. The sending message is normally executed after handling the received message. As shown in Figure 2(b), in the yellow background highlighted code snippet, workers can send any message type and related message parameters using the train() function.

Topology management.

As demonstrated in Figure 3, FL has various topology definitions, such as vertical FL , split learning , decentralized FL , and Hierarchical FL . In order to meet such diverse requirements, FedML provides TopologyManager to manage the topology and allows users to send messages to arbitrary neighbors during training. Specifically, after the initial setting of TopologyManager is completed, for each trainer in the network, the neighborhood worker ID can be queried via the TopologyManager. In line 26 of Figure 2(b), we see that the trainer can query its neighbor nodes through the TopologyManager before sending its message.

Trainer and coordinator.

We also need the coordinator to complete the training (e.g., in FedAvg, the central worker is the coordinator while the others are trainers). For the trainer and coordinator, FedML does not over-design. Rather, it gives the implementation completely to the developers, reflecting the flexibility of our framework. The implementation of the trainer and coordinator is similar to the process in Figure 2(a), which is consistent with the training implementation of a standalone version training. We provide some reference implementations of different trainers and coordinators in our source code (Section 4.1).

Privacy, security, and robustness.

While the FL framework facilitates data privacy by keeping data locally available to the users and only requiring communication for model updates, users may still be concerned about partial leakage of their data which may be inferred from the communicated model (e.g., ). Aside from protecting the privacy of users’ data, another critical security requirement for the FL platform, especially when operating over mobile devices, is the robustness towards user dropouts. Specifically, to accomplish the aforementioned goals of achieving security, privacy, and robustness, various cryptography and coding-theoretic approaches have been proposed to manipulate intermediate model data .

To facilitate rapid implementation and evaluation of data manipulation techniques to enhance security, privacy, and robustness, we include low-level APIs that implement common cryptographic primitives such as secrete sharing, key agreement, digital signature, and public key infrastructure. We also plan to include an implementation of Lagrange Coded Computing (LCC) . LCC is a recently developed coding technique on data that achieves optimal resiliency, security (against adversarial nodes), and privacy for any polynomial evaluations on the data. Finally, we plan to provide a sample implementation of the secure aggregation algorithm using the above APIs.

In standard FL settings, it is assumed that there is no single central authority that owns or verifies the training data or user hardware, and it has been argued by many recent studies that FL lends itself to new adversarial attacks during decentralized model training . Several robust aggregation methods have been proposed to enhance the robustness of FL against adversaries .

To accelerate generating benchmark results on new types of adversarial attacks in FL, we include the latest robust aggregation methods presented in literature including (i) norm difference clipping ; weak differential private (DP) ; (ii) RFA (geometric median) ; (iii) Krum and (iv) Multi-Krum . Our APIs are easily extendable to support newly developed types of robust aggregation methods. On the attack end, we observe that most of the existing attacks are highly task-specific. Thus, it is challenging to provide general adversarial attack APIs. Our APIs support the backdoor with model replacement attack presented in and the edge-case backdoor attack presented in to provide a reference for researchers to develop new attacks.

FedML Benchmark: Algorithms, Models, and Datasets

As shown in Figure 4, FedML is capable of supporting FL algorithms that are diverse in network topology, exchanged information, and training procedures. These supported algorithms can be used as implementation examples and baselines to help users develop and evaluate their own algorithms. Currently, FedML includes the standard implementations of multiple status quo FL algorithms: Federated Averaging (FedAvg) , Decentralized FL , Vertical Federated Learning (VFL) , Split learning , Federated Neural Architecture Search (FedNAS) , and Turbo-Aggregate . Fore more details of these algorithms, please refer to Appendix B.1.

We will keep following the latest algorithm to be published at top-tier machine learning conferences, and will continuously add new FL algorithms such as Adaptive Federated Optimizer , FedNova , FedProx , and FedMA in near future.

2 Models and Datasets

Inconsistent usage of datasets, models, and non-I.I.D. partition methods makes it difficult to fairly compare the performance of FL algorithms (in Table 8, we summarize the non-I.I.D. datasets and models used in existing work published at the top tier machine learning venues in the past two years). To enforce fair comparison, FedML benchmark explicitly specifies the combinations of datasets, models, and non-I.I.D. partition methods to be used for experiments. In particular, we divide the benchmark into three categories: 1) linear models (convex optimization), 2) lightweight shallow neural networks (non-convex optimization), and 3) deep neural networks (non-convex optimization).

The linear model category is used for convex optimization experiments such as the ones in and . In this category, we include three datasets (Table 2): MNIST , Federated EMNIST , and Synthetic (α\alpha, β\beta) , with the logistic regression as the baseline model.

Federated datasets for lightweight shallow neural networks (non-convex optimization).

Due to resource constraints of edge devices, shallow neural networks are commonly used in existing work for experiments. In this category, we include four datasets (Table 3): Federated EMNIST , CIFAR-100 , Shakespeare , and StackOverflow . Please refer to Appendix B.2 for more details.

Federated datasets for deep neural networks (non-convex optimization).

Given the resource constraints of edge devices, large DNN models are usually trained under the cross-organization FL (also called cross-silo FL) setting. For example, has studied large DNN models for cross-silo FL in the hospital scenario. However, large DNN models dominate the accuracy in most learning tasks. Pushing FL of large DNN models on edge devices is challenging but a meaningful endeavor, which motivates us to make this benchmark category. For example, has proposed an efficient training algorithm for large CNN models on edge devices. Table 4 shows datasets and models we include in this category. Please refer to Appendix B.2 for more details.

Experiments

FedML provides benchmark experimental results as references for newly developed FL algorithms. To ensure real-time updates, we maintain benchmark experimental results using Weight and Biashttps://www.wandb.com/. The web link to view all the benchmark experimental results can be found at our GitHub repository.

To demonstrate the capability of FedML, we ran experiments in a real distributed computing environment. We trained two CNNs (ResNet-56 and MobileNet) using the standard FedAvg algorithm. Table 5 shows the experimental results, and Figure 5 shows the corresponding test accuracy during training. As shown, the accuracy of the non-I.I.D. setting is lower than that of the I.I.D. setting, which is consistent with findings reported in prior work .

We also compared the training time of distributed computing with that of standalone simulation. The result in Table 6 reveals that when training large CNNs, the standalone simulation is about 8 times slower than distributed computing with 10 parallel workers. Therefore, when training large DNNs, we suggest using FedML’s distributed computing paradigm, which is not supported by existing FL libraries such as PySyft , LEAF , and TTF . Moreover, FedML supports multiprocessing in a single GPU card which enables FedML to run a large number of training workers by using only a few GPU cards. As an example, when training ResNet on CIFAR-10, FedML can run 112 workers in a server with 8 GPUs.

Conclusion

FedML is a research-oriented federated learning library and benchmark. We hope it could provide researchers and engineers with an end-to-end toolkit to facilitate developing FL algorithms and fairly comparing with existing algorithms. We welcome any useful feedback from the readers, and will continuously update FedML to support the research of the federated learning community.

References

Appendix A The Taxonomy of Research Areas and a Comprehensive Publication List

Appendix B Benchmark

Federated Averaging (FedAvg). FedAvg is a standard federated learning algorithm that is normally used as a baseline for advanced algorithm comparison. We summarize the algorithm message flow in Figure 4(a). Each worker trains its local model for several epochs, then updates its local model to the server. The server aggregates the uploaded client models into a global model by weighted coordinate-wise averaging (the weights are determined by the number of data points on each worker locally), and then synchronizes the global model back to all workers. In our FedML library, based on the worker-oriented programming, we can implement this algorithm in a distributed computing manner. We suggest that users start from FedAvg to learn using FedML.

Decentralized FL. We use , a central server free FL algorithm, to demonstrate how FedML supports decentralized topology with directed communication. As Figure 4(b) shows, such an algorithm uses a decentralized topology, and more specifically, some workers do not send messages (model) to all of their neighbors. The worker-oriented programming interface can easily meet this requirement since it allows users to define any behavior for each worker.

Vertical Federated Learning (VFL). VFL or feature-partitioned FL is applicable to the cases where all participating parties share the same sample space but differ in the feature space. As illustrated in Figure 4(c), VFL is the process of aggregating different features and computing the training loss and gradients in a privacy-preserving manner to build a model with data from all parties collaboratively . The FedML library currently supports the logistic regression model with customizable local feature extractors in the vertical FL setting, and it provides NUS-WIDE and lending club loan datasets for the experiments.

Split Learning. Split learning is computing and memory-efficient variant of FL introduced in where the model is split at a layer and the parts of the model preceding and succeeding this layer are shared across the worker and server, respectively. Only the activations and gradients from a single layer are communicated in split learning, as against that the weights of the entire model are communicated in federated learning. Split learning achieves better communication-efficiency under several settings, as shown in . Applications of this model to wireless edge devices are described in . Split learning also enables matching client-side model components with the best server-side model components for automating model selection as shown in work on ExpertMatcher .

Federated Neural Architecture Search (FedNAS). FedNAS is a federated neural architecture search algorithm that enables scattered clients to collaboratively search for a neural architecture. FedNAS differs from other FL algorithms in that it exchanges information beyond gradient even though it has a centralized topology similar to FedAvg.

B.2 Details of Datasets

Federated EMNIST: EMNIST consists of images of digits and upper and lower case English characters, with 62 total classes. The federated version of EMNIST partitions the digits by their author. The dataset has natural heterogeneity stemming from the writing style of each person.

CIFAR-100: Google introduced a federated version of CIFAR-100 by randomly partitioning the training data among 500 clients, with each client receiving 100 examples . The partition method is Pachinko Allocation Method (PAM) .

Shakespeare: first introduced this dataset to FL community. It is a dataset built from The Complete Works of William Shakespeare. Each speaking role in each play is considered a different device.

StackOverflow : Google TensorFlow Federated (TFF) team maintains this federated dataset, which is derived from the Stack Overflow Data hosted by kaggle.com. We integrate this dataset into our benchmark.

CIFAR-10 and CIFAR-100. CIFAR-10 and CIFAR-100 both consists of 32×\times32 colorr images. CIFAR-10 has 10 classes, while CIFAR-100 has 100 classes. Following and , we use latent Dirichlet allocation (LDA) to partition the dataset according to the number of workers involved in training in each round.

CINIC-10. CINIC-10 has 4.5 times as many images as that of CIFAR-10. It is constructed from two different sources: ImageNet and CIFAR-10. It is not guaranteed that the constituent elements are drawn from the same distribution. This characteristic fits for federated learning because we can evaluate how well models cope with samples drawn from similar but not identical distributions.

B.3 Lack of Fair Comparison: Diverse Non-I.I.D. Datasets and Models

Appendix C IoT Devices

Currently, we support two IoT devices: Raspberry PI 4 (Edge CPU Computing) and NVIDIA Jetson Nano (Edge GPU Computing).

Raspberry Pi 4 Desktop kit is supplied with:

Raspberry Pi 4 Model B (2GB, 4GB or 8GB version)

2 × micro HDMI to Standard HDMI (A/M) 1m Cables

16GB NOOBS with Raspberry Pi OS microSD card

For more details, please check this link: https://www.raspberrypi.org/products/raspberry-pi-4-desktop-kit.

C.2 NVIDIA Jetson Nano (Edge GPU Computing)

NVIDIA® Jetson Nano™ Developer Kit is a small, powerful computer that lets you run multiple neural networks in parallel for applications like image classification, object detection, segmentation, and speech processing. All in an easy-to-use platform that runs in as little as 5 watts.

For more details, please check this link: https://developer.nvidia.com/embedded/jetson-nano-developer-kit.