Peng Sun 0006

dblp:88/619-6 · DBLP profile ↗
← Back
39ranked-venue papers
7as first author
29since 2021 · last 2026
0000-0001-8456-0491ORCID · conflict

Domains — the database's venue-derived domains; a paper can count in several

Systems, architecture and hardware · 21 · 3 first-author · 17 since 2021Computer networks · 6 · 3 since 2021Software engineering, systems software and programming languages · 6 · 6 since 2021Applied, interdisciplinary, general and emerging computing · 5 · 4 first-author · 2 since 2021Databases, data management, data science and information retrieval · 4 · 4 since 2021Artificial intelligence and machine learning · 3 · 1 first-author · 2 since 2021
YearPublicationVenuePosition
2026 Bat: Efficient Generative Recommender Serving with Bipartite Attention
abstract
Generative Recommenders (GRs) have recently emerged as promising alternatives to traditional Deep Learning Recommendation Models (DLRMs). Despite their potential, GRs remain computationally expensive in inference, exhibiting compute-bound characteristics similar to the prefill stage of Large Language Model (LLM) inference. Prefix caching can reduce redundant computation by reusing previously constructed KV caches. However, the unique properties of GRs, i.e., highly personalized user profiles and real-time item retrieval, make cache reuse across queries challenging, resulting in limited computational savings.
Jie Sun 0017, Shaohang Wang, Zimo Zhang, Peng Sun 0006, Bo Zhao 0019, Bingsheng He, Fei Wu 0001, Zeke Wang
ASPLOS (2)6
2026 Zeppelin: Balancing Variable-length Workloads in Data Parallel Large Model Training
abstract
Training large language models (LLMs) with increasingly long and varying sequence lengths introduces severe load imbalance challenges in large-scale data-parallel training. Recent frameworks attempt to mitigate these issues through data reorganization or hybrid parallel strategies. However, they often overlook how computational and communication costs scale with sequence length, resulting in suboptimal performance. We identify three critical challenges: (1) varying computation-to-communication ratios across sequences of different lengths in distributed attention, (2) mismatch between static NIC-GPU affinity and dynamic parallel workloads, and (3) distinct optimal partitioning strategies required for quadratic attention versus linear components.
Chang Chen 0001, Tiancheng Chen, Jiangfei Duan, Qianchao Zhu, Zerui Wang, Qinghao Hu 0004, Peng Sun 0006, Chao Yang 0002, Torsten Hoefler
EuroSys7
2026 AdaCheck: An Adaptive Checkpointing System for Efficient LLM Training with Redundancy Utilization
Zhiquan Lai, Ke-shi Ge, Qiaoling Chen, Peng Sun 0006, Dongsheng Li 0001, Kai Lu 0001
FAST6
2026 SPPO: Making Million-Token LLM Training Practical on Modest GPU Clusters
abstract
In recent years, Large Language Models (LLMs) have exhibited remarkable capabilities, driving advancements in real-world applications. However, training LLMs on increasingly long input sequences imposes significant challenges due to high GPU memory and computational demands. Existing solutions face two key limitations: (1) memory reduction techniques, such as activation recomputation and CPU offloading, compromise training efficiency; (2) distributed parallelism strategies require excessive GPU resources, limiting the scalability of input sequence length.
Qiaoling chen, Shenggui Li, Wei Gao 0064, Peng Sun 0006, Yonggang Wen 0001, Tianwei Zhang 0004
ICS4
2026 Di-PS: System-Algorithm Co-Design for Asynchronous and Heterogeneous Cross-cluster LLM Training at Scale
Qiaoling Chen, Zhiquan Lai, Penglong Jiao, Wenwen Qu, Peng Sun 0006, Xingcheng Zhang, Xiaoge Deng, Dongsheng Li 0001, Kai Lu 0001, Tianwei Zhang 0004
NSDI8
2026 DBAIOps: A Reasoning LLM-Enhanced Database Operation and Maintenance System using Knowledge Graphs
Wei Zhou 0053, Peng Sun 0006, Xuanhe Zhou, Qianglei Zang, Tieying Zhang, Guoliang Li 0001, Fan Wu 0006
Proc. VLDB Endow.2
2025 Parrot: A Training Pipeline Enhances Both Program CoT and Natural Language CoT for Reasoning
abstract
Senjie Jin, Lu Chen, Zhiheng Xi, Yuhui Wang, Sirui Song, Yuhao Zhou, Xinbo Zhang, Peng Sun, Hong Lu, Tao Gui, Qi Zhang, Xuanjing Huang. Proceedings of the 2025 Conference on Empirical Methods in Natural Language Processing. 2025.
Senjie Jin, Lu Chen 0001, Zhiheng Xi, Sirui Song, Yuhao Zhou 0005, Xinbo Zhang, Peng Sun 0006, Tao Gui, Qi Zhang 0001, Xuanjing Huang 0001
EMNLP8
2025 IceFrog: A Layer-Elastic Scheduling System for Deep Learning Training in GPU Clusters
abstract
The high resource demand of deep learning training (DLT) workloads necessitates the design of efficient schedulers. While most existing schedulers expedite DLT workloads by considering GPU sharing and elastic training, they neglectlayer elasticity, which dynamically freezes certain layers of a network. This technique has been shown to significantly speed up individual workloads. In this paper, we explore how to incorporatelayer elasticityinto DLT scheduler designs to achieve higher cluster-wide efficiency. A key factor that hinders the application of layer elasticity in GPU clusters is the potential loss in model accuracy, making users reluctant to enable layer elasticity for their workloads. It is necessary to have an efficient layer-elastic system, which can well balance training accuracy and speed for layer elasticity. We introduceIceFrog, the first scheduling system that utilizes layer elasticity to improve the efficiency of DLT workloads in GPU clusters. It achieves this goal with superior algorithmic designs and intelligent resource management. In particular, (1) we model the frozen penalty and layer-aware throughput to measure the effective progress metric of layer-elastic workloads. (2) We design a novel scheduler to further improve the efficiency of layer elasticity. We implement and deployIceFrogin a physical cluster of 48 GPUs. Extensive evaluations and large-scale simulations show thatIceFrogreduces average job completion times by 36-48% relative to state-of-the-art DL schedulers.
Wei Gao 0064, Zhuoyuan Ouyang, Peng Sun 0006, Tianwei Zhang 0004, Yonggang Wen 0001
IEEE Trans. Parallel Distributed Syst.3
2025 GreenFlow: A Carbon-Efficient Scheduler for Deep Learning Workloads
abstract
Deep learning (DL) has become a key component of modern software. Training DL models leads to huge carbon emissions. In data centers, it is important to reduce carbon emissions while completing DL training jobs early. In this article, we propose GreenFlow, a GPU cluster scheduler that reduces the average Job Completion Time (JCT) under a carbon emission budget. We first present performance models for DL training jobs to predict the throughput and energy consumption performance under different configurations. Based on the performance models and the carbon intensity of the grid, GreenFlow dynamically allocates GPUs, and adjusts the GPU-level and job-level configurations of DL training jobs. GreenFlow applies network packing and buddy allocation to job placement, thus avoiding extra carbon incurred by resource fragmentations. Evaluations on a real testbed show that when emitting the same amount of carbon, GreenFlow can improve the average JCT by up to 2.15×, compared to competitive baselines.
Diandian Gu, Peng Sun 0006, Xin Jin 0008, Xuanzhe Liu
IEEE Trans. Parallel Distributed Syst.3
2024 Centauri: Enabling Efficient Scheduling for Communication-Computation Overlap in Large Model Training via Communication Partitioning
abstract
Efficiently training large language models (LLMs) necessitates the adoption of hybrid parallel methods, integrating multiple communications collectives within distributed partitioned graphs. Overcoming communication bottlenecks is crucial and is often achieved through communication and computation overlaps. However, existing overlap methodologies tend to lean towards either fine-grained kernel fusion or limited operation scheduling, constraining performance optimization in heterogeneous training environments.
Chang Chen 0001, Qianchao Zhu, Jiangfei Duan, Peng Sun 0006, Xingcheng Zhang, Chao Yang 0002
ASPLOS (3)5
2024 Sylvie: 3D-Adaptive and Universal System for Large-Scale Graph Neural Network Training
abstract
Distributed full-graph training of Graph Neural Networks (GNNs) has been widely adopted to learn large-scale graphs. While recent system advancements can improve the training throughput of GNNs, their practical adoption is limited by the potential accuracy decline. This concern is particularly prominent in deeper and more intricate GNN architectures, where noticeable performance degradation becomes apparent. Moreover, existing works fail to comprehensively consider diverse opportunities for acceleration. Motivated by these deficiencies, we propose Sylvie,a full-graph training system that not only improves the training throughput substantially but also maintains the model quality for universal GNNs. By harnessing the inherent information embedded in the graph data and model structure, Sylvie intelligently optimizes GNN training across three key dimensions: data, time, and execution. It identifies performance-relevant features of the input graph offline as subsequent optimization guidance. Subsequently, Sylvie devises an online convergence-maintenance strategy that adaptively integrates and aligns GNN-specific quantization and inter-epoch asynchronous training with the real-time training characteristics. Extensive experiments demonstrate that Sylvie surpasses existing GNN training systems by up to 17.2× speedup for both shallow and deep GNNs, without compromising the model accuracy.
Meng Zhang 0045, Qinghao Hu 0004, Cheng Wan 0005, Haozhao Wang, Peng Sun 0006, Yonggang Wen 0001, Tianwei Zhang 0001
ICDE5
2024 Training Large Language Models for Reasoning through Reverse Curriculum Reinforcement Learning
abstract
In this paper, we propose R$^3$: Learning Reasoning through Reverse Curriculum Reinforcement Learning (RL), a novel method that employs only outcome supervision to achieve the benefits of process supervision for large language models. The core challenge in applying RL to complex reasoning is to identify a sequence of actions that result in positive rewards and provide appropriate supervision for optimization. Outcome supervision provides sparse rewards for final results without identifying error locations, whereas process supervision offers step-wise rewards but requires extensive manual annotation. R$^3$ overcomes these limitations by learning from correct demonstrations. Specifically, R$^3$ progressively slides the start state of reasoning from a demonstration’s end to its beginning, facilitating easier model exploration at all stages. Thus, R$^3$ establishes a step-wise curriculum, allowing outcome supervision to offer step-level signals and precisely pinpoint errors. Using Llama2-7B, our method surpasses RL baseline on eight reasoning tasks by $4.1$ points on average. Notably, in program-based reasoning, 7B-scale models perform comparably to larger models or closed-source models with our R$^3$.
Zhiheng Xi, Wenxiang Chen, Boyang Hong, Senjie Jin, Wei He 0024, Yiwen Ding, Shichun Liu, Junzhe Wang 0001, Honglin Guo, Xiaoran Fan, Yuhao Zhou 0005, Shihan Dou, Xiao Wang 0001, Xinbo Zhang, Peng Sun 0006, Tao Gui, Qi Zhang 0001, Xuanjing Huang 0001
ICML18
2024 AutoSched: An Adaptive Self-configured Framework for Scheduling Deep Learning Training Workloads
abstract
Modern Deep Learning Training (DLT) schedulers in GPU datacenters are designed to be very sophisticated with many configurations. These configurations need to be adjusted delicately as they can significantly affect the scheduling performance. Existing schedulers require the datacenter operator to tune the configurations only once before they are deployed, based on the historical workload traces. Unfortunately, workloads in a datacenter would experience dynamic changes and deviate a lot from the historical ones over time, making the pre-determined configurations less effective.
Wei Gao 0064, Shangwei Guo, Peng Sun 0006, Yonggang Wen 0001, Tianwei Zhang 0004
ICS5
2024 Ymir: A Scheduler for Foundation Model Fine-tuning Workloads in Datacenters
abstract
The breakthrough of foundation models makes foundation model fine-tuning (FMF) workloads prevalent in modern GPU datacenters. However, existing schedulers tailored for model training do not consider the unique characteristics of FMs, making them inefficient in handling FMF workloads. To bridge the gap, we propose Ymir, a scheduler to improve the efficiency of FMF workloads in GPU datacenters. Ymir leverages the shared FM backbone architecture to expedite FMF workloads from two aspects: (1) Ymir investigates the task transferability among different FMF workloads and automatically merges FMF workloads with the same FM into one to improve the cluster-wide efficiency via transfer learning. (2) Ymir reuses the fine-tuning runtime of FMF workloads to reduce the significant context switch overhead. We conduct 32-GPU physical experiments and 240-GPU trace-driven simulations to validate the effectiveness of Ymir. Ymir can reduce the average job completion time by up to 4.3 × compared with existing state-of-the-art schedulers. It also promotes scheduling fairness by fully exploiting the task transferability. More supplementary materials can be found on our project website https://sites.google.com/view/ymir-project.
Wei Gao 0064, Weiming Zhuang, Minghao Li 0005, Peng Sun 0006, Yonggang Wen 0001, Tianwei Zhang 0004
ICS4
2024 Lins: Reducing Communication Overhead of ZeRO for Efficient LLM Training
abstract
Training large language models (LLMs) encounters challenges in GPU memory consumption due to the high memory requirements of model states. The widely used Zero Redundancy Optimizer (ZeRO) addresses this issue through strategic sharding but introduces communication challenges at scale. To tackle this problem, we propose Lins, a system designed to optimize ZeRO for scalable LLM training. Lins incorporates three flexible sharding strategies: Full-Replica, Full-Sharding, and Partial-Sharding, and allows each component within the model states (Parameters, Gradients, Optimizer States) to independently choose a sharding strategy as well as the device mesh. We conduct a thorough analysis of communication costs, formulating an optimization problem to discover the optimal sharding strategy. Evaluations demonstrate up to 52% Model FLOPs Utilization (MFU) when training the LLaMA-based model on 1024 GPUs, resulting in a 1.56 times improvement in training throughput compared to newly proposed systems like MiCS and ZeRO++.
Qiaoling Chen, Qinghao Hu 0004, Guoteng Wang, Yingtong Xiong, Yang Gao 0042, Hang Yan 0001, Yonggang Wen 0001, Tianwei Zhang 0004, Peng Sun 0006
IWQoS11
2024 Characterization of Large Language Model Development in the Datacenter
Qinghao Hu 0004, Zhisheng Ye 0002, Zerui Wang, Guoteng Wang, Meng Zhang 0045, Qiaoling Chen, Peng Sun 0006, Dahua Lin, Xiaolin Wang 0001, Yingwei Luo, Yonggang Wen 0001, Tianwei Zhang 0004
NSDI7
2024 dLoRA: Dynamically Orchestrating Requests and Adapters for LoRA LLM Serving
Bingyang Wu, Ruidong Zhu, Peng Sun 0006, Xuanzhe Liu, Xin Jin 0008
OSDI4
2024 TorchGT: A Holistic System for Large-Scale Graph Transformer Training
abstract
Graph Transformer is a new architecture that surpasses GNNs in graph learning. While there emerge inspiring algorithm advancements, their practical adoption is still limited, particularly on real-world graphs involving up to millions of nodes. We observe existing graph transformers fail on large-scale graphs mainly due to heavy computation, limited scalability and inferior model quality. Motivated by these observations, we propose TORCHGT, the first efficient, scalable, and accurate graph transformer training system. TORCHGT optimizes training at three different levels. At algorithm level, by harnessing the graph sparsity, TORCHGT introduces a Dual-interleaved Attention which is computation-efficient and accuracy-maintained. At runtime level, TORCHGT scales training across workers with a communicationlight Cluster-aware Graph Parallelism. At kernel level, an Elastic Computation Reformation further optimizes the computation by reducing memory access latency in a dynamic way. Extensive experiments demonstrate that TORCHGT boosts training by up to 62.7× and supports graph sequence lengths of up to 1M.
Meng Zhang 0045, Jie Sun 0017, Qinghao Hu 0004, Peng Sun 0006, Zeke Wang, Yonggang Wen 0001, Tianwei Zhang 0004
SC4
2024 LoongServe: Efficiently Serving Long-Context Large Language Models with Elastic Sequence Parallelism
abstract
The context window of large language models (LLMs) is rapidly increasing, leading to a huge variance in resource usage between different requests as well as between different phases of the same request. Restricted by static parallelism strategies, existing LLM serving systems cannot efficiently utilize the underlying resources to serve variable-length requests in different phases. To address this problem, we propose a new parallelism paradigm, elastic sequence parallelism (ESP), to elastically adapt to the variance across different requests and phases. Based on ESP, we design and build LoongServe, an LLM serving system that (1) improves computation efficiency by elastically adjusting the degree of parallelism in real-time, (2) improves communication efficiency by reducing key-value cache migration overhead and overlapping partial decoding communication with computation, and (3) improves GPU memory efficiency by reducing key-value cache fragmentation across instances. Our evaluation under diverse real-world datasets shows that LoongServe improves the throughput by up to 3.85× compared to chunked prefill and 5.81× compared to prefill-decoding disaggregation.
Bingyang Wu, Yinmin Zhong, Peng Sun 0006, Xuanzhe Liu, Xin Jin 0008
SOSP4
2024 FedDSE: Distribution-aware Sub-model Extraction for Federated Learning over Resource-constrained Devices
abstract
Sub-model extraction based federated learning has emerged as a popular strategy for training models on resource-constrained devices. However, existing methods treat all clients equally and extract sub-models using predetermined rules, which disregard the statistical heterogeneity across clients and may lead to fierce competition among them. Specifically, this paper identifies that when making predictions, different clients tend to activate different neurons of the entire model related to their respective distributions. If highly activated neurons from some clients with one distribution are incorporated into the sub-model allocated to other clients with different distributions, they will be forced to fit the new distributions, which can hinder their activation over the previous clients and result in a performance reduction. Motivated by this finding, we propose a novel method called FedDSE, which can reduce the conflicts among clients by extracting sub-models based on the data distribution of each client. The core idea of FedDSE is to empower each client to adaptively extract neurons from the entire model based on their activation over the local dataset. We theoretically show that FedDSE can achieve an improved classification score and convergence over general neural networks with the ReLU activation function. Experimental results on various datasets and models show that FedDSE outperforms all state-of-the-art baselines.
Haozhao Wang, Yabo Jia, Meng Zhang 0045, Qinghao Hu 0004, Hao Ren 0001, Peng Sun 0006, Yonggang Wen 0001, Tianwei Zhang 0004
WWW6
2024 UniSched: A Unified Scheduler for Deep Learning Training Jobs With Different User Demands
abstract
The growth of deep learning training (DLT) jobs in modern GPU clusters calls for efficient deep learning (DL) scheduler designs. Due to the extensive applications of DL technology, developers may have different demands for their DLT jobs. It is important for a GPU cluster to support all these demands and efficiently execute those DLT jobs. Unfortunately, existing DL schedulers mainly focus on part of those demands, and cannot provide comprehensive scheduling services.In this work, we present UniSched, a unified scheduler to optimize different types of scheduling objectives (e.g., guaranteeing the deadlines of SLO jobs, minimizing the latency of best-effort jobs). Meanwhile, UniSchedsupports different job stopping criteria (e.g., iteration-based, performance-based). UniSched includes two key components: Estimator for estimating the job duration, and Selector for selecting jobs and allocating resources. We perform large-scale simulations over the job traces from the production clusters. Compared to state-of-the-art schedulers, UniSchedcan significantly decrease the deadline miss rate of SLO jobs by up to 6.84×, and the latency of best-effort jobs by up to 4.02×, To demonstrate the practicality of UniSched, we implement and deploy a prototype on Kubernetes in a physical cluster consisting of 64 GPUs.
Wei Gao 0064, Zhisheng Ye 0002, Peng Sun 0006, Tianwei Zhang 0004, Yonggang Wen 0001
IEEE Trans. Computers3
2023 Lucid: A Non-intrusive, Scalable and Interpretable Scheduler for Deep Learning Training Jobs
abstract
While recent deep learning workload schedulers exhibit excellent performance, it is arduous to deploy them in practice due to some substantial defects, including inflexible intrusive manner, exorbitant integration and maintenance cost, limited scalability, as well as opaque decision processes. Motivated by these issues, we design and implement Lucid, a non-intrusive deep learning workload scheduler based on interpretable models. It consists of three innovative modules. First, a two-dimensional optimized profiler is introduced for efficient job metric collection and timely debugging job feedback. Second, Lucid utilizes an indolent packing strategy to circumvent interference. Third, Lucid orchestrates resources based on estimated job priority values and sharing scores to achieve efficient scheduling. Additionally, Lucid promotes model performance maintenance and system transparent adjustment via a well-designed system optimizer. Our evaluation shows that Lucid reduces the average job completion time by up to 1.3× compared with state-of-the-art preemptive scheduler Tiresias. Furthermore, it provides explicit system interpretations and excellent scalability for practical deployment.
Qinghao Hu 0004, Meng Zhang 0045, Peng Sun 0006, Yonggang Wen 0001, Tianwei Zhang 0004
ASPLOS (2)3
2023 Hydro: Surrogate-Based Hyperparameter Tuning Service in Datacenters
Qinghao Hu 0004, Zhisheng Ye 0002, Meng Zhang 0045, Qiaoling Chen, Peng Sun 0006, Yonggang Wen 0001, Tianwei Zhang 0004
OSDI5
2022 Titan: a scheduler for foundation model fine-tuning workloads
abstract
The recent breakthrough of foundation model (FM) research raises a new trend to acquire efficient DL models by fine-tuning FMs with low-resource datasets. Current GPU clusters are mainly established to develop DL models by training from scratch. How to tailor a GPU cluster scheduler for FM fine-tuning workloads is still not explored.
Wei Gao 0064, Peng Sun 0006, Yonggang Wen 0001, Tianwei Zhang 0004
SoCC2
2022 Primo: Practical Learning-Augmented Systems with Interpretable Models
Qinghao Hu 0004, Harsha Nori, Peng Sun 0006, Yonggang Wen 0001, Tianwei Zhang 0004
USENIX ATC3
2022 GradientFlow: Optimizing Network Performance for Large-Scale Distributed DNN Training
abstract
It is important to scale out deep neural network (DNN) training for reducing model training time. The high communication overhead is one of the major performance bottlenecks for distributed DNN training across multiple GPUs. Our investigations have shown that popular open-source DNN systems could only achieve 2.5 speedup ratio on 64 GPUs connected by 56 Gbps network. To address this problem, we propose a communication backend named GradientFlow for distributed DNN training, and employ a set of network optimization techniques. First, we integrate ring-based allreduce, mixed-precision training, and computation/communication overlap into GradientFlow. Second, we propose lazy allreduce to improve network throughput by fusing multiple communication operations into a single one, and design coarse-grained sparse communication to reduce network traffic by only transmitting important gradient chunks. When training AlexNet and ResNet-50 on the ImageNet dataset using 512 GPUs, our approach could achieve 410.2 and 434.1 speedup ratio, respectively.
Peng Sun 0006, Yonggang Wen 0001, Ruobing Han, Wansen Feng, Shengen Yan
IEEE Trans. Big Data1
2022 Astraea: A Fair Deep Learning Scheduler for Multi-Tenant GPU Clusters
abstract
Modern GPU clusters are designed to support distributed Deep Learning jobs from multiple tenants concurrently. Each tenant may have varied and dynamic resource demands. Unfortunately, existing GPU schedulers fail to thoroughly consider the fairness among the tenants and jobs, which can result in unbalanced resource allocation and unfair user experience. In this article, we present an efficient solution to provide strong fairness while maintaining high scheduling effectiveness in multi-tenant GPU clusters. First, we introduce a novel Long-Term GPU-time Fairness metric, which can comprehensively evaluate the fairness at both the tenant and job levels, based on both the temporal and spatial impacts of resource allocation. Second, we design a new and practical GPU scheduler,Astraea, to enforce the desired fairness among tenants and jobs. Large-scale evaluations show thatAstraeacan improve tenant fairness by up to 9.42× compared to state-of-the-art schedulers, without sacrificing the average job completion time.
Zhisheng Ye 0002, Peng Sun 0006, Wei Gao 0064, Tianwei Zhang 0004, Xiaolin Wang 0001, Shengen Yan, Yingwei Luo
IEEE Trans. Parallel Distributed Syst.2
2021 Chronus: A Novel Deadline-aware Scheduler for Deep Learning Training Jobs
abstract
Modern GPU clusters support Deep Learning training (DLT) jobs in a distributed manner. Job scheduling is the key to improve the training performance, resource utilization and fairness across users. Different training jobs may require various objectives and demands in terms of completion time. How to efficiently satisfy all these requirements is not extensively studied.
Wei Gao 0064, Zhisheng Ye 0002, Peng Sun 0006, Yonggang Wen 0001, Tianwei Zhang 0004
SoCC3
2021 Characterization and prediction of deep learning workloads in large-scale GPU datacenters
abstract
Modern GPU datacenters are critical for delivering Deep Learning (DL) models and services in both the research community and industry. When operating a datacenter, optimization of resource scheduling and management can bring significant financial benefits. Achieving this goal requires a deep understanding of the job features and user behaviors. We present a comprehensive study about the characteristics of DL jobs and resource management. First, we perform a large-scale analysis of real-world job traces from SenseTime. We uncover some interesting conclusions from the perspectives of clusters, jobs and users, which can facilitate the cluster system designs. Second, we introduce a general-purpose framework, which manages resources based on historical data. As case studies, we design (1) a Quasi-Shortest-Service-First scheduling service, which can minimize the cluster-wide average job completion time by up to 6.5×; (2) a Cluster Energy Saving service, which improves overall cluster utilization by up to 13%.
Qinghao Hu 0004, Peng Sun 0006, Shengen Yan, Yonggang Wen 0001, Tianwei Zhang 0004
SC2
2020 Elan: Towards Generic and Efficient Elastic Training for Deep Learning
abstract
Showing a promising future in improving resource utilization and accelerating training, elastic deep learning training has been attracting more and more attention recently. Nevertheless, existing approaches to provide elasticity have certain limitations. They either fail to fully explore the parallelism of deep learning training when scaling out or lack an efficient mechanism to replicate training states among different devices.To address these limitations, we design Elan, a generic and efficient elastic training system for deep learning. In Elan, we propose a novel hybrid scaling mechanism to make a good trade-off between training efficiency and model performance when exploring more parallelism. We exploit the topology of underlying devices to perform concurrent and IO-free training state replication. To avoid the high overhead of start and initialization, we further propose an asynchronous coordination mechanism. Powered by the above innovations, Elan can provide high-performance (~1s) migration, scaling in and scaling out support with negligible runtime overhead (<3‰). For elastic training of ResNet-50 on ImageNet, Elan improves the time to solution by 20%. For elastic scheduling, with the help of Elan, resource utilization is improved by 21%+ and job pending time is reduced by 43%+.
Jidong Zhai, Baodong Wu, Xingcheng Zhang, Peng Sun 0006, Shengen Yan
ICDCS6
2020 GraphMP: I/O-Efficient Big Graph Analytics on a Single Commodity Machine
abstract
Recent studies showed that single-machine graph processing systems can be as highly competitive as cluster-based approaches on large-scale problems. While several out-of-core graph processing systems and computation models have been proposed, the high disk I/O overhead could significantly reduce performance in many practical cases. In this paper, we propose GraphMP to tackle big graph analytics on a single machine. GraphMP achieves low disk I/O overhead with three techniques. First, we design a vertex-centric sliding window (VSW) computation model to avoid reading and writing vertices on disk. Second, we propose a selective scheduling method to skip loading and processing unnecessary edge shards on disk. Third, we use a compressed edge cache mechanism to fully utilize the available memory of a machine to reduce the amount of disk accesses for edges. Extensive evaluations have shown that GraphMP could outperform existing single-machine out-of-core systems such as GraphChi, X-Stream and GridGraph by up to 30, and can be as highly competitive as distributed graph engines like Pregel+, PowerGraph and Chaos.
Peng Sun 0006, Yonggang Wen 0001, Ta Nguyen Binh Duong, Xiaokui Xiao
IEEE Trans. Big Data1
2018 Speeding-Up Age Estimation in Intelligent Demographics System via Network Optimization
abstract
Age estimation is a difficult task which requires the automatic detection and interpretation of facial features. Recently, Convolutional Neural Networks (CNNs) have made remarkable improvement on learning age patterns from benchmark datasets. However, for a face ``in the wild'' (from a video frame or Internet), the existing algorithms are not as accurate as for a frontal and neutral face. In addition, with the increasing number of in-the- wild aging data, the computation speed of existing deep learning platforms becomes another crucial issue. In this paper, we propose a high-efficient age estimation system with joint optimization of age estimation algorithm and deep learning system. Cooperated with the city surveillance network, this system can provide age group analysis for intelligent demographics. First, we build a three- tier fog computing architecture including an edge, a fog and a cloud layer, which directly processes age estimation from raw videos. Second, we optimize the age estimation algorithm based on CNNs with label distribution and K-L divergence distance embedded in the fog layer and evaluate the model on the latest wild aging dataset. Experimental results demonstrate that: 1. our system collects the demographics data dynamically at far-distance without contact, and makes the city population analysis automatically; and 2. the age model training has been speed-up without losing training progress or model quality. To our best knowledge, this is the first intelligent demographics system which has potential applications in improving the efficiency of smart cities and urban living.
Zhenzhen Hu 0004, Peng Sun 0006, Yonggang Wen 0001
ICC2
2018 MetaFlow: A Scalable Metadata Lookup Service for Distributed File Systems in Data Centers
abstract
In large-scale distributed file systems, efficient metadata operations are critical since most file operations have to interact with metadata servers first. In existing distributed hash table (DHT) based metadata management systems, the lookup service could be a performance bottleneck due to its significant CPU overhead. Our investigations showed that the lookup service could reduce system throughput by up to 70 percent, and increase system latency by a factor of up to 8 compared to ideal scenarios. In this paper, we present MetaFlow, a scalable metadata lookup service utilizing software-defined networking (SDN) techniques to distribute lookup workload over network components. MetaFlow tackles the lookup bottleneck problem by leveraging B-tree, which is constructed over the physical topology, to manage flow tables for SDN-enabled switches. Therefore, metadata requests can be forwarded to appropriate servers using only switches. Extensive performance evaluations in both simulations and testbed showed that MetaFlow increases system throughput by a factor of up to 3.2, and reduce system latency by a factor of up to 5 compared to DHT-based systems. We also deployed MetaFlow in a distributed file system, and demonstrated significant performance improvement.
Peng Sun 0006, Yonggang Wen 0001, Ta Nguyen Binh Duong, Haiyong Xie 0001
IEEE Trans. Big Data1
2017 GraphH: High Performance Big Graph Analytics in Small Clusters
abstract
It is common for real-world applications to analyze big graphs using distributed graph processing systems. Popular in-memory systems require an enormous amount of resources to handle big graphs. While several out-of-core approaches have been proposed for processing big graphs on disk, the high disk I/O overhead could significantly reduce performance. In this paper, we propose GraphH to enable high-performance big graph analytics in small clusters. Specifically, we design a two-stage graph partition scheme to evenly divide the input graph into partitions, and propose a GAB (Gather-Apply-Broadcast) computation model to make each worker process a partition in memory at a time. We use an edge cache mechanism to reduce the disk I/O overhead, and design a hybrid strategy to improve the communication performance. GraphH can efficiently process big graphs in small clusters or even a single commodity server. Extensive evaluations have shown that GraphH could be up to 7.8x faster compared to popular in-memory systems, such as Pregel+ and PowerGraph when processing generic graphs, and more than 100x faster than recently proposed out-of-core systems, such as GraphD and Chaos when processing big graphs.
Peng Sun 0006, Yonggang Wen 0001, Ta Nguyen Binh Duong, Xiaokui Xiao
CLUSTER1
2017 GraphMP: An Efficient Semi-External-Memory Big Graph Processing System on a Single Machine
abstract
Recent studies showed that single-machine graph processing systems can be as highly competitive as clusterbased approaches on large-scale problems. While several out-of-core graph processing systems and computation models have been proposed, the high disk I/O overhead could significantly reduce performance in many practical cases. In this paper, we propose GraphMP to tackle big graph analytics on a single machine. GraphMP achieves low disk I/O overhead with three techniques. First, we design a vertex-centric sliding window (VSW) computation model to avoid reading and writing vertices on disk. Second, we propose a selective scheduling method to skip loading and processing unnecessary edge shards on disk. Third, we use a compressed edge cache mechanism to fully utilize the available memory of a machine to reduce the amount of disk accesses for edges. Extensive evaluations have shown that GraphMP could outperform state-of-the-art systems such as GraphChi, X-Stream and GridGraph by 31.6x, 54.5x and 23.1x respectively, when running popular graph applications on a billion-vertex graph.
Peng Sun 0006, Yonggang Wen 0001, Ta Nguyen Binh Duong, Xiaokui Xiao
ICPADS1
2017 Towards Distributed Machine Learning in Shared Clusters: A Dynamically-Partitioned Approach
abstract
Many cluster management systems (CMSs) have been proposed to share a single cluster with multiple distributed computing systems. However, none of the existing approaches can handle distributed machine learning (ML) workloads given the following criteria: high resource utilization, fair resource allocation and low sharing overhead. To solve this problem, we propose a new CMS named Dorm, incorporating a dynamically-partitioned cluster management mechanism and an utilization-fairness optimizer. Specifically, Dorm uses the container-based virtualization technique to partition a cluster, runs one application per partition, and can dynamically resize each partition at application runtime for resource efficiency and fairness. Each application directly launches its tasks on the assigned partition without petitioning for resources frequently, so Dorm imposes flat sharing overhead. Extensive performance evaluations showed that Dorm could simultaneously increase the resource utilization by a factor of up to 2.32, reduce the fairness loss by a factor of up to 1.52, and speed up popular distributed ML applications by a factor of up to 2.72, compared to existing approaches. Dorm's sharing overhead is less than 5% in most cases.
Peng Sun 0006, Yonggang Wen 0001, Ta Nguyen Binh Duong, Shengen Yan
SMARTCOMP1
2016 Timed Dataflow: Reducing Communication Overhead for Distributed Machine Learning Systems
abstract
Many distributed machine learning (ML) systems exhibit high communication overhead when dealing with big data sets. Our investigations showed that popular distributed ML systems could spend about an order of magnitude more time on network communication than computation to train ML models containing millions of parameters. Such high communication overhead is mainly caused by two operations: pulling parameters and pushing gradients. In this paper, we propose an approach called Timed Dataflow (TDF) to deal with this problem via reducing network traffic using three techniques: a timed parameter storage system, a hybrid parameter filter and a hybrid gradient filter. In particular, the timed parameter storage technique and the hybrid parameter filter enable servers to discard unchanged parameters during the pull operation, and the hybrid gradient filter allows servers to drop gradients selectively during the push operation. Therefore, TDF could reduce the network traffic and communication time significantly. Extensive performance evaluations in a real testbed showed that TDF could reduce up to 77% and 79% of network traffic for the pull and push operations, respectively. As a result, TDF could speed up model training by a factor of up to 4 without sacrificing much accuracy for some popular ML models, compared to systems not using TDF.
Peng Sun 0006, Yonggang Wen 0001, Ta Nguyen Binh Duong, Shengen Yan
ICPADS1
2014 CREATE: Correlation enhanced traffic matrix estimation in Data Center Networks
abstract
Understanding the pattern of end-to-end traffic flows in Data Center Networks (DCNs) is essential to many DCN designs and operations (e.g., traffic engineering and load balancing). However, little research work has been done to obtain traffic information efficiently and yet accurately. Researchers often assume the availability of traffic tracing tools (e.g., OpenFlow) when their proposals require traffic information as input, but these tools may generate high monitoring overhead and consume significant switch resources even if they are available in a DCN. Although estimating the traffic matrix between origin-destination pairs using only basic switch SNMP counters is a mature practice in IP networks, traffic flows in DCNs are notoriously more irregular and volatile, while the large number of redundant routes in a DCN further complicates the situation. To this end, we propose to utilize the service placement logs for deducing the correlations among top-of-rack switches, and to leverage the uneven traffic distribution in DCNs for reducing the number of routes potentially used by a flow. These allow us to develop an efficient CoRrelation Enhanced trAffic maTrix Estimation (CREATE) method that achieves high accuracy. We compare CREATE with two existing representative methods through both experiments and simulations; the results strongly confirm the promising performance of CREATE.
Zhiming Hu 0001, Jun Luo 0001, Peng Sun 0006, Yonggang Wen 0001
Networking4
2013 Cloud3DView: an interactive tool for cloud data center operations
abstract
The emergence of cloud computing has promoted growing demand and rapid deployment of data centers. However, data center operations require a set of sophisticated skills (e.g., command-line-interface), resulting in a high operational cost. In this demo, to reduce the data center operational cost, we design and build a novel cloud data center management system, based on the concept of 3D gamification. In particular, we apply data visualization techniques to overlay operational status upon a data center 3D model, allowing the operators to monitor the real-time situation and control the data center from a friendly user interface. This demo highlights: (1)a data center 3D view from a First Person Shooter (FPS) camera, (2)a run-time presentation of visualized infrastructures information. Moreover, to improve the user experience, we employ cutting-edge HCI technologies from multi-touch, for remote access to Cloud3DView.
Jianxiong Yin, Peng Sun 0006, Yonggang Wen 0001, Hai-gang Gong, Ming Liu 0002, Xuelong Li 0001, Haipeng You, Jinqi Gao, Cynthia Lin
SIGCOMM2