Qsparse-local-SGD: Distributed SGD with Quantization, Sparsification, and Local Computations

Debraj Basu, Deepesh Data, Can Karakus, Suhas Diggavi

Introduction

Stochastic Gradient Descent (SGD) [HM51] and its many variants have become the workhorse for modern large-scale optimization as applied to machine learning [Bot10, BM11]. We consider a setup, in which SGD is applied to the distributed setting, where RR different nodes compute local stochastic gradients on their own datasets Dr\mathcal{D}_{r}. Co-ordination between them is done by aggregating these local computations to update the overall parameter xt\mathbf{x}_{t} as,

Training of high dimensional models is typically performed at a large scale over bandwidth limited networks. Therefore, despite the distributed processing gains, it is well understood by now that exchange of full-precision gradients between nodes causes communication to be the bottleneck for many large-scale models [AHJ+18, WXY+17, BWAA18, SYKM17]. For example, consider training the BERT architecture for language models [DCLT18] which has about 340 million parameters, implying that each full precision exchange between workers is over 1.3GB. Such a communication bottleneck could be significant in emerging edge computation architectures suggested by federated learning [Kon17, MMR+17, ABC+16]. In such an architecture, data resides on and can even be generated by personal devices such as smart phones, and other edge (IoT) devices, in contrast to data-center architectures. Learning is envisaged with such an ultra-large scale, heterogeneous environment, with potentially unreliable or limited communication. These and other applications have led to many recently proposed methods, which are broadly based on three major approaches:

Quantization of gradients, where nodes locally quantize the gradient (perhaps with randomization) to a small number of bits [AGL+17, BWAA18, WHHZ18, WXY+17, SYKM17].

Skipping communication rounds, whereby nodes average their models after locally updating their models for several steps [YYZ19, Cop15, ZDW13, Sti19, CH16, WJ18].

In this work we propose a Qsparse-local-SGD algorithm, which combines aggressive sparsification with quantization and local computations, along with error compensation, by keeping track of the difference between the true and compressed gradients. We propose both synchronous and asynchronous implementations of Qsparse-local-SGD in a distributed setting, where the nodes perform computations on their local datasets. In our asynchronous model, the distributed nodes’ iterates evolve at the same rate, but update the gradients at arbitrary times; see Section 4 for more details. We analyze convergence for Qsparse-local-SGD in the distributed case, for smooth non-convex and smooth strongly-convex objective functions. We demonstrate that Qsparse-local-SGD converges at the same rate as vanilla distributed SGD for many important classes of sparsifiers and quantizers. We implement Qsparse-local-SGD for ResNet-50 using the ImageNet dataset, and for a softmax multiclass classifier using the MNIST dataset, and we achieve target accuracies with about a factor of 15-20 savings over the state-of-the-art [AHJ+18, SCJ18, Sti19], in the total number of bits transmitted.

2 Contributions

We study a distributed set of RR worker nodes, each of which perform computations on locally stored data, denoted by Dr\mathcal{D}_{r}. Consider the empirical-risk minimization of the loss function

Our main theoretical results are the convergence analyses of Qsparse-local-SGD for both non-convex as well as convex objectives; see Theorem 1 and Theorem 3 for the synchronous case, as well as Theorem 4 and Theorem 6, for the asynchronous operation. Our analyses also demonstrate natural gains in convergence that distributed, mini-batch operation affords, and has convergence similar to equivalent vanilla SGD with local iterations (see Corollary 2 and Corollary 3), for both the non-convex case (with convergence rate ∼1T\sim\frac{1}{\sqrt{T}} for fixed learning rate) as well as the convex case (with convergence rate ∼1T\sim\frac{1}{T}, for diminishing learning rate). We also demonstrate that quantizing and sparsifying the gradient, even after local iterations asymptotically yields an almost “free” efficiency gain (also observed numerically in Section 5 non-asymptotically). The numerical results on ImageNet dataset implemented for a ResNet-50 architecture and for the convex case for multi-class logistic classification on MNIST [LBBH98] dataset demonstrates that one can get significant communication savings, while retaining equivalent state-of-the-art performance. The combination of quantization, sparsification, and local computations poses several challenges for theoretical analyses, including the analyses of impact of local iterations (block updates) of parameters on quantization and sparsification (see Lemma 4-5 in Section 3), as well as asynchronous updates and its combination with distributed compression (see Lemma 9-12 in Section 4).

3 Paper Organization

In Section 2, we demonstrate that composing certain classes of quantizers with sparsifiers satisfies a certain regularity condition that is needed for several convergence proofs for our algorithms. We describe the synchronous implementation of Qsparse-local-SGD in Section 3, and outline the main convergence results for it in Section 3.3, briefly giving the proof ideas in Section 3.4. We describe our asynchronous implementation of Qsparse-local-SGD and provide the theoretical convergence results in Section 4. The experimental results are given in Section 5. Many of the proof details are given in the appendices, given as part of the supplementary material.

Communication Efficient Operators

Traditionally, distributed stochastic gradient descent affords to send full precision (32 or 64 bit) unbiased gradient updates across workers to peers or to a central server that helps with aggregation. However, communication bottlenecks that arise in bandwidth limited networks limit the applicability of such an algorithm at a large scale when the parameter size is massive or the data is widely distributed on a very large number of worker nodes. In such settings, one could think of updates which not only result in convergence, but also require less bandwidth thus making the training process faster. In the following sections we discuss several useful operators from literature and enhance their use by proposing a novel class of composed operators.

SGD computes an unbiased estimate of the gradient, which can be used to update the model iteratively and is extremely useful in large scale applications. It is well known that the first order terms in the rate of convergence are affected by the variance of the gradients. While stochastic quantization of gradients could result in a variance blow up, it preserves the unbiasedness of the gradients at low precision; and, therefore, when training over bandwidth limited networks, the convergence would be much faster; see [AGL+17, WXY+17, SYKM17, ZDJW13].

Examples of randomized quantizers include

Stochastic Rotated Quantization [SYKM17], which is a stochastic quantization, preprocessed by a random rotation, with βd,s=2log⁡2(2d)s2\beta_{d,s}=\frac{2\log_{2}(2d)}{s^{2}}.

Instead of quantizing randomly into ss levels, we can take a deterministic approach and round off each component of the vector to the nearest level. In particular, we can just take the sign, which has shown promise in [BWAA18, KRSJ19].

Such methods drew interest since Rprop [RB93], which only used the temporal behavior of the sign of the gradient. This is an example where the biased 1-bit quantizer as in Definition 2 is used. This further inspired optimizers, such as RMSprop [TH12], Adam [KB15], which incorporate appropriate adaptive scaling with momentum acceleration and have demonstrated empirical superiority in non-convex applications.

2 Sparsification

where expectation is taken over the randomness of the compression operator CompkComp_{k}.

Note that stochastic quantizers, as defined in Definition 1, also satisfy this regularity condition in Definition 3 for βd,s≤1\beta_{d,s}\leq 1. Now we give a simple but important corollary, which allows us to apply different compression operators to different coordinates of a vector. As an application, in the case of training neural networks, we can apply different operators to different layers.

Corollary 1 allows us to apply different compression operators to different coordinates of the updates which can based upon their dimensionality and sparsity patterns.

3 Composition of Quantization and Sparsification

Now we show that we can compose deterministic/randomized quantizers with sparsifiers and the resulting operator is a compression operator. First we compose a general stochastic quantizer with an explicit sparsifier, such as Topk(x)\textrm{Top}_{k}({\bf x}) and Randk(x)\textrm{Rand}_{k}({\bf x}), and show that the resulting operator is a “compression” operator. A proof is provided in Appendix A.1.

where expectation is taken over the randomness of the compression operator CompkComp_{k} as well as of the quantizer QsQ_{s}.

For the different quantizers mentioned earlier, the conditions when their composition with CompkComp_{k} gives βk,s<1\beta_{k,s}<1 are:

QSGD: for k<s2k<s^{2}, we get. γ=(1−ks2)kd\gamma=\left(1-\frac{k}{s^{2}}\right)\frac{k}{d}

Stochastic k-level Quantization: for k<2s2k<2s^{2}, we get γ=(1−k2s2)kd\gamma=\left(1-\frac{k}{2s^{2}}\right)\frac{k}{d}.

Stochastic Rotated Quantization: for k<2s2/2−1k<2^{s^{2}/2-1}, we get γ=(1−2log⁡2(2k)s2)kd\gamma=\left(1-\frac{2\log_{2}(2k)}{s^{2}}\right)\frac{k}{d}.

Observe that for a given stochastic quantizer that satisfies Definition 1, we have a prescribed operating regime of βk,s<1\beta_{k,s}<1. This results in an upper bound on the coarseness of the quantizer, which happens because the quantization leads to a blow-up of the second moment; see condition (ii) of Definition 1. However, by employing Corollary 1, we show that this can be alleviated to some extent via an example.

As discussed above, stochastic quantization results in a variance blow-up which limits our regime of operation, when we combine that with sparsification. However, it turns out that, we can expand our regime of operation unrestrictedly by scaling the vector QsCompk(x)Q_{s}Comp_{k}({\bf x}) appropriately. We summarize the result in the following lemma, which is proved in Appendix A.2.

Note that, unlike QsCompk(x)Q_{s}Comp_{k}(\mathbf{x}), the scaled version QsCompk(x)1+βk,s\frac{Q_{s}Comp_{k}(\mathbf{x})}{1+\beta_{k,s}} is always a compression operator for all values of βk,s>0\beta_{k,s}>0. Furthermore, observe that, if βk,s<1\beta_{k,s}<1, then we have (1−βk,s)kd<kd(1+βk,s)(1-\beta_{k,s})\frac{k}{d}<\frac{k}{d(1+\beta_{k,s})}, which implies that even in the operating regime of βk,s<1\beta_{k,s}<1, which is required in Lemma 1, the scaled composed operator QsCompk(x)1+βk,s\frac{Q_{s}Comp_{k}(\mathbf{x})}{1+\beta_{k,s}} of Lemma 2 gives better compression than what we get from the unscaled composed operator QsCompk(x)Q_{s}Comp_{k}(\mathbf{x}) of Lemma 1. So, appropriately scaled composed operator is always a better choice for compression.

In the following lemma we show that SignCompkSignComp_{k} is a compression operator; a proof of which is provided in Appendix A.3.

Observe that for m=1m=1, depending on the value of kk, either of the terms inside the max can be bigger than the other term. For example, if k=1k=1, then ∥Compk(x)∥1=∥Compk(x)∥2\|Comp_{k}(\mathbf{x})\|_{1}=\|Comp_{k}(\mathbf{x})\|_{2}, which implies that the second term inside the max is equal to 1/d21/d^{2}, which is much smaller than the first term. On the other hand, if k=dk=d and the vector x\mathbf{x} is dense, then the second term may be much bigger than the first term.

Distributed Synchronous Operation

Let IT(r)⊆[T]:={1,…,T}\mathcal{I}_{T}^{(r)}\subseteq[T]:=\{1,\ldots,T\} with T∈IT(r)T\in\mathcal{I}_{T}^{(r)} denote a set of indices for which worker r∈[R]r\in[R] synchronizes with the master. In a synchronous setting, IT(r)\mathcal{I}_{T}^{(r)} is same for all the workers. Let IT:=IT(r)\mathcal{I}_{T}:=\mathcal{I}_{T}^{(r)} for any r∈[R]r\in[R]. Every worker r∈[R]r\in[R] maintains a local parameter vector x^t(r)\widehat{\bf x}_{t}^{(r)} which is updated in each iteration tt. If t∈ITt\in\mathcal{I}_{T}, every worker r∈[R]r\in[R] sends the compressed and error-compensated update gt(r)g_{t}^{(r)} computed on the net progress made since the last synchronization to the master node, and updates its local memory mt(r)m_{t}^{(r)}. Upon receiving gt(r),r=1,2.…,Rg_{t}^{(r)},r=1,2.\ldots,R, master aggregates them, updates the global parameter vector, and sends the new model xt+1{\bf x}_{t+1} to all the workers; upon receiving which, they set their local parameter vector x^t+1(r)\widehat{\bf x}_{t+1}^{(r)} to be equal to the global parameter vector xt+1{\bf x}_{t+1}. Our algorithm is summarized in Algorithm 1.

All results in this paper use the following two standard assumptions.

In this section we present our main convergence results with synchronous updates, obtained by running Algorithm 1 for smooth functions, both non-convex and strongly convex. To state our results, we need the following definition from [Sti19].

Let IT={t0,t1,…,tk}\mathcal{I}_{T}=\{t_{0},t_{1},\ldots,t_{k}\}, where ti<ti+1t_{i}<t_{i+1} for i=0,1,…,k−1i=0,1,\ldots,k-1. The gap of IT\mathcal{I}_{T} is defined as gap(IT):=max⁡i∈[k]{(ti−ti−1)}gap(\mathcal{I}_{T}):=\max_{i\in[k]}\{(t_{i}-t_{i-1})\}, which is equal to the maximum difference between any two consecutive synchronization indices.

2 Error Compensation

Sparsified gradient methods, where workers send the largest kk coordinates of the updates based on their magnitudes have been investigated in the literature and serves as a communication efficient strategy for distributed training of learning models. However, the convergence rates are subpar to distributed vanilla SGD. Together with some form of error compensation, these methods have been empirically observed to converge as fast as vanilla SGD in [Str15, AH17, LHM+18, AHJ+18, SCJ18]. In [AHJ+18, SCJ18], sparsified SGD with such feedback schemes has been carefully analyzed. Under analytic assumptions, [AHJ+18] proves the convergence of distributed Topk\textrm{Top}_{k} SGD with error feedback. The net error in the system is accumulated by each worker locally on a per iteration basis and this is used as feedback for generating the future updates. [SCJ18] did the analysis for the centralized Topk\textrm{Top}_{k} SGD for strongly convex objectives.

In Algorithm 1, the error introduced in every iteration is accumulated into the memory of each worker, which is compensated for in the future rounds of communication. This feedback is the key to recovering the convergence rates matching vanilla SGD. The operators employed provide a controlled way of using both the current update as well as the compression errors from the previous rounds of communication. Under the assumption of the uniform boundedness of the gradients, we analyze the controlled evolution of memory through the optimization process; the results are summarized in Lemma 4 and Lemma 5 below.

Here we show that if we run Algorithm 1 with a decaying learning rate ηt\eta_{t}, then the local memory at each worker contracts and goes to zero as O(ηt)2\mathcal{O}(\eta_{t})^{2}.

We prove Lemma 4 in Appendix B.1. Note that for fixed γ,H\gamma,H, the memory decays as O(ηt2)\mathcal{O}(\eta_{t}^{2}). This implies that the net error in the algorithm from the compression of updates in each round of communication is compensated for in the end.

2.2 Fixed Learning Rate

Note that, for fixed γ,H\gamma,H, the memory is upper bounded by a constant O(η2)\mathcal{O}(\eta^{2}). Observe that since the memory accumulates the past errors due to compression and local computation, in order to asymptotically reduce the memory to zero, the learning rate would have to be reduced once in a while throughout the training process.

3 Main Results

We leverage the perturbed iterate analysis as in [MPP+17, SCJ18] to provide convergence guarantees for Qsparse-local-SGD. Under the assumptions of Section 3.1, the following theorems hold when Algorithm 1 is run with any compression operator (including our composed operators).

Here zT\mathbf{z}_{T} is a random variable which samples a previous parameter x^t(r)\widehat{\mathbf{x}}_{t}^{(r)} with probability 1/RT1/RT.

In order to ensure that the compression does not affect the dominating terms while converging at a rate of O(1/bRT)\mathcal{O}\left(1/\sqrt{bRT}\right), we would require H=O(γT1/4/(bR)3/4)H=\mathcal{O}\left(\gamma T^{1/4}/(bR)^{3/4}\right).

Theorem 1 is proved in Appendix B.6 and provides non-asymptotic guarantees, where we observe that compression does not affect the first order term. Here, we are required to decide the horizon TT before running the algorithm. Therefore, in order to converge to a fixed point, the learning rate needs to follow a piecewise schedule (i.e., the learning rate would have to be reduced once in a while throughout the training process), which is also the case in our numerics in Section 5.1. The corresponding asymptotic result (with decaying learning rate) is given below.

Here (i) δt:=ηt4R\delta_{t}:=\frac{\eta_{t}}{4R}; (ii) PT:=∑t=0T−1∑r=1RδtP_{T}:=\sum_{t=0}^{T-1}\sum_{r=1}^{R}\delta_{t}, which is lower bounded as PT≥ξ4ln⁡(T+a−1a)P_{T}\geq\frac{\xi}{4}\ln{\left(\frac{T+a-1}{a}\right)}; and (iii) zT\mathbf{z}_{T} is a random variable which samples a previous parameter x^t(r)\widehat{\mathbf{x}}_{t}^{(r)} with probability δt/PT\delta_{t}/P_{T}.

Note that Theorem 2 gives a convergence rate of O(1log⁡T)\mathcal{O}(\frac{1}{\log T}). We prove it in Appendix B.7.

Here (i) A=∑r=1Rσr2bR2A=\frac{\sum_{r=1}^{R}\sigma_{r}^{2}}{bR^{2}} , B=4((3μ2+3L)CG2H2γ2+3L2G2H2)B=4\left(\left(\frac{3\mu}{2}+3L\right)\frac{CG^{2}H^{2}}{\gamma^{2}}+3L^{2}G^{2}H^{2}\right), where C≥4aγ(1−γ2)aγ−4HC\geq\frac{4a\gamma(1-\gamma^{2})}{a\gamma-4H}; (ii) x‾T:=1ST∑t=0T−1[wt(1R∑r=1Rx^t(r))]\overline{\mathbf{x}}_{T}:=\frac{1}{S_{T}}\sum_{t=0}^{T-1}\left[w_{t}\left(\frac{1}{R}\sum_{r=1}^{R}\widehat{\mathbf{x}}_{t}^{\left(r\right)}\right)\right], where wt=(a+t)2w_{t}=\left(a+t\right)^{2}; and (iii) ST=∑t=oT−1wt≥T33S_{T}=\sum_{t=o}^{T-1}w_{t}\geq\frac{T^{3}}{3}.

In order to ensure that the compression does not affect the dominating terms while converging at a rate of O(1/(bRT))\mathcal{O}\left(1/(bRT)\right), we would require H=O(γT/(bR))H=\mathcal{O}\left(\gamma\sqrt{T/(bR)}\right).

Theorem 3 is proved in Appendix B.8. For no compression and only local computations, i.e., for γ=1\gamma=1, and under the same assumptions, we recover/generalize a few recent results from literature with similar convergence rates:

We recover [YYZ19, Theorem 1], which does local SGD for the non-convex case;

We generalize [Sti19, Theorem 2.2], which does local SGD for a strongly convex case and requires the unbiasedness assumption of gradients,The unbiasedness of gradients at every worker can be ensured by assuming that each worker samples data points from the entire dataset. to the distributed case.

We emphasize that unlike [YYZ19, Sti19], which only consider local computation, we combine quantization and sparsification with local computation, which poses several technical challenges; e.g., see proofs of Lemma 4, 5, 6.

4 Proof Outlines

In order to prove our results, we define virtual sequences for every worker r∈[R]r\in[R] and for all t≥0t\geq 0 as follows:

Here ηt\eta_{t} can be taken to be decaying or fixed, depending on the result that we are proving. Let iti_{t} be the set of random sampling of the mini-batches at each worker {it(1),it(2),…,it(R)}\{i_{t}^{(1)},i_{t}^{(2)},\ldots,i_{t}^{(R)}\}. We define

x~t+1:=1R∑r=1Rx~t+1(r)=x~t−ηtpt\widetilde{\mathbf{x}}_{t+1}:=\frac{1}{R}\sum_{r=1}^{R}\widetilde{\mathbf{x}}_{t+1}^{\left(r\right)}=\widetilde{\mathbf{x}}_{t}-\eta_{t}\mathbf{p}_{t}, x^t:=1R∑r=1Rx^t(r)\quad\widehat{\mathbf{x}}_{t}:=\frac{1}{R}\sum_{r=1}^{R}\widehat{\mathbf{x}}_{t}^{\left(r\right)}.

Since ff is LL-smooth, we have from (4) (with fixed learning rate ηt=η\eta_{t}=\eta) that

With some algebraic manipulations provided in Appendix B.6, for η≤\nicefrac12L\eta\leq\nicefrac{{1}}{{2L}}, we arrive at

Under the Assumption 2, stated in Section 3.1, we have

Let x~t(r),mt(r)\widetilde{\bf x}_{t}^{(r)},m_{t}^{(r)}, r∈[R]r\in[R], t≥0t\geq 0 be generated according to Algorithm 1 and let x^t(r)\widehat{\bf x}_{t}^{(r)} be as defined in (4). Let x~t=1R∑r=1Rx~t(r)\widetilde{\bf x}_{t}=\frac{1}{R}\sum_{r=1}^{R}\widetilde{\bf x}_{t}^{(r)} and x^t=1R∑r=1Rx^t(r)\widehat{\bf x}_{t}=\frac{1}{R}\sum_{r=1}^{R}\widehat{\bf x}_{t}^{(r)}. Then we have

i.e., the difference of the true and the virtual sequence is equal to the average memory.

The last term on the RHS of (6) depicts the deviation of the local sequences x~t(r)\widetilde{\bf x}_{t}^{(r)} from the global sequence x~t\widetilde{\bf x}_{t} which can be bounded as shown in Lemma 7. The details are provided in Appendix B.4.

Let gap(IT)≤Hgap(\mathcal{I}_{T})\leq H. For x^t(r)\widehat{{\bf x}}_{t}^{(r)} generated according to Algorithm 1 with a fixed learning rate η\eta and letting x^t=1R∑r=1Rx^t(r)\widehat{\bf x}_{t}=\frac{1}{R}\sum_{r=1}^{R}\widehat{\bf x}_{t}^{(r)}, we have the following bound on the deviation of the local sequences:

Substituting the bounds from (7)-(9) into (6) yields

Performing a telescopic sum from t=0t=0 to T−1T-1 and dividing by ηT4\frac{\eta T}{4} gives

By letting η=C^/T\eta=\widehat{C}/\sqrt{T}, where C^\widehat{C} is a constant such that C^T≤12L\frac{\widehat{C}}{\sqrt{T}}\leq\frac{1}{2L}, we arrive at bound stated in Theorem 1. ∎

4.2 Proof Outline of Theorem 2

Observe that (6) holds irrespective of the learning rate schedule, as long as learning rate is at most \nicefrac12L\nicefrac{{1}}{{2L}}; see Appendix B.7 for details. Here ηt≤12L\eta_{t}\leq\frac{1}{2L} follows from our assumption that a≥2ξLa\geq 2\xi L. Substituting a decaying learning rate ηt\eta_{t} (such that ηt≤\nicefrac12L\eta_{t}\leq\nicefrac{{1}}{{2L}} holds for every t≥0t\geq 0) in (6) gives

The last term on the RHS of (12) is the deviation of local sequences and we bound it in Lemma 8 for decaying learning rates. The details are provided in Appendix B.5.

Let gap(IT)≤Hgap(\mathcal{I}_{T})\leq H. By running Algorithm 1 with a decaying learning rate ηt\eta_{t}, we have

Let δt:=ηt4R\delta_{t}:=\frac{\eta_{t}}{4R} and PT:=∑t=0T−1∑r=1RδtP_{T}:=\sum_{t=0}^{T-1}\sum_{r=1}^{R}\delta_{t}. Performing a telescopic sum from t=0t=0 to T−1T-1 and dividing by PTP_{T} gives

In (15), we used the following bounds, which are shown in Appendix B.7: PT≥ξ4ln⁡(T+a−1a)P_{T}\geq\frac{\xi}{4}\ln{\left(\frac{T+a-1}{a}\right)}, ∑t=0T−1ηt2≤ξ2a−1\sum_{t=0}^{T-1}\eta_{t}^{2}\leq\frac{\xi^{2}}{a-1}, and ∑t=0T−1ηt3≤ξ32(a−1)2\sum_{t=0}^{T-1}\eta_{t}^{3}\leq\frac{\xi^{3}}{2(a-1)^{2}}. This completes the proof of Theorem 2. ∎

4.3 Proof Outline of Theorem 3

Using the definition of virtual sequences (4) that, we have

Note that the bounds in (13) and (14) hold irrespective of whether the function is convex or not. So, we can use them here as well in (17), which gives

For ηt=8μ(a+t)\eta_{t}=\frac{8}{\mu\left(a+t\right)} and wt=(a+t)2w_{t}=\left(a+t\right)^{2}, ST=∑t=oT−1≥T33S_{T}=\sum_{t=o}^{T-1}\geq\frac{T^{3}}{3}, we have

Where x‾T:=1ST∑t=0T−1[wt(1R∑r=1Rx^t(r))]=1ST∑t=0T−1wtx^t\overline{\mathbf{x}}_{T}:=\frac{1}{S_{T}}\sum_{t=0}^{T-1}\left[w_{t}\left(\frac{1}{R}\sum_{r=1}^{R}\widehat{\mathbf{x}}_{t}^{\left(r\right)}\right)\right]=\frac{1}{S_{T}}\sum_{t=0}^{T-1}w_{t}\widehat{\mathbf{x}}_{t}. This completes the proof of Theorem 3. ∎

Distributed Asynchronous Operation

We propose and analyze a particular form of asynchronous operation, where the workers synchronize with the master at arbitrary times decided locally or by master picking a subset of nodes as in federated learning [Kon17, MMR+17]. However, the local iterates evolve at the same rate, i.e., each worker takes the same number of steps per unit time according to a global clock. The asynchrony is therefore that updates occur after different number of local iterations but the local iterations are in synchrony with respect to the global clock. This is different from asynchronous algorithms studied for stragglers [WYL+18, RRWN11], where only one gradient step is taken but occurs at different times due to delays.

In this asynchronous setting, IT(r)\mathcal{I}_{T}^{(r)}’s may be different for different workers. However, we assume that gap(IT(r))≤Hgap(\mathcal{I}_{T}^{(r)})\leq H holds for every r∈[R]r\in[R], which means that there is a uniform bound on the maximum delay in each worker’s update times. The algorithmic difference from Algorithm 1 is that, in this case, a subset of workers (including a single worker) can send their updates to the master at their synchronization time steps; master aggregates them, updates the global parameter vector, and sends that only to those workers. Our algorithm is summarized in Algorithm 2

In this section we present our main convergence results with asynchronous updates, obtained by running Algorithm 2 for smooth objectives, both non-convex and strongly convex. Under the same assumptions as in the synchronous setting of Section 3.1, the following theorems hold even if Algorithm 2 is run with an arbitrary compression operators (including our composed operators from Section 2.3), whose compression coefficient is equal to γ\gamma.

Under the same conditions as in Theorem 1 with gap(IT(r))≤Hgap(\mathcal{I}_{T}^{(r)})\leq H, if {x^t(r)}t=0T−1\{\widehat{x}_{t}^{(r)}\}_{t=0}^{T-1} is generated according to Algorithm 2, the following holds.

Here (i) C1=(8γ2−6)(4−2γ)C_{1}=(\frac{8}{\gamma^{2}}-6)(4-2\gamma); (ii) zT\mathbf{z}_{T} is a random variable which samples a previous parameter x^t(r)\widehat{\mathbf{x}}_{t}^{(r)} with probability 1/RT1/RT; and (iii) C^\widehat{C} is a constant such that C^T≤12L\frac{\widehat{C}}{\sqrt{T}}\leq\frac{1}{2L}.

In order to ensure that the compression does not affect the dominating terms while converging at a rate of O(1/bRT)\mathcal{O}\left(1/\sqrt{bRT}\right), we would require H=O(γT1/8/(bR)3/8)H=\mathcal{O}\left(\sqrt{\gamma}T^{1/8}/(bR)^{3/8}\right).

Theorem 4 provides non asymptotic guarantees where we also observe that the compression comes for “free". The corresponding asymptotic result is given below.

Under the same conditions as in Theorem 2 with gap(IT(r))≤Hgap(\mathcal{I}_{T}^{(r)})\leq H, if {x^t(r)}t=0T−1\{\widehat{x}_{t}^{(r)}\}_{t=0}^{T-1} is generated according to Algorithm 2, the following holds.

Here (i) δt:=ηt4R\delta_{t}:=\frac{\eta_{t}}{4R} and PT:=∑t=0T−1∑r=1RδtP_{T}:=\sum_{t=0}^{T-1}\sum_{r=1}^{R}\delta_{t}, which is lower bounded as PT≥ξ4ln⁡(T+a−1a)P_{T}\geq\frac{\xi}{4}\ln{\left(\frac{T+a-1}{a}\right)}; (ii) C′=(4−2γ)(1+Cγ2)C^{\prime}=(4-2\gamma)(1+\frac{C}{\gamma^{2}}); and (iii) zT\mathbf{z}_{T} is a random variable which samples a previous parameter x^t(r)\widehat{\mathbf{x}}_{t}^{(r)} with probability δt/PT\delta_{t}/P_{T}.

Under the same conditions as in Theorem 3 with gap(IT(r))≤Hgap(\mathcal{I}_{T}^{(r)})\leq H, if {x^t(r)}t=0T−1\{\widehat{x}_{t}^{(r)}\}_{t=0}^{T-1} is generated according to Algorithm 2, the following holds.

Here (i) C≥4aγ(1−γ2)aγ−4HC\geq\frac{4a\gamma(1-\gamma^{2})}{a\gamma-4H}, C1=192(4−2γ)(1+Cγ2)C_{1}=192(4-2\gamma)\left(1+\frac{C}{\gamma^{2}}\right), C2=8(4−2γ)(1+Cγ2)C_{2}=8(4-2\gamma)(1+\frac{C}{\gamma^{2}}); (ii) A=∑r=1Rσr2bR2A=\frac{\sum_{r=1}^{R}\sigma_{r}^{2}}{bR^{2}}, D=(3μ2+3L)(12CG2H2γ2+C1ηt2H4G2)+24(1+C2H2)LG2H2D=\left(\frac{3\mu}{2}+3L\right)(\frac{12CG^{2}H^{2}}{\gamma^{2}}+C_{1}\eta_{t}^{2}H^{4}G^{2})+24(1+C_{2}H^{2})LG^{2}H^{2}; and (iii) x‾T\overline{\mathbf{x}}_{T}, STS_{T} are as defined in Theorem 3.

Under the same conditions as in Theorem 3 with gap(IT(r))≤Hgap(\mathcal{I}_{T}^{(r)})\leq H, a>max⁡{4Hγ,32κ,H}a>\max\{\frac{4H}{\gamma},32\kappa,H\}, σmax=max⁡r∈[R]σr\sigma_{max}=\max_{r\in[R]}\sigma_{r}, if {x^t(r)}t=0T−1\{\widehat{{\bf x}}_{t}^{(r)}\}_{t=0}^{T-1} is generated according to Algorithm 2, the following holds:

where x‾T\overline{\mathbf{x}}_{T}, STS_{T} are as defined in Theorem 3. In order to ensure that the compression does not affect the dominating terms while converging at a rate of O(1/(bRT))\mathcal{O}\left(1/(bRT)\right), we would require H=O(γ(T/(bR))1/4)H=\mathcal{O}\left(\sqrt{\gamma}(T/(bR))^{1/4}\right).

2 Proof Outlines

Let gap(IT(r))≤Hgap(\mathcal{I}_{T}^{(r)})\leq H holds for every r∈[R]r\in[R]. For x^t(r)\widehat{{\bf x}}_{t}^{(r)} generated according to Algorithm 2 with decaying learning rate ηt\eta_{t} and letting x^t=1R∑r=1Rx^t(r)\widehat{\bf x}_{t}=\frac{1}{R}\sum_{r=1}^{R}\widehat{\bf x}_{t}^{(r)}, we have the following bound on the deviation of the local sequences:

where C′′=8(4−2γ)(1+Cγ2)C^{\prime\prime}=8(4-2\gamma)(1+\frac{C}{\gamma^{2}}) and CC is a constant satisfying C≥4aγ(1−γ2)aγ−4HC\geq\frac{4a\gamma(1-\gamma^{2})}{a\gamma-4H}.

Let gap(IT(r))≤Hgap(\mathcal{I}_{T}^{(r)})\leq H holds for every r∈[R]r\in[R]. By running Algorithm 2 with fixed learning rate η\eta, we have

where C′=(16γ2−12)(4−2γ)C^{\prime}=(\frac{16}{\gamma^{2}}-12)(4-2\gamma).

Let gap(IT(r))≤Hgap(\mathcal{I}_{T}^{(r)})\leq H holds for every r∈[R]r\in[R]. If we run Algorithm 2 with a decaying learning rate ηt\eta_{t}, then we have the following bound on the difference between the true and virtual sequences:

where C′=192(4−2γ)(1+Cγ2)C^{\prime}=192(4-2\gamma)\left(1+\frac{C}{\gamma^{2}}\right) and CC is a constant satisfying C≥4aγ(1−γ2)aγ−4HC\geq\frac{4a\gamma(1-\gamma^{2})}{a\gamma-4H}.

Let gap(IT(r))≤Hgap(\mathcal{I}_{T}^{(r)})\leq H holds for every r∈[R]r\in[R]. If we run Algorithm 2 with a fixed learning rate η\eta, we have

where C′=(4−2γ)(8γ2−6)C^{\prime}=(4-2\gamma)\left(\frac{8}{\gamma^{2}}-6\right).

Now we give a brief summary of our convergence results in the synchronous as well as asynchronous settings.

In the synchronous setting, Qsparse-local-SGD asymptotically converges as fast as distributed vanilla SGD for H=O(γT1/4/(bR)3/4)H=\mathcal{O}\left(\gamma T^{1/4}/(bR)^{3/4}\right) in the smooth and non-convex case and for H=O(γT/(bR))H=\mathcal{O}\left(\gamma\sqrt{T/(bR)}\right) in the strongly convex case.

In the asynchronous setting, Qsparse-local-SGD asymptotically converges as fast as distributed vanilla SGD for H=O(γT1/8/(bR)3/8)H=\mathcal{O}(\sqrt{\gamma}T^{1/8}/(bR)^{3/8}) in the smooth and non-convex case and for H=O(γ(T/(bR))1/4)H=\mathcal{O}(\sqrt{\gamma}(T/(bR))^{1/4}) in the strongly convex case.

Therefore, our algorithm provides a lot of flexibility in terms of different ways of mitigating the communication bottleneck. For example, by increasing the batch size on each node, or by increasing the maximum synchronization period HH up to allowable limits. Furthermore, one could also choose to opt for different values of kk for the Topk\textrm{Top}_{k} sparsifier, as well as adjust the configurations of the quantizer. We present numerics in Section 5 demonstrating significant savings in the number of bits exchanged over the state-of-the-art.

Experimental Results

In this section we give extensive experimental results for validating our theoretical findings.

We train ResNet-50 [HZRS16] (which has d=25,610,216d=25,610,216 parameters) on ImageNet dataset, using 8 NVIDIA Tesla V100 GPUs. We use a learning rate schedule consisting of 5 epochs of linear warmup, followed by a piecewise decay of 0.1 at epochs 30, 60 and 80, with a batch size of 256 per GPU. For the purpose of experiments, we focus on SGD with momentum of 0.9, applied on the local iterations of the workers. We build our compression scheme into the Horovod framework [SB18]. We use SignTopkSignTop_{k} as defined in Lemma 3 and QTopkQTop_{k} as defined in Lemma 1 (which has an operating regime βk,s<1\beta_{k,s}<1), where QQ is from [AGL+17], as our composed operators.Even though the “scaled” QTopkQTop_{k} from Lemma 2 (with a scaling factor of (1+βk,s)(1+\beta_{k,s})) works with all values of βk,s\beta_{k,s}, and also does better than the “unscaled” QTopkQTop_{k} from Lemma 1 even when βk,s<1\beta_{k,s}<1 (see Remark 2), we report our experimental results in the non-convex setting only with unscaled QTopkQTop_{k}. We give some plots for a comparison on both these operators in Appendix D and observe that our algorithm with the unscaled operator gives at least as good performance as it gives with the scaled operator. We can attribute this to the fact that scaling the composed operator is a sufficient condition to obtain better convergence results, which does not necessarily mean that in practice also it does better. In TopkTop_{k}, we only update kt=min⁡(dt,1000)k_{t}=\min(d_{t},1000) elements per step for each tensor tt, where dtd_{t} is the number of elements in the tensor. For ResNet-50 architecture, this amounts to updating a total of k=99,400k=99,400 elements per step.

1.2 Results

From Figure 1(a), we observe that quantization and sparsification, both individually and combined, when error compensation is enabled through accumulating errors, has almost no penalty in terms of convergence rate, with respect to vanilla SGD. We observe that both QTopkQTop_{k}-SGDSGD, which employs a 4 bit quantizer and the TopkTop_{k} sparsifier, as well as SignTopkSignTop_{k}-SGDSGD, which employs the 1 bit sign quantizer and the TopkTop_{k} sparsifier, demonstrate superior performance over other schemes, both in terms of the required number of communicated bits for achieving certain target loss as well as test accuracy. This is because, in QTopkQTop_{k}, the QQ operator from [AGL+17] further induces sparsity, which results in fewer than kk coordinates being transmitted, and in SignTopkSignTop_{k}, we send only 1 bit for each TopkTop_{k} coordinate.

In Figure 2(a)-2(d), we show how the performance of different methods (used in Figure 1(a)-1(d)) change when we incorporate local iterations on top of them. Observe that the incorporation of local iterations in Figure 2(a) and 2(c) has very little impact on the convergence rates, as compared to vanilla SGD with the corresponding number of local iterations. Furthermore, this provides an added advantage over the Qsparse operator, in terms of savings in communicated bits for achieving target loss as seen in Figure 2(b) and 2(d), by a factor of 6 to 8 times on average.

Figure 3(b), Figure 3(c), and Figure 3(d) show the training loss, top-1, and top-5 convergence ratesHere top-i refers to the accuracy of the top i predictions by the model from the list of possible classes, see [LHS15]. respectively, with respect to the total number of bits of communication used. We observe that Qsparse-local-SGD combines the bit savings of either the deterministic sign based operator or the stochastic quantizer (QSGD), and aggressive sparsifier along with infrequent communication, thereby, outperforming the cases where these techniques are individually used. In particular, the required number of bits to achieve the same loss or top-1 accuracy in the case of Qsparse-local-SGD is around 1/16 in comparison with TopkTop_{k}-SGDSGD and over 1000×\times less than vanilla SGD. This also verifies that error compensation through memory can be used to mitigate not only the missing components from updates in previous synchronization rounds, but also explicit quantization error.

2 Convex Objective

The experiments in Figure 4-6 are in a synchronous distributed setting with 15 worker nodes, each processing a mini-batch size of 8 samples per iteration using the MNIST [LBBH98] handwritten digits dataset. The corresponding experiments for the asynchronous operation (as in Algorithm 2) are shown in Figure 7.

and z(i)z^{(i)} for every i∈[L]i\in[L] are the biases to be learnt corresponding to every class. We set λ\lambda to be 1/n1/n.

2.2 Parameter Selection and Learning Rates

We use the deterministic operator as in Lemma 3 and the stochastic operator QSGDQSGD denoted by QQ, as defined in [AGL+17], as our quantizers and Topk{Top}_{k} with error compensation as the sparsifier. The schemes with which we compare our composed operators QTopkQTop_{k} Lemma 2, and SignTopkSign{Top}_{k} Lemma 3, are ef-QSGD[WHHZ18], ef-signSGD[KRSJ19], TopK-SGD [SCJ18, AHJ+18], and local SGD [Sti19]. The learning rate used for training is of the form cλ(a+t)\frac{c}{\lambda(a+t)}, where (i) λ\lambda is the regularization parameter; (ii) cc is set with a careful hyperparameter sweep; (iii) wt=(a+t)2w_{t}=(a+t)^{2} as in Theorem 3, where aa is set as dHk\frac{dH}{k} with dd being the dimension of the gradient vector (7850 for MNIST); (iv) k=40k=40 is the sparsity; (v) HH is the synchronization period; (vi) tt is the iteration index; (vii) b=8b=8 is the batch size; and (viii) R=15R=15 is the number of workers.

2.3 Results

In Figure 4(a), we observe that the composition of a quantizer with a sparsifier has very little effect on the rate of convergence as compared to when the techniques are used individually. Observe that the algorithm run with the 2 bit QSGDQSGD is slower than the 4 bit quantizer, both with or without sparsification, which can be attributed to the reduction in the compression coefficient γ\gamma in going from 4 to 2 bits; see Theorem 3. From Figure 4(b) and 4(c), we see that our composed operators achieve gains in communicated bits by a factor of 6-8 times over the state-of-the-art.

Figure 5(a), demonstrates the effect of incorporating local iterations together with Qsparse operators, and we see that the rate of convergence is not significantly affected as we go from 1 to 8 local iterations. Furthermore, observe that for a fixed number of local iterations, the Qsparse operator maintains the same rates as vanilla SGDSGD or TopkTop_{k}-SGDSGD. In doing so, it is able to achieve gains in communicated bits as seen in Figure 5(b), simply by communicating infrequently with the master. On comparing Figure 5(c) and 5(e), we observe that the QTopkQTop_{k} operator is more sensitive to the increase in local computations for coarser quantizers (smaller values of ss, in this case s=2#−bits−1=3s=2^{\#-bits}-1=3). This can be verified from Figure 5(e) which uses a 4 bit quantizer (which implies s=15s=15 instead), and the corresponding effect of local iterations on the convergence rate is less prominent. We make comparisons between vanilla SGDSGD, QSGDQSGD with error accumulation (ef-QSGD) and our QTopkQTop_{k} operator in Figure 5(c) and 5(d), for which we do not observe much difference in performance between a finer and coarser quantizer, even though the convergence rates with respect to iterations are affected. This can be attributed to the precision of the quantizer itself.

In Figure 6(a) and Figure 6(b), we compare the convergence of our proposed scheme in Algorithm 1 with QTopkQTop_{k} and SignTopkSign{Top}_{k} being the composed operators, with vanilla SGD (32 bit floating point), ef-QSGD, ef-signSGD [KRSJ19], and TopK-SGD [SCJ18, AHJ+18]. Both figures follow a similar trend, where we observe QTopkQTop_{k}, SignTopkSign{Top}_{k} and TopK-SGD to be converging at the same rate as that of vanilla SGD, which is similar to the observations in [SCJ18]. This implies that the composition of quantization with sparsification does not affect the convergence while achieving improved communication efficiency, as can be seen in Figure 6(c) and Figure 6(b). Figure 6(c) shows that for test error approximately 0.1, Qsparse-local-SGD combines the benefits of the composed operator SignTopkSign{Top}_{k} or QTopkQTop_{k}, with local computations, and needs 10-15 times less bits than TopK-SGD and 1000×\times less bits than vanilla SGD.

We observe similar trends in Figure 7(a)-7(b) for our asynchronous operation, where workers synchronize with the master at arbitrary time intervals as per Algorithm 2. Specifically, in our experiments, for each r∈[R]r\in[R], the time interval for the rrth worker is decided uniformly at random from [H][H] after every synchronization by that worker. This ensures that gap(IT(r))≤Hgap(\mathcal{I}_{T}^{(r)})\leq H holds for every worker r∈[R]r\in[R] and the schedule IT(r)\mathcal{I}_{T}^{(r)} is different for each of them.

Conclusion

In this paper, we propose a gradient compression scheme that composes both unbiased and biased quantization with aggressive sparsification. Furthermore, we incorporate local computations, which, when combined with quantization and explicit sparsification, results in a highly communication efficient distributed algorithm, which we call Qsparse-local-SGD. We developed convergence analyses of our scheme in both synchronous as well as asynchronous settings and for both convex and non-convex objectives, and we show that our proposed algorithm achieves the same rate as that of distributed vanilla SGD in each of these cases. Our schemes provide flexibility in terms of different options for mitigating the communication bottlenecks that arise in training high-dimensional learning models over bandwidth limited networks. When run without compression, this also subsumes/generalizes several recent results from the literature on local SGD, with similar convergence rates, as mentioned at the end of Section 3.3.

Our numerics incorporate momentum acceleration, whose analysis is a topic for future research (e.g., potentially by incorporating ideas from [YJY19]). Although we use momentum for each local iteration, our preliminary results suggest that our method works with momentum applied to a block of updates as well though it was not the main focus of this paper.

Acknowledgement

The authors gratefully thank Navjot Singh for his help with experiments in the early stages of this work. This work was partially supported by NSF grant #1514531, by UC-NL grant LFR-18-548554 and by Army Research Laboratory under Cooperative Agreement W911NF-17-2-0196. The views and conclusions contained in this document are those of the authors and should not be interpreted as representing the official policies, either expressed or implied, of the Army Research Laboratory or the U.S. Government. The U.S. Government is authorized to reproduce and distribute reprints for Government purposes notwithstanding any copyright notation here on.

Appendix A Omitted Details from Section 2

where expectation is taken over the randomness of the compression operator CompkComp_{k} as well as the quantizer QsQ_{s}.

A.2 Proof of Lemma 2

A.3 Proof of Lemma 3

For proving Lemma 3 we first state and prove Lemma 13 below.

Appendix B Omitted Details from Section 3

Here (a) is due to the compression property, (b) holds since the memory and master parameter remain unchanged between two rounds of synchronization, and in (c) we used that x^t(i)(r)=xt(i)\widehat{\bf x}_{t_{(i)}}^{(r)}={\bf x}_{t_{(i)}}, which holds for every rr. Using the inequality ∥a+b∥2≤(1+τ)∥a∥2+(1+1τ)∥b∥2\|{\bf a}+{\bf b}\|^{2}\leq(1+\tau)\|{\bf a}\|^{2}+(1+\tfrac{1}{\tau})\|{\bf b}\|^{2}, which holds for every τ>0\tau>0, in (30) gives (take any p>1p>1 in the following):

Base case (i=1)(i=1): Note that mt(1)−1(r)=m0(r)=0m_{t_{(1)}-1}^{(r)}=m_{0}^{(r)}={\bf 0}. Consider the following:

The last inequality holds whenever β≥2H\beta\geq 2H. ∎

B.2 Proof of Lemma 5

Observe that (31) holds irrespective of the learning rate schedule. In particular, using a fixed learning rate ηt=η\eta_{t}=\eta for every tt gives

When rolled out we see that the memory is upper bounded by a geometric sum.

Note that the last inequality holds for every p>1p>1, and is minimized when p=2p=2. By plugging p=2p=2, we get

B.3 Proof of Lemma 6

Let x~t(r),mt(r)\widetilde{\bf x}_{t}^{(r)},m_{t}^{(r)}, r∈[R]r\in[R], t≥0t\geq 0 be generated according to Algorithm 1 and let x^t(r)\widehat{\bf x}_{t}^{(r)} be as defined in (4). Let x~t=1R∑r=1Rx~t(r)\widetilde{\bf x}_{t}=\frac{1}{R}\sum_{r=1}^{R}\widetilde{\bf x}_{t}^{(r)} and x^t=1R∑r=1Rx^t(r)\widehat{\bf x}_{t}=\frac{1}{R}\sum_{r=1}^{R}\widehat{\bf x}_{t}^{(r)}. Then we have

i.e., the difference of the true and the virtual sequence is equal to the average memory.

Now consider x^t−x~t=1R∑r=1Rx^t(r)−x~t(r)\widehat{\mathbf{x}}_{t}-\widetilde{\mathbf{x}}_{t}=\frac{1}{R}\sum_{r=1}^{R}\widehat{\mathbf{x}}_{t}^{(r)}-\widetilde{\mathbf{x}}_{t}^{(r)}. For the nearest tr+1∈ITt_{r}+1\in\mathcal{I}_{T} such that tr+1≤tt_{r}+1\leq t and the nearest tr′+1∈ITt_{r}^{\prime}+1\in\mathcal{I}_{T} such that tr′+1≤trt_{r}^{\prime}+1\leq t_{r}

Here we used that x^tr′+1(r)−x^tr+12(r)=∑j=tr′+1trηj∇(r)f(ij)(x^j(r))\widehat{\mathbf{x}}_{t_{r}^{\prime}+1}^{(r)}-\widehat{\mathbf{x}}_{t_{r}+\frac{1}{2}}^{(r)}=\overset{t_{r}}{\underset{j=t_{r}^{\prime}+1}{\sum}}\eta_{j}\nabla^{\left(r\right)}f_{\left(i_{j}\right)}\left(\widehat{\mathbf{x}}_{j}^{\left(r\right)}\right). Substituting x^tr′+1(r)=xtr′+1\widehat{\mathbf{x}}_{t_{r}^{\prime}+1}^{(r)}=\mathbf{x}_{t_{r}^{\prime}+1} we get

Now since xtr′+1=xtr\mathbf{x}_{t_{r}^{\prime}+1}=\mathbf{x}_{t_{r}} we have

On rolling out the expression in (36) we get

Therefore x^t−x~t=1R∑r=1Rmt(r)\widehat{\mathbf{x}}_{t}-\widetilde{\mathbf{x}}_{t}=\frac{1}{R}\sum_{r=1}^{R}m_{t}^{\left(r\right)} is the average memory. This completes the proof of Lemma 6. ∎

B.4 Proof of Lemma 7

Let gap(IT)≤Hgap(\mathcal{I}_{T})\leq H. For x^t(r)\widehat{{\bf x}}_{t}^{(r)} generated according to Algorithm 1 with a fixed learning rate η\eta and letting x^t=1R∑r=1Rx^t(r)\widehat{\bf x}_{t}=\frac{1}{R}\sum_{r=1}^{R}\widehat{\bf x}_{t}^{(r)}, we have the following bound on the deviation of the local sequences:

B.5 Proof of Lemma 8

Let gap(IT)≤Hgap(\mathcal{I}_{T})\leq H. By running Algorithm 1 with a decaying learning rate ηt\eta_{t}, we have

The last inequality (39) uses ηtr≤2ηtr+H≤2ηt\eta_{t_{r}}\leq 2\eta_{t_{r}+H}\leq 2\eta_{t} and t−tr≤Ht-t_{r}\leq H. ∎

B.6 Proof of Theorem 1

Let x∗{\bf x}^{*} be the minimizer of f(x)f({\bf x}), therefore we denote f(x∗)f({\bf x}^{*}) by f∗f^{*}. For the purpose of reusing the proof later while proving Theorem 2, we start off with the decaying learning rate ηt\eta_{t} until (43) and then switch to the fixed learning rate η\eta. Note that the proof remains the same until (43) irrespective of the learning rate schedule; in particular, we can take ηt=η\eta_{t}=\eta and the same proof holds until (43).

By the definition of LL-smoothness, we have

Define iti_{t} as the set of random sampling of the mini-batches at each worker {it(1),it(2),…,it(R)}\{i_{t}^{(1)},i_{t}^{(2)},\ldots,i_{t}^{(R)}\}. Taking expectation w.r.t. the sampling at time tt (conditioned on the past) and using the lipschitz continuity of the gradients of local functions gives

We bound the first term in terms of ∥∇f(x^t(r))∥2\|\nabla f(\widehat{{\bf x}}_{t}^{(r)})\|^{2} as follows:

where the 2nd inequality follows from the smoothness (LL-Lipschitz gradient) assumption. Using this and that ηt≤12L\eta_{t}\leq\frac{1}{2L} in (40) and rearranging terms give

Taking expectation w.r.t. to the entire process and using the inequality ∥u+v∥2≤2∥u∥2+2∥v∥2\|{\bf u}+{\bf v}\|^{2}\leq 2\|{\bf u}\|^{2}+2\|{\bf v}\|^{2} gives

Observe that (43) holds irrespective of the learning rate schedule. In particular, if we take a fixed learning rate ηt=η≤12L\eta_{t}=\eta\leq\frac{1}{2L} in (43), we get

By taking a telescopic sum from t=0t=0 to t=T−1t=T-1, we get

Take η=C^T\eta=\frac{\widehat{C}}{\sqrt{T}}, where C^\widehat{C} is a constant (that satisfies C^<T2L\widehat{C}<\frac{\sqrt{T}}{2L}). For example, we can take C^=12L\widehat{C}=\frac{1}{2L}. This gives

B.7 Proof of Theorem 2

Observe that we can use the proof of Theorem 1 exactly until (43), for ηt≤12L\eta_{t}\leq\frac{1}{2L} (which follows from our assumption that a≥2ξLa\geq 2\xi L), which gives

Taking a telescopic sum from t=0t=0 to t=T−1t=T-1 gives

Let δt:=ηt4R\delta_{t}:=\frac{\eta_{t}}{4R} and PT:=∑t=0T−1∑r=1RδtP_{T}:=\sum_{t=0}^{T-1}\sum_{r=1}^{R}\delta_{t}. We show at the end of this proof that PT≥ξ4ln⁡(T+a−1a)P_{T}\geq\frac{\xi}{4}\ln{\left(\frac{T+a-1}{a}\right)}, ∑t=0T−1ηt2≤ξ2a−1\sum_{t=0}^{T-1}\eta_{t}^{2}\leq\frac{\xi^{2}}{a-1}, and that ∑t=0T−1ηt3≤ξ32(a−1)2\sum_{t=0}^{T-1}\eta_{t}^{3}\leq\frac{\xi^{3}}{2(a-1)^{2}}. Using these in (49) yields

We therefore can show a weak convergence result, i.e.,

Bounding the terms PTP_{T}, ∑t=0T−1ηt2\sum_{t=0}^{T-1}\eta_{t}^{2} and ∑t=0T−1ηt3\sum_{t=0}^{T-1}\eta_{t}^{3}:

B.8 Proof of Theorem 3

Let x∗{\bf x}^{*} be the minimizer of f(x)f({\bf x}), therefore we have ∇f(x∗)=0\nabla f({\bf x}^{*})=0. We denote f(x∗)f({\bf x}^{*}) by f∗f^{*}. By taking the average of the virtual sequences x~t+1(r)=x~t(r)−ηt∇fit(r)(x^t(r))\widetilde{\mathbf{x}}_{t+1}^{(r)}=\widetilde{\mathbf{x}}_{t}^{(r)}-\eta_{t}\nabla f_{i_{t}^{(r)}}\left(\widehat{\mathbf{x}}_{t}^{(r)}\right) for each worker r∈[R]r\in[R] and defining pt:=1R∑r=1R∇fit(r)(x^t(r))\mathbf{p}_{t}:=\frac{1}{R}\sum_{r=1}^{R}\nabla f_{i_{t}^{(r)}}\left(\widehat{\mathbf{x}}_{t}^{(r)}\right), we get

Taking the expectation w.r.t. the sampling iti_{t} at time tt (conditioning on the past) and noting that last term in (53) becomes zero gives:

If ηt≤14L\eta_{t}\leq\frac{1}{4L}, then we have

Using the definition of p‾t\overline{\mathbf{p}}_{t} we have

By the definition of smoothness, we have ∥∇f(x~t)−∇f(x∗)∥2≤2L(f(x~t)−f(x∗))\|\nabla f\left(\widetilde{{\bf x}}_{t}\right)-\nabla f\left(\mathbf{x}^{*}\right)\|^{2}\leq 2L\left(f\left(\widetilde{{\bf x}}_{t}\right)-f({\bf x}^{*})\right), where ∇f(x∗)=0\nabla f(\mathbf{x}^{*})=0. Substituting this in (58) gives

Now we bound the last term of (57). By definition, we have

For the first term on the RHS of (60), we can use strong convexity

For the second term on the RHS of (60), we can use the following by smoothness.

Adding (59) and (63) and using a≥32L/μa\geq 32L/\mu which implies ηt≤\nicefrac14L\eta_{t}\leq\nicefrac{{1}}{{4L}} yields

Since ∥x+y∥2≤2∥x∥2+2∥y∥2\|\mathbf{x}+\mathbf{y}\|^{2}\leq 2\|\mathbf{x}\|^{2}+2\|\mathbf{y}\|^{2}, we have

Using (65) in (64) and then substituting (64) in (57) gives

Now use −∥x~t−x∗∥2≤∥x~t−x^t∥2−12∥x^t−x∗∥2-\|\widetilde{{\bf x}}_{t}-{\bf x}^{*}\|^{2}\leq\|\widetilde{{\bf x}}_{t}-\widehat{{\bf x}}_{t}\|^{2}-\frac{1}{2}\|\widehat{{\bf x}}_{t}-{\bf x}^{*}\|^{2} We get

Using (68) in (55) and then taking the expectation over the entire process gives

Now using ηt≤\nicefrac14L\eta_{t}\leq\nicefrac{{1}}{{4L}} we have

For ηt=8μ(a+t)\eta_{t}=\frac{8}{\mu\left(a+t\right)} and wt=(a+t)2w_{t}=\left(a+t\right)^{2}, ST=∑t=oT−1≥T33S_{T}=\sum_{t=o}^{T-1}\geq\frac{T^{3}}{3} we have

Where x‾T:=1ST∑t=0T−1[wt(1R∑r=1Rx^t(r))]=1ST∑t=0T−1wtx^t\overline{\mathbf{x}}_{T}:=\frac{1}{S_{T}}\sum_{t=0}^{T-1}\left[w_{t}\left(\frac{1}{R}\sum_{r=1}^{R}\widehat{\mathbf{x}}_{t}^{\left(r\right)}\right)\right]=\frac{1}{S_{T}}\sum_{t=0}^{T-1}w_{t}\widehat{\mathbf{x}}_{t}. This completes the proof of Theorem 3. ∎

Appendix C Omitted Details from Section 4

As before, in order to prove our results in the asynchronous setting, we define virtual sequences for every worker r∈[R]r\in[R] and for all t≥0t\geq 0 as follows:

x~t+1:=1R∑r=1Rx~t+1(r)=x~t−ηtR∑r=1R∇fit(r)(x^t(r))\widetilde{\mathbf{x}}_{t+1}:=\frac{1}{R}\sum_{r=1}^{R}\widetilde{\mathbf{x}}_{t+1}^{\left(r\right)}=\widetilde{\mathbf{x}}_{t}-\frac{\eta_{t}}{R}\sum_{r=1}^{R}\nabla f_{i_{t}^{(r)}}\left(\widehat{\mathbf{x}}_{t}^{\left(r\right)}\right)

pt:=1R∑r=1R∇fit(r)(x^t(r))\mathbf{p}_{t}:=\frac{1}{R}\sum_{r=1}^{R}\nabla f_{i_{t}^{(r)}}\left(\widehat{\mathbf{x}}_{t}^{\left(r\right)}\right)

x^t=1R∑r=1Rx^t(r)\widehat{\mathbf{x}}_{t}=\frac{1}{R}\sum_{r=1}^{R}\widehat{\mathbf{x}}_{t}^{\left(r\right)}

Let gap(IT(r))≤Hgap(\mathcal{I}_{T}^{(r)})\leq H holds for every r∈[R]r\in[R]. For x^t(r)\widehat{{\bf x}}_{t}^{(r)} generated according to Algorithm 2 with decaying learning rate ηt\eta_{t} and letting x^t=1R∑r=1Rx^t(r)\widehat{\bf x}_{t}=\frac{1}{R}\sum_{r=1}^{R}\widehat{\bf x}_{t}^{(r)}, we have the following bound on the deviation of the local sequences:

where C′′=8(4−2γ)(1+Cγ2)C^{\prime\prime}=8(4-2\gamma)(1+\frac{C}{\gamma^{2}}) and CC is a constant satisfying C≥4aγ(1−γ2)aγ−4HC\geq\frac{4a\gamma(1-\gamma^{2})}{a\gamma-4H}.

We bound both the terms separately. For the first term:

The last inequality (76) uses ηtr≤2ηtr+H≤2ηt\eta_{t_{r}}\leq 2\eta_{t_{r}+H}\leq 2\eta_{t} and t−tr≤Ht-t_{r}\leq H. To bound the second term of (75), note that we have

Putting this and the bound from (76) back in (75) gives

C.2 Proof of Lemma 10

Let gap(IT(r))≤Hgap(\mathcal{I}_{T}^{(r)})\leq H holds for every r∈[R]r\in[R]. By running Algorithm 2 with fixed learning rate η\eta, we have

where C′=(16γ2−12)(4−2γ)C^{\prime}=(\frac{16}{\gamma^{2}}-12)(4-2\gamma).

For a fixed learning rate η\eta, using (83) and following similar analysis as in (76) we can bound the first term in (75) as follows

Similarly as in (77)-(81) we can bound the second term in (75) as follows

Using (84) and (85) in (75) we can show that

C.3 Proof of Lemma 11

Let gap(IT(r))≤Hgap(\mathcal{I}_{T}^{(r)})\leq H holds for every r∈[R]r\in[R]. If we run Algorithm 2 with a decaying learning rate ηt\eta_{t}, then we have the following bound on the difference between the true and virtual sequences:

where C′=192(4−2γ)(1+Cγ2)C^{\prime}=192(4-2\gamma)\left(1+\frac{C}{\gamma^{2}}\right) and CC is a constant satisfying C≥4aγ(1−γ2)aγ−4HC\geq\frac{4a\gamma(1-\gamma^{2})}{a\gamma-4H}.

Applying Jensen’s inequality and taking expectation gives

We bound each of the three terms of (88) separately. We have upper-bounded the first term earlier in (82), which is

where B=(4−2γ)B=(4-2\gamma). To bound the second term of (88), note that

To bound the last term of (88), note that

Let tr(1)t_{r}^{(1)} and tr(2)t_{r}^{(2)} be two consecutive synchronization steps in IT(r)\mathcal{I}_{T}^{(r)}. Then, by the update rule of x^t(r)\widehat{\bf x}_{t}^{(r)}, we have x^tr(1)(r)−x^tr(2)−12(r)=∑j=tr(1)tr(2)−1∇fij(r)(x^j(r))\widehat{\bf x}_{t_{r}^{(1)}}^{(r)}-\widehat{\bf x}_{t_{r}^{(2)}-\frac{1}{2}}^{(r)}=\sum_{j=t_{r}^{(1)}}^{t_{r}^{(2)}-1}\nabla f_{i_{j}^{(r)}}\left(\widehat{\bf x}_{j}^{(r)}\right). Since xtr(1)(r)=x^tr(1)(r){\bf x}_{t_{r}^{(1)}}^{(r)}=\widehat{\bf x}_{t_{r}^{(1)}}^{(r)} and the workers do not modify their local xt(r){\bf x}_{t}^{(r)}’s in between the synchronization steps, we have xtr(2)−1(r)=xtr(1)(r)=x^tr(1)(r){\bf x}_{t_{r}^{(2)}-1}^{(r)}={\bf x}_{t_{r}^{(1)}}^{(r)}=\widehat{\bf x}_{t_{r}^{(1)}}^{(r)}. Therefore, we can write

Using (95) for every consecutive synchronization steps, we can equivalently write (94) as

Putting the bounds from (89), (92), and (97) in (88) and using B=(4−2γ)B=(4-2\gamma) give

C.4 Proof of Lemma 12

Let gap(IT(r))≤Hgap(\mathcal{I}_{T}^{(r)})\leq H holds for every r∈[R]r\in[R]. If we run Algorithm 2 with a fixed learning rate η\eta, we have

where C′=(4−2γ)(8γ2−6)C^{\prime}=(4-2\gamma)\left(\frac{8}{\gamma^{2}}-6\right).

For a constant learning rate the first term in (88) has been bounded earlier in (85). Following similar steps as in (91) we would have

Finally, using (85),(96), Lemma 5 and (98) in (88) we have

where B=(4−2γ)B=(4-2\gamma). This completes the proof of Lemma 12. ∎

Appendix D Omitted Details from Section 5

As mentioned in Footnote 4, here we compare the performance of Qsparse-local-SGD with scaled and unscaled composed operator QTopkQTop_{k} in the non-convex setting. We will see that even though the scaled QTopkQTop_{k} from Lemma 2 works better than unscaled QTopkQTop_{k} from Lemma 1 theoretically (see Remark 2), our experiments show the opposite phenomena, that the unscaled QTopkQTop_{k} works at least as good as the scaled QTopkQTop_{k}, and strictly better in some cases. We can attribute this to the fact that scaling the composed operator is a sufficient condition to obtain better convergence results, which does not necessarily mean that in practice also it does better. Therefore, we perform our experiments in the non-convex setting Section 5 with unscaled QTopkQTop_{k}.

References