The Convergence of Stochastic Gradient Descent in Asynchronous Shared Memory

Dan Alistarh, Christopher De Sa, Nikola Konstantinov

Introduction

The tremendous recent progress in machine learning can be explained in part through the availability of extremely vast amounts of data, but also through improved computational support for executing machine learning tasks at scale. Perhaps surprisingly, some of the main algorithms driving this progress are have been known in some form or another for a very long time. One such tool is the classic stochastic gradient descent (SGD) optimization algorithm , introduced by Robbins and Munro in the 1950s, variants of which are currently the tool of choice for optimization in a variety of settings, such as image classification and speech recognition via neural networks, but also in fundamental data processing tools such as regression.

where xtx_{t} is a dd-dimensional vector which we will henceforth call the model, encoding our current beliefs about the data, and αt\alpha_{t} is the step size, which controls how “aggressive” updates should be.

A classic instance of this method, which we consider in this paper, is the following. We are given a large set of mm data points X1,…,XmX_{1},\ldots,X_{m}, where to each point XiX_{i}, with i∈{1,2,…,m}i\in\{1,2,\ldots,m\}, we associate a loss function Li(x)L_{i}(x), measuring the loss of any model xx at the data point XiX_{i} (mapping the difference between the prediction of our model on XiX_{i} with the true label to a real value), and we wish to identify a model x⋆x^{\star} which minimizes the average loss over the dataset, that is, minimize the function

This setting can be mapped to fundamental optimization tasks, such as linear regression over a given set of points, or training neural networks through backpropagation . In both settings, we are given a large dataset, from which we pick a sample point uniformly at random in every iteration tt. In every such iteration, a process reads the current version of the model xtx_{t}, computes the stochastic gradient g~(xt)\widetilde{g}(x_{t}) of this model with respect to the current sample by examining the derivative of the loss function at the point, given its label, and updates the model’s components correspondingly. Thus, the stochastic gradients g~(xt)\widetilde{g}(x_{t}) correspond to the gradient of the model xtx_{t} taken at a uniformly random data point XiX_{i}. Indeed, since the data point at which we compute the gradient is chosen uniformly at random, we have

Parallel SGD.

The ability to parallelize this process at extremely large scale, up to thousands of nodes, e.g. , has enabled researchers to approach new problems, and to reach super-human accuracy for several classic problems, e.g. . The standard way to distribute SGD to multiple compute nodes is data parallelism: given nn nodes, which we will abstract as parallel processes, we split the dataset into nn partitions. Nodes process samples from their partitions in parallel, and synchronize by using a shared model xx. In this paper, we will focus on the case where the nodes are threads, communicating through asynchronous shared memory.

In theory, data parallelism would allow the system to perform nn times more iterations (and model updates) per unit of time. The catch is that threads will need to synchronize on the shared model, reducing scalability. In fact, early studies proposed using coarse-grained locking to keep the process consistent to a sequential execution. As expected, coarse-grained locking leads to significant loss of performance.

In this context, a breakthrough by Niu et al. proved the unintuitive result that SGD should be able to converge even without coarse-grained synchronization to maintain model consistency, under some technical conditions, including high sparsity of the gradient updates, and under the assumption that individual updates are applied via fetch-and-add synchronization operations. (This latter assumption appears necessary in general, since otherwise a delayed thread could completely obliterate all progress made up to some point, by overwriting the entire model, which resides in shared memory.) This work has sparked significant research on the convergence properties of asynchronous SGD, e.g. , improving theoretical bounds and considering more general settings.

In a nutshell, these results show, under various delay models and analytic assumptions, that SGD can still converge if iterations are asynchronous, that is, they cannot be “serialized” in any meaningful way. At the same time, all known convergence upper bounds, e.g. , have a linear dependence in τmax⁡\tau_{\max}, the maximum delay between the time where a gradient update is generated by any thread, and the point where the gradient is applied to the model.In shared-memory parlance, τmax⁡\tau_{\max} is upper bounded by the interval contention, where we define SGD iterations as individual operations. Intuitively, τmax⁡\tau_{\max} is upper bounded by the number of iterations which can take steps between the start and end points of any fixed iteration. This means that asynchronous SGD will take τmax⁡\tau_{\max} times more iterations to converge, compared to the synchronous variant. Since τmax⁡\tau_{\max} is technically only upper bounded by the length of the execution, this upper bound appears extremely harsh. It is therefore natural to ask if this dependency is inherent, or whether it can be improved, yielding superior convergence rates for this classical algorithm.

Contribution.

In this paper, we address this problem. Our approach is to express data-parallel SGD in the classic asynchronous shared memory model, against a strong (adaptive) adversarial scheduler, which designs schedules to delay the algorithm from converging, with full knowledge of the algorithm and random coin flips. Under this formulation, our main results are as follows:

We show that, under standard analytic assumptions, for convex objective functions ff, there exists a simple variant of the classic SGD algorithm which still converges under this strong adversarial model.

We prove that, under reasonable parameter settings, the convergence slowdown of SGD iterations due to asynchrony can be upper bounded by τmax⁡n\sqrt{\tau_{\max}n}, where τmax⁡\tau_{\max} is the maximum interval contention over all operations and nn is the number of threads. This result shows for the first time that the runtime dependence need not be linear in τmax⁡\tau_{\max}, even against a strong adversary, and is our main technical contribution.

We prove that, in general, the adversary can cause a convergence slowdown that is linear in τmax⁡\tau_{\max}, if the algorithm does not decrease the step size αt\alpha_{t} to offset the influence of stale gradient updates. This shows for the first time that the adversary can consistently and significantly slow down convergence by inducing delays, and that decreasing the step size (which is done by our algorithm) is in fact necessary for good convergence under asynchrony.

In sum, our results give new and improved upper and lower bounds on the “price of asynchrony” when executing the fundamental SGD algorithm in a concurrent setting. They show that this classic optimization tool can converge faster and with a wider range of parameters than previously known, under asynchronous iterations. At the same time, we exhibit a simple yet fundamental trade-off between the maximum delay in the system and the convergence rate of SGD, which governs the set of parameters under which SGD can still work efficiently.

Techniques and Related Work.

From the technical perspective, our results build upon martingale-based approaches for bounding the convergence of SGD, e.g. . These techniques complement the classic “regret” bounds for characterizing the convergence of SGD, e.g. . We exploit martingale-based techniques in the asynchronous setting. To our knowledge, the only other work to employ such techniques for convex SGD is , whose results we significantly extend. Specifically, with respect to this reference, the main departures are that 1) we consider a more challenging adaptive adversarial model, as opposed to a stochastic scheduling model; 2) eliminate the requirement that gradients contain a single non-zero entry, thereby significantly expanding the applicability of the framework; 3) reduce the linear convergence dependency in τmax⁡\tau_{\max} to one of the form τmax⁡n\sqrt{\tau_{\max}n}; 4) prove lower bounds on the slowdown due to asynchrony.

There is an extremely vast literature studying the convergence properties of asynchronous optimization methods , as well as efficient parallel implementations, e.g. , starting with seminal work by Bertsekas and Tsitsiklis . A complete survey is beyond the scope of this paper, and we therefore focus on work that is directly related to ours. Reference showed for the first time that, under strong analytical assumptions on sparsity and on the target loss function, asynchronous SGD can still converge, and that, moreover, the convergence rate can be similar to that of the baseline under further assumptions on the parameters. Agarwal and Duchi showed that, under strong ordering assumptions, delayed gradient computation does not affect the convergence of SGD. Lian et al. provided ergodic convergence rates for asynchronous SGD for non-convex objectives. Duchi et al. considered a model similar to ours, and showed that the impact of any asynchrony on the rate at which the algorithm converges is negligible, under strong technical assumptions on the convex function ff to be optimized, on the structure of its optimum, and on the sampling noise. Concurrent work provides significantly more general analyses of iterative processes under asynchrony, covering several important optimization algorithms. With the exception of , which makes strong technical assumptions, all previous results for asynchronous SGD had a linear dependence in the maximum delay τmax⁡\tau_{\max}. We improve this dependency in this work.

There exists significant work on mitigating the effects of asynchrony in applied settings, e.g. . A subset of these works are designed for a distributed shared memory setting, where it may be possible to examine the “staleness” of an update immediately before applying it, and adjust hyperparameters accordingly, and validated empirically. By contrast, we consider an adversarial setting, where the scheduler actively attempts to thwart the algorithm’s progress. Our lower bound applies to these works as well.

There has recently been significant work connecting machine learning and optimization with distributed computing. References consider distributed SGD in a Byzantine adversarial setting, but in a message-passing system. In a series of papers , Su and Vaidya have considered the problem of adding fault-tolerance to the problem of multi-agent optimization, as well as non-Bayesian optimization under asynchrony and crash failures .

Model

We consider a standard asynchronous shared-memory model , in which nn threads (or processes) P1,…,PnP_{1},\ldots,P_{n}, communicate through atomic memory locations called registers, on which they perform atomic operations such as read\mathsf{read}, write\mathsf{write}, compare&swap\mathsf{compare\&swap} and fetch&add\mathsf{fetch\&add}. In particular, the algorithm we consider employs atomic read\mathsf{read} and fetch&add\mathsf{fetch\&add} operations, which are now standard in mass-produced multiprocessors. The fetch&add\mathsf{fetch\&add} operation takes one argument, and returns the value of the register before the increment was performed, incrementing its value by the corresponding operand.

As is usual, we will assume a sequentially consistent memory model, in which once a thread returns from its invocation of a primitive (for example, a fetch&add\mathsf{fetch\&add}), the value written by the thread is immediately applied to shared memory, and henceforth visible by other processors.

The Adversarial Scheduler.

Threads follow an algorithm composed of shared-memory steps and local computation, including random coin flips. The order of process steps is controlled by an adversarial entity we call the scheduler. Time is measured in terms of the number of shared-memory steps scheduled by the adversary. The adversary may choose to crash a set of at most n−1n-1 processes by not scheduling them for the rest of the execution. A process that is not crashed at a certain step is correct, and if it never crashes then it takes an infinite number of steps in the execution. In the following, we assume a standard strong adversarial scheduler, which can see the results of the threads’ local coins when deciding the scheduling.

Contention Bound.

In the following, fix a (concurrent) SGD iteration θ\theta, and let ρ(θ)\rho(\theta) be the interval contention of the iteration θ\theta, defined as the number of SGD iterations which can execute concurrently with θ\theta. Let τmax⁡\tau_{\max} be the upper bound over all these interval contention values, i.e. τmax⁡=max⁡θρ(θ)\tau_{\max}=\max_{\theta}\rho(\theta). Let τavg\tau_{avg} be an upper bound on the average interval contention over all iterations during the (finite) execution of a program, i.e. τavg=1/T∑1≤θ≤Tρ(θ)\tau_{avg}=1/T\sum_{1\leq\theta\leq T}\rho(\theta), where TT is the total number of iterations of the algorithm. It is known that τavg≤2n\tau_{avg}\leq 2n, where nn is the number of threads .

Background on Stochastic Gradient Descent

Unless otherwise noted, we consider SGD for convex optimization and with a constant learning rate, that is αt=α\alpha_{t}=\alpha for all tt. We also make the following standard assumptions:

Classic approaches for analyzing the convergence of SGD attempt to bound the distance between the expected value of ff at the average of the currently generated iterates and the optimal value of the function (e.g. Theorem 6.3 in ), showing that this distance decreases linearly with the number of iterations.

Here we consider a different approach that aims at estimating the probability that the algorithm has failed to converge to a success region around the optimal parameter value after TT steps. To this end, we employ a martingale-based analysis of the algorithm, an approach that has recently become a popular tool for analyzing asynchronous optimization algorithms . Let x∗x^{*} be the minimizer of the function ff. Given an ϵ>0\epsilon>0, we denote by

the success region around this minimizer, to which we want to converge. Our analysis aims to bound the probability of the event FTF_{T} that xi∉Sx_{i}\not\in S for all i≤Ti\leq T, i.e. the event that the algorithm has failed to hit the success region by time TT.

An existing result of this type about the convergence of parallel SGD was derived in . Under the non-standard additional assumption that each gradient update on the parameter value only effects a single entry of xtx_{t} (i.e., that the stochastic gradients contain a single non-zero entry),Our analysis eliminates this assumption, a result which may be of independent interest. one can obtain the following:

Consider the SGD algorithm under the assumptions above, run for TT steps, with a success region S={x ∣ ∥x−x∗∥2≤ϵ}S=\{x\text{ }|\text{ }\|x-x^{*}\|^{2}\leq\epsilon\} and with learning rate α=cϵϑM2\alpha=\frac{c\epsilon\vartheta}{M^{2}} for some constant ϑ∈(0,1)\vartheta\in\left(0,1\right). Then the probability of the event FTF_{T} that xi∉Sx_{i}\not\in S for all i≤Ti\leq T is:

Note that this bound also decreases linearly with the number of iterations.

Lock-Free SGD in Shared-Memory

A standard way of parallelizing the SGD algorithm is to have multiple parallel threads execute the procedure in Equation 1. We assume a lock-free setting, in which threads share the set of parameters (model) X[d]X[d], which they can read and update entry-wise (through read\mathsf{read} and fetch&add\mathsf{fetch\&add} operations) concurrently, without additional synchronization. Each thread executes the steps in Algorithm 1:

We emphasize that this modeling of the algorithm is standard: virtually all papers which analyze asynchronous SGD consider this formulation, e.g. . Updates are assumed to occur via fetch&add, to avoid complete resets of the state by a delayed thread.

A Slowdown Lower Bound via Adversarial Delays

We now provide a simple argument which yields a lower bound on the achievable speedup if the adversary can delay gradients by a large τmax⁡\tau_{\max}. This argument might also serve as a brisk hands-on introduction to SGD. We consider a standard setting, where two threads each have access to local gradient samples, and share the model. Assume we are trying to minimize the (convex) objective function

Adversarial Strategy.

Suppose that the adversary is executing the following strategy. First, both threads generate a gradient with respect to x0x_{0}. The first thread then executes for τ\tau consecutive iterations, and then the second writes the state with a gradient from the initial value. Let us now analyze the rate at which this algorithm can converge.

Analysis.

After the first thread runs, the output will be

The second term of this last expression will be a zero-mean Gaussian with a variance of

For a fixed α\alpha, if we choose τ\tau large enough that 2(1−α)τ≤α2(1-\alpha)^{\tau}\leq\alpha, and suppose for simplicity that σ=0\sigma=0, then we can get

in the case with no adversary. To compare these rates, we take the logarithm, and obtain a slowdown factor of

which implies an Ω(τ)\Omega(\tau) factor slowdown is possible from a delay of τ\tau. This shows that with a maximum delay τmax⁡=τ\tau_{\max}=\tau, the adversary can achieve an asymptotic slowdown that is linear in τ\tau. We conclude as follows.

Given an instance of the lock-free SGD Algorithm in 1 with fixed learning rate α\alpha, there exists an adversarial strategy with maximum delay τmax⁡=O(log⁡(α)/log⁡(1−α))\tau_{\max}=O(\log(\alpha)/\log(1-\alpha)) such that the algorithm converges τmax⁡\tau_{\max} times slower than the sequential variant.

Convergence Upper Bounds in Asynchronous Shared Memory

We now introduce some notation, and a few basic claims about the above concurrent process. First, we define an order on the above iterations, performed possibly by distinct threads, by the time at which the iteration performs its first fetch&add operation, on the first model component XX. (Here we are using the sequential consistency property of the memory model.) This ordering induces a useful total order between iterations: for any integer t≥1t\geq 1, iteration tt is the ttth iteration to complete its fetch&add on XX. We now note that all of these iterations up to iteration tt must have completed their computation, but may not have completed writing their updates to XX by the time when iteration t+1t+1 starts. At most nn of these iterations may be incomplete at any given time. We formalize this as follows.

Let iteration tt be the ttth iteration to update XX. This is a total order on the iterations. We say that an iteration is incomplete at a given point in the execution if it has performed its first update (on XX), but has not completed its last update (on X[d]X[d]). For any t≥1t\geq 1, at most nn iterations with indices ≤t\leq t can be incomplete.

By assumption the maximal interval contention is bounded by τmax⁡\tau_{\max}. However, since at most nn threads can run at the same time, it is intuitive that the average contention should be O(n)\mathcal{O}\left(n\right). We formalize this via the following:

Fix a parameter KK, and an arbitrary time interval II during which exactly KnKn consecutive SGD iterations start. We call an SGD iteration θ\theta bad if more than KnKn iterations start between its start time and end time. Otherwise, an SGD iteration is good. Then, the number of bad iterations which complete during II is less than nn.

Assume for contradiction that the number of bad iterations is nn or larger. Then, there must exist a thread pp which completes two of these bad iterations during this interval. Denote the second iteration by θ\theta. Since the two iterations by pp cannot be concurrent with each other, the maximum number of iterations which can be concurrent with θ\theta (except itself) is at most Kn−2Kn-2. Hence, θ\theta cannot be bad, a contradiction. ∎

SGD Convergence.

We now analyze the convergence of SGD under this lock-free model. Following , we employ a martingale approach for proving convergence rates of SGD.

where expectation is taken with respect to the randomness at time tt and conditional on the past. Secondly, for any time T≤BT\leq B and any sequence xT,...,x0x_{T},...,x_{0}, if the algorithm has not succeeded by time TT (i.e. xt∉Sx_{t}\not\in S for all t≤Tt\leq T), then:

The main result in shows how constructing such a supermartingale for an optimization algorithm can be used to obtain a bound on the probability that the algorithm has not visited the success region after a certain number of iterations. Under the considered stochastic scheduling model, the authors employ a parameter τ\tau, that denotes the worst-case expected delay caused by the parallel updates. Under the additional assumption that the stochastic gradients contain a single non-zero entry, the following result is derived:

Consider the SGD algorithm for optimizing a convex function ff that satisfies the assumptions above and under the asynchronous model of , with a success region S={x ∣ ∥x−x∗∥2≤ϵ}S=\{x\text{ }|\text{ }\|x-x^{*}\|^{2}\leq\epsilon\} and with learning rate α=cϵϑM2+2LMτϵ\alpha=\frac{c\epsilon\vartheta}{M^{2}+2LM\tau\sqrt{\epsilon}} for some constant ϑ∈(0,1)\vartheta\in\left(0,1\right). Then the probability of the event FTF_{T} that xi∉Sx_{i}\not\in S for all i≤Ti\leq T is:

Again, the bound decreases linearly with the number of iterations. However, there is an additive term that increases with τ\tau. In particular, the bound on the failure probability is worse than in the sequential SGD case described previously.

2 Convergence Analysis

We now apply a martingale analysis similar to the one in to obtain results about the rate of convergence of SGD under the Asynchronous Shared Memory model. We denote the maximum delay at time tt by τt≥0\tau_{t}\geq 0. We also assume that all τt\tau_{t} are bounded by some maximum τmax\tau_{\textrm{max}}, i.e. τt≤τmax⁡\tau_{t}\leq\tau_{\max} for all tt.

Note that at any time tt, the gradient is computed based on a view vtv_{t} that might be missing updates from only the last τt\tau_{t} iterations. Therefore,

For the subsequent analysis we will also need the following:

By Lemma 6.2, we know that for any constant KK and for any KnKn consecutive steps t+1,...,t+Knt+1,...,t+Kn, τt+i>Kn\tau_{t+i}>Kn for at most nn indexes. Hence,

This holds for any positive KK. The bound is minimized at K=τmaxnK=\sqrt{\frac{\tau_{\textrm{max}}}{n}}, which yields the result. ∎

Next, we obtain a bound on the probability that the algorithm has not visited a given success region S={x ∣ ∥x−x∗∥2≤ϵ}S=\{x\text{ }|\text{ }\|x-x^{*}\|^{2}\leq\epsilon\}. To this end, we show a result similar to the one in Theorem 1 in . We will assume the existence of a rate supermartingale with respect to the underlying sequential SGD process that is Lipschitz in its first coordinate and show that this can be used to obtain a bound on the failure probability. The exact assumptions on WW are as follows:

WW is a supermartingale with horizon BB with respect to the sequential SGD process xt+1=xt−αg~(xt)x_{t+1}=x_{t}-\alpha\widetilde{g}\left(x_{t}\right). Note that WW need not be a supermartingale with respect to the lock-free SGD algorithm.

For any T>0T>0, if xi∉Sx_{i}\not\in S for all i≤Ti\leq T, then: WT(xT,…,x0)≥T.W_{T}\left(x_{T},\ldots,x_{0}\right)\geq T. Otherwise, we say that the algorithm has succeeded at time TT.

WW is Lipschitz continuous in the current iterate with parameter HH, i.e. for all t,u,vt,u,v and any xt−1,...,x0x_{t-1},...,x_{0}:

Under these assumptions, we can prove our main technical claim, whose proof is deferred to the Appendix.

Assume that WW is a rate supermartingale with horizon BB for the sequential SGD algorithm and that WW is HH-Lipschitz in the first coordinate. Assume further that α2HLMCd<1\alpha^{2}HLMC\sqrt{d}<1, where C=2τmaxnC=2\sqrt{\tau_{\textrm{max}}n}. Then for any T≤BT\leq B, the probability that the lock-free SGD algorithm has not succeeded at time TT (that is, the probability of the event FTF_{T} that xi∉Sx_{i}\not\in S for all i≤Ti\leq T) is:

We now apply the result with a particular choice for the martingale WtW_{t}. We use the process proposed in in the case of convex optimization, which they show is a rate supermartingale for the sequential SGD process. More precisely, we have the following:

Define the piecewise logarithm function to be

If the algorithm has not succeeded by timestep tt (i.e. xi∉Sx_{i}\not\in S for all i≤ti\leq t) and by Wt=Wu−1W_{t}=W_{u-1} whenever xi∈Sx_{i}\in S for some i≤ti\leq t and uu is the minimal index with this property. Then WtW_{t} is a rate supermartingale for sequential SGD with horizon B=∞B=\infty. It is also HH-Lipschitz in the first coordinate, with H=2ϵ(2αcϵ−α2M2)−1H=2\sqrt{\epsilon}\left(2\alpha c\epsilon-\alpha^{2}M^{2}\right)^{-1}, that is for any t,u,vt,u,v and any sequence xt−1,…,x0x_{t-1},\ldots,x_{0}:

Using this martingale in Theorem 6.5, we obtain the following result, whose proof is deferred to the Appendix:

Assume that we run the lock-free SGD algorithm under the Asynchronous Shared Memory model for minimizing a convex function ff satisfying the listed assumptions. Set the learning rate to:

for some constant ϑ∈(0,1]\vartheta\in(0,1]. Then for any T>0T>0 the probability that xi∉Sx_{i}\not\in S for all i≤Ti\leq T is:

Choosing ϑ=1\vartheta=1 gives the smallest value for the upper bound under this setting. We can impose an arbitrary small learning rate by selecting a small ϑ\vartheta, while still ensuring convergence.

Iterated Algorithm with Guaranteed Convergence

The procedure in Algorithm 1 has the property that, for an appropriately chosen stopping time TT, it will eventually reach the “success” region, where the distance to the optimal parameter value falls below ϵ\epsilon. However, due to asynchrony and adversarial updates, threads might perform updates which cause them to leave the success region: a delayed thread might apply stale gradients to the model, overwriting the progress.

We now present an extension of the algorithm which deals with this problem, and allows us to converge to a success region of radius ϵ\epsilon around the optimum x∗x^{*}, for any ϵ>0\epsilon>0. The algorithm will run a series of epochs, each of which is a series of TT SGD iterations, executed using the procedure in Algorithm 1. The only difference between epochs is that they are executed with an exponentially decreasing learning rate α\alpha. (We note that this epoch pattern is already used in many settings, such as neural network training.) The epochs share the model XX, with the critical note that we require that a gradient update can only be applied to XX in the same epoch when it was generated. This condition can be enforced either by maintaining an epoch counter, on which threads condition their update via double-compare-single-swap (DCAS), or by having a distinct model allocated for each epoch.

The only distinct epoch is the last, in which threads each aggregate the gradients they produced locally. At the end of the epoch, the threads will collect all local gradients locally. As we show in the Appendix, the model will be guaranteed to be close to optimal in expectation.

We can characterize the probability of success by reducing the target ϵ\epsilon, and applying Markov’s inequality.

Discussion

An immediate consequence of our lower bound in Theorem 5.1 is that the algorithm must either have a low initial learning rate, or lower the learning rate across multiple iterations (as in Algorithm 2) in order to be able to withstand adversarial delays. Otherwise, the adversary can always apply stale gradients generated at a far enough time in the past to nullify progress. An alternative approach, which we did not consider here, would be to introduce a “momentum” term by which the current model value is multiplied .

Lower Bounds versus Upper Bounds.

The attentive reader might find it curious that our lower bound suggests a linear slowdown in τmax⁡\tau_{\max}, whereas our upper bound suggests that the slowdown is linear in O(τmax⁡n)O(\sqrt{\tau_{\max}{}n}). However, a close examination of the preconditions for these two results will reveal that they are in fact complementary: in the lower bound, given fixed learning rate α\alpha, the adversary needs to set a large delay τmax⁡≥(log⁡(α/2))/log⁡(1−α)\tau_{\max}{}\geq(\log(\alpha/2))/\log(1-\alpha) in order to slow down convergence. At the same time, the upper bound in Theorem 6.5 requires that 2α2HLMdτmax⁡n<12\alpha^{2}HLM\sqrt{d\tau_{\max}n}<1, which is incompatible with the above condition. Specifically, our improved convergence bound shows that asynchronous SGD converges faster and for a wider range of parameters than previously known.

Why is Asynchronous SGD Fast in Practice?

In a nutshell, Theorem 6.5 shows that the gap in convergence between asynchronous SGD and the sequential variant becomes negligible if α2HLMCdτmax⁡n≪1\alpha^{2}HLMC\sqrt{d\tau_{\max}n}\ll 1. Intuitively, this condition holds in practice since gradients are often sparse, meaning that dd is low, the delay factors τmax⁡\tau_{\max} and τavg\tau_{avg} are not set adversarially, and the learning rate α\alpha can be set by the user to be small enough to offset any increase in the other terms. In particular, τmax⁡\tau_{\max} is limited by the staleness of updates in the write buffer at each core, which is well bounded in practice .

At the same time, it is important to note that, while we “sequentialize” iterations in our analysis, up to nn iterations may happen in parallel at any time, reducing the wall-clock convergence time by up to a factor of nn. Thus, the practical trade-off is between any slow-down caused by asynchrony, and the parallelism due to multiple computation threads.

Acknowledgements

We would like to thank Martin Jaggi for useful discussions. This project has received funding from the European Union’s Horizon 2020 research and innovation programme under the Marie Skłodowska-Curie Grant Agreement No. 665385.

References

Appendix A Deferred Proofs

Assume that WW is a rate supermartingale with horizon BB for the sequential SGD algorithm and that WW is HH-Lipschitz in the first coordinate. Assume further that α2HLMCd<1\alpha^{2}HLMC\sqrt{d}<1, where C=2τmaxnC=2\sqrt{\tau_{\textrm{max}}n}. Then for any T≤BT\leq B, the probability that the lock-free SGD algorithm has not succeeded at time TT (that is, xi∉Sx_{i}\not\in S for all i≤Ti\leq T) is:

Consider a process defined by V0=W0V_{0}=W_{0} and by:

whenever xi∉Sx_{i}\not\in S for all 0≤i≤t0\leq i\leq t. Finally, if xi∈Sx_{i}\in S for some 0≤i≤t0\leq i\leq t, then define

where uu is the minimal index, such that xu∈Sx_{u}\in S. Assume that the algorithm has not succeeded at time tt. Using the Lipschitzness of WW:

This inequality also trivially holds in the case when the algorithm has succeeded at time tt. Hence, the process VtV_{t} is a supermartingale for the lock-free SGD process.

Note also that if the algorithm has not succeeded at time TT, then WT≥TW_{T}\geq T and hence:

It follows that Vt≥0V_{t}\geq 0 for all tt. We also have that V0=W0V_{0}=W_{0}. Now for any T>0T>0:

A.2 Proof of Corollary 6.7

Assume that we run the lock-free SGD algorithm under the Asynchronous Shared Memory model for minimizing a convex function ff satisfying the listed assumptions. Set the learning rate to:

for some constant ϑ∈(0,1]\vartheta\in(0,1]. Then for any T>0T>0 the probability that xi∉Sx_{i}\not\in S for all i≤Ti\leq T is:

Substituting and using the result from that

Substituting the suggested value for the learning rate:

A.3 Analysis of Algorithm 2

The main idea behind the analysis is as follows. By Theorem 6.5, we know that, with high probability, there exists a time tt in each epoch where the aggregated gradients xtx_{t} enter the success region, i.e. ∥xt−x∗∥2≤ϵ\|x_{t}-x^{*}\|^{2}\leq\epsilon, for a fixed parameter ϵ\epsilon. The algorithm will guarantee that the model will not leave the success region before the end of the last epoch. This is ensured via our choice of the learning rate.

Let us now focus on the last epoch. Fix an ϵ>0\epsilon>0; we wish to prove that, ∥xT−x⋆∥2≤ϵ\|x_{T}-x^{\star}\|^{2}\leq\epsilon at the end of this epoch, in expectation. We fix the success condition of EpochSGD\mathsf{EpochSGD} such that that there exists an iteration tt in the epoch such that ∥xt−x⋆∥≤ϵ/2\|x_{t}-x^{\star}\|\leq\sqrt{\epsilon}/2. We note that the adversary may attempt to schedule “stale” updates, generated earlier in the execution, to cause the algorithm to leave the success region. However, we notice that there can be at most n−1n-1 gradients generated before time tt, which have not been applied yet. Denote these gradients by (G(vθi))i=1,…,n−1(G(v_{\theta_{i}}))_{i=1,\ldots,n-1}.

Finally, we claim that, in expectation, the distance between the final model xTx_{T} and the optimum is upper bounded by