Decentralized Deep Learning with Arbitrary Communication Compression

Anastasia Koloskova, Tao Lin, Sebastian U. Stich, Martin Jaggi

Introduction

Distributed machine learning—i.e. the training of machine learning models using distributed optimization algorithms—has recently enabled many successful applications in research and industry. Such methods offer two of the key success factors: 1) computational scalability by leveraging the simultaneous computational power of many devices, and 2) data-locality, the ability to perform joint training while keeping each part of the training data local to each participating device. Recent theoretical results indicate that decentralized schemes can be as efficient as the centralized approaches, at least when considering convergence of training loss vs. iterations (Scaman et al., 2017; 2018; Lian et al., 2017; Tang et al., 2018; Koloskova et al., 2019; Assran et al., 2019).

Gradient compression techniques have been proposed for the standard distributed training case (Alistarh et al., 2017; Wen et al., 2017; Lin et al., 2018; Wangni et al., 2018; Stich et al., 2018), to reduce the amount of data that has to be sent over each communication link in the network. For decentralized training of deep neural networks, Tang et al. (2018) introduce two algorithms (DCD, ECD) which allow for communication compression. However, both these algorithms are restrictive with respect to the used compression operators, only allowing for unbiased compressors and—more significantly—so far not supporting arbitrarily high compression ratios. We here study Choco-SGD—recently introduced for convex problems only (Koloskova et al., 2019)—which overcomes these constraints.

For the evaluation of our algorithm we in particular focus on the generalization performance (on the test-set) on standard machine learning benchmarks, hereby departing from previous work such as e.g. (Tang et al., 2018; Wang et al., 2019; Tang et al., 2019; Reisizadeh et al., 2019) that mostly considered training performance (on the train-set). We study two different scenarios: firstly, (i) training on a challenging peer-to-peer setting, where the training data is distributed over the training devices (and not allowed to move), similar to the federated learning setting (McMahan et al., 2017; Kairouz et al., 2019). We are again able to show speed-ups for Choco-SGD over the decentralized baseline (Lian et al., 2017) with much less communication overhead. Secondly, (ii) training in a datacenter setting, where decentralized communication patterns allow better scalability than centralized approaches. For this setting we show that communication efficient Choco-SGD can improve time-to-accuracy on large tasks, such as e.g. ImageNet training. However, when investigating the scaling of decentralized algorithms to larger number of nodes we observe that (all) decentralized schemes encounter difficulties and often do not reach the same (test and train) performance as centralized schemes. As these findings point out some deficiencies of current decentralized training schemes (and are not particular to our scheme) we think that reporting these results is a helpful contribution to the community to spur further research on decentralized training schemes that scale to large number of peers.

On the theory side, we are the first to show that Choco-SGD converges at rate \smash{\mathcal{O}\big{(}\nicefrac{{1}}{{\sqrt{nT}}}+\nicefrac{{1}}{{(\rho^{2}\delta T)^{2/3}}}\bigr{)}} on non-convex smooth functions, where nn denotes the number of nodes, TT the number of iterations, ρ\rho the spectral gap of the mixing matrix and δ\delta the compression ratio. The main term, \smash{\mathcal{O}\big{(}\nicefrac{{1}}{{\sqrt{nT}}}\big{)}}, matches with the centralized baselines with exact communication and shows a linear speedup in the number of workers nn. Both ρ\rho and δ\delta only affect the asymptotically smaller second term.

On the practical side, we present a version of Choco-SGD with momentum and analyze its practical performance on two relevant scenarios:

for on-device training over a realistic peer-to-peer social network, where lowering the bandwidth requirements of joint training is especially impactful

in a datacenter setting for computational scalability of training deep learning models for resource efficiency and improved time-to-accuracy

Lastly, we systematically investigate performance of the decentralized schemes when scaling to larger number of nodes and we point out some (shared) difficulties encountered by current decentralized learning approaches.

Related Work

For the training in communication restricted settings a variety of methods have been proposed. For instance, decentralized schemes (Lian et al., 2017; Nedić et al., 2018; Koloskova et al., 2019), gradient compression (Seide et al., 2014; Strom, 2015; Alistarh et al., 2017; Wen et al., 2017; Lin et al., 2018; Wangni et al., 2018; Bernstein et al., 2018; Lin et al., 2018; Alistarh et al., 2018; Stich et al., 2018; Karimireddy et al., 2019), asynchronous methods (Recht et al., 2011; Assran et al., 2019), coordinate updates Nesterov (2012); Richtárik & Takáč (2016); Stich et al. (2017a; b); He et al. (2018), or performing multiple local SGD steps before averaging (Zhang et al., 2016; McMahan et al., 2017; Stich, 2019; Lin et al., 2020). This especially covers learning over decentralized data, as extensively studied in the federated learning literature for the centralized algorithms (McMahan et al., 2016; Kairouz et al., 2019). In this paper we advocate for combining decentralized SGD schemes with gradient compression.

Decentralized SGD. We in particular focus on approaches based on gossip averaging (Kempe et al., 2003; Xiao & Boyd, 2004; Boyd et al., 2006) whose convergence rate typically depends on the spectral gap ρ≥0\rho\geq 0 of the mixing matrix (Xiao & Boyd, 2004). Lian et al. (2017) combine SGD with gossip averaging and show that the leading term in the convergence rate \smash{\mathcal{O}\big{(}\nicefrac{{1}}{{\sqrt{nT}}}\bigr{)}} is consistent with the convergence of the centralized mini-batch SGD (Dekel et al., 2012) and the spectral gap only affects the asymptotically smaller terms. Similar results have been observed very recently for related schemes (Scaman et al., 2017; 2018; Koloskova et al., 2019; Yu et al., 2019).

Quantization. Communication compression with quantization has been popularized in the deep learning community by the reported successes in (Seide et al., 2014; Strom, 2015). Theoretical guarantees were first established for schemes with unbiased compression (Alistarh et al., 2017; Wen et al., 2017; Wangni et al., 2018) but soon extended to biased compression (Bernstein et al., 2018) as well. Schemes with error correction work often best in practice and give the best theoretical gurantees (Lin et al., 2018; Alistarh et al., 2018; Stich et al., 2018; Karimireddy et al., 2019; Stich & Karimireddy, 2019). Recently, also proximal updates and variance reduction have been studied in combination with quantized updates (Mishchenko et al., 2019; Horváth et al., 2019).

Decentralized Optimization with Quantization. It has been observed that gossip averaging can diverge (or not converge to the correct solution) in the presence of quantization noise (Xiao et al., 2005; Carli et al., 2007; Nedić et al., 2008; Dimakis et al., 2010; Carli et al., 2010b; Yuan et al., 2012). Reisizadeh et al. (2018) propose an algorithm that can still converge, though at a slower rate than the exact scheme. Another line of work proposed adaptive schemes (with increasing compression accuracy) that converge at the expense of higher communication cost (Carli et al., 2010a; Doan et al., 2018; Berahas et al., 2019). For deep learning applications, Tang et al. (2018) proposed the DCD and ECD algorithms that converge at the same rate as the centralized baseline though only for constant compression ratio. The Choco-SGD algorithm that we consider in this work can deal with arbitrary high compression, and has been introduced in (Koloskova et al., 2019) but only been analyzed for convex functions. For non-convex functions we show a rate of \smash{\mathcal{O}\big{(}\nicefrac{{1}}{{\sqrt{nT}}}+\nicefrac{{1}}{{(\rho^{2}\delta T)^{\frac{2}{3}}}}\bigr{)}}, where δ>0\delta>0 measures the compression quality. Simultaneous work of Tang et al. (2019) introduced DeepSqueeze, an alternative method which also converges with arbitrary compression ratio. In our experiments, under the same amount of tuning, Choco-SGD achieves higher test accuracy.

Choco-SGD

In this section we formally introduce the decentralized optimization problem, compression operators, and the gossip-based stochastic optimization algorithm Choco-SGD from (Koloskova et al., 2019).

Distributed Setup. We consider optimization problems distributed across nn nodes of the form

Communication. Every device is only allowed to communicate with its local neighbours defined by the network topology, given as a weighted graph G=([n],E)G=([n],E), with edges EE representing the communication links along which messages (e.g. model updates) can be exchanged. We assign a positive weight wijw_{ij} to every edge (wij=0w_{ij}=0 for disconnected nodes {i,j}∉E\{i,j\}\notin E).

We assume that W∈n×nW\in^{n\times n}, (W)ij=wij(W)_{ij}=w_{ij} is a symmetric (W=W⊤W=W^{\top}) doubly stochastic (W1=1W\mathbf{1}=\mathbf{1},1⊤W=1⊤\mathbf{1}^{\top}W=\mathbf{1}^{\top}) matrix with eigenvalues 1=∣λ1(W)∣>∣λ2(W)∣≥⋯≥∣λn(W)∣1=|\lambda_{1}(W)|>|\lambda_{2}(W)|\geq\dots\geq|\lambda_{n}(W)| and spectral gap ρ:=1−∣λ2(W)∣∈(0,1] .\rho:=1-|\lambda_{2}(W)|\in(0,1]\,.

In our experiments we set the weights based on the local node degrees: wij=max⁡{deg⁡(i),deg⁡(j)}−1w_{ij}=\max\{\deg(i),\deg(j)\}^{-1} for {i,j}∈E\{i,j\}\in E. This will not only guarantee ρ>0\rho>0 but these weights can easily be computed in a local fashion on each node (Xiao & Boyd, 2004).

Compression. We aim to only transmit compressed (e.g. quantized or sparsified) messages. We formalized this through the notion of compression operators that was e.g. also used in (Tang et al., 2018; Stich et al., 2018).

In contrast to the quantization operators used in e.g. (Alistarh et al., 2017; Horváth et al., 2019), compression operators defined as in (2) are not required to be unbiased and therefore supports a larger class of compression operators. Some examples can be found in (Koloskova et al., 2019) and we further discuss specific compression schemes in Section 5.

Algorithm. Choco-SGD is summarized in Algorithm 1.

From an implementation aspect, it is worth highlighting that the communication part \raisebox{-.9pt} {1}⃝ and the gradient computation part \raisebox{-.9pt} {2}⃝ can both be executed in parallel because they are independent. Moreover, each node only needs to store 3 vectors at most, independent of the number of neighbors (this might not be obvious from the notation used here for additinal clarity, for further details c.f. (Koloskova et al., 2019)). We further propose a momentum-version of Choco-SGD in Algorithm 2 (see Section D for further details).

Convergence of Choco-SGD on Smooth Non-Convex Problems

As the first main contribution, we extend the analysis of Choco-SGD to non-convex problems. For this we make the following technical assumptions:

Under Assumptions 1–2 there exists a constant stepsize η\eta and the consensus stepsize from (Koloskova et al., 2019), γ:=ρ2δ16ρ+ρ2+4β2+2ρβ2−8ρδ\gamma:=\frac{\rho^{2}\delta}{16\rho+\rho^{2}+4\beta^{2}+2\rho\beta^{2}-8\rho\delta} with β=∥I−W∥2∈\beta=\left\lVert I-W\right\rVert_{2}\in, such that the averaged iterates x‾(t):=1n∑i=1nxi(t)\overline{\mathbf{x}}^{(t)}:=\tfrac{1}{n}\sum_{i=1}^{n}\mathbf{x}_{i}^{(t)} of Algorithm 1 satisfy:

where c:=ρ2δ82c:=\tfrac{\rho^{2}\delta}{82} denotes the convergence rate of the underlying consensus averaging scheme of (Koloskova et al., 2019), F0:=f(x‾(0))−f⋆F_{0}:=f(\overline{\mathbf{x}}^{(0)})-f^{\star}.

This result shows that Choco-SGD converges as \smash{\mathcal{O}\big{(}\nicefrac{{1}}{{\sqrt{nT}}}+\nicefrac{{1}}{{(\rho^{2}\delta T)^{2/3}}}\bigr{)}}. The first term shows a linear speed-up compared to SGD on a single node, while compression and graph topology affect only the higher order second term. In the special case when exact averaging without compression is used (δ=1)\delta=1) , then c=ρc=\rho and the rate improves to \smash{\mathcal{O}\big{(}\nicefrac{{1}}{{\sqrt{nT}}}+\nicefrac{{1}}{{(\rho T)^{2/3}}}\bigr{)}}, recovering the rate in (Wang & Joshi, 2018). This upper bound improves slightly over (Lian et al., 2017) that shows \smash{\mathcal{O}\big{(}\nicefrac{{1}}{{\sqrt{nT}}}+\nicefrac{{n}}{{(n\rho T)^{2/3}}}\bigr{)}}.Theorem 1 of Lian et al. (2017) and stepsize tuned with Lemma A.4. For the proofs and convergence of the individual iterates xi\mathbf{x}_{i} we refer to Appendix A.

Comparison to Baselines for Various Compression Schemes

In this section we experimentally compare Choco-SGD to the relevant baselines for a selection of commonly used compression operators. For the experiments we further leverage momentum in all implemented algorithms. The newly developed momentum version of Choco-SGD is given as Algorithm 2.

In order to match the setting in (Tang et al., 2018) for our first set of experiments, we use a ring topology with n=8n=8 nodes and train the ResNet20 architecture (He et al., 2016) on the Cifar10 dataset (50K/10K training/test samples) (Krizhevsky, 2012). We randomly split the training data between workers and shuffle it after every epoch, following standard procedure as e.g. in (Goyal et al., 2017). We implement DCD and ECD with momentum (Tang et al., 2018), DeepSqueeze with momentum (Tang et al., 2019), Choco-SGD with momentum (Algorithm 2) and standard (all-reduce) mini-batch SGD with momentum and without compression (Dekel et al., 2012). Our implementations are open-source and available at https://github.com/epfml/ChocoSGD. The momentum factor is set to 0.90.9 without dampening. For all algorithms we fine-tune the initial learning rate and gradually warm it up from a relative small value (0.1) (Goyal et al., 2017) for the first 55 epochs. The learning rate is decayed by 1010 twice, at 150150 and 225225 epochs, and stop training at 300 epochs. For Choco-SGD and DeepSqueeze the consensus learning rate γ\gamma is also tuned. The detailed hyper-parameter tuning procedure refers to Appendix F. Every compression scheme is applied to every layer of ResNet20 separately. We evaluate the top-1 test accuracy on every node separately over the whole dataset and report the average performance over all nodes.

Compression Schemes.

We implement two unbiased compression schemes: (i) gsgd⁡b\operatorname{gsgd}_{b} quantization that randomly rounds the weights to bb-bit representations (Alistarh et al., 2017), and (ii) random⁡a\operatorname{random}_{a} sparsification, which preserves a randomly chosen aa fraction of the weights and sets the other ones to zero (Wangni et al., 2018). Further two biased compression schemes: (iii) top⁡a\operatorname{top}_{a}, which selects the aa fraction of weights with the largest magnitude and sets the other ones to zero (Alistarh et al., 2018; Stich et al., 2018), and (iv) sign⁡\operatorname{sign} compression, which compresses each weight to its sign scaled by the norm of the full vector (Bernstein et al., 2018; Karimireddy et al., 2019). We refer to Appendix C for exact definitions of the schemes.

DCD and ECD have been analyzed only for unbiased quantization schemes, thus the combination with the two biased schemes is not supported by theory. In converse, Choco-SGD and DeepSqueeze has been studied only for biased schemes according to Definition 2. However, both unbiased compression schemes can be scaled down in order to meet the specification (cf. discussions in (Stich et al., 2018; Koloskova et al., 2019)) and we adopt this for the experiments.

Results.

The results are summarized in Tab. 1. For unbiased compression schemes, ECD and DCD only achieve good performance when the compression ratio is small, and sometimes even diverge when the compression ratio is high. This is consistent Tang et al. (2018) only consider absolute bounds on the quantization error. Such bounds might be restrictive (i.e. allowing only for low compression) when the input vectors are unbounded. This might be the reason for the instabilities observed here and also in (Tang et al., 2018, Fig. 4), (Koloskova et al., 2019, Figs. 5–6). with the theoretical and experimental results in (Tang et al., 2018). We further observe that the performance of DCD with the biased top⁡a\operatorname{top}_{a} sparsification is much better than with the unbiased random⁡a\operatorname{random}_{a} counterpart, though this operator is not yet supported by theory.

Choco-SGD can generalize reasonably well in all scenarios (at most 1.65% accuracy drop) for fixed training budget. The sign⁡\operatorname{sign} compression achieves state-of-the-art accuracy and requires approximately 32×32\times less bits per weight than the full precision baseline.

Use case I: On-Device Peer-to-Peer Learning

We now shift our focus to challenging real-world scenarios which are intrinsically decentralized, i.e. each part of the training data remains local to each device, and thus centralized methods either fail or are inefficient to implement. Typical scenarios comprise e.g. sensor networks, or mobile devices or hospitals which jointly train a machine learning model. Common to these applications is that i) each device has only access to locally stored or acquired data, ii) communication bandwidth is limited (either physically, or artificially for e.g. metered connections), iii) the global network topology is typically unknown to a single device, and iv) the number of connected devices is typically large. Additionally, this fully decentralized setting is also strongly motivated by privacy aspects, enabling to keep the training data private on each device at all times.

To simulate this scenario, we permanently split the training data between the nodes, i.e. the data is never shuffled between workers during training, and every node has distinct part of the dataset. To the best of our knowledge, no prior works studied this scenario for decentralized deep learning. For the centralized approach, gathering methods such as all-reduce are not efficiently implementable in this setting, hence we compare to the centralized baseline where all nodes route their updates to a central coordinator for aggregation. For the comparison we consider Choco-SGD with sign⁡\operatorname{sign} compression (this combination achieved the compromise between accuracy and compression level in Tab. 1)), decentralized SGD without compression (Lian et al., 2017), and centralized SGD without compression.

Scaling to Large Number of Nodes.

To study the scaling properties of Choco-SGD, we train on 4,16,364,16,36 and 6464 number of nodes. We compare decentralized algorithms on two different topologies: ring as the worst possible topology, and on the torus with much larger spectral gap. The corresponding parameters are listed in Table 2.

We train ResNet8 (He et al., 2016) (7878K parameters), on Cifar10 dataset (50K/10K training/test samples) (Krizhevsky, 2012). For simplicity, we keep the learning rate constant and separately tune it for all methods. We further tune the consensus learning rate for Choco-SGD.

Fix budget of 300 epochs Fixed budget of communication size (1000 MB)

The results are summarized in Fig. 1 (and Fig. 6, Tabs. 7–8 in Appendix G). First we compare the testing accuracy reached after 300 epochs (Fig. 1, left). CentralizedSGD has a good performance for all the considered number of nodes. Choco-SGD slows down due to the influence of the graph topology (Decentralized curve), which is consistent with the spectral gaps order (see Tab. 2), and also influenced by the communication compression (CHOCO curve), which slows down training uniformly for both topologies. We observed that the train performance is similar to the test on Fig. 1, therefore the performance degradation is explained by the slower convergence (Theorem 4.1) and is not a generalization issue. Increasing the number of epochs improves the performance of the decentralized schemes. However, even using 10 times more epochs, we were not able to perfectly close the gap between centralized and decentralized algorithms for both train and test performance.

In the real decentralized scenario, the interest is not to minimize the epochs number, but the amount of communication to reduce the cost of the user’s mobile data. We therefore fix the number of transmitted bits to 1000 MB and compare the best testing accuracy reached (Fig. 1, right). Choco-SGD performs the best while having slight degradation due to increasing number of nodes. It is beneficial to use torus topology when the number of nodes is large because it has good mixing properties, for small networks there is not much difference between these two topologies—the benefit of a large spectral gap is canceled by the increased communication due larger node degree for torus topology. Both Decentralized and Centralized SGD requires significantly larger number of bits to reach reasonable accuracy.

Experiments on a Real Social Network Graph.

We simulate training models on user devices (e.g. mobile phones), connected by a real social network. We chosen Davis Southern women social network (Davis et al., 1941) with 32 nodes. We train ResNet20 (0.270.27 million parameters) model on the Cifar10 dataset (50K/10K training/test samples) (Krizhevsky, 2012) for image classification and a three-layer LSTM architecture (Hochreiter & Schmidhuber, 1997) (28.9528.95 million parameters) for a language modeling task on WikiText-2 (600 training and 60 validation articles with a total of 2′088′6282^{\prime}088^{\prime}628 and 217′646217^{\prime}646 tokens respectively) (Merity et al., 2016). The depicted curves of the training loss are the averaged local loss over all workers (local model with fixed local data); the test performance uses the mean of the evaluations for local models on whole test dataset. For more detailed experimental setup we refer to Appendix F.

The results are summarized in Figs. 2–3 and in Tab. 3. For the image classification task, when comparing the training accuracy reached after the same number of epochs, we observe that the decentralized algorithm performs best, follows by the centralized and lastly the quantized decentralized. However, the test accuracy is highest for the centralized scheme. When comparing the test accuracy reached for the same transmitted data The figure reports the transmitted data on the busiest node, i.e on the max-degree node (degree 14) node for decentralized schemes, and degree 32 for the centralized one. , Choco-SGD significantly outperforms the exact decentralized scheme, with the centralized performing worst. We note a slight accuracy drop, i.e. after the same number of epochs (but much less transmitted data), Choco-SGD does not reach the same level of test accuracy than the baselines.

For the language modeling task, both decentralized schemes suffer a drop in the training loss when the evaluation reaching the epoch budget; while our Choco-SGD outperforms the centralized SGD in test perplexity. When considering perplexity for a fixed data volume (middle and right subfigure of Fig. 3), Choco-SGD performs best, followed by the exact decentralized and centralized algorithms.

On Figure 4 we additionally depict the test accuracy of the averaged model x‾(t)=1n∑i=1nxi(t)\overline{\mathbf{x}}^{(t)}=\frac{1}{n}\sum_{i=1}^{n}\mathbf{x}_{i}^{(t)} (left) and averaged distance of the local models from the averaged model (right), for Choco-SGD on image classification task. Towards the end of the optimization the local models reach consensus (Figure 4, right), and their individual test performances are the same as performance of averaged model. Interestingly, before decreasing the stepsize at the epoch 225, the local models are in general diverging from the averaged model, while decreasing only when the stepsize decreases. A similar behavior was also reported in (Assran et al., 2019).

Use case II: Efficient Large-Scale Training in a Datacenter

Decentralized optimization methods offer a way to address scaling issues even for well connected devices, such as e.g. in datacenter with fast InfiniBand (100Gbps) or Ethernet (10Gbps) connections. Lian et al. (2017) describe scenarios when decentralized schemes can outperform centralized ones, and recently, Assran et al. (2019) presented impressive speedups for training on 256 GPUs, for the setting when all nodes can access all training data. The main differences of their algorithm to Choco-SGD are the asynchronous gossip updates, time-varying communication topology and most importantly exact communication, making their setup not directly comparable to ours. We note that these properties of asynchronous communication and changing topology for faster mixing are orthogonal to our contribution, and offer promise to be combined.

We train ImageNet-1k (1.281.28M/5050K training/validation) (Deng et al., 2009) with Resnet-50 (He et al., 2016). We perform our experiments on 88 machines (n1-standard-32 from Google Cloud with Intel Ivy Bridge CPU platform), where each of machines has 44 Tesla P100 GPUs and each machine interconnected via 10Gbps Ethernet. Within one machine communication is fast and we rely on the local data parallelism to aggregate the gradients for the later gradients communication (over the machines). Between different machines we consider centralized (fully connected topology) and decentralized (ring topology) communication, with and without compressed communication (sign⁡\operatorname{sign} compression). Several methods categorized by communication schemes are evaluated: (i) centralized SGD (full-precision communication), (ii) error-feedback centralized SGD with compressed communications Karimireddy et al. (2019) through sign⁡\operatorname{sign} compression, (iii) decentralized SGD (Lian et al., 2017) with parallelized forward pass and gradients communication (full-precision communication), and (iv) Choco-SGD with sign⁡\operatorname{sign} compressed communications. The mini-batch size on each GPU is 128128, and we follow the general SGD training scheme in (Goyal et al., 2017) and directly use all their hyperparameters for all evaluated methods. Due to the limitation of the computational resource, we did not heavily tune the consensus stepsize for Choco-SGD We estimate the consensus stepsize by running Choco-SGD with different values for the first 3 epochs. .

Results.

We depict the training loss and top-1 test accuracy in terms of epochs and time in Fig. 5. Choco-SGD benefits from its decentralized and parallel structure and takes less time than all-reduce to perform the same number of epochs, while having only a slight 1.5%1.5\% accuracy loss Centralized SGD with full precision gradients achieved test accuracy of 76.37%76.37\%, v.s. 76.03%76.03\% for centralized SGD (with sign⁡\operatorname{sign} compression), v.s. 74.92%74.92\% for plain decentralized SGD, and vs. 75.15%75.15\% for Choco-SGD (with sign⁡\operatorname{sign} compression). . In terms of time per epoch, our speedup does not match that of (Assran et al., 2019), as the used hardware and the communication pattern We consider undirected communication, contrary to the directed 1-peer communication (every node sends and receives one message at every iteration) in Assran et al. (2019). are very different. Their scheme is orthogonal to our approach and could be integrated for better training efficiency. Nevertheless, we still demonstrate a time-wise 20% gain over the common all-reduce baseline, on our used commodity hardware cluster.

Conclusion

We propose the use of Choco-SGD (and its momentum version) for enabling decentralized deep learning training in bandwidth-constrained environments. We provide theoretical convergence guarantees for the non-convex setting and show that the algorithm enjoys linear speedup in the number of nodes. We empirically study the performance of the algorithm in a variety of settings on the image classification (ImageNet-1k, Cifar10) and on the language modeling task (WikiText-2). Whilst previous work successfully demonstrated that decentralized methods can be a competitive alternative to centralized training schemes when no communication constraints are present (Lian et al., 2017; Assran et al., 2019), our main contribution is to enable training in strongly communication-restricted environments, and while respecting the challenging constraint of locality of the training data. We theoretically and practically demonstrate the performance of decentralized schemes for arbitrary high communication compression, and under data-locality, and thus significantly expand the reach of potential applications of fully decentralized deep learning.

Acknowledgements

We acknowledge funding from SNSF grant 200021_175796, as well as a Google Focused Research Award.

References

Appendix A Convergence of Choco-SGD

In this section we present the proof of Theorem 4.1. For this, we will first derive a slightly more general statement: in Theorem A.3 we analyze Choco-SGD for arbitrary stepsizes η\eta, and then derive Theorem 4.1 as a special case.

The structure of the proof follows Koloskova et al. (2019). That is, we first show that Algorithm 1 is a special case of a more general class of algorithms (given in Algorithm 3): Observe that Algorithm 1 consists of two main components: \raisebox{-.9pt} {2}⃝ the stochastic gradient update, performed locally on each node, and \raisebox{-.9pt} {1}⃝ the (quantized) averaging among the nodes. We can show convergence of all algorithms of this type—i.e. stochastic gradient updates \raisebox{-.9pt} {2}⃝ followed by an arbitrary averaging step \raisebox{-.9pt} {1}⃝—as long as the averaging scheme exhibits linear convergence. For the specific averaging used in Choco-SGD, linear convergence has been shown in (Koloskova et al., 2019) and we will use their estimate of the convergence rate of the averaging scheme.

For convenience, we use the following matrix notation in this subsection.

Decentralized SGD with arbitrary averaging is given in Algorithm 3.

Setting X+=XWX^{+}=XW and Y+=X+Y^{+}=X^{+} gives an exact consensus averaging algorithm with mixing matrix WW (Xiao & Boyd, 2004). It converges at the rate c=ρc=\rho, where ρ\rho is an eigengap of mixing matrix WW, defined in Assumption 1. Substituting it into the Algorithm 3 we recover D-PSGD algorithm, analyzed in Lian et al. (2017).

Example: Choco-SGD.

To recover Choco-SGD, we need to choose Choco-Gossip (Koloskova et al., 2019) as consensus averaging scheme, which is defined as X+=X+γY(W−I)X^{+}=X+\gamma Y(W-I) and Y+=Y+Q(X+−Y)Y^{+}=Y+Q(X^{+}-Y) (in the main text we write X^\hat{X} instead of YY). This scheme converges with c=ρ2δ82c=\tfrac{\rho^{2}\delta}{82}. The results from the main part can be recovered by substituting this c=ρ2δ82c=\tfrac{\rho^{2}\delta}{82} in the more general results below. It is important to note that for Algorithm 1 given in the main text, the order of the communication part \raisebox{-.9pt} {1}⃝ and the gradient computation part \raisebox{-.9pt} {2}⃝ is exchanged. We did this to better illustrate that both these parts are independent and that they can be executed in parallel. The effect of this change can be captured by changing the initial values but does not affect the convergence rate.

A.2 Proofs

where σ‾2=∑i=1nσi2n\overline{\sigma}^{2}=\frac{\sum_{i=1}^{n}\sigma_{i}^{2}}{n}.

for Yi=fi(xi(t))−∇Fi(xi(t),ξi(t))Y_{i}=f_{i}(\mathbf{x}_{i}^{(t)})-\nabla F_{i}(\mathbf{x}_{i}^{(t)},\xi_{i}^{(t)}). Expectation of scalar product is equal to zero because ξi\xi_{i} is independent of ξj\xi_{j} since i≠ji\neq j. ∎

Under Assumptions 1–3 the iterates of the Algorithm 3 with constant stepsize η\eta satisfy

Indeed, r0=0≤η24Ac2r_{0}=0\leq\eta^{2}\frac{4A}{c^{2}} as X(0)=X‾(0)X^{(0)}=\overline{X}^{(0)} and Y(0)=0Y^{(0)}=0

Under Assumptions 1–3 with constant stepsize η<14L\eta<\frac{1}{4L}, the averaged iterates x‾(t)=1n∑i=1nxi(t)\overline{\mathbf{x}}^{(t)}=\frac{1}{n}\sum_{i=1}^{n}\mathbf{x}_{i}^{(t)} of Algorithm 3 satisfy:

where cc denotes convergence rate of underlying averaging scheme.

To estimate the second term, we add and subtract ∇f(x‾(t))\nabla f(\overline{\mathbf{x}}^{(t)})

For the last term, we add and subtract ∇f(x‾(t))\nabla f(\overline{\mathbf{x}}^{(t)}) and the sum of ∇fi(xi(t))\nabla f_{i}(\mathbf{x}_{i}^{(t)})

Combining this together and using LL-smoothness to estimate ∥∇f(x‾(t))−∇fi(xi(t))∥22\left\lVert\nabla f(\overline{\mathbf{x}}^{(t)})-\nabla f_{i}(\mathbf{x}_{i}^{(t)})\right\rVert_{2}^{2},

Using Lemma A.2 to bound the third term and using that η≤14L\eta\leq\frac{1}{4L} in the second and in the third terms

A.3 Corollaries

To obtain final convergence rate we carefully tune the stepsize. For this we consider first an auxiliary lemma.

For any parameters r0≥0,b≥0,e≥0,d≥0r_{0}\geq 0,b\geq 0,e\geq 0,d\geq 0 there exists constant stepsize η≤1d\eta\leq\frac{1}{d} such that

Choosing η=min⁡{(r0b(T+1))12,(r0e(T+1))13,1d}≤1d\eta=\min\left\{\left(\frac{r_{0}}{b(T+1)}\right)^{\frac{1}{2}},\left(\frac{r_{0}}{e(T+1)}\right)^{\frac{1}{3}},\frac{1}{d}\right\}\leq\frac{1}{d} we have three cases

η=1d\eta=\frac{1}{d} and is smaller than both (r0b(T+1))12\left(\frac{r_{0}}{b(T+1)}\right)^{\frac{1}{2}} and (r0e(T+1))13\left(\frac{r_{0}}{e(T+1)}\right)^{\frac{1}{3}}, then

η=(r0b(T+1))12<(r0e(T+1))13\eta=\left(\frac{r_{0}}{b(T+1)}\right)^{\frac{1}{2}}<\left(\frac{r_{0}}{e(T+1)}\right)^{\frac{1}{3}}, then

The last case, η=(r0e(T+1))13<(r0b(T+1))12\eta=\left(\frac{r_{0}}{e(T+1)}\right)^{\frac{1}{3}}<\left(\frac{r_{0}}{b(T+1)}\right)^{\frac{1}{2}}

Under Assumptions 1–3 with constant stepsize η\eta tuned as in Lemma A.4, the averaged iterates x‾(t)=1n∑i=1nxi(t)\overline{\mathbf{x}}^{(t)}=\frac{1}{n}\sum_{i=1}^{n}\mathbf{x}_{i}^{(t)} of Algorithm 3 satisfy:

where cc denotes convergence rate of underlying averaging scheme, F0=f(x‾(0))−f⋆F_{0}=f(\overline{\mathbf{x}}^{(0)})-f^{\star}.

The result follows from Theorem A.3 and Lemma A.4 with r0=4(f(x‾(0))−f⋆)r_{0}=4\left(f(\overline{\mathbf{x}}^{(0)})-f^{\star}\right), b=2σ‾2Lnb=\frac{2\overline{\sigma}^{2}L}{n}, e=36G2L2c2e=\frac{36G^{2}L^{2}}{c^{2}} and d=4Ld=4L. ∎

The first term shows a linear speed up compared to SGD on one node, whereas the underlying averaging scheme affects only the second-order term. Substituting the convergence rate for exact averaging with WW (c=ρc=\rho) gives the rate O(\nicefrac1nT+\nicefrac1(Tρ)23)\mathcal{O}(\nicefrac{{1}}{{\sqrt{nT}}}+\nicefrac{{1}}{{(T\rho)^{\frac{2}{3}}}}).

Choco-SGD with the underlying Choco-Gossip averaging scheme converges at the rate O(\nicefrac1nT+\nicefrac1(Tρ2δ)23)\mathcal{O}(\nicefrac{{1}}{{\sqrt{nT}}}+\nicefrac{{1}}{{(T\rho^{2}\delta)^{\frac{2}{3}}}}). The dependence on ρ\rho (eigengap of the mixing matrix WW) is worse than in the exact case. This might either just be an artifact of our proof technique or a consequence of supporting arbitrary high compression.

The corollary gives guarantees for the averaged vector of parameters x‾\overline{\mathbf{x}}, however in a decentralized setting it is very expensive and sometimes impossible to average all the parameters distributed across several machines, especially when the number of machines and the model size is large. We can get similar guarantees on the individual iterates xi\mathbf{x}_{i} as e.g. in (Assran et al., 2019). We summarize these briefly below.

Under the same setting as in Corollary A.5,

where we used LL-smoothness of ff. Using Theorem A.3 and tuning the stepsize as in Lemma A.4 we get the statement of the corollary. ∎

Choosing the stepsize differently, we can also get the following convergence rate for T=Ω(nL2)T=\Omega(nL^{2}):

Under Assumptions 1–3 with constant stepsize η=nT+1\eta=\sqrt{\frac{n}{T+1}} for T≥16nL2T\geq 16nL^{2}, the averaged iterates x‾(t)=1n∑i=1nxi(t)\overline{\mathbf{x}}^{(t)}=\frac{1}{n}\sum_{i=1}^{n}\mathbf{x}_{i}^{(t)} of Algorithm 3 satisfy:

where cc denotes convergence rate of underlying averaging scheme.

Appendix B Useful Inequalities

Appendix C Compression Schemes

We implement the compression schemes detailed below.

where u∼u.a.r.d\mathbf{u}\sim_{u.a.r.}^{d} is a random dithering vector and sig⁡(x)\operatorname{sig}(\mathbf{x}) assigns the element-wise sign: (sig⁡(x))i=1(\operatorname{sig}(\mathbf{x}))_{i}=1 if (x)i≥0(\mathbf{x})_{i}\geq 0 and (sig⁡(x))i=−1(\operatorname{sig}(\mathbf{x}))_{i}=-1 if (x)i<0(\mathbf{x})_{i}<0. As the value in the right bracket will be rounded to an integer in {0,…,2(b−1)−1}\{0,\dots,2^{(b-1)}-1\}, each coordinate can be encoded with at most (b−1)+1(b-1)+1 bits (1 for the sign). For more efficent encoding schemes cf. Alistarh et al. (2017).

for τ=1+min⁡{d22(b−1),d2(b−1)}\tau=1+\min\left\{\frac{d}{2^{2(b-1)}},\frac{\sqrt{d}}{2^{(b-1})}\right\} and is a δ=1τ\delta=\frac{1}{\tau} compression operator (Koloskova et al., 2019).

and is a δ=a\delta=a compression operator (Stich et al., 2018).

Only 32⌊ad⌋32{\lfloor ad\rfloor} bits are required to send random⁡a(x)\operatorname{random}_{a}(\mathbf{x}) to another node—all the values of non-zero entries (we assume that entries are represented as float32 numbers). Receiver can recover positions of these entries if it knows the random seed of uniform sampling operator used to select these entries. This random seed could be communicated once on preprocessing stage (before starting the algorithm).

where u(x)∈{0,1}d\mathbf{u}(\mathbf{x})\in\{0,1\}^{d}, ∥u∥1=⌊ad⌋\left\lVert\mathbf{u}\right\rVert_{1}=\lfloor ad\rfloor is a masking vector with (u)i=1(\mathbf{u})_{i}=1 for indices i∈π−1({1,…,⌊ad⌋})i\in\pi^{-1}(\{1,\dots,\lfloor ad\rfloor\}) where the permutation π\pi is such that ∣(x)π(1)∣≥∣(x)π(2)∣≥⋯≥∣(x)π(d)∣\left\lvert(\mathbf{x})_{\pi(1)}\right\rvert\geq\left\lvert(\mathbf{x})_{\pi(2)}\right\rvert\geq\cdots\geq\left\lvert(\mathbf{x})_{\pi(d)}\right\rvert. The top⁡a\operatorname{top}_{a} operator is a δ=a\delta=a compression operator (Stich et al., 2018).

In the case of top⁡a\operatorname{top}_{a} compression 2⋅32⌊ad⌋2\cdot 32{\lfloor ad\rfloor} bits are required because along with the values we need to send positions of these values.

The sign⁡\operatorname{sign} operator is a δ=∥x∥12d∥x∥22\delta=\frac{\left\lVert\mathbf{x}\right\rVert_{1}^{2}}{d\left\lVert\mathbf{x}\right\rVert_{2}^{2}} compression operator (Karimireddy et al., 2019).

In total for the sign⁡\operatorname{sign} compression we need to send only d+32d+32 bits—one bit for every entry in x\mathbf{x} and 32 bits for ∥x∥1\left\lVert\mathbf{x}\right\rVert_{1}.

Appendix D Choco-SGD with Momentum

Algorithm 2 demonstrates how to combine Choco-SGD with weight decay and momentum. Nesterov momentum can be analogously adapted for our decentralized setting.

Appendix E Error Feedback Interpretation of Choco-SGD

To better understand how does Choco-SGD work, we can interpret it as an error feedback algorithm (Stich et al., 2018; Karimireddy et al., 2019; Stich & Karimireddy, 2019). We can equivalently rewrite Choco-SGD (Algorithm 1) as Algorithm 4. The common feature of error feedback algorithms is that quantization errors are saved into the internal memory, which is added to the compressed value at the next iteration. In Choco-SGD the value we want to transmit is the difference xi(t)−xi(t−1)\mathbf{x}_{i}^{(t)}-\mathbf{x}_{i}^{(t-1)}, which represents the evolution of local variable xi\mathbf{x}_{i} at step tt. Before compressing this value on line 4, the internal memory is added on line 3 to correct for the errors. Then, on line 5 internal memory is updated. Note that mi(t)=xi(t−1)−x^i(t)\mathbf{m}_{i}^{(t)}=\mathbf{x}^{(t-1)}_{i}-\hat{\mathbf{x}}_{i}^{(t)} in the old notation.

Appendix F Detailed Experimental Setup and Tuned Hyperparameters

We precise the procedure of model training as well as the hyper-parameter tuning in this section.

For the comparison we consider Choco-SGD with sign⁡\operatorname{sign} compression (this combination achieved the compromise between accuracy and compression level in Table 1)), decentralized SGD without compression, and centralized SGD without compression. We train two models, firstly ResNet20 (He et al., 2016) (0.270.27 million parameters) for image classification on the Cifar10 dataset (50K/10K training/test samples) (Krizhevsky, 2012) and secondly, a three-layer LSTM architecture (Hochreiter & Schmidhuber, 1997) (28.9528.95 million parameters) for a language modeling task on WikiText-2 (600 training and 60 validation articles with a total of 2′088′6282^{\prime}088^{\prime}628 and 217′646217^{\prime}646 tokens respectively) (Merity et al., 2016). For the language modeling task, we borrowed and adapted the general experimental setup of Merity et al. (2017), where we use a three-layer LSTM with hidden dimension of size 650650. The loss is averaged over all examples and timesteps. The BPTT length is set to 3030. We fine-tune the value of gradient clipping (0.40.4), and the dropout (0.40.4) is only applied on the output of LSTM.

We train both of ResNet20 and LSTM for 300300 epochs, unless mentioned specifically. The per node mini-batch size is 3232 for both datasets. The momentum (with factor 0.90.9) is only applied on the ResNet20 training.

Social Network and a Datacenter details.

For all algorithms, we gradually warmup (Goyal et al., 2017) the learning rate from a relative small value (0.1) to the fine-tuned initial learning rate for the first 55 training epochs. During the training procedure, the tuned initial learning rate is decayed by the factor of 1010 when accessing 50%50\% and 75%75\% of the total training epochs. The learning rate is tuned by finding the optimal initial learning rate (after the scaling).

The optimal η^\hat{\eta} is searched in a pre-defined grid and we ensure that the best performance was contained in the middle of the grids. For example, if the best performance was ever at one of the extremes of the grid, we would try new grid points. Same searching logic applies to the consensus stepsize.

Table 4 demonstrates the fine-tuned hpyerparameters of Choco-SGD for training ResNet-20 on Cifar10, while Table 6 reports our fine-tuned hpyerparameters of our baselines. Table 5 demonstrates the fine-tuned hpyerparameters of Choco-SGD for training ResNet-20/LSTM on a social network topology.

We estimate the runtime information (depicted in Figure 5) of different methods from three trials of the evaluation on Google Cloud (Kubernetes Engine). More precisely, we create the cluster on Google Cloud for three times and each time we estimate the time per mini-batch of different methods (through the first two training epochs).

Appendix G Additional Plots

To complement our results for scaling to a large number of nodes, we here additionally depict the learning curves (e.g. test accuracy) for the training on 64 nodes. We also mark the levels used for Fig. 1.

We additionally visualize the learning curves for the social network topology in Fig. 7 and Fig. 8.

We additionally provide the learning curves of training top-1, top-5 accuracy and test top-5 accuracy for the datacenter experiment in Fig. 9.