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 is a -dimensional vector which we will henceforth call the model, encoding our current beliefs about the data, and 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 data points , where to each point , with , we associate a loss function , measuring the loss of any model at the data point (mapping the difference between the prediction of our model on with the true label to a real value), and we wish to identify a model 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 . In every such iteration, a process reads the current version of the model , computes the stochastic gradient 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 correspond to the gradient of the model taken at a uniformly random data point . 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 nodes, which we will abstract as parallel processes, we split the dataset into partitions. Nodes process samples from their partitions in parallel, and synchronize by using a shared model . 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 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 , 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, is upper bounded by the interval contention, where we define SGD iterations as individual operations. Intuitively, 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 times more iterations to converge, compared to the synchronous variant. Since 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 , 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 , where is the maximum interval contention over all operations and is the number of threads. This result shows for the first time that the runtime dependence need not be linear in , 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 , if the algorithm does not decrease the step size 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 to one of the form ; 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 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 . 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 threads (or processes) , communicate through atomic memory locations called registers, on which they perform atomic operations such as , , and . In particular, the algorithm we consider employs atomic and operations, which are now standard in mass-produced multiprocessors. The 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 ), 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 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 , and let be the interval contention of the iteration , defined as the number of SGD iterations which can execute concurrently with . Let be the upper bound over all these interval contention values, i.e. . Let be an upper bound on the average interval contention over all iterations during the (finite) execution of a program, i.e. , where is the total number of iterations of the algorithm. It is known that , where 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 for all . 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 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 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 be the minimizer of the function . Given an , 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 that for all , i.e. the event that the algorithm has failed to hit the success region by time .
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 (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 steps, with a success region and with learning rate for some constant . Then the probability of the event that for all 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) , which they can read and update entry-wise (through and 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 . 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 . The first thread then executes for 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 , if we choose large enough that , and suppose for simplicity that , 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 factor slowdown is possible from a delay of . This shows that with a maximum delay , the adversary can achieve an asymptotic slowdown that is linear in . We conclude as follows.
Given an instance of the lock-free SGD Algorithm in 1 with fixed learning rate , there exists an adversarial strategy with maximum delay such that the algorithm converges 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 . (Here we are using the sequential consistency property of the memory model.) This ordering induces a useful total order between iterations: for any integer , iteration is the th iteration to complete its fetch&add on . We now note that all of these iterations up to iteration must have completed their computation, but may not have completed writing their updates to by the time when iteration starts. At most of these iterations may be incomplete at any given time. We formalize this as follows.
Let iteration be the th iteration to update . 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 ), but has not completed its last update (on ). For any , at most iterations with indices can be incomplete.
By assumption the maximal interval contention is bounded by . However, since at most threads can run at the same time, it is intuitive that the average contention should be . We formalize this via the following:
Fix a parameter , and an arbitrary time interval during which exactly consecutive SGD iterations start. We call an SGD iteration bad if more than iterations start between its start time and end time. Otherwise, an SGD iteration is good. Then, the number of bad iterations which complete during is less than .
Assume for contradiction that the number of bad iterations is or larger. Then, there must exist a thread which completes two of these bad iterations during this interval. Denote the second iteration by . Since the two iterations by cannot be concurrent with each other, the maximum number of iterations which can be concurrent with (except itself) is at most . Hence, 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 and conditional on the past. Secondly, for any time and any sequence , if the algorithm has not succeeded by time (i.e. for all ), 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 , 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 that satisfies the assumptions above and under the asynchronous model of , with a success region and with learning rate for some constant . Then the probability of the event that for all is:
Again, the bound decreases linearly with the number of iterations. However, there is an additive term that increases with . 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 by . We also assume that all are bounded by some maximum , i.e. for all .
Note that at any time , the gradient is computed based on a view that might be missing updates from only the last iterations. Therefore,
For the subsequent analysis we will also need the following:
By Lemma 6.2, we know that for any constant and for any consecutive steps , for at most indexes. Hence,
This holds for any positive . The bound is minimized at , which yields the result. ∎
Next, we obtain a bound on the probability that the algorithm has not visited a given success region . 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 are as follows:
is a supermartingale with horizon with respect to the sequential SGD process . Note that need not be a supermartingale with respect to the lock-free SGD algorithm.
For any , if for all , then: Otherwise, we say that the algorithm has succeeded at time .
is Lipschitz continuous in the current iterate with parameter , i.e. for all and any :
Under these assumptions, we can prove our main technical claim, whose proof is deferred to the Appendix.
Assume that is a rate supermartingale with horizon for the sequential SGD algorithm and that is -Lipschitz in the first coordinate. Assume further that , where . Then for any , the probability that the lock-free SGD algorithm has not succeeded at time (that is, the probability of the event that for all ) is:
We now apply the result with a particular choice for the martingale . 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 (i.e. for all ) and by whenever for some and is the minimal index with this property. Then is a rate supermartingale for sequential SGD with horizon . It is also -Lipschitz in the first coordinate, with , that is for any and any sequence :
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 satisfying the listed assumptions. Set the learning rate to:
for some constant . Then for any the probability that for all is:
Choosing gives the smallest value for the upper bound under this setting. We can impose an arbitrary small learning rate by selecting a small , while still ensuring convergence.
Iterated Algorithm with Guaranteed Convergence
The procedure in Algorithm 1 has the property that, for an appropriately chosen stopping time , it will eventually reach the “success” region, where the distance to the optimal parameter value falls below . 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 around the optimum , for any . The algorithm will run a series of epochs, each of which is a series of 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 . (We note that this epoch pattern is already used in many settings, such as neural network training.) The epochs share the model , with the critical note that we require that a gradient update can only be applied to 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 , 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 , whereas our upper bound suggests that the slowdown is linear in . 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 , the adversary needs to set a large delay in order to slow down convergence. At the same time, the upper bound in Theorem 6.5 requires that , 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 . Intuitively, this condition holds in practice since gradients are often sparse, meaning that is low, the delay factors and are not set adversarially, and the learning rate can be set by the user to be small enough to offset any increase in the other terms. In particular, 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 iterations may happen in parallel at any time, reducing the wall-clock convergence time by up to a factor of . 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 is a rate supermartingale with horizon for the sequential SGD algorithm and that is -Lipschitz in the first coordinate. Assume further that , where . Then for any , the probability that the lock-free SGD algorithm has not succeeded at time (that is, for all ) is:
Consider a process defined by and by:
whenever for all . Finally, if for some , then define
where is the minimal index, such that . Assume that the algorithm has not succeeded at time . Using the Lipschitzness of :
This inequality also trivially holds in the case when the algorithm has succeeded at time . Hence, the process is a supermartingale for the lock-free SGD process.
Note also that if the algorithm has not succeeded at time , then and hence:
It follows that for all . We also have that . Now for any :
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 satisfying the listed assumptions. Set the learning rate to:
for some constant . Then for any the probability that for all 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 in each epoch where the aggregated gradients enter the success region, i.e. , for a fixed parameter . 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 ; we wish to prove that, at the end of this epoch, in expectation. We fix the success condition of such that that there exists an iteration in the epoch such that . 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 gradients generated before time , which have not been applied yet. Denote these gradients by .
Finally, we claim that, in expectation, the distance between the final model and the optimum is upper bounded by