scieee AI-readable full text Open interactive document viewer

Communication-Efficient Distributed Deep Learning via Federated Dynamic Averaging

Theologitis, Michael; Frangias, Georgios; Anestis, Georgios; Samoladas, Vasilis; Deligiannakis, Antonios

Abstract

Driven by the ever-growing volume and decentralized nature of data, coupled with the need to harness this data and generate knowledge from it, has led to the extensive use of distributed deep learning (DDL) techniques for training. These techniques rely on local training that is performed at the distributed nodes based on locally collected data, followed by a periodic synchronization process that combines these models to create a global model. However, frequent synchronization of DL models, encompassing millions to many billions of parameters, creates a communication bottleneck, severely hindering scalability. Worse yet, DDL algorithms typically waste valuable bandwidth, and make themselves less practical in bandwidth-constrained federated settings, by relying on overly simplistic, periodic, and rigid synchronization schedules. These drawbacks also have a direct impact on the time required for the training process, necessitating excessive time for data communication. To address these shortcomings, we propose Federated Dynamic Averaging (FDA), a communication-efficient DDL strategy that dynamically triggers synchronization based on the value of the model variance. In essence, the costly synchronization step is triggered only if the local models, which are initialized from a common global model after each synchronization, have significantly diverged. This decision is facilitated by the communication of a small local state from each distributed node/worker. Through extensive experiments across a wide range of learning tasks we demonstrate that FDA reduces communication cost by orders of magnitude, compared to both traditional and cutting-edge communication-efficient algorithms. Additionally, we show that FDA maintains robust performance across diverse data heterogeneity settings.

Full text

Communication-Efficient Distributed Deep Learning via Federated Dynamic Averaging Michail Theologitis Technical University of Crete Chania, Greece [email protected] Georgios Frangias Technical University of Crete Chania, Greece [email protected] Georgios Anestis Technical University of Crete Chania, Greece [email protected] Vasilis Samoladas Technical University of Crete Chania, Greece [email protected] Antonios Deligiannakis Technical University of Crete Chania, Greece [email protected] ABSTRACT Driven by the ever-growing volume and decentralized nature of data, coupled with the need to harness this data and generate knowledge from it, has led to the extensive use of distributed deep learning (DDL) techniques for training. These techniques rely on local training that is performed at the distributed nodes based on locally collected data, followed by a periodic synchronization process that combines these models to create a global model. However, frequent synchronization of DL models, encompassing millions to many billions of parameters, creates a communication bottleneck, severely hindering scalability. Worse yet, DDL algorithms typically waste valuable bandwidth, and make themselves less practical in bandwidth-constrained federated settings, by relying on overly simplistic, periodic, and rigid synchronization schedules. These drawbacks also have a direct impact on the time required for the training process, necessitating excessive time for data communication. To address these shortcomings, we propose Federated Dynamic Averaging (FDA), a communication-efficient DDL strategy that dynamically triggers synchronization based on the value of the model variance. In essence, the costly synchronization step is triggered only if the local models, which are initialized from a common global model after each synchronization, have significantly diverged. This decision is facilitated by the communication of a small local state from each distributed node/worker. Through extensive experiments across a wide range of learning tasks we demonstrate that FDA reduces communication cost by orders of magnitude, compared to both traditional and cutting-edge communicationefficient algorithms. Additionally, we show that FDA maintains robust performance across diverse data heterogeneity settings. 1 INTRODUCTION The big data era has been marked by an unprecedented scale of training datasets [ 41 , 67 ]. These datasets are not only growing in size, but are often physically distributed and cannot be easily centralized due to business considerations, privacy concerns, bandwidth limitations (especially in federated settings, such as drones collecting and collaboratively building a global model/view of an area), and data sovereignty laws [ 9 , 23 , 64 ]. Such constraints complicate the use of Deep Learning (DL) techniques in the aforementioned scenarios. Β©2025 Copyright held by the owner/author(s). Published in Proceedings of the 28th International Conference on Extending Database Technology (EDBT), 25th March-28th March, 2025, ISBN 978-3-89318-098-1 on OpenProceedings.org. Distribution of this paper is permitted under the terms of the Creative Commons license CC-by-nc-nd 4.0. Distributed Deep Learning (DDL) has emerged as an alternative paradigm to the traditional centralized approach [ 6 , 69 ], offering efficient learning over large-scale data across multiple worker-nodes, enhancing the speed of training DL models and paving the way for more scalable and resilient DL applications [ 10 , 28 , 35 , 55 , 68 ]. Most DDL methods are iterative, where, in each iteration, some amount of local training is followed by synchronization of the local models with the global one. The predominant method, based on the bulk synchronous parallel (BSP) approach [ 56 ], is to average the local model updates and then apply the average update to each local model [ 69 ]. Less synchronized variants have also been proposed, to ameliorate the effect of straggler workers [ 14 , 37 ] but compromise convergence speed and model quality. A significant challenge inherent in the traditional techniques, especially in federated DL settings, where models are huge and worker interconnections are slow, is the communication bottleneck, restricting system scalability [ 53 , 60 ]. Specifically, the communication bottleneck arises from the frequent exchange (synchronization) of model parameters, often in the range of billions, across distributed workers. The synchronization process entails substantial data volume transfer and generally dominates the overall training time, leading to a low computationto-communication ratio [ 14 , 46 ]. Addressing this challenge to expedite DDL algorithms has been a focal point of research for many years; speeding-up SGD is arguably among the most impactful and transformative problems in machine learning [58]. The most direct method to alleviate the communication burden is to reduce the frequency of communication rounds. Local-SGD is the prime example of this approach. It allows workers to perform 𝜏 local update steps on their models before aggregating them, as opposed to averaging the updates in every step [ 17 , 66 ]. Although Local-SGD is effective in reducing communication while maintaining comparable model quality [ 58 ], determining the optimal value of 𝜏 presents a critical challenge, with only a handful of studies offering theoretical insights into its influence on convergence [50, 58, 66]. To further reduce communication costs of Local-SGD, more sophisticated communication strategies introduce varying sequences of local update steps {𝜏0, ...,πœπ‘…} , instead of a fixed 𝜏 . In [ 57 ], in order to minimize convergence error with respect to wall-time, the authors proposed a decreasing sequence of local update steps. Conversely, the focus in [ 17 ] was on reducing the number of communication rounds for a fixed number of model updates and an increasing sequence emerged. These contrasting approaches underscore the multifaceted nature of communication strategies in distributed deep learning, highlighting not only Series ISSN: 2367-2005 411 10.48786/edbt.2025.33 the absence of a one-size-fits-all solution but also the growing need for dynamic, context-aware strategies that can continuously adapt to the specific intricacies of the learning task. Main Idea and Contributions. Our work addresses critical efficiency challenges in DDL, particularly in communicationconstrained environments, such as the ones encountered in Federated Learning (FL) applications [ 23 ]. We introduce Federated Dynamic Averaging (FDA), a novel, adaptive distributed deep learning strategy that massively improves communication efficiency over previous work. FDA utilizes a novel 2-action, conditional synchronization protocol, designed to avoid the need to decide or guess the proper values of local update steps, or to synchronize after each training step, but rather only performs the costly synchronization process when needed. Our FDA algorithm dynamically triggers synchronization based on the value of model variance across worker-nodes. In a nutshell, the costly synchronization step is only triggered if the local models have diverged significantly, which implies that the global model may no longer be accurate. As Figure 1 demonstrates, at the start, workers enter the local training step with the same global model (Figure 1.A). Then, local training commences and each distributed worker-node computes its local state, which encapsulates helpful information for estimating the model variance (Figure 1.B). This is followed by the transmission (Figure 1.C) of these small-size local states, an operation that is bandwidthand time-efficient because of their small size. During transmission, the local states are aggregated and their average is made available to all workersβ€”an operation known as AllReduce. This operation does not require (or prohibit) the use of a central node. Based on the aggregated state, the workers can estimate (Figure 1.D) whether the variance of the local models may have exceeded a threshold. If this is not the case, the costly synchronization step (Figure 1.E) is avoided and local training continues. What is important is how to properly pick these local states computed at, and then transmitted by, the local workers. To address this problem, we propose two variants of our FDA algorithm. Our contributions can be summarized as follows: β€’ We propose FDA, an algorithm that dynamically decides to synchronize local workers when model variance across workers exceeds a threshold. This strategy drastically reduces communication, while preserving cohesive progress towards the shared training objective. β€’ We propose two variants of FDA, which differ in the amount of information preserved in the local states that are transmitted by each worker and aggregated for subsequent estimation of model variance. These two variants, termed SketchFDA and LinearFDA, offer a different balance between communication efficiency and approximation accuracy. β€’ We evaluate and compare FDA with other DDL algorithms through a comprehensive suite of experiments with diverse datasets, models, and tasks. Our experiments demonstrate that FDA outperforms traditional and contemporary FL algorithms by 1-2 orders of magnitude in communication savings, while maintaining equivalent model performance. Furthermore, it effectively balances the competing demands of communication and computation, providing greatly improved trade-offs. β€’ We demonstrate FDA’s robustness in various challenging NonIID settings, common in real-world Federated Learning applications. While state-of-the-art methods typically require substantially more resources to converge under Non-IID conditions, Nodes start training step Local training step / compute local state Estimate if synchronization is needed. If not go to Step (B) Aggregate local state using AllReduce Synchronize models using AllReduce. Go to Step (A) (A) (B) (C) (D) (E) Figure 1: FDA. The local training step is followed by the computation of a local state by all worker-nodes. Then, the (small in size) local states are aggregated. Based on the aggregated result, all workers estimate if synchronization is required. In most cases, the expensive synchronization step of the models is avoided and local training continues FDA maintains consistent and comparable performance across both IID and Non-IID settings. Outline. The remainder of this paper is organized as follows: Section 2 reviews related work. Section 3 introduces our DDL technique, Federated Dynamic Averaging (FDA), and its two variants. Section 4 details the experimental setup, and discusses the insights and conclusions drawn from our empirical investigation. Lastly, Section 5 contains concluding remarks. 2 RELATED WORK Problem formulation. Consider distributed training of deep neural networks over multiple workers [ 11 , 31 ]. In this setting, each worker represents a data owner (equivalently, a local model owner) and has access to its own set of training data Dπ‘˜ . Workers can utilize any available hardware they possess (e.g., GPUs, CPUs) to perform learning steps. The collective goal is to find a common model w ∈R𝑑 by minimizing the overall training loss. This scenario can be effectively modeled as a distributed optimization problem, formulated as follows: minimize w∈R𝑑𝐹(w)β‰œ1 𝐾 𝐾 βˆ‘οΈ π‘˜=1 πΉπ‘˜(w)(1) where 𝐾 is the number of workers and πΉπ‘˜( w )β‰œ Eπœπ‘˜βˆΌDπ‘˜[β„“(w;πœπ‘˜)] is the local objective function for worker π‘˜ . Function β„“( w; πœπ‘˜) represents the loss for data sample πœπ‘˜given model w. Solution direction. As noted in the seminal work [ 23 ], research in FL should focus primarily on synchronous solutions. This allows different lines of research (e.g., compression, privacy, etc.) to be developed independently and then combined seamlessly. Our work, along with most communication-efficient FL strategies, adheres to this synchronous paradigm. However, such approaches may be less effective in environments where each communication operation incurs significant overhead regardless of the size of the data being transmitted (e.g., high-latency). In these scenarios, asynchronous mechanisms become necessary, though they 412 typically fall outside the primary focus of contemporary FL research. That said, FDA can be modified to work asynchronously (as explained in Section 3.3). Communication efficient Local-SGD. The work in [ 31 ] decomposes each round into two phases. In the first phase, each worker runs Local-SGD with 𝜏=𝐼1 , while the second phase runs 𝐼2 steps with 𝜏= 1; [ 31 ] proposes to exponentially decay 𝐼1 every 𝑀 rounds. In the heterogeneous setting, the work in [ 40 ], by analysing the convergence rate, proposes an increasing sequence of local update steps for strongly-convex local objectives and fixed local update steps for other types of local objectives. The study in [ 65 ] dynamically increases batch sizes to reduce communication rounds, maintaining the same convergence rate as SSP-SGD. However, the large-batch approach leads to poor generalization [ 20 ], a challenge addressed by the post-local SGD method [ 32 ], which divides training into two phases: BSP-SGD followed by Local-SGD with a fixed number of steps. In the Lazily Aggregated Algorithm (LAG) [ 5 ], a different approach was taken, using only new gradients from some selected workers and reusing the outdated gradients from the rest, which essentially skips communication rounds. Federated Averaging (FedAvg) [ 36 ] is another representative of communication efficient Local-SGD algorithms, which is a pivotal method in Federated Learning (FL) [ 23 ]. In the FL setting with edge computing systems, the work in [ 59 ] tries to find the optimal synchronization period 𝜏 subject to local computation and aggregation constraints. Recently [ 38 ], in the FL setting with the assumption of strongly-convex objectives, by analysing the balance between fast convergence and higher-round completion rate, a decaying local update step scheme emerged. Unlike previous approaches that rely on predetermined synchronization schedules (fixed, decaying, or otherwise), our work introduces a dynamic synchronization strategy. FDA adapts continuously during the training process, basing synchronization decisions on a real-time metric: the model variance across workers. Accelerating convergence. An indirect, yet highly effective way to mitigate the communication burden in DDL, is to speed up convergence. Consequently, recent works have built upon communication efficient Local-SGD methods by deploying accelerated versions of SGD to the distributed setting. Specifically, FedAdam [ 42 ] extends Adam [ 26 ] and FedAvgM [ 21 ] extends SGD with momentum (SGD-M) [ 51 ]. Recently, Mime [ 24 ] provides a framework to adapt arbitrary centralized optimization algorithms to the FL setting. However, these methods still suffer from the model divergence problem, particularly in heterogeneous settings. When solving (1) , the disparity between each worker’s optimal solution w βˆ— π‘˜ for their objective πΉπ‘˜ , and the global optimum w βˆ— for 𝐹 , can potentially cause worker models to diverge (drift) towards their disparate minima [ 25 , 42 , 63 ]. The result is slow and unstable convergence with significant communication overhead. To address this problem, the SCAFFOLD algorithm [ 25 ] used control-variates (in the same spirit to SVRG), with significant speed-up. FedProx [ 45 ] re-parameterized FedAvg [ 36 ] by adding 𝐿2 regularization in the workers’ objectives to be near the global model. Lastly, FedDyn [ 2 ] improved upon these ideas with a dynamic regularizer making sure that if local models converge to a consensus, this consensus point aligns with the stationary point of the global objective function. While these approaches primarily focus on enhancing the optimization process and typically employ fixed synchronization intervals (e.g., every local epoch), our work addresses a complementary aspect: determining the optimal timing for synchronization. FDA’s dynamic synchronization strategy is orthogonal to these optimization techniques and can be integrated with them by simply adjusting the synchronization decision. Compression. To reduce communication overhead in DDL, significant efforts have been directed towards minimizing message sizes. Key strategies include sparsification, where only crucial components of information are transmitted, as explored in [ 3 ], and quantization techniques, which involve transmitting only quantized gradients, as detailed in [ 47 ]. These techniques can be combined with Local-SGD methods to enhance communicationefficiency further. An example is Qsparse-local-SGD [ 4 ], which integrates aggressive sparsification and quantization with LocalSGD, achieving substantial communication savings. Crucially, FDA is fully compatible with any technique that reduces the cost of synchronization (e.g. model compression). Our approach simply adjusts the timing of the synchronization decision without altering the data being synchronized. This ensures that any compression technique effective in traditional methods (BSP, Local-SGD, etc.) will be equally effective when deployed with FDA. Therefore, the communication savings demonstrated in the relevant literature [ 61 ] can be safely expected to carry over to our approach as well. Additionally, sketching emerges as another fundamental tool in large-scale machine learning. It effectively compresses highdimensional problems into lower dimensions to save runtime and memory, typically utilizing hash-based probabilistic data structures. For instance, [ 49 ] use Count Sketches to compress auxiliary variables in optimization algorithms, significantly freeing up memory. Similarly, FetchSGD [ 43 ] employs Count Sketches to compress model updates and leverages their linearity for efficient merging. In contrast to these applications, our approach utilizes sketches not for compression but to estimate local state information, and based on this to decide whether a synchronization is requiredβ€”an orthogonal application to traditional use cases. A comprehensive survey of compression techniques in DDL can be found in [61]. 3 FEDERATED DYNAMIC AVERAGING We now present our algorithms, based on our notion of Federated Dynamic Averaging (FDA). Our algorithms deviate from prior work in these two key ways: (1) The decision on when to synchronize. (2) The actual synchronization process. To the best of our knowledge, this is the first Distributed Deep Learning algorithm that dynamically decides when to synchronize based on the current collective state of the training progressβ€” whether it is advancing well or poorly. Notation. At each time step 𝑑 , each worker π‘˜ independently maintains its own vector of model parameters 1 , denoted as w (π‘˜) π‘‘βˆˆR𝑑 . Let w 𝑑 represent the 𝐾×𝑑 tensor of all local model vectors, and w𝑑 be the average model vector (this notation applies to all vector quantities): w𝑑=hw(1) 𝑑, . . . , w(𝐾) 𝑑i,w𝑑= 1 𝐾 𝐾 βˆ‘οΈ π‘˜=1 w(π‘˜) 𝑑 1 The terms β€œmodel” and "model parameters" are used interchangeably, as is common in the literature. 413 Table 1: Notation Symbol Meaning ⟨·,·⟩Dot product 𝑑Time step index 𝐾Number of workers 𝑑Model dimension Dπ‘˜Training data of worker π‘˜ B(π‘˜) 𝑑A batch sampled from Dπ‘˜ w(π‘˜) π‘‘βˆˆR𝑑Model of worker π‘˜ w𝑑=[w(1) 𝑑, . . . , w(𝐾) 𝑑]Tensor of local models w𝑑=1 𝐾Í𝐾 π‘˜=1w(π‘˜) 𝑑Average model (global model) w𝑑0Model after most recent sync. wπ‘‘βˆ’1Model after 2nd most recent sync. u(π‘˜) 𝑑=w(π‘˜) π‘‘βˆ’w𝑑0Local model drift u𝑑=1 𝐾Í𝐾 π‘˜=1u(π‘˜) 𝑑Average model drift (global drift) Var (w𝑑)Model variance ΘModel variance threshold S(π‘˜) 𝑑State of worker π‘˜ S𝑑=1 𝐾Í𝐾 π‘˜=1S(π‘˜) 𝑑Average state 𝐻(Β·) Function for variance estimation sk(Β·) :R𝑑→Rπ‘™Γ—π‘šAMS sketch operator (Β§3.1) M2(Β·) :Rπ‘™Γ—π‘šβ†’R𝐿2norm squared estimate (Β§3.1) πœ–Error of sketch estimate (Β§3.1) (1βˆ’π›Ώ)Confidence of approximation (Β§3.1) 𝑙=O(log 1/𝛿)#Rows of sketch matrix (Β§3.1) π‘š=O(1/πœ–2)#Columns of sketch matrix (Β§3.1) πœ‰= w𝑑0βˆ’wπ‘‘βˆ’1 βˆ₯w𝑑0βˆ’wπ‘‘βˆ’1βˆ₯2 Heuristic vec. for LinearFDA (Β§3.2) Furthermore, let Optimize( w ,B) be the updated model [ 16 ] computed by some optimization algorithm (e.g., SGD, Adam) using the model w, and the batch B of training data. It incorporates the learning rate, loss function and relevant gradients. During step 𝑑, each worker π‘˜first applies the update: w(π‘˜) 𝑑=Optimize(w(π‘˜) π‘‘βˆ’1,B(π‘˜) 𝑑) Moreover, operation AllReduce( w (π‘˜) 𝑑) computes and returns the average model vector [30]: w𝑑=AllReduce(w(π‘˜) 𝑑) Workers synchronize by executing AllReduce( w (π‘˜) 𝑑) , thereby setting w (π‘˜) 𝑑 : =w𝑑 . If synchronization is not performed at step 𝑑 , each worker continues training with its locally updated model. A comprehensive list of the notation used throughout this section is provided in Table 1. Model Variance and FDA. The model variance quantifies the dispersion or spread of worker models around the average model: Var (w𝑑)= 1 𝐾 𝐾 βˆ‘οΈ π‘˜=1w(π‘˜) π‘‘βˆ’wπ‘‘ξ˜ξ˜ξ˜ 2 2(2) This measure provides insight into how closely aligned the workers’ models are at any given time. High variance indicates that the models are widely spread out, essentially drifting apart, leading to a lack of cohesion in the aggregated model. Conversely, a moderate or low variance suggests that the workers’ models are closely aligned, working collectively towards the shared objective. The FDA algorithm (Algorithm 1) is based on the premise that, as long as the variance is below a threshold Θ , synchronization is not needed. Thus, we introduce the Round Invariant (RI): Var (w𝑑)β‰€Ξ˜(3) To preserve the RI, our FDA algorithm maintains (Lines 4-6 of Algorithm 1) at each worker π‘˜ a local (low-dimensional) statevector S (π‘˜) 𝑑 , which is computed based on w (π‘˜) 𝑑 . These state vectors are vital for the subsequent estimation of the model variance, and underpin the two variants of the FDA algorithm (provided in Sections 3.1 and 3.2, respectively). Our estimation techniques begin by performing AllReduce on the states S (π‘˜) 𝑑 , consolidating them into the average state S𝑑 (Line 7). Importantly, this communication step requires significantly less bandwidth and resources than transmitting the full models w(π‘˜) 𝑑. For each FDA variant, we also define a (different) function 𝐻(S𝑑) that overestimates the variance, i.e., it ensures that as long as 𝐻(S𝑑) ≀ Θ then the variance is bounded by Θ . This guarantee is probabilistic for the Sketch-based variant of FDA, and deterministic for its Linear counterpart. Consequently, if 𝐻(S𝑑)>Θthen synchronization is performed (Lines 8-9) β€” the RI invariant cannot be guaranteed. After synchronization, the model variance is zero. Efficiently Monitoring the RI. Estimating model variance efficiently is at the heart of FDA. To this end, we first introduce the local model drift,u(π‘˜) 𝑑, and average drift,u𝑑, defined as follows: u(π‘˜) 𝑑=w(π‘˜) π‘‘βˆ’w𝑑0,u𝑑= 1 𝐾 𝐾 βˆ‘οΈ π‘˜=1 u(π‘˜) 𝑑 Here, w𝑑0 denotes the model vector after the most recent synchronization. Subsequently, the model variance can be written as: Var (w𝑑)= 1 𝐾 𝐾 βˆ‘οΈ π‘˜=1u(π‘˜) π‘‘ξ˜ξ˜ξ˜ 2 2!βˆ’βˆ₯u𝑑βˆ₯2 2(4) Proof. Adding an offset ( βˆ’w𝑑0 ) to each w (π‘˜) 𝑑 does not alter the variance, therefore: Var (w𝑑)=Var ξ˜€wπ‘‘βˆ’w𝑑0=Var (u𝑑)= 1 𝐾 𝐾 βˆ‘οΈ π‘˜=1u(π‘˜) π‘‘βˆ’uπ‘‘ξ˜ξ˜ξ˜ 2 2 = 1 𝐾 𝐾 βˆ‘οΈ π‘˜=1ξ˜’ξ˜ξ˜ξ˜u(π‘˜) π‘‘ξ˜ξ˜ξ˜ 2 2βˆ’2Du(π‘˜) 𝑑,u𝑑E+βˆ₯u𝑑βˆ₯2 2ξ˜“ = 1 𝐾 𝐾 βˆ‘οΈ π‘˜=1u(π‘˜) π‘‘ξ˜ξ˜ξ˜ 2 2!βˆ’2 1 𝐾 𝐾 βˆ‘οΈ π‘˜=1Du(π‘˜) 𝑑,u𝑑E!+ 1 𝐾 𝐾 βˆ‘οΈ π‘˜=1 βˆ₯u𝑑βˆ₯2 2! = 1 𝐾 𝐾 βˆ‘οΈ π‘˜=1u(π‘˜) π‘‘ξ˜ξ˜ξ˜ 2 2!βˆ’2* 1 𝐾 𝐾 βˆ‘οΈ π‘˜=1 u(π‘˜) 𝑑!,u𝑑++βˆ₯u𝑑βˆ₯2 2 = 1 𝐾 𝐾 βˆ‘οΈ π‘˜=1u(π‘˜) π‘‘ξ˜ξ˜ξ˜ 2 2!βˆ’2⟨u𝑑,uπ‘‘βŸ©+βˆ₯u𝑑βˆ₯2 2 = 1 𝐾 𝐾 βˆ‘οΈ π‘˜=1u(π‘˜) π‘‘ξ˜ξ˜ξ˜ 2 2!βˆ’2βˆ₯u𝑑βˆ₯2 2+βˆ₯u𝑑βˆ₯2 2 = 1 𝐾 𝐾 βˆ‘οΈ π‘˜=1u(π‘˜) π‘‘ξ˜ξ˜ξ˜ 2 2!βˆ’βˆ₯u𝑑βˆ₯2 2 β–‘ 414 Algorithm 1 Federated Dynamic Averaging - FDA Require: 𝐾: The number of workers indexed by π‘˜ Require: Θ: The model variance threshold Require: 𝑏: The local mini-batch size 1: Initialize w(π‘˜) 0=w0∈R𝑑 2: for each step 𝑑=1,2, . . . do 3: for each worker π‘˜=1, . . . , 𝐾 in parallel do 4: B(π‘˜) 𝑑←(sample a batch of size 𝑏from Dπ‘˜) 5: w(π‘˜) 𝑑←Optimize(w(π‘˜) π‘‘βˆ’1,B(π‘˜) 𝑑) 6: Update S(π‘˜) 𝑑 7: S𝑑←AllReduce(S(π‘˜) 𝑑) 8: if 𝐻(S𝑑)>Θthen 9: w(π‘˜) 𝑑←AllReduce(w(π‘˜) 𝑑)⊲In-place Conceptually, following Eq (4) , to precisely monitor the variance, we need to calculate two quantities: (1) 1 𝐾Í𝐾 π‘˜=1βˆ₯ u (π‘˜) 𝑑βˆ₯2 2 , and (2) βˆ₯u𝑑βˆ₯2 2 . The first quantity requires an AllReduce operation on the squared norm of the worker drifts, which involves minimal overhead since these values are scalar. In contrast, the second quantity necessitates an AllReduce operation on the worker drifts themselves, which are of model dimension, thus incurring a high communication cost. In fact, this operation is equivalent to synchronization, which is exactly what we aim to avoid in the first place. Thus, it becomes evident that communicationefficient model variance estimation hinges on estimating βˆ₯u𝑑βˆ₯2 2 efficiently. Upcoming sections will detail two techniques for communication efficient variance estimation (which primarily involves estimating βˆ₯u𝑑βˆ₯2 2 ): SketchFDA and LinearFDA. To present them uniformly, we introduce the local state S (π‘˜) 𝑑 , a tensor which contains: (1) the scalar value βˆ₯ u (π‘˜) 𝑑βˆ₯2 2 for precisely calculating the first quantity, and (2) a low-dimensional summary of u (π‘˜) 𝑑 , different for each technique, for estimating the second quantity. For each technique we define an estimation function 𝐻(Β·) that calculates the current variance estimate from average state S𝑑=1 𝐾Í𝐾 π‘˜=1 S (π‘˜) 𝑑 (obtained via AllReduce). 3.1 SketchFDA: Sketch-based Estimation An optimal estimator for βˆ₯u𝑑βˆ₯2 2 can be obtained through the utilization and properties of AMS sketches, as detailed in [ 8 ]. An AMS sketch of a vector v∈R𝑑is an π‘™Γ—π‘šreal matrix: sk (v)=ξ˜‚πœ“1πœ“2. . . πœ“π‘™ξ˜ƒβŠ€βˆˆRπ‘™Γ—π‘š, 𝑙 Β·π‘šβ‰ͺ𝑑 An estimate for squared-norm βˆ₯vβˆ₯2 2is provided by the formula M2(sk(v))=median βˆ₯πœ“π‘–βˆ₯2 2, 𝑖 =1, . . . ,π‘™ξ˜‰ The quality of estimation depends on the size of the sketch. For chosen πœ–, 𝛿 > 0, where sketch dimensions are given by 𝑙=O(log 1/𝛿) and π‘š=Oξ˜€1/πœ–2 , we have the following probabilistic guarantee: with confidence at least 1βˆ’π›Ώ, M2(sk(v)) ∈ (1Β±πœ–)βˆ₯vβˆ₯2 2 Notably, observe that the accuracy ( πœ– ) and confidence (1 βˆ’π›Ώ ) only depend on the size of the sketch and not on the dimensionality of vector v. Two crucial properties of the AMS sketch are that (a) it is a linear transformation, i.e., for 𝛼1, 𝛼2∈Rand v1,v2∈R𝑑, sk(𝛼1v1+𝛼2v2)=𝛼1sk(v1) + 𝛼2sk(v2) and (b) can be computed efficiently in time 𝑂(𝑙·𝑑). In the SketchFDA approach, the salient idea is to employ AMS sketches sk( u (π‘˜) 𝑑) ∈ Rπ‘™Γ—π‘š as a low-dimensional representation of the local drifts u(π‘˜) 𝑑. Theorem 3.1. Let 𝑙=O(log 1 𝛿) and π‘š=O( 1 πœ–2) . Define the local state as S(π‘˜) 𝑑=ξ˜’ξ˜ξ˜ξ˜u(π‘˜) π‘‘ξ˜ξ˜ξ˜ 2 2,sk u(π‘˜) π‘‘ξ˜‘ξ˜“βˆˆRΓ—Rπ‘™Γ—π‘š and the approximation function as 𝐻Sπ‘‘ξ˜‘= 1 πΎβˆ‘οΈ π‘˜ξ˜ξ˜ξ˜u(π‘˜) π‘‘ξ˜ξ˜ξ˜ 2 2βˆ’1 1+πœ–M2 1 𝐾 𝐾 βˆ‘οΈ π‘˜=1 sk u(π‘˜) π‘‘ξ˜‘!. Then, the condition 𝐻(S𝑑) ≀ Θ implies Var (w𝑑)β‰€Ξ˜ with probability at least (1βˆ’π›Ώ). Proof. 𝐻Sπ‘‘ξ˜‘= 1 𝐾 𝐾 βˆ‘οΈ π‘˜=1u(π‘˜) π‘‘ξ˜ξ˜ξ˜ 2 2βˆ’1 1+πœ–M2 1 𝐾 𝐾 βˆ‘οΈ 𝑖=1 sk u(π‘˜) π‘‘ξ˜‘! (lin.) = 1 𝐾 𝐾 βˆ‘οΈ π‘˜=1u(π‘˜) π‘‘ξ˜ξ˜ξ˜ 2 2βˆ’1 1+πœ–M2 sk 1 𝐾 𝐾 βˆ‘οΈ 𝑖=1 u(π‘˜) 𝑑!! = 1 𝐾 𝐾 βˆ‘οΈ π‘˜=1u(π‘˜) π‘‘ξ˜ξ˜ξ˜ 2 2βˆ’1 1+πœ–M2(sk (u𝑑)) (πœ–-err.) β‰₯1 𝐾 𝐾 βˆ‘οΈ π‘˜=1u(π‘˜) π‘‘ξ˜ξ˜ξ˜ 2 2βˆ’βˆ₯u𝑑βˆ₯2 2with prob. at least (1βˆ’π›Ώ) =Var (w𝑑) We proved that 𝐻(S𝑑) β‰₯ Var (w𝑑) with probability at least (1 βˆ’π›Ώ ), i.e., we overestimate the model variance with probability at least (1βˆ’π›Ώ), completing the proof. β–‘ In Section 3.3, we discuss the empirical basis for choosing the values of 𝑙 and π‘š , and how they practically impact the quality of the sketch approximation. 3.2 LinearFDA: Linear Approximation Although AMS sketches provide good estimates for variance, their dimension is in the several hundreds, and the communication cost of AllReduce on sketches, performed at each step, may be non-negligible. Therefore, we also introduce a low-cost, ad-hoc estimation variant. In this approach, instead of an AMS sketch, each local state contains the scalar value βŸ¨πœ‰ , u (π‘˜) π‘‘βŸ© ∈ R , where πœ‰βˆˆR𝑑 is a unit vector, known to all workers. Theorem 3.2. Define the local state as S(π‘˜) 𝑑=ξ˜’ξ˜ξ˜ξ˜u(π‘˜) π‘‘ξ˜ξ˜ξ˜ 2 2,Dπœ‰ , u(π‘˜) 𝑑Eξ˜“βˆˆRΓ—R,βˆ₯πœ‰βˆ₯2=1 and the approximation function as 𝐻Sπ‘‘ξ˜‘= 1 𝐾 𝐾 βˆ‘οΈ π‘˜=1u(π‘˜) π‘‘ξ˜ξ˜ξ˜ 2 2βˆ’ξ˜Œξ˜Œξ˜Œξ˜Œξ˜Œ 1 𝐾 𝐾 βˆ‘οΈ 𝑖=1Dπœ‰ , u(π‘˜) 𝑑E 2 Then, the condition 𝐻(S𝑑) ≀ Θimplies Var (w𝑑)β‰€Ξ˜. 415 Proof. 𝐻Sπ‘‘ξ˜‘= 1 𝐾 𝐾 βˆ‘οΈ π‘˜=1u(π‘˜) π‘‘ξ˜ξ˜ξ˜ 2 2βˆ’ξ˜Œξ˜Œξ˜Œξ˜Œξ˜Œ 1 𝐾 𝐾 βˆ‘οΈ 𝑖=1Dπœ‰ , u(π‘˜) 𝑑E 2 = 1 𝐾 𝐾 βˆ‘οΈ π‘˜=1u(π‘˜) π‘‘ξ˜ξ˜ξ˜ 2 2βˆ’ξ˜Œξ˜Œξ˜Œξ˜Œξ˜Œ*πœ‰ , 1 𝐾 𝐾 βˆ‘οΈ 𝑖=1 u(π‘˜) 𝑑+ 2 = 1 𝐾 𝐾 βˆ‘οΈ π‘˜=1u(π‘˜) π‘‘ξ˜ξ˜ξ˜ 2 2βˆ’|βŸ¨πœ‰ , uπ‘‘βŸ©|2 β‰₯1 𝐾 𝐾 βˆ‘οΈ π‘˜=1u(π‘˜) π‘‘ξ˜ξ˜ξ˜ 2 2βˆ’βˆ₯πœ‰βˆ₯2 2βˆ₯u𝑑βˆ₯2 2 = 1 𝐾 𝐾 βˆ‘οΈ π‘˜=1u(π‘˜) π‘‘ξ˜ξ˜ξ˜ 2 2βˆ’βˆ₯u𝑑βˆ₯2 2 =Var (w𝑑) We proved that 𝐻(S𝑑) β‰₯ Var (w𝑑), i.e., we always overestimate the model variance, completing the proof. β–‘ An arbitrary choice of πœ‰ (e.g., a random vector) is likely to estimate βˆ₯u𝑑βˆ₯2 2poorly; if πœ‰is uncorrelated to u𝑑, then |βŸ¨πœ‰ , uπ‘‘βŸ©|2 will likely be close to zero. A heuristic choice that might be correlated to u𝑑 is the (normalized) value of u𝑑0 , the global drift vector right at the time of last synchronization. All nodes can compute it independently without extra communication, if they take the difference of the models of the last two synchronizations: πœ‰= u𝑑0 u𝑑02 = w𝑑0βˆ’wπ‘‘βˆ’1 w𝑑0βˆ’wπ‘‘βˆ’12 3.3 Discussion FDA: Intuition. The main intuition for FDA is summarized in making the decision to synchronize dynamic, based on model variance during training. This metric is designed to capture the collective state of the training process. In what follows, we provide intuition on why this is the case. It is important to remember that the global model w𝑑 and, by extension, the global drift u𝑑 , are ultimately what we care about and evaluate. Model variance, as defined in Equation (4) , is the difference between the average of the squared local drifts 1 𝐾Íβˆ₯ u (π‘˜) 𝑑βˆ₯2 2 and the squared global drift βˆ₯u𝑑βˆ₯2 2 . The first term reflects how far the individual worker models have moved–essentially, how much each worker has learned. The second term indicates how much of this learning is retained in the global model after aggregation. The interplay between these two quantities is crucial. For example, when the local drifts are high but the global drift is low, the variance increases, signaling the need for synchronization. This scenario suggests that while individual workers have made significant progress (as indicated by high local drifts), this progress is not being effectively captured in the global model (indicated by the low global drift). In other words, the worker models have moved significantly, but the global model has remained relatively stationary in this high-dimensional space. This misalignment indicates that training is no longer progressing optimally, as the workers are moving towards disparate and conflicting local minima, making it crucial to synchronize and realign them. Conversely, when both the local and global drifts are either low or high, synchronization is not necessary, and the variance naturally remains low. ( ( ) ) , , Local StateLocal Drift LΙͺɴᴇᴀʀFDA Sα΄‹α΄‡α΄›α΄„ΚœFDA Figure 2: SketchFDA &LinearFDA: Local State structure. Neither the average of the local drifts nor the global drift alone provides a complete picture of the collective training progress. Relying solely on one or the other would lead to suboptimal synchronization decisions and likely prove ineffective. In FDA, it is the relationship between these quantities, as captured by the model variance, that offers valuable insights and guides the crucial decision of when to synchronize. SketchFDA vs. LinearFDA: Both methods send the squared norm of the drift βˆ₯ u (π‘˜) 𝑑βˆ₯2 2 , but differ in the additional accompanying lower-dimensional representation they transmit (Figure 2): (1) SketchFDA: An AMS sketch of the local drift. (2) LinearFDA: The dot product of a vector and the local drift. The key difference between these two variants lies in the fidelity of approximation of the model variance. While both methods conservatively overestimate the variance, SketchFDA provides a provably accurate estimation, which is expected to lead to fewer synchronizations. LinearFDA requires less computational effort and bandwidth to create and communicate the local states, but may overestimate variance by too much, causing unnecessary synchronizations. SketchFDA: Choice of 𝑙 and π‘š . We empirically measured the approximation achieved with sketch dimensions of 𝑙= 5rows and π‘š= 250 columns (as defined in Section 3.1): these settings yield an error bound of πœ–β‰ˆ 6% and a probabilistic confidence of ( 1 βˆ’π›Ώ) β‰ˆ 95%. Based on our experiments, we have adopted these values in our experiments and recommend them. Using these values, the byte-size of a sketch is π‘™Β·π‘šΒ· 4 bytes = 5 kB , significantly smaller than the size of all our models. Sketches of smaller size could be used, albeit weakening the approximation of the variance. However, given that LinearFDA similarly weakens approximation and avoids using AMS sketches, in the interest of space we do not explore varying AMS sketch sizes in this paper. FDA: Asynchronous Operation. As mentioned in Section 2, FDA can be readily modified to operate asynchronously. In this setup, one worker-node acts as a coordinator, aggregating local states and determining whether synchronization is needed each time a local state is received. This decision is based on the most recent local states from all workers. It is important to note that, since local states are small in size, asynchronous operation is unlikely to alleviate bandwidth issues. The primary advantage is that it allows training to continue even in the presence of stragglers. Asynchronous operation might also be beneficial in rare cases where the overhead of initializing communication dominates the actual transmission time. 416 4 EXPERIMENTS 4.1 Setup Table 2 provides a comprehensive overview of our experiments. For each experiment, we detail the Neural Network (NN) architecture, its parameter count ( 𝑑 ), and the dataset used for training. The table also specifies key hyper-parameters: the batch size ( 𝑏 ), the number of workers ( 𝐾 ), and the FDA-specific variance threshold ( Θ ). Additionally, we indicate the chosen optimizer (as detailed in Section 3) and the training algorithms employed for each configuration. Platform. We employ TensorFlow [ 1 ], integrated with Keras [ 7 ], as the platform for conducting our experiments. We used TensorFlow to implement our FDA variants and all competitive algorithms. All relevant code, figures, and data of this study are available in https://github.com/miketheologitis/FedL-Sync-FDA. Hardware & Infrastructure. We conducted our experiments on the ARIS High performance computing (HPC) environment 2 , utilizing a cluster of 44 GPU-accelerated worker-nodes. Each worker is equipped with two NVIDIA Tesla K40m GPUs and interconnected via an InfiniBand FDR14 network, providing up to 56 GB/s of bandwidth. Crucially, our evaluation remains agnostic to the underlying infrastructure of the specific workers. Datasets & Models. The core experiments involve training Convolutional Neural Networks (CNNs) of varying sizes and complexities on two datasets: MNIST [ 12 ] and CIFAR-10 [ 27 ]. For the MNIST dataset, we employ LeNet-5 [ 29 ], composed of approximately 62 thousand parameters, and a modified version of VGG16 [ 48 ], denoted as VGG16*, consisting of 2.6 million parameters. VGG16* was specifically adapted for the MNIST dataset, a less demanding learning problem compared to ImageNet [ 44 ], for which VGG16 was designed. In VGG16*, we omitted the 512channel convolutional blocks and downscaled the final two fully connected (FC) layers from 4096 to 512 units each. Both models use Glorot uniform initialization [ 15 ]. For CIFAR-10, we utilize DenseNet121 and DenseNet201 [ 22 ], as implemented in Keras [ 7 ], with the addition of dropout regularization layers at rate 0.2 and weight decay of 10 βˆ’4 , as prescribed in [ 22 ]. The DenseNet121 and DenseNet201 models have 6.9 million and 18 million parameters, respectively, and are both initialized with He normal [19]. Lastly, we explore a transfer learning scenario on the dataset CIFAR-100 [ 27 ], a choice reflecting the DL community’s growing preference of using pre-trained models in such downstream tasks [ 18 ]. For example, a pre-trained visual transformer (ViT) on ImageNet, transferred to classify CIFAR-100, is currently on par with the state-of-the-art results for this task [ 13 ]. We adopt this exact transfer learning scenario, leveraging the more powerful ConvNeXtLarge model, pre-trained on ImageNet, with 198 million parameters [ 7 , 33 ]. Following the feature extraction step [ 16 ], the testing accuracy on CIFAR-100 stands at 60%. Subsequently, we employ and evaluate our FDA algorithms in the arduous fine-tuning stage, where the entirety of the model is trained [ 39 ]. Algorithms. We consider five distributed deep learning algorithms: LinearFDA,SketchFDA,Synchronous 3 ,FedAdam [ 42 ], and FedAvgM [ 21 ]; the first three are standard in all experiments. Depending on the local optimizer, Adam [ 26 ] or SGD with Nesterov momentum (SGD-NM) [ 52 ], we also include their 2https://www.hpc.grnet.gr/en/hardware-2/ 3 The name was derived from the Bulk Synchronous Parallel approach; can be understood as a special case of the FDA Algorithm 1 where Θis set to zero. communication-efficient federated counterparts FedAdam or FedAvgM, respectively. Evaluation Methodology. Comparing DDL algorithms is not straightforward. For example, comparing DDL algorithms based on the average cost of a training epoch can be misleading, as it does not consider the effects on the trained model’s quality. To achieve a comprehensive performance assessment of FDA, we define a training run as the process of executing the DDL algorithm under evaluation, on (a) a specific DL model and training dataset, and (b) until a final epoch in which the trained model achieves a specific testing accuracy (termed as Accuracy Target in figures). Based on this definition, we focus on two performance metrics: (1) Communication cost, which is the total data (in bytes) transmitted by all workers. Notably, communication cost is unaffected by the training data volume since only model updates (when synchronizing) and local states (at each step), but not training data, are transmitted. Thus, the communication cost mainly depends on the complexity (number of parameters) of the used model. Translating the communication cost to wall-clock time (i.e., the total time required for the computation and communication of the DDL) depends on the network infrastructure connecting the workers and on the overhead of establishing and initializing communication. Its impact is larger in FL scenarios, where workers often use slower Wi-Fi connections. (2) Computation cost, which is the number of mini-batch steps (termed as In-Parallel Learning Steps in figures) performed by each worker. Translating this cost to wall-clock time is determined by the mini-batch size and the computational resources of the worker-nodes. Its impact is larger for workers with lower computational resources. Hyper-Parameters & Optimizers. Hyper-parameters unique to each training dataset and model are detailed in Table 2; Θ is pertinent to FDA algorithms and not applicable to others. Notably, a guideline for setting the parameter Θ is provided in Section 4.3. For experiments involving FedAvgM and FedAdam, we use 𝐸= 1 local epochs, following [ 42 ]. For experiments with LeNet-5 and VGG16*, local optimization employs Adam, using the default settings as per [ 26 ]. In these cases, FedAdam also adheres to the default settings for both local and server optimization [ 7 , 42 ]. For DenseNet121 and DenseNet201, local optimization is performed using SGD with Nesterov momentum (SGD-NM), setting the momentum parameter at 0 . 9and learning rate at 0 . 1[ 22 ]. For FedAvgM, local optimization is conducted with default settings [ 7 , 21 ], while server optimization employs SGD with momentum, setting the momentum parameter and learning rate to 0 . 9and 0 . 316, respectively [ 42 ]. Lastly, for the transfer learning experiments, local optimization leverages AdamW [ 34 ], with the hyper-parameters used for fine-tuning ConvNeXtLarge in the original study [33]. Data Distribution. In all experiments, the training dataset is divided into approximately equal parts among the workers. To assess the impact of data heterogeneity, we explore three scenarios: (1) IID β€” Independent and identically distributed. (2) Non-IID: 𝑋 %β€” A portion 𝑋 %of the dataset is sorted by label and sequentially allocated to workers, with the remainder distributed in an IID fashion. 417 Table 2: Summary of Experiments Hyper-Parameters Training NN d Dataset Θb K Optimizer Algorithms LeNet-5 62K MNIST {0.5,1,1.5,2,3,5,7}32 { 5, 10, ..., 60 } Adam FDA,Synchronous,FedAdam VGG16* 2.6M MNIST {20,25,30,50,75,90,100}32 { 5, 10, ..., 60 } Adam FDA,Synchronous,FedAdam DenseNet121 6.9M CIFAR-10 {200,250,275,300,325,350,400}32 { 5, 10, ..., 30 } SGD-NM FDA,Synchronous,FedAvgM DenseNet201 18M CIFAR-10 {350,500,600,700,800,850,900}32 { 5, 10, ..., 30 } SGD-NM FDA,Synchronous,FedAvgM (fine-tuning) ConvNeXtLarge 198M CIFAR-100 {25,50,100,150}32 { 3, 5 } AdamW FDA,Synchronous 10 1100101102 Communication (GB) 104 105 In-Parallel Learning Steps IID, Accuracy Target: 0.985 LinearFDA SketchFDA FedAdam Synchronous 10 1100101102 Communication (GB) 104 105 In-Parallel Learning Steps Non-IID: Label " 0 ", Accuracy Target: 0.985 LinearFDA SketchFDA FedAdam Synchronous 10 1100101102 Communication (GB) 104 105 In-Parallel Learning Steps Non-IID: 60% , Accuracy Target: 0.985 LinearFDA SketchFDA FedAdam Synchronous Figure 3: LeNet-5 on MNIST. At Non-IID: Label "0", the samples of Label "0" are assigned to few workers. At Non-IID: 60%, 60% of the dataset is sorted and allocated to workers, causing some workers to receive many samples from the same label (3) Non-IID: Label π‘Œ β€” All samples from label π‘Œ are assigned to a few workers, while the rest are distributed in an IID manner. 4.2 Main Findings The main findings of our experimental analyses are: (1) LinearFDA and SketchFDA outperform the Synchronous,FedAdam and FedAvgM techniques (their use depends on the local optimizer choice) by 1-2 orders of magnitude in communication, while maintaining equivalent model performance. (2) LinearFDA and SketchFDA also significantly outperform the FedAdam and FedAvgM techniques in terms of computation. (3) The performance of LinearFDA and SketchFDA is comparable in most experiments. SketchFDA provides a more accurate estimator of the variance and leads to fewer synchronizations than LinearFDA, but has a larger communication overhead for its local state (a sketch, compared to two numbers). SketchFDA significantly outperforms LinearFDA at the transfer learning scenario. (4) The FDA variants remain robust at various data heterogeneity settings, maintaining comparable performance to the IID case. 4.3 Results Due to the extensive set of unique experiments (over 1000), as detailed in Table 2, we leverage Kernel Density Estimation (KDE) plots [ 62 ] to visualize the bivariate distribution of computation and communication costs incurred by each strategy for attaining the Accuracy Target. These KDE plots provide a high-level overview of the cost trade-off for training accurate models. The varying levels of opacity in the filled areas of the KDE plots represent the density of the underlying data points: higher opacity indicates areas with a greater concentration of data, whereas lower opacity signifies less dense areas. As an illustrative example, Figure 3 depicts the strategies’ bivariate distribution for the LeNet-5 model trained on MNIST with different data heterogeneity setups. In these plots, the SketchFDA distribution is generated from experiments across all hyperparameter combinations ( Θ and 𝐾 in Table 2) that attained the Accuracy Target of 0 . 985. The observed high variance in the method’s distribution stems from the varying 𝐾 and Θ values. In subsequent subsections, we elucidate how these hyper-parameters influence the communication and computation costs. FDA balances Communication vs. Computation. DDL algorithms face a fundamental challenge: balancing the competing demands of computation and communication. Frequent communication accelerates convergence and potentially improves model performance, but incurs higher network overhead, an overhead that may be prohibitive when workers communicate through lower speed connections. Conversely, reducing communication saves bandwidth but risks hindering, or even stalling, convergence. Traditional DDL approaches, like Synchronous, require synchronizing model parameters after every learning step, leading to significant communication overhead but facilitating faster convergence (lower computation cost). This is evident in Figures 3, 4, 5, and 6 (where Synchronous appears in the bottom right β€” low computation, very high communication). Conversely, Federated Optimization (FedOpt) methods [ 42 ] are designed to be communication-efficient, reducing communication between 418 100101102103 Communication (GB) 103 104 In-Parallel Learning Steps IID, Accuracy Target: 0.994 LinearFDA SketchFDA FedAdam Synchronous 101102103 Communication (GB) 104 105 In-Parallel Learning Steps IID, Accuracy Target: 0.995 LinearFDA SketchFDA FedAdam Synchronous 100101102103 Communication (GB) 103 104 In-Parallel Learning Steps Non-IID: Label " 0 ", Accuracy Target: 0.994 LinearFDA SketchFDA FedAdam Synchronous 101102103 Communication (GB) 104 105 In-Parallel Learning Steps Non-IID: Label " 0 ", Accuracy Target: 0.995 LinearFDA SketchFDA FedAdam Synchronous 101102103 Communication (GB) 103 104 In-Parallel Learning Steps Non-IID: Label " 8 ", Accuracy Target: 0.994 LinearFDA SketchFDA FedAdam Synchronous 101102103 Communication (GB) 104 105 In-Parallel Learning Steps Non-IID: Label " 8 ", Accuracy Target: 0.995 LinearFDA SketchFDA FedAdam Synchronous Figure 4: VGG16* on MNIST 101102103 Communication (GB) 103 104 In-Parallel Learning Steps IID, Accuracy Target: 0.78 LinearFDA SketchFDA FedAvgM Synchronous 101102103104 Communication (GB) 104 105 In-Parallel Learning Steps IID, Accuracy Target: 0.81 LinearFDA SketchFDA FedAvgM Synchronous Figure 5: DenseNet121 on CIFAR-10 101102103104 Communication (GB) 104 In-Parallel Learning Steps IID, Accuracy Target: 0.78 LinearFDA SketchFDA FedAvgM Synchronous 102103104 Communication (GB) 104 In-Parallel Learning Steps IID, Accuracy Target: 0.8 LinearFDA SketchFDA FedAvgM Synchronous Figure 6: DenseNet201 on CIFAR-10 devices (workers) at the expense of increased local computation. Indeed, as shown in Figures 3-6, FedAvgM and FedAdam reduce communication by orders of magnitude but at the price of a corresponding increase in computation. Our two proposed FDA strategies achieve the best of both worlds: the low computation cost of traditional methods and the communication efficiency of FedOpt approaches, as seen in Figures 3, 4, 5, and 6. In fact, they significantly outperform FedAvgM and FedAdam in their element, that is, communication-efficiency. Across all experiments, 419