VLDB 2026 Research / reviewers in the wild / expert
Shaohuai Shi
dblp:79/8378
· DBLP profile ↗
47ranked-venue papers
13as first author
35since 2021 · last 2026
0000-0002-1418-5160ORCID · verified
Domains — the database's venue-derived domains; a paper can count in several
Systems, architecture and hardware · 25 · 7 first-author · 18 since 2021Computer networks · 12 · 4 first-author · 10 since 2021Artificial intelligence and machine learning · 11 · 2 first-author · 7 since 2021Graphics, computer vision, multimedia, augmented reality and games · 6 · 2 first-author · 4 since 2021Software engineering, systems software and programming languages · 2 · 2 since 2021Databases, data management, data science and information retrieval · 1Applied, interdisciplinary, general and emerging computing · 1
| Year | Publication | Venue | Position |
|---|---|---|---|
| 2026 | PipeDiT: Accelerating Diffusion Transformers in Video Generation with Task Pipelining and Model DecouplingabstractVideo generation has been advancing rapidly, and diffusion transformer (DiT) based models have demonstrated remarkable capabilities. However, their practical deployment is often hindered by slow inference speeds and high memory consumption. In this paper, we propose a novel pipelining framework named PipeDiT to accelerate video generation, which is equipped with three main innovations. First, we design a pipelining algorithm (PipeSP) for sequence parallelism (SP) to enable the computation of latent generation and communication among multiple GPUs to be pipelined, thus reducing the inference latency. Second, we propose DeDiVAE to decouple the diffusion module and the VAE module into two GPU groups whose executions can also be pipelined to reduce the memory consumption and inference latency. Third, to better utilize the GPU resources in the VAE group, we propose an attention co-processing (Aco) method to further reduce the overall video generation latency. We integrate our PipeDiT into both OpenSoraPlan and HunyuanVideo, two state-of-the-art open-source video generation frameworks, and conduct extensive experiments on two 8-GPU systems. Experimental results show that, under many common resolution and timestep configurations, our PipeDiT achieves 1.06× to 4.02× speedups over OpenSoraPlan and HunyuanVideo. Qiang Wang 0022, Shaohuai Shi |
AAAI | 3 |
| 2026 | SALR: Sparsity-Aware Low-Rank Representation for Efficient Fine-Tuning of Large Language ModelsabstractAdapting large pre-trained language models to downstream tasks often entails fine-tuning millions of parameters or deploying costly dense weight updates, which hinders their use in resource-constrained environments. Low-rank Adaptation (LoRA) reduces trainable parameters by factorizing weight updates, yet the underlying dense weights still impose high storage and computation costs. Magnitude-based pruning can yield sparse models but typically degrades LoRA’s performance when applied naively. In this paper, we introduce SALR (Sparsity-Aware Low-Rank Representation), a novel fine-tuning paradigm that unifies low-rank adaptation with sparse pruning under a rigorous mean-squared-error framework. We prove that statically pruning only the frozen base weights minimizes the pruning error bound, and we recover the discarded residual information via a truncated-SVD low-rank adapter, which provably reduces per-entry MSE by a factor of (1 - r/min(d, k)). To maximize hardware efficiency, we fuse multiple low-rank adapters into a single concatenated GEMM, and we adopt a bitmap-based encoding with a two-stage pipelined decoding + GEMM design to achieve true model compression and speedup. Empirically, SALR attains 50% sparsity on various LLMs while matching the performance of LoRA on GSM8K and MMLU, reduces model size by 2x, and delivers up to a 1.7x inference speedup. Longteng Zhang, Sen Wu 0001, Zhengyu Qing, Zhuo Zheng, Danning Ke, Qihong Lin, Qiang Wang 0022, Shaohuai Shi, Xiaowen Chu 0001 |
AAAI | 9 |
| 2026 | HierMoE: Accelerating MoE Training with Hierarchical Token Deduplication and Expert Swap
Wenxiang Lin, Xinglin Pan, Lin Zhang 0059, Shaohuai Shi, Xuan Wang 0002, Xiaowen Chu 0001 |
INFOCOM | 4 |
| 2026 | Compass: Dissecting Communication and Computation Operators for Efficient LLM TrainingabstractOverlapping communication and computation operators is a common practice to hide communication overheads, accelerating large language models (LLMs) training on GPU clusters. Existing systems achieve this through either intra-operator fusion (IntraFusion), which packs operators into a single large kernel, or inter-operator decomposition (InterDecom), which splits a tensor into multiple parts for pipelined execution. However, current IntraFusion methods underutilize network topology, causing suboptimal bandwidth usage on multi-GPU systems, while InterDecom struggles to determine the optimal number of decomposed parts for peak performance. To address these issues, we introduce Compass, which employs systematic optimization and comprehensive modeling. First, we design a novel IntraFusion algorithm leveraging double-ring communications to maximize bandwidth utilization in hybrid NVLink-PCIe systems, achieving 1.5x-2.5x speedups. Second, we develop a decomposition model that mathematically derives the optimal tensor decomposition degree for InterDecom, improving performance by up to 1.3x. Finally, we develop a unified performance framework that accurately determines the best strategy for different scenarios. We validate Compass through extensive evaluation across 288 configurations and end-to-end experiments on real-world applications. The results demonstrate that Compass consistently selects the optimal strategy, achieving up to a 1.42x end-to-end speedup compared to the Megatron-LM baseline. Guangyu Xiang, Lin Zhang 0059, Haoxuan Yu, Xinglin Pan, Shaohuai Shi, Xiaowen Chu 0001 |
INFOCOM | 5 |
| 2026 | Accelerating Multi-modal LLM Training with Adaptive Model Placement and Parallelization
Yiming Yin, Shaohuai Shi, Qiang Wang 0022, Xiaowen Chu 0001 |
INFOCOM | 2 |
| 2026 | ZipCCL: Efficient Lossless Data Compression of Communication Collectives for Accelerating LLM TrainingabstractCommunication has emerged as a critical bottleneck in the distributed training of large language models (LLMs). While numerous approaches have been proposed to reduce communication overhead, the potential of lossless compression has remained largely underexplored since compression and decompression typically consume larger overheads than the benefits of reduced communication traffic. We observe that the communication data, including activations, gradients and parameters, during training often follows a near-Gaussian distribution, which is a key feature for data compression. Thus, we introduce ZipCCL, a lossless compressed communication library of collectives for LLM training. ZipCCL is equipped with our novel techniques: (1) theoretically grounded exponent coding that exploits the Gaussian distribution of LLM tensors to accelerate compression without expensive online statistics, (2) GPU-optimized compression and decompression kernels that carefully design memory access patterns and pipeline using communication-aware data layout, and (3) adaptive communication strategies that dynamically switch collective operations based on workload patterns and system characteristics. Evaluated on a 64-GPU cluster using both mixture-of-experts and dense transformer models, ZipCCL reduces communication time by up to 1.35X and achieves end-to-end training speedups of up to 1.18X without any impact on model quality. Wenxiang Lin, Xinglin Pan, Ruibo Fan, Shaohuai Shi, Xiaowen Chu 0001 |
SIGCOMM | 4 |
| 2026 | Castor: Optimizing Deep Learning Job Scheduling in Multi-Tenant GPU Clusters via Intelligent ColocationabstractDeep learning (DL) has achieved significant success across a wide range of domains, prompting the widespread deployment of GPU clusters equipped with specialized accelerators to support high-performance training workloads. To minimize operational costs while maximizing resource utilization, efficient job scheduling in these clusters is essential. Although recent schedulers have improved cluster efficiency through periodic reallocation or selection of GPU resources, they still face challenges such as preemption and migration overheads, along with the risk of degrading model accuracy. Despite these limitations, the potential of GPUsharing remains largely underexplored. Few existing studies have systematically examined GPU sharing as a strategy to enhance resource utilization and reduce job queuing delays in multi-tenant DL clusters. Motivated by these insights, we propose a job scheduling model that enables multiple jobs to share the same set of GPUs without modifying their original training configurations. We introduce Castor, a simple yet efficient scheduling system, to achieve intelligent GPU colocation for multiple DL jobs. Castor intelligently selects job pairs for GPU sharing and determines runtime parameters (sub-batch size and scheduling time point) to optimize overall system performance while preserving the accuracy of DL convergence through gradient accumulation. Through a combination of physical DL workloads and trace-driven simulations across various configurations, we demonstrate that Castor reduces average job completion time by 26–52% compared to state-of-the-art preemptive DL schedulers, despite operating under a preemption-free policy. Furthermore, Castor effectively identifies optimal resource-sharing configurations, outperforming the baseline first-fit sharing policy (SJF-FFS) by up to 20% on large-scale workload traces. Yizhou Luo, Jiaxin Lai, Shaohuai Shi, Chen Chen 0067, Shuhan Qi, Jiajia Zhang 0001, Qiang Wang 0022 |
IEEE Trans. Cloud Comput. | 3 |
| 2025 | FSMoE: A Flexible and Scalable Training System for Sparse Mixture-of-Experts ModelsabstractRecent large language models (LLMs) have tended to leverage sparsity to reduce computations, employing the sparsely activated mixture-of-experts (MoE) technique. MoE introduces four modules, including token routing, token communication, expert computation, and expert parallelism, that impact model quality and training efficiency. To enable ver- satile usage of MoE models, we introduce FSMoE, a flexible training system optimizing task scheduling with three novel techniques: 1) Unified abstraction and online profiling of MoE modules for task scheduling across various MoE implementations. 2) Co-scheduling intra-node and inter-node communications with computations to minimize communication overheads. 3) To support near-optimal task scheduling, we design an adaptive gradient partitioning method for gradient aggregation and a schedule to adaptively pipeline communications and computations. We conduct extensive experiments with configured MoE layers and real-world MoE models on two GPU clusters. Experimental results show that 1) our FSMoE supports four popular types of MoE routing functions and is more efficient than existing implementations (with up to a 1.42× speedup), and 2) FSMoE outperforms the state-of-the-art MoE training systems (DeepSpeed-MoE and Tutel) by 1.18×-1.22× on 1458 MoE layers and 1.19×-3.01× on real-world MoE models based on GPT-2 and Mixtral using a popular routing function. In this work, we present a flexible training system named FSMoE to optimize task scheduling. To achieve this goal: 1) we design unified abstraction and online profiling of MoE modules across various MoE implementations, 2) we co-schedule intra-node and inter-node communications with computations to minimize communication overhead, and 3) we design an adaptive gradient partitioning method for gradient aggregation and a schedule to adaptively pipeline communications and computations. Experimental results on two clusters up to 48 GPUs show that our FSMoE outperforms the state-of-the-art MoE training systems (DeepSpeed-MoE and Tutel) with speedups of 1.18x-1.22x on 1458 customized MoE layers and 1.19x-3.01x on real-world MoE models based on GPT-2 and Mixtral. Xinglin Pan, Wenxiang Lin, Lin Zhang 0059, Shaohuai Shi, Zhenheng Tang, Rui Wang 0172, Bo Li 0001, Xiaowen Chu 0001 |
ASPLOS (1) | 4 |
| 2025 | ScheInfer: Efficient Inference of Large Language Models with Task Scheduling on Moderate GPUs
Wenxiang Lin, Xinglin Pan, Shaohuai Shi, Xuan Wang 0002, Xiaowen Chu 0001 |
Euro-Par (3) | 3 |
| 2025 | SQ-DeAR: Sparsified and Quantized Gradient Compression for Distributed Training
Shaohuai Shi |
Euro-Par (3) | 2 |
| 2025 | DeepFill: Accelerating MLLM Training by Filling Bubbles with Frozen Encoders
Zhengyu Qing, Shaohuai Shi, Qiang Wang 0022 |
ICA3PP (6) | 2 |
| 2025 | Mast: Efficient Training of Mixture-of-Experts Transformers with Task Pipelining and OrderingabstractThe utilization of the sparsely activated mixture-of-experts (MoE) technique has enabled the expansion of modern large language models (LLMs) to trillion-level sizes while maintaining a sub-linear increase in computations. This involves equipping an MoE layer with multiple experts, where only one or two experts are activated for each input data. However, the dynamic activation of MoE experts introduces extensive communications, limiting the scaling efficiency of distributed systems. In this work, we propose Mast to efficiently train MoE models by pipelining and re-ordering communication and computation tasks to effectively hide communication costs. Specifically, we first propose to overlap tasks in both attention layers and MoE layers. Then we theoretically analyze the task overlaps between communications and computations, identifying the inefficiencies of existing schedules. We then develop an optimization formulation to determine a near-optimal order for task pipelining with the objective of minimizing iteration time. We conduct extensive experiments on two 32-GPU clusters employing 432 configured MoE layers and three real-world MoE models based on BERT, GPT-2 and Mistral. The experimental results demonstrate that Mast outperforms state-of-the-art MoE training systems (DeepSpeed-MoE, Tutel, PipeMoE and CoCoNet) with an average speedup 1.13 ×-1.43 × on the MoE models. Wenxiang Lin, Xinglin Pan, Shaohuai Shi, Xuan Wang 0002, Bo Li 0001, Xiaowen Chu 0001 |
ICDCS | 3 |
| 2025 | Mitigating Contention in Stream Multiprocessors for Pipelined Mixture of Experts: An SM-Aware Scheduling ApproachabstractSparsely activated Mixture-of-Experts (MoEs) models have become prominent in Large Language Models (LLMs) due to their ability to expand model capacity without proportional increases in computation. MoE layers feature multiple experts, with only a few activated per sample, enhancing model performance across various domains such as natural language generation and translation. The dynamic activation of MoE experts introduces extensive communications in distributed training. However, this dynamic activation creates communication challenges in distributed training. While previous work attempted to pipeline computation and communication through input chunking, we found that these tasks compete for Stream Multiprocessors (SMs) on GPUs, making the scheduling ineffective. In this paper, we update the optimization problem to minimize training time while accounting for SM contentions. We develop performance models for computation and communication tasks to identify MoE layer bottlenecks. By delaying GEMM launching and splitting GEMM operations, we enable communication to preempt SMs, enhancing overall efficiency and more stability. Xinglin Pan, Rui Wang 0172, Wenxiang Lin, Shaohuai Shi, Xiaowen Chu 0001 |
ICDCS | 4 |
| 2025 | QPO: Accelerating Memory-Efficient DNN Training with Quantization and PipeliningabstractDeep neural network (DNN) training demands significant computational power and also necessitates costly device memory for the storage of parameters, gradients, optimizer states, and activations. Activations generally take up the largest portion of device memory, and their usage increases linearly with the mini-batch size and sequence length, which are two main hyper-parameters in training language models. Offloading is one of the widely used memory-efficient techniques by transferring temporarily data from the device memory to CPU memory to save device memory. However, the offloading and uploading operations easily cause significant data transfer overheads. In this paper, we propose QPO (Quantized and Pipelined Offloading), which combines 1) compressing the activation data with low-bit quantization to alleviate the transfer overhead, and 2) pipelining communication tasks of offloading/uploading with computation tasks of feed-forward/backpropagation to reduce the iteration time. Experiments on GPT-2 and LLaMA-2 models using A100 and RTX 3090 GPUs show a speedup of up to 15% while approaching minimal memory requirements over existing offloading approaches. Shaohuai Shi |
ICPADS | 2 |
| 2025 | SP-MoE: Expediting Mixture-of-Experts Training with Optimized Pipelining PlanningabstractSparsely activated Mixture-of-Experts (MoE) has emerged as a key technique to expand the size of Transformer-based large language models (LLMs) while maintaining low computational costs. However, MoE layers require to route the input data to distributed devices, incurring significant communication latency. Existing studies have primarily focused on alleviating this problem by overlapping computation and communication tasks within a single MoE layer, which fails to achieve sufficient overlap and results in limited performance gains. In this work, we introduce an orthogonal partitioning dimension from existing task-parallel methods by leveraging the autoregressive nature of causal Transformer-based LLMs, i.e. partitioning tasks along the sequence dimension. This provides more flexible and efficient overlaps among tasks from both non-MoE and MoE layers. To this end, we propose an efficient MoE training approach, SP-MoE, with two innovative designs. 1) It incorporates non-MoE layers into the overlapping with not only the current MoE layer but also the preceding MoE layer, thereby facilitating more efficient training; 2) It identifies the optimal combination of pipeline degrees for non-MoE and MoE layers and devises the best scheduling plans for load-imbalanced non-MoE and uniform MoE layers to achieve the goal of minimizing the total training latency. Extensive experiments conducted on two GPU clusters demonstrate that SP-MoE can effectively identify the optimal combination of pipeline degrees and achieve 16.1% - 34.3% reduction in training latency compared to three state-of-the-art MoE systems. Ne Wang, Wenxiang Lin, Lin Zhang 0059, Shaohuai Shi, Ruiting Zhou, Bo Li 0001 |
INFOCOM | 4 |
| 2024 | Performance Analysis and Optimizations of Matrix Multiplications on ARMv8 ProcessorsabstractGeneral matrix multiplication (GEMM) as a fundamental subroutine has been widely used in many applications like scientific computing, machine learning, etc. Although many studies are dedicated to optimizing its performance, they mainly focus on matrices with regular shapes or x86 platforms. The irregularly shaped matrices on GEMM running on modern ARMv8 processors are under-explored. In this paper, we provide a thorough performance analysis of the general block-panel multiplication (GEBP) kernel of GEMM that has irregular shapes. Based on our analysis, we propose a new GEMM algorithm named EPPA with three novel schemes to improve GEMM performance on ARMv8 processors: i) eliminating packing to reduce Ll cache contention, ii) avoiding data eviction and pre-fetching data to reduce the Ll cache miss penalty, and iii) an adaptive selection strategy of the above two and original schemes. We conduct extensive experiments with a large range of irregular matrices on three popular ARMv8 processors compared to seven state-of-the-art GEMM libraries. The experimental results show that our EPPA algorithm outperforms existing ones across workloads and processors and accelerates real-world applications. Hucheng Liu, Shaohuai Shi, Xuan Wang 0002, Zoe Lin Jiang, Qian Chen 0028 |
DATE | 2 |
| 2024 | ScheMoE: An Extensible Mixture-of-Experts Distributed Training System with Tasks SchedulingabstractIn recent years, large-scale models can be easily scaled to trillions of parameters with sparsely activated mixture-of-experts (MoE), which significantly improves the model quality while only requiring a sub-linear increase in computational costs. However, MoE layers require the input data to be dynamically routed to a particular GPU for computing during distributed training. The highly dynamic property of data routing and high communication costs in MoE make the training system low scaling efficiency on GPU clusters. In this work, we propose an extensible and efficient MoE training system, ScheMoE, which is equipped with several features. 1) ScheMoE provides a generic scheduling framework that allows the communication and computation tasks in training MoE models to be scheduled in an optimal way. 2) ScheMoE integrates our proposed novel all-to-all collective which better utilizes intra- and inter-connect bandwidths. 3) ScheMoE supports easy extensions of customized all-to-all collectives and data compression approaches while enjoying our scheduling algorithm. Extensive experiments are conducted on a 32-GPU cluster and the results show that ScheMoE outperforms existing state-of-the-art MoE systems, Tutel and Faster-MoE, by 9%-30%. Shaohuai Shi, Xinglin Pan, Qiang Wang 0022, Chengjian Liu, Xiaozhe Ren, Zhongzhe Hu, Bo Li 0001, Xiaowen Chu 0001 |
EuroSys | 1 |
| 2024 | FedImpro: Measuring and Improving Client Update in Federated LearningabstractFederated Learning (FL) models often experience client drift caused by heterogeneous data, where the distribution of data differs across clients. To address this issue, advanced research primarily focuses on manipulating the existing gradients to achieve more consistent client models. In this paper, we present an alternative perspective on client drift and aim to mitigate it by generating improved local models. First, we analyze the generalization contribution of local training and conclude that this generalization contribution is bounded by the conditional Wasserstein distance between the data distribution of different clients. Then, we propose FedImpro, to construct similar conditional distributions for local training. Specifically, FedImpro decouples the model into high-level and low-level components, and trains the high-level portion on reconstructed feature distributions. This approach enhances the generalization contribution and reduces the dissimilarity of gradients in FL. Experimental results show that FedImpro can help FL defend against data heterogeneity and enhance the generalization performance of the model. Zhenheng Tang, Yonggang Zhang 0003, Shaohuai Shi, Xinmei Tian 0001, Tongliang Liu, Bo Han 0003, Xiaowen Chu 0001 |
ICLR | 3 |
| 2024 | Sparse Gradient Communication with AlltoAll for Accelerating Distributed Deep LearningabstractSynchronous stochastic gradient descent (S-SGD) with data parallelism has become a de-facto approach in training large-scale deep neural networks (DNNs) on multi-GPU systems. However, S-SGD requires iteratively synchronizing gradients from all workers, which incurs excessive communication costs and limits the scaling efficiency of GPU clusters. Gradient sparsification like top-k sparsification techniques has been shown to be potentially effective in reducing the communication volume. Yet, existing approaches with top-k sparsification still suffer from high communication complexity in that they often need a very low density to achieve better performance while easily sacrificing the model accuracy. To this end, we propose a novel sparse communication approach called TopKA2A, which integrates top-k sparsification with AlltoAll communication to exchange sparse tensors among GPUs, resulting in a notable reduction in communication complexity. With rigorous theoretical analysis on the conditions that TopKA2A can be applied, we design a simple yet effective tensor fusion algorithm based on binary search. We perform in-depth analysis and evaluation of communication efficiency by comparing TopKA2A with state-of-the-art solutions over popular models without compromising model accuracy. Experimental results demonstrate that TopKA2A yields substantial communication efficiency gains and runs up to 73% faster than existing algorithms on a 32-GPU cluster. Jing Peng 0003, Shaohuai Shi, Bo Li 0001 |
ICPP | 3 |
| 2024 | Bandwidth-Aware and Overlap-Weighted Compression for Communication-Efficient Federated LearningabstractCurrent data compression methods, such as sparsification in Federated Averaging (FedAvg), effectively enhance the communication efficiency of Federated Learning (FL). However, these methods encounter challenges such as the straggler problem and diminished model performance due to heterogeneous bandwidth and non-IID (Independently and Identically Distributed) data. To address these issues, we introduce a bandwidth-aware compression framework for FL, aimed at improving communication efficiency while mitigating the problems associated with non-IID data. First, our strategy dynamically adjusts compression ratios according to bandwidth, enabling clients to upload their models at a close pace, thus exploiting the otherwise wasted time to transmit more data. Second, we identify the non-overlapped pattern of retained parameters after compression, which results in diminished client update signals due to uniformly averaged weights. Based on this finding, we propose a parameter mask to adjust the client-averaging coefficients at the parameter level, thereby more closely approximating the original updates, and improving the training convergence under heterogeneous environments. Our evaluations reveal that our method significantly boosts model accuracy, with a maximum improvement of 13% over the uncompressed FedAvg. Moreover, it achieves a 3.37 × speedup in reaching the target accuracy compared to FedAvg with a Top-K compressor, demonstrating its effectiveness in accelerating convergence with compression. The integration of common compression techniques into our framework further establishes its potential as a versatile foundation for future cross-device, communication-efficient FL research, addressing critical challenges in FL and advancing the field of distributed machine learning. Zichen Tang, Rudan Yan, Yuxin Wang 0003, Zhenheng Tang, Shaohuai Shi, Amelie Chi Zhou, Xiaowen Chu 0001 |
ICPP | 6 |
| 2024 | Parm: Efficient Training of Large Sparsely-Activated Models with Dedicated SchedulesabstractSparsely-activated Mixture-of-Expert (MoE) layers have found practical applications in enlarging the model size of large-scale foundation models, with only a sub-linear increase in computation demands. Despite the wide adoption of hybrid parallel paradigms like model parallelism, expert parallelism, and expert-sharding parallelism (i.e., MP+EP+ESP) to support MoE model training on GPU clusters, the training efficiency is hindered by communication costs introduced by these parallel paradigms. To address this limitation, we propose Parm, a system that accelerates MP+EP+ESP training by designing two dedicated schedules for placing communication tasks. The proposed schedules eliminate redundant computations and communications and enable overlaps between intra-node and inter-node communications, ultimately reducing the overall training time. As the two schedules are not mutually exclusive, we provide comprehensive theoretical analyses and derive an automatic and accurate solution to determine which schedule should be applied in different scenarios. Experimental results on an 8-GPU server and a 32-GPU cluster demonstrate that Parm outperforms the state-of-the-art MoE training system, DeepSpeed-MoE, achieving 1.13× to 5.77× speedup on 1296 manually configured MoE layers and approximately 3× improvement on two real-world MoE models based on BERT and GPT-2. Xinglin Pan, Wenxiang Lin, Shaohuai Shi, Xiaowen Chu 0001, Weinong Sun, Bo Li 0001 |
INFOCOM | 3 |
| 2024 | Scheduling Deep Learning Jobs in Multi-Tenant GPU Clusters via Wise Resource SharingabstractDeep learning (DL) has demonstrated significant success across diverse fields, leading to the construction of dedicated GPU accelerators within GPU clusters for high-quality training services. Efficient scheduler designs for such clusters are vital to reduce operational costs and enhance resource utilization. While recent schedulers have shown impressive performance in optimizing DL job performance and cluster utilization through periodic reallocation or selection of GPU resources, they also encounter challenges such as preemption and migration overhead, along with potential DL accuracy degradation. Nonetheless, few explore the potential benefits of GPU sharing to improve resource utilization and reduce job queuing times.Motivated by these insights, we present a job scheduling model allowing multiple jobs to share the same set of GPUs without altering job training settings. We introduce SJF-BSBF (shortest job first with best sharing benefit first), a straightforward yet effective heuristic scheduling algorithm. SJF-BSBF intelligently selects job pairs for GPU resource sharing and runtime settings (sub-batch size and scheduling time point) to optimize overall performance while ensuring DL convergence accuracy through gradient accumulation. In experiments with both physical DL workloads and trace-driven simulations, even as a preemptionfree policy, SJF-BSBF reduces the average job completion time by 27-33% relative to the state-of-the-art preemptive DL schedulers. Moreover, SJF-BSBF can wisely determine the optimal resource sharing settings, such as the sharing time point and sub-batch size for gradient accumulation, outperforming the aggressive GPU sharing approach (baseline SJF-FFS policy) by up to 17% in large-scale traces. Yizhou Luo, Qiang Wang 0022, Shaohuai Shi, Jiaxin Lai, Shuhan Qi, Jiajia Zhang 0001, Xuan Wang 0002 |
IWQoS | 3 |
| 2023 | DeAR: Accelerating Distributed Deep Learning with Fine-Grained All-Reduce PipeliningabstractCommunication scheduling has been shown to be effective in accelerating distributed training, which enables all-reduce communications to be overlapped with backpropagation computations. This has been commonly adopted in popular distributed deep learning frameworks. However, there exist two fundamental problems: (1) excessive startup latency proportional to the number of workers for each all-reduce operation; (2) it only achieves sub-optimal training performance due to the dependency and synchronization requirement of the feed-forward computation in the next iteration. We propose a novel scheduling algorithm, DeAR, that decouples the all-reduce primitive into two continuous operations, which overlaps with both backpropagation and feed-forward computations without extra communications. We further design a practical tensor fusion algorithm to improve the training performance. Experimental results with five popular models show that DeAR achieves up to 83% and 15% training speedup over the state-of-the-art solutions on a 64-GPU cluster with 10Gb/s Ethernet and 100Gb/s InfiniBand interconnects, respectively. Lin Zhang 0059, Shaohuai Shi, Xiaowen Chu 0001, Wei Wang 0030, Bo Li 0001, Chengjian Liu |
ICDCS | 2 |
| 2023 | Evaluation and Optimization of Gradient Compression for Distributed Deep LearningabstractTo accelerate distributed training, many gradient compression methods have been proposed to alleviate the communication bottleneck in synchronous stochastic gradient descent (S-SGD), but their efficacy in real-world applications still remains unclear. In this work, we first evaluate the efficiency of three representative compression methods (quantization with Sign-SGD, sparsification with Top-k SGD, and low-rank with Power-SGD) on a 32-GPU cluster. The results show that they cannot always outperform well-optimized S-SGD or even worse due to their incompatibility with three key system optimization techniques (all-reduce, pipelining, and tensor fusion) in S-SGD. To this end, we propose a novel gradient compression method, called alternate compressed Power-SGD (ACP-SGD), which alternately compresses and communicates low-rank matrices. ACP-SGD not only significantly reduces the communication volume, but also enjoys the three system optimizations like S-SGD. Compared with Power-SGD, the optimized ACP-SGD can largely reduce the compression and communication overheads, while achieving similar model accuracy. In our experiments, ACP-SGD achieves an average of 4.06× and 1.43× speedups over S-SGD and Power-SGD, respectively, and it consistently outperforms other baselines across different setups (from 8 GPUs to 64 GPUs and from 1Gb/s Ethernet to 100Gb/s InfiniBand). Lin Zhang 0059, Longteng Zhang, Shaohuai Shi, Xiaowen Chu 0001, Bo Li 0001 |
ICDCS | 3 |
| 2023 | Eva: Practical Second-order Optimization with Kronecker-vectorized Approximation
Lin Zhang 0059, Shaohuai Shi, Bo Li 0001 |
ICLR | 2 |
| 2023 | PipeMoE: Accelerating Mixture-of-Experts through Adaptive PipeliningabstractLarge models have attracted much attention in the AI area. The sparsely activated mixture-of-experts (MoE) technique pushes the model size to a trillion-level with a sub-linear increase of computations as an MoE layer can be equipped with many separate experts, but only one or two experts need to be trained for each input data. However, the feature of dynamically activating experts of MoE introduces extensive communications in distributed training. In this work, we propose PipeMoE to adaptively pipeline the communications and computations in MoE to maximally hide the communication time. Specifically, we first identify the root reason why a higher pipeline degree does not always achieve better performance in training MoE models. Then we formulate an optimization problem that aims to minimize the training iteration time. To solve this problem, we build performance models for computation and communication tasks in MoE and develop an optimal solution to determine the pipeline degree such that the iteration time is minimal. We conduct extensive experiments with 174 typical MoE layers and two real-world NLP models on a 64-GPU cluster. Experimental results show that our PipeMoE almost always chooses the best pipeline degree and outperforms state-of-the-art MoE training systems by 5%-77% in training time. Shaohuai Shi, Xinglin Pan, Xiaowen Chu 0001, Bo Li 0001 |
INFOCOM | 1 |
| 2023 | Accelerating Distributed K-FAC with Efficient Collective Communication and SchedulingabstractDiDistributed training with synchronous stochastic gradient descent (SGD) on GPU clusters has been widely used to accelerate the training process of deep models. However, SGD only utilizes the first-order gradient in model parameter updates, which may take days or weeks. Recent studies have successfully exploited approximate second-order information to speed up the training process, in which the Kronecker-Factored Approximate Curvature (KFAC) emerges as one of the most efficient approximation algorithms for training deep models. Yet, when leveraging GPU clusters to train models with distributed KFAC (D-KFAC), it incurs extensive computation as well as introduces extra communications during each iteration. In this work, we propose D-KFAC (SPD-KFAC) with smart parallelism of computing and communication tasks to reduce the iteration time. Specifically, 1) we first characterize the performance bottlenecks of D-KFAC, 2) we design and implement a pipelining mechanism for Kronecker factors computation and communication with dynamic tensor fusion, and 3) we develop a load balancing placement for inverting multiple matrices on GPU clusters. We conduct real-world experiments on a 64-GPU cluster with 100Gb/s InfiniBand interconnect. Experimental results show that our proposed SPD-KFAC training scheme can achieve 10%-35% improvement over state-of-the-art algorithms. Lin Zhang 0059, Shaohuai Shi, Bo Li 0001 |
INFOCOM | 2 |
| 2023 | Scalable K-FAC Training for Deep Neural Networks With Distributed PreconditioningabstractThe second-order optimization methods, notably the D-KFAC (Distributed Kronecker Factored Approximate Curvature) algorithms, have gained traction on accelerating deep neural network (DNN) training on GPU clusters. However, existing D-KFAC algorithms require to compute and communicate a large volume of second-order information, i.e., Kronecker factors (KFs), before preconditioning gradients, resulting in large computation and communication overheads as well as a high memory footprint. In this paper, we propose DP-KFAC, a novel distributed preconditioning scheme that distributes the KF constructing tasks at different DNN layers to different workers. DP-KFAC not only retains the convergence property of the existing D-KFAC algorithms but also enables three benefits: reduced computation overhead in constructing KFs, no communication of KFs, and low memory footprint. Extensive experiments on a 64-GPU cluster show that DP-KFAC reduces the computation overhead by 1.55×-1.65×, the communication cost by 2.79×-3.15×, and the memory footprint by 1.14×-1.47× in each second-order update compared to the state-of-the-art D-KFAC methods. Our codes are available athttps://github.com/lzhangbv/kfac\_pytorch. Lin Zhang 0059, Shaohuai Shi, Wei Wang 0030, Bo Li 0001 |
IEEE Trans. Cloud Comput. | 2 |
| 2023 | GossipFL: A Decentralized Federated Learning Framework With Sparsified and Adaptive CommunicationabstractRecently, federated learning (FL) techniques have enabled multiple users to train machine learning models collaboratively without data sharing. However, existing FL algorithms suffer from the communication bottleneck due to network bandwidth pressure and/or low bandwidth utilization of the participating clients in both centralized and decentralized architectures. To deal with the communication problem while preserving the convergence performance, we introduce a communication-efficient decentralized FL framework GossipFL. In GossipFL, we 1) design a novel sparsification algorithm to enable that each client only needs to communicate with one peer with a highly sparsified model, and 2) propose a new and novel gossip matrix generation algorithm that can better utilize the bandwidth resources while preserving the convergence property. We also theoretically prove that GossipFL has convergence guarantees. We conduct experiments with three convolutional neural networks on two datasets (IID and non-IID) under two distributed environments (14 clients and 100 clients) to verify the effectiveness of GossipFL. Experimental results show that GossipFL takes less communication traffic for 38.5% and less communication time for$49.8$% than state-of-the-art solutions while achieving comparative model accuracy. Zhenheng Tang, Shaohuai Shi, Bo Li 0001, Xiaowen Chu 0001 |
IEEE Trans. Parallel Distributed Syst. | 2 |
| 2022 | EASNet: Searching Elastic and Accurate Network Architecture for Stereo Matching
Qiang Wang 0022, Shaohuai Shi, Kaiyong Zhao, Xiaowen Chu 0001 |
ECCV (32) | 2 |
| 2022 | Virtual Homogeneity Learning: Defending against Data Heterogeneity in Federated LearningabstractIn federated learning (FL), model performance typically suffers from client drift induced by data heterogeneity, and mainstream works focus on correcting client drift. We propose a different approach named virtual homogeneity learning (VHL) to directly “rectify” the data heterogeneity. In particular, VHL conducts FL with a virtual homogeneous dataset crafted to satisfy two conditions: containing no private information and being separable. The virtual dataset can be generated from pure noise shared across clients, aiming to calibrate the features from the heterogeneous clients. Theoretically, we prove that VHL can achieve provable generalization performance on the natural distribution. Empirically, we demonstrate that VHL endows FL with drastically improved convergence speed and generalization performance. VHL is the first attempt towards using a virtual dataset to address data heterogeneity, offering new and effective means to FL. Zhenheng Tang, Yonggang Zhang 0003, Shaohuai Shi, Xin He 0019, Bo Han 0003, Xiaowen Chu 0001 |
ICML | 3 |
| 2021 | Automated Model Design and Benchmarking of Deep Learning Models for COVID-19 Detection with Chest CT ScansabstractThe COVID-19 pandemic has spread globally for several months. Because its transmissibility and high pathogenicity seriously threaten people's lives, it is crucial to accurately and quickly detect COVID-19 infection. Many recent studies have shown that deep learning (DL) based solutions can help detect COVID-19 based on chest CT scans. However, most existing work focuses on 2D datasets, which may result in low quality models as the real CT scans are 3D images. Besides, the reported results span a broad spectrum on different datasets with a relatively unfair comparison. In this paper, we first use three state-of-the-art 3D models (ResNet3D101, DenseNet3D121, and MC3\_18) to establish the baseline performance on three publicly available chest CT scan datasets. Then we propose a differentiable neural architecture search (DNAS) framework to automatically search the 3D DL models for 3D chest CT scans classification and use the Gumbel Softmax technique to improve the search efficiency. We further exploit the Class Activation Mapping (CAM) technique on our models to provide the interpretability of the results. The experimental results show that our searched models (CovidNet3D) outperform the baseline human-designed models on three datasets with tens of times smaller model size and higher accuracy. Furthermore, the results also verify that CAM can be well applied in CovidNet3D for COVID-19 datasets to provide interpretability for medical diagnosis. Code: https://github.com/HKBU-HPML/CovidNet3D. Xin He 0019, Xiaowen Chu 0001, Shaohuai Shi, Jiangping Tang, Xin Liu 0027, Chenggang Yan 0001, Jiyong Zhang 0001, Guiguang Ding |
AAAI | 4 |
| 2021 | Accelerating Distributed K-FAC with Smart Parallelism of Computing and Communication TasksabstractDistributed training with synchronous stochastic gradient descent (SGD) on GPU clusters has been widely used to accelerate the training process of deep models. However, SGD only utilizes the first-order gradient in model parameter updates, which may take days or weeks. Recent studies have successfully exploited approximate second-order information to speed up the training process, in which the Kronecker-Factored Approximate Curvature (KFAC) emerges as one of the most efficient approximation algorithms for training deep models. Yet, when leveraging GPU clusters to train models with distributed KFAC (D-KFAC), it incurs extensive computation as well as introduces extra communications during each iteration. In this work, we propose D-KFAC (SPD-KFAC) with smart parallelism of computing and communication tasks to reduce the iteration time. Specifically, 1) we first characterize the performance bottlenecks of D-KFAC, 2) we design and implement a pipelining mechanism for Kronecker factors computation and communication with dynamic tensor fusion, and 3) we develop a load balancing placement for inverting multiple matrices on GPU clusters. We conduct realworld experiments on a 64-GPU cluster with 100Gb/s InfiniBand interconnect. Experimental results show that our proposed SPD-KFAC training scheme can achieve 10%-35% improvement over state-of-the-art algorithms. Shaohuai Shi, Lin Zhang 0059, Bo Li 0001 |
ICDCS | 1 |
| 2021 | Exploiting Simultaneous Communications to Accelerate Data Parallel Distributed Deep LearningabstractSynchronous stochastic gradient descent (S-SGD) with data parallelism is widely used for training deep learning (DL) models in distributed systems. A pipelined schedule of the computing and communication tasks of a DL training job is an effective scheme to hide some communication costs. In such pipelined S-SGD, tensor fusion (i.e., merging some consecutive layers' gradients for a single communication) is a key ingredient to improve communication efficiency. However, existing tensor fusion techniques schedule the communication tasks sequentially, which overlooks their independence nature. In this paper, we expand the design space of scheduling by exploiting simultaneous All-Reduce communications. Through theoretical analysis and experiments, we show that simultaneous All-Reduce communications can effectively improve the communication efficiency of small tensors. We formulate an optimization problem of minimizing the training iteration time, in which both tensor fusion and simultaneous communications are allowed. We develop an efficient optimal scheduling solution and implement the distributed training algorithm ASC-WFBP with Horovod and PyTorch. We conduct real-world experiments on an 8-node GPU cluster of 32 GPUs with 10Gbps Ethernet. Experimental results on four modern DNNs show that ASC-WFBP can achieve about 1.09 × -2.48× speedup over the baseline without tensor fusion, and 1.15× -1.35× speedup over the state-of-the-art tensor fusion solution. Shaohuai Shi, Xiaowen Chu 0001, Bo Li 0001 |
INFOCOM | 1 |
| 2021 | MG-WFBP: Merging Gradients Wisely for Efficient Communication in Distributed Deep LearningabstractDistributed synchronous stochastic gradient descent has been widely used to train deep neural networks (DNNs) on computer clusters. With the increase of computational power, network communications generally limit the system scalability. Wait-free backpropagation (WFBP) is a popular solution to overlap communications with computations during the training process. In this article, we observe that many DNNs have a large number of layers with only a small amount of data to be communicated at each layer in distributed training, which could make WFBP inefficient. Based on the fact that merging some short communication tasks into a single one can reduce the overall communication time, we formulate an optimization problem to minimize the training time in pipelining communications and computations. We derive an optimal solution that can be solved efficiently without affecting the training performance. We then apply the solution to propose a distributed training algorithm named merged-gradient WFBP (MG-WFBP) and implement it in two platforms Caffe and PyTorch. Extensive experiments in three GPU clusters are conducted to verify the effectiveness of MG-WFBP. We further exploit trace-based simulations of 4 to 2048 GPUs to explore the potential scaling efficiency of MG-WFBP. Experimental results show that MG-WFBP achieves much better scaling performance than existing methods. Shaohuai Shi, Xiaowen Chu 0001, Bo Li 0001 |
IEEE Trans. Parallel Distributed Syst. | 1 |
| 2020 | Benchmarking the Performance and Energy Efficiency of AI Accelerators for AI TrainingabstractDeep learning has become widely used in complex AI applications. Yet, training a deep neural network (DNNs) model requires a considerable amount of calculations, long running time, and much energy. Nowadays, many-core AI accelerators (e.g., GPUs and TPUs) are designed to improve the performance of AI training. However, processors from different vendors perform dissimilarly in terms of performance and energy consumption. To investigate the differences among several popular off-the-shelf processors (i.e., Intel CPU, NVIDIA GPU, AMD GPU, and Google TPU) in training DNNs, we carry out a comprehensive empirical study on the performance and energy efficiency of these processors1by benchmarking a representative set of deep learning workloads, including computation-intensive operations, classical convolutional neural networks (CNNs), recurrent neural networks (LSTM), Deep Speech 2, and Transformer. Different from the existing end-to-end benchmarks which only present the training time, We try to investigate the impact of hardware, vendor's software library, and deep learning framework on the performance and energy consumption of AI training. Our evaluation methods and results not only provide an informative guide for end users to select proper AI accelerators, but also expose some opportunities for the hardware vendors to improve their software library. Yuxin Wang 0003, Qiang Wang 0022, Shaohuai Shi, Xin He 0019, Zhenheng Tang, Kaiyong Zhao, Xiaowen Chu 0001 |
CCGRID | 3 |
| 2020 | Layer-Wise Adaptive Gradient Sparsification for Distributed Deep Learning with Convergence GuaranteesabstractTo reduce the long training time of large deep neural network (DNN) models, distributed synchronous stochastic gradient descent (S-SGD) is commonly used on a cluster of workers. However, the speedup brought by multiple workers is limited by the communication overhead. Two approaches, namely pipelining and gradient sparsification, have been separately proposed to alleviate the impact of communication overheads. Yet, the gradient sparsification methods can only initiate the communication after the backpropagation, and hence miss the pipelining opportunity. In this paper, we propose a new distributed optimization method named LAGS-SGD, which combines S-SGD with a novel layer-wise adaptive gradient sparsification (LAGS) scheme. In LAGS-SGD, every worker selects a small set of 'significant' gradients from each layer independently whose size can be adaptive to the communication-to-computation ratio of that layer. The layer-wise nature of LAGS-SGD opens the opportunity of overlapping communications with computations, while the adaptive nature of LAGS-SGD makes it flexible to control the communication time. We prove that LAGS-SGD has convergence guarantees and it has the same order of convergence rate as vanilla S-SGD under a weak analytical assumption. Extensive experiments are conducted to verify the analytical assumption and the convergence performance of LAGS-SGD. Experimental results on a 16-GPU cluster show that LAGS-SGD outperforms the original S-SGD and existing sparsified S-SGD without losing obvious model accuracy. Shaohuai Shi, Zhenheng Tang, Qiang Wang 0022, Kaiyong Zhao, Xiaowen Chu 0001 |
ECAI | 1 |
| 2020 | Communication-Efficient Decentralized Learning with Sparsification and Adaptive Peer SelectionabstractThe increasing size of machine learning models, especially deep neural network models, can improve the model generalization capability. However, large models require more training data and more computing resources (such as GPU clusters) to train. In distributed training, the communication overhead of exchanging gradients or models among workers becomes a potential system bottleneck that limits the system scalability. Recently, many research works aim to reduce communication time of two types of distributed deep learning architectures, centralized and decentralized. Zhenheng Tang, Shaohuai Shi, Xiaowen Chu 0001 |
ICDCS | 2 |
| 2020 | Efficient Sparse-Dense Matrix-Matrix Multiplication on GPUs Using the Customized Sparse Storage FormatabstractMultiplication of a sparse matrix to a dense matrix (SpDM) is widely used in many areas like scientific computing and machine learning. However, existing work under-looks the performance optimization of SpDM on modern manycore architectures like GPUs. The storage data structures help sparse matrices store in a memory-saving format, but they bring difficulties in optimizing the performance of SpDM on modern GPUs due to irregular data access of the sparse structure, which results in lower resource utilization and poorer performance. In this paper, we refer to the roofline performance model of GPUs to design an efficient SpDM algorithm called GCOOSpDM, in which we exploit coalescent global memory access, fast shared memory reuse, and more operations per byte of global memory traffic. Experiments are evaluated on three Nvidia GPUs (i.e., GTX 980, GTX Titan X Pascal, and Tesla P100) using a large number of matrices including a public dataset and randomly generated matrices. Experimental results show that GCOOSpDM achieves 1.5-8x speedup over Nvidia's library cuSPARSE in many matrices. Shaohuai Shi, Qiang Wang 0022, Xiaowen Chu 0001 |
ICPADS | 1 |
| 2020 | FADNet: A Fast and Accurate Network for Disparity EstimationabstractDeep neural networks (DNNs) have achieved great success in the area of computer vision. The disparity estimation problem tends to be addressed by DNNs which achieve much better prediction accuracy in stereo matching than traditional hand-crafted feature based methods. On one hand, however, the designed DNNs require significant memory and computation resources to accurately predict the disparity, especially for those 3D convolution based networks, which makes it difficult for deployment in real-time applications. On the other hand, existing computation-efficient networks lack expression capability in large-scale datasets so that they cannot make an accurate prediction in many scenarios. To this end, we propose an efficient and accurate deep network for disparity estimation named FADNet with three main features: 1) It exploits efficient 2D based correlation layers with stacked blocks to preserve fast computation; 2) It combines the residual structures to make the deeper model easier to learn; 3) It contains multi-scale predictions so as to exploit a multi-scale weight scheduling training technique to improve the accuracy. We conduct experiments to demonstrate the effectiveness of FADNet on two popular datasets, Scene Flow and KITTI 2015. Experimental results show that FADNet achieves state-of-the-art prediction accuracy, and runs at a significant order of magnitude faster speed than existing 3D models. The codes of FADNet are available at https://github.com/HKBU-HPML/FADNet. Qiang Wang 0022, Shaohuai Shi, Shizhen Zheng, Kaiyong Zhao, Xiaowen Chu 0001 |
ICRA | 2 |
| 2020 | Communication-Efficient Distributed Deep Learning with Merged Gradient Sparsification on GPUsabstractDistributed synchronous stochastic gradient descent (SGD) algorithms are widely used in large-scale deep learning applications, while it is known that the communication bottleneck limits the scalability of the distributed system. Gradient sparsification is a promising technique to significantly reduce the communication traffic, while pipelining can further overlap the communications with computations. However, gradient sparsification introduces extra computation time, and pipelining requires many layer-wise communications which introduce significant communication startup overheads. Merging gradients from neighbor layers could reduce the startup overheads, but on the other hand it would increase the computation time of sparsification and the waiting time for the gradient computation. In this paper, we formulate the trade-off between communications and computations (including backward computation and gradient sparsification) as an optimization problem, and derive an optimal solution to the problem. We further develop the optimal merged gradient sparsification algorithm with SGD (OMGS-SGD) for distributed training of deep learning. We conduct extensive experiments to verify the convergence properties and scaling performance of OMGS-SGD. Experimental results show that OMGS-SGD achieves up to 31% end-to-end time efficiency improvement over the state-of-the-art sparsified SGD while preserving nearly consistent convergence performance with original SGD without sparsification on a 16-GPU cluster connected with 1Gbps Ethernet. Shaohuai Shi, Qiang Wang 0022, Xiaowen Chu 0001, Bo Li 0001, Yang Qin 0001, Ruihao Liu, Xinxiao Zhao |
INFOCOM | 1 |
| 2019 | Computer-Aided Clinical Skin Disease Diagnosis Using CNN and Object Detection ModelsabstractSkin disease is one of the most common types of human diseases, which may happen to everyone regardless of age, gender or race. Due to the high visual diversity, human diagnosis highly relies on personal experience; and there is a serious shortage of experienced dermatologists in many countries. To alleviate this problem, computer-aided diagnosis with state-of-the-art (SOTA) machine learning techniques would be a promising solution. In this paper, we aim at understanding the performance of convolutional neural network (CNN) based approaches. We first build two versions of skin disease datasets from Internet images: (a) Skin -10, which contains 10 common classes of skin disease with a total of 10,218 images; (b) Skin -100, which is a larger dataset that consists of 19,807 images of 100 skin disease classes. Based on these datasets, we benchmark several SOTA CNN models and show that the accuracy of skin -100 is much lower than the accuracy of skin -10. We then implement an ensemble method based on several CNN models and achieve the best accuracy of 79.01% for Skin -10 and 53.54% for Skin -100. We also present an object detection based approach by introducing bounding boxes into the Skin -10 dataset. Our results show that object detection can help improve the accuracy of some skin disease classes. Xin He 0019, Zhi-Li Wu, Wu Yu, Xiaowen Chu 0001, Shaohuai Shi, Zhenheng Tang, Yuxin Wang 0003, Ronghao Ni, Xiaofeng Zhang 0002 |
IEEE BigData | 7 |
| 2019 | A Distributed Synchronous SGD Algorithm with Global Top-k Sparsification for Low Bandwidth NetworksabstractDistributed synchronous stochastic gradient descent (S-SGD) with data parallelism has been widely used in training large-scale deep neural networks (DNNs), but it typically requires very high communication bandwidth between computational workers (e.g., GPUs) to exchange gradients iteratively. Recently, Top-k sparsification techniques have been proposed to reduce the volume of data to be exchanged among workers and thus alleviate the network pressure. Top-k sparsification can zero-out a significant portion of gradients without impacting the model convergence. However, the sparse gradients should be transferred with their indices, and the irregular indices make the sparse gradients aggregation difficult. Current methods that use AllGather to accumulate the sparse gradients have a communication complexity of O(kP), where P is the number of workers, which is inefficient on low bandwidth networks with a large number of workers. We observe that not all top-k gradients from P workers are needed for the model update, and therefore we propose a novel global Top-k (gTop-k) sparsification mechanism to address the difficulty of aggregating sparse gradients. Specifically, we choose global top-k largest absolute values of gradients from P workers, instead of accumulating all local top-k gradients to update the model in each iteration. The gradient aggregation method based on gTop-k sparsification, namely gTopKAllReduce, reduces the communication complexity from O(kP) to O(k log P). Through extensive experiments on different DNNs, we verify that gTop-k S-SGD has nearly consistent convergence performance with S-SGD, and it has only slight degradations on generalization performance. In terms of scaling efficiency, we evaluate gTop-k on a cluster with 32 GPU machines which are interconnected with 1 Gbps Ethernet. The experimental results show that our method achieves 2.7-12× higher scaling efficiency than S-SGD with dense gradients and 1.1-1.7× improvement than the existing Top-k S-SGD. Shaohuai Shi, Qiang Wang 0022, Kaiyong Zhao, Zhenheng Tang, Yuxin Wang 0003, Xiaowen Chu 0001 |
ICDCS | 1 |
| 2019 | A Convergence Analysis of Distributed SGD with Communication-Efficient Gradient SparsificationabstractGradient sparsification is a promising technique to significantly reduce the communication overhead in decentralized synchronous stochastic gradient descent (S-SGD) algorithms. Yet, many existing gradient sparsification schemes (e.g., Top-k sparsification) have a communication complexity of O(kP), where k is the number of selected gradients by each worker and P is the number of workers. Recently, the gTop-k sparsification scheme has been proposed to reduce the communication complexity from O(kP) to O(k logP), which significantly boosts the system scalability. However, it remains unclear whether the gTop-k sparsification scheme can converge in theory. In this paper, we first provide theoretical proofs on the convergence of the gTop-k scheme for non-convex objective functions under certain analytic assumptions. We then derive the convergence rate of gTop-k S-SGD, which is at the same order as the vanilla mini-batch SGD. Finally, we conduct extensive experiments on different machine learning models and data sets to verify the soundness of the assumptions and theoretical results, and discuss the impact of the compression ratio on the convergence performance. Shaohuai Shi, Kaiyong Zhao, Qiang Wang 0022, Zhenheng Tang, Xiaowen Chu 0001 |
IJCAI | 1 |
| 2019 | MG-WFBP: Efficient Data Communication for Distributed Synchronous SGD AlgorithmsabstractDistributed synchronous stochastic gradient descent has been widely used to train deep neural networks on computer clusters. With the increase of computational power, network communications have become one limiting factor on the system scalability. In this paper, we observe that many deep neural networks have a large number of layers with only a small amount of data to be communicated. Based on the fact that merging some short communication tasks into a single one may reduce the overall communication time, we formulate an optimization problem to minimize the training iteration time. We develop an optimal solution named merged-gradient wait-free backpropagation (MG-WFBP) and implement it in our open-source deep learning platform B-Caffe. Our experimental results on an 8-node GPU cluster with 10GbE interconnect and trace-based simulation results on a 64-node cluster both show that the MG-WFBP algorithm can achieve much better scaling efficiency than existing methods WFBP and SyncEASGD. Shaohuai Shi, Xiaowen Chu 0001, Bo Li 0001 |
INFOCOM | 1 |
| 2018 | A DAG Model of Synchronous Stochastic Gradient Descent in Distributed Deep LearningabstractWith huge amounts of training data, deep learning has made great breakthroughs in many artificial intelligence (AI) applications. However, such large-scale data sets present computational challenges, requiring training to be distributed on a cluster equipped with accelerators like GPUs. With the fast increase of G PU computing power, the data communications among GPUs have become a potential bottleneck on the overall training performance. In this paper, we first propose a general directed acyclic graph (DAG) model to describe the distributed synchronous stochastic gradient descent (S-SG D) algorithm, which has been widely used in distributed deep learning frameworks. To understand the practical impact of data communications on training performance, we conduct extensive empirical studies on four state-of-the-art distributed deep learning frameworks (i.e., Caffe-MPI, CNTK, MXNet and TensorFlow) over multi-GPU and multi-node environments with different data communication techniques, including PCIe, NVLink, 10GbE, and InfiniBand. Through both analytical and experimental studies, we identify the potential bottlenecks and overheads that could be further optimized. At last, we make the data set of our experimental traces publicly available, which could be used to support simulation-based studies. Shaohuai Shi, Qiang Wang 0022, Xiaowen Chu 0001, Bo Li 0001 |
ICPADS | 1 |
| 2017 | Supervised Learning Based Algorithm Selection for Deep Neural NetworksabstractMany recent deep learning platforms rely on thirdparty libraries (such as cuBLAS) to utilize the computing power of modern hardware accelerators (such as GPUs). However, we observe that they may achieve suboptimal performance because the library functions are not used appropriately. In this paper, we target at optimizing the operations of multiplying a matrix with the transpose of another matrix (referred to as NT operation hereafter), which contribute half of the training time of fully connected deep neural networks. Rather than directly calling the library function, we propose a supervised learning based algorithm selection approach named MTNN, which uses a gradient boosted decision tree to select one from two alternative NT implementations intelligently: (1) calling the cuBLAS library function; (2) calling our proposed algorithm TNN that uses an efficient out-of-place matrix transpose. We evaluate the performance of MTNN on two modern GPUs: NVIDIA GTX 1080 and NVIDIA Titan X Pascal. MTNN can achieve 96% of prediction accuracy with very low computational overhead, which results in an average of 54% performance improvement on a range of NT operations. To further evaluate the impact of MTNN on the training process of deep neural networks, we have integrated MTNN into a popular deep learning platform Caffe. Our experimental results show that the revised Caffe can outperform the original one by an average of 28%. Both MTNN and the revised Caffe are open-source. Shaohuai Shi, Pengfei Xu 0006, Xiaowen Chu 0001 |
ICPADS | 1 |