ARock: an Algorithmic Framework for Asynchronous Parallel Coordinate Updates
Zhimin Peng, Yangyang Xu, Ming Yan, Wotao Yin
Introduction
Technological advances in data gathering and storage have led to a rapid proliferation of big data in diverse areas such as climate studies, cosmology, medicine, the Internet, and engineering house2014big . The data involved in many of these modern applications are large and grow quickly. Therefore, parallel computational approaches are needed. This paper introduces a new approach to asynchronous parallel computing with convergence guarantees.
In a synchronous(sync) parallel iterative algorithm, the agents must wait for the slowest agent to finish an iteration before they can all proceed to the next one (Figure 1(a)). Hence, the slowest agent may cripple the system. In contract, the agents in an asynchronous(async) parallel iterative algorithm run continuously with little idling (Figure 1(b)). However, the iterations are disordered, and an agent may carry out an iteration without the newest information from other agents.
Asynchrony has other advantages bertsekas1991some : the system is more tolerant to computing faults and communication glitches; it is also easy to incorporate new agents.
On the other hand, it is more difficult to analyze asynchronous algorithms and ensure their convergence. It becomes impossible to find a sequence of iterates that one completely determines the next. Nonetheless, we let any update be a new iteration and propose an async-parallel algorithm (ARock) for the generic fixed-point iteration. It converges if the fixed-point operator is nonexpansive (Def. 1) and has a fixed point.
Let be Hilbert spaces and be their Cartesian product. For a nonexpansive operator , our problem is to
Finding a fixed point to is equivalent to finding a zero of denoted by such that . Hereafter, we will use both and for convenience.
Problem (1) is widely applicable in linear and nonlinear equations, statistical regression, machine learning, convex optimization, and optimal control. A generic framework for problem (1) is the Krasnosel’skiĭ–Mann (KM) iteration krasnosel1955two :
where is the step size. If — the set of fixed points of (zeros of ) — is nonempty, then the sequence converges weakly to a point in and converges strongly to 0. The KM iteration generalizes algorithms in convex optimization, linear algebra, differential equations, and monotone inclusions. Its special cases include the following iterations: alternating projection, gradient descent, projected gradient descent, proximal-point algorithm, Forward-Backward Splitting (FBS) passty1979ergodic , Douglas-Rachford Splitting (DRS) lions1979splitting , a three-operator splitting davis2015three , and the Alternating Direction Method of Multipliers (ADMM) lions1979splitting ; glowinski1975approximation .
In ARock, a set of agents, , solve problem (1) by updating the coordinates , , in a random and asynchronous fashion. Algorithm 1 describes the framework. Its special forms for several applications are given in Section 2 below.
Whenever an agent updates a coordinate, the global iteration counter increases by one. The th update is applied to , where is an independent random variable. Each coordinate update has the form:
where is a scalar whose range will be set later, and is used to normalize nonuniform selection probabilities. In the uniform case, namely, for all , we have , which simplifies the update (3) to
Here, the point is what an agent reads from global memory to its local cache and to which is applied, and denotes the state of in global memory just before the update (3) is applied. In a sync-parallel algorithm, we have , but in ARock, due to possible updates to by other agents, can be different from . This is a key difference between sync-parallel and async-parallel algorithms. In Subsection 1.2 below, we will establish the relationship between and as
The update (3) is only computationally worthy if is much cheaper to compute than . Otherwise, it is more preferable to apply the full KM update (2). In Section 2, we will present several applications that have the favorable structures for ARock. The recent work peng2016coordinate studies coordinate friendly structures more thoroughly.
The convergence of ARock (Algorithm 1) is stated in Theorems 3.2 and 3.3. Here we include a shortened version, leaving detailed bounds to the full theorems:
Let be a nonexpansive operator that has a fixed point. Let be the sequence generated by Algorithm 1 with properly bounded step sizes . Then, with probability one, converges weakly to a fixed point of . This convergence becomes strong if has a finite dimension.
In addition, if is demicompact (see Definition 2 below), then with probability one, converges strongly to a fixed point of .
In the theorem, the weak convergence result only requires to be nonexpansive and has a fixed point. In addition, the computation requires: (a) bounded step sizes; (b) random coordinate selection; and (c) a finite maximal delay . Assumption (a) is standard, and we will see the bound can be . Assumption (b) is essential to both the analysis and the numerical performance of our algorithms. Assumption (c) is not essential; an infinite delay with a light tail is allowed (but we leave it to future work). The strong convergence result applies to all the examples in Section 2, and the linear convergence result applies to Examples 2.2 and 2.4 when the corresponding operator is quasi-strongly monotone. Step sizes are discussed in Remarks 2 and 4.
ARock employs random coordinate selection. This subsection discusses its advantages and disadvantages.
Its main disadvantage is that an agent cannot caching the data associated with a coordinate. The variable and its related data must be either stored in global memory or passed through communication. A secondary disadvantage is that pseudo-random number generation takes time, which becomes relatively significant if each coordinate update is cheap. (The network optimization examples in Subsections 2.3 and 2.6.2 are exceptions, where data are naturally stored in a distributed fashion and random coordinate assignments are the results of Poisson processes.)
There are several advantages of random coordinate selection. It realizes the user-specified update frequency for every component , , even when different agents have different computing powers and different coordinate updates cost different amounts of computation. Therefore, random assignment ensures load balance. The algorithm is also fault tolerant in the sense that if one or more agents fail, it will still converge to a fixed-point of . In addition, it has been observed numerically on certain problems chang2008coordinate that random coordinate selection accelerates convergence.
2 Uncoordinated memory access
In ARock, since multiple agents simultaneously read and update in global memory, — the result of that is read from global memory by an agent to its local cache for computation — may not equal for any , that is, may never be consistent with a state of in global memory. This is known as inconsistent read. In contrast, consistent read means that for some , i.e., is consistent with a state of that existed in global memory.
Even with inconsistent read, each component is consistent under the atomic coordinate update assumption, which will be defined below. Therefore, we can express what has been read in terms of the changes of individual coordinates. In the above example, the first change is , which is added to just before time by agent 2, and the second change is , added to just before time by agent 3. The inconsistent read by agent 1, which gives the result , equals .
We have demonstrated that can be inconsistent, but each of its coordinates is consistent, that is, for each , is an ever-existed state of among . Suppose that , where . Therefore, can be related to through the interim changes applied to . Let be the index set of these interim changes. If , then ; otherwise, . In addition, we have . Since the global counter is increased after each coordinate update, updates to and , , must occur at different ’s and thus . Therefore, by letting and noticing for where , we have , which is equivalent to (5). Here, we have made two assumptions:
atomic coordinate update: a coordinate is not further broken to smaller components during an update; they are all updated at once.
bounded maximal delay : during any update cycle of an agent, in global memory is updated at most times by other agents.
When each coordinate is a single scalar, updating the scalar is a single atomic instruction on most modern hardware, so the first assumption naturally holds, and our algorithm is lock-free. The case where a coordinate is a block that includes multiple scalars is discussed in the next subsection.
In the “block coordinate” case (updating a block of several coordinates each time), the atomic coordinate update assumption can be met by either employing a per-coordinate memory lock or taking the following dual-memory approach: Store two copies of each coordinate in global memory, denoting them as and ; let a bit point to the active copy; an agent will only read from the active copy ; before an agent updates the components of , it obtains a memory lock to the inactive copy to prevent other agents from simultaneously updating it; then after it finishes updating , flip the bit so that other agents will begin reading from the updated copy. This approach never blocks any read of , yet it eliminates inconsistency.
3 Straightforward generalization
Our async-parallel coordinate update scheme (3) can be generalized to (overlapping) block coordinate updates after a change to the step size. Specifically, the scheme (3) can be generalized to
where is randomly drawn from a set of operators (), , following the probability , (, and ). The operators must satisfy and for some .
Let , which has ; then (6) reduces to (3). If is endowed with a metric such that (e.g., the metric in the Condat-Vũ primal-dual splitting condat2013primal ; vu2013splitting ), then we have
In general, multiple coordinates can be updated in (6). Consider linear , , where for each . Then, for , we have
4 Special cases
If there is only one agent (), ARock (Algorithm 1) reduces to randomized coordinate update, which includes the special case of randomized coordinate descent nesterov2012rcd for convex optimization. Sync-parallel coordinate update is another special case of ARock corresponding to . In both cases, there is no delay, i.e., and . In addition, the step size can be more relaxed. In particular, if , , then we can let , , for any , or when is -averaged (see Definition 2 for the definition of an -averaged operator).
5 Related work
Chazan and Miranker chazan1969chaotic proposed the first async-parallel method in 1969. The method was designed for solving linear systems. Later, async-parallel methods have been successful applied in many fields, e.g., linear systems avron2014revisiting ; bethune2014performance ; FSS1997asyn-addSch ; rosenfeld1969case , nonlinear problems BMR1997asyn-multisplit ; baudet1978asynchronous , differential equations aharoni2000parallel ; AAI1998implicit ; Chau20081126 ; donzis2014asynchronous , consensus problems LMS1986asynchronous ; leifang_information_2005 , and optimization hsieh2015passcode ; liu2014asynchronous ; liu2013asynchronous ; tai2002convergence ; zhang2014asynchronous . We review the theory for async-parallel fixed-point iteration and its applications.
The above works assign coordinates in a deterministic manner. Different from them, ARock is stochastic, works for nonexpansive operators, and is more applicable.
Linear, nonlinear, and differential equations. The first async-parallel method for solving linear equations was introduced by Chazan and Miranker in chazan1969chaotic . They proved that on solving linear systems, P-contraction was necessary and sufficient for convergence. The performance of the algorithm was studied by Iain et al. bethune2014performance ; rosenfeld1969case on different High Performance Computing (HPC) architectures. Recently, Avron et al. avron2014revisiting revisited the async-parallel coordinate update and showed its linear convergence for solving positive-definite linear systems. Tarazi and Nabih el1982some extended the poineering work chazan1969chaotic to solving nonlinear equations, and the async-parallel methods have also been applied for solving differential equations, e.g., in aharoni2000parallel ; AAI1998implicit ; Chau20081126 ; donzis2014asynchronous . Except for avron2014revisiting , all these methods are totally async-parallel with the P-contraction condition or its variants. On solving a positive-definite linear system, avron2014revisiting made assumptions similar to ours, and it obtained better linear convergence rate on that special problem.
Our framework differs from the recent surge of the aforementioned sync-parallel and async-parallel coordinate descent algorithms (e.g., peng2013parallel ; kyrola2011parallel ; liu2013asynchronous ; liu2014asynchronous ; hsieh2015passcode ; richtarik2015parallel ). While they apply to convex function minimization, ARock covers more cases (such as ADMM, primal-dual, and decentralized methods) and also provides sequence convergence. In Section 2, we will show that some of the existing async-parallel coordinate descent algorithms are special cases of ARock, through relating their optimality conditions to nonexpansive operators. Another difference is that the convergence of ARock only requires a nonexpansive operator with a fixed point, whereas properties such as strong convexity, bounded feasible set, and bounded sequence, which are seen in some of the recent literature for async-parallel convex minimization, are unnecessary.
Others. Besides solving equations and optimization problems, there are also applications of async-parallel algorithms to optimal control problems LMS1986asynchronous , network flow problems ESMG1996asyn-flex , and consensus problems of multi-agent systems leifang_information_2005 .
6 Contributions
Our contributions and techniques are summarized below:
ARock is the first async-parallel coordinate update framework for finding a fixed point to a nonexpansive operator.
By introducing a new metric and establishing stochastic Fejér monotonicity, we show that, with probability one, ARock converges to a point in the solution set; linear convergence is obtained for quasi-strongly monotone operators.
Based on ARock, we introduce an async-parallel algorithm for linear systems, async-parallel ADMM algorithms for distributed or decentralized computing problems, as well as async-parallel operator-splitting algorithms for nonsmooth minimization problems. Some problems are treated in they async-parallel fashion for the first time in history. The developed algorithms are not straightforward modifications to their serial versions because their underlying nonexpansive operators must be identified before applying ARock.
7 Notation, definitions, background of monotone operators
Throughout this paper, denotes a separable Hilbert space equipped with the inner product and norm , and denotes the underlying probability space, where , , and are the sample space, -algebra, and probability measure, respectively. The map , where is the Borel -algebra, is an -valued random variable. Let denote either a sequence of deterministic points in or a sequence of -valued random variables, which will be clear from the context, and let denote the th coordinate of . In addition, we let denote the smallest -algebra generated by . “Almost surely” is abbreviated as “a.s.”, and the product space of is denoted by . We use and for strong convergence and weak convergence, respectively.
We define as the set of fixed points of operator , and, in the product space, we let .
An operator is -Lipschitz, where , if it satisfies , . In particular, is nonexpansive if , and contractive if .
Consider an operator .
is -averaged with , if there is a nonexpansive operator such that , where is the identity operator.
is -cocoercive with , if
is -strongly monotone, where , if it satisfies When the inequality holds for , is monotone.
is quasi--strongly monotone, where , if it satisfies . When the inequality holds for , is quasi-monotone.
is demicompact petryshyn1966construction at if for every bounded sequence in such that , there exists a strongly convergent subsequence.
Averaged operators are nonexpansive. By the Cauchy-Schwarz inequality, a -cocoercive operator is -Lipschitz; the converse is generally untrue, but true for the gradients of convex differentiable functions. Examples are given in the next section.
Applications
In this section, we provide some applications that are special cases of the fixed-point problem (1). For each application, we identify its nonexpansive operator (or the corresponding operator ) and implement the conditions in Theorem 1.1. For simplicity, we use the uniform distribution, , and apply the simpler update (4) instead of (3).
(bauschke2011convex, , Example 22.5) Suppose that is -Lipschitz continuous with . Then, is -strongly monotone.
Suppose . Since is -Lipschitz continuous, by Proposition 1, is -strongly monotone. By Theorem 3.3, Algorithm 2 converges linearly.
2 Minimize convex smooth function
where is a closed proper convex differentiable function and is -Lipschitz continuous, . Let . As is convex and differentiable, is a minimizer of if and only if is a zero of . Note that is -cocoercive. By Lemma 1, is nonexpansive. Applying ARock, we have the following iteration:
where . Note that needs a structure that makes it cheap to compute . Let us give two such examples: (i) quadratic programming: , where and only depends on a part of and ; (ii) sum of sparsely supported functions: and , where each depends on just a few variables.
Theorem 3.2 below guarantees the convergence of if . In addition, If is restricted strongly convex, namely, for any and , where is the solution set to (7), we have for some , then is quasi-strongly monotone with modulus . According to Theorem 3.3, iteration (8) converges at a linear rate if the step size meets the condition therein.
3 Decentralized consensus optimization
If the agents are independent Poisson processes and that each agent has activation rate , then the probability that agent activates before other agents is equal to larson1981urban and therefore our random sample scheme holds and ARock applies naturally. The algorithm is summarized as follows:
4 Minimize smooth ++ nonsmooth functions
+ nonsmooth functions Consider the problem
where is closed proper convex and is convex and -Lipschitz differentiable with . Problems in the form of (10) arise in statistical regression, machine learning, and signal processing and include well-known problems such as the support vector machine, regularized least-squares, and regularized logistic regression. For any and scalar , define the proximal operator and the reflective-proximal operator as
Assume that is a closed proper convex function, and is -Lipschitz differentiable and strongly convex with modulus . Let . Then, both and are quasi-contractive operators.
We first show that is a quasi-contractive operator. Note
where the first inequality follows from the Baillon-Haddad theoremLet be a convex differentiable function. Then, is -Lipschitz if and only if it is -cocoercive. and the second one from the strong convexity of . Hence, is quasi-contractive if . Since is convex, is firmly nonexpansive, and thus we immediately have the quasi-contractiveness of from that of .
5 Minimize nonsmooth ++ nonsmooth functions
where both and are closed proper convex and their maps are easy to compute. Define the Peaceman-Rachford lions1979splitting operator:
where and are intermediate variables. Note that the order in which the proximal operators are applied to and affects both YanYin2014 and whether coordinate-wise updates can be efficiently computed. Next, we present two special cases of (13) in Subsections 2.5.1 and 2.6 and discuss how to efficiently implement the update (15).
Suppose that are closed convex subsets of with a nonempty intersection. The problem is to find a point in the intersection. Let be the indicator function of the set , that is, if and otherwise. The feasibility problem can be formulated as the following
Let , , and . We can implement (15) as follows (see Appendix A for the step-by-step derivation):
The update (16) can be implemented as follows. Let global memory hold , as well as . At the th update, an agent independently generates a random number , then reads as and as , and finally computes and updates in global memory according to (16). Since is maintained in global memory, the agent updates according to . This implementation saves each agent from computing (16a) or reading all . Each agent only reads and , executes (16b), and updates (16c) and .
6 Async-parallel ADMM
This is another application of (15). Consider
where and are Hilbert spaces, and are bounded linear operators. We apply the update (15) to the Lagrange dual of (17) (see gabay1983chapter for the derivation):
where , , and and denote the convex conjugates of and , respectively. The proximal maps induced by and can be computed via solving subproblems that involve only the original terms in (17): can be computed by (see Appendix A for the derivation)
and by
Plugging (19) and (20) into (15) yields the following naive implementation
Note that in (15c) becomes in (21e) because ADMM is equivalent to the Douglas-Rachford operator, which is the average of the Peaceman-Rachford operator and the identity operator lions1979splitting . Under favorable structures, (21) can be implemented efficiently. For instance, when and are block diagonal matrices and are corresponding block separable functions, steps (21a)–(21d) reduce to independent computation for each . Since only and are needed to update the main variable , we only need to compute (21a)–(21d) for the th block. This is exploited in distributed and decentralized ADMM in the next two subsections.
Consider the consensus optimization problem:
where are proper close convex functions. Rewrite (22) to the ADMM form:
where . Now apply the async-parallel ADMM (21) to (2.6.1) with dual variables . In particular, the update (21a), (21b), (21c), (21d) reduce to
Therefore, we obtain the following async-parallel ADMM algorithm for the problem (22). This algorithm applies to all the distributed applications in boyd2011distributed .
6.2 Async-parallel ADMM for decentralized optimization
Let be a set of agents and E=\{(i,j)~{}|~{}\text{if agentij},i<j\} be the set of undirected links between the agents. Consider the following decentralized consensus optimization problem on the graph :
for proper matrices and . Applying the async-parallel ADMM (21) to (29) gives rise to the following simplified update: Let be the set of edges connected with agent and be its cardinality. Let and . To every pair of constraints and , , we associate the dual variables and , respectively. Whenever some agent is activated, it calculates
We present the algorithm based on (30) for problem (27) in Algorithm 5.
Algorithm 5 activates one agent at each iteration and updates all the dual variables associated with the agent. In this case, only one-sided communication is needed, for sending the updated dual variables in the last step. We allow this communication to be delayed in the sense that agent ’s neighbors may be activated and start their computation before receiving the latest dual variables from agent .
Our algorithm is different from the asynchronous ADMM algorithm by Wei and Ozdaglar wei2013on . Their algorithm activates an edge and its two associated agents at each iteration and thus requires two-sided communication at each activation. We can recover their algorithm as a special case by activating an edge and its associated agents and at each iteration, updating the dual variables and associated with the edge, as well as computing the intermediate variables , , and . The updates are derived from (29) with the orders of and swapped. Note that wei2013on does not consider the situation that adjacent edges are activated in a short period of time, which may cause overlapped computation and delay communication. Indeed, their algorithm corresponds to and the corresponding stepsize . Appendix B presents the steps to derive the algorithms in this subsection.
Convergence
We establish weak and strong convergence in Subsection 3.1 and linear convergence in Subsection 3.2. Step size selection is also discussed.
Throughout the our analysis, we assume and
We let be the number of elements in (see Subsection 1.2). Only for the purpose of analysis, we define the (never computed) full update at th iteration:
Lemma 1 below shows that is nonexpansive if and only if is -cocoercive.
Operator is nonexpansive if and only if is -cocoercive, i.e., .
See textbook (bauschke2011convex, , Proposition 4.33) for the proof of the “if” part, and the “only if” part, though missing there, follows by just reversing the proof.
The lemma below develops an an upper bound for the expected distance between and any .
Let be the sequence generated by Algorithm 1. Then for any and (to be optimized later), we have
where the first inequality follows from the Young’s inequality. Plugging (35) and (36) into (34) gives the desired result.
We need the following lemma on nonnegative almost supermartingales robbins1985convergence .
Let be a product space and be the induced inner product:
Let be a symmetric tri-diagonal matrix with its main diagonal as and first off-diagonal as , and let . Here represents the Kronecker product. For a given , is given by:
Then is a self-adjoint and positive definite linear operator since is symmetric and positive definite, and we define as the -weighted inner product and the induced norm. Let
where we set for . With
we have the following fundamental inequality:
Let be the sequence generated by ARock. Then for any , it holds that
Let . Since , then (33) indicates
Let us check our step size bound . Consider the uniform case: . Then, the bound simplifies to . If the max delay is no more than the square root of coordinates, i.e., , then the bound is . In general, depends on several factors such as problem structure, system architecture, load balance, etc. If all updates and agents are identical, then is proportional to , the number of agents. Hence, ARock takes an step size for solving a problem with coordinates by agents under balanced loads.
The next lemma is a direct consequence of the invertibility of the metric .
A sequence (weakly) converges to under the metric if and only if it does so under the metric .
In light of Lemma 4, the metric of the inner product for weak convergence in the next lemma is not specified. The lemma and its proof are adapted from combettes2014stochastic .
Let be the sequence generated by ARock with for any and . Then we have:
a.s..
a.s. and a.s..
The sequence is bounded a.s..
Let be the set of weakly convergent cluster points of . Then, a.s..
(i): Note that . Also note that, in (42), is -measurable. Hence, applying Lemma 3 with and to (42) gives this result directly.
(ii) From (i), we have a.s.. Since , we have a.s.. Then from (5), we have a.s..
(iii): From Lemma 3, we have that converges a.s. and so does , i.e., a.s., where is a -valued random variable. Hence, must be bounded a.s. and so is .
(v): By (ii), there exists such that and
For any , let be a weakly convergent subsequence of , i.e., , where and . Note that implies Therefore, , for any because .
Furthermore, observing , we have
From the triangle inequality and the nonexpansiveness of , it follows that
From (44), (45), and the above inequality, it follows Finally, the demiclosedness principle (bauschke2011convex, , Theorem 4.17) implies .
Under the assumptions of Lemma 5, the sequence weakly converges to an -valued random variable a.s.. In addition, if is demicompact at 0, strongly converges to an -valued random variable a.s..
For the generalization in Section 1.3, we need to replace (35) by
and update the step size condition to . Then the proofs of Theorem 3.2 and Lemma 5 will go through and yield the same convergence result.
2 Linear convergence
In this section, we establish linear convergence under the assumption that is quasi-strongly monotone. We first present a key lemma.
Assume that the step size is fixed, i.e., , and satisfies
for some . Then we have, for all ,
We prove (47) by induction. First, based on the inequality we observe that, for any ,
Applying the triangle inequality and (5) yields
For the basic case, we have , . Letting in (48) gets us
For the induction step, applying Young’s inequality gives us
Taking the expectation on (51) and combining it with (48) yield
With this lemma, we are ready to derive the linear convergence rate of ARock.
Assume that is quasi--strongly monotone with . Let and be the sequence generated by ARock with a constant stepsize , where is given in (46) and
Following the proof of Lemma 2 and starting from (36), we have
where the second inequality holds because is -cocoercive and also quasi--strongly monotone, and the last one comes from the Cauchy-Schwartz inequality. Plugging the above inequality and (35) into (34) and noting gives
where we have let in the second equality, and the last inequality holds because of the choice of . Therefore, (55) holds.
Assume is chosen uniformly at random, so . We consider the case when and are large. Let . Then from the fact that increasingly converges to the natural number , we have from (46) that . In addition, note from (54) that , and thus . Therefore, if , then the stepsize in Theorem 3.3 can be . Hence, linear speedup can be achieved.
Experiments
Our experiments run on 1 to 32 threads on a machine with eight Quad-Core AMD OpteronTM Processors (32 cores in total) and Gigabytes of RAM. All of the experiments were coded in C++ and OpenMP. We use the Eigen libraryhttp://eigen.tuxfamily.org for sparse matrix operations. Our codes as well as numerical results for other applications will be publicly available on the authors’ website.
The running times and speedup ratios of both sync-parallel and async-parallel algorithms are sensitive to a number of factors, such as the size of each coordinate update (granularity), sparsity of the problem data, compiler optimization flags, and operations that affect cache performance and memory access contention. In addition, since all agents in the sync-parallel implementation must wait for the last agent to finish an iteration, a large load imbalance will significantly degrade the performance. We do not have the space in this paper to present numerical results under all variations of these cases.
where is the set of sample-label pairs with , , and and represent the numbers of features and samples, respectively. This test uses the datasetshttp://www.csie.ntu.edu.tw/~cjlin/libsvmtools/datasets/: rcv1 and news20, which are summarized in Table 1.
We let each coordinate hold roughly 50 features. Since the total number of features is not divisible by 50, some coordinates have 51 features. We let each agent draw a coordinate uniformly at random at each iteration. We stop all the tests after 100 epochs since they have nearly identical progress per iteration. The step size is set to . Let and . In global memory, we store , and . We also store the product in global memory so that the forward step can be efficiently computed. Whenever a coordinate of gets updated, is immediately updated at a low cost. Note that if is not stored in global memory, every coordinate update will have to compute from scratch, which involves the entire and will be very expensive.
Table 2 gives the running times of the sync-parallel and ARock (async-parallel) implementations on the two datasets. We can observe that ARock achieves almost-linear speedup, but sync-parallel scales very poorly as we explain below.
In the sync-parallel implementation, all the running cores have to wait for the last core to finish an iteration, and therefore if a core has a large load, it slows down the iteration. Although every core is (randomly) assigned to roughly the same number of features (either 50 or 51 components of ) at each iteration, their ’s have very different numbers of nonzeros (see Figure 3 for the distribution), and the core with the largest number of nonzeros is the slowest (Sparse matrix computation is used for both datasets, which are very large.) As more cores are used, despite that they altogether do more work at each iteration, the per-iteration time reduces as the slowest core tends to be slower. The very large imbalance of load explains why the 32 cores only give speedup ratios of 4.0 and 1.3 in Table 2.
On the other hand, being asynchronous, ARock does not suffer from the load imbalance. Its performance grows nearly linear with the number of cores. In theory, a large load imbalance may cause a large , and thus a small . However, the uniform works well in all the tests, possibly because the ’s are sparse.
Finally, we have observed that the progress toward solving (56) is mainly a function of the number of epochs and does not change appreciably when the number of cores increases or between sync-parallel and async-parallel. Therefore, we always stop at 100 epochs.
Conclusion
We have proposed an async-parallel framework, ARock, for finding a fixed-point of a nonexpansive operator by coordinate updates. We establish the almost sure weak and strong convergence, linear convergence rate and almost-linear speedup of ARock under certain assumptions. Preliminary numerical results on real data illustrate the high efficiency of the proposed framework compared to the traditional parallel (sync-parallel) algorithms.
Acknowledgements
We would like to thank Brent Edmunds for offering invaluable suggestions on the organization and writing of this paper. We would also like to thank Robert Hannah for coming up with the dual-memory approach. The authors are grateful to Kun Yuan for helpful discussions on decentralized optimization.
References
Appendix A Derivation of certain updates
We show in details how to obtain the updates in (16) and (21).
Let ,
where equals if and otherwise. Then (15a) reduces to
where the last equality is obtained by noting that is the unique minimizer of . Next, (15b) reduces to
Since (15c) only updates the th coordinate of , we only need and , and thus in (16a) and (16b), we only compute and . Plugging the above and into (15c) gives (16c) directly.
A.2 Derivation of (21)
We first show how to get (18). The Lagrangian of (17) is and the Lagrange dual function is
where the last equality is from the definition of convex conjugate: Hence, the dual problem is , which is equivalent to (18).
Secondly, we show why is given by (20). Note
where the fifth equality holds because Hence, by the definition of the proximal operator and the above arguments, we have that can be obtained from (20). Then (19) is from (20) through replacing to , to , and to .
Finally, it is straightforward to have (21) by plugging (19) and (20) into (15).
Appendix B Derivation of async-parallel ADMM for decentralized optimization
This section describes how to implement the updates (21) for the model (29).
In (29), and vanish and, corresponding to the two constraints and , the two rows of matrices and are where are zeros, the two coefficients 1 correspond to and , and the two coefficients correspond to . Then, (21a) and (21b) can be calculated as
In addition, can be obtained by solving (30a), and both and can be updated from (30b) and (30c).
Furthermore, as mentioned in Section 2.6.2, we can derive another version of async-parallel ADMM for decentralized optimization, which reduces to the algorithm in wei2013on , by activating an edge instead of an agent each time. In this version, the agents and associated with the edge must also be activated. Here we derive the update (21) for the model (29) with the update order of and swapped. Following (21) we obtain the following steps whenever an edge is activated:
Every agent in the network maintains the dual variables , , and , , and the variables are intermediate and do not need to be maintained between the activations. When an edge is activated, the agents and first compute their and independently and respectively, then they collaboratively compute , and finally they update their own and , respectively. We allow adjacent edges (which share agents) to be activated in a short period of time when their updates are possibly overlapped in time. When , i.e., there is no simultaneous activation or overlap, it reduces to the algorithm in wei2013on .