CoCoA: A General Framework for Communication-Efficient Distributed Optimization
Virginia Smith, Simone Forte, Chenxin Ma, Martin Takac, Michael I. Jordan, Martin Jaggi
Introduction
Distributed computing architectures have come to the fore in modern machine learning, in response to the challenges arising from a wide range of large-scale learning applications. Distributed architectures offer the promise of scalability by increasing both computational and storage capacities. A critical challenge in realizing this promise of scalability is to develop efficient methods for communicating and coordinating information between distributed machines, taking into account the specific needs of machine-learning algorithms.
On most distributed systems, the communication of data between machines is vastly more expensive than reading data from main memory and performing local computation. Moreover, the optimal trade-off between communication and computation can vary widely depending on the dataset being processed, the system being used, and the objective being optimized. It is therefore essential for distributed methods to accommodate flexible communication-computation profiles while still providing convergence guarantees.
Although numerous distributed optimization methods have been proposed, the mini-batch optimization approach has emerged as one of the most popular paradigms for tackling this communication-computation tradeoff (e.g., Dekel et al., 2012; Shalev-Shwartz and Zhang, 2013b; Shamir and Srebro, 2014; Qu et al., 2015; Richtárik and Takáč, 2016). Mini-batch methods are often developed by generalizing stochastic methods to process multiple data points at a time, which helps to alleviate the communication bottleneck by enabling more distributed computation per round of communication. However, while the need to reduce communication would suggest large mini-batch sizes, the theoretical convergence rates of these methods tend to degrade with increased mini-batch size, reverting to the rates of classical (batch) gradient methods. Empirical results corroborate these theoretical rates, and in practice, mini-batch methods have limited flexibility to adapt to the communication-computation tradeoffs that would maximally leverage parallel execution. Moreover, because mini-batch methods are typically derived from a specific single-machine solver, these methods and their associated analyses are often tailored to specific problem instances and can suffer both theoretically and practically when applied outside of their restricted setting.
In this work, we propose a framework, CoCoA CoCoA-v1 (Jaggi et al., 2014) and CoCoA (Ma et al., 2017b, 2015a) are predecessors of this work. We continue to use the name CoCoA for the more general framework proposed here, and show how earlier work can be derived as a special case (Section 4). Portions of this newer work additionally appear in SF’s master’s thesis (Forte, 2015) and Smith et al. (2015). , that addresses these two fundamental limitations. First, we allow arbitrary local solvers to be used on each machine in parallel. This allows the framework to directly incorporate state-of-the-art, application-specific single-machine solvers in the distributed setting. Second, the framework shares information between machines through a highly flexible communication scheme. This allows the amount of communication to be easily tailored to the problem and system at hand, in particular allowing for the case of significantly reduced communication in the distributed environment.
A key step in providing these features in the framework is to first define meaningful subproblems for each machine to solve in parallel, and to then combine updates from the subproblems in an efficient manner. Our method and convergence results rely on noting that, depending on the distribution of the data (e.g., by feature or by training point), and whether we solve the problem in the primal or the dual, certain machine learning objectives can be more easily decomposed into subproblems in the distributed setting. In particular, we categorize common machine learning objectives into several cases, and use duality to help decompose these objectives. Using primal-dual information in this manner not only allows for efficient methods (achieving, e.g., up to 50x speedups compared to state-of-the-art), but also allows for strong primal-dual convergence guarantees and practical benefits such as computation of the duality gap for use as an accuracy certificate and stopping criterion.
We develop a communication-efficient primal-dual framework that is applicable to a broad class of convex optimization problems. Notably, in contrast to earlier work of Yang (2013), Jaggi et al. (2014), Ma et al. (2017b) and Ma et al. (2015a), our generalized, cohesive framework: (1) specifically incorporates difficult cases of regularization and other non-strongly-convex regularizers; (2) allows for the flexibility of distributing the data by either feature or training point; and (3) can be run in either a primal or dual formulation, which we show to have significant theoretical and practical implications.
Two key advantages of the proposed framework are its communication efficiency and ability to employ off-the-shelf single-machine solvers internally. On real-world systems, the cost of communication versus computation can vary widely, and it is thus advantageous to permit a flexible amount of communication depending on the setting at hand. Our framework provides exactly such control. Moreover, we allow arbitrary solvers to be used on each machine, which permits the reuse of existing code and the benefits from multi-core or other optimizations therein. We note that beyond the selection of the local solver and communication vs. computation profile, there are no required hyperparameters to tune; the provided default parameters ensure convergence and are used throughout our experiments to achieve state-of-the-art performance.
We derive convergence rates for the framework, guaranteeing, e.g., a rate of convergence in terms of communication rounds for convex objectives with Lipschitz continuous losses, and a faster linear rate for strongly convex losses. Importantly, our convergence guarantees do not degrade with the number of machines, , and allow for subproblems to be solved to arbitrary accuracies, which allows for highly flexible computation vs. communication profiles. Additionally, we leverage a novel approach in the analysis of primal-dual rates for non-strongly-convex regularizers. The proposed technique is an improvement over simple smoothing techniques used in, e.g., Nesterov (2005), Shalev-Shwartz and Zhang (2014) and Zhang and Lin (2015) that enforce strong convexity by adding a small term to the objective. Our results include primal-dual rates and certificates for the general class of linear regularized loss minimization, and we show how earlier work can be derived as a special case of our more general approach.
The proposed framework yields order-of-magnitude speedups (as much as 50 faster) compared to state-of-the-art methods for large-scale machine learning. We demonstrate these gains with an extensive experimental comparison on real-world distributed datasets. We additionally explore properties of the framework itself, including the effect of running the framework in the primal or the dual, and the impact of subproblem accuracy on convergence. All algorithms for comparison are implemented in Apache Spark and run on Amazon EC2 clusters. Our code is available at: gingsmith.github.io/cocoa/.
Background and Setup
In this paper we develop a general framework for minimizing problems of the following form:
The following standard definitions will be used throughout the paper.
for any , where denotes the subdifferential of at .
2 Primal-Dual Setting
Numerous methods have been proposed to solve (I), and these methods generally fall into two categories: primal methods, which run directly on the primal objective, and dual methods, which instead run on the dual formulation of the primal objective. In developing our framework, we present an abstraction that allows for either a primal or a dual variant of our framework to be run. In particular, to solve the input problem (I), we consider mapping the problem to one of the following two general problems:
In general, our aim will be to compute a minimizer of the problem (A) in a distributed fashion; the main difference will be whether we initially map the primal (I) to (A) or (B).
The dual relationship between problems (A) and (B) is known as Fenchel-Rockafellar duality (Borwein and Zhu, 2005, Theorem 4.4.2). We provide a self-contained derivation of the duality in Appendix B. Note that while dual problems are typically presented as a pair of (min, max) problems, we have equivalently reformulated (A) and (B) to both be minimization problems in accordance with their roles in our framework.
This mapping arises from first-order optimality conditions on the -part of the objective. The duality gap, given by:
is always non-negative, and under strong duality, the gap will reach zero only for an optimal pair . The duality gap at any point provides a practically computable upper bound on the unknown primal as well as dual optimization error (suboptimality), since
In developing the proposed framework, noting the duality between (A) and (B) has many benefits, including the ability to compute the duality gap, which acts as a certificate of the approximation quality. It is also useful as an analysis tool, helping us to present a cohesive framework and relate this work to the prior work of Yang (2013), Jaggi et al. (2014) and Ma et al. (2015a, 2017b). As a word of caution, note that we avoid prescribing the name “primal” or “dual” directly to either of the problems (A) or (B), as we demonstrate below that their role as primal or dual can change depending on the application problem of interest.
3 Assumptions and Problem Cases
Our main assumptions on problem (A) are that is -smooth, and the function is separable, i.e., , with each having -bounded support. Given the duality between the problems (A) and (B), this can be equivalently stated as assuming that in problem (B), is -strongly convex, and the function is separable with each being -Lipschitz.
In Section 3, we will see that different variants of the framework may be realized depending on which of these three cases we consider when solving the input problem (I).
4 Running Examples
To illustrate the cases in Table 1, we consider several examples below. Those interested in details of the framework itself may skip to Section 3. These applications will serve as running examples throughout the paper, and we will revisit them in our experiments (Section 6). For further applications and details, see Section 5.
Elastic Net Regression (Case I: map to either (A) or (B)). We can map elastic-net regularized least squares regression,
to either objective (A) or (B). To map to objective (A), we let: and , setting to be the number of features and the number of training points. To map to (B), we let: and , setting to be the number of features and the number of training points. We discuss in Section 3 how the choice of mapping to either (A) or to (B) can have implications on the distribution scheme and overall performance of the framework.
Lasso (Case II: map to (A)). We can represent -regularized least squares regression by mapping the model:
to objective (A), letting and . In this mapping, represents the number of features, and the number of training points. Note that we cannot map the lasso objective to (B) directly, as must be -strongly convex and the -norm is non-strongly convex.
Support Vector Machine (Case III: map to (B)). We can represent a hinge loss support vector machine (SVM) by mapping the model:
to objective (B), letting and . In this mapping, represents the number of features, and the number of training points. Note that we cannot map the hinge loss SVM primal to objective (A) directly, as must be -smooth and the hinge loss is non-smooth.
5 Data Partitioning
The CoCoA Method
In the following sections, we describe the proposed framework, CoCoA, at a high level, and then discuss two approaches for using the framework in practice: CoCoA in the primal, where we consider (A) to be the primal objective and run the framework on this problem directly, and CoCoA in the dual, where we instead consider (B) to be the primal objective, and then run the framework on the dual (A).
Note that in both approaches, the aim will be to compute a minimizer of the problem (A) in a distributed fashion; the main difference will be whether we view (A) as the primal objective or as the dual objective.
The goal of the CoCoA framework is to find a global minimizer of the objective (A), while distributing computation based on the partitioning of the dataset across machines (Section 2.5). As a first step, note that distributing the update to the function in objective (A) is straightforward, as we have required that this term is separable according to the partitioning of our data, i.e., . However, the same does not hold for the term . To minimize this part of the objective in a distributed fashion, we propose minimizing a quadratic approximation of the function, which allows the minimization to separate across machines. We make this approximation precise in the following subsection.
In the general CoCoA framework (Algorithm 1),
and . Here we let denote the change of local variables for indices , and we set for all . It is important to note that the subproblem (10) is simple in the sense that it is always a quadratic objective (apart from the term). The subproblem does not depend on the function itself, but only its linearization at the fixed shared vector . This property additionally simplifies the task of the local solver, especially for cases of complex functions .
There are two parameters that must be set in the framework: , the aggregation parameter, which controls how the updates from each machine are combined, and , the subproblem parameter, which is a data-dependent term measuring the difficulty of the data partitioning . These terms play a crucial role in the convergence of the method, as we demonstrate in Section 4. In practice, we provide a simple and robust way to set these parameters: For a given aggregation parameter , the subproblem parameter will be set as , but can also be improved in a data-dependent way as we discuss below. In general, as we show in Section 4, setting and will guarantee convergence while delivering our fastest convergence rates.
In Algorithm 1, the aggregation parameter controls the level of adding () versus averaging () of the partial solutions from all machines. For our convergence results (Section 4) to hold, the subproblem parameter must be chosen not smaller than
The simple choice of is valid for (11), i.e.,
In some cases, it will be possible to give a better (data-dependent) choice for , closer to the actual bound given in .
Here we provide further intuition behind the data-local subproblems (10). The local objective functions are defined to closely approximate the global objective in (A) as the “local” variable varies, which we will see in the analysis (Appendix D, Lemma ‣ 4.1). In fact, if the subproblem were solved exactly, this could be interpreted as a data-dependent, block-separable proximal step, applied to the part of the objective (A) as follows:
where
However, note that in contrast to traditional proximal methods, CoCoA does not assume that this subproblem is solved to high accuracy, as we instead allow the use of local solvers of any approximation quality .
The local subproblems (10) have the appealing property of being very similar in structure to the global problem (A), with the main difference being that they are defined on a smaller (local) subset of the data, and are simpler because they are not dependent on the shape of . For a user of CoCoA, this presents a significant advantage in that existing single machine-solvers can be directly re-used in our distributed framework (Algorithm 1) by employing them on the subproblems .
Therefore, problem-specific tuned solvers which have already been developed, along with associated speed improvements (such as multi-core implementations), can be easily leveraged in the distributed setting. We quantify the dependence on local solver performance with the following assumption and remark, and relate this performance to our global convergence rates in Section 4.
We assume that there exists such that , the local solver at any outer iteration produces a (possibly) randomized approximate solution , which satisfies
In practice, the time spent solving the local subproblems in parallel should be chosen comparable to the time of a communication round, for best overall efficiency on a given system. We study this trade-off in theory (Section 4) and experiments (Section 6).
Note that the accuracy parameter does not have to be chosen a priori: Our convergence results (Section 4) are valid if is an upper bound on the actual empirical values in Algorithm 1. This allows for some of the machines to at times deliver better or worse accuracy (e.g., this would allow a slow local machine to be stopped early during a specific round in order to avoid stragglers). See (Smith et al., 2017) for more details.
From a theoretical perspective, the multiplicative notion of accuracy is advantageous over classical additive accuracy as existing convergence results for first- and second-order optimization methods typically appear in multiplicative form, i.e., relative to the error at the initialization point (here ). This accuracy notion is also useful beyond the distributed setting (see, e.g., Karimireddy et al., 2018a, b). We discuss local solvers and associated rates to achieve accuracy for particular applications in Section 5.
With this general framework in place, we next discuss two variants of our framework, CoCoA-Primal and CoCoA-Dual. In running either the primal or dual variant of the framework, the goal will always be to solve objective (A) in a distributed fashion. The main difference will be whether this objective is viewed as the primal or dual of the input problem (I). We make this mapping technique precise and discuss its implications in the following subsections (Sections 3.2–3.4).
2 Primal Distributed Optimization
In the primal distributed version of the framework (Algorithm 2), the framework is run by mapping the initial problem (I) directly to objective (A) and then applying the generalized CoCoA framework described in Algorithm 1. In other words, we view problem (A) as the primal objective, and solve this problem directly.
From a theoretical perspective, viewing (A) as the primal will allow us to consider non-strongly convex regularizers, since we allow the terms to be non-strongly convex. This setting was not covered in earlier work of Yang (2013); Jaggi et al. (2014); Ma et al. (2015a); and Ma et al. (2017b), and we discuss it in detail in Section 4, as additional machinery must be introduced to develop primal-dual rates for this setting.
Running the primal version of the framework has important practical implications in the distributed setting, as it typically implies that the data is distributed by feature rather than by training point. In this setting, the amount of communication at every outer iteration will be of training points. When the number of features is high (as is common when using sparsity-inducing regularizers) this can help to reduce communication and improve overall performance, as we demonstrate in Section 6.
3 Dual Distributed Optimization
In the dual distributed version of the framework (Algorithm 3), we run the framework by mapping the original problem (I) to objective (B), and then solve the problem by running Algorithm 1 on the dual (A). In other words, we view problem (B) as the primal, and solve this problem via the dual (A).
This version of the framework will allow us to consider non-smooth losses, such as the hinge loss or absolute deviation loss, since the terms can be non-smooth. From a practical perspective, this version of the framework will typically imply that the data is distributed by training point, and for a vector of features to be communicated at every outer iteration. This variant may therefore be preferable when the number of training points exceeds the number of features.
4 Primal vs. Dual
In the following subsection, we provide greater insight into the CoCoA framework and its relation to prior work. An extended discussion on related work is available in Section 7.
5 Interpretations of CoCoA
There are numerous methods that have been developed to solve (A) and (B) in parallel and distributed environments. We describe related work in detail in Section 7, and here briefly position CoCoA and in relation to other widely-used parallel and distributed methods.
We first contrast CoCoA with common distributed mini-batch and batch methods, such as mini-batch stochastic gradient descent or coordinate descent, gradient descent, and quasi-Newton methods.
The Jacobi method does not require information from the other coordinates to update coordinate , which makes this style of method well-suited for parallelization. However, the sequential Gauss-Seidel-style method tends to converge faster in terms of iterations, as it is able to incorporate information from the updates of other coordinates more quickly. This difference is well-known and evident in single machine solvers, where stochastic methods (benefiting from fresh updates) tend to outperform their batch counterparts.
Typical mini-batch methods, e.g., mini-batch coordinate descent, perform a Jacobi-style update on a subset of the coordinates at each iteration. This makes these methods amenable to high levels of parallelization. However, they are unable to incorporate information as quickly as their serial counterparts in terms of number of data points accessed, as synchronization is required before updating the coordinates. As the size of the mini-batch grows, this can increase the runtime and even lead to divergence (Richtárik and Takáč, 2016).
CoCoA instead aims to combine attractive properties of both of these update paradigms. In CoCoA, Jacobi-style updates are applied in parallel to blocks of the coordinates of to distribute the method, while allowing for (though not necessarily requiring) faster Gauss-Seidel-style updates on each machine. This change in parallelization scheme is one of the key reasons for improved performance over simpler mini-batch or batch style methods.
In addition to the parallel block-Jacobi updating scheme described above, CoCoA incorporates an additional level of flexibility by allowing for an arbitrary number of sequential Gauss-Seidel iterations (or any local solver for that matter) to be performed locally on each machine. This flexibility is critical in the distributed setting, as one of the key indicators of parallel efficiency is the time spent on local computation vs. communication. In particular, the flexibility to solve each subproblem to arbitrary accuracy, , allows CoCoA to scale from low-communication environments, where more iterations can be performed before communicating, to high communication environments, where fewer local iterations are necessary.
In comparison with other distributed methods, this flexibility also affords an explanation of CoCoA as a method that can freely move between two extremes. On one extreme, if the subproblems (10) are solved exactly, CoCoA recovers block coordinate descent, where the coordinate updates are applied as part of a block-separable proximal step (12). If only one outer round of communication is performed, this is similar in spirit to one-shot communication schemes, which attempt to completely solve for and then combine locally-computed models (see, e.g., Mann et al., 2009; Zhang et al., 2013; Heinze et al., 2016). While these one-shot communication schemes are ideal in terms of reducing communication, they are, in contrast to CoCoA, generally not guaranteed to converge to the optimal solution.
On the other extreme, if just a single update (i.e., with respect to one coordinate ) is performed at each communication round, this recovers traditional distributed coordinate descent. In comparison to CoCoA, vanilla distributed coordinate descent can suffer from a high communication bottleneck due to the low relative amount of local computation. Even in the case of mini-batch coordinate descent, the most amount of work that can be performed locally at each round includes a single pass through the data, whereas CoCoA has the flexibility to take multiple passes. We empirically compare to mini-batch distributed coordinate descent in Section 6 to demonstrate the effect of this issue in practice.
Finally, we provide a direct comparison between CoCoA and the alternating direction method of multipliers (ADMM), a well-established framework for distributed optimization (Boyd et al., 2010). Similar to CoCoA, ADMM defines a subproblem for each machine to solve in parallel, rather than parallelizing a mini-batch update. ADMM also leverages duality structure, similar to that presented in Section 2. For consensus ADMM (Mota et al., 2013), (B) is decomposed with a re-parameterization:
This problem is then solved by constructing the augmented Lagrangian, which yields the following decomposable updates:
where is a penalty parameter that must be tuned for best performance. It can be shown that the update to can be reformulated in terms of the conjugate functions as:
Thus, we see that the update to closely matches the CoCoA subproblem (10), where . This is intuitive as both methods use proximal steps to derive the subproblem, and a similar result can be shown when applying ADMM to the (A) formulation, which can be seen as an instantiation of the sharing variant of ADMM (Boyd et al., 2010, Section 7.3).
However, there remain major differences between the methods despite this connection. First, CoCoA has a more direct and simplified scheme for updating the global weight vector , as the additional proximal step is not required. Second, in CoCoA, there is no need to tune any parameters such as , as the method can be run simply using . Finally, in the CoCoA method and theory, the subproblem can be solved approximately to any accuracy , rather than requiring a full batch update as in ADMM. We will see in our experiments that these differences have a substantial impact in practice (Section 6). We provide a full derivation of the comparison to ADMM for reference in Appendix C.
Convergence Analysis
In this section, we provide convergence rates for the proposed framework and introduce a key theoretical technique in analyzing non-strongly convex terms in the primal-dual setting.
For simplicity of presentation, we assume in the analysis that the data partitioning is balanced, i.e., for all . Furthermore, we assume that the columns of A satisfy for all , and contains the average term , as is common in ERM-type problems. We present rates for the case where in Algorithm 1, and where the subproblems (10) are defined using the corresponding safe bound . This case will guarantee convergence while delivering our fastest rates in the distributed setting, which in particular do not degrade as the number of machines increases and remains fixed. More general rates and all proof details can be found in the appendix.
To guarantee convergence, it is critical to show how progress made on the local subproblems (10) relates to the global objective . Our first lemma provides exactly this information. In particular, we see that if the aggregation and subproblem parameters are selected according to Definition ‣ 3.1, the sum of the subproblem objectives, , will form a block-separable upper bound on the global objective .
A proof of Lemma ‣ 4.1 is provided in Appendix D. We use this main lemma, in combination with our measure of quality of the subproblem approximations (Assumption 1), to deliver global convergence rates.
Our first main theorem provides convergence guarantees for objectives with general convex (or, equivalently, -Lipschitz ), including models with non-strongly convex regularizers such as lasso and sparse logistic regression, or models with non-smooth losses, such as the hinge loss support vector machine.
Consider Algorithm 1 with , and let be the quality of the local solver as in Assumption 1. Let have -bounded support, and let be -smooth. Then after iterations, where
we have that the expected duality gap satisfies
where is the averaged iterate: .
Providing primal-dual rates and globally defined primal-dual accuracy certificates for these objectives may require a theoretical technique that we introduce below, in which we show how to satisfy the notion of -bounded support for , as stated in Definition ‣ 2.1.
To address this problem, existing approaches typically use a simple smoothing technique (e.g., Nesterov, 2005; Shalev-Shwartz and Zhang, 2014): by adding a small amount of regularization, the functions become strongly convex. Following this change, the methods are run on the dual instead of the original primal problem. While this modification satisfies the necessary assumptions for convergence, this smoothing technique is often undesirable in practice, as it changes the iterates, the algorithms at hand, the convergence rate, and the tightness of the resulting duality gap compared to the original objective. Further, the amount of smoothing can be difficult to tune and has a large impact on empirical performance. We perform experiments to highlight these issues in practice in Section 6.
In contrast to smoothing, our approach preserves all solutions of the original objective, leaves the iterate sequence unchanged, and allows for direct reusability of existing solvers for the original objectives (such as solvers). It also removes the need for tuning a smoothing parameter. To achieve this, we modify the function by imposing an additional weak constraint that is inactive in our region of interest. Formally, we replace by the following modified function:
For large enough , this problem yields the same solution as the original objective. Note also that this only affects convergence theory, in that it allows us to present a strong primal-dual rate (Theorem ‣ 4.2 for =). The modification of does not affect the algorithms for the original problems. Whenever a monotone optimizer is used, we will never leave the level set defined by the objective at the starting point.
Using the resulting modified function will allow us to apply the results of Theorem ‣ 4.2 for general convex functions . This technique can also be thought of as “Lipschitzing” the dual , because of the general result that is -Lipschitz if and only if has -bounded support (Rockafellar, 1997, Corollary 13.3.3). We derive the conjugate function for completeness in Appendix B (Lemma ‣ B.2). In Section 5, we show how to leverage this technique for a variety of application input problems. See also Dünner et al. (2016) for a follow-up discussion of this technique in the non-distributed case.
For the case of objectives with strongly convex (or, equivalently, smooth ), e.g., elastic net regression or logistic regression, we obtain the following faster linear convergence rate.
Consider Algorithm 1 with , and let be the quality of the local solver as in Assumption 1. Let be -strongly convex , and let be -smooth. Then after iterations where
We provide proofs of both Theorem ‣ 4.2 and Theorem ‣ 4.3 in Appendix D.
4 Convergence Cases
Revisiting Table 1 from Section 2, we summarize our convergence guarantees for the three cases of input problems (I) in the following table. In particular, we see that for cases II and III, we obtain a sublinear convergence rate, whereas for case I we can obtain a faster linear rate, as provided in Theorem ‣ 4.3.
5 Recovering Earlier Work as a Special Case
As a special case, the proposed framework and rates directly apply to -regularized loss-minimization problems, including those presented in the earlier work of Jaggi et al. (2014) and Ma et al. (2015a).
These cases follow since is -Lipschitz if and only if has -bounded support (Rockafellar, 1997, Corollary 13.3.3), and is -strongly convex if and only if is -smooth (Hiriart-Urruty and Lemaréchal, 2001, Theorem 4.2.2).
Applications
For the examples in this subsection, we use to represent the number of training points and the number of features. Note that these definitions may change in the following subsections: this flexibility is useful so that we can present both the primal and dual variations of our framework (Algorithms 2 and 3) via a single abstract method (Algorithm 1).
Observing that the gradient of is , the primal-dual mapping is given by: , which is well known as the residual vector in least-squares regression.
An application we can consider for a strongly convex regularizer, in (A) or in (B), is elastic net regularization, , for fixed parameter . This can be obtained in (A) by setting
For the special case , we obtain the -norm, and for , we obtain the -norm. The conjugate of is given by: g^{*}_{i}(x):=\textstyle\frac{1}{2(1-\eta)}\big{(}\big{[}|x|-\eta\big{]}_{+}\big{)}^{2}, where is the positive part operator, for , and zero otherwise.
regularization is obtained in objective (A) by letting . However, an additional modification is necessary to obtain primal-dual convergence and certificates for this setting. In particular, we employ the modification introduced in Section 4, which will guarantee -bounded support. Formally, we replace by
For large enough , this problem yields the same solution as the original -regularized objective. Note that this only affects convergence theory, in that it allows us to present a strong primal-dual rate (Theorem ‣ 4.2 for =). With this modified regularizer, the optimization problem (A) with regularization parameter becomes
For large enough choice of the value , this problems yields the same solution as the original objective: . The modified is simply a constrained version of the absolute value to the interval . Therefore by setting to a large enough value that the values of will never reach it, will be continuous and at the same time make (23) equivalent to the original objective.
Formally, a simple way to obtain a large enough value of , so that all solutions are unaffected, is the following: If we start the algorithm at , for every point encountered during execution of a monotone optimizer, the objective values will never become worse than . Formally, under the assumption that is non-negative, we will have that (for each ):
We can therefore safely set the value of as . For the modified , the conjugate is given by:
We provide a proof of this in Appendix B (Lemma ‣ B.2).
The group lasso penalty can be mapped to objective (A), with:
For details, see, e.g., Dünner et al. (2016) or Boyd and Vandenberghe (2004, Example 3.26).
The conjugate function of the hinge loss is given by if , else . When using the norm for regularization in this problem: , a primal-dual mapping is given by: .
The absolute deviation loss, used, e.g., in quantile regression or least absolute deviation regression, can be realized in objective (B) by setting:
The conjugate function of the absolute deviation loss is given by , with .
4 Local Solvers
As discussed in Section 3, the subproblems solved on each machine in the CoCoA framework are appealing in that they are very similar in structure to the global problem (A), with the main difference being that they are defined on a smaller (local) subset of the data, and have a simpler dependence on the term . Therefore, solvers which have already proven their value in the single machine or multicore setting can be easily leveraged within the framework. We discuss some specific examples of local solvers below, and point the reader to Ma et al. (2017b) for an empirical exploration of these choices.
In the primal setting (Algorithm 2), the local subproblem (10) becomes a simple quadratic problem on the local data, with regularization applied only to local variables . For the -regularized examples discussed, existing fast -solvers for the single-machine case, such as glmnet variants (Friedman et al., 2010) or blitz (Johnson and Guestrin, 2015) can be directly applied to each local subproblem within Algorithm 1. The sparsity induced on the subproblem solutions of each machine naturally translates into the sparsity of the global solution, since the local variables will be concatenated.
In terms of the approximation quality parameter for the local problems (Assumption 1), we can apply existing recent convergence results from the single machine case. For example, for randomized coordinate descent (as part of glmnet), Lu and Xiao (2013, Theorem 1) gives a approximation quality for any separable regularizer, including and elastic net; see also Tappenden et al. (2015) and Shalev-Shwartz and Tewari (2011).
In the dual setting (Algorithm 3) for the discussed examples, the losses are applied only to local variables , and the regularizer is approximated via a quadratic term. Current state of the art for the problems of the form in (B) are variants of randomized coordinate ascent—Stochastic Dual Coordinate Ascent (SDCA) (Shalev-Shwartz and Zhang, 2013a). This algorithm and its variants are increasingly used in practice (Wright, 2015), and extensions such as accelerated and parallel versions can directly be applied (Shalev-Shwartz and Zhang, 2014; Fan et al., 2008) in our framework. For non-smooth losses such as SVMs, the analysis of Shalev-Shwartz and Zhang (2013a) provides a rate, and for smooth losses, a faster linear rate. There have also been recent efforts to derive a linear convergence rate for problems like the hinge-loss support vector machine that could be applied, e.g., by using error bound conditions (Necoara and Nedelcu, 2014; Wang and Lin, 2014), weak strong convexity conditions (Ma et al., 2015b; Necoara, 2015) or by considering Polyak-Łojasiewicz conditions (Karimi et al., 2016).
Experiments
In this section we demonstrate the empirical performance of CoCoA in the distributed setting. We first compare CoCoA to competing methods for two common machine learning applications: lasso regression (Section 6.1) and support vector machine (SVM) classification (Section 6.2). We then explore the performance of CoCoA in the primal versus the dual directly by solving an elastic net regression model with both variants (Section 6.3). Finally, we illustrate general properties of the CoCoA method empirically in Section 6.4.
We compare CoCoA to numerous state-of-the-art general-purpose methods for large-scale optimization, including:
Mb-SGD: Mini-batch stochastic gradient. For our experiments with lasso, we compare against Mb-SGD with an -prox.
GD: Full gradient descent. For lasso we use the proximal version, Prox-GD.
L-BFGS: Limited-memory quasi-Newton method. For lasso, we use OWL-QN (orthant-wise limited quasi-Newton).
ADMM: Alternating direction method of multipliers. We use conjugate gradient internally for the lasso experiments, and SDCA for SVM experiments.
Mb-CD: Mini-batch parallel coordinate descent. For SVM experiments, we implement Mb-SDCA (mini-batch stochastic dual coordinate ascent).
The first three methods are optimized and implemented in Apache Spark’s MLlib (v1.5.0) (Meng et al., 2016). We test the performance of each method in large-scale experiments fitting lasso, elastic net regression, and SVM models to the datasets shown in Table 5. In comparing to other methods, we plot the distance to the optimal primal solution. This optimal value is calculated by running all methods for a large number of iterations (until progress has stalled), and then selecting the smallest primal value amongst the results. All code is written in Apache Spark and experiments are run on public-cloud Amazon EC2 m3.xlarge machines with one core per machine. Our code is publicly available at gingsmith.github.io/cocoa/.
We carefully tune each competing method in our experiments for best performance. ADMM requires the most tuning, both in selecting the penalty parameter and in solving the subproblems. Solving the subproblems to completion for ADMM is prohibitively slow, and we thus use an iterative method internally and improve performance by allowing early stopping. We also use a varying penalty parameter — practices described in Boyd et al. (2010, Sections 4.3, 8.2.3, 3.4.1). For Mb-SGD, we tune the step size and mini-batch size parameters. For Mb-CD and Mb-SDCA, we scale the updates at each round by for mini-batch size and , and tune both parameters and . Further implementation details for all methods are given in Section 6.5. For simplicity of presentation and comparison, in all of the following experiments, we restrict CoCoA to only use simple coordinate descent as the local solver. We note that even stronger empirical results for CoCoA could be obtained by plugging in state-of-the-art local solvers for each application at hand.
1 CoCoA in the Primal: An Application to Lasso Regression
We first demonstrate the performance of CoCoA in the primal (Algorithm 2) by applying CoCoA to a lasso regression model (8) fit to the datasets in Table 5. We use stochastic coordinate descent as a local solver for CoCoA, and select the number of local iterations (a proxy for subproblem approximation quality, ) from several options with best performance.
We compare CoCoA to the general methods listed above, including Mb-SGD with an -prox, Prox-GD, OWL-QN, ADMM, and Mb-CD. We provide a comparison with Shotgun (Bradley et al., 2011), a popular method for solving -regularized problems in the multicore environment, as an extreme case to highlight the detrimental effects of frequent communication in the distributed environment. For Mb-CD, Shotgun, and CoCoA in the primal, datasets are distributed by feature, whereas for Mb-SGD, Prox-GD, OWL-QN and ADMM they are distributed by training point.
In analyzing the performance of each algorithm (Figure 1), we measure the improvement to the primal objective given in (A) in terms of wall-clock time in seconds. We see that both Mb-SGD and Mb-CD are slow to converge, and come with the additional burden of having to tune extra parameters (though Mb-CD makes clear improvements over Mb-SGD). As expected, naively distributing Shotgun (single coordinate updates per machine) does not perform well, as it is tailored to shared-memory systems and requires communicating too frequently. OWL-QN performs the best of all compared methods, but is still much slower to converge than CoCoA, and converges, e.g., 50 more slowly for the webspam dataset. The optimal performance of CoCoA is particularly evident in datasets with large numbers of features (e.g., url, kddb, webspam), which are exactly the datasets where regularization would most typically be applied.
Results are shown for regularization parameters such that the resulting weight vector is sparse. However, our results are robust to varying values of as well as to various problem settings, as we illustrate in Figure 2.
We additionally motivate the use of CoCoA in the primal by showing how it improves upon CoCoA in the dual (Yang, 2013; Jaggi et al., 2014; Ma et al., 2015a, 2017b) for non-strongly convex regularizers. First, CoCoA in the dual cannot be included in the set of experiments in Figure 1 because it cannot be directly applied to the lasso objective (recall that Algorithm 3 only allows for strongly convex regularizers).
To get around this requirement, previous work has suggested implementing the smoothing technique used in, e.g., Shalev-Shwartz and Zhang (2014); Zhang and Lin (2015) — adding a small amount of strong convexity to the objective for lasso regression. In Figure 3 we demonstrate the issues with this approach, comparing CoCoA in the primal on a pure -regularized regression problem to CoCoA in the dual for decreasing levels of . The smaller we set , the less smooth the problem becomes. As decreases, the final sparsity of running CoCoA in the dual starts to match that of running pure (Table 6), but the performance also degrades (Figure 3). We note that by using CoCoA in the primal with the modification presented in Section 4, we can deliver strong rates without having to make any compromises in terms of the training speed or accuracy.
2 CoCoA in the Dual: An Application to SVM Classification
Next we present results on CoCoA in the dual against competing methods, for a hinge loss support vector machine model (9) on the datasets in Table 5. We use stochastic dual coordinate ascent (SDCA) as a local solver for CoCoA in this setting, again selecting the number of local iterations from several options with best performance. We compare CoCoA to the general methods listed above, including Mb-SGD, GD, L-BFGS, ADMM, and Mb-SDCA. All datasets are distributed by training point for these methods.
In comparing methods in this setting (Figure 4), we measure the improvement to the primal objective in terms of wall-clock time in seconds. We see again that Mb-SGD and Mb-CD are slow to converge, and come with the additional burden of having to tune extra parameters. ADMM performs the best of the methods other than CoCoA, followed by L-BFGS. However, both are still much slower to converge than CoCoA in the dual. ADMM in particular is affected by the fact that many internal iterations of SDCA are necessary in order to guarantee convergence. In contrast, CoCoA can incorporate arbitrary amounts of work locally and still converge. We note that although CoCoA, ADMM, and Mb-SDCA run in the dual, Figure 4 tracks progress towards the primal objective, .
3 Primal vs. Dual: An Application to Elastic Net Regression
To compare primal vs. dual optimization for CoCoA, we explore both variants by fitting an elastic net regression model (7) to two datasets. We use coordinate descent (with closed-form updates) as the local solver in both variants. From the results in Figure 5, we see that CoCoA in the dual tends to perform better on datasets with a large number of training points (relative to the number of features), and that the performance deteriorates as the strong convexity in the problem disappears. In contrast, CoCoA in the primal performs well on datasets with a large number of features relative to training points, and is robust to changes in strong convexity. These changes in performance are to be expected, as we have discussed that CoCoA in the primal is more suited for non-strongly convex regularizers (Section 6.1), and that the feature size dominates communication for CoCoA in the dual, as compared to the training point size for CoCoA in the primal (Section 3.4).
4 General Properties: Effect of Communication
Finally, we note that in contrast to the compared methods from Sections 6.1 and 6.2, CoCoA comes with the benefit of having only a single parameter to tune: the subproblem approximation quality, , which we control in our experiments via the number of local subproblem iterations, , for the example of local coordinate descent. We further explore the effect of this parameter in Figure 6, and provide a general guideline for choosing it in practice (see Remark ‣ 3.1). In particular, we see that while increasing always results in better performance in terms of the number of communication rounds, smaller or larger values of may result in better performance in terms of wall-clock time, depending on the cost of communication and computation. The flexibility to fine-tune is one of the reasons for CoCoA’s significant performance gains.
5 Experiment Details
In this subsection we provide thorough details on the experimental setup and methods used in our comparison. All experiments are run on Amazon EC2 clusters of m3.xlarge machines, with one core per machine. The code for each method is written in Apache Spark, v1.5.0. Our code is open source and publicly available at gingsmith.github.io/cocoa/.
Mini-batch SGD is a standard and widely used method for parallel and distributed optimization. We use the optimized code provided in Spark’s machine learning library, MLlib, v1.5.0 (Meng et al., 2016). We tune both the size of the mini-batch and SGD step size using grid search. For lasso, we use the proximal version of the method. Full gradient descent can be seen as a specific setting of mini-batch SGD, where the mini-batch size is equal to the total number of training points. We thus also use the implementation in MLlib for full GD, and tune the step size parameter using grid search.
Mini-batch CD (for lasso) and SDCA (for SVM) aim to improve mini-batch SGD by employing coordinate descent, which has theoretical and practical justifications (Shalev-Shwartz and Tewari, 2011; Takáč et al., 2013; Fercoq and Richtárik, 2015). We implement mini-batch CD and SDCA in Spark and scale the updates made at each round by for mini-batch size and , tuning both parameters and via grid search. For the case of lasso regression, we implement Shotgun (Bradley et al., 2011), which is a popular method for parallel optimization. Shotgun can be seen an extreme case of mini-batch CD where the mini-batch is set to , i.e., there is a single update made by each machine per round. We see in the experiments that communicating this frequently becomes prohibitively slow in the distributed environment.
OWN-QN (Yu et al., 2010) is a quasi-Newton method optimized in Spark’s spark.ml package (Meng et al., 2016). Outer iterations of OWL-QN make significant progress towards convergence, but the iterations themselves can be slow as they require processing the entire dataset. CoCoA, the mini-batch methods, and ADMM with early stopping all improve on this by allowing the flexibility to process only a subset of the dataset at each iteration. CoCoA and ADMM have even greater flexibility by allowing internal methods to process the dataset more than once. CoCoA makes this approximation quality explicit, both in theoretical convergence rates and via guidelines for setting the parameter.
We implement CoCoA with coordinate descent as the local solver. We note that since the framework and theory allow any internal solver to be used, CoCoA could benefit even beyond the results shown, e.g., by using existing fast -solvers for the single-machine case, such as glmnet variants (Friedman et al., 2010) or blitz (Johnson and Guestrin, 2015) or SVM solvers like liblinear (Fan et al., 2008). The only parameter influencing the overall performance of CoCoA is the level of approximation quality, which we parameterize in the experiments through , the number of local iterations of the iterative method run locally. Our theory relates local approximation quality to global convergence (Section 4), and we provide a guideline for how to choose this value in practice that links the parameter to the systems environment at hand (Remark ‣ 3.1).
Related Work
There exist myriad optimization methods for the distributed setting; the following section is not meant to be wholly comprehensive, but to provide an overview of the most prevalent and related approaches. We additionally note that many new methods have been proposed since the time of submission of this manuscript in October 2016, including several extensions of the presented CoCoA framework—e.g., for federated learning (Smith et al., 2017), computing over heterogeneous systems (Dünner et al., 2017), second-order algorithm extensions (Gargiani, 2017; Lee and Chang, 2017; Dünner et al., 2018; Lee et al., 2018), and accelerated methods (Ma et al., 2017a; Zheng et al., 2017). We defer the readers to these follow-up works for the most current literature review.
For strongly convex regularizers, state-of-the-art for empirical loss minimization is randomized coordinate ascent on the dual (SDCA) (Shalev-Shwartz and Zhang, 2013a) and accelerated variants (e.g., Shalev-Shwartz and Zhang, 2014). In contrast to primal stochastic gradient descent (SGD) methods, the SDCA family is often preferred as it is free of learning rate parameters and has faster (geometric) convergence guarantees. Interestingly, a similar trend in coordinate solvers exists in recent lasso literature, but with the roles of primal and dual reversed. For those problems, primal-based coordinate descent methods are state-of-the-art, as in glmnet (Friedman et al., 2010) and extensions (Yuan et al., 2012); see, e.g., the overview in Yuan et al. (2010). However, primal-dual rates for unmodified coordinate methods have to our knowledge only been obtained for strongly convex regularizers to date (Shalev-Shwartz and Zhang, 2014; Zhang and Lin, 2015).
Coordinate descent on -regularized problems (i.e., (A) with ) can be interpreted as the iterative minimization of a quadratic approximation of the smooth part of the objective, followed by a shrinkage step. In the single-coordinate update case, this is at the core of glmnet (Friedman et al., 2010; Yuan et al., 2010), and widely used in, e.g., solvers based on the primal formulation of -regularized objectives (Shalev-Shwartz and Tewari, 2011; Yuan et al., 2012; Bian et al., 2013; Fercoq and Richtárik, 2015; Tappenden et al., 2015). When changing more than one coordinate at a time, again employing a quadratic upper bound on the smooth part, this results in a two-loop method as in glmnet for the special case of logistic regression. In the distributed setting, when the set of active coordinates coincides with the ones on the local machine, these single-machine approaches closely resemble the distributed framework proposed here.
For the general regularized loss minimization problems of interest, methods based on stochastic subgradient descent (SGD) are well-established. Several variants of SGD have been proposed for parallel computing, many of which build on the idea of asynchronous communication (Niu et al., 2011; Duchi et al., 2013). Despite their simplicity and competitive performance on shared-memory systems, the downside of this approach in the distributed environment is that the amount of required communication is equal to the amount of data read locally, since one data point is accessed per machine per round (e.g., mini-batch SGD with a batch size of one per worker). These variants are in practice not competitive with the more communication-efficient methods considered in this work, which allow more local updates per communication round.
For the specific case of -regularized objectives, parallel coordinate descent (with and without using mini-batches) was proposed in Bradley et al. (2011) (Shotgun) and generalized in Bian et al. (2013) , and is among the best performing solvers in the parallel setting. Our framework reduces to Shotgun as a special case when the internal solver is a single-coordinate update on the subproblem (10), , and for a suitable . However, Shotgun is not covered by our convergence theory, since it uses a potentially unsafe upper bound of instead of , which isn’t guaranteed to satisfy our condition for convergence (11). We compare empirically with Shotgun in Section 6 to highlight the detrimental effects of running this high-communication method in the distributed environment.
At the other extreme, there are distributed methods that use only a single round of communication, such as Mann et al. (2009); Zinkevich et al. (2010); Zhang et al. (2013); McWilliams et al. (2014); and Heinze et al. (2016). These methods require additional assumptions on the partitioning of the data, which are usually not satisfied in practice if the data are distributed “as is”, i.e., if we do not have the opportunity to distribute the data in a specific way beforehand. Furthermore, some cannot guarantee convergence rates beyond what could be achieved if we ignored data residing on all but a single computer, as shown in Shamir et al. (2014). Additional relevant lower bounds on the minimum number of communication rounds necessary for a given approximation quality are presented in Balcan et al. (2012) and Arjevani and Shamir (2015).
Mini-batch methods (which use updates from several training points or features per round) are more flexible and lie within the two extremes of parallel and one-shot communication schemes. However, mini-batch versions of both SGD and coordinate descent (CD) (e.g., Dekel et al., 2012; Takáč et al., 2013; Shalev-Shwartz and Zhang, 2013b; Shamir and Srebro, 2014; Qu et al., 2015; Richtárik and Takáč, 2016; Défossez and Bach, 2017) suffer from their convergence rate degrading towards the rate of batch gradient descent as the size of the mini-batch is increased. This follows because mini-batch updates are made based on the outdated previous parameter vector , in contrast to methods that allow immediate local updates like CoCoA.
Another disadvantage of mini-batch methods is that the aggregation parameter is more difficult to tune, as it can lie anywhere in the order of mini-batch size. The optimal choice is often either unknown or too challenging to compute in practice. In the CoCoA framework there is no need to tune parameters, as the aggregation parameter and subproblem parameters can be set directly using the safe bound discussed in Section 3 (Definition ‣ 3.1).
ADMM (Boyd et al., 2010), gradient descent, and quasi-Newton methods such as L-BFGS and are also often used in distributed settings because of their relatively low communication requirements. However, they require at least a full (distributed) batch gradient computation at each round, and therefore do not allow the gradual trade-off between communication and computation provided by CoCoA. In Section 6, we include experimental comparisons with ADMM, gradient descent, and L-BFGS variants, including orthant-wise limited memory quasi-Newton (OWL-QN) for the setting (Andrew and Gao, 2007).
Finally, we note that while the convergence rates provided for CoCoA mirror the convergence class of classical batch gradient methods in terms of the number of outer rounds, existing batch gradient methods come with a weaker theory, as they do not allow general inexactness for the local subproblem (10). In contrast, our convergence rates incorporate this approximation directly, and, moreover, hold for arbitrary local solvers of much cheaper cost than batch methods (where in each round, every machine has to process exactly a full pass through the local data). This makes CoCoA more flexible in the distributed setting, as it can adapt to varied communication costs on real systems. We have seen in Section 6 that this flexibility results in significant performance gains over the competing methods.
By making use of the primal-dual structure in the line of work of Yu et al. (2012); Pechyony et al. (2011); Yang (2013); Yang et al. (2013) and Lee and Roth (2015), the CoCoA-v1 and CoCoA frameworks (which are special cases of the presented framework, CoCoA) are the first to allow the use of any local solver—of weak local approximation quality—in each round in the distributed setting. The practical variant of the DisDCA (Yang, 2013), called DisDCA-p, allows for additive updates in a similar manner to CoCoA, but is restricted to coordinate decent (CD) being the local solver, and was initially proposed without convergence guarantees. DisDCA-p, CoCoA-v1, and CoCoA are all limited to strongly convex regularizers, and therefore are not as general as the CoCoA framework discussed in this work.
In the -regularized setting, an approach related to our framework includes distributed variants of glmnet as in Mahajan et al. (2017). Inspired by glmnet and Yuan et al. (2012), the works of Bian et al. (2013) and Mahajan et al. (2017) introduced the idea of a block-diagonal Hessian upper approximation in the distributed context. The later work of Trofimov and Genkin (2014, 2016) specialized this approach to sparse logistic regression.
If hypothetically each of our quadratic subproblems as defined in (10) were to be minimized exactly, the resulting steps could be interpreted as block-wise Newton-type steps on each coordinate block , where the Newton-subproblem is modified to also contain the -regularizer (Mahajan et al., 2017; Yuan et al., 2012; Qu et al., 2016). While Mahajan et al. (2017) allows a fixed accuracy for these subproblems, but not arbitrary approximation quality as in our framework, the works of Trofimov and Genkin (2016); Yuan et al. (2012); and Yen et al. (2015) assume that the quadratic subproblems are solved exactly. Therefore, these methods are not able to freely trade off communication and computation. Also, they do not allow the re-use of arbitrary local solvers. On the theoretical side, the convergence rate results provided by Mahajan et al. (2017); Trofimov and Genkin (2016); and Yuan et al. (2012) are not explicit convergence rates but only asymptotic, as the quadratic upper bounds are not explicitly controlled for safety as with our .
Discussion
To enable large-scale machine learning and signal processing, we have developed, analyzed, and evaluated a general-purpose framework for communication-efficient primal-dual optimization in the distributed environment. Our framework, CoCoA, takes a unique approach by using duality to derive subproblems for each machine to solve in parallel. These subproblems closely match the global problem of interest, which allows for state-of-the-art single-machine solvers to easily be re-used in the distributed setting. Further, by allowing the local solvers to find solutions of arbitrary approximation quality to the subproblems on each machine, our framework permits a highly flexible communication scheme. In particular, as the local solvers make updates directly to their local parameters, the need to communicate reduces and can be adapted to the system at hand, which helps to manage the communication bottleneck in the distributed setting.
We analyzed the impact of the local solver approximation quality and derived global primal-dual convergence rates for our framework that are agnostic to the specifics of the local solvers. We have taken particular care in extending our framework to the case of non-strongly convex regularizers, where we introduced a bounded-support modification technique to provide robust convergence guarantees. Finally, we demonstrated the efficiency of our framework in an extensive experimental comparison with state-of-the-art distributed solvers. Our framework achieves up to a 50 speedup over other widely-used methods on real-world distributed datasets.
We thank Michael P. Friedlander, Matilde Gargiani, Sai Praneeth Karimireddy, Jakub Konečný, Ching-pei Lee, and Peter Richtárik for their help and for fruitful discussions. We are additionally grateful to the reviewers for their valuable comments. We wish to acknowledge support from the U.S. National Science Foundation, under award number NSF:CCF:1618717, NSF:CMMI:1663256 and NSF:CCF:1740796; the Swiss National Science Foundation, under grant number 175796; and the Mathematical Data Science program of the Office of Naval Research, under grant number N00014-15-1-2670.
A Convex Conjugates
Below we list several useful properties of conjugates (see, e.g., Boyd and Vandenberghe, 2004, Section 3.3.2):
Double conjugate: if is closed and convex.
Value Scaling: (for )
Argument Scaling: (for )
Conjugate of a separable sum:
Given a proper convex function , it holds that is -Lipschitz if and only if has -bounded support.
Given a closed convex function , it holds that is strongly convex w.r.t. the norm if and only if is -smooth w.r.t. the dual norm .
B Proofs of Primal-Dual Relationships
In the following subsections we provide derivations of the primal-dual relationship of the general objectives (A) and (B), and then show how to derive the conjugate of the modified -norm, as an example of the bounded-support modification introduced in Section 4.
The relation of the original formulation (A) to its dual formulation (B) is standard in convex analysis. Using the linear map as in our case, the relationship is an instance of Fenchel-Rockafellar Duality, see e.g. Borwein and Zhu (2005, Theorem 4.4.2) or Bauschke and Combettes (2011, Proposition 15.18). For completeness, we illustrate this correspondence with a self-contained derivation of the duality.
The dual problem of (A) follows by taking the infimum with respect to both and :
We change signs and turn the maximization of the dual problem (29) into a minimization, thereby arriving at the dual formulation as claimed:
B.2 Continuous Conjugate Modification for Indicator Functions
The convex conjugate of the bounded support modification of the -norm, as defined in (18), is:
We start by applying the definition of convex conjugate:
We begin by looking at the case in which ; in this case it’s easy to see that when , we have:
as . The case holds analogously. We’ll now look at the case ; in this case it is clear we must have . It also must hold that , since
for every . Therefore the maximization becomes
which has maximum at . The remaining case follows in similar fashion.
Lipschitz continuity of follows directly, or alternatively also from the general result that is -Lipschitz if and only if has -bounded support (Rockafellar, 1997, Corollary 13.3.3) or (Dünner et al., 2016, Lemma 5). ∎
C Comparison to ADMM
Here we compare consensus ADMM (Mota et al., 2013) applied to the problem (B) to the CoCoA framework, as discussed in Section 3.5. For consensus ADMM, the objective can be decomposed using the following re-parameterization:
To solve this problem, we construct the augmented Lagrangian:
which yields the following decomposable updates:
These updates can be further simplified by using the scaled form of and combining terms for using the averages and :
To compare this to the CoCoA subproblems (10), we will derive the dual form of the update to . Suppressing the iteration counter for simplicity, the minimization is of the form:
Solving the minimization yields: . Plugging this back in, we have:
We therefore see that the update to has a similar form to the CoCoA subproblem (10), where .
We can also compare CoCoA to ADMM as applied to the (A) problem. For consensus ADMM, the objective can be decomposed using the following re-parametrization, which introduces local copies of the global variable , and a set of consensus constraints to achieve the equality between them:
To solve this problem, we construct the augmented Lagrangian with a penalty parameter :
which yields to the following decomposable updates:
The first minimization is solved locally in a distributed manner by the partitions. By setting and applying the change of variables , we can obtain the CoCoA local subproblems (10):
Note that the update to in its current form is not separable and must be solved via some distributed optimization procedure. The comparison to the formulation is more natural in this sense, as it captures a setting and formulation in which distributed ADMM would more commonly be applied. However, the above formulation is closely related to the sharing variant of ADMM (Boyd et al., 2010, Section 7.3), where data for canonical regularized loss minimization problems is assumed to be distributed via features (see, e.g., Boyd et al., 2010, Section 8.3).
D Convergence Proofs
In this section we provide proofs of our main convergence results. The arguments follow the reasoning in Ma et al. (2015a, 2017b), but where we have generalized them to be applicable directly to (A). We provide full details of Lemma ‣ 4.1 as a proof of concept, but omit details in later proofs that can be derived using the arguments in Ma et al. (2015a) or earlier work of Shalev-Shwartz and Zhang (2013a), and instead outline the proof strategy and highlight sections where the theory deviates.
Our first lemma in the overall proof of convergence helps to relate progress on the local subproblems to the global objective .
In this proof we follow the line of reasoning in Ma et al. (2015a, Lemma 4) with a more general smoothness assumption on . An outer iteration of CoCoA performs the following update:
We bound and separately. First we bound A using -smoothness of :
Next we use Jensen’s inequality to bound B:
Plugging and back into (31) yields:
where the last equality is by the definition of the subproblem objective as in (10). ∎
D.2 Proof of Main Convergence Result (Theorem ‣ 4.2)
Before proving the main convergence results, we introduce several useful quantities, and establish the following lemma, which characterizes the effect of iterations of Algorithm 1 on the duality gap for any chosen local solver of approximation quality .
Let be strongly convex Note that the case of weakly convex is explicitly allowed here as well, as the Lemma holds for the case . with convexity parameter with respect to the norm , . Then at each iteration of Algorithm 1 under Assumption 1, and any , it holds that
This proof is motivated by Shalev-Shwartz and Zhang (2013a, Lemma 19) and follows Ma et al. (2015a, Lemma 5), with a difference being the extension to our generalized subproblems along with the mappings with .
For simplicity, we write instead of , instead of , instead of and instead of . We can estimate the expected change of the objective as follows. Starting from the definition of the update from Algorithm 1, we apply Lemma ‣ 4.1, which relates the local approximation to the global objective , and then bound this using the notion of quality of the local solver (), as in Assumption 1. This gives us:
We next upper bound the term, denoting . We first plug in the definition of the objective in (A) and the local subproblems (10), and then substitute for and apply the -strong convexity of the terms. This gives us:
From the definition of the optimization problems (A) and (B), and definition of convex conjugates, we can write the duality gap as:
The convex conjugate maximal property from (34) implies that
The claimed improvement bound (32) then follows by plugging (39) into (35). ∎
The following Lemma provides a uniform bound on .
If are -Lipschitz continuous for all , then
(Ma et al., 2015a, Lemma 6). For general convex functions, the strong convexity parameter is , and hence the definition (33) of the complexity constant becomes
(Ma et al., 2015a, Remark 7) If the data points are normalized such that , , then . Furthermore, if we assume that the data partition is balanced, i.e., that for all , then . This can be used to bound the constants , above, as
Consider Algorithm 1, using a local solver of quality (See Assumption 1). Let be -Lipschitz continuous, and be the desired duality gap (and hence an upper-bound on suboptimality ). Then after iterations, where
we have that the expected duality gap satisfies
We begin by estimating the expected change of feasibility for . We can bound this above by using Lemma ‣ D.2 and the fact that the is always a lower bound for , and then applying (40) to find:
Choosing and leads to
Clearly, (46) implies that (47) holds for . Assuming that it holds for any , we show that it must also hold for . Indeed, using
by applying the bounds (44) and (47), plugging in the definition of (48), and simplifying. We upper bound the term using the fact that geometric mean is less or equal to arithmetic mean:
If is defined as (43), we apply the results of Lemma ‣ D.2 and Lemma ‣ D.2 to obtain
If such that we have
To have right hand side of (52) smaller then it is sufficient to choose and such that
Hence if and then (53) and (54) are satisfied. ∎
The following main theorem simplifies the results of Theorem ‣ D.2 and is a generalization of Ma et al. (2015a, Corollary 9) for general functions:
Consider Algorithm 1 with , using a local solver of quality (see Assumption 1). Let be -Lipschitz continuous, and assume that the columns of satisfy , and is of the form , as is common in ERM-type problems. Let be the desired duality gap (and hence an upper-bound on primal sub-optimality). Then after iterations, where
we have that the expected duality gap satisfies
where is the averaged iterate returned by Algorithm 1.
Our second main theorem follows reasoning in Shalev-Shwartz and Zhang (2013a) and is a generalization of Ma et al. (2015a, Corollary 11). We first introduce a lemma to simplify the proof.
Since , and given our initial assumption on , the duality gap reduces to:
Assume that are -strongly convex . We define . Then after iterations of Algorithm 1, with
Given that is -strongly convex, we can apply (33) and the definition of to find:
where . If we plug the following value of
into (57) we obtain that . Putting the same into (32) will give us
Therefore if we denote we have recursively that
The right hand side will be smaller than some if
Moreover, to bound the duality gap, we have
Thus, . Hence if then . Therefore after
iterations we have obtained a duality gap less than . ∎
Consider Algorithm 1 with , using a local solver of quality (see Assumption 1). Let be -strongly convex, , and assume that the columns of satisfy and and is of the form , as is common in ERM-type problems. Then we have that iterations are sufficient for suboptimality , with