Li Chen 0019

dblp:c/LiChen19 · DBLP profile ↗
← Back
52ranked-venue papers
10as first author
33since 2021 · last 2026
0000-0002-2300-6996ORCID · conflict

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

Systems, architecture and hardware · 26 · 4 first-author · 18 since 2021Computer networks · 12 · 3 first-author · 8 since 2021Artificial intelligence and machine learning · 4 · 2 since 2021Applied, interdisciplinary, general and emerging computing · 3 · 1 first-author · 2 since 2021Security and privacy · 2 · 1 since 2021Software engineering, systems software and programming languages · 1 · 1 since 2021Databases, data management, data science and information retrieval · 1 · 1 since 2021Graphics, computer vision, multimedia, augmented reality and games · 1 · 1 since 2021Human-computer interaction and ubiquitous computing · 1 · 1 since 2021
YearPublicationVenuePosition
2026 FedACT: Concurrent Federated Intelligence across Heterogeneous Data Sources
Md Sirajul Islam, Isabelle G. Chapman, N. I Md Ashafuddula, Xu Yuan 0001, Li Chen 0019, Nian-Feng Tzeng, Klara Nahrstedt
IPDPS5
2026 Resource Heterogeneity-Aware and Utilization-Enhanced Scheduling for Deep Learning Clusters
abstract
Scheduling deep learning (DL) models to train on powerful clusters with accelerators like GPUs and TPUs, presently falls short, either lacking fine-grained heterogeneity awareness or leaving resources substantially under-utilized. To fill this gap, we propose a novel task-level heterogeneity-aware scheduler for DL clusters, <italic xmlns:mml="http://www.w3.org/1998/Math/MathML" xmlns:xlink="http://www.w3.org/1999/xlink">Hadar</i>, based on an optimization framework able to boost cluster resource utilization. <italic xmlns:mml="http://www.w3.org/1998/Math/MathML" xmlns:xlink="http://www.w3.org/1999/xlink">Hadar</i> leverages the performance traits of DL jobs on a heterogeneous DL cluster to make scheduling decisions across both spatial and temporal dimensions. It characterizes the task-level performance heterogeneity for optimization and involves the primal-dual framework employing a dual subroutine, to solve the optimization problem and guide the scheduling design. Our trace-driven simulation with representative DL model training workloads demonstrates that <italic xmlns:mml="http://www.w3.org/1998/Math/MathML" xmlns:xlink="http://www.w3.org/1999/xlink">Hadar</i> accelerates the total training time duration by 1.20× when compared with its state-of-the-art heterogeneity-aware counterpart, Gavel. Further, our <italic xmlns:mml="http://www.w3.org/1998/Math/MathML" xmlns:xlink="http://www.w3.org/1999/xlink">Hadar</i> scheduler is enhanced to <italic xmlns:mml="http://www.w3.org/1998/Math/MathML" xmlns:xlink="http://www.w3.org/1999/xlink">Hadar</i>E by forking each job into multiple copies to let a job train concurrently on heterogeneous GPUs resided on separate available cluster nodes (i.e., machines or servers) for resource utilization enhancement. <italic xmlns:mml="http://www.w3.org/1998/Math/MathML" xmlns:xlink="http://www.w3.org/1999/xlink">Hadar</i>E is evaluated extensively on physical DL clusters for comparison with <italic xmlns:mml="http://www.w3.org/1998/Math/MathML" xmlns:xlink="http://www.w3.org/1999/xlink">Hadar</i> and Gavel. With substantial enhancement in cluster resource utilization (by 1.45×), <italic xmlns:mml="http://www.w3.org/1998/Math/MathML" xmlns:xlink="http://www.w3.org/1999/xlink">Hadar</i>E exhibits considerable speed-ups in DL model training, reducing the total training time duration by 50% (or 80%) on an Amazon’s AWS (or our lab) cluster, while producing trained DL models with consistently better inference quality than those trained by <italic xmlns:mml="http://www.w3.org/1998/Math/MathML" xmlns:xlink="http://www.w3.org/1999/xlink">Hadar</i>.
Abeda Sultana, Nabin Pakka, Fei Xu 0009, Xu Yuan 0001, Li Chen 0019, Nian-Feng Tzeng
IEEE Trans. Computers5
2026 SSE-TSR: An Approach to Integrate Secondary Structure Elements Into Triangular Spatial Relationships for Protein Classification
abstract
Protein structures are fundamental to understanding biological function, yet many detailed similarities remain hidden from conventional alignment-based or 3D superposition methods. Triangular Spatial Relationship (TSR) offers an alignment-free encoding of backbone geometry; however, classical TSR ignores the context of secondary structure elements (SSEs), such as helices, strands, and coils. To address this, we introduce SSE-TSR, which enriches each TSR key by categorizing it into one of 18 helix-strand-coil combination labels derived from DSSP-style annotations in PDB HELIX/SHEET records. By mapping the protein representation involving SSE-TSR keys into a sparse tensor, SSE-TSR compactly captures both tertiary geometry and local secondary motifs. We evaluated SSE-TSR on four datasets, two structural (CATH-based, 9.2K; SCOP-based, 7.0K) and two functional (published, 7.8K; new, 7.2K), using a 3D convolutional neural network. On structure-based tasks, SSETSR noticeably boosts accuracy from 96.00% to 98.33% (CATHbased) and from 95.46% to 99.00% (SCOP-based). On functional tasks, it yields modest yet consistent gains (e.g., from 99.41% to 99.50% and 95.83% to 98.83%). Comparisons to Foldseek confirm competitive accuracy across diverse tasks. Additionally, the sparse tensor representation enables memory-efficient handling of large-scale datasets, making SSE-TSR practical for extensive bioinformatics analyses. These results demonstrate SSE-TSR as a scalable, interpretable, and robust method, enhancing protein classification and structural bioinformatics.
Poorya Khajouie, Titli Sarkar, Krishna Rauniyar, Li Chen 0019, Wu Xu, Vijay Raghavan 0001
IEEE Trans. Comput. Biol. Bioinform.4
2025 Espresso: Cost-Efficient Large Model Training by Exploiting GPU Heterogeneity in the Cloud
Qiannan Zhou, Fei Xu 0009, Lingxuan Weng, Ruixing Li, Li Chen 0019, Zhi Zhou 0006, Fangming Liu
INFOCOM6
2025 SEAFL: Enhancing Efficiency in Semi-Asynchronous Federated Learning Through Adaptive Aggregation and Selective Training
abstract
Federated Learning (FL) is a promising distributed machine learning framework that allows collaborative learning of a global model across decentralized devices without uploading their local data. However, in real-world FL scenarios, the conventional synchronous FL mechanism suffers from inefficient training caused by slow-speed devices, commonly known as stragglers, especially in heterogeneous communication environments. Though asynchronous FL effectively tackles the efficiency challenge, it induces substantial system overheads and model degradation. Striking for a balance, semi-asynchronous FL has gained increasing attention, while still suffering from the open challenge of stale models, where newly arrived updates are calculated based on outdated weights that easily hurt the convergence of the global model. In this paper, we present SEAFL, a novel FL framework designed to mitigate both the straggler and the stale model challenges in semi-asynchronous FL. SEAFL dynamically assigns weights to uploaded models during aggregation based on their staleness and importance to the current global model. We theoretically analyze the convergence rate of SEAFL and further enhance the training efficiency with an extended variant that allows partial training on slower devices, enabling them to contribute to global aggregation while reducing excessive waiting times. We evaluate the effectiveness of SEAFL through extensive experiments on three benchmark datasets. The experimental results demonstrate that SEAFL outperforms its closest counterpart by up to$\sim 22 \%$in terms of the wall-clock training time required to achieve target accuracy.
Md Sirajul Islam, Sanjeev Panta, Fei Xu 0009, Xu Yuan 0001, Li Chen 0019, Nian-Feng Tzeng
IPDPS5
2025 Opara: Exploiting Operator Parallelism for Expediting DNN Inference on GPUs
abstract
GPUs have become thedefactohardware devices for accelerating Deep Neural Network (DNN) inference workloads. However, the conventionalsequential execution mode of DNN operatorsin mainstream deep learning frameworks cannot fully utilize GPU resources, even with the operator fusion enabled, due to the increasing complexity of model structures and a greater diversity of operators. Moreover, theinadequate operator launch orderin parallelized execution scenarios can lead to GPU resource wastage and unexpected performance interference among operators. In this paper, we proposeOpara, a resource- and interference-aware DNNOperatorparallel scheduling framework to accelerate DNN inference on GPUs. Specifically,Oparafirst employsCUDA StreamsandCUDA Graphtoparallelizethe execution of multiple operators automatically. To further expedite DNN inference,Oparaleverages the resource demands of operators to judiciously adjust the operator launch order on GPUs, overlapping the execution of compute-intensive and memory-intensive operators. We implement and open source a prototype ofOparabased on PyTorch in anon-intrusivemanner. Extensive prototype experiments with representative DNN and Transformer-based models demonstrate thatOparaoutperforms the default sequentialCUDA Graphin PyTorch and the state-of-the-art operator parallelism systems by up to$1.68\boldsymbol{\times}$and$1.29\boldsymbol{\times}$, respectively, yet with acceptable runtime overhead.
Aodong Chen, Fei Xu 0009, Li Han 0001, Li Chen 0019, Zhi Zhou 0006, Fangming Liu
IEEE Trans. Computers5
2025 Regional Weather Variable Predictions by Machine Learning With Near-Surface Observational and Atmospheric Numerical Data
abstract
Accurate and timely regional weather prediction is vital for sectors dependent on weather-related decisions. Traditional prediction methods, based on atmospheric equations, often struggle with coarse temporal resolutions and inaccuracies. This article presents a novel machine learning (ML) model, called Micro-Macro (MiMa), that integrates both near-surface observational data from Kentucky Mesonet stations (collected every 5 min, known as Micro data) and hourly atmospheric numerical outputs (termed as Macro data) for fine-resolution weather forecasting. The MiMa model employs an encoder-decoder transformer structure, with two encoders for processing multivariate data from both datasets and a decoder for forecasting weather variables over short time horizons. Each instance of the MiMa model, called a modelet, predicts the values of a specific weather parameter at an individual mesonet station. The approach is extended with Regional MiMa (Re-MiMa) modelets, which are designed to predict weather variables at ungauged locations by training on multivariate data from a few representative stations in a region, tagged with their elevations. Re-MiMa can provide highly accurate predictions across an entire region, even in areas without observational stations. Experimental results show that MiMa significantly outperforms current models, with Re-MiMa offering precise short-term forecasts for ungauged locations, marking a significant advancement in weather forecasting accuracy and applicability.
Yihe Zhang 0001, Bryce Turney, Purushottam Sigdel, Xu Yuan 0001, Eric Rappin, Adrian Lago, Sytske K. Kimball, Li Chen 0019, Paul J. Darby, Lu Peng 0001, Sercan Aygün, Yazhou Tu, M. Hassan Najafi, Nian-Feng Tzeng
IEEE Trans. Geosci. Remote. Sens.8
2024 FedFair3: Unlocking Threefold Fairness in Federated Learning
abstract
Federated Learning (FL) is an emerging paradigm in machine learning without exposing clients' raw data. In practical scenarios with numerous clients, encouraging fair and efficient client participation in federated learning is of utmost importance, which is also challenging given the heterogeneity in data distribution and device properties. Existing works have proposed different client-selection methods that consider fairness; however, they fail to select clients with high utilities while simultaneously achieving fair accuracy levels. In this paper, we propose a fair client-selection approach that unlocks threefold fairness in federated learning. In addition to having a fair client-selection strategy, we enforce an equitable number of rounds for client participation and ensure a fair accuracy distribution over the clients. The experimental results demonstrate that FedFair3, in comparison to the state-of-the-art baselines, achieves 18.15% less accuracy variance on the IID data and 54.78 % on the non-IID data, without decreasing the global accuracy. Furthermore, it shows 24.36% less wall-clock training time on average.
Simin Javaherian, Sanjeev Panta, Shelby Williams, Md Sirajul Islam, Li Chen 0019
ICC5
2024 FedClust: Tackling Data Heterogeneity in Federated Learning through Weight-Driven Client Clustering
abstract
Federated learning (FL) is an emerging distributed machine learning paradigm that enables collaborative training of machine learning models over decentralized devices without exposing their local data. One of the major challenges in FL is the presence of uneven data distributions across client devices, violating the well-known assumption of independent-and-identically-distributed (IID) training samples in conventional machine learning. To address the performance degradation issue incurred by such data heterogeneity, clustered federated learning (CFL) shows its promise by grouping clients into separate learning clusters based on the similarity of their local data distributions. However, state-of-the-art CFL approaches require a large number of communication rounds to learn the distribution similarities during training until the formation of clusters is stabilized. Moreover, some of these algorithms heavily rely on a predefined number of clusters, thus limiting their flexibility and adaptability. In this paper, we propose FedClust, a novel approach for CFL that leverages the correlation between local model weights and the data distribution of clients. FedClust groups clients into clusters in a one-shot manner by measuring the similarity degrees among clients based on the strategically selected partial weights of locally trained models. We conduct extensive experiments on four benchmark datasets with different non-IID data settings. Experimental results demonstrate that FedClust achieves higher model accuracy up to ∼ 45% as well as faster convergence with a significantly reduced communication cost up to 2.7 × compared to its state-of-the-art counterparts.
Md Sirajul Islam, Simin Javaherian, Fei Xu 0009, Xu Yuan 0001, Li Chen 0019, Nian-Feng Tzeng
ICPP5
2024 Periscoping: Private Key Distribution for Large-Scale Mixnets
abstract
Mix networks, or mixnets, are one of the fundamental building blocks of anonymity systems. To defend against epistemic attacks, existing free-route mixnet designs require all clients to maintain a consistent, up-to-date view of the entire key directory. This, however, inevitably raises the performance concern under system scale-out: in a larger mixnet, a client will consume more bandwidth for updating keys in the background.This paper presents Periscoping, a key distribution protocol for mixnets at scale. Periscoping relaxes the download-all requirement for clients. Instead, it allows a client to selectively download a constant number of entries of the key directory, while guaranteeing the privacy of selections. Periscoping achieves this goal via a novel Private Information Retrieval scheme, constructed based on constrained Pseudorandom Functions. Moreover, the protocol is integrated seamlessly into the mixnet operations, readily applicable to existing mixnet systems as an extension at a minimal cost. Our experiments show that, with millions of mixes, it can reduce the traffic load of a mixnet by orders of magnitude, at a minor computational and bandwidth overhead.
Shuhao Liu 0001, Li Chen 0019, Yuanzhong Fu
INFOCOM2
2024 Hadar: Heterogeneity-Aware Optimization-Based Online Scheduling for Deep Learning Cluster
abstract
With the wide adoption of deep neural network (DNN) models for various applications, enterprises, and cloud providers have built deep learning clusters and increasingly deployed specialized accelerators, such as GPUs and TPUs, for DNN training jobs. To arbitrate cluster resources among multi-user jobs, existing schedulers fall short, either lacking fine-grained heterogeneity awareness or hardly generalizable to various scheduling policies. To fill this gap, we propose a novel design of a task-level heterogeneity-aware scheduler, Hadar, based on an online optimization framework that can express other scheduling algorithms. Hadar leverages the performance traits of DNN jobs on a heterogeneous cluster, characterizes the task-level performance heterogeneity in the optimization problem, and makes scheduling decisions across both spatial and temporal dimensions. The primal-dual framework is employed, with our design of a dual subroutine, to solve the optimization problem and guide the scheduling design. Extensive trace-driven simulations with representative DNN models have been conducted to demonstrate that Hadar improves the average job completion time (JCT) by 3× over an Apache YARN-based resource manager used in production. Moreover, Hadar outperforms Gavel [1], the state-of-the-art heterogeneity-aware scheduler, by 2.5× for the average JCT, shortens the queuing delay by 13%, and improves FTF (Finish-Time-Fairness) by 1.5%.
Abeda Sultana, Fei Xu 0009, Xu Yuan 0001, Li Chen 0019, Nian-Feng Tzeng
IPDPS4
2024 HarmonyBatch: Batching multi-SLO DNN Inference with Heterogeneous Serverless Functions
abstract
Deep Neural Network (DNN) inference on serverless functions is gaining prominence due to its potential for substantial budget savings. Existing works on serverless DNN inference solely optimize batching requests from one application with a single Service Level Objective (SLO) on CPU functions. However, production serverless DNN inference traces indicate that the request arrival rate of applications is surprisingly low, which inevitably causes a long batching time and SLO violations. Hence, there is an urgent need for batching multiple DNN inference requests with diverse SLOs (i.e., multi-SLO DNN inference) in serverless platforms. Moreover, the potential performance and cost benefits of deploying heterogeneous (i.e., CPU and GPU) functions for DNN inference have received scant attention.In this paper, we present HarmonyBatch, a cost-efficient resource provisioning framework designed to achieve predictable performance for multi-SLO DNN inference with heterogeneous serverless functions. Specifically, we construct an analytical performance and cost model of DNN inference on both CPU and GPU functions, by explicitly considering the GPU time-slicing scheduling mechanism and request arrival rate distribution. Based on such a model, we devise a two-stage merging strategy in HarmonyBatch to judiciously batch the multi-SLO DNN inference requests into application groups. It aims to minimize the budget of function provisioning for each application group while guaranteeing diverse performance SLOs of inference applications. We have implemented a prototype of HarmonyBatch on Alibaba Cloud Function Compute. Extensive prototype experiments with representative DNN inference workloads demonstrate that HarmonyBatch can provide predictable performance to serverless DNN inference workloads while reducing the monetary cost by up to 82.9% compared to the state-of-the-art methods.
Fei Xu 0009, Yikun Gu, Li Chen 0019, Fangming Liu, Zhi Zhou 0006
IWQoS4
2024 EchoSensor: Fine-grained Ultrasonic Sensing for Smart Home Intrusion Detection
abstract
This article presents the design and implementation of a novel intrusion detection system, called EchoSensor, which leverages speakers and microphones in smart home devices to capture human gait patterns for individual identification. EchoSensor harnesses the speaker to send inaudible acoustic signals (around 20 kHz) and utilizes the microphone to capture the reflected signals. As the reflected signals have unique variations in the Doppler shift respective to the gaits of different people, EchoSensor is able to profile human gait patterns from the generated spectrograms. To mine the gait information, we first propose a two-stage interference cancellation scheme to remove the background noise and environmental interference, followed by a new method to detect the starting point of walking and estimate the gait cycle time. We then perform the fine-grained analysis of the spectrograms to extract a series of features. In the end, machine learning is employed to construct an identifier for individual recognition. We implement the EchoSensor system and deploy it under different household environments to conduct intrusion detection tasks. Extensive experimental results have demonstrated that EchoSensor can achieve the averaged Intruder Gait Detection Rate (IDR) and True Family Member Gait Detection Rate (TFR) of 92.7% and 91.9%, respectively.
Changlai Du, Jiadong Lou, Li Chen 0019, Xu Yuan 0001
ACM Trans. Sens. Networks4
2024 Room-scale Location Trace Tracking via Continuous Acoustic Waves
abstract
The increasing prevalence of smart devices spurs the development of emerging indoor localization technologies for supporting diverse personalized applications at home. Given marked drawbacks of popular chirp signal-based approaches, we aim at developing a novel device-free localization system via the continuous wave of the inaudible frequency. To achieve this goal, solutions are developed for fine-grained analyses, able to precisely locate moving human traces in the room-scale environment. In particular, a smart speaker is controlled to emit continuous waves at inaudible 20kHz , with a co-located microphone array to record their Doppler reflections for localization. We first develop solutions to remove potential noises and then propose a novel idea by slicing signals into a set of narrowband signals, each of which is likely to include at most one body segment’s reflection. Different from previous studies, which take original signals themselves as the baseband, our solutions employ the Doppler frequency of a narrowband signal to estimate the velocity first and apply it to get the accurate baseband frequency, which permits a precise phase measurement after I-Q (i.e., in-phase and quadrature) decomposition. A signal model is then developed, able to formulate the phase with body segment’s velocity, range, and angle. We next develop novel solutions to estimate the motion state in each narrowband signal, cluster the motion states for different body segments corresponding to the same person, and locate the moving traces while mitigating multi-path effects. Our system is implemented with commodity devices in room environments for performance evaluation. The experimental results exhibit that our system can conduct effective localization for up to three persons in a room, with the average errors of 7.49 cm for a single person, with 24.06 cm for two persons, with 51.15 cm for three persons.
Xu Yuan 0001, Jiadong Lou, Li Chen 0019, Hao Wang 0022, Nian-Feng Tzeng
ACM Trans. Sens. Networks4
2024 Tetris: Proactive Container Scheduling for Long-Term Load Balancing in Shared Clusters
abstract
Long-running containerized workloads (e.g., machine learning), which typically showtime-varyingpatterns, are increasingly prevailing in shared production clusters. To improve workload performance, current schedulers mainly focus on optimizingshort-termbenefits of cluster load balancing orinitial container placementon servers. However, this would inevitably bring manyinvalid migrations(i.e., containers are migrated back and forth among servers over a short time window), leading to significant service level objective (SLO) violations. This paper introducesTetris, amodel predictive control(MPC)-based container scheduling strategy to proactively migrate long-running workloads for cluster load balancing. Specifically, we first build a discrete-time dynamic model forlong-termoptimization of container scheduling. To solve such an optimization problem,Tetristhen employs two main components: (1) a container resource predictor, which leverages time-series analysis approaches to accurately predict the container resource consumption; (2) an MPC-based container scheduler that jointly optimizes the cluster load balancing and container migration costover a certain sliding time window. We implement and open source a prototype ofTetrisbased on K8s. Extensive prototype experiments and trace-driven simulations demonstrate thatTetriscan improve the cluster load balancing degree by up to 77.8% without incurring any SLO violations, compared to the state-of-the-art container scheduling strategies.
Fei Xu 0009, Xiyue Shen, Shuohao Lin, Li Chen 0019, Zhi Zhou 0006, Fen Xiao, Fangming Liu
IEEE Trans. Serv. Comput.4
2023 MMST-ViT: Climate Change-aware Crop Yield Prediction via Multi-Modal Spatial-Temporal Vision Transformer
abstract
Precise crop yield prediction provides valuable information for agricultural planning and decision-making processes. However, timely predicting crop yields remains challenging as crop growth is sensitive to growing season weather variation and climate change. In this work, we develop a deep learning-based solution, namely Multi-Modal Spatial-Temporal Vision Transformer (MMST-ViT), for predicting crop yields at the county level across the United States, by considering the effects of short-term meteorological variations during the growing season and the long-term climate change on crops. Specifically, our MMST-ViT consists of a Multi-Modal Transformer, a Spatial Transformer, and a Temporal Transformer. The Multi-Modal Transformer leverages both visual remote sensing data and short-term meteorological data for modeling the effect of growing season weather variations on crop growth. The Spatial Transformer learns the high-resolution spatial dependency among counties for accurate agricultural tracking. The Temporal Transformer captures the long-range temporal dependency for learning the impact of long-term climate change on crops. Meanwhile, we also devise a novel multi-modal contrastive learning technique to pre-train our model without extensive human supervision. Hence, our MMST-ViT captures the impacts of both short-term weather variations and long-term climate change on crops by leveraging both satellite images and meteorological data. We have conducted extensive experiments on over 200 counties in the United States, with the experimental results exhibiting that our MMST-ViT outperforms its counterparts under three performance metrics of interest. Our dataset and code are available at https://github.com/fudong03/MMST-ViT.
Fudong Lin, Summer Crawford, Kaleb Guillot, Yihe Zhang 0001, Xu Yuan 0001, Li Chen 0019, Shelby Williams, Robert Minvielle, Xiangming Xiao, Drew Gholson, Nicolas Ashwell, Tri Setiyono, Brenda Tubana, Lu Peng 0001, Magdy A. Bayoumi, Nian-Feng Tzeng
ICCV7
2023 ACTS: Autonomous Cost-Efficient Task Orchestration for Serverless Analytics
abstract
Serverless computing has become increasingly popular for cloud applications, due to its compelling properties of high-level abstractions, lightweight runtime, high elasticity and pay-per-use billing. In this revolutionary computing paradigm shift, challenges arise when adapting data analytics applications to the serverless environment, due to the lack of support for efficient state sharing, which attract ever-growing research attention. In this paper, we aim to exploit the advantages of task-level orchestration and fine-grained resource provisioning for data analytics on serverless platforms, with the hope of fulfilling the promise of serverless deployment to the maximum extent. To this end, we present ACTS, an autonomous cost-efficient task orchestration framework for serverless analytics. ACTS judiciously schedules and coordinates function tasks to mitigate cold-start latency and state sharing overhead. In addition, ACTS explores the optimization space of fine-grained workload distribution and function resource configuration for cost efficiency. We have deployed and implemented ACTS on AWS Lambda, evaluated with various data analytics workloads. Results from extensive experiments demonstrate that ACTS achieves up to 98% monetary cost reduction while maintaining superior job completion time performance, in comparison with the state-of-the-art baselines.
Jananie Jarachanthan, Li Chen 0019, Fei Xu 0009
IWQoS2
2023 spotDNN: Provisioning Spot Instances for Predictable Distributed DNN Training in the Cloud
abstract
Distributed Deep Neural Network (DDNN) training on cloud spot instances is increasingly compelling as it can significantly save the user budget. To handle unexpected instance revocations, provisioning a heterogeneous cluster using the asynchronous parallel mechanism becomes the dominant method for DDNN training with spot instances. However, blindly provisioning a cluster of spot instances can easily result in unpre-dictable DDNN training performance, mainly because bottlenecks occur on the parameter server network bandwidth and PCIe bandwidth resources, as well as the inadequate cluster heterogeneity. To address the challenges above, we propose spotDNN, a heterogeneity-aware spot instance provisioning framework that provides predictable performance for DDNN training in the cloud. By explicitly considering the contention for bottle-neck resources, we first build an analytical performance model of DDNN training in heterogeneous clusters. It leverages the weighted average batch size and convergence coefficient to quantify the DDNN training loss in heterogeneous clusters. Through a lightweight workload profiling, we further design a cost-efficient instance provisioning strategy which incorporates the bounds calculation and sliding window techniques to effectively guarantee the training performance service level objectives (SLOs). We have implemented a prototype of spotDNN and conducted extensive experiments on Amazon EC2. Experiment results show that spotDNN can deliver predictable DDNN training performance while reducing the monetary cost by up to 68.1% compared to the existing solutions, yet with acceptable runtime overhead.
Ruitao Shang, Fei Xu 0009, Zhuoyan Bai, Li Chen 0019, Zhi Zhou 0006, Fangming Liu
IWQoS4
2023 iGniter: Interference-Aware GPU Resource Provisioning for Predictable DNN Inference in the Cloud
abstract
GPUs are essential to accelerating the latency-sensitive deep neural network (DNN) inference workloads in cloud datacenters. To fully utilize GPU resources,spatial sharingof GPUs among co-located DNN inference workloads becomes increasingly compelling. However, GPU sharing inevitably bringssevere performance interferenceamong co-located inference workloads, as motivated by an empirical measurement study of DNN inference on EC2 GPU instances. While existing works on guaranteeing inference performance service level objectives (SLOs) focus on eithertemporal sharingof GPUs orreactiveGPU resource scaling and inference migration techniques, how toproactivelymitigate such severe performance interference has received comparatively little attention. In this paper, we proposeiGniter, aninterference-awareGPU resource provisioning framework for cost-efficiently achieving predictable DNN inference in the cloud.iGniteris comprised of two key components: (1) alightweightDNN inference performance model, which leverages the system and workload metrics that are practically accessible to capture the performance interference; (2) Acost-efficientGPU resource provisioning strategy thatjointlyoptimizes the GPU resource allocation and adaptive batching based on our inference performance model, with the aim of achieving predictable performance of DNN inference workloads. We implement a prototype ofiGniterbased on the NVIDIA Triton inference server hosted on EC2 GPU instances. Extensive prototype experiments on four representative DNN models and datasets demonstrate thatiGnitercan guarantee the performance SLOs of DNN inference workloads with practically acceptable runtime overhead, while saving the monetary cost by up to$25\%$in comparison to the state-of-the-art GPU resource provisioning strategies.
Fei Xu 0009, Jianian Xu, Li Chen 0019, Ruitao Shang, Zhi Zhou 0006, Fangming Liu
IEEE Trans. Parallel Distributed Syst.4
2022 An Interactive Visualization System for Streaming Data Online Exploration
Fengzhou Liang, Fang Liu 0002, Tongqing Zhou, Yunhai Wang, Li Chen 0019
MobiQuitous5
2022 λDNN: Achieving Predictable Distributed DNN Training With Serverless Architectures
abstract
Serverless computing is becoming a promising paradigm for Distributed Deep Neural Network (DDNN) training in the cloud, as it allows users to decompose complex model training into a number offunctionswithout managing virtual machines or servers. Though provided with a simpler resource interface (i.e., function number and memory size), inadequate function resource provisioning (either under-provisioning or over-provisioning) easily leads tounpredictableDDNN training performance in serverless platforms. Our empirical studies on AWS Lambda indicate that, suchunpredictable performanceof serverless DDNN training is mainly caused by the resource bottleneck of Parameter Servers (PS) and small local batch size. In this article, we design and implement$\lambda$λDNN, a cost-efficient function resource provisioning framework to provide predictable performance for serverless DDNN training workloads, while saving the budget of provisioned functions. Leveraging the PS network bandwidth and function CPU utilization, we build alightweightanalytical DDNN training performance model to enable our design of$\lambda$λDNNresource provisioning strategy, so as to guarantee DDNN training performance with serverless functions. Extensive prototype experiments on AWS Lambda and complementary trace-driven simulations demonstrate that,$\lambda$λDNNcan deliver predictable DDNN training performance and save the monetary cost of function resources by up to 66.7 percent, compared with the state-of-the-art resource provisioning strategies, yet with an acceptable runtime overhead.
Fei Xu 0009, Yiling Qin, Li Chen 0019, Zhi Zhou 0006, Fangming Liu
IEEE Trans. Computers3
2022 Optimizing Network Transfers for Data Analytic Jobs Across Geo-Distributed Datacenters
abstract
It has become a recent trend that large volumes of data are generated, stored, and processed across geographically distributed datacenters. When popular data parallel frameworks, such as MapReduce and Spark, are employed to process such geo-distributed data, optimizing the network transfer in communication stages becomes increasingly crucial to application performance, as the inter-datacenter links have much lower bandwidth than intra-datacenter links. In this article, we focus on exploiting the flexibility of multi-path routing for inter-datacenter flows of data analytic jobs, with the hope of better utilizing inter-datacenter links and thus improve job performance. We design an optimal multi-path routing and scheduling strategy to achieve the best possible network performance for all concurrent jobs, based on our formulation of an optimization problem that can be transformed into an equivalent linear programming (LP) problem to be efficiently solved. As a highlight of this article, we have implemented our proposed algorithm in the controller of an application-layer software-defined inter-datacenter overlay testbed, designed to provide transfer optimization service for Spark jobs. With extensive evaluations of our real-world implementation on Google Cloud, we have shown convincing evidence that our optimal multi-path routing and scheduling strategies have achieved significant improvements in terms of job performance.
Li Chen 0019, Shuhao Liu 0001, Baochun Li
IEEE Trans. Parallel Distributed Syst.1
2022 Astrea: Auto-Serverless Analytics Towards Cost-Efficiency and QoS-Awareness
abstract
With the ability to simplify the code deployment with one-click upload and lightweight execution, serverless computing has emerged as a promising paradigm with increasing popularity. However, there remain open challenges when adapting data-intensive analytics applications to the serverless context, in which users ofserverless analyticsencounter the difficulty in coordinating computation across different stages and provisioning resources in a large configuration space. This paper presents our design and implementation ofAstrea, which configures and orchestrates serverless analytics jobs in an autonomous manner, while taking into account flexibly-specified user requirements.Astrearelies on the modeling of performance and cost which characterizes the intricate interplay among multi-dimensional factors (e.g., function memory size, degree of parallelism at each stage). We formulate an optimization problem based on user-specific requirements towards performance enhancement or cost reduction, and develop a set of algorithms based on graph theory to obtain the optimal job execution. We deployAstreain the AWS Lambda platform and conduct real-world experiments over representative benchmarks, including Big Data analytics and machine learning workloads, at different scales. Extensive results demonstrate thatAstreacan achieve the optimal execution decision for serverless data analytics, in comparison with various provisioning and deployment baselines. For example, when compared with three provisioning baselines,Astreamanages to reduce the job completion time by 21% to 69% under a given budget constraint, while saving cost by 20% to 84% without violating performance requirements.
Jananie Jarachanthan, Li Chen 0019, Fei Xu 0009, Bo Li 0001
IEEE Trans. Parallel Distributed Syst.2
2022 Eiffel: Efficient and Fair Scheduling in Adaptive Federated Learning
abstract
Emerging machine learning (ML) technologies, in combination with the increasing computational power of mobile devices, lead to the extensive adoption of ML-based applications. Different from conventional model training that needs to collect all the user data in centralized cloud servers, federated learning (FL) has recently drawn increasing research attention as it enables privacy-preserving model training. With FL, decentralized edge devices in participation, train their model copies locally over their siloed datasets, and periodically synchronize the model parameters. However, model training is computationally extensive which easily drains the battery of mobile devices. In addition, due to the uneven distribution of siloed datasets, the shared model may become biased. To address theefficiencyandfairnessconcerns in a resource-constrained federated learning setting, in this paper, we proposeEiffelto judiciously select mobile devices to participate in the global model aggregation, and adaptively adjust the frequency of local and global model updates.Eiffelaims to make scheduling and coordination for the federated learning towards both resource efficiency and model fairness. We have conducted theoretical analysis ofEiffelfrom the perspectives of fairness and convergence. Extensive experiments with a wide variety of real-world datasets and models, both on a networked prototype system and in a larger-scale simulated environment, have demonstrated that while maintaining similar accuracy performance,Eiffeloutperforms existing baselines with respect to reducing communication overhead by up to 6× for higher efficiency and improving the fairness metric by up to 57% compared to the state-of-the-art algorithms.
Abeda Sultana, Md. Mainul Haque, Li Chen 0019, Fei Xu 0009, Xu Yuan 0001
IEEE Trans. Parallel Distributed Syst.3
2021 Reverse Attack: Black-box Attacks on Collaborative Recommendation
abstract
Collaborative filtering (CF) recommender systems have been extensively developed and widely deployed in various social websites, promoting products or services to the users of interest. Meanwhile, work has been attempted at poisoning attacks to CF recommender systems for distorting the recommend results to reap commercial or personal gains stealthily. While existing poisoning attacks have demonstrated their effectiveness with the offline social datasets, they are impractical when applied to the real setting on online social websites. This paper develops a novel and practical poisoning attack solution toward the CF recommender systems without knowing involved specific algorithms nor historical social data information a priori. Instead of directly attacking the unknown recommender systems, our solution performs certain operations on the social websites to collect a set of sampling data for use in constructing a surrogate model for deeply learning the inherent recommendation patterns. This surrogate model can estimate the item proximities, learned by the recommender systems. By attacking the surrogate model, the corresponding solutions (for availability and target attacks) can be directly migrated to attack the original recommender systems. Extensive experiments validate the generated surrogate model's reproductive capability and demonstrate the effectiveness of our attack upon various CF recommender algorithms.
Yihe Zhang 0001, Xu Yuan 0001, Jin Li 0002, Jiadong Lou, Li Chen 0019, Nian-Feng Tzeng
CCS5
2021 AMPS-Inf: Automatic Model Partitioning for Serverless Inference with Cost Efficiency
abstract
The salient pay-per-use nature of serverless computing has driven its continuous penetration as an alternative computing paradigm for various workloads. Yet, challenges arise and remain open when shifting machine learning workloads to the serverless environment. Specifically, the restriction on the deployment size over serverless platforms combining with the complexity of neural network models makes it difficult to deploy large models in a single serverless function. In this paper, we aim to fully exploit the advantages of the serverless computing paradigm for machine learning workloads targeting at mitigating management and overall cost while meeting the response-time Service Level Objective (SLO). We design and implement AMPS-Inf, an autonomous framework customized for model inferencing in serverless computing. Driven by the cost-efficiency and timely-response, our proposed AMPS-Inf automatically generates the optimal execution and resource provisioning plans for inference workloads. The core of AMPS-Inf relies on the formulation and solution of a Mixed-Integer Quadratic Programming problem for model partitioning and resource provisioning with the objective of minimizing cost without violating response time SLO. We deploy AMPS-Inf on the AWS Lambda platform, evaluate with the state-of-the-art pre-trained models in Keras including ResNet50, Inception-V3 and Xception, and compare with Amazon SageMaker and three baselines. Experimental results demonstrate that AMPS-Inf achieves up to 98% cost saving without degrading response time performance.
Jananie Jarachanthan, Li Chen 0019, Fei Xu 0009, Bo Li 0001
ICPP2
2021 Accelerated Device Placement Optimization with Contrastive Learning
abstract
With the ever-increasing size and complexity of deep neural network models, it is difficult to fit and train a complete copy of the model on a single computational device with limited capability. Therefore, large neural networks are usually trained on a mixture of devices, including multiple CPUs and GPUs, of which the computational speed and efficiency are drastically affected by how these models are partitioned and placed on the devices. In this paper, we propose Mars, a novel design to find efficient placements for large models. Mars leverages a self-supervised graph neural network pre-training framework to generate node representations for operations, which is able to capture the topological properties of the computational graph. Then, a sequence-to-sequence neural network is applied to split large models into small segments so that Mars can predict the placements sequentially. Novel optimizations have been applied in the placer design to achieve the best possible performance in terms of the time needed to complete training the agent for placing models with very large sizes. We deployed and evaluated Mars on benchmarks involving Inception-V3, GNMT, and BERT models. Extensive experimental results show that Mars can achieve up to 27.2% and 2.7% speedup of per-step training time than the state-of-the-art for GNMT and BERT models, respectively. We also show that with self-supervised graph neural network pre-training, our design achieves the fastest speed in discovering the optimal placement for Inception-V3.
Li Chen 0019, Baochun Li
ICPP2
2021 Prophet: Speeding up Distributed DNN Training with Predictable Communication Scheduling
abstract
Optimizing performance for Distributed Deep Neural Network (DDNN) training has recently become increasingly compelling, as the DNN model gets complex and the training dataset grows large. While existing works on communication scheduling mostly focus on overlapping the computation and communication to improve DDNN training performance, the GPU and network resources are still under-utilized in DDNN training clusters. To tackle this issue, in this paper, we design and implement a predictable communication scheduling strategy named Prophet to schedule the gradient transfer in an adequate order, with the aim of maximizing the GPU and network resource utilization. Leveraging our observed stepwise pattern of gradient transfer start time, Prophet first uses the monitored network bandwidth and the profiled time interval among gradients to predict the appropriate number of gradients that can be grouped into blocks. Then, these gradient blocks can be transferred one by one to guarantee high utilization of GPU and network resources while ensuring the priority of gradient transfer (i.e., low-priority gradients cannot preempt high-priority gradients in the network transfer). Prophet can make the forward propagation start as early as possible so as to greedily reduce the waiting (idle) time of GPU resources during the DDNN training process. Prototype experiments with representative DNN models trained on Amazon EC2 demonstrate that Prophet can improve the DDNN training performance by up to 40% compared with the state-of-the-art priority-based communication scheduling strategies, yet with negligible runtime performance overhead.
Qiang Qi, Ruitao Shang, Li Chen 0019, Fei Xu 0009
ICPP4
2021 Astra: Autonomous Serverless Analytics with Cost-Efficiency and QoS-Awareness
abstract
With the ability to simplify the code deployment with one-click upload and lightweight execution, serverless computing has emerged as a promising paradigm with increasing popularity. However, there remain open challenges when adapting data-intensive analytics applications to the serverless context, in which users of serverless analytics encounter with the difficulty in coordinating computation across different stages and provisioning resources in a large configuration space. This paper presents our design and implementation of Astra, which configures and orchestrates serverless analytics jobs in an autonomous manner, while taking into account flexibly-specified user requirements. Astra relies on the modeling of performance and cost which characterizes the intricate interplay among multi-dimensional factors (e.g., function memory size, degree of parallelism at each stage). We formulate an optimization problem based on user-specific requirements towards performance enhancement or cost reduction, and develop a set of algorithms based on graph theory to obtain optimal job execution. We deploy Astra in the AWS Lambda platform and conduct real-world experiments over three representative benchmarks with different scales. Results demonstrate that Astra can achieve the optimal execution decision for serverless analytics, by improving the performance of 21% to 60% under a given budget constraint, and resulting in a cost reduction of 20% to 80% without violating performance requirement, when compared with three baseline configuration algorithms.
Jananie Jarachanthan, Li Chen 0019, Fei Xu 0009, Bo Li 0001
IPDPS2
2021 EAGLE: Expedited Device Placement with Automatic Grouping for Large Models
abstract
Advanced deep neural networks with large sizes are usually trained on a mixture of devices, including multiple CPUs and GPUs. The model training speed and efficiency are drastically impacted by the placement of operations on devices. To identify the optimal device placement, the state-of-the-art method is based on reinforcement learning with a hierarchical model, which partitions the operations into groups and then assigns each group to specific devices. However, due to the additional dimension of grouping decisions coupled with the placement, the reinforcement learning efficiency is greatly reduced. With modern neural networks growing in size and complexity, the issue of low efficiency and high cost in device placement is further aggravated. In this paper, we propose our design of EAGLE (Expedited Automatic Grouping for Large modEls), which integrates automatic grouping into reinforcement learning-based placement in an optimal way, to achieve the best possible training time performance for very large models. An extra RNN is introduced to transform parameters of the grouper into inputs of the placer, linking the originally separated parts together. Further optimizations have also been made in the network inputs. We have deployed and extensively evaluated EAGLE on InceptionV3, GNMT and BERT benchmarks. Compared with the state-of-the-art, the performance achieved by our design, measured by the per-step time with the resulted placement, is 2.7% and 18.7% better for GNMT and BERT, respectively. For Inception-V3, our design achieves the fastest speed in discovering the optimal placement.
Li Chen 0019, Baochun Li
IPDPS2
2021 Precise Weather Parameter Predictions for Target Regions via Neural Networks
Yihe Zhang 0001, Xu Yuan 0001, Sytske K. Kimball, Eric Rappin, Li Chen 0019, Paul J. Darby III, Tom Johnsten, Lu Peng 0001, Boisy Pitre, David M. Bourrie, Nian-Feng Tzeng
ECML/PKDD (5)5
2021 Rationing bandwidth resources for mitigating network resource contention in distributed DNN training clusters
Qiang Qi, Fei Xu 0009, Li Chen 0019, Zhi Zhou 0006
CCF Trans. High Perform. Comput.3
2021 A Case for Pricing Bandwidth: Sharing Datacenter Networks With Cost Dominant Fairness
abstract
Unlike other resources such as CPU or memory in a virtual machine, inter-virtual-machine (inter-VM) bandwidth has not been explicitly priced in datacenter networks. In this article, we argue that tenants of an IaaS cloud computing platform should be given the flexibility to pay more for explicitly priced datacenter bandwidth beyond traditional virtual machines, in order to achieve better (or more predictable) application performance. We show that a much simpler design principle can be followed to allocate bandwidth fairly, and desirable properties related to fairness can be more easily achieved, compared with state-of-the-art proposals. We call such a design principle cost dominant fairness, which stipulates that bandwidth should be allocated based on the total cost that a tenant incurs for running its applications in the cloud. Guided by the principle of cost dominant fairness, we explore the design space of pricing inter-VM bandwidth, as well as achieving fair bandwidth sharing among multiple tenants. Through our study, we believe that it is best to assign per-VM-pair weights based on individualized prices. We present a distributed bandwidth allocation algorithm that is theoretically supported by a network utility maximization formulation, and practically implemented as a shim layer at each virtual machine. We are also concerned with practical issues of billing, where discounts are needed to ensure that a tenant only pays for the bandwidth share that it is allocated. Finally, we have evaluated our pricing framework and per-VM-pair weighted fair bandwidth allocation in the Mininet emulation testbed and simulations.
Li Chen 0019, Yuan Feng 0005, Baochun Li, Bo Li 0001
IEEE Trans. Parallel Distributed Syst.1
2020 Towards Poisoning the Neural Collaborative Filtering-Based Recommender Systems
Yihe Zhang 0001, Jiadong Lou, Li Chen 0019, Xu Yuan 0001, Jin Li 0002, Tom Johnsten, Nian-Feng Tzeng
ESORICS (1)3
2020 FEEL: A Federated Edge Learning System for Efficient and Privacy-Preserving Mobile Healthcare
abstract
With the prosperity of artificial intelligence, neural networks have been increasingly applied in healthcare for a variety of tasks for medical diagnosis and disease prevention. Mobile wearable devices, widely adopted by hospitals and health organizations, serve as emerging sources of medical data and participate in the training of neural network models for accurate model inference. Since the medical data are privacy-sensitive and non-shareable, federated learning has been proposed to train a model across decentralized data, which involves each mobile device running a training task with its own data in parallel. However, due to the ever-increasing size and complexity of modern neural network models, it becomes inefficient, and may even infeasible, to perform training tasks on wearable devices that are resource-constrained. In this paper, we propose a FEderated Edge Learning system, FEEL, for efficient privacy-preserving mobile healthcare. Specifically, we design an edge-based training task offloading strategy to improve the training efficiency. Further, we build our system on the basis of federated learning to make use of distributed user data to improve the inference performance. In addition, during model training, we provide a differential privacy scheme to strengthen the privacy protection. A prototype system has been implemented to evaluate the training efficiency, inference performance and noise sensitivity, respectively. And the results have demonstrated that our proposal could train models in an efficient and privacy-preserving way.
Yeting Guo, Fang Liu 0002, Zhiping Cai, Li Chen 0019, Nong Xiao 0001
ICPP4
2020 E-LAS: Design and Analysis of Completion-Time Agnostic Scheduling for Distributed Deep Learning Cluster
abstract
With the prosperity of deep learning, enterprises, and large platform providers, such as Microsoft, Amazon, and Google, have built and provided GPU clusters to facilitate distributed deep learning training. As deep learning training workloads are heterogeneous, with a diverse range of characteristics and resource requirements, it becomes increasingly crucial to design an efficient and optimal scheduler for distributed deep learning jobs in the GPU cluster. This paper aims to propose a simple and yet effective scheduler, called E-LAS, with the objective of reducing the averaged training completion time of deep learning jobs. Without relying on the estimation or prior knowledge of the job running time, E-LAS leverages the real-time epoch progress rate, unique for distributed deep learning training jobs, as well as the attained services from temporal and spatial domains, to guide the scheduling decisions. The theoretical analysis for E-LAS is conducted to offer a deeper understanding on the components of scheduling criteria. Furthermore, we present a placement algorithm to achieve better resource utilization without involving much implementation overhead, complementary to the scheduling algorithm. Extensive simulations have been conducted, demonstrating that E-LAS improves the averaged job completion time (JCT) by 10 × over an Apache YARN-based resource manager used in production. Moreover, E-LAS outperforms Tiresias, the state-of-the-art scheduling algorithm customized for deep learning jobs, by almost 1.5 × for the average JCT as well as queuing time.
Abeda Sultana, Li Chen 0019, Fei Xu 0009, Xu Yuan 0001
ICPP2
2019 Stage Delay Scheduling: Speeding up DAG-style Data Analytics Jobs with Resource Interleaving
abstract
To increase the resource utilization of datacenters, big data analytics jobs are commonly running stages in parallel which are organized into and scheduled according to the Directed Acyclic Graph (DAG). Through an in-depth analysis of the latest Alibaba cluster trace and our motivation experiments on Amazon EC2, however, we show that the CPU and network resources are still under-utilized due to the unwise stage scheduling, thereby prolonging the completion time of a DAG-style job (e.g., Spark). While existing works on reducing the job completion time focus on either task scheduling or job scheduling, stage scheduling has received comparably little attention. In this paper, we design and implement DelayStage, a simple yet effective stage delay scheduling strategy to interleave the cluster resources across the parallel stages, so as to increase the cluster resource utilization and speed up the job performance. With the aim of minimizing the makespan of parallel stages, DelayStage judiciously arranges the execution of stages in a pipelined manner to maximize the performance benefits of resource interleaving. Extensive prototype experiments on 30 Amazon EC2 instances and complementary trace-driven simulations show that DelayStage can improve the cluster resource utilization by up to 81.8% and reduce the job completion time by up to 41.3%, in comparison to the stock Spark and the state-of-the-art stage scheduling strategies, yet with acceptable runtime overhead.
Wujie Shao, Fei Xu 0009, Li Chen 0019, Haoyue Zheng, Fangming Liu
ICPP3
2019 Cynthia: Cost-Efficient Cloud Resource Provisioning for Predictable Distributed Deep Neural Network Training
abstract
It becomes an increasingly popular trend for deep neural networks with large-scale datasets to be trained in a distributed manner in the cloud. However, widely known as resource-intensive and time-consuming, distributed deep neural network (DDNN) training suffers from unpredictable performance in the cloud, due to the intricate factors of resource bottleneck, heterogeneity and the imbalance of computation and communication which eventually cause severe resource under-utilization. In this paper, we propose Cynthia, a cost-efficient cloud resource provisioning framework to provide predictable DDNN training performance and reduce the training budget. To explicitly explore the resource bottleneck and heterogeneity, Cynthia predicts the DDNN training time by leveraging a lightweight analytical performance model based on the resource consumption of workers and parameter servers. With an accurate performance prediction, Cynthia is able to optimally provision the cost-efficient cloud instances to jointly guarantee the training performance and minimize the training budget. We implement Cynthia on top of Kubernetes by launching a 56-docker cluster to train four representative DNN models. Extensive prototype experiments on Amazon EC2 demonstrate that Cynthia can provide predictable training performance while reducing the monetary cost for DDNN workloads by up to 50.6%, in comparison to state-of-the-art resource provisioning strategies, yet with acceptable runtime overhead.
Haoyue Zheng, Fei Xu 0009, Li Chen 0019, Zhi Zhou 0006, Fangming Liu
ICPP3
2019 Promenade: Proportionally Fair Multipath Rate Control in Datacenter Networks with Random Network Coding
abstract
In today's datacenter topologies, there exist multiple equal-cost paths between each pair of communicating virtual machines. Yet, splitting flows and routing them along multiple paths may lead to packet reordering, which may affect the performance of TCP. In this paper, we propose Promenade, a new protocol that uses random network coding to mitigate the negative effects of packet reordering, while at the same time achieving weighted proportional fairness in bandwidth allocation across different tenants. To achieve weighted proportional fairness when allocating bandwidth to tenants, the problem of rate control is formulated as a convex optimization problem, and Promenade uses its distributed solution as a theoretical foundation to design its bandwidth allocation protocol. With our real-world implementation of Promenade in the Mininet testbed, we are able to show that Promenade is able to achieve weighted proportional fairness in its rate control when individual flows are split into multiple paths.
Li Chen 0019, Yuan Feng 0005, Baochun Li, Bo Li 0001
IEEE Trans. Parallel Distributed Syst.1
2018 Spotlight: Optimizing Device Placement for Training Deep Neural Networks
abstract
Training deep neural networks (DNNs) requires an increasing amount of computation resources, and it becomes typical to use a mixture of GPU and CPU devices. Due to the heterogeneity of these devices, a recent challenge is how each operation in a neural network can be optimally placed on these devices, so that the training process can take the shortest amount of time possible. The current state-of-the-art solution uses reinforcement learning based on the policy gradient method, and it suffers from suboptimal training times. In this paper, we propose Spotlight, a new reinforcement learning algorithm based on proximal policy optimization, designed specifically for finding an optimal device placement for training DNNs. The design of our new algorithm relies upon a new model of the device placement problem: by modeling it as a Markov decision process with multiple stages, we are able to prove that Spotlight achieves a theoretical guarantee on performance improvements. We have implemented Spotlight in the CIFAR-10 benchmark and deployed it on the Google Cloud platform. Extensive experiments have demonstrated that the training time with placements recommended by Spotlight is 60.9% of that recommended by the policy gradient method.
Yuanxiang Gao, Li Chen 0019, Baochun Li
ICML2
2018 A Hierarchical Synchronous Parallel Model for Wide-Area Graph Analytics
abstract
Graph analytics has emerged as one of the fundamental techniques to support modern Internet applications. As real-world graph data is generated and stored globally, the scale of the graph that needs to be processed keeps growing. It is critical to efficiently process graphs across multiple geographically distributed datacenters, running wide-area graph analytics. Existing graph analytics frameworks are not designed to run across multiple datacenters well, as they implement a Bulk Synchronous Parallel model that requires excessive wide-area data transfers. In this paper, we present a new Hierarchical Synchronous Parallel model designed and implemented for synchronization across datacenters with a much improved efficiency in inter-datacenter communication. Our new model requires no modifications to graph analytics applications, yet guarantees their convergence and correctness. Our prototype implementation on Apache Spark can achieve up to 32% lower WAN bandwidth usage, 49% faster convergence, and 30% less total cost for benchmark graph algorithms, with input data stored across five geographically distributed datacenters.
Shuhao Liu 0001, Li Chen 0019, Baochun Li, Aiden Carnegie
INFOCOM2
2018 Post: Device Placement with Cross-Entropy Minimization and Proximal Policy Optimization
abstract
Training deep neural networks requires an exorbitant amount of computation resources, including a heterogeneous mix of GPU and CPU devices. It is critical to place operations in a neural network on these devices in an optimal way, so that the training process can complete within the shortest amount of time. The state-of-the-art uses reinforcement learning to learn placement skills by repeatedly performing Monte-Carlo experiments. However, due to its equal treatment of placement samples, we argue that there remains ample room for significant improvements. In this paper, we propose a new joint learning algorithm, called Post, that integrates cross-entropy minimization and proximal policy optimization to achieve theoretically guaranteed optimal efficiency. In order to incorporate the cross-entropy method as a sampling technique, we propose to represent placements using discrete probability distributions, which allows us to estimate an optimal probability mass by maximal likelihood estimation, a powerful tool with the best possible efficiency. We have implemented Post in the Google Cloud platform, and our extensive experiments with several popular neural network training benchmarks have demonstrated clear evidence of superior performance: with the same amount of learning time, it leads to placements that have training times up to 63.7% shorter over the state-of-the-art.
Yuanxiang Gao, Li Chen 0019, Baochun Li
NeurIPS2
2018 Siphon: Expediting Inter-Datacenter Coflows in Wide-Area Data Analytics
Shuhao Liu 0001, Li Chen 0019, Baochun Li
USENIX ATC2
2018 Efficient Performance-Centric Bandwidth Allocation with Fairness Tradeoff
abstract
Fair bandwidth allocation in datacenter networks has received a substantial amount of research attention, as multiple tenants are hosted by virtual machines in a public cloud. In the context of private datacenters, link bandwidth is shared among applications running data parallel frameworks, such as MapReduce, instead. In this paper, we introduce the rigorous definition of performance-centric fairness, with the guiding principle that the performance that data parallel applications will enjoy should be proportional to their weights. We first investigate the problem of maximizing application performance while maintaining strict performance-centric fairness. We then present an inherent tradeoff between fairness and efficiency, which is interpreted from the perspectives of bandwidth utilization and social welfare, respectively. From the first perspective, we propose an algorithm to improve bandwidth utilization by introducing an extended version of fairness. From the second perspective, we formulate an optimization problem of bandwidth allocation that maximizes the social welfare across all the applications, allowing a tunable degree of relaxation on performance-centric fairness. A distributed algorithm is then presented to solve the problem, based on dual based decomposition. With extensive simulations, we demonstrate the effectiveness of our algorithms in improving efficiency and application performance (by up to 1.4X), with flexible degree of relaxation on the performance-centric fairness.
Li Chen 0019, Yuan Feng 0005, Baochun Li, Bo Li 0001
IEEE Trans. Parallel Distributed Syst.1
2017 Siphon: a high-performance substrate for inter-datacenter transfers in wide-area data analytics
abstract
Large volumes of raw data are naturally generated in multiple geographical regions. Modern data parallel frameworks may suffer from a degraded performance, due to the much lower inter-datacenter WAN bandwidth. We present Siphon, a new high-performance substrate that improves the ability of modern data parallel frameworks to optimize WAN transfers. Siphon aggregates inter-datacenter flows and routes them with the awareness of runtime network capacities. Following the principles of software-defined networking, Siphon is designed to accommodate a variety of flow scheduling algorithms that can improve the job level performance.
Shuhao Liu 0001, Li Chen 0019, Baochun Li
SoCC2
2017 Scheduling jobs across geo-distributed datacenters with max-min fairness
abstract
It has become routine for large volumes of data to be generated, stored, and processed across geographically distributed datacenters. To run a single data analytic job on such geo-distributed data, recent research proposed to distribute its tasks across datacenters, considering both data locality and network bandwidth across datacenters. Yet, it remains an open problem in the more general case, where multiple analytic jobs need to fairly share the resources at these geo-distributed data-centers. In this paper, we focus on the problem of assigning tasks belonging to multiple jobs across datacenters, with the specific objective of achieving max-min fairness across jobs sharing these datacenters, in terms of their job completion times. We formulate this problem as a lexicographical minimization problem, which is challenging to solve in practice due to its inherent multi-objective and discrete nature. To address these challenges, we iteratively solve its single-objective subproblems, which can be transformed to equivalent linear programming (LP) problems to be efficiently solved, thanks to their favorable properties. As a highlight of this paper, we have designed and implemented our proposed solution as a fair job scheduler based on Apache Spark, a modern data processing framework. With extensive evaluations of our real-world implementation on Amazon EC2, we have shown convincing evidence that max-min fairness has been achieved using our new job scheduler.
Li Chen 0019, Shuhao Liu 0001, Baochun Li, Bo Li 0001
INFOCOM1
2016 Surviving Failures with Performance-Centric Bandwidth Allocation in Private Datacenters
abstract
In the context of private datacenters that are operated by Web service providers such as Google, multiple applications using data parallel frameworks, such as MapReduce, coexist and share a limited supply of link bandwidth capacities. It has been shown that failures are the norm, rather than the exception, in datacenters, and will negatively affect the performance of data parallel applications as failed tasks need to be relaunched and placed on newly selected servers. In this paper, we argue that even with the presence of failures, link bandwidth should be allocated to competing applications with performance-centric fairness, in that the performance that applications enjoy should be proportional to their weights. We formulate and solve the open challenge of jointly optimizing placement decisions for relaunched tasks and bandwidth allocation, so that the adverse effects of failures on application performance are minimized. With our proposed algorithm implemented in the Mininet emulation testbed, our experiments show the effectiveness of our solutions towards minimizing the negative effects of failures, while still achieving performance-centric fairness.
Li Chen 0019, Baochun Li, Bo Li 0001
IC2E1
2016 Barrier-Aware Max-Min Fair Bandwidth Sharing and Path Selection in Datacenter Networks
abstract
In production datacenters operated by Web service providers such as Google, multiple data parallel applications, such as MapReduce, are employed to facilitate data processing at a large scale, with a strong demand for intra-datacenter bandwidth in their communication stages. A noteworthy phenomenon in these applications is the presence of barriers, which implies that a job will not finish until the last task completes. Existing flow-level sharing in datacenter networks is not designed and optimized to meet such application-level needs. In this paper, we promote the awareness of application barriers in the design of both bandwidth allocation and path selection strategies. In particular, we propose the notion of application-level fairness when bandwidth is allocated, with favorable properties of performance-centric max-min fairness and Pareto efficiency. Further, we show that both application-level performance and resource utilization can be further improved by considering path selection as well. With our implementation in the Mininet emulation testbed and large-scale simulations, we demonstrate that our new barrier-aware strategy for fair bandwidth sharing and path selection significantly outperforms barrier-agnostic strategy when application performance is concerned.
Li Chen 0019, Baochun Li, Bo Li 0001
IC2E1
2016 Optimizing coflow completion times with utility max-min fairness
abstract
In data parallel frameworks such as MapReduce and Spark, a coflow represents a set of network flows used to transfer intermediate data between successive computation stages for a job. The completion time of a job is then determined by the collective behavior of such a coflow, rather than any individual flow within, and influenced by the amount of network bandwidth allocated to it. Different jobs in a shared cluster have different degrees of sensitivity to their completion times, modeled by their respective utility functions. In this paper, we focus on the design and implementation of a new utility optimal scheduler across competing coflows, in order to provide differential treatment to coflows with different degrees of sensitivity, yet still satisfying max-min fairness across these coflows. Though this objective can be formulated as a lexicographical maximization problem, it is challenging to solve in practice due to its inherent multi-objective and discrete nature. To address this challenge, we first divide the problem into iterative steps of single-objective subproblems; and in each of these steps, we then perform a series of transformations to obtain an equivalent linear programming (LP) problem, which can be efficiently solved in practice. To demonstrate that our solutions are practically feasible, we have implemented it as a real-world coflow scheduler based on the Varys open-source framework to evaluate its effectiveness.
Li Chen 0019, Baochun Li, Bo Li 0001
INFOCOM1
2014 Towards performance-centric fairness in datacenter networks
abstract
Fair bandwidth allocation in datacenter networks has been a focus of research recently, yet this has not received adequate attention in the context of private cloud, where link bandwidth is often shared among applications running data parallel frameworks, such as MapReduce. In this paper, we introduce a rigorous definition of performance-centric fairness, with the guiding principle that the performance of data parallel applications should be proportional to their weights. We investigate the problem of maximizing application performance while maintaining strict performance-centric fairness and present the inherent tradeoff between resource utilization and fairness. We then formulate the link bandwidth allocation problem with the objective of maximizing social welfare across all applications, so that resource utilization can be manipulated and improved by allowing a tunable degree of relaxation on performance-centric fairness. Based on dual based decomposition, we present a distributed algorithm to solve this problem, and evaluate its performance with extensive simulations.
Li Chen 0019, Yuan Feng 0005, Baochun Li, Bo Li 0001
INFOCOM1
2014 Allocating Bandwidth in Datacenter Networks: A Survey
Li Chen 0019, Baochun Li, Bo Li 0001
J. Comput. Sci. Technol.1
2012 Lifetime or energy: Consolidating servers with reliability control in virtualized cloud datacenters
abstract
Server consolidation using virtualization technologies allow cloud-scale datacenters to improve resource utilization and energy efficiency. However, most existing consolidation strategies solely focused on balancing the tradeoff between service-level-agreements (SLAs) desired by cloud applications and energy costs consumed by hosting servers. With the presence of fluctuating workloads in datacenters, the lifetime and reliability of servers under dynamic power-aware consolidation could be adversely impacted by repeated on-off thermal cycles, wear-and-tear and temperature rise. In this paper, we propose a Reliability-Aware server Consolidation stratEgy, named RACE, to address when and how to perform energy-efficient server consolidation in a reliability-friendly and profitable way. The focus is on the characterization and analysis of this problem as a multi-objective optimization, by developing an utility model that unifies multiple constraints on performance SLAs, reliability factors, and energy costs in a holistic manner. An improved grouping genetic algorithm is proposed to search the global optimal solution, which takes advantage of a collection of reliability-aware resource buffering, and virtual machines-to-servers re-mapping heuristics for generating good initial solutions and improving the convergence rate. Extensive simulations are conducted to validate the effectiveness, scalability and overhead of RACE in improving the overall utility of datacenters while avoiding unprofitable consolidation in the long term - compared with pMapper and PADD strategies for server consolidation.
Fangming Liu, Hai Jin 0001, Xiaofei Liao, Haikun Liu, Li Chen 0019
CloudCom6