Communication trade-offs for synchronized distributed SGD with large step size
Kumar Kshitij Patel, Aymeric Dieuleveut
Introduction
where is a deterministic sequence of positive scalars, referred to as the learning rate and is an oracle on the gradient of the function at . We focus on objective functions that are both smooth and strongly convex (Bach and Moulines, 2011). While these assumptions might be restrictive in practice, they enable to provide a tight analysis of the error of SGD. In such a setting, two types of proofs have been used traditionally. On one hand, Lyapunov-type proofs rely on controlling the expected squared distance to the optimal point (Zhao and Zhang, 2015). Such analysis suggests using small decaying steps, inversely proportional to the number of iterations (). On the other hand, studying the recursion as a stochastic process (Polyak and Juditsky, 1992) enables to better capture the reduction of the noise through averaging. It results in optimal convergence rates for larger steps, typically scaling as , (Bach and Moulines, 2011).
Over the past decade, the amount of available data has steadily increased: to adapt SGD to such situations, it has become necessary to distribute the workload between several machines, also referred to as workers (Delalleau and Bengio, 2007; Zinkevich et al., 2010; Recht et al., 2011). For SGD, two extreme approaches have received attention: 1) workers run SGD independently and at the end aggregate their results, called one-shot averaging (OSA) (Zinkevich et al., 2010; Godichon and Saadane, 2017) or parameter mixing, and 2) mini-batch averaging (MBA) (Dekel et al., 2012a; Takáč et al., 2013; Li et al., 2014c; Goyal et al., 2017; Jain et al., 2016), where workers communicate after every iteration: all gradients are thus computed at the same support point (iterate) and the algorithm is equivalent to using mini-batches of size , with the number of workers. While OSA requires only a single communication step, it typically does not perform very well in practice (Zhang et al., 2016). At the other extreme, MBA performs better in practice, but the number of communications equals the number of steps, which is a major burden, as communication is highly time consuming (Zhang et al., 2016). To optimize this computation-communication-convergence trade-off, we consider the Local-SGD framework: workers run SGD iterations in parallel and communicate periodically. This framework encompasses one-shot averaging and mini-batch averaging as special cases (see Figure 1).
We make the following contributions: 1) We provide the first non-asymptotic analysis for local-SGD with large step sizes (typically scaling as , for ), in both on-line and finite horizon settings. Our assumptions encompass the ubiquitous least-squares regression and logistic regression. 2) Our comparison of the two extreme cases, OSA and MBA, underlines the communication trade-offs. While both of these algorithms are asymptotically equivalent for a fixed number of machines, mini-batch theoretically outperforms one-shot averaging when we consider the precise bias-variance split. In the regime where both the number of machines and gradients grow simultaneously we show that mini-batch SGD outperforms one-shot averaging. 3) Under three different sets of assumptions, we quantify the frequency of communication necessary for Local SGD to be optimal (i.e., as good as mini-batch). Precisely, we show that the communication frequency can be reduced by as much as O\Big{(}\frac{\sqrt{T}}{P^{3/2}}\Big{)}, with gradients and workers. Moreover, our bounds suggest an adaptive communication frequency for logistic regression, which depending on the expected distance to the optimal point (a phenomenon observed by Zhang et al. (2016)). 4) We support our analysis by experiments illustrating the behavior of the algorithms.
The paper is organized as follows: in Section 2, we introduce the general setting, notations and algorithms, then in Section 2.2, we describe the related literature. Next, in Section 2.3, we describe assumptions made on the objective function. In Section 3, we provide our main results, their interpretation, consequence and comparison with other results. Results in the on-line setting and experiments are presented in the Section 4 and Appendix A.
Algorithms and setting
We consider machines, each of them running SGD. Periodically, workers aggregate (i.e., average) their models and restart from the resulting model. We denote by the number of communication steps. We define a phase as the time between two communication rounds. At phase , for any worker , we perform local steps of SGD. Iterations are thus naturally indexed by . We consider the lexicographic order on such pairs, which matches the order in which iterations are processed. Note that we assume the number of local steps is the same over all machines . While this assumption can be relaxed in practice, is facilitates our proof technique and notation. At any , we denote by the model proposed by worker , at phase , after local iterations. All machines initially start from the same point , that is for any , . The update rule is thus the following, for any :
Aggregation steps consist in averaging the final local iterates of a phase: for any , At phase , every worker restarts from the averaged model: . Eventually, we are interested in controlling the excess risk of the Polyak-Ruppert averaged iterate:
with . We use the notation to underline the fact that iterates are averaged over one phase and when averaging is made over all iterations. All averaged iterates can be computed on-line.
The algorithm, called local-SGD, is thus parameterized by the number of machines , communication steps , local iterations , the starting point , the learning rate , and the first order oracle on the gradient. Pseudo-code of the algorithm is given in Table 2.
Link with classical algorithms. Special cases of Local-SGD correspond to one-shot averaging or mini-batch averaging, as summarized in Table 2. More precisely, for a total number of gradients , with workers, communication rounds, and , we realize an instance of P-mini-batch averaging (P-MBA). On the other hand, with workers, communication, and , we realize an instance of one shot-averaging. Our goal is to get general convergence bounds for Local-SGD that recover classical bounds for both these settings when we choose the correct parameters. While comparing to Serial-SGD (which is also a particular case of the algorithm), would also be interesting, we focus here on the comparison between Local-SGD, one-shot averaging and mini-batch averaging. Indeed, the step size is generally increased for mini-batch with respect to Serial SGD, and the running efficiency of algorithms is harder to compare: we only focus on different algorithms that use the same number of machines.
2 Related Work
For quadratic functions, larger steps can be used, as pointed by Bach and Moulines (2013). Indeed, even with non-decaying step size, the averaged process converges to the optimal point. Several studies focus on understanding properties of SGD for quadratic functions: a detailed non-asymptotic analysis is provided by Défossez and Bach (2015), acceleration under the additive noise oracle (see Assumption 4 below) is studied by Dieuleveut et al. (2016) (without this assumption by Jain et al. (2017)), and Jain et al. (2016) analyze the effects of mini-batch and tail averaging.
Mini-batch averaging. Mini-batch averaging has been studied by Dekel et al. (2012a); Takáč et al. (2013). These papers show an improvement in the variance of the process, and make comparisons to SGD. It has been found that increasing the mini-batch size often leads to increasing generalization errors, which limits their distributivity (Li et al., 2014d). Jain et al. (2016) have provided upper bounds on learning-rate and mini-batch size for optimal performance. Recently, large mini-batches have been leveraged successfully in deep learning as in (Shirish Keskar et al., 2016; You et al., 2017; Goyal et al., 2017) by properly tuning learning rates, etc.
Local-SGD. Zhang et al. (2016) empirically show that local SGD performs well. They also provide a theoretical guarantee on the variance of the process, however, they assume the variance of the estimated gradients to be uniformly upper bounded (Assumption 4 below). Such an assumption is restrictive in practice, for example it is not satisfied for least squares regression. In a simultaneous work, Stich (2018) has provided an analysis for local-SGD. The limitation with their analysis is that they also assume bounded gradients and use a small step size scaling as . More importantly, their analysis doesn’t extend to the extreme case of one-shot averaging like ours. Lin et al. (2018) have experimentally shown that Local-SGD is better than the synchronous mini-batch techniques, in terms of overcoming the large communication bottleneck. Recently, Yu et al. (2018) have given convergence rates for the non-convex synchronous and a stale synchronous settings.
We have summarized the major limitations of some of these analyses in Table S3, given in Appendix H. Our motivation is to get away with some of these restrictive assumptions, and provide tight upper bounds for the above three averaging schemes. In the following section, we present the set of assumptions under which our analysis is conducted.
3 Assumptions
The function is strongly-convex with convexity constant .
If Q 1 is satisfied, then Assumptions 1, 2 are satisfied, and and are respectively the largest and smallest eigenvalues of . At any iteration , any machine can query an unbiased estimator of the gradient at a point . Formally, we make the following assumption:
A filtration is an increasing (i.e., for all , ), sequence of -algebras. 3 expresses that we have access to an i.i.d. sequence of unbiased estimators of . Remark that with such notations, for any , is -measurable. In Proposition 3.3, we make the additional, stronger assumption that the variance of gradient estimates is uniformly upper bounded, a standard assumption in the SGD literature, see e.g. Zhang et al. (2016):
Assumption 4 is for example true if the sequence of random vectors is i.i.d.. This setting is referred to as the semi-stochastic setting in Dieuleveut et al. (2016).
We also consider the following conditions on the regularity of the gradients, for :
Almost sure -co-coercivity (Zhu and Marcotte, 1996) is for example satisfied if for any , there exist a random function such that and which is a.s. convex and -smooth. Finally, we assume the fourth order moment of the random gradients at to be well defined:
It must be noted that 6 is a much weaker assumption than 4, for e.g., least-square regression satisfies former but not latter. Most of these assumptions are classical in machine learning. SGD for least squares regression satisfies Q 1, 3, 5 and 6. On the other hand, SGD for logistic regression satisfies 1, 2, 3 and 4. Our main result Theorem 3.6 (lower bounding the frequency of communications) applies to both these sets of assumptions. In Section B.3 we further detail how these assumptions apply in machine learning.
Though both of these approaches are often considered to be nearly equivalent (Bach, 2014; Dieuleveut and Bach, 2016), fundamental differences exist in their convergence properties. The on-line case is harder to analyze, but ultimately provides a better convergence rate. However as the behavior is easier to interpret in the finite horizon case, we postpone results for on-line setting to Section 4.
Moreover, we always assume that for any , the learning rate satisfies . In the following section, we present our main results.
Main Results
Sketch of the proof. We follow the approach by Polyak and Juditsky, which relies on the following decomposition: for any , Equation (2) is trivially equivalent to:
We have used a first order Taylor expansion around the optimal value of the gradient. Thus, using the definition of :
We first compare the results for Mini-batch averaging and One-shot averaging for finite horizon (FH) setting, and then provide these results for local-SGD.
First we assume the step size to be a constant at every iteration for any . Our first contribution is to provide non-asymptotic convergence rates for mini-batch SGD and one shot averaging, that allow a simple comparison. For the benefit of presentation, we define following quantities:
We have the following result for mini-batch averaging:
Under Assumptions 1, 2, 3, 5, 6, we have the following bound for mini-batch SGD: for any ,
The notation denotes inequality up to an absolute constant. Recall that for mini-batch, the total number of gradients processed is .
On the other hand, we also have the following result for one-shot averaging:
Under Assumptions 1, 2, 3, 5, 6, we have the following bound for one shot averaging: ,
Note that for one-shot averaging, the total number of gradients used is .
Interpretation, fixed . Using mini-batch naturally reduces the variance of the process . Equations 4 and 6 show that the speed at which the initial condition is forgotten remains the same, but that the variance of the local process is reduced by a factor .
Equations 5 and 7 show that the convergence depends on an initial condition term and a variance term. For a fixed number of machines , and a step size scaling as , , , the speed at which the initial condition is forgotten is asymptotically dictated by where , for both algorithms (if we use the same number of gradients for both algorithms, naturally, .) As for the variance term, it scales as as , as the remaining terms asymptotically vanish for . It reduces with the total number of gradients used in the process. Interestingly, this term is the same for the two extreme cases (MBA and OSA): it does not depend on the number of communication rounds. This phenomenon is often described as “the noise is the noise and SGD doesn’t care” (for asynchronous SGD, (Duchi et al., 2015)). Though we recover this asymptotic equivalence here, our belief is that this asymptotic point of view is typically misleading as the asymptotic regime is not always reached, and the residual terms do then matter.
Indeed, the lower order terms do have a dependence on the number of communication rounds: when the number of communications increases, the overall effect of the noise is reduced. More precisely, since the remaining terms are respectively or times smaller for mini-batch. This provides a theoretical explanation of why mini-batch SGD outperforms one shot averaging in practice. It also highlights the weakness of an asymptotic analysis: the dominant term might be equivalent, without reflecting the actual behavior of the algorithm. Disregarding communication aspects, mini-batch SGD is in that sense optimal.
Note that for quadratic functions, as . The conditions on the step size can thus be relaxed, and the asymptotic rates described above would be valid for any step size satisfying (Jain et al., 2016).
Extension to the on-line setting, eventually leading to a better convergence rate, is given in Proposition 4.1 in Section 4.
Interpretation, . When both the total number of gradients used and the number of machines are allowed to grow simultaneously, the asymptotic regime is not necessarily the same for MBA and OSA, as remaining terms are not always negligible. For example, if fixing , (we chose to balance and ), the variance term would be controlled by . Thus, unless , MBA could outperform OSA by a factor as large as .
Novelty and proofs. Both Propositions 3.1 and 3.2 are proved in the Appendix F. Importantly, Equations 4 and 6 respectively imply Equations 5 and 7 under the stated conditions: this is the reason why we only focus on proving equations similar to Equations 4 and 6 for Local-SGD.
Proposition 3.1 is similar to the analysis of Serial-SGD for large step size, but with a reduction in the variance proportional to the number of machines. Such a result is derived from the analysis by Dieuleveut et al. (2017), combining the approach of Bach and Moulines (2013) with the correct upper bound for smooth strongly convex SGD (Needell et al., 2014), and controlling similarly higher order moments. While this result is expected, we have not found it under such a simple form in the literature. Proposition 3.2 follows a similar approach, we combine the proof for mini-batch with a control of the iterates of each of the machines. This is closely related to Godichon and Saadane (2017), but we preserve a non-asymptotic approach.
Remark: link with convergence in function values. We mainly focus on proving convergence results on the Mahalanobis distance , which is the natural quantity in such a setting (Bach and Moulines, 2011, 2013; Godichon and Saadane, 2017). These results could be translated into function value convergence , using the inequality but the dependence on would be pessimistic and sub-optimal. However, a similar approach has been used by Bach (2014), under a slightly different set of assumptions (including self-concordance, e.g., for logistic regression), recovering optimal rates. Extension to such a set of assumptions, which relies on tracking other quantities, is an important direction.
While the “classical proof”, which provides rates for function values directly (with smoothness, or with uniformly bounded gradients) has a better dependence on , one cannot easily obtain a noise reduction when averaging between machines. Similarly, there is no proof showing that one-shot averaging is asymptotically optimal that relies only on function values. In other words, these proofs do not adequately capture the noise reduction due to averaging. Moreover, such proof techniques relying on function values typically involve a small step size (because the noise reduction is captured inefficiently). Such step size performs poorly in practice (initial condition is forgotten slowly), and is unknown.
In conclusion, though they do not directly result in optimal dependence on for function values, we believe our approach allows to correctly capture the effect of the noise, and is thus suitable for capturing the effect of local SGD.
Conclusion: for a fixed or limited number of machines, asymptotically, the convergence rate is similar for OSA and MBA. However, non-asymptotically, or when the number of machines also increases, the dominant terms can be as much as times smaller for MBA. In the following we provide conditions for Local-SGD to perform as well as MBA (while requiring much fewer communication rounds).
2 Convergence of Local-SGD, FH setting
For local-SGD we first consider the case of a quadratic function, under the assumption that the noise has a uniformly upper bounded variance. While this set of assumptions is not realistic, it allows an intuitive presentation of the results. Similar results for settings encompassing LSR and LR follow. We provide a bound on the moment of an iterate after the communication step (i.e., the restart point of the next phase), and on the second order moment of any iterate.
For , we denote .
Under Assumptions Q 1, 3, 4, we have the following bound for Local-SGD: for any ,
To prove such a result, we use the classical technique, and introduce a ghost sequence , and recursively control . We conclude by remarking that . This proof is given in Section C.2.
Interpretation. The variance bound for the iterates after communication, exactly behaves as in mini-batch case: the initialization term decays linearly with the number of local steps, and the variance is reduced proportionally to the number of workers . On the other hand, the bound on the iterates shows that the variance of this process is composed of a “long term” reduced variance, that accumulates through phases, and is increasingly converging to and of an extra variance , that increases within the phase, and is upper bounded by .
In the case of constant step size, the iterates of serial SGD converge to a limit distribution that depends on the step size (Dieuleveut et al., 2017). Here, the iterates after communication (or the mini-batch iterates) converge to a distribution with reduced variance , thus local iterates periodically restart from a distribution with reduced variance, then slowly “diverge” to the distribution with large variance. If the number of local iterations is small enough, the iterates keep a reduced variance. More precisely, we have the following result.
If for all , , then the second order moment of admits the same upper bound as the mini-batch iterate (Equation (4)) up to a constant factor of 2. As a consequence, Equation (5) is still valid, and Local-SGD performs optimally.
Interpretation. This result shows that if the algorithm communicates often enough, the convergence of the Polyak Ruppert iterate is as good as in the mini-batch case, thus it is “optimal”. Moreover, the minimal number of communication rounds is easy to define: the maximal number of local steps decays as the number of workers and the step size increases. This bound implies that more communication steps are necessary when more machines are used. Note that is a large number, as a typical value for is inversely proportional to (a power of) the number of local steps for e.g., , .
With constant number of local steps , and learning rate in order to obtain an optimal parallel convergence rate, local-SGD communicates times less as compared to mini-batch averaging.
We believe that this is the first result (with Stich (2018)) that shows a communication reduction proportional to a power of the number of local steps of a local solver (i.e., ), compared to mini-batch averaging.
In the following, we alternatively relax the bounded variance assumption 4 and the quadratic assumption Q 1, and show similar results for local SGD. This allows us to successively cover the cases of least squares regression (LSR) and logistic regression (LR).
Under either of the following sets of assumptions, the convergence of the Polyak Ruppert iterate is as good as in the mini-batch case, up to a constant:
Assume Q 1, 3, 5, 6, and for any , and .
These results are derived from Proposition S9 and Proposition S13 which generalize Proposition 3.3. Those results are proved in Appendix C and D and constitute the main technical challenge of the paper.
Interpretation. We note that in both of these situations, the optimal rates can be achieved if the communications happen often enough, and beyond such a number of communication rounds, there is no substantial improvement in the convergence. This result corresponds to the effect observed in practice (Zhang et al., 2016). The first set of assumption is valid for LSR, the second for LR. In the first case, the maximal number of local steps before communication is upper bounded by the same ratio as in Corollary 3.4, but the “constant” that appears is , so we need this quantity to be small (which is typically always satisfied in practice) in order to be optimal w.r.t. mini-batch averaging. A similar result as Example 3.5 can be provided reducing the communication by a factor of .
Main results: On-line Setting
In the on-line setting we consider the particular case of a decaying sequence , for some . The analysis is slightly more involved as Equation 3 results in more terms than in the finite horizon setting (sums do not directly telescope). While the decaying step-size case enables to improve some terms with respect to the finite horizon case (e.g. the speed at which one forgets the initial condition), the trade-offs concerning communication remain unchanged. We define the following constants to make the presentation clear, for :
Now we present a result similar to Proposition 3.1 for mini-batch averaging and one shot averaging:
Under the Assumptions 1, 2, 3, 5, 6 using we have for respectively mini-batch averaging and one-shot averaging:
with respectively and for one-shot averaging, and and for mini-batch averaging.
This proposition is directly derived from Lemma S15 in Appendix G. This proposition is similar to Propositions 3.1 and 3.2, but the overall convergence rate is better as using decaying step size eventually performs better. For example, the bias term mainly decays as instead of . This underlines why in practice decaying step size is often preferable. Asymptotically, the variance term is now dominant, and as before, MBA and OSA have similar performance as .
For a fixed number of machine , the bias is asymptotically vanishing, and if we ignore the linearly decaying terms and the dependence on , the resulting dominating term in is controlled by , which would result in an optimal choice of .
In the non asymptotic regime, where the total number of iterations and grow simultaneously, the variance of OSA scales as as long as . In other words, for , we need : the number of machines as to be smaller than the number of iterations to the power , in other words, for 1000 iterations, one could only use 10 machines to reach the asymptotic regime where OSA performs similarly to MBA.
Conclusion
Stochastic approximation and distributed optimization are both very densely studied research areas. However, in practice most distributed applications stick to bulk synchronous mini-batch SGD. While the algorithm has desirable convergence properties, it suffers from a huge communication bottleneck. In this paper we have analyzed a natural generalization of mini-batch averaging, Local SGD. Our analysis is non-asymptotic, which helps us to better understand the exact communication trade-offs. We give feasible lower bounds on communication frequency which significantly reduce the need for communication, while providing similar non-asymptotic convergence as mini-batch averaging. Our results apply to common loss functions, and use large step sizes. Further, our analysis unifies and extends all the scattered results for one-shot averaging, mini-batch averaging and local SGD, providing an intuitive understanding of their behavior.
Some important future directions are obtaining lower bounds, studying observable quantities to predict an adaptive communication frequency and relaxing some of the technical assumptions required by the analysis. The on-line case, experiments, proofs, additional materials and a review of distributed optimization follow in the appendix.
Acknowledgments
We thank Martin Jaggi, Sebastian Stichs, and Sai Praneeth Reddy for helpful discussions.
References
Appendix A Experimental results
We perform experiments for three different data-sets\hrefhttps://www.csie.ntu.edu.tw/ cjlin/libsvmtools/datasets/https://www.csie.ntu.edu.tw/ cjlin/libsvmtools/datasets/, two for least-squares regression and one for logistic regression Table S1. For all the curves we use v/s plots unless explicitly mentioned. Moreover, to elucidate the theory we use the same learning rates for all the algorithms in an experiment. The number of workers is set to every where, and plots are labeled w.r.t. the number of local steps which we don’t change along the different phases. We do the following experiments:
Performance of local SGD with different number of local steps spanning OSA to MBA (Figure S1). We globally find MBA to perform the best. Besides, as we increase the number of local steps the performance gets closer to OSA. This observation aligns with our theoretical guarantees. We use the averaged iterate (i.e., , the average over all the iterates till that point) for reporting the performance. The current iterate (i.e., , the ghost iterate for the current iteration) is omitted as the graphs are too noisy to be interpreted, and a variance of the loss is used instead.
Performance of local SGD with different number of local steps when started at the optimal point (Figure S2). We expect that if we start at then the bias term goes to zero and the difference between the algorithms becomes sharper. This is because our results predict that for constant learning rate, the initial conditions are forgotten at the same rate. We see that mini-batch outperforms OSA no the first iterations, but not asymptotically.
Variance of the estimators, for loss (Figure S3) and iterate values (Figure S4). We expect that a larger mini-batch size predicts a lower variance for these cases, and we observe the same through our experiments. In fact, the mean squared error of the parameters at the optimal is observed to be following a periodic curve. The value on an individual worker rises until it communicates, but always remains lower than a single SGD process run for the same number of iterations. This, verifies our theory and results for iterate convergence. Moreover, the variance at the loss function follows a similar pattern which elucidates the fact the intuitions developed in the paper also hold for functional convergence.
Appendix B Some Additional Material
Pseudo codes of both algorithms are given in Figure S7.
B.2 Summary of Results
In the table below, we specify for which algorithm our results apply (mini batch, one shot, or local SGD), under which assumptions they are proved and if they apply to the on-line setting(OL) or just the finite horizon(FH) case.
B.3 Example: Learning from i.i.d. observations
Stochastic Approximation: We use independent observations at each iteration. The total number of iterations is thus at most the number of observations we access. SGD then corresponds to following the gradient of the loss on a single independent observation . As the gradients we use are then unbiased gradients of the generalization error, this means that SGD directly minimizes this (unknown) function.
In practice, this means that in the first situation, we want to optimize the precision of the algorithm for a limited number of oracle calls, while in the second situation one would rather optimize the number of outer iterations of the algorithm (i.e. its running time). In both these assumptions, Assumption 3 is satisfied for the filtration generated by all the observations before time (respectively all the indices sampled before time ).
Then, Assumption 5 and 6 are satisfied, if is bounded and has finite variance.
SGD for least squares regression typically satisfies Q 1, 3, 5 and 6. On the other hand, SGD for logistic regression satisfies 1, 2, 3 and 4.
Appendix C Convergence guaranties for the second order moment.
In this section, we prove several Lemmas that allow to control the second order moment for the iterate. We first recall a few useful inequalities that will be used in the following. See for example Nesterov (2004).
We first recall the proof of the convergence for inner iterates. This proof corresponds to what happens on one machine, and can be found in the literature Bach and Moulines (2011); Dieuleveut et al. (2017) for example.
For any , under Assumptions 1, 2, 3, 5, 6, we have
Using the second equation recursively results in:
More precisely, for precise reference in the following proofs, we referenced this inequality with the following specific cases:
Under Assumptions 1, 2, 3, 5, 6, for mini-batch SGD with batch-size and step-size we have,
Such a result on reduced variance for mini-batch SGD () can be found in many previous works like Dekel et al. (2012b). Since mini-batch SGD is trivial to parallelize, this result also holds for the averaged iterate for outer iteration while using mini-batch averaging. Similarly, for decaying step sizes,
Similarly, in the case of one-shot averaging,
Under Assumptions 1, 2, 3, 5, 6 and a constant step-size using one-shot averaging, for any and we have,
C.2 Proof of Proposition 3.3
In this Section we prove Proposition 3.3. In order to provide a bound on the mean squared distance to the optimum of the outer iterates, we introduce a ghost sequence Mania et al. (2015), i.e., a sequence of iterates which is not actually computed. For any , we define
Under Assumptions Q 1, 3 and 4, for any , we have:
Remarking that for any , this implies the first inequality of Proposition 3.3. Note that this Lemma is valid for both decaying steps and and a constant learning rate. Especially, for a constant step size , and :
Under Assumptions Q 1, 3 and 4, for any , we have:
By induction, Lemma S5 implies that for any
To prove the second inequality of Proposition 3.3, we combine Lemma S5 and Equation (S5), using the fact that .
This results means that for a quadratic function with gradients having uniformly bounded variance, the outer iteration decay is the same as for mini-batch iterations (but for mini-batch, it is true under the weaker set of Assumptions 1, 2, 3, 5, 6).
By definition of , we have for any , using the linearity of (Assumption Q 1):
Under the independence of the noises (Assumption 3), then the uniform upper bound on the variance (Assumption 4), we have the following upper bound :
Under Assumption Q 1, is co-coercive, thus using Equation (S2), we have the following upper bound:
And using strong convexity (esp. Equation (S3)), and the fact that :
By recursion, we then have, for any :
C.3 Proof of Proposition S9
Under Assumptions Q 1,3,5,6, we have the following bound for one shot averaging: ,
with, for , , and .
When considering a constant step size , we have the following corollary.
Under Assumptions Q 1,3,5,6, we have the following bound for one shot averaging: , constant learning rate ,
Where we have and . Under the latter requirement (for optimality) that for any , , we have , thus this is generally a small constant. This result is a consequence of Lemma S11.
As before, the first bound shows that the variance of the iterates after communication is reduced by a factor of w.r.t. the serial case, thus almost as good as mini-batch averaging. However, the constants involved are worse than in the additive noise setting (Proposition 3.3). Consequently, and similarly to Proposition 3.3, the bound for the current iterates is composed of two terms for the variance: a “reduced variance” coming from the communication step, and a “inner loop” variance, that does not benefit from the number of machines.
Finally, we provide a convergence result in the most general case, removing the quadratic assumption. For the sake of concision, we skip the bound for the averaged iterate after a communication round, and directly give the result for the inner process.
C.3.2 Proof
This result is a consequence of Lemma S11, which implies Equation S15. Indeed, using it recursively, and using , we get:
Under Assumptions Q 1, 3, 5, 6, for any , we have:
The proof is a bit technical, so we summarize here the 2 main steps:
We prove an inequality (namely Equation S20) that is comparable to Equation S12, but with an extra term.
We use the control on the inner process (Section C.1) to control the extra term.
We consider again the ghost process defined at Equation (S6). Equations (S10) and (S11) are still valid. We now use the following decompositionIn the following, , etc. are used as symbolic notations to ease presentation.:
Using the independence of the noises (Assumption 3) we have,
Using Assumption 5 (co-coercivity for -s and ) we obtain,
This leads to, combining Equations (S10) and (S19), and the upper bound on the variance of the noise at the optimum (Assumption 6)
Using , and strong-convexity (Assumption 1)
This inequality should be compared to Equation (S12). It is interesting to remark that the last term is not an artifact of the proof: this is easy to check for least-squares regression.
Using recursively the above inequality and using the definition of , and taking expectation on the historical randomness we have, for any
Especially, for , , and moreover :
To upper bound the last term in the above equation, we use Equation S4,
Note that since the mean squared distance doesn’t depend on the machine, we can assume to be working on machine . This leads to, using an Abel transform:
We now use Equation S5. It leads to the following,
Combining Equations S21, S22 and S23, we get, denoting :
This concludes the proof of the Lemma, using :
This result can be used recursively. It implies that if , then the upper bound on the outer iterates is as good as the one for mini-batch, up to a constant.
C.4 Proof of Proposition S13
In this Section we prove the first upper bound of Corollary S14.
Finally, we provide a convergence result in the most general case, removing the quadratic assumption.
with .
if is uniformly bounded, we perform as well as minibatch SGD for the outer iterations (up to a constant).
For a constant step size , the proposition has the following corollary:
In the following, we alternatively relax the bounded variance assumption 4 and the quadratic assumption Q 1, and show similar results for local SGD. This allows us to successively cover the cases of least squares regression (LSR) and logistic regression (LR).
C.4.2 Proof
Proposition S13 follows from Lemma S15. We have for any ,
with .
For any , under Assumptions 1, 2, 3, 4 we have:
C.4.3 Proof of Lemma S15
We rely on the following decomposition. Almost surely, we have:
The first two lines correspond to the quadratic case (Equation S10), that has been analyzed in Lemma S11. The third term accounts for the difference between the mean gradient and the gradient at the mean point. We use Assumption 2 to control this term.
We then use the following Lemma, which control how the inner iterates deviate from their average :
For any , under Assumptions 1, 2, 3, 4 we have a.s.:
The proof of this Lemma is postponed to Section C.4.4.
Using Cauchy-Schwarz inequality and the bound on the third order derivative of , we have:
and, using a second order expansion of the gradient at together with Assumption 2 we have:
Using the proof of Equation S12, and combining Equations S24, S26 and S25 and Lemma S16, we have, for any :
Thus by induction, for any :
In the following section, we proved the auxiliary Lemma that was used in the proof.
C.4.4 Proof of Lemma S16
We now study as increases. Note that initially (), this quantity is 0. For any :
Thus, expanding and using cocoercivity Assumption:
Note that everything is tight until the last line for ( then for all , ). Under Assumption 4, we thus have:
Appendix D Convergence guaranties for the fourth order moment.
In this section, we prove several Lemmas that allow to control the fourth order moment of the iterate. While controlling the second order moment is sufficient for quadratic functions as no “residual” term appears in Equation 3 (the “residual” corresponds to the rest of a linear expansion of the gradient, which is thus exact for a quadratic function), in the general case, we also need to control the 4th order moment.
We first give guarantees for the inner iterates (within a phase) in Section D.1, then in the local SGD framework in Section D.2.
Here, we can use the following Lemma from Dieuleveut et al. (2017), that gives a recursion for the 4th order moment.
Under the Assumptions 1, 2, 3, 5 for th -order moment, assuming we have,
In the mini-batch setting, we have of course the same result with a variance reduction:
Under the Assumptions 1, 2, 3, 5 for th -order moment for mini-batch averaging we have, assuming we have,
Analogous to Lemma S1 we have the following result for fourth order moments,
Under the Assumptions 1, 2, 3, 5 for th -order moment, assuming we have,
Similarly for mini-batch analogous to Lemma S2,
Under the Assumptions 1, 2, 3, 5 for th -order moment for mini-batch averaging and decreasing step size we have, assuming we have,
The proof is included for completeness and because the same proof technique is used afterwards in Section D.2.
For , and we define the notation . We have that,
Above we have used Cauchy Schwartz inequality several times for the second inequality and equation (S28) for the third one.
Above we used in the last line. Finally, using strong convexity, we have:
D.2 Proof of Lemma S6
In this section, we prove the following Lemma, which is necessary to conclude the proof for the second set of Assumptions in Theorem 3.6. Indeed, we need to control the moment of order 4 to be able to control the residual term that arises from linear expansion of the gradient around .
There exist absolute constants , such that if :
This proof combines element from the classical bound for the fourth order moment, and from the proof of Lemma S15, which addresses the similar setting but only for the second order moment. We start from the definition of :
Thus, squaring this equation we get, denoting :
formally, we have used .
That is, conditioning on the past, and using Assumption 5 (cocoercivity and the fact that is a.s. -Lipshitz):
Rearranging terms and using the uniform upper bound on the 4-th moment of the noise 6, we have:
The first 2 lines of Equation S31 correspond to the expansion in S5 (the constants are slightly different because we use a uniform bound on the gradient instead of co-coercivity). The last two lines correspond to the residual term, for which we will use Lemma S16.
As a result, there exist absolute constants (“numbers”) , such that if :
Appendix E Main error decomposition
In this section, we prove the following decomposition for the on-line setting.
Under the differentiability of 2 we haveNote that after the final iteration of the phase the learning rate (which the algorithm uses nowhere) corresponds to the first learning rate for the next phase. This anomaly in notation is a direct result of us considering the ghost process, which runs continuously till the end.,
where and .
Below, we have as the stochastic gradient at step on machine for communication phase . After adding and subtracting few quantities and rearranging we have,
where and are respectively terms related to stochastic noise and quadratic residual. Obtaining the horizontal average over all the machines and recalling the definition of the ghost process as defined above we have,
Obtaining the vertical average over all the machines first within a communication phase and then among different phases we have,
Now recalling the definitions for the overall iterate , , the initial point , and the total number of gradients as we have defined above. After making these changes and on rearranging we obtain,
Thus we have obtained the required result as,
E.2 Bounding the noise term
The stochastic noise term which appears above can be bounded using the following lemma,
Using Assumptions 3, 5, 6 respectively we prove the result
Appendix F Proofs for OSA, MBA and Local-SGD in the finite horizon setting
Polyak and Judisky were mainly interested in the asymptotic analysis, and the set of assumptions considered was different.
In Bach and Moulines (2011), the authors prove comparable bounds in the case of bounded gradients. However, their analysis in the smooth and strongly convex setting is not optimal. Precisely, they use a sub-optimal upper bound when controlling the second order moments, that significantly worsens the subsequent proof. This point was underlined in Needell et al. (2014); Dieuleveut et al. (2017). The result they provide under our set of assumptions is eventually 1) not optimal, 2) uselessly complex, and 3) only for serial-SGD.
The result is an application of Jensen’s inequality with the convex function .
F.2 Proof of Proposition 3.1 (Mini-batch case)
Under the Assumptions 1, 2, 3, 5, 6 we have,
In order to upper bound the expectation we need to separately upper bound all the terms that appear in the result for Lemma S1. But before that we can actually simplify the result with constant step size and using as follows,
Now we bound each of the terms in the above decomposition one by one. For the first term,
For the third term using Lemma S1 and Lemma S3 we get,
Now using the upper bound from 2 followed by Lemma S2 we get,
For the fourth term, note that we are sampling i.i.d observations and thus the stochastic noise across all machines and iterations is independent and equal to zero in expectation (see 3). This implies the first equation below while the second inequality is obtained using Lemma S3,
Now using Lemma S1, we have proved the lemma.
It can be seen in the above lemma that there are two kinds of terms: one that depend on the history or initialization and second the ones that depend on the variance bound. This implies that it would be possible to restate Lemma S5 as follows,
Under the assumptions 1, 2, 3, 5, 6 we have,
Ignoring constants the above constants can be upper bounded as follows,
F.3 Proof Proposition 3.2 (One-shot averaging case)
To prove the proposition we need to prove a bound on second moment of the inner iterations followed by a bound on the final average outer iteration. For inner iterations we follow the result from Moulines and Bach (2011) as the process on a single worker is completely independent of any other worker. We have the following lemma,
Under the Assumptions 1, 2, 3, 5, 6 for constant step size for one shot averaging we have,
We follow the same line of proof as before. We can use the decomposition from Lemma S1 with constant step size and , which results in the following simpler decomposition,
For the second term using Lemma S3 and rearranging we have,
For the third term using Lemma S1 and Lemma S3 we obtain,
Now first using the upper bound of 2, followed by Lemma S1 and some rearranging we can obtain the following,
For the fourth term, using the fact that on different machines noise of the gradient is i.i.d. over different iterations and zero in expectation (3) we obtain,
Finally using Lemma S3, concludes the proof.
Similar to the mini-batch case, there are two kinds of terms one that depend on the history or initialization and second that depend on the variance bound of the functions. This implies that it would be possible to restate Lemma S8 as follows,
Under the Assumptions 3, 2, 1, 5, 6 we have,
On upper-bounding the above two terms while ignoring the constants,
Appendix G Proofs for OSA, MBA and Local-SGD in the online setting
Recall that the step size at iteration is defined as where . Though our results can be extended for the entire range of learning rates, we prove results only for .
We first state a few technical results which are helpful in the following proofs.
The proof simply follows from applying the inequality , followed by an integral bound over the series as . Note that it is possible to consider but the integral bound changes. For brevity we don’t include it here.
First we decompose the term, then use , followed by a series of integral bounds like Lemma S1,
For the gamma function we have, .
First we use an integral bound as , followed by the integral substitution after which the proof follows from the definition of the gamma function.
For the gamma function we have, .
First we use an integral bound as , followed by the integral substitution after which the proof follows from the definition of the gamma function.
For , .
It is a simple application of the integral bound on a decreasing function, .
G.2 Proof of Proposition 4.1 (Mini-batch Averaging Case)
We have the following lemma for mini-batch averaging for the decreasing step-size case,
Under the Assumptions 1, 2, 3, 5, 6 we have for mini-batch averaging,
Using again the decomposition in Lemma S1, we can obtain the following simpler version for mini-batch averaging,
Note again that we assume , just for the sake of brevity. For the first term,
For the second term using Lemma S2, followed by Lemma S1 and Lemma S3 we obtain,
For the third term using Lemma S11 and ,
Now using Lemma S2, Lemma S1, Lemma S3 and we get,
Now using Lemma S7 (with , and ), followed by using Lemma S7 again (with , and ) and Lemma S9 (with ) we get,
Finally using Lemma S1 and re-organizing with constants defined as above,
For the fourth term first proceeding as in Lemma S5 with Lemma S1 and Lemma S3 we can obtain,
Now using Lemma S3, followed by Lemma S1 and Lemma S3 we getNote that we ignore t=1 in second inequality for second term as we have already incorporated it in the first term,
Now using Lemma S5 (with and ), followed by Lemma S5 again (with and ), followed by Lemma S9 (with ) and Lemma S1 we get,
Bounding again with the constants defined above,
For the fifth term, proceeding as in Lemma S5,
Now using Lemma S2, Lemma S1 and Lemma S3 like before,
Further using Lemma S5 (with and ), followed by Lemma S5 again (with and ), followed by Lemma S9 (with ) and the constants as used above we get,
Finally using Lemma S1 we have proved the lemma.
The following lemma separates the terms above into bias and variance terms, following which we can easily prove Proposition 4.1,
Under the Assumptions 1, 2, 3, 5, 6 we have for mini-batch averaging,
Where for constants defined as above the terms are,
To get Proposition 4.1, we upper bound every term up to constants depending only on . Specifically, we use , , and .
G.3 Proof of Proposition 4.1 (One-shot Averaging case)
The analysis for the one-shot case is very similar to the mini-batch case, just like the constant step-size case. In fact at many place the communications of MBA get replaced by and the form of the bound remains the same. This intuitive conversion strengthens our analysis, which smoothly extends to both the extreme cases.
Under the Assumptions 1, 2, 3, 5, 6 for decreasing step size, for one shot averaging we have,
And the constants are and .
We follow an analysis similar to Godichon and Saadane (2017). We can simplify the decomposition from Lemma S1 for one outer phase as follows,
For the second term note that the inner iterate bound is independent for different machines using Lemma S4 for say machine , followed by Lemma S1 and Lemma S3 we get,
For the third term using , Lemma S11, and noting that the individual bounds on inner iterates for different machines are the same, thus using machine for brevity we can obtain,
Now using Lemma S4, Lemma S1, Lemma S3 and we get,
Now using Lemma S7 again with and with defined as above and Lemma S9 we get,
Now for the fourth term proceeding as in Lemma S8 with Lemma S1 and Lemma S3 we can obtain ,
Now first using the upper bound of 2, followed by Lemma S3, Lemma S1, Lemma S3 and Lemma S5 we can obtain the following,
For the fifth term, using the fact that for different machines noise is independent, zero in expectation (3) we obtain,
Now using Lemma S4, followed by Lemma S1, Lemma S3 and Lemma S5 with definition of as before, and we have,
Thus using Lemma S1 we have proved the lemma.
We can get the following lemma combining the bias and variance terms separately,
Under the Assumptions 1, 2, 3, 5, 6 for decreasing step size, for one shot averaging we have,
Where for constants defined as above the terms are,
Appendix H Brief overview of distributed optimization
The above three schemes (OSA, MBA, Local-SGD) are the most studied synchronous parallel schemes. However, communication latencies often make it difficult to use these algorithms for large-scale problems. Thus many alternative parallelization schemes which minimize communication or perform better have been studied. The major problem with some of these variants is that they are often difficult to tune, are not as stable and don’t scale well to non-convex optimization problems. Result-wise, most of the machine learning packages use centralized mini-batch synchronous SGD.
Asynchronous SGD: These techniques are characterized by avoiding a centralized synchronization, using delayed updates, maintaining parameter server estimates and being fault tolerant. Some of the notable references in a chronological order are Langford et al. (2009); Niu et al. (2011); Agarwal and Duchi (2011); Paine et al. (2013); Li et al. (2014b); Zhang et al. (2014); Keuper and Pfreundt (2015); De and Goldstein (2015); Feyzmahdavian et al. (2015); Lian et al. (2015); Mania et al. (2015); Zhao and Li (2015); Duchi et al. (2015); Chen et al. (2016a); Lian et al. (2017a); Pedregosa et al. (2017); Lian et al. (2017b); Leblond et al. (2018); Alistarh et al. (2018).
Federated optimization: This setting is characterized by a huge number of mobile user devices, which run their local model in a decentralized manner with often unbalanced data, but aim to train jointly. Many research questions still remain open but the direction is very relevant for distributed AI. Some references are Konecný et al. (2015, 2016); McMahan et al. (2016).
Compressed Communication: A common strategy to combat the communication overhead is to introduce lossless or lossy compression of exchanged information, often the gradients. Some of the work in this direction can be found in Zhang et al. (2017); Wen et al. (2017); Wangni et al. (2017); Sa et al. (2015); Na et al. (2017); Gupta et al. (2015); Alistarh et al. (2016); Khirirat et al. (2018).
Non-SGD methods: Many other optimization algorithms (coordinate descent, quasi newton, etc.) have also been studied in the parallel setting, owing to their better distributivity or convergence for some applications compared to the SGD algorithm. Some of them are Boyd et al. (2011) (ADMM), Shamir et al. (2014) (DANE), Zhang and Xiao (2015) (DiSCO), Reddi et al. (2016) (AIDE), Ma et al. (2017); Smith et al. (2016); Ma et al. (2015) (COCOA) and some of the references therein. Recently Scaman et al. (2017) gave provably optimal algorithms for the strongly convex and smooth functions for both synchronous and asynchronous cases. More broadly speaking, variance reduction methods are often the methods of choice in better understood, convex optimization problems [add reference]. Yet, their usage in the deep learning community has been relatively scarce, and often they are more difficult to parallelize [add reference]. Some of the works for instance are Reddi et al. (2015); Zhao and Li (2016); De and Goldstein (2015); Lee et al. (2015). Among second order methods, quasi newton methods like distributed L-BFGS Najafabadi et al. (2017); Şimşekli et al. (2018) are also widely popular among the machine learning community.
Communication Lower Bounds: On a broader level our work is related to communication lower bounds which arise from information and learning-theoretic considerations. Unfortunately, these bounds are difficult to match for convex optimization as they are provided in Arjevani and Shamir (2015). Similar bounds have also been provided for the generally easier statistical estimation setting in Duchi et al. (2014); Braverman et al. (2015); Zhang et al. (2013).
Feature Distribution: As clearly evident training data is not the only element of our optimization scheme which can be parallelized. Often in many problems in natural language processing and linear estimation, the features number in hundreds of thousands, and it might be of some interest to distribute the features alongside or beside training data. Some relevant references are Lee et al. (2014); Ma and Takáč (2015); Smith et al. (2016); Chen et al. (2016b); Fang and Klabjan (2018).
There has also been work in parallelizing stochastic optimization algorithms for specific problems (like PCA) in the past, for e.g., Mcdonald et al. (2009); McDonald et al. (2010); Meng et al. (2012); Zhuang et al. (2013); Li et al. (2014a); Chin et al. (2015); Oh et al. (2015).
We also provide a brief overview of some other techniques in distributed optimization in Appendix H.