Communication Compression for Decentralized Training

Hanlin Tang, Shaoduo Gan, Ce Zhang, Tong Zhang, Ji Liu

Introduction

When training machine learning models in a distributed fashion, the underlying constraints of how workers (or nodes) communication have a significant impact on the training algorithm. When workers cannot form a fully connected communication topology or the communication latency is high (e.g., in sensor networks or mobile networks), decentralizing the communication comes to the rescue. On the other hand, when the amount of data sent through the network is an optimization objective (maybe to lower the cost or energy consumption), or the network bandwidth is low, compressing the traffic, either via sparsification (Wangni et al., 2017; Konečnỳ and Richtárik, 2016) or quantization (Zhang et al., 2017a; Suresh et al., 2017) is a popular strategy. In this paper, our goal is to develop a novel framework that works robustly in an environment that both decentralization and communication compression could be beneficial. In this paper, we focus on quantization, the process of lowering the precision of data representation, often in a stochastically unbiased way. But the same techniques would apply to other unbiased compression schemes such as sparsification.

Both decentralized training and quantized (or compressed more generally) training have attracted intensive interests recently (Yuan et al., 2016; Zhao and Song, 2016; Lian et al., 2017a; Konečnỳ and Richtárik, 2016; Alistarh et al., 2017). Decentralized algorithms usually exchange local models among nodes, which consumes the main communication budget; on the other hand, quantized algorithms usually exchange quantized gradient, and update an un-quantized model. A straightforward idea to combine these two is to directly quantize the models sent through the network during decentralized training. However, this simple strategy does not converge to the right solution as the quantization error would accumulate during training. The technical contribution of this paper is to develop novel algorithms that combine both decentralized training and quantized training together.

We consider the following decentralized optimization:

where nn is the number of node and Di\mathcal{D}_{i} is the local data distribution for node ii. nn nodes form a connected graph and each node can only communicate with its neighbors. Here we only assume fi(x)f_{i}(x)’s are with L-Lipschitzian gradients.

Summary of Technical Contributions.

In this paper, we propose two decentralized parallel stochastic gradient descent algorithms (D-PSGD): extrapolation compression D-PSGD (ECD-PSGD) and difference compression D-PSGD (DCD-PSGD). Both algorithms can be proven to converge in the rate roughly O(1/nT)O(1/\sqrt{nT}) where TT is the number of iterations. The convergence rates are consistent with two special cases: centralized parallel stochastic gradient descent (C-PSGD) and D-PSGD. To the best of our knowledge, this is the first work to combine quantization algorithms and decentralized algorithms for generic optimization.

The key difference between ECD-PSGD and DCD-PSGD is that DCD-PSGD quantizes the difference between the last two local models, and ECD-PSGD quantizes the extrapolation between the last two local models. DCD-PSGD admits a slightly better convergence rate than ECD-PSGD when the data variation among nodes is very large. On the other hand, ECD-PSGD is more robust to more aggressive quantization, as extremely low precision quantization can cause DCD-PSGD to diverge, since DCD-PSGD has strict constraint on quantization. In this paper, we analyze both algorithms, and empirically validate our theory. We also show that when the underlying network has both high latency and low bandwidth, both algorithms outperform state-of-the-arts significantly. We present both algorithm because we believe both of them are theoretically interesting. In practice, ECD-PSGD could potentially be a more robust choice.

Definitions and notations

Throughout this paper, we use following notations and definitions:

∇f(⋅)\nabla f(\cdot) denotes the gradient of a function ff.

f∗f^{*} denotes the optimal solution of (1).

λi(⋅)\lambda_{i}(\cdot) denotes the ii-th largest eigenvalue of a matrix.

∥⋅∥\|\cdot\| denotes the l2l_{2} norm for vector.

∥⋅∥F\|\cdot\|_{F} denotes the vector Frobenius norm of matrices.

C(⋅)\bm{C}(\cdot) denotes the compressing operator.

Related work

The Stocahstic Gradient Descent (SGD) (Ghadimi and Lan, 2013; Moulines and Bach, 2011; Nemirovski et al., 2009) - a stochastic variant of the gradient descent method - has been widely used for solving large scale machine learning problems (Bottou, 2010). It admits the optimal convergence rate O(1/T)O(1/\sqrt{T}) for non-convex functions.

Centralized algorithms

The centralized algorithms is a widely used scheme for parallel computation, such as Tensorflow (Abadi et al., 2016), MXNet (Chen et al., 2015), and CNTK (Seide and Agarwal, 2016). It uses a central node to control all leaf nodes. For Centralized Parallel Stochastic Gradient Descent (C-PSGD), the central node performs parameter updates and leaf nodes compute stochastic gradients based on local information in parallel. In Agarwal and Duchi (2011); Zinkevich et al. (2010), the effectiveness of C-PSGD is studied with latency taken into consideration. The distributed mini-batches SGD, which requires each leaf node to compute the stochastic gradient more than once before the parameter update, is studied in Dekel et al. (2012). Recht et al. (2011) proposed a variant of C-PSGD, HOGWILD, and proved that it would still work even if we allow the memory to be shared and let the private mode to be overwriten by others. The asynchronous non-convex C-PSGD optimization is studied in Lian et al. (2015). Zheng et al. (2016) proposed an algorithm to improve the performance of the asynchronous C-PSGD. In Alistarh et al. (2017); De Sa et al. (2017), a quantized SGD is proposed to save the communication cost for both convex and non-convex object functions. The convergence rate for C-PSGD is O(1/Tn)O(1/\sqrt{Tn}). The tradeoff between the mini-batch number and the local SGD step is studied in Lin et al. (2018); Stich (2018).

Decentralized algorithms

Recently, decentralized training algorithms have attracted significantly amount of attentions. Decentralized algorithms are mostly applied to solve the consensus problem (Zhang et al., 2017b; Lian et al., 2017a; Sirb and Ye, 2016), where the network topology is decentralized. A recent work shows that decentralized algorithms could outperform the centralized counterpart for distributed training (Lian et al., 2017a). The main advantage of decentralized algorithms over centralized algorithms lies on avoiding the communication traffic in the central node. In particular, decentralized algorithms could be much more efficient than centralized algorithms when the network bandwidth is small and the latency is large. The decentralized algorithm (also named gossip algorithm in some literature under certain scenarios (Colin et al., 2016)) only assume a connect computational network, without using the central node to collect information from all nodes. Each node owns its local data and can only exchange information with its neighbors. The goal is still to learn a model over all distributed data. The decentralized structure can applied in solving of multi-task multi-agent reinforcement learning (Omidshafiei et al., 2017; Mhamdi et al., 2017). Boyd et al. (2006) uses a randomized weighted matrix and studied the effectiveness of the weighted matrix in different situations. Two methods (Li et al., 2017; Shi et al., 2015) were proposed to blackuce the steady point error in decentralized gradient descent convex optimization. Dobbe et al. (2017) applied an information theoretic framework for decentralize analysis. The performance of the decentralized algorithm is dependent on the second largest eigenvalue of the weighted matrix. In He et al. (2018), In He et al. (2018), they proposed the gradient descent based algorithm (CoLA) for decentralized learning of linear classification and regression models, and proved the convergence rate for strongly convex and general convex cases.

Decentralized parallel stochastic gradient descent

The Decentralized Parallel Stochastic Gradient Descent (D-PSGD) (Nedic and Ozdaglar, 2009; Yuan et al., 2016) requires each node to exchange its own stochastic gradient and update the parameter using the information it receives. In Nedic and Ozdaglar (2009), the convergence rate for a time-varying topology was proved when the maximum of the subgradient is assumed to be bounded. In Lan et al. (2017), a new decentralized primal-dual type method is proposed with a computational complexity of O(n/T)O(\sqrt{n/T}) for general convex objectives. The linear speedup of D-PSGD is proved in Lian et al. (2017a), where the computation complexity is O(1/nT)O(1/\sqrt{nT}). The asynchronous variant of D-PSGD is studied in Lian et al. (2017b).

Compression

To guarantee the convergence and correctness, this paper only considers using the unbiased stochastic compression techniques. Existing methods include randomized quantization (Zhang et al., 2017a; Suresh et al., 2017) and randomized sparsification (Wangni et al., 2017; Konečnỳ and Richtárik, 2016). Other compression methods can be found in Kashyap et al. (2007); Lavaei and Murray (2012); Nedic et al. (2009). In Drumond et al. (2018), a compressed DNN training algorithm is proposed. In Stich et al. (2018), a centralized biased sparsified parallel SGD with memory is studied and proved to admits an factor of acceleration.

Preliminary: decentralized parallel stochastic gradient descent (D-PSGD)

Unlike the traditional (centralized) parallel stochastic gradient descent (C-PSGD), which requires a central node to compute the average value of all leaf nodes, the decentralized parallel stochastic gradient descent (D-PSGD) algorithm does not need such a central node. Each node (say node ii) only exchanges its local model x(i)\bm{x}^{(i)} with its neighbors to take weighted average, specifically, x(i)=∑j=1nWijx(j)\bm{x}^{(i)}=\sum_{j=1}^{n}W_{ij}\bm{x}^{(j)} where Wij≥0W_{ij}\geq 0 in general and Wij=0W_{ij}=0 means that node ii and node jj is not connected. At ttth iteration, D-PSGD consists of three steps (ii is the node index):

1. Each node computes the stochastic gradient ∇Fi(xt(i);ξt(i))\nabla F_{i}(\bm{x}^{(i)}_{t};\xi_{t}^{(i)}), where ξt(i)\xi^{(i)}_{t} is the samples from its local data set and xt(i)\bm{x}^{(i)}_{t} is the local model on node ii.

2. Each node queries its neighbors’ variables and updates its local model using x(i)=∑j=1nWijx(j)\bm{x}^{(i)}=\sum_{j=1}^{n}W_{ij}\bm{x}^{(j)}.

3. Each node updates its local model xt(i)←xt(i)−γt∇Fi(xt(i);ξt(i))\bm{x}^{(i)}_{t}\leftarrow\bm{x}^{(i)}_{t}-\gamma_{t}\nabla F_{i}\left(\bm{x}_{t}^{(i)};\xi^{(i)}_{t}\right) using stochastic gradient, where γt\gamma_{t} is the learning rate.

To look at the D-PSGD algorithm from a global view, by defining

the D-PSGD can be summarized into the form Xt+1=XtW−γtG(Xt;ξt)X_{t+1}=X_{t}W-\gamma_{t}G(X_{t};\xi_{t}).

The convergence rate of D-PSGD can be shown to be O(σnT+n13ζ23T23)O\left(\frac{\sigma}{\sqrt{nT}}+\frac{n^{\frac{1}{3}}\zeta^{\frac{2}{3}}}{T^{\frac{2}{3}}}\right) (without assuming convexity) where both σ\sigma and ζ\zeta are the stochastic variance (please refer to Assumption 1 for detailed definitions), if the learning rate is chosen appropriately.

Quantized, Decentralized Algorithms

We introduce two quantized decentralized algorithms that compress information exchanged between nodes. All communications for decentralized algorithms are exchanging local models x(i)\bm{x}^{(i)}.

To reduce the communication cost, a straightforward idea is to compress the information exchanged within the decentralized network just like centralized algorithms sending compressed stochastic gradient (Alistarh et al., 2017). Unfortunately, such naive combination does not work even using the unbiased stochastic compression and diminishing learning rate as shown in Figure 1. The reason can be seen from the detailed derivation (please find it in Supplement).

Before propose our solutions to this issue, let us first make some common optimization assumptions for analyzing decentralized stochastic algorithms (Lian et al., 2017b).

Throughout this paper, we make the following commonly used assumptions:

Lipschitzian gradient: All function fi(⋅)f_{i}(\cdot)’s are with LL-Lipschitzian gradients.

Symmetric double stochastic matrix: The weighted matrix WW is a real double stochastic matrix that satisfies W=W⊤W=W^{\top} and W1=WW\bm{1}=W.

Spectral gap: Given the symmetric doubly stochastic matrix WW, we define ρ:=max⁡{∣λ2(W)∣,∣λn(W)∣}\rho:=\max\{|\lambda_{2}(W)|,|\lambda_{n}(W)|\} and assume ρ<1\rho<1.

Bounded variance: Assume the variance of stochastic gradient to be bounded

The last assumption essentially restricts the compression to be lossy but unbiased. Biased stochastic compression is generally hard to ensure the convergence and lossless compression can combine with any algorithms. Both of them are beyond of the scope of this paper. The commonly used stochastic unbiased compression include random quantizationA real number is randomly quantized into one of closest thresholds, for example, givens the thresholds {0,0.3,0.8,1}\{0,0.3,0.8,1\}, the number “0.50.5” will be quantized to 0.30.3 with probability 40%40\% and to 0.80.8 with probability 60%60\%. Here, we assume that all numbers have been normalized into the range $.(Zhangetal.,2017a)andsparsificationArealnumber. (Zhang et al., 2017a) and sparsificationA real numberzissettowithprobabilityis set to with probability1-pandtoand toz/pwithprobabilitywith probabilityp$. (Wangni et al., 2017; Konečnỳ and Richtárik, 2016).

In this section, we introduces a difference based approach, namely, difference compression D-PSGD (DCD-PSGD), to ensure efficient convergence.

The DCD-PSGD basically follows the framework of D-PSGD, except that nodes exchange the compressed difference of local models between two successive iterations, instead of exchanging local models. More specifically, each node needs to store its neighbors’ models in last iteration \{{\color[rgb]{0,0,0}{\hat{\bm{x}}}^{(j)}_{t}}:j~{}\text{is nodei's neighbor}\} and follow the following steps:

take the weighted average and apply stochastic gradient descent step: \bm{x}_{t+\frac{1}{2}}^{(i)}=\sum_{j=1}^{n}W_{ij}{{\color[rgb]{0,0,0}\hat{\bm{x}}}^{(j)}_{t}}-\gamma\nabla F_{i}(\bm{x}^{(i)}_{t};\xi^{(i)}_{t}), where x^t(j)\hat{\bm{x}}^{(j)}_{t} is just the replica of xt(j)\bm{x}^{(j)}_{t} but is stored on node iiActually each neighbor of node jj maintains a replica of xt(j)\bm{x}^{(j)}_{t}.;

compress the difference between xt(i)\bm{x}^{(i)}_{t} and xt+12(i)\bm{x}_{t+\frac{1}{2}}^{(i)} and update the local model:zt(i)=xt+12(i)−xt(i),xt+1(i)=xt(i)+C(zt(i))\bm{z}_{t}^{(i)}=\bm{x}_{t+\frac{1}{2}}^{(i)}-\bm{x}_{t}^{(i)},\quad\bm{x}^{(i)}_{t+1}=\bm{x}^{(i)}_{t}+\bm{C}(\bm{z}_{t}^{(i)});

send C(zt(i))\bm{C}(\bm{z}_{t}^{(i)}) and query neighbors’ C(zt)\bm{C}(\bm{z}_{t}) to update the local replica: \forall j~{}\text{is nodei's neighbor}, x^t+1(j)=x^t(j)+C(zt(j))\hat{\bm{x}}_{t+1}^{(j)}=\hat{\bm{x}}^{(j)}_{t}+\bm{C}(\bm{z}_{t}^{(j)}).

The full DCD-PSGD algorithm is described in Algorithm 1.

To ensure convergence, we need to make some restriction on the compression operator C(⋅)\bm{C}(\cdot). Again this compression operator could be random quantization or random sparsification or any other operators. We introduce the definition of the signal-to-noise related parameter α\alpha. Let α:=sup⁡Z≠0∥Q∥F2/∥Z∥F2\alpha:=\sqrt{\sup_{Z\neq 0}{\|Q\|^{2}_{F}/\|Z\|^{2}_{F}}}, where Q=Z−C(Z)Q=Z-C(Z). We have the following theorem.

Under the Assumption 1, if α\alpha satisfies (1−ρ)2−4μ2α2>0(1-\rho)^{2}-4\mu^{2}\alpha^{2}>0, choosing γ\gamma satisfying 1−3D1L2γ2>01-3D_{1}L^{2}\gamma^{2}>0, we have the following convergence rate for Algorithm 1

where μ:=max⁡i∈{2,⋯ ,n}∣λi−1∣\mu:=\max_{i\in\{2,\cdots,n\}}|\lambda_{i}-1|, and

To make the result more clear, we appropriately choose the steplength in the following:

Choose γ=(6D1L+6D2L+σnT12+ζ23T13)−1\gamma=\left(6\sqrt{D_{1}}L+6\sqrt{D_{2}L}+\frac{\sigma}{\sqrt{n}}T^{\frac{1}{2}}+\zeta^{\frac{2}{3}}T^{\frac{1}{3}}\right)^{-1} in Algorithm 1. If α\alpha is small enough that satisfies (1−ρ)2−4μ2α2>0(1-\rho)^{2}-4\mu^{2}\alpha^{2}>0, then we have

where D1D_{1}, D2D_{2} follow to same definition in Theorem 1 and we treat f(0)−f∗f(0)-f^{*}, LL, and ρ\rho constants.

Since the leading term of the convergence rate is O(1/Tn)O\left(1/\sqrt{Tn}\right) when TT is large, which is consistent with the convergence rate of C-PSGD, this indicates that we would achieve a linear speed up with respect to the number of nodes.

Consistence with D-PSGD

Setting α=0\alpha=0 to match the scenario of D-PSGD, ECD-PSGD admits the rate O(σnT+ζ23T23)O\left(\frac{\sigma}{\sqrt{nT}}+\frac{\zeta^{\frac{2}{3}}}{T^{\frac{2}{3}}}\right), that is slightly better the rate of D-PSGD proved in Lian et al. (2017b) O(σnT+n23ζ23T23)O\left(\frac{\sigma}{\sqrt{nT}}+\frac{n^{\frac{2}{3}}\zeta^{\frac{2}{3}}}{T^{\frac{2}{3}}}\right). The non-leading terms’ dependence of the spectral gap (1−ρ)(1-\rho) is also consistent with the result in D-PSDG.

2 Extrapolation compression approach

From Theorem 1, we can see that there is an upper bound for the compressing level α\alpha in DCD-PSGD. Moreover, since the spectral gap (1−ρ)(1-\rho) would decrease with the growth of the amount of the workers, so DCD-PSGD will fail to work under a very aggressive compression. So in this section, we propose another approach, namely ECD-PSGD, to remove the restriction of the compressing degree, with a little sacrifice on the computation efficiency.

For ECD-PSGD, we make the following assumption that the noise brought by compression is bounded.

The node jj, computes the zz-value that is obtained through extrapolation

Using this way to estimate the neighbors’ local models leads to the following equivalent updating form

The full extrapolation compression D-PSGD (ECD-PSGD) algorithm is summarized in Algorithm 2.

Below we will show that EDC-PSGD algorithm would admit the same convergence rate and the same computation complexity as D-PSGD.

Under Assumptions 1 and 2, choosing γt\gamma_{t} in Algorithm 2 to be constant γ\gamma satisfying 1−6C1L2γ2>01-6C_{1}L^{2}\gamma^{2}>0, we have the following convergence rate for Algorithm 2

where C1:=1(1−ρ)2C_{1}:=\frac{1}{(1-\rho)^{2}}, C2:=11−6ρ−2C1L2γ2C_{2}:=\frac{1}{1-6\rho^{-2}C_{1}L^{2}\gamma^{2}}, C3:=12L2C2C1γ2C_{3}:=12L^{2}C_{2}C_{1}\gamma^{2}, and C4:=1−LγC_{4}:=1-L\gamma.

To make the result more clear, we choose the steplength in the following:

In Algorithm 2 choose the steplength γ=(12C1L+σnT12+ζ23T13)−1\gamma=\left(12\sqrt{C_{1}}L+\frac{\sigma}{\sqrt{n}}T^{\frac{1}{2}}+\zeta^{\frac{2}{3}}T^{\frac{1}{3}}\right)^{-1}. Then it admits the following convergence rate (with f(0)−f∗f(0)-f^{*}, LL, and ρ\rho treated as constants).

Since the leading term of the convergence rate is O(1/nT)O(1/\sqrt{nT}) when TT is large, which is consistent with the convergence rate of C-PSGD, this indicates that we would achieve a linear speed up with respect to the number of nodes.

Consistence with D-PSGD

Comparison between DCD-PSGD and ECD-PSGD

Experiments

In this section we evaluate two decentralized algorithms by comparing with an Allreduce implementation of centralized SGD. We run experiments under diverse network conditions and show that, decentralized algorithms with low precision can speed up training without hurting convergence.

We choose the image classification task as a benchmark to evaluate our theory. We train ResNet-20 (He et al., 2016) on CIFAR-10 dataset which has 50,000 images for training and 10,000 images for testing. Two proposed algorithms are implemented in Microsoft CNTK and compablack with CNTK’s original implementation of distributed SGD:

Centralized: This implementation is based on MPI Allreduce primitive with full precision (32 bits). It is the standard training method for multiple nodes in CNTK.

Decentralized_32bits/8bits: The implementation of the proposed decentralized approach with OpenMPI. The full precision is 32 bits, and the compressed precision is 8 bits.

In this paper, we omit the comparison with quantized centralized training because the difference between Decentralized 8bits and Centralized 8bits would be similar to the original decentralized training paper Lian et al. (2017a) – when the network latency is high, decentralized algorithm outperforms centralized algorithm in terms of the time for each epoch.

We run experiments on 8 Amazon p2.xlargep2.xlarge EC2 instances, each of which has one Nvidia K80 GPU. We use each GPU as a node. In decentralized cases, 8 nodes are connected as a ring topology, which means each node just communicates with its two neighbors. The batch size for each node is same as the default configuration in CNTK. We also tune learning rate for each variant.

2 Convergence and Run Time Performance

We first study the convergence of our algorithms. Figure 2(a) shows the convergence w.r.t # epochs of centralized and decentralized cases. We only show ECD-PSGD in the figure (and call it Decentralized) because DCD-PSGD has almost identical convergence behavior in this experiment. We can see that with our algorithms, decentralization and compression would not hurt the convergence rate.

We then compare the runtime performance. Figure 2(b, c, d) demonstrates how training loss decreases with the run time under different network conditions. We use tctc command to change bandwidth and latency of the underlying network. By default, 1.4 Gbps bandwidth and 0.13 ms latency is the best network condition we can get in this cluster. On this occasion, all implementations have a very similar runtime performance because communication is not the bottleneck for system. When the latency is high, as shown in 2(c), decentralized algorithms in both low and full precision can outperform the Allreduce method because of fewer number of communications. However, in low bandwidth case, training time is mainly dominated by the amount of communication data, so low precision method can be obviously faster than these full precision methods.

3 Speedup in Diverse Network Conditions

To better understand the influence of bandwidth and latency on speedup, we compare the time of one epoch under various of network conditions. Figure 3(a, b) shows the trend of epoch time with bandwidth decreasing from 1.4 Gbps to 5 Mbps. When the latency is low (Figure 3(a)), low precision algorithm is faster than its full precision counterpart because it only needs to exchange around one fourth of full precision method’s data amount. Note that although in a decentralized way, full precision case has no advantage over Allreduce in this situation, because they exchange exactly the same amount of data. When it comes to high latency shown in Figure 3(b), both full and low precision cases are much better than Allreduce in the beginning. But also, full precision method gets worse dramatically with the decline of bandwidth.

Figure 3(c, d) shows how latency influences the epoch time under good and bad bandwidth conditions. When bandwidth is not the bottleneck (Figure 3(c)), decentralized approaches with both full and low precision have similar epoch time because they have same number of communications. As is expected, Allreduce is slower in this case. When bandwidth is very low (Figure 3(d)), only decentralized algorithm with low precision can achieve best performance among all implementations.

4 Discussion

Our previous experiments validate the efficiency of the decentralized algorithms on 8 nodes with 8 bits. However, we wonder if we can scale it to more nodes or compress the exchanged data even more aggressively. We firstly conducted experiments on 16 nodes with 8 bits as before. According to Figure 4(a), Alg. 1 and Alg. 2 on 16 nodes can still achieve basically same convergence rate as Allreduce, which shows the scalability of our algorithms. However, they can not be comparable to Allreduce with 4 bits, as is shown in 4(b). What is noteworthy is that these two compression approaches have quite different behaviors in 4 bits. For Alg. 1, although it converges much slower than Allreduce, its training loss keeps reducing. However, Alg. 2 just diverges in the beginning of training. This observation is consistent with our theoretical analysis.

Conclusion

In this paper, we studied the problem of combining two tricks of training distributed stochastic gradient descent under imperfect network conditions: quantization and decentralization. We developed two novel algorithms or quantized, decentralized training, analyze the theoretical property of both algorithms, and empirically study their performance in a various settings of network conditions. We found that when the underlying communication networks has both high latency and low bandwidth, quantized, decentralized algorithm outperforms other strategies significantly.

References

Proof organization

In section A, we first provide the properties of updating rule (7). Sections C and B will next specify QtQ_{t} and show the upper bounds for ∥Qt∥2\|Q_{t}\|^{2} in Algorithms 2 and 1 respectively.

Notations

We define some additional notations throughout the following proof

Appendix A General bound with compression noise

In this section, to see the influence of the compression more clearly, we are going to prove two general bounds (see see Lemma 7 and Lemma 8) for compressed D-PSGD that has the same updating rule like (7). Those bounds are very helpful for the following proof of our algorithms.

The most challenging part of a decentralized algorithm, unlike the centralized algorithm, is that we need to ensure the local model on each node to converge to the average value X‾t\overline{X}_{t}. So we start with an analysis of the quantity ∥X‾t−xt(i)∥2\left\|\overline{X}_{t}-\bm{x}_{t}^{(i)}\right\|^{2} and its influence on the final convergence rate. For both ECD-PSGD and DCD-PSGD, we are going to prove that

From the above two inequalities, we can see that the extra noise term decodes the convergence efficiency of xt(i)\bm{x}_{t}^{(i)} to the average X‾t\overline{X}_{t}.

The proof of the general bound for (7) is divided into two parts. In subsection A.1, we provide a new perspective in understanding decentralization, which can be very helpful for the simplicity of our following proof. In subsection A.2, we give the detail proof for the general bound.

To have a better understanding of how decentralized algorithms work, and how can we ensure a consensus from all local variables on each node. We provide a new perspective to understand decentralization using coordinate transformation, which can simplify our analysis in the following.

The confusion matrix WW satisfies W=∑i=1nλivi(vi)⊤W=\sum_{i=1}^{n}\lambda_{i}\bm{v}^{i}\left(\bm{v}^{i}\right)^{\top} is doubly stochastic, so we can decompose it into W=PΛP⊤W=P\Lambda P^{\top}, where P=(v1,v2,⋯ ,vn)P=\left(\bm{v}^{1},\bm{v}^{2},\cdots,\bm{v}^{n}\right) that satisfies P⊤P=PP⊤=IP^{\top}P=PP^{\top}=I. Without the loss of generality, we can assume λ1≥λ2≥⋯ ,≥λn\lambda_{1}\geq\lambda_{2}\geq\cdots,\geq\lambda_{n}. Then we have the following equalities:

Consider the coordinate transformation using PP as the base change matrix, and denote Yt=XtPY_{t}=X_{t}P, H(Xt;ξt)=G(Xt;ξt)PH(X_{t};\xi_{t})=G(X_{t};\xi_{t})P, Rt=QtPR_{t}=Q_{t}P. Then the above equation can be rewritten as

Since Λ\Lambda is a diagonal matrix, so we use yt(i)\bm{y}_{t}^{(i)}, ht(i)\bm{h}_{t}^{(i)}, rt(i)\bm{r}_{t}^{(i)} to indicate the ii-th column of YtY_{t}, H(Xt;ξt)H(X_{t};\xi_{t}), RtR_{t}. Then (8) becomes

(9) offers us a much intuitive way to analysis the algorithm. Since all eigenvalues of WW, except λ1\lambda_{1}, satisfies ∣λi∣<1|\lambda_{i}|<1, so the corresponding yt(i)\bm{y}_{t}^{(i)} would “decay to zero” due to the scaling factor λi\lambda_{i}.

Moreover, since the eigenvector corresponding to λ1\lambda_{1} is 1n(1,1,⋯ ,1)\frac{1}{\sqrt{n}}(1,1,\cdots,1), then we have yt(1)=X‾tn\bm{y}_{t}^{(1)}=\overline{X}_{t}\sqrt{n}. So, if t→∞t\to\infty, intuitively we can set yt(i)→0\bm{y}_{t}^{(i)}\to\bm{0} for i≠1i\neq 1, then Yt→(X‾tn,0,⋯ ,0)Y_{t}\to(\overline{X}_{t}\sqrt{n},0,\cdots,0) and Xt→X‾t1nX_{t}\to\overline{X}_{t}\frac{\bm{1}}{n}. This whole process shows how the confusion works under a coordinate transformation.

A.2 Analysis for the general updating form in (7)

where ρ\rho follows the defination in Theorem 3.

Since Wt=PΛtP⊤W^{t}=P\Lambda^{t}P^{\top}, we have

where yt(i)=XtPe(i)\bm{y}_{t}^{(i)}=X_{t}P\bm{e^{(i)}}. ∎

Given two non-negative sequences {at}t=1∞\{a_{t}\}_{t=1}^{\infty} and {bt}t=1∞\{b_{t}\}_{t=1}^{\infty} that satisfying

Lemma 5 shows us an overall understanding about how the confusion matrix works, while Lemma 6 is a very important tool for analyzing the sequence in Lemma 5. Next we are going to give a upper bound for the difference between the local modes and the global mean mode.

We can see that ∑s=1t−1ρ2(t−s)∥Qs∥F2\sum_{s=1}^{t-1}\rho^{2(t-s)}\left\|Q_{s}\right\|^{2}_{F} and ∑s=1t−1γsρt−s−1∥G(Xs;ξs)∥F\sum_{s=1}^{t-1}\gamma_{s}\rho^{t-s-1}\left\|G(X_{s};\xi_{s})\right\|_{F} has the same structure with the sequence in Lemma 6, which leads to

From the Lipschitzian condition for the objective function fif_{i}, we know that ff also satisfies the Lipschitzian condition. Then we have

Combining (14) and (15) together, we have

Appendix B Analysis for Algorithm 1

D1D_{1}, μ\mu and ρ\rho are defined in Theorem 1.

Under Assumption 1, when using Algorithm 1, we have

when (1−ρ)2−4μ2α2>0(1-\rho)^{2}-4\mu^{2}\alpha^{2}>0, where

In the proof below, we use [A](i,j)[A]^{(i,j)} to indicate the (i,j)(i,j) element of matrix AA.

For the noise induced by quantization, we have

As for Qt⊤QtQ_{t}^{\top}Q_{t}, the expectation of non-diagonal elements would be zero because the compression noise on node ii is independent on node jj, which leads to

where τij=1\tau_{ij}=1 if i=ji=j, else τij=0\tau_{ij}=0. Then

where vj(i)\bm{v}_{j}^{(i)} is the jjth element of v(i)\bm{v}^{(i)}. So we have

Denote mt(i)=∑s=1t−1∣λi∣t−1−s∥−hs(i)+rs(i)∥m_{t}^{(i)}=\sum_{s=1}^{t-1}|\lambda_{i}|^{t-1-s}\left\|-\bm{h}_{s}^{(i)}+\bm{r}_{s}^{(i)}\right\|, we can see that mt(i)m_{t}^{(i)} has the same structure of the sequence in Lemma 6, Therefore

Combining (19) and (20) together, we have

If α\alpha is small enough that satisfies (1−ρ)2−4μ2α2>0(1-\rho)^{2}-4\mu^{2}\alpha^{2}>0, then we have

Moreover, setting γt=γ\gamma_{t}=\gamma and denote 2α2(2μ2(1+2α2)(1−ρ)2−4μ2α2+1)=D22\alpha^{2}\left(\frac{2\mu^{2}(1+2\alpha^{2})}{(1-\rho)^{2}-4\mu^{2}\alpha^{2}}+1\right)=D_{2}. Applying Lemma 12 to bound ∥G(Xt;ξt)∥F2\|G\left(X_{t};\xi_{t}\right)\|^{2}_{F}, we have

when 1−2α2μ2Cλ>01-2\alpha^{2}\mu^{2}C_{\lambda}>0, where

Denote 2α21−ρ2(2μ2(1+2α2)(1−ρ)2−4μ2α2+1)+1(1−ρ)2=D1\frac{2\alpha^{2}}{1-\rho^{2}}\left(\frac{2\mu^{2}(1+2\alpha^{2})}{(1-\rho)^{2}-4\mu^{2}\alpha^{2}}+1\right)+\frac{1}{(1-\rho)^{2}}=D_{1},then from Lemma 12, we have

If γ\gamma is not too large that satisfies 1−3D1L2γ2>01-3D_{1}L^{2}\gamma^{2}>0, we have

Summarizing both sides of (23) and applying (24) yields

Denote D3=(4L2+3L3D2γ2)3D1γ21−3D1L2γ2+3LD2γ22D_{3}=\frac{\left(4L^{2}+3L^{3}D_{2}\gamma^{2}\right)3D_{1}\gamma^{2}}{1-3D_{1}L^{2}\gamma^{2}}+\frac{3LD_{2}\gamma^{2}}{2} and D4=1−LγD_{4}=1-L\gamma, we have

Setting γ=16D1L+6D2L+σnT12+ζ23T13\gamma=\frac{1}{6\sqrt{D_{1}}L+6\sqrt{D_{2}L}+\frac{\sigma}{\sqrt{n}}T^{\frac{1}{2}}+\zeta^{\frac{2}{3}}T^{\frac{1}{3}}}, then we can verify

So we can remove the ∥∇f‾(Xt)∥2\|\overline{\nabla f}(X_{t})\|^{2} on the LHS and substitute (1−D3)(1-D_{3}) with 12\frac{1}{2}. Therefore (2) becomes

If α2≤min⁡{(1−ρ)28μ2,14}\alpha^{2}\leq\min\left\{\frac{(1-\rho)^{2}}{8\mu^{2}},\frac{1}{4}\right\}, then

which means D2=O(α2)D_{2}=O\left(\alpha^{2}\right) and D1=O(α2+1)D_{1}=O\left(\alpha^{2}+1\right). So (26) becomes

Combing the inequality above and (26), we have

Appendix C Analysis for Algorithm 2

We are going to prove that by using (3) and (4) in Algorithm 2, the upper bound for the compression noise would be

Therefore, combing this with Lemma 8 and Lemma 7, we would be able to prove the convergence rate for ALgorithm 2.

which ensures that all nodes would converge to the same value.

For any non-negative sequences {an}n=1+∞\{a_{n}\}_{n=1}^{+\infty} and {bn}n=1+∞\{b_{n}\}_{n=1}^{+\infty} that satisfying

Under the Assumption 1, when using Algorithm 2, we have

Under Assumption 1, when using Algorithm 2, we have

Setting γt=γ\gamma_{t}=\gamma, then from Lemma 8, we have

If γ\gamma is not too large that satisfies 1−6C1L2γ2>01-6C_{1}L^{2}\gamma^{2}>0, we have

Summarizing both sides of (33) and applying (34) yields

where C2=11−6C1L2γ2C_{2}=\frac{1}{1-6C_{1}L^{2}\gamma^{2}}. It implies

where C3=12L2C2C1γ2C_{3}=12L^{2}C_{2}C_{1}\gamma^{2} and C4=(1−Lγ)C_{4}=\left(1-L\gamma\right). It completes the proof. ∎

when γ=112C1L+σnT12+ζ23T13\gamma=\frac{1}{12\sqrt{C_{1}}L+\frac{\sigma}{\sqrt{n}}T^{\frac{1}{2}}+\zeta^{\frac{2}{3}}T^{\frac{1}{3}}}, we have

Appendix D Why naive combination between compression and D-PSGD does not work?

Consider combine compression with the D-PSGD algorithm. Let the compression of exchanged models XtX_{t} be

This naive combination does not work, because the compression error QtQ_{t} does not diminish unlike the stochastic gradient variance that can be controlled by γt\gamma_{t} either decays to zero or is chosen to be small enough. ∎