CYCLADES: Conflict-free Asynchronous Machine Learning

Xinghao Pan, Maximilian Lam, Stephen Tu, Dimitris Papailiopoulos, Ce Zhang, Michael I. Jordan, Kannan Ramchandran, Chris Re, Benjamin Recht

Introduction

Following the seminal work of Hogwild! [NRRW11], many studies have demonstrated that near-linear speedups are achievable on a variety of machine learning tasks via asynchronous, lock-free implementations [RRTB12, ZCJL13, YYH+13, LWR+14, DJM13, WSD+14, HYD15, MPP+15]. In all of these studies, classic algorithms are parallelized by simply running parallel and asynchronous model updates without locks. These lock-free, asynchronous algorithms exhibit speedups even when applied to large, non-convex problems, as demonstrated by deep learning systems such as Google’s Downpour SGD [DCM+12] and Microsoft’s Project Adam [CSAK14].

While these techniques have been remarkably successful, many of the papers cited above require delicate and tailored analyses to understand the benefits of asynchrony for the particular learning task at hand. Moreover, in non-convex settings, we currently have little quantitative insight into how much speedup is gained from asynchrony and how much accuracy may be lost.

In this work, we present Cyclades, a general framework for lock-free, asynchronous machine learning algorithms that obviates the need for specialized analyses. Cyclades runs asynchronously and maintains serial equivalence, i.e., it produces the same outcome as the serial algorithm. Since it returns exactly the same output as a serial implementation, any algorithm parallelized by our framework inherits the correctness proof of the serial counterpart without modifications. Additionally, if a particular heuristic serial algorithm is popular, but does not have a rigorous analysis, such as backpropagation on neural networks, Cyclades still guarantees that its execution will return a serially equivalent output.

Cyclades achieves serial equivalence by partitioning updates among cores, in a way that ensures that there are no conflicts across partitions. Such a partition can always be found efficiently by leveraging a powerful result on graph phase transitions [Kri16]. When applied to our setting, this result guarantees that a sufficiently small sample of updates will have only a logarithmic number of conflicts. This allows us to evenly partition model updates across cores, with the guarantee that all conflicts are localized within each core. Given enough problem sparsity, Cyclades guarantees a nearly linear speedup, while inheriting all the qualitative properties of the serial counterpart of the algorithm, e.g., proofs for rates of convergence. Enforcing a serially equivalent execution in Cyclades comes with additional practical benefits. Serial equivalence is helpful for hyperparameter tuning, or locating the best model produced by the asynchronous execution, since experiments are reproducible, and solutions are easily verifiable. Moreover, a Cyclades program is easy to debug, because bugs are repeatable and we can examine the step-wise execution to localize them.

A significant benefit of the update partitioning in Cyclades is that it induces considerable access locality compared to the more unstructured nature of the memory accesses during Hogwild!. Cores will access the same data points and read/write the same subset of model variables. This has the additional benefit of reducing false sharing across cores. Because of these gains, Cyclades can actually outperform Hogwild! in practice on sufficiently sparse problems, despite appearing to require more computational overhead. Remarkably, because of the added locality, even a single threaded implementation of Cyclades can actually be faster than serial SGD. In our SGD experiments for matrix completion and word embedding problems, Cyclades can offer a speedup gain of up to 40%40\% compared to that of Hogwild!. Furthermore, for variance reduction techniques such as SAGA [DBLJ14] and SVRG [JZ13], Cyclades yields better accuracy and more significant speedups, with up to 5×5\times performance gains over Hogwild!-type implementations.

The remainder of our paper is organized as follows. Section 2 establishes some preliminaries. Details and theory of Cyclades are presented in Section 3. We present our experiments in Section 4, we discuss related work in Section 5, and then conclude with Section 6.

The Algorithmic Family of Stochastic-Updates

We study parallel asynchronous iterative algorithms on the computational model used by [NRRW11], and similar to the partially asynchronous model of [BT89]: a number of cores have access to the same shared memory, and each of them can read and update components of x{\bf x} in parallel from the shared memory.

In this work, we consider a large family of randomized algorithms that we will refer to as Stochastic Updates (SU). The main algorithmic component of SU focuses on updating small subsets of a model variable x{\bf x}, that lives in shared memory, according to prefixed access patterns, as sketched by Alg. 1.

In Alg. 1 each set Si\mathcal{S}_{i} is a subset of the coordinate indices of x{\bf x}, each function fif_{i} only operates on the subset Si\mathcal{S}_{i} of coordinates (i.e., both its domain and co-domain are inside Si\mathcal{S}_{i}), and uiu_{i} is a local update function that computes a vector with support on Si\mathcal{S}_{i} using as input xSi{\bf x}_{\mathcal{S}_{i}} and fif_{i}. Moreover, TT is the total number of iterations, and D\mathcal{D} is the distribution with support {1,…,n}\{1,\ldots,n\} from which we draw ii. As we explain in Appendix A, several machine learning and optimization algorithms belong to the SU algorithmic family, such as stochastic gradient descent (SGD), with or without weight decay and regularization, variance-reduced learning algorithms like SAGA and SVRG, and even some combinatorial graph algorithms.

A useful construction for our developments is the conflict graph between updates, which can be generated from the bipartite graph between the updates and the model variables. We define these graphs below, and provide an illustrative sketch in Fig. 1.

We denote as GuG_{u} the bipartite update-variable graph between the updates u1,…,unu_{1},\ldots,u_{n} and the dd model variables. In GuG_{u} an update uiu_{i} is linked to a variable xjx_{j}, if uiu_{i} requires to read or write xjx_{j}. We let EuE_{u} denote the number of edges in the bipartite graph, and also denote as ΔL\Delta_{L} the left max vertex degree of GuG_{u}, and as Δ‾L\overline{\Delta}_{L} its average left degree.

We denote by GcG_{c} a conflict graph on nn vertices, each corresponding to an update uiu_{i}. Two vertices of GcG_{c} are linked with an edge, if and only if the corresponding updates share at least one variable in the bipartite-update graph GuG_{u}. We also denote as Δ\Delta the max vertex degree of GcG_{c}.

We stress that the conflict graph is never actually constructed, but serves as a useful concept for understanding Cyclades.

Our Main Result

By exploiting the structure of the above graphs and through a light-weight and careful sampling and allocation of updates, Cyclades is able to guarantee the following result for SU algorithms, which we establish in the following sections.

We will now provide some simple examples of how the above parameters, and guarantees translate for specific problem cases.

Many machine learning applications often seek to minimize the empirical risk

Consider the following generic minimization problem

where ϕi,j\phi_{i,j} is a convex function of a scalar. The above generic formulation captures several problems like matrix completion and matrix factorization [RR13] (where ϕi,j=(Ai,j−xiTyj)2\phi_{i,j}=(A_{i,j}-{\bf x}_{i}^{T}{\bf y}_{j})^{2}), word embeddings [ALL+15] (where ϕi,j=Ai,j(log⁡(Ai,j)−∥xi+xj∥22−C)2\phi_{i,j}=A_{i,j}(\log(A_{i,j})-\|{\bf x}_{i}+{\bf x}_{j}\|_{2}^{2}-C)^{2}), graph kk-way cuts [NRRW11] (where ϕi,j=Ai,j∥xi−xj∥1\phi_{i,j}=A_{i,j}\|{\bf x}_{i}-{\bf x}_{j}\|_{1}), and others. Let m1=m2=mm_{1}=m_{2}=m for simplicity, and assume that we aim to minimize the above by sampling a single function ϕi,j\phi_{i,j} and then updating xi{\bf x}_{i} and yj{\bf y}_{j} using SGD. Here, the number of update functions is proportional to n=m2n=m^{2}, and for the above setup each gradient update with respect to the sampled function ϕi,j(xi,yj)\phi_{i,j}({\bf x}_{i},{\bf y}_{j}) is only interacting with the variables xi{\bf x}_{i} and yj{\bf y}_{j}, i.e., only two variable vectors out of the 2m2m many (i.e., ΔL=2\Delta_{L}=2). Moreover, the previous imply a conflict degree of at most Δ=2m\Delta=2m. In this case, Cyclades can provably guarantee an Ω~(P)\widetilde{\Omega}(P) speedup for up to P=O(m)P=O(m) cores.

In our experiments we test Cyclades on several problems including least squares, classification with logistic models, matrix factorization, and word embeddings, and several algorithms including SGD, SVRG, and SAGA. We show that in most cases it can significantly outperform the Hogwild! implementation of these algorithms, if the data is sparse.

We would like to note, that there are several cases where there might be a few outlier updates with extremely high conflict degree. In Appendix E, we prove that if there are no more than O(nδ)O(n^{\delta}) vertices of high conflict degree Δo\Delta_{o}, and the rest of the vertices have max degree at most Δ\Delta, then the result of Theorem 1 still holds in expectation.

In the following section, we establish the technical results behind Cyclades and provide the details behind our parallelization framework.

Cyclades: Shattering Dependencies

Cyclades consists of three computational components as shown in Figure 2.

It starts by sampling (according to a distribution D\mathcal{D}) a number of BB updates from the graph shown in Fig. 1, and assigns a label to each of them (a processing order). We note that in practice the sampling is done on the bipartite graph, which avoids the need to actually construct the conflict graph. After sampling, it computes the connected components of the sampled subgraph induced by the BB sampled updates, to determine the conflict groups. Once the conflicts groups are formed, it allocates them across PP cores. Finally, each core processes locally the conflict groups of updates that it has been assigned, following the order that each update has been labeled with. The above process is then repeated, for as many iterations as needed.

The key component of Cyclades is to carry out the sampling in such a way that we have as many connected components as possible, and all of them of small size, provably. In the next subsections, we explain how each part is carried out, and provide theoretical guarantees for each of them individually, which we combine at the end of this section for our main theorem.

A key technical aspect that we exploit in Cyclades is that appropriate sampling and allocation of updates can lead to near optimal parallelization of sparse SU algorithms. To do that we expand upon the following result established in [Kri16].

Let GG be a graph on nn vertices, with maximum vertex degree Δ\Delta. Let us sample each vertex independently with probability p=1−ϵΔp=\frac{1-\epsilon}{\Delta} and define as G′G^{\prime} the induced subgraph on the sampled vertices. Then, the largest connected component of G′G^{\prime} has size at most 4ϵ2log⁡n\frac{4}{\epsilon^{2}}\log n, with high probability.

The above result pays homage to the giant component phase transition phenomena in random Erdos-Renyi graphs. What is surprising is that a similar phase transition can apply for any given graph!

Adapting to ML-friendly sampling procedures.

In practice, for most SU algorithms of interest, the sampling distribution of updates is either with or without replacement from the nn updates. As it turns out, morphing Theorem 2 into a with-/without-replacement result is not straightforward. We defer the analysis needed to Appendix B, and present our main theorem about graph sampling here.

Let GG be a graph on nn vertices, with maximum vertex degree Δ\Delta. Let us sample B=(1−ϵ)nΔB=(1-\epsilon)\frac{n}{\Delta} vertices with or without replacement, and define as G′G^{\prime} the induced subgraph on the sampled vertices. Then, the largest connected component of G′G^{\prime} has size at most O(log⁡nϵ2)O(\frac{\log n}{\epsilon^{2}}), with high probability.

The key idea from the above theorem is that if one samples no more than B=(1−ϵ)nΔB=(1-\epsilon)\frac{n}{\Delta} vertices, then there will be at least O(\nicefracϵ2Blog⁡n)O\left(\nicefrac{{\epsilon^{2}B}}{{\log n}}\right) conflict groups to allocate across cores, all of size at most O(log⁡n/ϵ2)O\left(\log n/\epsilon^{2}\right). Moreover, since there are no conflicts between different conflict-groups, the processing of updates per any single group will never interact with the variables corresponding to the updates of another conflict group.

The next step of Cyclades is to form and allocate the connected components (CCs) across cores, and do so efficiently. We address this in the following subsection. In the following, for simplicity we carry our analysis for the with-replacement sampling case, but it can be readily extended to the without-replacement sampling case.

Identifying Groups of Conflict via CCs

In Cyclades, we sample batches of updates of size B=(1−ϵ)nΔB=(1-\epsilon)\frac{n}{\Delta} multiple times, and for each batch we need to identify the conflict groups across the updates. Let us refer to GuiG_{u}^{i} as the subgraph induced by the iith sampled batch of updates on the update-variable bipartite graph GuG_{u}. In the following we always assume that we sample at most nb=c⋅Δ1−ϵn_{b}=c\cdot\frac{\Delta}{1-\epsilon} batches, where c≥1c\geq 1 is a constant that does not depend on nn. This number of batches results in a constant number of passes over the dataset.

Identifying the conflict groups in GuiG_{u}^{i} can be done with a connected components (CC) algorithm. The main question we need to address is what is the best way to parallelize this graph partitioning part. There are two avenues that we can take for this, depending on the number of cores PP at our disposal. We can either parallelize the computation of the CCs of a single batch (i.e., compute the CCs of GuiG_{u}^{i} on PP cores), or we can compute in parallel the CCs of all nbn_{b} batches, by allocating the sampled graphs GuiG_{u}^{i} to cores, so that each of them can compute the CCs of its allocated subgraphs. Depending on the number of available cores, one technique can be better than the other. In Appendix C we provide the details of this part, and prove the following result:

Let the number of cores by bounded as P=O(nΔΔL)P=O(\frac{n}{\Delta\Delta_{\text{L}}}), and let ΔLΔ‾L≤n\frac{\Delta_{\text{L}}}{\overline{\Delta}_{\text{L}}}\leq\sqrt{n}. Then, the overall computation of CCs for nb=c⋅Δ1−ϵn_{b}=c\cdot\frac{\Delta}{1-\epsilon} batches, each of size B=(1−ϵ)nΔB=(1-\epsilon)\frac{n}{\Delta}, costs no more than O(Eulog⁡2nP)O(\frac{E_{u}\log^{2}n}{P}).

Allocating Updates to Cores

Once we compute the CCs (i.e., the conflicts groups of the sampled updates), we have to allocate them across cores. Once a core has been assigned with CCs, it will process the updates included in these CCs, according to the order that each update has been labeled with. Due to Theorem 3, each connected component will contain at most O(log⁡nϵ2)O(\frac{\log n}{\epsilon^{2}}) updates. Assuming that the cost of the jj-th update in the batch is wjw_{j}, the cost of a single connected component C\mathcal{C} will be wC=∑j∈Cwjw_{\mathcal{C}}=\sum_{j\in\mathcal{C}}w_{j}. To proceed with characterizing the maximum load among the PP cores, we assume that the cost of a single update uiu_{i}, for i∈{1,…,n}i\in\{1,\ldots,n\}, is proportional to the out-degree of that update —according to the update-variable graph GuG_{u}— times a constant cost which we shall refer to as κ\kappa. Hence, wj=O(dL,j⋅κ)w_{j}=O(d_{L,j}\cdot\kappa), where dL,jd_{L,j} is the degree of the jj-th left vertex of GuG_{u}. In Appendix D we establish that a near-uniform allocation of CCs according to their weights leads to the following guarantee.

Let the number of cores by bounded as P=O(nΔΔL)P=O(\frac{n}{\Delta\Delta_{\text{L}}}), and let ΔLΔ‾L≤n\frac{\Delta_{\text{L}}}{\overline{\Delta}_{\text{L}}}\leq\sqrt{n}. Then, computing the stochastic updates across all nb=c⋅Δ1−ϵn_{b}=c\cdot\frac{\Delta}{1-\epsilon} batches can be performed in time O(Elog⁡2nP⋅κ)O(\frac{E\log^{2}n}{P}\cdot\kappa), with high probability, where κ\kappa is the per edge cost for computing one of the nn updates defined on GuG_{u}.

Stitching the pieces together

Now that we have described the sampling, conflict computation, and allocation strategies, we are ready to put all the pieces together and detail Cyclades in full. Let us assume that we sample a total number of nb=c⋅Δ1−ϵn_{b}=c\cdot\frac{\Delta}{1-\epsilon} batches of size B=(1−ϵ)nΔB=(1-\epsilon)\frac{n}{\Delta}, and that each update is sampled uniformly at random. For the ii-th batch let us denote as C1i,…Cmii\mathcal{C}_{1}^{i},\ldots\mathcal{C}_{m_{i}}^{i} the connected components on the induced subgraph GuiG_{u}^{i}. Due to Theorem 3, each connected component C\mathcal{C} contains a number of at most O(log⁡nϵ2)O(\frac{\log n}{\epsilon^{2}}) updates, and each update carries an ID (the order of which it would have been sampled by the serial algorithm). Using the above notation, we give the pseudocode for Cyclades in Alg. 2.

Note that the inner loop that is parallelized (i.e., the SU processing loop in lines 6 – 9), can be performed asynchronously; cores do not have to synchronize, and do not need to lock any memory variables, as they are all accessing non-overlapping subset of x{\bf x}. This also provides for better cache coherence. Moreover, each core potentially accesses the same coordinates several times, leading to good cache locality. These improved cache locality and coherence properties experimentally lead to substantial performance gains as we see in the next section.

We can now combine the results of the previous subsection to obtain our main theorem for Cyclades.

Let us assume any given update-variable graph GuG_{u} with average, and max left degree Δ‾L\overline{\Delta}_{L} and ΔL{\Delta}_{L}, such that ΔLΔ‾L≤n\frac{\Delta_{\text{L}}}{\overline{\Delta}_{\text{L}}}\leq\sqrt{n}, and with induced max conflict degree Δ\Delta. Then, Cyclades on P=O(nΔ⋅ΔL)P=O(\frac{n}{\Delta\cdot\Delta_{L}}) cores, with batch sizes B=(1−ϵ)nΔB=(1-\epsilon)\frac{n}{\Delta} can execute T=c⋅nT=c\cdot n updates, for any constant c≥1c\geq 1, selected uniformly at random with replacement, in time

Observe that Cyclades bypasses the need to establish convergence guarantees for the parallel algorithm. Hence, it could be the case for many applications of interest that although we might not be able to analyze how “well” the serial SU algorithm might perform in terms of the accuracy of the solution, Cyclades can provide black box guarantees for speedup, since our analysis is completely oblivious to the qualitative performance of the serial algorithm. This is in contrast to recent studies similar to [DSZOR15], where the authors provide speedup guarantees via a convergence-to-optimal proof for an asynchronous SGD on a nonconvex problem. Unfortunately these proofs can become complicated especially on a wider range of nonconvex objectives.

In the following section we show that Cyclades is not only useful theoretically, but can consistently outperform Hogwild! on sufficiently sparse datasets.

Evaluation

We implemented Cyclades in C++ and tested it on a variety of problems and datasets described below. We tested a number of stochastic updates algorithms, and compared against their Hogwild! (i.e., asynchronous and lock-free) implementations — in some cases, there are no theoretical foundations for these Hogwild! implementations, even if they work reasonably well in practice. Since Cyclades is intended to be a general approach for parallelization of stochastic updates algorithms, we do not compare against algorithms designed and tailored for specific applications, nor do we expect Cyclades to outperform every such highly-tuned, well-designed, specific algorithms.

Our experiments were conducted on a machine with 72 CPUs (Intel(R) Xeon(R) CPU E7-8870 v3, 2.10 GHz) on 4 NUMA nodes, each with 18 CPUs, and 1TB of memory. We ran both Cyclades and Hogwild! with 1, 4, 8, 16 and 18 threads pinned to CPUs on a single NUMA node (i.e., the maximum physical number of cores possible, for a single node), so that we can avoid well-known cache coherence and scaling issues across different nodes [ZR14]. We note that distributing threads across NUMA nodes significantly increased running times for both Cyclades and Hogwild!, but was relatively worse for Hogwild!. We believe this is due to the poorer locality of Hogwild!, which results in more cross-node communication. In this paper, we exclusively focus our study and experiments on parallelization within a single NUMA node, and leave cross-NUMA node parallelization for future work, while referring the interested reader to a recent study of the various tradeoffs of ML algorithms on NUMA aware architectures [ZR14].

In our experiments, we measure overall running times which include the overheads for computing connected components and allocating work in Cyclades. Separately, we also measure running times for performing the stochastic updates by excluding the Cyclades coordination overheads. We also compute the objective value at the end of each epoch (i.e., one full pass over the data). We measure the speedups for each algorithm as

where ϵ\epsilon was chosen to be the smallest objective value that is achievable by all parallel algorithms on every choice of number of threads. That is, ϵ=max⁡A,Tmin⁡ef(XA,T,e)\epsilon=\max_{\mathcal{A},T}\min_{e}f(X_{\mathcal{A},T,e}) where XA,T,eX_{\mathcal{A},T,e} is the model learned by algorithm A\mathcal{A} on TT threads after ee epochs. The serial algorithm used for comparison is Hogwild! running serially on one thread.

In Table 1 we list some details of the datasets that we use in our experiments. The stepsizes and batch sizes used for each problem are listed in Table 2, along with dataset and problem details. In general, we chose the stepsizes to maximize convergence without diverging. Batch sizes were picked to optimize performance for Cyclades.

2 Learning tasks and algorithmic setup

The first problem we consider is least squares:

which we will solve using the SAGA algorithm [DBLJ14], an incrimental gradient algorithm with faster than SGD rates on convex, or strongly convex functions. In SAGA, we initialize gi=∇fi(x0){\bf g}_{i}=\nabla f_{i}({\bf x}_{0}) and iterate the following two steps

Graph eigenvector via SVRG

Given an adjacency matrix A{\bf A}, the top eigenvector of ATA{\bf A}^{T}{\bf A} is useful in several applications such as spectral clustering, principle component analysis, and others. In a recent work, [JKM+15] proposes an algorithm for computing the top eigenvector of ATA{\bf A}^{T}{\bf A} by running intermediate SVRG steps to approximate the shift-and-invert iteration. Specifically, at each step SVRG is used to solve

According to [JKM+15], if we initialize y=x0{\bf y}={\bf x}_{0} and assume ∥ai∥=1\|{\bf a}_{i}\|=1, we have to iterate the following updates

where after every TT iterations we update y=xk{\bf y}={\bf x}_{k}, and the stochastic gradients are of the form ∇fi(x)=(λnI−aiaiT)x−1nb.\nabla f_{i}({\bf x})=\left(\frac{\lambda}{n}{\bf I}-{\bf a}_{i}{\bf a}_{i}^{T}\right){\bf x}-\frac{1}{n}{\bf b}.

Matrix completion via SGD

In the matrix completion problem, we are given a partially observed n×mn\times m matrix M{\bf M}, and wish to factorize it as M≈UV{\bf M}\approx{\bf U}{\bf V} where U{\bf U} and V{\bf V} are low rank matrices with dimensions n×rn\times r and r×mr\times m respectively. This may be achieved by optimizing

where Ω\Omega is the set of observed entries, which can be approximated by SGD on the observed samples. The objective can also be regularized as:

The regularized objective can be optimized by weighted SGD, which samples (i,j)∈Ω(i,j)\in\Omega and updates

and analogously for V⋅,j{\bf V}_{\cdot,j} In our experiments, we chose a rank of r=100r=100, and ran SGD and weighted SGD for 200 epochs. We used the MovieLens 10M dataset [Gro09] containing 10M ratings for 10,000 movies by 72,000 users.

Word embedding via SGD

where Aw,w′A_{w,w^{\prime}} is the number of times words ww and w′w^{\prime} co-occur within τ\tau words in the corpus. In our experiments we set τ=10\tau=10 following the suggested recipe of the aforementioned paper. We can approximate the solution to the above problem by SGD: we can repeatedly sample entries Aw,w′A_{w,w^{\prime}} from A{\bf A} and update the corresponding vectors vw,vw′{\bf v}_{w},{\bf v}_{w^{\prime}}. In this case the update is of the form as:

and identically for vw′{\bf v}_{w^{\prime}} Then, at the end of each full pass over the data, we update the constant CC by its locally optimal value, which can be calculated in closed form:

In our experiments, we optimized for a word embedding of dimension r=100r=100, and tested on a 8080MB subset of the English Wikipedia dump available at [Mah06]. The dataset contains 213213K words and A{\bf A} has 20M non-zero entries. For our experiments, we run SGD for 200 epochs.

3 Speedup and Convergence Results

In this subsection, we present the bulk of our experimental findings. Our extended and complete set of results can be found in Appendix F.

When running SAGA for least squares, we found that Hogwild! was divergent with the large stepsizes that we were using for Cyclades (Fig. 5). Thus, in the multi-thread setting, we were only able to use smaller stepsizes for Hogwild!, which resulted in slower convergence than Cyclades, as seen in Fig. 3(a). The effects of a smaller stepsize for Hogwild! are also manisfested in terms of speedups in Fig. 4(a), since Hogwild! takes a longer time to converge to an ϵ\epsilon objective value.

Graph eigenvector

The convergence of SVRG for graph eigenvectors is shown in Fig. 3(b). Cyclades starts off slower than Hogwild!, but always produces results equivalent to the convergence on a single thread. Conversely, Hogwild! does not exhibit the same behavior on multiple threads as it does serially—in fact, the error due to asynchrony causes Hogwild! to converge slower on multiple threads. This effect is clearly seen on Figs. 4(b), where Hogwild! fails to converge faster than the serial counterpart, and Cyclades attains a significantly better speedup on 16 threads.

Matrix completion and word embeddings

Figures 3(c) and 3(d) show the convergence for the matrix completion and word embeddings problems. Cyclades is initially slower than Hogwild! due to the overhead of computing connected components. However, due to better cache locality and convergence properties, Cyclades is able to reach a lower objective value in less time than Hogwild!. In fact, we observe that Cyclades is faster than Hogwild! when both are run serially, demonstrating that the gains from (temporal) cache locality outweigh the coordination overhead of Cyclades. These results are reflected in the speedups of Cyclades and Hogwild! (Figs. 4(c) and 4(d)). Cyclades consistently achieves a better speedup (9−10×9-10\times on 1818 threads) compared to that of Hogwild! (7−9×7-9\times on 1818 threads).

4 Runtime Breakdown

The cost of partitioning and allocation for Cyclades is given in Table 3, relatively to the time that Hogwild! takes to complete one epoch of stochastic updates (i.e., a single pass over the dataset). For matrix completion and the graph eigenvector problem, on 18 threads, Cyclades takes the equivalent of 4-6 epochs of Hogwild! to complete its partitioning, as the problem is either very sparse or the updates are expensive. For solving least squares using SAGA and word embeddings using SGD, the cost of partitioning is equivalent to 11-14 epochs of Hogwild! on 18 threads. However, we point out that partitioning and allocation is a one-time cost which becomes cheaper with more stochastic update epochs. Additionally, we note that this cost can become amortized quickly due to the extra experiments one has to run for hyperparameter tuning, since the graph partitioning is identical across different stepsizes one might want to test.

Stochastic updates running time

When performing stochastic updates, Cyclades has better cache locality and coherence, but requires synchronization after each batch. Table 4 shows the time for each method to complete a single pass over the dataset, only with respect to stochastic updates (i.e., here we factor out the partitioning time). In most cases, Cyclades is faster than Hogwild!. In the cases where Cyclades is not faster, the overheads of synchronizing outweigh the gains from better cache locality and coherency. However, in some of these cases, synchronization can help by preventing errors due to asynchrony that lead to worse convergence, thus allowing Cyclades to use larger stepsizes and maximize convergence speed.

5 Diminishing stepsizes

In the previous experiments we used constant stepsizes. Here, we investigate the behavior of Cyclades and Hogwild! in the regime where we decrease the stepsize after each epoch. In particular, we ran the matrix completion experiments with SGD (with and without regularization), where we multiplicatively updated the stepsize by 0.95 after each epoch. The convergence and speedup plots are given in Figure 6. Cyclades is able to achieve a speedup of up to 6−7×6-7\times on 16−1816-18 threads. On the other hand, Hogwild! is performing worse comparatively to its performance with constant stepsizes (Figure 4(c)). The difference is more significant on regularized SGD, where we have to perform lazy updates (Appendix A.1), and Hogwild! fails to achieve the same optimum as Cyclades with multiple threads. Thus, on 18 threads, Hogwild! obtains a maximum speedup of 3×3\times, whereas Cyclades attains a speedup of 6.5×6.5\times.

6 Binary Classification and Dense Coordinates

In addition to the above experiments, here we explore settings where Cyclades is expected to perform poorly due to the inherent density of updates (i.e., for data sets with dense features). In particular, we test Cyclades on a classification problem for text based data, where a few features appear in most data points. Specifically, we run classification for the URL dataset [MSSV09] contains ∼\sim 2.4M URLs, labeled as either benign or malicious, and 3.2M features, including bag-of-words representation of tokens in the URL.

For this classification task, we used a logistic regression model, trained using SGD. By its power-law nature, the dataset consists of a small number of extremely dense features which occur in nearly all updates. Since Cyclades explicitly avoids all conflicts, for these dense cases it will have a schedule of SGD updates that leads to poor speedups. However, we observe that most conflicts are caused by a small percentage of the densest features. If these features are removed from the dataset, Cyclades is able to obtain much better speedups. To that end, we ran Cyclades and Hogwild! after filtering the densest 0.016%0.016\% to 0.048%0.048\% of features. The number of features that are filtered are shown in Table 7.

The speedups that are obtained by Cyclades and Hogwild! on 16 threads for different filtering percentages are shown in Figure 8. Full results of the experiment are presented in Figure 9 and in App. F. Cyclades fails to get much speedup when nearly all the features are used, however, as more dense features are removed, Cyclades obtains a better speedup, almost equalling Hogwild!’s speedup when 0.048%0.048\% of the densest features are filtered.

Related work

Parallel stochastic optimization has been studied under various guises, with literature stretching back at least to the late 60s [CM69]. The end of Moore’s Law coupled with recent advances in parallel and distributed computing technologies have triggered renewed interest [ZLS09, ZWLS10, GNHS11, AD11, RT12, JST+14] in the theory and practice of this field. Much of this contemporary work is built upon the foundational work of Bertsekas, Tsitsiklis et al. [BT89, TBA86].

[NRRW11] introduced Hogwild!, a completely lock-free and asynchronous parallel stochastic gradient descent (SGD), in shared-memory multicore systems. Inspired by Hogwild!’s success at achieving nearly linear speedups for a variety of machine learning tasks, several authors developed other lock-free and asynchronous optimization algorithms, such as parallel stochastic coordinate descent [LWR+14, LW15]. Additional work in first order optimization and beyond [DJM13, WSD+14, Hon14, HYD15, FAJ15], extending to parallel iterative linear solvers [LWS14, ADG14], has further demonstrated that linear speedups are generically possible in the asynchronous shared-memory setting. Moreover, [RHS+15] proposes an analysis for asynchronous parallel, but dense, SVRG, under assumptions similar to those found in [NRRW11]. The authors of [DSZOR15] offer a new analysis for the “coordinate-wise” update version of Hogwild! using martingales, with similar assumptions to [NRRW11], that however can be applied to some non-convex problems. Furthermore, [LHLL15] presents an analysis for stochastic gradient methods on smooth, potentially nonconvex functions. Finally, [PXYY15] introduces a new framework for analyzing coordinate-wise fixed point stochastic iterations.

Recently, [MPP+15] provided a new analysis for asynchronous stochastic optimization by interpreting the effects of asynchrony as noise in the iterates. This perspective of asycnhrony as noise in the iterates was used in [PPO+15] to analyze a combinatorial graph clustering algorithm.

Parallel optimization algorithms have also been developed for specific subclasses of problems, including L1-regularized loss minimization [BKBG11] and matrix factorization [RR13].

Other machine learning algorithms that have been parallelized using coordination and concurrency control, including non-parametric clustering [PGJ+13], submodular maximization [PJG+14], and correlation clustering [PPO+15].

Sparse, graph-based parallel computation are supported by systems like GraphLab [LGK+14] and PowerGraph [GLG+12]. These frameworks require computation to be written in a specific programming model with associative, commutative operations. GraphLab and PowerGraph support serializable execution via locking mechanisms This is in contrast to our partition-and-allocate coordination which allows us to provide guarantees on speedup.

Parameter servers [HCC+13, LAP+14] are frameworks supporting distributed, asynchronous computation. Convergence is only guaranteed for some specific algorithms, namely SGD, but not in general as execution is not serializable. The focus of parameter servers is in the distributed setting, where the challenges are different than in the multicore setting.

Conclusion

We presented Cyclades, a general framework for lock-free parallelization of stochastic optimization algorithms, while maintaining serial equivalence. Our framework can be used to parallelize a large family of stochastic updates algorithms in a conflict-free manner, thereby ensuring the parallelized algorithm produces the same result as its serial counterpart. Theoretical properties, such as convergence rates, are therefore preserved by the Cyclades-parallelized algorithm, and we provide a single unified theoretical analysis that guarantees near linear speedups.

By eliminating conflicts across processors within each batch of updates, Cyclades is able to avoid all asynchrony errors and conflicts, and leads to better cache locality and cache coherence than Hogwild!. These features of Cyclades translate to near linear speedups in practice, where it can outperform Hogwild!-type of implementations by up to a factor of 5×5\times, in terms of speedups

In the future, we intend to explore hybrids of Cyclades with Hogwild!, pushing the boundaries of what is possible in a shared-memory setting. We are also considering solutions for scaling out in a distributed setting, where the cost of communication is significantly higher.

Acknowledgements

BR is generously supported by ONR awards N00014-14-1-0024, N00014-15-1-2620, and N00014-13-1-0129, and NSF awards CCF-1148243 and CCF-1217058. XP’s work is also supported by a DSO National Laboratories Postgraduate Scholarship. This research is supported in part by NSF CISE Expeditions Award CCF-1139158, LBNL Award 7076018, and DARPA XData Award FA8750-12-2-0331, and gifts from Amazon Web Services, Google, SAP, The Thomas and Stacey Siebel Foundation, Adatao, Adobe, Apple, Inc., Blue Goji, Bosch, C3Energy, Cisco, Cray, Cloudera, EMC2, Ericsson, Facebook, Guavus, HP, Huawei, Informatica, Intel, Microsoft, NetApp, Pivotal, Samsung, Schlumberger, Splunk, Virdata and VMware.

References

Appendix A Algorithms in the Stochastic Updates family

Here we show that several algorithms belong to the Stochastic Updates (SU) family. These include the well-known stochastic gradient descent, iterative linear solvers, stochastic PCA and others, as well as combinations of weight decay updates, variance reduction methods, and more. Interestingly, even some combinatorial graph algorithms fit under the SU umbrella, such as the maximal independent set, greedy correlation clustering, and others. We visit some of these algorithms below.

Given nn functions f1,…,fnf_{1},\ldots,f_{n}, one often wishes to minimize the average of these functions:

A popular algorithm to do so —even in the case of non-convex losses— is the stochastic gradient descent:

In this case, the distribution D\mathcal{D} for each sample iki_{k} is usually a with or without replacement uniform sampling among the nn functions. For this algorithm the conflict graph between the nn possible different updates is completely determined by the support of the gradient vector ∇fik(xk)\nabla f_{i_{k}}({\bf x}_{k}).

Weight decay and regularization

In this case, the update is a weighted version of the “non-regularized" SGD:

The above stochastic update algorithm can be also be written in the SU language. Although here for each single update the entire model has to be updated with the new weights, we show below that with a simple technique it can be equivalently expressed so that each update is sparse and the support is again determined by the gradient vector ∇fik(xk)\nabla f_{i_{k}}({\bf x}_{k}).

First order techniques with variance reduction

Variance reduction is a technique that is usually employed for (strongly) convex problems, where we wish to minimize the variance of SGD in order to achieve better rates. A popular way to do variance reduction is either through SVRG or SAGA, where a “memory” factor is used in the computation of each gradient update rule. For SAGA we have

where g=∇f(y){\bf g}=\nabla f({\bf y}), and y{\bf y} is updated every TT iterations of the previous form to be equal to the last xk{\bf x}_{k} iterate. Again, although at first sight the above updates seem to be dense, we show below how we can equivalently rewrite them so that the update-conflict graph is completely determined by the support of the gradient vector ∇fik(xk)\nabla f_{i_{k}}({\bf x}_{k}).

Combinatorial graph algorithms

Interestingly, even some combinatorial graph problems can be phrased in the SU language: greedy correlation clustering (The Pivot Algorithm [ACN08]) and the maximal independent set (these are in fact identical algorithms). In the case of correlation clustering, we are given a graph GG vertices joined with either positive or negative edges. Here the objective is to create a number of clusters so that the number of vertex pairs that are sharing negative edges within clusters, plus the number of pairs that are sharing positive edges across clusters, is minimized. For these cases, there exists a very simple algorithm that obtains a 33 approximation for the above objective: randomly sample a vertex, create a cluster with that vertex and its neighborhood, remove that cluster from the graph, and repeat. The above procedure is amenable to the following update rule:

where x{\bf x} is intialized to ∞\infty, and at each iteration we sample vv uniformly among those with label xv=∞x_{v}=\infty, and N(v)N(v) denotes the neighborhood of a vertex in GG. Interestingly, we can directly apply the same guarantees of the main Cyclades theorem here. An optimized implementation of Cyclades for correlation clustering was developed in [PPO+15].

To reiterate, all of the above algorithms, various combinations of them, and further extensions can be written in the language of SU, as presented in Alg. 2.

A.1 Lazy Updates

For the cases of weight decay/regularization, and variance reduction, we can reinterpret their inherently dense updates in an equivalent sparse form. Let us consider the following generic form of updates:

where hij(xSi)=0h_{ij}({\bf x}_{\mathcal{S}_{i}})=0 for all j∉Sij\not\in\mathcal{S}_{i}. Each stochastic update therefore reads from the set Si\mathcal{S}_{i} but writes to every coordinate. However, it is possible to make updates lazily only when they are required. Observe that if τj\tau_{j} updates are made, each of which have hij(xSi)=0h_{ij}({\bf x}_{\mathcal{S}_{i}})=0, then we could rewrite these τj\tau_{j} updates in closed form as

This allows the stochastic updates to only write to coordinates in Si\mathcal{S}_{i} and defer writes to other coordinates. This procedure is described in Algorithm 3. With Cyclades it is easy to keep track of τj\tau_{j}, since we know the serially equivalent order of each stochastic update. On the other hand, it is unclear how a Hogwild! approach would behave with additional noise in τj\tau_{j} due to asynchrony. In fact, Hogwild! could possibly result in negative values of τj\tau_{j}, and in practice, we find that it is often useful to threshold τj\tau_{j} by max⁡(0,τj)\max(0,\tau_{j}).

The weighted decay SGD update is a special case of Eq 1, with μj=ηγ\mu_{j}=\eta\gamma and νj=0\nu_{j}=0. Eq 3 becomes xj←(1−ηγ)τjxjx_{j}\leftarrow(1-\eta\gamma)^{\tau_{j}}x_{j}.

Variance reduction with sparse gradients

Suppose ∇fi(x)\nabla f_{i}({\bf x}) is sparse, such that [∇fi(x)]j=0[\nabla f_{i}({\bf x})]_{j}=0 for all x\bf x and j∉Sij\not\in\mathcal{S}_{i}. Then we can perform SVRG and SAGA using lazy updates, with μj=0\mu_{j}=0. Just-in-time updates for SAG (a similar algorithm to SAGA) were introduced in [SRB13]. For SAGA, the update Eq 2 becomes

where gj=[1n∑i=1nyk,i]jg_{j}=\left[\frac{1}{n}\sum_{i=1}^{n}\mathbf{y}_{k,i}\right]_{j} is the jjth coordinate of the average gradient. For SVRG, we instead use gj=[∇f(y)]jg_{j}=[\nabla f(\mathbf{y})]_{j}.

SVRG with dense linear gradients

Appendix B With and Without Replacement Proofs

In this Appendix, we show how the sampling and shattering Theorem 2 can be restated for sampling with, or without replacement to establish Theorem 3.

Let us define three sequences of binary random variables {Xi}i=1n,{Yi}i=1n,\{X_{i}\}_{i=1}^{n},\{Y_{i}\}_{i=1}^{n}, and {Zi}i=1n\{Z_{i}\}_{i=1}^{n}. {Xi}i=1n\{X_{i}\}_{i=1}^{n} consists of nn i.i.d. Bernoulli random variables, each with probability pp. In the second sequence {Yi}i=1n\{Y_{i}\}_{i=1}^{n}, a random subset of BB random variables is set to 11 without replacement. Finally, in the third sequence {Zi}i=1n\{Z_{i}\}_{i=1}^{n}, we draw BB variables with replacement, and we set them to 11. Here, BB is integer that satisfies the following bounds

Now, consider any function ff, that has a “monotonicity" property:

for some number CC, and let us further assume that we have an upper bound on the above probability

Our goal is to bound ρY\rho_{Y} and ρZ\rho_{Z}. By expanding ρX\rho_{X} using the law of total probability we have

where qb=Pr⁡(f(X1,…,Xn)>C∣∑i=1nXi=b)q_{b}=\Pr\left(f(X_{1},\ldots,X_{n})>C\left|\sum_{i=1}^{n}X_{i}=b\right.\right), denotes the probability that f(X1,…,Xn)>Cf(X_{1},\ldots,X_{n})>C given that a uniformly random subset of bb variables was set to 11. Moreover, we have

where (i)(i) comes form the fact that Pr⁡(f(Y1,…,Yn)>C∣∑i=1nYi=b)\Pr\left(f(Y_{1},\ldots,Y_{n})>C\left|\sum_{i=1}^{n}Y_{i}=b\right.\right) is the same as the probability that that f(X1,…,Xn)>Cf(X_{1},\ldots,X_{n})>C given that a uniformly random subset of bb variables where set to 11, and (ii)(ii) comes from the fact that since we sample without replacement in YY, we have that ∑inYi=B\sum_{i}^{n}Y_{i}=B always.

In the expansion of ρX\rho_{X}, we can keep the b=Bb=B term, and lower bound the probability to obtain:

since all terms in the sum are non-negative numbers. Moreover, since XiX_{i}s are Bernoulli random variables, their sum ∑i=1nXi\sum_{i=1}^{n}X_{i} is Binomially distributed with parameters nn and pp. We know that the maximum of the Binomial pmf with parameters nn and pp occurs at Pr⁡(∑iXi=B)\Pr\left(\sum_{i}X_{i}=B\right) where BB is the integer that satisfies the upper bound mentioned above: (n+1)⋅p−1≤B<(n+1)⋅p(n+1)\cdot p-1\leq B<(n+1)\cdot p. Furthermore, the maximum value of the Binomial pmf, with parameters nn and pp, cannot be less than the corresponding probability of a uniform element:

The above establish a relation between the without replacement sampling sequence {Yi}i=1n\{Y_{i}\}_{i=1}^{n}, and the i.i.d. uniform sampling sequence {Xi}i=1n\{X_{i}\}_{i=1}^{n}.

Then, for the last sequence {Zi}i=1n\{Z_{i}\}_{i=1}^{n} we have

where (i)(i) comes from the fact that Pr⁡(∑i=1nZi=b)\Pr\left(\sum_{i=1}^{n}Z_{i}=b\right) is zero for b=0b=0 and b>Bb>B, (ii)(ii) comes by applying Hölder’s Inequality, and (iii)(iii) holds since ff is assumed to have the monotonicity property:

for any sequence of variables x1,…,xnx_{1},\ldots,x_{n}. Hence, for any b1≥b2b_{1}\geq b_{2}

In conclusion, we have upper bounded ρZ\rho_{Z} and ρY\rho_{Y} by

Application to Theorem 3: For our purposes, the above bound Eq. (10) allows us to assert Theorem 3 for with replacement, without replacement, and i.i.d. sampling, with different constants. Specifically, for any graph GG, the size of the largest connected component in the sampled subgraph can be expressed as a function fG(x1,…,xn)f_{G}(x_{1},\dots,x_{n}), where each xix_{i} is an indicator for whether the iith vertex was chosen in the sampling process. Note that fGf_{G} is a monotone function, i.e., fG(x1,…,xi,…,xn)≥fG(x1,…,0,…,xn)f_{G}(x_{1},\dots,x_{i},\dots,x_{n})\geq f_{G}(x_{1},\dots,0,\dots,x_{n}) since adding vertices to the sampled subgraph may only increase (or keep constant) the size of the largest connected component. We note that the high probability statement of Theorem 2, can be restated so that the constants in front of the size of the connected components accomodate for a statement that is true with probability 1−1/nζ1-1/n^{\zeta}, for any constant ζ>1\zeta>1. This is required to take care of the extra nn factor that appears in the bound of Eq. 10, and to obtain Theorem 3.

Appendix C Parallel Connected Components Computation

As we will see in the following, the cost of computing CCs in parallel will depend on the number of cores so that uniform allocation across them is possible, and the number of edges that are induced by the sampled updates on the bipartite update-variable graph GuG_{u} is bounded. As a reminder we denote as GuiG_{u}^{i} the bipartite subgraphs of the update-variable graph GuG_{u}, that is induced by the sampled updates of the ii-th batch. Let us denote as EuiE_{u}^{i} the number of edges in GuiG_{u}^{i}.

Following the sampling recipe of our main text (i.e., sampling each update per batch uniformly and with replacement), let us assume here that we are sampling c⋅nc\cdot n updates in total, for some constant c≥1c\geq 1. Assuming that the size of each batch is B=(1−ϵ)nΔB=(1-\epsilon)\frac{n}{\Delta}, the total number of sampled batches will be nb=c1−ϵΔn_{b}=\frac{c}{1-\epsilon}\Delta. The total number of edges in the induced sampled bipartite graphs is a random variable that we denote as

where ΔL\Delta_{\text{L}} is the max left degree of the bipartite graph GuG_{u} and Δ‾L\overline{\Delta}_{\text{L}} is its average left degree. Now assuming that

Hence, we get the following simple lemma:

Let ΔLΔ‾L≤n\frac{\Delta_{\text{L}}}{\overline{\Delta}_{\text{L}}}\leq\sqrt{n}. Then, the total number of edges Z=∑i=1nbEuiZ=\sum_{i=1}^{n_{b}}E_{u}^{i} across the nb=c1−ϵΔn_{b}=\frac{c}{1-\epsilon}\Delta sampled subgraphs Gu1,…,GunbG_{u}^{1},\ldots,G_{u}^{n_{b}} is less than O(Eulog⁡n)O(E_{u}\log n) with probability 1−n−Ω(log⁡n)1-n^{-\Omega(\log n)}.

Now that we have a bound on the number of edges in the sampled subgraphs, we can derive the complexity bounds for computing CCs in parallel. We will break the discussion into the not-too-many- and many-core regime.

In this case, we sample nbn_{b} subgraphs, allocate them across PP cores, and let each core compute CCs on its allocated subgraphs using BFS or DFS. Since each batch is of size B=(1−ϵ)nΔB=(1-\epsilon)\frac{n}{\Delta}, we need nb=⌊\nicefracc⋅nB⌋=c⋅⌊Δ1−ϵ⌋n_{b}=\lfloor\nicefrac{{c\cdot n}}{{B}}\rfloor=c\cdot\lfloor\frac{\Delta}{1-\epsilon}\rfloor batches to cover c⋅nc\cdot n updates in total. If the number of cores is

then the max cost of a single CC computation on a subgraph (which is O(BΔL)O(B\Delta_{\text{L}})) is smaller than the average cost across PP cores, which is O(Z/P)O(Z/P). This implies that a uniform allocation is possible, so that PP cores can share the computational effort of computing CCs. Hence, we can get the following lemma:

Let the number of cores be P=O(Δ‾LΔL⋅Δ)P=O\left(\frac{\overline{\Delta}_{\text{L}}}{\Delta_{\text{L}}}\cdot\Delta\right), and let us sample O(Δ)O(\Delta) batches, where each batch is of size O(nΔ)O(\frac{n}{\Delta}). Then, each core will not spend more than O(Eulog⁡nP)O\left(\frac{E_{u}\log n}{P}\right) time in computing CCs, with high probability.

The many-cores regime.

When P>>Δ‾LΔLP>>\frac{\overline{\Delta}_{\text{L}}}{\Delta_{\text{L}}} the uniform balancing of the above method will break, leaving no room for further parallelization. In that case, we can use a very simple “push-label” CC algorithm, whose cost on PP cores and arbitrary graphs with EE edges and max degree Δ\Delta is O(max⁡{EP,Δ}⋅Cmax⁡)O(\max\{\frac{E}{P},\Delta\}\cdot C_{\max}), where Cmax⁡C_{\max} is the size of the longest-shortest path, or the diameter of the graph [KTF09]. This parallel CC algorithm is given below, where each core is allocated with a number of vertices

The above simple algorithm can be significantly slow for graphs where the longest-shortest path is large. Observe, that in the sampled subgraphs GuiG_{u}^{i} the size of the shortest-longest path is always bounded by the size of the largest connected component. By Theorem 3 that is bounded by O(log⁡nϵ2)O\left(\frac{\log n}{\epsilon^{2}}\right). Hence, we obtain the following lemma.

For any number of cores P=O(nΔ⋅ΔL)P=O(\frac{n}{\Delta\cdot\Delta_{\text{L}}}), computing the connected component of a single sampled graph GuiG_{u}^{i} can be performed in time O(Euilog⁡nP)\mathcal{O}(\frac{E^{i}_{u}\log n}{P}), with high probability.

Since, we are interested in the overall running time for nbn_{b} batches of the CC algorithm, we can see that the above lemma simply boils down to the following:

For any number of cores P=O(nΔ⋅ΔL)P=O(\frac{n}{\Delta\cdot\Delta_{\text{L}}}), computing the connected component for all sampled graph Gu1,…,GunbG_{u}^{1},\ldots,G_{u}^{n_{b}} can be performed in time O(Elog⁡2nP)\mathcal{O}(\frac{E\log^{2}n}{P}).

In practice it seems to be that parallelizing the CC computation using the not-too-many core regime policy is significantly more scalable.

Appendix D Allocating the Conflict Groups

After we have sampled a single batch (i.e., a subgraph GuiG_{u}^{i}), and computed the CCs for it, we have to allocate the connected components of that sampled subgraph across cores. Observe that each connected component will contain at most log⁡n\log n updates, each ordered according to the a serial predetermined order. Once a core has been assigned all the CCs, it will process all the updates included in the CCs according to the order that each update has been labeled with.

Now assuming that the cost of the ii-th update is wiw_{i}, the cost of a single connected component C\mathcal{C} will be wC=∑i∈Cwiw_{\mathcal{C}}=\sum_{i\in\mathcal{C}}w_{i}. We can now allocate the CCs accross cores so that the maximum core load is minimized, using the following 4/34/3-approximation algorithm (i.e., an allocation with max load that is at most 4/34/3 times the maximum between the max weight, and the sum of weights divided by PP):

To proceed with characterizing the maximum load among the PP cores, we assume that the cost of a single update UiU_{i} is proportional to the out-degree of that update —according to the update-variable graph GuG_{u}— times a constant cost which we shall refer to as κ\kappa. Hence, wi=O(dL,i⋅κ)w_{i}=O(d_{L,i}\cdot\kappa), where dL,id_{L,i} is the degree of the ii-th left vertex of GuG_{u}.

Observe that the total cost of computing the updates in a single sampled subgraph GuiG_{u}^{i} is proportional to O(Eui⋅κ)\mathcal{O}(E_{u}^{i}\cdot\kappa). Moreover, observe that the maximum weight among all CCs cannot be more than O(ΔLlog⁡nκ)O(\Delta_{L}\log n\kappa) where ΔL\Delta_{L} is the max left degree of the bipartite update-variable graph GuG_{u}.

We can allocate CCs such that the maximum load among cores is O(max⁡{EuiP,ΔLlog⁡n}⋅κ),\mathcal{O}\left(\max\left\{\frac{E^{i}_{u}}{P},\Delta_{L}\log n\right\}\cdot\kappa\right), with high probability, where κ\kappa is the per edge cost for computing one of the nn updates defined on GuG_{u}.

If P=O(nΔ⋅ΔL)P=O\left(\frac{n}{\Delta\cdot\Delta_{\text{L}}}\right) then the average weight will be larger than the maximum divided by a log⁡n\log n factor, and a near-uniform allocation of CCs according to their weights possible. Since, we are interested in the overall running time for nbn_{b} batches, we can see that the above lemma simply boils down to the following:

For any number of cores P=O(nΔ⋅ΔL)P=O(\frac{n}{\Delta\cdot\Delta_{\text{L}}}), computing the stochastic updates of the allocated connected component for all sampled graphs (i.e., batches) Gu1,…,GunbG_{u}^{1},\ldots,G_{u}^{n_{b}} can be performed in time O(Elog⁡2nP⋅κ)\mathcal{O}(\frac{E\log^{2}n}{P}\cdot\kappa).

Appendix E Robustness against High-degree Outliers

Here, we discuss how Cyclades can guarantee nearly linear speedups when there is a sublinear O(nδ)O(n^{\delta}) number of high-conflict updates, as long as the remaining updates have small degree.

Assume that our conflict graph GcG_{c} defined between the nn update functions has a very high maximum degree Δo\Delta_{o}. However, consider the case where there are only O(nδ)O(n^{\delta}) nodes that are of that high-degree, while the rest of the vertices have degree much smaller (on the induced subgraph by the latter vertices), say Δ\Delta. According to our main analysis, our prescribed batch sizes cannot be greater than B=(1−ϵ)(1−ϵ)nΔoB=(1-\epsilon)\frac{(1-\epsilon)n}{\Delta_{o}}. However, if say Δo=Θ(n)\Delta_{o}=\Theta(n), then that would imply that B=O(1)B=O(1), hence there is not room for parallelization by Cyclades. What we will show, is that by sampling according to B=(1−ϵ)n−O(nδ)ΔB=(1-\epsilon)\frac{n-O(n^{\delta})}{\Delta}, we can on average expect a parallelization that is similar to the case where the outliers are not present in the conflict graph. For a toy example see Figure 10.

Our main result for the outlier case follows:

Let us assume that there are O(nδ)O(n^{\delta}) outlier vertices in the original conflict graph GG with degree at most Δo\Delta_{o}, and let the remaining vertices have degree (induced on the remaining graph) at most Δ\Delta. Let the induced update-variable graph on these low degree vertices abide to the same graph assumptions as those of Theorem 4. Moreover, let the batch size be bounded as

Then, the expected runtime of Cyclades will be O(Eu⋅κP⋅log⁡2n).O\left(\frac{E_{u}\cdot\kappa}{P}\cdot\log^{2}n\right).

Let wsiw^{i}_{s} denote the total work required for batch ii if that batch contains no outlier notes, and woiw^{i}_{o} otherwise. It is not hard to see that ws=∑iwsi=O(Eu⋅κP⋅log⁡2n)w_{s}=\sum_{i}w^{i}_{s}=O\left(\frac{E_{u}\cdot\kappa}{P}\cdot\log^{2}n\right) and wo=∑iwsi=O(Eu⋅κ⋅log⁡2n)w_{o}=\sum_{i}w^{i}_{s}=O\left(E_{u}\cdot\kappa\cdot\log^{2}n\right) Hence, the expected computational effort by Cyclades will be

Hence the expected running time will be proportional to O(Eu⋅κP⋅log⁡2n)O\left(\frac{E_{u}\cdot\kappa}{P}\cdot\log^{2}n\right), if O(Eu⋅κP⋅log⁡2n)=O(Eu⋅κ⋅log⁡2n⋅Bn1−δ),O(\frac{E_{u}\cdot\kappa}{P}\cdot\log^{2}n)=O(E_{u}\cdot\kappa\cdot\log^{2}n\cdot\frac{B}{n^{1-\delta}}), which holds when B=O(n1−δP).B=O\left(\frac{n^{1-\delta}}{P}\right). ∎

Appendix F Complete Experiment Results

In this section, we present the remaining experimental results that were left out for brevity from our main experimental section. In Figures 11 and 12, we show the convergence behaviour of our algorithms, as a function of the overall time, and then as a function of the time that it takes to perform only the stochastic updates (i.e., in Fig. 12 we factor out the graph partitioning, and allocation time). In Figure 13, we provide the complete set of speedup results for all algorithms and data sets we tested, in terms of the number of cores. In Figure 14, we provide the speedups in terms of the the computation of the stochastic updates, as a function of the number of cores. Then, in Figures 15 – 18, we present the convergence, and speedups of the overal computation, and then of the stochastic updates part, for our dense feature URL data set. Finally, in Figure 19 we show the divergent behavior of Hogwild! for the least square experiments with SAGA, on the NH2010 and DBLP datasets.

Our overall observations here are similar to the main text. One additional note to make is that when we take a closer look to the figures relative to the times and speedups of the stochastic updates part of Cyclades (i.e., when we factor out the time of the graph partitioning part), we see that Cyclades is able to perform stochastic updates faster than Hogwild! due to its superior spatial and temporal access locality patterns. If the coordination overheads for Cyclades are excluded, we are able to improve speedups, in some cases by up to 20-70% (Table 5). This suggests that by further optimizing the computation of connected components, we can hope for better overall speedups of Cyclades.