Can Decentralized Algorithms Outperform Centralized Algorithms? A Case Study for Decentralized Parallel Stochastic Gradient Descent

Xiangru LianCe ZhangHuan ZhangCho-Jui HsiehWei ZhangJi Liu

article2017NeurIPS1,542 citations

Demonstrates theoretically and empirically that decentralized parallel stochastic gradient descent can outperform centralized distributed training by up to an order of magnitude by eliminating central-node communication bottlenecks on low-bandwidth networks.

Listen

Distributed machine learning relies heavily on parallel stochastic gradient descent algorithms to train models across multiple computing nodes. The standard industry approach relies on centralized architectures, such as parameter server frameworks or collective communication primitives. However, centralized systems suffer from severe network traffic jams at the central node or high synchronization overhead, creating substantial performance bottlenecks when network bandwidth is limited or network latency is high.

The article evaluates whether decentralized parallel stochastic gradient descent, where nodes communicate only with direct logical neighbors rather than a central coordinator, can outperform standard centralized architectures. Specifically, the authors analyze the theoretical convergence rate and computational complexity of decentralized algorithms and conduct empirical evaluations across diverse computing frameworks and network conditions.

The research evaluated algorithm performance through theoretical proofs and empirical benchmarks using vision models (such as 20-layer, 32-layer, and 56-layer residual networks on the CIFAR-10 dataset) and industry natural language processing tasks. Experiments were conducted using Microsoft CNTK and Torch across setups ranging from 7 to 112 graphics processing units (GPUs), testing various physical configurations and synthetic network constraints in bandwidth and latency.

The analysis produced several key findings. First, decentralized parallel stochastic gradient descent achieves the exact same theoretical computational complexity and linear speedup as centralized methods, requiring only a fraction of the per-node communication cost. Second, on networks with low bandwidth or high latency, the decentralized algorithm ran up to ten times faster than well-optimized centralized alternatives while converging to equivalent loss levels. Third, in distributed multi-machine tests over standard gigabit Ethernet, the decentralized approach outperformed elasticity-based centralized baselines and reduced communication overhead by two to three times on real-world natural language processing workloads. Finally, the decentralized model exhibited no penalty in generalization ability, matching or exceeding published residual network accuracy benchmarks.

These findings demonstrate that decentralized training is not merely a fallback for restricted network topologies, but a superior architecture for distributed machine learning in bandwidth-constrained and latency-sensitive environments. Adopting decentralized coordination can substantially decrease cloud infrastructure costs, eliminate the requirement for costly high-speed interconnects in certain training pipelines, and accelerate machine learning development cycles without degrading model quality.

Engineering and infrastructure teams should consider adopting decentralized parallel gradient descent frameworks for distributed model training, especially when deploying workloads across standard Ethernet clusters, distributed data centers, or edge devices. For organizations running large GPU clusters, decentralized communication can serve as an effective bridge between distributed sub-clusters to eliminate parameter server bottlenecks.

The study's primary limitation is its reliance on synchronous iterations, which require all nodes to finish local computation before exchanging updates with neighbors; performance can degrade if node computing speeds vary significantly. The findings are strongly supported by rigorous mathematical proofs and multi-framework testing up to 112 GPUs. Future work should focus on developing asynchronous decentralized variants and validating the approach on extreme-scale supercomputers and mobile edge networks.

Cover for Can Decentralized Algorithms Outperform Centralized Algorithms? A Case Study for Decentralized Parallel Stochastic Gradient Descent

Abstract

Most distributed machine learning systems nowadays, including TensorFlow and CNTK, are built in a centralized fashion. One bottleneck of centralized algorithms lies on high communication cost on the central node. Motivated by this, we ask, can decentralized algorithms be faster than its centralized counterpart?

Although decentralized PSGD (D-PSGD) algorithms have been studied by the control community, existing analysis and theory do not show any advantage over centralized PSGD (C-PSGD) algorithms, simply assuming the application scenario where only the decentralized network is available. In this paper, we study a D-PSGD algorithm and provide the first theoretical analysis that indicates a regime in which decentralized algorithms might outperform centralized algorithms for distributed stochastic gradient descent. This is because D-PSGD has comparable total computational complexities to C-PSGD but requires much less communication cost on the busiest node. We further conduct an empirical study to validate our theoretical analysis across multiple frameworks (CNTK and Torch), different network configurations, and computation platforms up to 112 GPUs. On network configurations with low bandwidth or high latency, D-PSGD can be up to one order of magnitude faster than its well-optimized centralized counterparts.

Table of Contents

  • 1 Introduction
  • 2 Related work
  • 3 Decentralized parallel stochastic gradient descent (D-PSGD)
  • 4 Convergence rate analysis
  • 5 Experiments
  • 5.1 Experiment setting
  • 5.2 Results on CNTK
  • 5.3 Results on Torch
  • 6 Conclusion
  • References
  • Supplemental Materials: More Experiments
  • Industrial benchmark
  • Supplemental Materials: Proofs

Knowls

  1. Knowl 1 — Asymptotic Linear Speedup of Decentralized Parallel SGD

    theoretical result

    Under standard smoothness and variance assumptions, Decentralized Parallel Stochastic Gradient Descent (D-PSGD) achieves an asymptotic linear speedup with respect to the number of compute nodes nn when optimizing a general, possibly non-convex objective f(x):=1n∑i=1nfi(x)f(x) := \frac{1}{n} \sum_{i=1}^n f_i(x) on an undirected network.

    Setting the stepsize to γ=12L+σK/n,\gamma = \frac{1}{2L + \sigma \sqrt{K/n}}, where LL is the Lipschitz constant of the gradients ∇fi(x)\nabla f_i(x), σ2\sigma^2 bounds the stochastic gradient variance on each node, and KK is the total iteration count, the average gradient norm of the iterates across all nn nodes satisfies 1K∑k=0K−1E[∥∇f(1n∑i=1nxk,i)∥2]≤8(f(0)−f∗)LK+(8(f(0)−f∗)+4L)σKn,\frac{1}{K} \sum_{k=0}^{K-1} \mathbb{E} \left[ \left\| \nabla f \left( \frac{1}{n} \sum_{i=1}^n x_{k,i} \right) \right\|^2 \right] \le \frac{8(f(0) - f^*)L}{K} + \frac{(8(f(0) - f^*) + 4L)\sigma}{\sqrt{Kn}}, provided the total number of iterations KK is sufficiently large, specifically satisfying K>4L4n5σ6(f(0)−f∗+L)2(σ21−ρ+9ς2(1−ρ)2)2andK>72L2n2σ2(1−ρ)2,K > \frac{4L^4 n^5}{\sigma^6 (f(0) - f^* + L)^2} \left( \frac{\sigma^2}{1 - \rho} + \frac{9\varsigma^2}{(1 - \sqrt{\rho})^2} \right)^2 \quad \text{and} \quad K > \frac{72L^2 n^2}{\sigma^2 (1 - \sqrt{\rho})^2}, where ρ:=(max⁡{∣λ2(W)∣,∣λn(W)∣})2<1\rho := (\max\{|\lambda_2(W)|, |\lambda_n(W)|\})^2 < 1 represents the spectral gap of the symmetric doubly stochastic mixing matrix WW, f∗f^* is the infimum of ff, and ς2\varsigma^2 bounds the gradient variance across different local node functions (with ς=0\varsigma = 0 if all nodes sample from identical distributions).

    As K→∞K \to \infty, the O(1/Kn)\mathcal{O}(1/\sqrt{Kn}) term dominates the O(1/K)\mathcal{O}(1/K) term, matching the convergence rate of centralized mini-batch SGD with aggregate batch size proportional to nn.

  2. Knowl 2 — Decentralized Parallel Stochastic Gradient Descent (D-PSGD) Algorithm

    algorithm

    Decentralized Parallel Stochastic Gradient Descent (D-PSGD) optimizes a distributed objective f(x)=1n∑i=1nEξ∼Di[Fi(x;ξ)]f(x) = \frac{1}{n} \sum_{i=1}^n \mathbb{E}_{\xi \sim \mathcal{D}_i}[F_i(x; \xi)] over nn interconnected computational nodes without a central parameter server. Communication is structured by an undirected graph parameterized by a symmetric doubly stochastic weight matrix W∈Rn×nW \in \mathbb{R}^{n \times n}, where Wij∈[0,1]W_{ij} \in [0, 1] specifies the mixing weight between node ii and neighbor jj, and Wij=0W_{ij} = 0 if nodes ii and jj are disconnected.

    Input: Initial iterate x0∈RNx_{0} \in \mathbb{R}^N (set x0,i=x0x_{0,i} = x_0 on every node i∈{1,…,n}i \in \{1, \ldots, n\}), stepsize γ>0\gamma > 0, symmetric doubly stochastic matrix W∈Rn×nW \in \mathbb{R}^{n \times n}, total iterations KK
    Output: Approximate consensus solution xˉK=1n∑i=1nxK,i\bar{x}_K = \frac{1}{n} \sum_{i=1}^n x_{K,i}
    for k=0,1,…,K−1k = 0, 1, \ldots, K-1 do
        for each node i∈{1,…,n}i \in \{1, \ldots, n\} concurrently do
            Sample mini-batch / data point ξk,i\xi_{k,i} from local distribution Di\mathcal{D}_i
            Compute local stochastic gradient gk,i=∇Fi(xk,i;ξk,i)g_{k,i} = \nabla F_i(x_{k,i}; \xi_{k,i})
            Fetch iterates xk,jx_{k,j} from neighbors {j:Wij>0}\{j : W_{ij} > 0\} and compute neighborhood average:
                xk+1/2,i=∑j=1nWijxk,jx_{k+1/2,i} = \sum_{j=1}^n W_{ij} x_{k,j}
            Update local model parameter:
                xk+1,i=xk+1/2,i−γgk,ix_{k+1,i} = x_{k+1/2,i} - \gamma g_{k,i}
        end for
    end for
    return 1n∑i=1nxK,i\frac{1}{n} \sum_{i=1}^n x_{K,i}

    The local gradient computation and neighborhood parameter communication can execute concurrently in parallel threads to overlap communication latency with computation time.

  3. Knowl 3 — Non-Convex Convergence Bound for D-PSGD

    theoretical result

    Let Xk=[xk,1,…,xk,n]∈RN×nX_k = [x_{k,1}, \ldots, x_{k,n}] \in \mathbb{R}^{N \times n} be the concatenation of local parameters across nn workers at iteration kk, and let ∂f(Xk)=[∇f1(xk,1),…,∇fn(xk,n)]∈RN×n\partial f(X_k) = [\nabla f_1(x_{k,1}), \ldots, \nabla f_n(x_{k,n})] \in \mathbb{R}^{N \times n}. Under LL-Lipschitz continuous local gradients, symmetric doubly stochastic mixing matrix WW with spectral parameter ρ=(max⁡{∣λ2(W)∣,∣λn(W)∣})2<1\rho = (\max\{|\lambda_2(W)|, |\lambda_n(W)|\})^2 < 1, bounded intra-node variance σ2\sigma^2, and cross-node variance ς2\varsigma^2, the iterates of D-PSGD satisfy:

    1K(1−γL2∑k=0K−1E[∥∂f(Xk)1nn∥2]+D1∑k=0K−1E[∥∇f(Xk1nn)∥2])≤f(0)−f∗γK+γLσ22n+γ2L2nσ2(1−ρ)D2+9γ2L2nς2(1−ρ)2D2,\frac{1}{K} \left( \frac{1 - \gamma L}{2} \sum_{k=0}^{K-1} \mathbb{E}\left[\left\| \frac{\partial f(X_k) \mathbf{1}_n}{n} \right\|^2\right] + D_1 \sum_{k=0}^{K-1} \mathbb{E}\left[\left\| \nabla f\left(\frac{X_k \mathbf{1}_n}{n}\right) \right\|^2\right] \right) \le \frac{f(0) - f^*}{\gamma K} + \frac{\gamma L \sigma^2}{2n} + \frac{\gamma^2 L^2 n \sigma^2}{(1 - \rho) D_2} + \frac{9\gamma^2 L^2 n \varsigma^2}{(1 - \sqrt{\rho})^2 D_2},

    where the auxiliary constants D1D_1 and D2D_2 are defined as D1:=12−9γ2L2n(1−ρ)2D2,D2:=1−18γ2nL2(1−ρ)2.D_1 := \frac{1}{2} - \frac{9\gamma^2 L^2 n}{(1 - \sqrt{\rho})^2 D_2}, \qquad D_2 := 1 - \frac{18\gamma^2 n L^2}{(1 - \sqrt{\rho})^2}.

    When stepsize γ\gamma satisfies γ2≤(1−ρ)236nL2\gamma^2 \le \frac{(1 - \sqrt{\rho})^2}{36 n L^2} and γ2≤(1−ρ)272L2n\gamma^2 \le \frac{(1 - \sqrt{\rho})^2}{72 L^2 n}, the constants satisfy D2≥1/2D_2 \ge 1/2 and D1≥1/4D_1 \ge 1/4.

  4. Knowl 4 — Complexity and Communication Comparison Between C-PSGD and D-PSGD

    data/table

    Centralized Parallel SGD (C-PSGD via a parameter server) and Decentralized Parallel SGD (D-PSGD) exhibit identical asymptotic computational complexities to obtain an ϵ\epsilon-approximate stationary point (where 1K∑k=0K−1E[∥∇f(xˉk)∥2]≤ϵ\frac{1}{K}\sum_{k=0}^{K-1} \mathbb{E}[\|\nabla f(\bar{x}_k)\|^2] \le \epsilon), but D-PSGD substantially reduces the communication burden on the busiest network node.

    Algorithm Communication Complexity on Busiest Node Computational Complexity
    C-PSGD (mini-batch SGD) O(n)O(n) O(nϵ+1ϵ2)O\left(\frac{n}{\epsilon} + \frac{1}{\epsilon^2}\right)
    D-PSGD O(Deg(network))O(\text{Deg}(\text{network})) O(nϵ+1ϵ2)O\left(\frac{n}{\epsilon} + \frac{1}{\epsilon^2}\right)

    In this comparison, nn denotes the number of compute nodes. The communication cost is measured in units of transmitted optimization vectors or gradients per iteration. While a centralized parameter server suffers from an O(n)O(n) communication bottleneck due to concurrent synchronization with all workers, D-PSGD bounds the per-iteration communication cost of each node to O(Deg(network))O(\text{Deg}(\text{network})), which simplifies to O(1)O(1) on topologies such as rings or sparse regular graphs.

  5. Knowl 5 — Consensus Error Rate in D-PSGD

    theoretical result

    Under the standard smoothness and variance assumptions with learning rate γ=12L+σK/n\gamma = \frac{1}{2L + \sigma \sqrt{K/n}}, the average variance between the local model parameters xk,ix_{k,i} across all nn nodes and their global network average xˉk=1n∑j=1nxk,j\bar{x}_k = \frac{1}{n} \sum_{j=1}^n x_{k,j} converges to zero at an O(1/K)\mathcal{O}(1/K) rate:

    1Kn∑k=0K−1∑i=1nE[∥1n∑j=1nxk,j−xk,i∥2]≤nγ2AD2=O(1K),\frac{1}{K n} \sum_{k=0}^{K-1} \sum_{i=1}^n \mathbb{E} \left[ \left\| \frac{1}{n} \sum_{j=1}^n x_{k,j} - x_{k,i} \right\|^2 \right] \le \frac{n \gamma^2 A}{D_2} = \mathcal{O}\left(\frac{1}{K}\right),

    where AA is an explicit constant given by A:=2σ21−ρ+18ς2(1−ρ)2+L2D1(σ21−ρ+9ς2(1−ρ)2)+18(1−ρ)2(f(0)−f∗γK+γLσ22nD1),A := \frac{2\sigma^2}{1 - \rho} + \frac{18\varsigma^2}{(1 - \sqrt{\rho})^2} + \frac{L^2}{D_1} \left( \frac{\sigma^2}{1 - \rho} + \frac{9\varsigma^2}{(1 - \sqrt{\rho})^2} \right) + \frac{18}{(1 - \sqrt{\rho})^2} \left( \frac{f(0) - f^*}{\gamma K} + \frac{\gamma L \sigma^2}{2 n D_1} \right), and D1,D2D_1, D_2 are positive constants. This ensures that every individual worker's local parameter vector tracks the global consensus iterate asymptotically.

  6. Knowl 6 — Maximum Scalable Worker Regimes for Ring Network Topology

    theoretical result

    For a ring network topology where each node communicates exclusively with its two immediate neighbors, the symmetric doubly stochastic mixing matrix W∈Rn×nW \in \mathbb{R}^{n \times n} has entries Wi,i=Wi,i−1=Wi,i+1=1/3W_{i,i} = W_{i,i-1} = W_{i,i+1} = 1/3 (with periodic boundary conditions). As the node count nn grows large, the spectral properties satisfy ρ≈1−16π23n2\rho \approx 1 - \frac{16\pi^2}{3n^2} and ρ≈1−8π23n2\sqrt{\rho} \approx 1 - \frac{8\pi^2}{3n^2}.

    Under stepsize γ=12L+σK/n\gamma = \frac{1}{2L + \sigma \sqrt{K/n}}, D-PSGD achieves linear speedup (O(1/Kn)O(1/\sqrt{Kn}) convergence rate) provided the number of participating nodes nn scales within the following limits:

    • n=O(K1/9)n = \mathcal{O}(K^{1/9}) under Strategy 1 (all nodes sample uniformly from a shared global database, so ς=0\varsigma = 0).
    • n=O(K1/13)n = \mathcal{O}(K^{1/13}) under Strategy 2 (nodes partition data locally and draw from heterogeneous local distributions, so ς>0\varsigma > 0).
  7. Knowl 7 — Analytical Assumptions for Decentralized Parallel SGD

    assumption

    The theoretical convergence analysis of D-PSGD relies on the following four assumptions on the objective function and communication topology:

    1. Lipschitzian Gradients: Each local objective fi(x):=Eξ∼Di[Fi(x;ξ)]f_i(x) := \mathbb{E}_{\xi \sim \mathcal{D}_i}[F_i(x; \xi)] is continuously differentiable with LL-Lipschitz continuous gradient: ∥∇fi(x)−∇fi(y)∥≤L∥x−y∥,  ∀x,y∈RN,  ∀i∈{1,…,n}\|\nabla f_i(x) - \nabla f_i(y)\| \le L \|x - y\|, \; \forall x, y \in \mathbb{R}^N, \; \forall i \in \{1, \ldots, n\}.
    2. Spectral Gap: The communication weight matrix W∈Rn×nW \in \mathbb{R}^{n \times n} is symmetric and doubly stochastic (i.e., Wij∈[0,1]W_{ij} \in [0, 1], W=W⊤W = W^\top, and W1n=1nW \mathbf{1}_n = \mathbf{1}_n). The spectral gap parameter satisfies ρ:=(max⁡{∣λ2(W)∣,∣λn(W)∣})2<1\rho := (\max\{|\lambda_2(W)|, |\lambda_n(W)|\})^2 < 1.
    3. Bounded Variance: The intra-node stochastic gradient variance and cross-node distribution discrepancy are bounded by constants σ2\sigma^2 and ς2\varsigma^2: Eξ∼Di[∥∇Fi(x;ξ)−∇fi(x)∥2]≤σ2,∀i,∀x,\mathbb{E}_{\xi \sim \mathcal{D}_i} \left[ \|\nabla F_i(x; \xi) - \nabla f_i(x)\|^2 \right] \le \sigma^2, \quad \forall i, \forall x, Ei∼U([n])[∥∇fi(x)−∇f(x)∥2]≤ς2,∀x.\mathbb{E}_{i \sim \mathcal{U}([n])} \left[ \|\nabla f_i(x) - \nabla f(x)\|^2 \right] \le \varsigma^2, \quad \forall x. If all nodes access the same shared database (homogeneous data), ς=0\varsigma = 0.
    4. Zero Initialization: Iterates start from X0=[x0,1,…,x0,n]=0X_0 = [x_{0,1}, \ldots, x_{0,n}] = 0 without loss of generality.
  8. Knowl 8 — Empirical Speedup of D-PSGD in Bandwidth-Constrained and High-Latency Networks

    empirical result

    In distributed experiments on clusters scaling up to 112 GPUs using Microsoft CNTK training ResNet-20 and ResNet-56 on CIFAR-10, D-PSGD achieved wall-clock speed improvements up to 10×10\times compared to centralized parameter server and MPI AllReduce baselines under low bandwidth or high network latency.

    Key observations:

    • Per-epoch time under network constraints: When network bandwidth is constrained (e.g., dialed down via the tc command from 100 Mbps to 1 Mbps) or network latency is increased (from 0 ms up to 10 ms), centralized parameter server SGD slows down due to bandwidth saturation at the root node, and MPI AllReduce slows down due to latency scaling with communication frequency. D-PSGD maintains flat, low per-epoch times because communication is strictly localized and pairwise.
    • Epoch-to-accuracy invariance: In terms of convergence per epoch, D-PSGD matched the convergence trajectory of centralized SGD across 7, 10, 16, and 112 GPUs, confirming that decentralized communication does not degrade statistical efficiency.
  9. Knowl 9 — Empirical Comparison Between D-PSGD and Elastic Averaging SGD (EASGD)

    empirical result

    When evaluated on a 9-node cluster (8 worker GPUs and 1 parameter server) connected via commodity 1 Gbps Gigabit Ethernet training ResNet-32 on CIFAR-10 in Torch, D-PSGD with a logical ring topology outperformed Elastic Averaging SGD (EASGD / EAMSGD with momentum):

    • Wall-Clock Speed and Saturation: EASGD with communication period τ=1\tau = 1 converged rapidly per epoch but severely saturated the 1 Gbps Ethernet, causing large wall-clock delays. When communication frequency was reduced (τ=4,16\tau = 4, 16), EASGD avoided network saturation but degraded in convergence speed per epoch. D-PSGD bypassed the bandwidth bottleneck while communicating every iteration, reaching target training loss (0.20.2) significantly faster in wall-clock time.
    • Scaling and Generalization: D-PSGD achieved near-linear scaling when scaling across 1, 4, 8, and 16 machines (reaching training loss 0.20.2 in approximately 80, 20, 10, and 5 epochs, respectively, with only a 3% increase in per-epoch execution time over single-GPU execution). Final test error after 160 epochs was 0.07150.0715 (4 machines), 0.07460.0746 (8 machines), and 0.07350.0735 (16 machines), outperforming the baseline literature benchmark of 0.07510.0751 for ResNet-32.
  10. Knowl 10 — Synchronous Barrier Limitation of D-PSGD

    limitation

    The standard D-PSGD algorithm relies on a synchronous barrier at each iteration where each node exchanges iterates with its immediate graph neighbors. In heterogeneous environments where compute nodes possess differing hardware speeds or fluctuating network connection qualities, the execution time of each iteration is gated by the slowest neighbor (straggler effect).

    While partial mitigations include adjusting local mini-batch sizes proportionally or computing multiple local mini-batches per synchronization step on faster workers, fully resolving this limitation requires designing and analyzing asynchronous decentralized stochastic optimization protocols.

Coverage note — None was omitted; all primary theoretical proofs, algorithmic formulations, complexity comparisons, scalability analyses on graph topologies, empirical benchmarks across CNTK/Torch, and system limitations from the paper are represented.

References

  1. 1.A. Agarwal and J. C. Duchi. Distributed delayed stochastic optimization. NIPS, 2011.
  2. 2.N. S. Aybat, Z. Wang, T. Lin, and S. Ma. Distributed linearized alternating direction method of multipliers for composite convex consensus optimization. arXiv preprint arXiv:1512.08122, 2015.
  3. 3.T. C. Aysal, M. E. Yildiz, A. D. Sarwate, and A. Scaglione. Broadcast gossip algorithms for consensus. IEEE Transactions on Signal processing, 57(7):2748–2761, 2009.
  4. 4.P. Bianchi, G. Fort, and W. Hachem. Performance of a distributed stochastic approximation algorithm. IEEE Transactions on Information Theory, 59(11):7405–7418, 2013.
  5. 5.S. Boyd, A. Ghosh, B. Prabhakar, and D. Shah. Gossip algorithms: Design, analysis and applications. In INFOCOM 2005. 24th Annual Joint Conference of the IEEE Computer and Communications Societies. Proceedings IEEE, volume 3, pages 1653–1664. IEEE, 2005.
  6. 6.R. Carli, F. Fagnani, P. Frasca, and S. Zampieri. Gossip consensus algorithms via quantized communication. Automatica, 46(1):70–80, 2010.
  7. 7.J. Chen, R. Monga, S. Bengio, and R. Jozefowicz. Revisiting distributed synchronous sgd. arXiv preprint arXiv:1604.00981, 2016.
  8. 8.K. Crammer, O. Dekel, J. Keshet, S. Shalev-Shwartz, and Y. Singer. Online passive-aggressive algorithms. Journal of Machine Learning Research, 7:551–585, 2006.
  9. 9.J. Dean, G. Corrado, R. Monga, K. Chen, M. Devin, M. Mao, A. Senior, P. Tucker, K. Yang, Q. V. Le, et al. Large scale distributed deep networks. In Advances in neural information processing systems, pages 1223–1231, 2012.
  10. 10.O. Dekel, R. Gilad-Bachrach, O. Shamir, and L. Xiao. Optimal distributed online prediction using mini-batches. Journal of Machine Learning Research, 13(Jan):165–202, 2012.
  11. 11.F. Fagnani and S. Zampieri. Randomized consensus algorithms over large scale networks. IEEE Journal on Selected Areas in Communications, 26(4), 2008.
  12. 12.M. Feng, B. Xiang, and B. Zhou. Distributed deep learning for question answering. In Proceedings of the 25th ACM International on Conference on Information and Knowledge Management, pages 2413–2416. ACM, 2016.
  13. 13.H. R. Feyzmahdavian, A. Aytekin, and M. Johansson. An asynchronous mini-batch algorithm for regularized stochastic optimization. arXiv, 2015.
  14. 14.S. Ghadimi and G. Lan. Stochastic first-and zeroth-order methods for nonconvex stochastic programming. SIAM Journal on Optimization, 23(4):2341–2368, 2013.
  15. 15.K. He, X. Zhang, S. Ren, and J. Sun. Deep Residual Learning for Image Recognition. ArXiv e-prints, Dec. 2015.
  16. 16.K. He, X. Zhang, S. Ren, and J. Sun. Deep residual learning for image recognition. In Proceedings of the IEEE Conference on Computer Vision and Pattern Recognition, pages 770–778, 2016.
  17. 17.A. Krizhevsky. Learning multiple layers of features from tiny images. In Technical Report, 2009.
  18. 18.G. Lan, S. Lee, and Y. Zhou. Communication-efficient algorithms for decentralized and stochastic optimization. arXiv preprint arXiv:1701.03961, 2017.
  19. 19.Y. LeCun, Y. Bengio, and G. Hinton. Deep learning. Nature, 521(7553):436–444, 2015.
  20. 20.M. Li, D. G. Andersen, J. W. Park, A. J. Smola, A. Ahmed, V. Josifovski, J. Long, E. J. Shekita, and B.-Y. Su. Scaling distributed machine learning with the parameter server. In OSDI, volume 14, pages 583–598, 2014.
  21. 21.X. Lian, Y. Huang, Y. Li, and J. Liu. Asynchronous parallel stochastic gradient for nonconvex optimization. In Advances in Neural Information Processing Systems, pages 2737–2745, 2015.
  22. 22.X. Lian, H. Zhang, C.-J. Hsieh, Y. Huang, and J. Liu. A comprehensive linear speedup analysis for asynchronous stochastic parallel optimization from zeroth-order to first-order. In Advances in Neural Information Processing Systems, pages 3054–3062, 2016.
  23. 23.Z. Lin, M. Feng, C. N. d. Santos, M. Yu, B. Xiang, B. Zhou, and Y. Bengio. A structured self-attentive sentence embedding. 5th International Conference on Learning Representations, 2017.
  24. 24.J. Lu, C. Y. Tang, P. R. Regier, and T. D. Bow. A gossip algorithm for convex consensus optimization over networks. In American Control Conference (ACC), 2010, pages 301–308. IEEE, 2010.
  25. 25.A. Mokhtari and A. Ribeiro. Dsa: decentralized double stochastic averaging gradient algorithm. Journal of Machine Learning Research, 17(61):1–35, 2016.
  26. 26.E. Moulines and F. R. Bach. Non-asymptotic analysis of stochastic approximation algorithms for machine learning. NIPS, 2011.
  27. 27.A. Nedic and A. Ozdaglar. Distributed subgradient methods for multi-agent optimization. IEEE Transactions on Automatic Control, 54(1):48–61, 2009.
  28. 28.A. Nemirovski, A. Juditsky, G. Lan, and A. Shapiro. Robust stochastic approximation approach to stochastic programming. SIAM Journal on Optimization, 19(4):1574–1609, 2009.
  29. 29.Nvidia. Nccl: Optimized primitives for collective multi-gpu communication. https://github.com/NVIDIA/nccl.
  30. 30.R. Olfati-Saber, J. A. Fax, and R. M. Murray. Consensus and cooperation in networked multi-agent systems. Proceedings of the IEEE, 95(1):215–233, 2007.
  31. 31.S. S. Ram, A. Nedić, and V. V. Veeravalli. Asynchronous gossip algorithms for stochastic optimization. In Decision and Control, 2009 held jointly with the 2009 28th Chinese Control Conference. CDC/CCC 2009. Proceedings of the 48th IEEE Conference on, pages 3581–3586. IEEE, 2009a.
  32. 32.S. S. Ram, A. Nedic, and V. V. Veeravalli. Distributed subgradient projection algorithm for convex optimization. In Acoustics, Speech and Signal Processing, 2009. ICASSP 2009. IEEE International Conference on, pages 3653–3656. IEEE, 2009b.
  33. 33.S. S. Ram, A. Nedić, and V. V. Veeravalli. Asynchronous gossip algorithm for stochastic optimization: Constant stepsize analysis. In Recent Advances in Optimization and its Applications in Engineering, pages 51–60. Springer, 2010.
  34. 34.B. Recht, C. Re, S. Wright, and F. Niu. Hogwild: A lock-free approach to parallelizing stochastic gradient descent. In Advances in Neural Information Processing Systems, pages 693–701, 2011.
  35. 35.L. Schenato and G. Gamba. A distributed consensus protocol for clock synchronization in wireless sensor network. In Decision and Control, 2007 46th IEEE Conference on, pages 2289–2294. IEEE, 2007.
  36. 36.S. Shalev-Shwartz. Online learning and online convex optimization. Foundations and Trends in Machine Learning, 4(2):107–194, 2011.
  37. 37.W. Shi, Q. Ling, K. Yuan, G. Wu, and W. Yin. On the linear convergence of the admm in decentralized consensus optimization.
  38. 38.W. Shi, Q. Ling, G. Wu, and W. Yin. Extra: An exact first-order algorithm for decentralized consensus optimization. SIAM Journal on Optimization, 25(2):944–966, 2015.
  39. 39.B. Sirb and X. Ye. Consensus optimization with delayed and stochastic gradients on decentralized networks. In Big Data (Big Data), 2016 IEEE International Conference on, pages 76–85. IEEE, 2016.
  40. 40.K. Srivastava and A. Nedic. Distributed asynchronous constrained stochastic optimization. IEEE Journal of Selected Topics in Signal Processing, 5(4):772–790, 2011.
  41. 41.S. Sundhar Ram, A. Nedić, and V. Veeravalli. Distributed stochastic subgradient projection algorithms for convex optimization. Journal of optimization theory and applications, 147(3):516–545, 2010.
  42. 42.T. Wu, K. Yuan, Q. Ling, W. Yin, and A. H. Sayed. Decentralized consensus optimization with asynchrony and delays. arXiv preprint arXiv:1612.00150, 2016.
  43. 43.F. Yan, S. Sundaram, S. Vishwanathan, and Y. Qi. Distributed autonomous online learning: Regrets and intrinsic privacy-preserving properties. IEEE Transactions on Knowledge and Data Engineering, 25(11): 2483–2493, 2013.
  44. 44.T. Yang, M. Mahdavi, R. Jin, and S. Zhu. Regret bounded by gradual variation for online convex optimization. Machine learning, 95(2):183–223, 2014.
  45. 45.K. Yuan, Q. Ling, and W. Yin. On the convergence of decentralized gradient descent. SIAM Journal on Optimization, 26(3):1835–1854, 2016.
  46. 46.R. Zhang and J. Kwok. Asynchronous distributed admm for consensus optimization. In International Conference on Machine Learning, pages 1701–1709, 2014.
  47. 47.S. Zhang, A. E. Choromanska, and Y. LeCun. Deep learning with elastic averaging sgd. In Advances in Neural Information Processing Systems, pages 685–693, 2015.
  48. 48.Y. Zhuang, W.-S. Chin, Y.-C. Juan, and C.-J. Lin. A fast parallel sgd for matrix factorization in shared memory systems. In Proceedings of the 7th ACM conference on Recommender systems, pages 249–256. ACM, 2013.
  49. 49.M. Zinkevich, M. Weimer, L. Li, and A. J. Smola. Parallelized stochastic gradient descent. In Advances in neural information processing systems, pages 2595–2603, 2010.

Citation

MLA
Lian, X., et al. “Can Decentralized Algorithms Outperform Centralized Algorithms? A Case Study for Decentralized Parallel Stochastic Gradient Descent”. arXiv, 2017, http://arxiv.org/abs/1705.09056v5.
APA
Lian, X., Zhang, C., Zhang, H., Hsieh, C.-J., Zhang, W., & Liu, J. (2017). Can Decentralized Algorithms Outperform Centralized Algorithms? A Case Study for Decentralized Parallel Stochastic Gradient Descent. arXiv. http://arxiv.org/abs/1705.09056v5
Chicago
Lian, X., C. Zhang, H. Zhang, C.-J. Hsieh, W. Zhang, and J. Liu. 2017. “Can Decentralized Algorithms Outperform Centralized Algorithms? A Case Study for Decentralized Parallel Stochastic Gradient Descent”. arXiv. http://arxiv.org/abs/1705.09056v5.
Harvard
Lian, X. et al. (2017) “Can Decentralized Algorithms Outperform Centralized Algorithms? A Case Study for Decentralized Parallel Stochastic Gradient Descent”, arXiv [Preprint]. Available at: http://arxiv.org/abs/1705.09056v5.
Vancouver
1. Lian X, Zhang C, Zhang H, Hsieh C-J, Zhang W, Liu J (2017) Can Decentralized Algorithms Outperform Centralized Algorithms? A Case Study for Decentralized Parallel Stochastic Gradient Descent. arXiv

BibTeX

@article{lian2017can,
  title = {Can Decentralized Algorithms Outperform Centralized Algorithms? A Case Study for Decentralized Parallel Stochastic Gradient Descent},
  author = {Lian, Xiangru and Zhang, Ce and Zhang, Huan and Hsieh, Cho-Jui and Zhang, Wei and Liu, Ji},
  year = {2017},
  journal = {arXiv},
  url = {http://arxiv.org/abs/1705.09056v5},
  eprint = {1705.09056}
}
Metadata:arXiv

Access the Paper

This paper is available from its original source. Click below to access the PDF.

Open PDF
License: Authors