Wei Wang 0030

dblp:35/7092-30 · DBLP profile ↗
← Back
112ranked-venue papers
19as first author
65since 2021 · last 2026
—ORCID · conflict

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

Systems, architecture and hardware · 61 · 7 first-author · 36 since 2021Computer networks · 41 · 12 first-author · 19 since 2021Security and privacy · 4 · 4 since 2021Software engineering, systems software and programming languages · 4 · 1 first-author · 3 since 2021Artificial intelligence and machine learning · 2 · 2 since 2021Applied, interdisciplinary, general and emerging computing · 2 · 2 since 2021Databases, data management, data science and information retrieval · 1 · 1 since 2021Human-computer interaction and ubiquitous computing · 1 · 1 since 2021
YearPublicationVenuePosition
2026 ZipServ: Fast and Memory-Efficient LLM Inference with Hardware-Aware Lossless Compression
abstract
Lossless model compression holds tremendous promise for alleviating the memory and bandwidth bottlenecks in bitexact Large Language Model (LLM) serving. However, existing approaches often result in substantial inference slowdowns due to fundamental design mismatches with GPU architectures: at the kernel level, variable-length bitstreams produced by traditional entropy codecs break SIMT parallelism; at the system level, decoupled pipelines lead to redundant memory traffic. We present ZipServ, a lossless compression framework co-designed for efficient LLM inference. ZipServ introduces Tensor-Core-Aware Triple Bitmap Encoding (TCA-TBE), a novel fixed-length format that enables constant-time, parallel decoding, together with a fused decompression-GEMM (ZipGEMM) kernel that decompresses weights on-the-fly directly into Tensor Core registers. This "load-compressed, compute-decompressed" design eliminates intermediate buffers and maximizes compute intensity. Experiments show that ZipServ reduces the model size by up to 30%, achieves up to 2.21× kernel-level speedup over NVIDIA’s cuBLAS, and expedites end-to-end inference by an average of 1.22× over vLLM. ZipServ is the first lossless compression system that provides both storage savings and substantial acceleration for LLM inference on GPUs.
Ruibo Fan, Xiangrui Yu, Xinglin Pan, Weile Luo, Qiang Wang 0022, Wei Wang 0030, Xiaowen Chu 0001
ASPLOS (2)7
2026 FlashPS: Efficient Generative Image Editing with Mask-aware Caching and Scheduling
abstract
Generative image editing using diffusion models has become a prevalent application in today's AI cloud services. In production environments, image editing typically involves a mask that specifies the regions of an image template to be edited. The use of mask provides direct control over the editing process and introduces sparsity in the model inference. In this paper, we present FlashPS, a system that efficiently serves image editing requests. The key insight behind FlashPS is that image editing only modifies the masked regions of image templates, while preserving the original content in the unmasked areas. Driven by this insight, FlashPS judiciously skips redundant computations associated with the unmask areas by reusing cached intermediate activations from previous inferences. To mitigate the high cache loading overhead, FlashPS employs a bubble-free pipeline scheme that overlaps computation with cache loading. Additionally, to reduce queuing latency in online serving while improving the GPU utilization, FlashPS proposes a novel continuous batching strategy for diffusion model serving, allowing newly arrived requests to join the running batch in just one step of denoising computation, without waiting for the entire batch to complete. As heterogenous masks induce imbalanced load, FlashPS also develops a load balancing strategy that takes into account the loads of both computation and cache loading. Collectively, FlashPS outperforms state-of-the-art diffusion serving systems for image editing, achieving up to 3× higher throughput and reducing average request latency by up to 14.7× while ensuring image quality.
Xiaoxiao Jiang, Suyi Li 0002, Lingyun Yang, Tianyu Feng, Zhipeng Di, Weiyi Lu, Guoxuan Zhu, Xiu Lin, Yinghao Yu, Tao Lan, Lin Qu, Liping Zhang 0013, Wei Wang 0030
EuroSys15
2026 Efficient Data Passing for Serverless Inference Workflows: A GPU-Centric Approach
abstract
Serverless computing offers a compelling paradigm for deploying machine learning inference workflows composed of heterogeneous CPU and GPU functions. However, existing data-passing solutions in serverless systems primarily rely on host memory for data exchange (host-centric), leading to substantial data movement and salient I/O overhead. Moreover, modern GPU communication libraries (e.g., NCCL, NVSHMEM, UCX) are ill-suited to serverless environments, suffering from redundant data copies, underutilized transfer bandwidth, and inefficient temporary GPU storage.
Hao Wu 0032, Yaochen Liu, Minchen Yu, Qizhen Weng 0001, Junxiao Deng, Hao Fan 0006, Song Wu 0001, Wei Wang 0030, Hai Jin 0001
EuroSys9
2026 ELORA: Efficient LoRA and KV Cache Management for Multi-LoRA LLM Serving
abstract
Multiple Low-Rank Adapters (Multi-LoRA) are gaining popularity for task-specific Large Language Model (LLM) applications. For Multi-LoRA serving, caching hot LoRAs and KV caches in the GPU memory can improve inference performance. However, existing Multi-LoRA inference systems fail to optimize serving performance like Time-To-First-Token (TTFT), neglecting usage dependencies when caching LoRAs and KV caches. We therefore propose ELORA, a Multi-LoRA caching system to optimize the serving performance. ELORA comprises a dependency-aware cache manager and a performancedriven cache swapper. The cache manager maintains the usage dependencies between LoRAs and KV caches during inference with a unified caching pool. The cache swapper determines the swap-in or swap-out of LoRAs and KV caches based on a unified cost model, when the GPU memory is idle or busy, respectively. Experimental results show that ELORA reduces the TTFT by$\mathbf{4 5. 7 \%}$on average, compared to state-of-the-art works.
Jiuchen Shi, Quan Chen 0002, Yizhou Shan, Kaihua Fu, Wei Wang 0030, Minyi Guo
HPCA7
2026 ZkChainDB: A Verifiable and Privacy-Preserving SQL Query Engine for Off-Chain Database Using zk-SNARKs and Blockchain
Wei Wang 0030, Jianan Hong
ICBC2
2026 Less is More: Persistent Low-Frequency Backdoor Injection in Federated Learning
abstract
Federated learning (FL) enables multiple clients to collaboratively train a machine learning model without sharing their local data. However, the distributed nature of FL makes it vulnerable to backdoor attacks from malicious clients. Most existing attack methods often assume that attackers can inject backdoors in every training round - a scenario that is both unrealistic and inefficient in real-world FL deployment. In this paper, we investigate why backdoor attacks become less effective under low-frequency injection and propose a novel attack paradigm for FL, called REinforced Memorization-based INterval backDoor attack (REMIND). REMIND optimizes the backdoor trigger via task alignment and feature alignment. Task alignment aligns backdoor and main task objectives to resist benign update suppression during non-attack rounds, while feature alignment guides poisoned samples to match the activation trajectory of target-class samples. This dual alignment enhances the backdoor's persistence and narrows the divergence between malicious and benign updates. With strong attack success rates established, we further analyze the advantages of low-frequency backdoor attacks, particularly their ability to improve robustness against defense mechanisms. Extensive evaluations on four benchmark datasets show that REMIND consistently outperforms eight state-of-the-art attack baselines under nine defense strategies.
Pei Ye, Yuqing Li 0001, Kun He 0008, Ruiying Du, Wei Wang 0030
INFOCOM6
2026 RollPacker: Taming Long-Tail Rollouts for RL Post-Training with Tail Batching
Yuheng Zhao, Dakai An, Tianyuan Wu, Lunxi Cao, Shaopan Xiong, Ju Huang, Weixun Wang, Siran Yang, Wenbo Su, Jiamang Wang, Lin Qu, Bo Zheng 0007, Wei Wang 0030
NSDI14
2026 Attack of the Bubbles: Straggler-Resilient Pipeline Parallelism for Large Model Training
Tianyuan Wu, Lunxi Cao, Hanfeng Lu, Xiaoxiao Jiang, Yinghao Yu, Siran Yang, Jiamang Wang, Lin Qu, Liping Zhang 0013, Wei Wang 0030
NSDI11
2026 Enabling Low-Latency, GPU-Efficient Serverless Inference with Model Swapping
abstract
Serverless computing offers a compelling cloud model for online inference services. However, existing serverless platforms lack efficient support for GPUs, hindering their ability to deliver high-performance inference. In this article, we present Torpor , a serverless platform for GPU-efficient, low-latency inference. To enable efficient sharing of a node’s GPUs among numerous inference functions, Torpor maintains models in main memory and dynamically swaps them onto GPUs upon request arrivals (i.e., late binding with model swapping). Torpor uses various techniques, including asynchronous API redirection, GPU runtime sharing, pipelined model execution, and efficient GPU memory management, to minimize latency overhead caused by model swapping. Additionally, we design an interference-aware request scheduling algorithm that utilizes high-speed GPU interconnects to meet latency service-level objectives (SLOs) for individual inference functions. We have implemented Torpor and evaluated its performance in a production environment. Utilizing late binding and model swapping, Torpor can concurrently serve hundreds of inference functions on a worker node with 4 GPUs, while achieving latency performance comparable to native execution, where each model is cached exclusively on a GPU. Pilot deployment in a leading commercial serverless cloud shows that Torpor reduces the GPU provisioning cost by 70% and 65% for users and the platform, respectively.
Minchen Yu, Bohui Wu, Haoxuan Yu, Wei Wang 0030, Ruichuan Chen, Dapeng Nie
ACM Trans. Archit. Code Optim.7
2026 QoS Awareness and Improved Throughput of Point Cloud Services With Dynamic Workloads
abstract
Deep learning on 3D point clouds plays a vital role in a wide range of applications such as AR/VR visualization, 3D cloth virtual try-on, and game rendering. As some applications require low latency, the point cloud services are also deployed on datacenter with powerful GPUs. While the queries of point cloud services show various workload change patterns due to different degrees of sparsity, current batching-based serving schemes result in either long latency or low throughput. We propose a scheme called Volans to address the above challenges and effectively support point cloud services. Volans comprises a workload predictor, a topology deployer, and a progress-aware scheduler. The predictor grids the input query and estimates the workload changes. Afterward, the deployer splits the model into several stages and determines the batch size for each stage based on the workload changes. The scheduler reduces the QoS violation when queries run slower due to unpredicted workload spikes. Experiments show that Volans enhances the peak supported throughput by up to 31.1% while maintaining the required 99%-ile latencies compared to state-of-the-art techniques.
Kaihua Fu, Jiuchen Shi, Yao Chen 0008, Quan Chen 0002, Weng-Fai Wong, Wei Wang 0030, Bingsheng He, Minyi Guo
IEEE Trans. Computers6
2026 Mitigating Server-Side Communication Bottlenecks in Distributed Learning With Round-Robin Participant Coordination
abstract
Deep neural networks are increasingly trained in a distributed manner—in either clusters or with federated devices, where the participants jointly refine the global model with their gradients calculated locally. More often than not, those gradients are collected to a central server in a synchronous manner to avoid the negative impact of stale updates. However, when all the participants communicate their gradients to the server in such a uniform pace, the network on the server side—under intense contention—often becomes a performance bottleneck. To address this problem, for the cluster environment we propose theRound-Robin Synchronous Parallel(R2SP) scheme, which coordinates the participants to make updates in anevenly-gapped,round-robinmanner. This way, we can minimize the network contention with a minimum cost of the update quality; we also propose to incorporate adaptive batch sizing in R2SP to address the hardware heterogeneity among workers. Moreover, for the federated learning (FL) scenarios, we note that it is necessary yet challenging to apply the insight of R2SP to mitigate the network bottleneck in the FL server, given that there are a huge number of participants with unstable resources and inconsistent data distributions. To tackle those challenges, we further propose FL-R2SP, which extents the coordination units from individual participants to participantgroups—with the resource instability and data heterogeneity tackled within each group. We have implemented R2SP and FL-R2SP respectively with TensorFlow and PyTorch, and extensive EC2 experiments show that R2SP and FL-R2SP can respectively speed up model convergence for clustered and federated scenarios by over 20%.
Jiayi Zhang 0006, Chen Chen 0067, Zuo Gan, Wei Wang 0030, Bo Li 0001, Minyi Guo
IEEE Trans. Netw.4
2025 ZipBatch: Multi-Tenant GPU Batching with Dual-Resource Regulation
abstract
GPU multiplexing is a widely-adopted strategy in GPU clusters for improving overall throughput and lowering the total cost of ownership. To mitigate inter-task interference in compute power and memory bandwidth on multiplexed GPUs, existing techniques divide a GPU into instances with limited predefined rigid configurations. Low utilization arises from the mismatch between heterogeneous burstiness and immutable resource configurations: 1) bursty inference traffic forces the scheduler to launch underfilled batches that cannot saturate the instance; 2) bursty kernel resource utilization leads to bubbles in compute power and memory bandwidth.
Haoxuan Yu, Sheng Yao 0006, Wei Wang 0030
SoCC3
2025 SpInfer: Leveraging Low-Level Sparsity for Efficient Large Language Model Inference on GPUs
abstract
Large Language Models (LLMs) have demonstrated remarkable capabilities, but their immense scale poses significant challenges in terms of both memory and computational costs. While unstructured pruning offers promising solutions by introducing sparsity to reduce resource requirements, realizing its benefits in LLM inference remains elusive. This is primarily due to the storage overhead of indexing non-zero elements and the inefficiency of sparse matrix multiplication (SpMM) kernels at low sparsity levels (around 50%). In this paper, we present SpInfer, a high-performance framework tailored for sparsified LLM inference on GPUs. SpInfer introduces Tensor-Core-Aware Bitmap Encoding (TCA-BME), a novel sparse format that minimizes indexing overhead by leveraging efficient bitmap-based indexing, optimized for GPU Tensor Core architectures. Furthermore, SpInfer integrates an optimized SpMM kernel with Shared Memory Bitmap Decoding (SMBD) and asynchronous pipeline design to enhance computational efficiency. Experimental results show that SpInfer significantly outperforms state-of-the-art SpMM implementations (up to 2.14× and 2.27× over Flash-LLM and SparTA, respectively) across a range of sparsity levels (30% to 70%), with substantial improvements in both memory efficiency and end-to-end inference speed (up to 1.58×). SpInfer outperforms highly optimized cuBLAS at sparsity levels as low as 30%, marking the first effective translation of unstructured pruning's theoretical advantages into practical performance gains for LLM inference.
Ruibo Fan, Xiangrui Yu, Peijie Dong, Gu Gong, Qiang Wang 0022, Wei Wang 0030, Xiaowen Chu 0001
EuroSys7
2025 Optimizing Distributed Deployment of Mixture-of-Experts Model Inference in Serverless Computing
Mengfan Liu, Wei Wang 0030
INFOCOM2
2025 GPU-Disaggregated Serving for Deep Learning Recommendation Models at Scale
Lingyun Yang, Yongchen Wang, Yinghao Yu, Qizhen Weng 0001, Jianbo Dong, Chi Zhang 0005, Yanyi Zi, Zechao Zhang, Menglei Zheng, Lanlan Xi, Binzhang Fu, Tao Lan, Liping Zhang 0013, Lin Qu, Wei Wang 0030
NSDI22
2025 Alibaba Stellar: A New Generation RDMA Network for Cloud AI
abstract
The rapid adoption of Large Language Models (LLMs) in cloud environments has intensified the demand for high-performance AI training and inference, where Remote Direct Memory Access (RDMA) plays a critical role. However, existing RDMA virtualization solutions, such as Single-Root Input/Output Virtualization (SR-IOV), face significant limitations in scalability, performance, and stability. These issues include lengthy container initialization times, hardware resource constraints, and inefficient traffic steering. To address these challenges, we propose Stellar, a new generation RDMA network for cloud AI. Stellar introduces three key innovations: Para-Virtualized Direct Memory Access (PVDMA) for on-demand memory pinning, extended Memory Translation Table (eMTT) for optimized GPU Direct RDMA (GDR) performance, and RDMA Packet Spray for efficient multi-path utilization. Deployed in our large-scale AI clusters, Stellar spins up virtual devices in seconds, reduces container initialization time by 15 times, and improves LLM training speed by up to 14%. Our evaluations demonstrate that Stellar significantly outperforms existing solutions, offering a scalable, stable, and high-performance RDMA network for cloud AI.
Menglei Zheng, Binbin Liao, Suwei Xu, Yongjia Mo, Qinghua Peng, Jilie Luo, Qingxu Li, Zishu Wang, Jianbo Dong, Kunling He, Sheng Cheng 0002, Jiamin Cao, Hairong Jiao, Lingjun Zhu, Yiquan Chen, Wei Wang 0030, Shuhong Zhu, Xingru Li, Qiang Wang 0022, Wei Lin 0016, Ennan Zhai, Jiesheng Wu, Qiang Liu 0036, Binzhang Fu, Dennis Cai
SIGCOMM28
2025 Toppings: CPU-Assisted, Rank-Aware Adapter Serving for LLM Inference
Suyi Li 0002, Hanfeng Lu, Tianyuan Wu, Minchen Yu, Qizhen Weng 0001, Xusheng Chen, Yizhou Shan, Binhang Yuan, Wei Wang 0030
USENIX ATC9
2025 Katz: Efficient Workflow Serving for Diffusion Models with Many Adapters
Suyi Li 0002, Lingyun Yang, Xiaoxiao Jiang, Hanfeng Lu, Dakai An, Zhipeng Di, Weiyi Lu, Yinghao Yu, Tao Lan, Lin Qu, Liping Zhang 0013, Wei Wang 0030
USENIX ATC15
2025 GREYHOUND: Hunting Fail-Slows in Hybrid-Parallel Training at Scale
Tianyuan Wu, Wei Wang 0030, Yinghao Yu, Siran Yang, Qinkai Duan, Jiamang Wang, Lin Qu, Liping Zhang 0013
USENIX ATC2
2025 Torpor: GPU-Enabled Serverless Computing for Low-Latency, Resource-Efficient Inference
Minchen Yu, Haoxuan Yu, Zhuohao Li, Wei Wang 0030, Ruichuan Chen, Dapeng Nie
USENIX ATC7
2025 SP-Chain: Boosting Intrashard and Cross-Shard Security and Performance in Blockchain Sharding
abstract
A promising way to overcome the scalability limitations of the current blockchain is to use sharding, which is to split the transaction processing among multiple, smaller groups of nodes. A well-performing blockchain sharding system requires both high performance and high security in both intra-and cross-shard perspectives. However, existing protocols either have issues in protecting security or trade off great performance for security. In this paper, we propose SP-Chain, a blockchain sharding system with enhanced Security and Performance for both intra-and cross-shard perspectives. For the intra-shard aspect, we design a pipelined two-phase concurrent voting scheme to provide high system throughput and low transaction confirmation latency. Moreover, we propose an efficient unbiased leader rotation scheme to ensure high performance under malicious behavior. For the cross-shard aspect, a proof-assisted efficient cross-shard transaction processing mechanism is proposed to guard cross-shard transactions with low overhead. We implement SP-Chain based on Harmony, and evaluate its performance via large-scale deployment. Extensive evaluations suggest that SP-Chain can process more than 10,000 tx/sec under malicious behaviors with a confirmation latency of 7.6s in a network of 4,000 nodes.
You Lin, Wei Wang 0030, Jin Zhang 0001
IEEE Internet Things J.3
2025 Guest Editorial Special Issue on Federated Learning for Big Data Applications
Xiaowen Chu 0001, Wei Wang 0030, Cong Wang 0001, Yang Liu 0165, Rongfei Zeng, Christopher G. Brinton
IEEE Trans. Big Data2
2025 FedPHE: A Secure and Efficient Federated Learning via Packed Homomorphic Encryption
abstract
Cross-silo federated learning (FL) enables multiple institutions (clients) to collaboratively build a global model without sharing private data. To prevent privacy leakage during aggregation, homomorphic encryption (HE) is widely used to encrypt model updates, yet incurs high computation and communication overheads. To reduce these overheads,packedHE (PHE) has been proposed to encrypt multiple plaintexts into a single ciphertext. However, the original design of PHE assumes all clients share a single private key, making the system vulnerable to security threats of ciphertexts being intercepted and decrypted byhonest-but-curious clients. Also, it does not consider theheterogeneityamong different clients, resulting in undermined training efficiency with slow convergence and stragglers. To address these challenges, we propose FedPHE, a secure and efficient FL framework with PHE by jointly exploiting contribution-aware secure aggregation and straggler-resistant client selection. Using CKKS with sparsification and blinding, FedPHE achieves efficient secure aggregation that allows clients to only provideobscuredencrypted updates while the server can perform aggregation by accounting forcontributionsof local updates. To mitigate the straggler effect, we devise aperturbed sketch-based selection to cherry-pick representative clients withheterogeneous models and computing capabilitiesin a communication-efficient and privacy-preserving manner. We show, through rigorous security analysis and extensive experiments, that FedPHE can efficiently safeguard clients' privacy, achieve$2.45-6.56\times$training speedup, cut the communication overhead by$1.32-24.85\times$, and reduce straggler effects by$1.89-2.78\times$.
Yuqing Li 0001, Nan Yan 0001, Jing Chen 0003, Xiong Wang 0006, Jianan Hong, Kun He 0008, Wei Wang 0030, Bo Li 0001
IEEE Trans. Dependable Secur. Comput.7
2025 Feature Reconstruction Attacks and Countermeasures of DNN Training in Vertical Federated Learning
abstract
Federated learning (FL) has increasingly been deployed, in its vertical form, among organizations to facilitate secure collaborative training. In vertical FL (VFL), participants hold disjoint features of the same set of sample instances. The one withlabels- theactive party, initiates training and interacts with other participants - thepassive parties. It remains largely unknownwhetherandhowan active party can extract private feature data owned by passive parties, especially when training deep neural network (DNN) models. This work examines the feature security problem of DNN training in VFL. We consider a DNN model partitioned between active and passive parties, where the passive party holds a subset of the input layer with some features of binary values. Though proved to be NP-hard. we demonstrate that, unless the feature dimension is exceedingly large, it remains feasible, both theoretically and practically, to launch a reconstruction attack with an efficient search-based algorithm that prevails over current feature protection. We propose a novel feature protection scheme by perturbing intermediate results and fabricated input features, which effectively misleads reconstruction attacks towards pre-specified random values. The evaluation shows it sustains feature reconstruction attack in various VFL applications with negligible impact on model performance.
Peng Ye 0005, Zhifeng Jiang 0001, Wei Wang 0030, Bo Li 0001, Baochun Li
IEEE Trans. Dependable Secur. Comput.3
2025 Pheromone: Restructuring Serverless Computing With Data-Centric Function Orchestration
abstract
Serverless applications are typically composed of function workflows in which multiple short-lived functions are triggered to exchange data in response to events or state changes. Current serverless platforms coordinate and trigger functions by following high-level invocation dependencies but are oblivious to the underlying data exchanges between functions. This design is neither efficient nor easy to use in orchestrating complex workflows – developers often have to manage complex function interactions by themselves, with customized implementation and unsatisfactory performance. Therefore, we argue that function orchestration should follow a data-centric approach. In our design, the platform provides a data bucket abstraction to hold the intermediate data generated by functions. Developers can use a rich set of data trigger primitives to control when and how the output of each function should be passed to the next functions in a workflow. By making data consumption explicit and allowing it to trigger functions and drive the workflow, complex function interactions can be easily and efficiently supported. We presentPheromone– a scalable, low-latency serverless platform following this data-centric design. Compared to well-established commercial and open-source platforms,Pheromonecuts the latencies of function interactions and data exchanges by orders of magnitude, scales to large workflows, and enables easy implementation of complex applications.
Minchen Yu, Tingjia Cao, Wei Wang 0030, Ruichuan Chen
IEEE Trans. Netw.3
2024 DTC-SpMM: Bridging the Gap in Accelerating General Sparse Matrix Multiplication with Tensor Cores
abstract
Sparse Matrix-Matrix Multiplication (SpMM) is a building-block operation in scientific computing and machine learning applications. Recent advancements in hardware, notably Tensor Cores (TCs), have created promising opportunities for accelerating SpMM. However, harnessing these hardware accelerators to speed up general SpMM necessitates considerable effort. In this paper, we undertake a comprehensive analysis of the state-of-the-art techniques for accelerating TC-based SpMM and identify crucial performance gaps. Drawing upon these insights, we propose DTC-SpMM, a novel approach with systematic optimizations tailored for accelerating general SpMM on TCs. DTC-SpMM encapsulates diverse aspects, including efficient compression formats, reordering methods, and runtime pipeline optimizations. Our extensive experiments on modern GPUs with a diverse range of benchmark matrices demonstrate remarkable performance improvements in SpMM acceleration by TCs in conjunction with our proposed optimizations. The case study also shows that DTC-SpMM speeds up end-to-end GNN training by up to 1.91× against popular GNN frameworks.
Ruibo Fan, Wei Wang 0030, Xiaowen Chu 0001
ASPLOS (3)2
2024 SAFE: Intelligent Online Scheduling for Collaborative DNN Inference in Vehicular Network
abstract
Recent years have witnessed a widespread use of deep neural networks (DNNs) in providing various intelligent services, and vehicular networks are no exception. Given the limited computing capabilities of vehicles, collaborative vehicle-edge DNN inference has emerged as a viable alternative. This approach employs DNN partitioning, where a part of DNN is computed on vehicles, and the other part on the edge, e.g., roadside unit (RSU), aiming to enhance the inference accuracy and reduce the inference latency. In this setting, deriving an optimal DNN partitioning scheme becomes critical, yet challenging given the constant movement of vehicles and the highly dynamic wireless connections. Furthermore, vehicles may move out of the signal coverage of an RSU, making it difficult to receive the inference results. To this end, we propose a two-stage intelligent scheduling framework named Soft Actor-critic for discrete actions (SAC-D) based collaborative DNN inference FramEwork (SAFE). SAFE engages multiple RSUs to assist vehicles in completing inference tasks sequentially and ensuring reliable data transmission. It can learn the dynamic vehicular network and make scheduling decisions to minimize the overall latency of vehicle inference tasks. Extensive experimental results show that SAFE can reduce up to 80% of the overall latency with a lower failure rate, compared to four baselines.
Ruiting Zhou, Ziyi Han, Zhi Zhou 0006, Wei Wang 0030
CSCWD6
2024 Dordis: Efficient Federated Learning with Dropout-Resilient Differential Privacy
abstract
Federated learning (FL) is increasingly deployed among multiple clients to train a shared model over decentralized data. To address privacy concerns, FL systems need to safeguard the clients' data from disclosure during training and control data leakage through trained models when exposed to untrusted domains. Distributed differential privacy (DP) offers an appealing solution in this regard as it achieves a balanced tradeoff between privacy and utility without a trusted server. However, existing distributed DP mechanisms are impractical in the presence of client dropout, resulting in poor privacy guarantees or degraded training accuracy. In addition, these mechanisms suffer from severe efficiency issues.
Zhifeng Jiang 0001, Wei Wang 0030, Ruichuan Chen
EuroSys2
2024 Improved Bounds for Pure Private Agnostic Learning: Item-Level and User-Level Privacy
abstract
Machine Learning has made remarkable progress in a wide range of fields. In many scenarios, learning is performed on datasets involving sensitive information, in which privacy protection is essential for learning algorithms. In this work, we study pure private learning in the agnostic model -- a framework reflecting the learning process in practice. We examine the number of users required under item-level (where each user contributes one example) and user-level (where each user contributes multiple examples) privacy and derive several improved upper bounds. For item-level privacy, our algorithm achieves a near optimal bound for general concept classes. We extend this to the user-level setting, rendering a tighter upper bound than the one proved by Ghazi et al. (2023). Lastly, we consider the problem of learning thresholds under user-level privacy and present an algorithm with a nearly tight user complexity.
Bo Li 0001, Wei Wang 0030, Peng Ye 0005
ICML2
2024 Efficient and Straggler-Resistant Homomorphic Encryption for Heterogeneous Federated Learning
abstract
Cross-silo federated learning (FL) enables multiple institutions (clients) to collaboratively build a global model without sharing their private data. To prevent privacy leakage during aggregation, homomorphic encryption (HE) is widely used to encrypt model updates, yet incurs high computation and communication overheads. To reduce these overheads, packed HE (PHE) has been proposed to encrypt multiple plaintexts into a single ciphertext. However, the original design of PHE does not consider the heterogeneity among different clients, an intrinsic problem in cross-silo FL, often resulting in undermined training efficiency with slow convergence and stragglers. In this work, we propose FedPHE, an efficiently packed homomorphically encrypted FL framework with secure weighted aggregation and client selection to tackle the heterogeneity problem. Specifically, using CKKS with sparsification, FedPHE can achieve efficient encrypted weighted aggregation by accounting for contributions of local updates to the global model. To mitigate the straggler effect, we devise a sketching-based client selection scheme to cherry-pick representative clients with heterogeneous models and computing capabilities. We show, through rigorous security analysis and extensive experiments, that FedPHE can efficiently safeguard clients’ privacy, achieve a training speedup of 1.85 − 4.44×, cut the communication overhead by 1.24 − 22.62× , and reduce the straggler effect by up to 1.71 − 2.39×.
Nan Yan 0001, Yuqing Li 0001, Jing Chen 0003, Xiong Wang 0006, Jianan Hong, Kun He 0008, Wei Wang 0030
INFOCOM7
2024 The Limits of Differential Privacy in Online Learning
abstract
Differential privacy (DP) is a formal notion that restricts the privacy leakage of an algorithm when running on sensitive data, in which privacy-utility trade-off is one of the central problems in private data analysis. In this work, we investigate the fundamental limits of differential privacy in online learning algorithms and present evidence that separates three types of constraints: no DP, pure DP, and approximate DP. We first describe a hypothesis class that is online learnable under approximate DP but not online learnable under pure DP under the adaptive adversarial setting. This indicates that approximate DP must be adopted when dealing with adaptive adversaries. We then prove that any private online learner must make an infinite number of mistakes for almost all hypothesis classes. This essentially generalizes previous results and shows a strong separation between private and non-private settings since a finite mistake bound is always attainable (as long as the class is online learnable) when there is no privacy requirement.
Bo Li 0001, Wei Wang 0030, Peng Ye 0005
NeurIPS2
2024 Lotto: Secure Participant Selection against Adversarial Servers in Federated Learning
Zhifeng Jiang 0001, Peng Ye 0005, Shiqi He, Wei Wang 0030, Ruichuan Chen, Bo Li 0001
USENIX Security Symposium4
2024 MorphDAG: A Workload-Aware Elastic DAG-Based Blockchain
abstract
Directed Acyclic Graph(DAG)-based blockchain represents a paradigm shift from conventional blockchains, which has the potential to drastically improve throughput performance through concurrent storage and executions. In practice, however, existing DAG-based blockchains fail to deliver such promises, often with limited throughput, high conflicts, and security vulnerabilities under dynamic workloads. The root causes are their unawareness of the workload characteristics of different workload sizes and skewed access patterns. In this paper, we propose MorphDAG, the first workload-aware DAG-based blockchain that can significantly enhance throughput without compromising security and achieve elastic scaling under realistic workloads. We derive the theoretically optimal degree of storage concurrency to achieve high throughput while retaining system security as the workload size changes, while enabling fine-grained concurrency adjustment that accommodates aProof-of-Stake(PoS)-based consensus protocol. We develop a dual-mode transaction processing mechanism that effectively resolves the conflicts brought by skewed access. We implement a prototype of MorphDAG and evaluate under real-world workloads. Extensive evaluations demonstrate that MorphDAG improves end-to-end throughput by up to 2.3× and 2.4× over state-of-the-art DAG-based blockchain systems AdaptChain and OHIE, respectively.
Jiang Xiao 0001, Enping Wu, Bo Li 0001, Wei Wang 0030, Hai Jin 0001
IEEE Trans. Knowl. Data Eng.6
2024 Towards Efficient and Deposit-Free Blockchain-Based Spatial Crowdsourcing
abstract
Spatial crowdsourcing leverages the widespread use of mobile devices to outsource tasks to a crowd of users based on their geographical location. Despite its growing popularity, current crowdsourcing systems often suffer from a lack of transparency, centralization, and other security issues. Blockchain technology has revolutionized this sector with its potential for decentralization, security, and transparency. However, existing blockchain-based crowdsourcing systems often overlook efficient task assignment mechanisms and expose users to potential losses due to the obligatory deposit payments to smart contracts, which might be vulnerable or untrustworthy. This article proposes EDF-Crowd, an E fficient and D eposit- F ree blockchain-based spatial crowdsoucing framework, to address these challenges. EDF-Crowd introduces an efficient, customizable task assignment mechanism based on smart contracts, operating periodically and batch-wise. We also design a fair compensation mechanism to compensate users for the extra overhead caused by invoking certain smart contracts. More importantly, we propose a series of linkage protocols. By linking users’ back-and-forth actions, EDF-Crowd can regulate user behavior without requiring users to deposit. The versatility of EDF-Crowd also allows its application to generic crowdsourcing systems with minimal modifications. We implement EDF-Crowd based on the EOS blockchain. Extensive evaluations show that EDF-Crowd achieves high task assignment efficiency and low cost.
Wei Wang 0030, Jin Zhang 0001
ACM Trans. Sens. Networks2
2024 Synchronize Only the Immature Parameters: Communication-Efficient Federated Learning By Freezing Parameters Adaptively
abstract
Federated learning allows edge devices to collaboratively train a global model without sharing their local private data. Yet, with limited network bandwidth at the edge, communication often becomes a severe bottleneck. In this paper, we find that it is unnecessary to always synchronize the full model in the entire training process, because many parameters already become mature (i.e., stable) prior to model convergence, and can thus be excluded from later synchronizations. This allows us to reduce the communication overhead without compromising the model accuracy. However, challenges are that the local parameters excluded from global synchronization may diverge on different clients, and meanwhile some parameters may stabilize only temporally. To address these challenges, we propose a novel scheme called Adaptive Parameter Freezing (APF), which fixes (freezes) the non-synchronized stable parameters in intermittent periods. Specifically, the freezing periods are tentatively adjusted in an additively-increase and multiplicatively-decrease manner—depending on whether the previously-frozen parameters remain stable in subsequent iterations. We also extend APF into APF# and APF++, which freeze parameters in a more aggressive manner to achieve larger performance benefit for large complex models. We implemented APF and its variants as Python modules with PyTorch, and extensive experiments show that APF can reduce data transfer amount by over 60%.
Chen Chen 0067, Hong Xu 0001, Wei Wang 0030, Baochun Li, Bo Li 0001, Li Chen 0008, Gong Zhang 0001
IEEE Trans. Parallel Distributed Syst.3
2023 Golgi: Performance-Aware, Resource-Efficient Function Scheduling for Serverless Computing
abstract
This paper introduces Golgi, a novel scheduling system designed for serverless functions, with the goal of minimizing resource provisioning costs while meeting the function latency requirements. To achieve this, Golgi judiciously over-commits functions based on their past resource usage. To ensure overcommitment does not cause significant performance degradation, Golgi identifies nine low-level metrics to capture the runtime performance of functions, encompassing factors like request load, resource allocation, and contention on shared resources. These metrics enable accurate prediction of function performance using the Mondrian Forest, a classification model that is continuously updated in real-time for optimal accuracy without extensive offline training. Golgi employs a conservative exploration-exploitation strategy for request routing. By default, it routes requests to non-overcommitted instances to ensure satisfactory performance. However, it actively explores opportunities for using more resource-efficient overcommitted instances, while maintaining the specified latency SLOs. Golgi also performs vertical scaling to dynamically adjust the concurrency of overcommitted instances, maximizing request throughput and enhancing system robustness to prediction errors. We have prototyped Golgi and evaluated it in both EC2 cluster and a small production cluster. The results show that Golgi can meet the SLOs while reducing the resource provisioning cost by 42% (30%) in EC2 cluster (our production cluster).
Suyi Li 0002, Wei Wang 0030, Guangzhen Chen, Daohe Lu
SoCC2
2023 DeAR: Accelerating Distributed Deep Learning with Fine-Grained All-Reduce Pipelining
abstract
Communication scheduling has been shown to be effective in accelerating distributed training, which enables all-reduce communications to be overlapped with backpropagation computations. This has been commonly adopted in popular distributed deep learning frameworks. However, there exist two fundamental problems: (1) excessive startup latency proportional to the number of workers for each all-reduce operation; (2) it only achieves sub-optimal training performance due to the dependency and synchronization requirement of the feed-forward computation in the next iteration. We propose a novel scheduling algorithm, DeAR, that decouples the all-reduce primitive into two continuous operations, which overlaps with both backpropagation and feed-forward computations without extra communications. We further design a practical tensor fusion algorithm to improve the training performance. Experimental results with five popular models show that DeAR achieves up to 83% and 15% training speedup over the state-of-the-art solutions on a 64-GPU cluster with 10Gb/s Ethernet and 100Gb/s InfiniBand interconnects, respectively.
Lin Zhang 0059, Shaohuai Shi, Xiaowen Chu 0001, Wei Wang 0030, Bo Li 0001, Chengjian Liu
ICDCS4
2023 CoChain: High Concurrency Blockchain Sharding via Consensus on Consensus
abstract
Sharding is an effective technique to improve the scalability of blockchain. It splits nodes into multiple groups so that they can process transactions in parallel. To achieve higher parallelism and concurrency at large scales, it is desirable to maintain a large number of small shards. However, simply configuring small shards easily results in a higher fraction of malicious nodes inside shards, causing shard corruption and compromising system security. Existing sharding techniques hence demand large shards, at the expense of limited concurrency. To address this limitation, we propose CoChain: a blockchain sharding system that can securely configure small shards for enhanced concurrency. CoChain allows some shards to be corrupted. For security, each shard is monitored by multiple other shards. The latter reach a cross-shard Consensus on the Consensus results of their monitored shard. Once a corrupted shard is found, its subsequent consensus will be taken over by another shard, hence recovering the system. Via Consensus on Consensus, CoChain allows the existence of shards with more fraction of malicious nodes (<2/3) while securing the system, thus reducing the shard size safely. We implement CoChain based on Harmony and conduct extensive experiments. Compared with Harmony, CoChain achieves 35x throughput gain with 6,000+ nodes.
You Lin, Jin Zhang 0001, Wei Wang 0030
INFOCOM4
2023 Fast Sparse GPU Kernels for Accelerated Training of Graph Neural Networks
abstract
Graph Neural Networks (GNNs) are gaining huge traction recently as they achieve state-of-the-art performance on various graph-related problems. GNN training typically follows the standard Message Passing Paradigm, in which SpMM and SDDMM are the two essential sparse kernels. However, existing sparse GPU kernels are inefficient and may suffer from load imbalance, dynamics in GNN computing, poor memory efficiency, and tail effect. We propose two new kernels, Hybrid-Parallel SpMM (HP-SpMM) and Hybrid-Parallel SDDMM (HP-SDDMM), that efficiently perform SpMM and SDDMM on GPUs with a unified hybrid parallel strategy of mixing nodes and edges. In view of the emerging graph-sampling training, we design the Dynamic Task Partition (DTP) method to minimize the tail effect by exposing sufficient parallelism. We further devise the Hierarchical Vectorized Memory Access scheme to achieve aligned global memory accesses and enable vectorized instructions for improved memory efficiency. We also propose to enhance data locality by reordering the graphs with the Graph Clustering method. Experiments on extensive sparse matrices collected from real GNN applications demonstrate that our kernels achieve significant performance improvements over state-of-the-art implementations. We implement our sparse kernels in popular GNN frameworks and use them to train various GNN models, including the GCN model in full-graph mode and the GraphSAINT model in graph-sampling mode. Evaluation results show that our kernels can accelerate GNN training by up to 1.72×.
Ruibo Fan, Wei Wang 0030, Xiaowen Chu 0001
IPDPS2
2023 Following the Data, Not the Function: Rethinking Function Orchestration in Serverless Computing
Minchen Yu, Tingjia Cao, Wei Wang 0030, Ruichuan Chen
NSDI3
2023 Beware of Fragmentation: Scheduling GPU-Sharing Workloads with Fragmentation Gradient Descent
Qizhen Weng 0001, Lingyun Yang, Yinghao Yu, Wei Wang 0030, Xiaochuan Tang, Liping Zhang 0013
USENIX ATC4
2023 Monocular 3-D Object Detection Based on Depth-Guided Local Convolution for Smart Payment in D2D Systems
abstract
3-D object detection from mobile phones in Device-to-Device (D2D) system provides a new smart payment tool for the next generation of fintech, which is more flexible and efficient than the traditional barcode. In this article, we propose a monocular 3-D object detection method based on depth-guided local convolution. The method combines the information of RGB image mode and depth mode by using a convolution kernel through depth image and works on a single RGB image locally. According to the multiscale input information, the convolution kernel is adaptively adjusted to capture the target objects of different scales, so as to improve the performance of 3-D object detection. In addition, we use the soft-non-maximum suppression algorithm instead of traditional non-maximum suppression to select the best prediction box. In order to further improve the accuracy of 3-D object detection, the depth estimation network and 3-D object detection network are jointly trained in this method to make the two networks constrain each other and achieve the best performance.
Jun Li 0036, Yongbin Gao, Huixing Wang, Yier Yan, Bo Huang 0014, Jun Zhang 0004, Wei Wang 0030
IEEE Internet Things J.8
2023 GIFT: Toward Accurate and Efficient Federated Learning With Gradient-Instructed Frequency Tuning
abstract
Federated learning (FL) enables distributed clients to collectively train a global model without revealing their private data, and for efficiency clients synchronize their gradients periodically. However, this can lead to the inaccuracy in model convergence due to inconsistent data distributions among clients. In this work, we find that there is a strong correlation between FL accuracy loss and the synchronization frequency, and seek to fine tune the synchronization frequency at training runtime to make FL accurate and also efficient. Specifically, aware that under the FL privacy requirement only gradients can be utilized for making frequency tuning decisions, we propose a novel metric called gradient consistency, which can effectively reflect the training status despite the instability of realistic FL scenarios. We further devise a feedback-driven algorithm called Gradient-Instructed Frequency Tuning (GIFT), which adaptively increases or decreases the synchronization frequency based on the gradient consistency metric. We have implemented GIFT in PyTorch, and large-scale evaluations show that it can improve FL accuracy by up to 10.7% with a time reduction of 58.1%.
Chen Chen 0067, Hong Xu 0001, Wei Wang 0030, Baochun Li, Bo Li 0001, Li Chen 0008, Gong Zhang 0001
IEEE J. Sel. Areas Commun.3
2023 Towards Efficient Synchronous Federated Training: A Survey on System Optimization Strategies
abstract
The increasing demand for privacy-preserving collaborative learning has given rise to a new computing paradigm called federated learning (FL), in which clients collaboratively train a machine learning (ML) model without revealing their private training data. Given an acceptable level of privacy guarantee, the goal of FL is to minimize thetime-to-accuracyof model training. Compared with distributed ML in data centers, there are four distinct challenges to achieving short time-to-accuracy in FL training, namely the lack of information for optimization, the tradeoff between statistical and system utility, client heterogeneity, and large configuration space. In this paper, we survey recent works in addressing these challenges and present them following a typical training workflow through three phases: client selection, configuration, and reporting. We also review system works including measurement studies and benchmarking tools that aim to support FL developers.
Zhifeng Jiang 0001, Wei Wang 0030, Bo Li 0001, Qiang Yang 0001
IEEE Trans. Big Data2
2023 Accelerating Distributed Learning in Non-Dedicated Environments
abstract
Machine learning (ML) models are increasingly trained with distributed workers possessing heterogeneous resources. In such scenarios, model training efficiency may be negatively affected bystragglers—workers that run much slower than others. Efficient model training requires eliminating such stragglers, yet for modern ML workloads, existing load balancing strategies are inefficient and even infeasible. In this article, we propose a novel strategy, calledsemi-dynamic load balancing, to eliminate stragglers of distributed ML workloads. The key insight is that ML workers shall be load-balanced atiteration boundaries, being non-intrusive to intra-iteration execution. Based on it we further develop LB-BSP, an integrated worker coordination mechanism that adapts workers’ load to their instantaneous processing capabilities—by right-sizing the sample batches at the synchronization barriers. We have designed distinct load tuning algorithms for ML in CPU clusters, in GPU clusters as well as in federated learning setups, based on their respective characteristics. LB-BSP has been implemented as a Python module for ML frameworks like TensorFlow and PyTorch. Our EC2 deployment confirms that LB-BSP is practical, effective and light-weight, and is able to accelerating distributed training by up to 54 percent.
Chen Chen 0067, Qizhen Weng 0001, Wei Wang 0030, Baochun Li, Bo Li 0001
IEEE Trans. Cloud Comput.3
2023 Scalable K-FAC Training for Deep Neural Networks With Distributed Preconditioning
abstract
The second-order optimization methods, notably the D-KFAC (Distributed Kronecker Factored Approximate Curvature) algorithms, have gained traction on accelerating deep neural network (DNN) training on GPU clusters. However, existing D-KFAC algorithms require to compute and communicate a large volume of second-order information, i.e., Kronecker factors (KFs), before preconditioning gradients, resulting in large computation and communication overheads as well as a high memory footprint. In this paper, we propose DP-KFAC, a novel distributed preconditioning scheme that distributes the KF constructing tasks at different DNN layers to different workers. DP-KFAC not only retains the convergence property of the existing D-KFAC algorithms but also enables three benefits: reduced computation overhead in constructing KFs, no communication of KFs, and low memory footprint. Extensive experiments on a 64-GPU cluster show that DP-KFAC reduces the computation overhead by 1.55×-1.65×, the communication cost by 2.79×-3.15×, and the memory footprint by 1.14×-1.47× in each second-order update compared to the state-of-the-art D-KFAC methods. Our codes are available athttps://github.com/lzhangbv/kfac\_pytorch.
Lin Zhang 0059, Shaohuai Shi, Wei Wang 0030, Bo Li 0001
IEEE Trans. Cloud Comput.3
2023 LB-Chain: Load-Balanced and Low-Latency Blockchain Sharding via Account Migration
abstract
Blockchain sharding has been increasingly used to improve blockchain systems’ performance, in which a blockchain is split into multiple smaller, disjoint shards. In practice, however, sharding can only achieve limited throughput and latency improvement, especially for theuser-perceived transaction confirmation delay.The performance degradation is believed to be caused by the cross-shard transactions. However, we show, through comprehensive system deployment and measurement studies, that the main culprit is theimbalanced transaction loadon different blockchain shards. To address this problem, we propose a novel sharding system, called LB-Chain, whichdynamicallybalances the transaction load on different shards by periodicallymigrating active accountsfrom heavily-loaded shards to less-loaded ones. We have implemented a prototype of LB-Chain, and evaluated its performance through large-scale blockchain deployment using real-world transaction traces. Extensive experiments confirm that LB-Chain significantly boosts sharding performance, reducing the transaction confirmation delays by up to 90% while increasing the transaction throughput by more than 10%. The delay difference between different accounts is also reduced dramatically, leading to improved fairness in the system.
Wei Wang 0030, Jin Zhang 0001
IEEE Trans. Parallel Distributed Syst.2
2022 Pisces: efficient federated learning via guided asynchronous training
abstract
Federated learning (FL) is typically performed in a synchronous parallel manner, and the involvement of a slow client delays the training progress. Current FL systems employ a participant selection strategy to select fast clients with quality data in each iteration. However, this is not always possible in practice, and the selection strategy has to navigate a knotty tradeoff between the speed and the data quality.
Zhifeng Jiang 0001, Wei Wang 0030, Baochun Li, Bo Li 0001
SoCC2
2022 Owl: performance-aware scheduling for resource-efficient function-as-a-service cloud
abstract
This work documents our experience of improving the scheduler in Alibaba Function Compute, a public FaaS platform. It commences with our observation that memory and CPU are under-utilized in most FaaS sandboxes. A natural solution is to overcommit VM resources when allocating sandboxes, whereas the ensuing contention may cause performance degradation and compromise user experience. To complicate matters, the degradation in FaaS can arise from external factors, such as failed dependencies of user functions.
Huangshi Tian, Suyi Li 0002, Wei Wang 0030, Tianlong Wu
SoCC4
2022 Workload consolidation in alibaba clusters: the good, the bad, and the ugly
abstract
Web companies typically run latency-critical long-running services and resource-intensive, throughput-hungry batch jobs in a shared cluster for improved utilization and reduced cost. Despite many recent studies on workload consolidation, the production practice remains largely unknown. This paper describes our efforts to efficiently consolidate the two types of workloads in Alibaba clusters to support the company's e-commerce businesses.
Yongkang Zhang 0003, Yinghao Yu, Wei Wang 0030, Qiukai Chen, Tianchen Ding, Qizhen Weng 0001, Lingyun Yang, Jian He 0004, Liping Zhang 0013
SoCC3
2022 Jenga: Orchestrating Smart Contracts in Sharding-Based Blockchain for Efficient Processing
abstract
Sharding is a promising way to achieve blockchain scalability, increasing the throughput by partitioning nodes into multiple smaller groups, splitting the workload. However, when tackling the increasingly important smart contracts, existing blockchain sharding protocols do not scale well. They usually require complex multi-round cross-shard consensus protocols for contract execution and extensive cross-shard communication during state transmission, mainly because that each shard stores and executes an isolated, disjoint subset of contracts. In this paper, we present Jenga, a novel sharding-based approach for efficient smart contract processing. Its main idea is to break the isolation between shards by orchestrating the logic storage, state storage, and execution of smart contracts. In Jenga, all shards share the logic for all contracts. Therefore, multiple contracts involved in a smart contract transaction can be executed together by the same shard within one round. Moreover, different shards store distinct states (named state shards), several "orthogonal" execution channels are established based on the state shards, where each channel overlaps with all shards. Each node simultaneously belongs to a shard and an "orthogonal" channel, different channels execute different contracts. Therefore, via the overlapped nodes, the contract states can be directly broadcast between the state shards and the execution channels without additional cross-shard communication. We implement Jenga and evaluation results show that it provides outstanding performance gains in terms of throughput and transaction confirmation latency.
You Li 0004, Jin Zhang 0001, Wei Wang 0030
ICDCS4
2022 MLaaS in the Wild: Workload Analysis and Scheduling in Large-Scale Heterogeneous GPU Clusters
Qizhen Weng 0001, Wencong Xiao, Yinghao Yu, Wei Wang 0030, Jian He 0004, Yong Li 0045, Liping Zhang 0013, Wei Lin 0016
NSDI4
2022 An LLVM-based open-source compiler for NVIDIA GPUs
abstract
We present GASS, an LLVM-based open-source compiler for NVIDIA GPU's SASS machine assembly. GASS is the first open-source compiler targeting SASS, and it provides a unified toolchain for currently fragmented low-level performance research on NVIDIA GPUs. GASS supports all recent architectures, including Volta, Turing, and Ampere.
Da Yan 0002, Wei Wang 0030, Xiaowen Chu 0001
PPoPP2
2022 Incentivizing WiFi-Based Multilateration Location Verification
abstract
Due to the proliferation of WiFi devices and the high verification precision, researchers have shown interests in WiFi-based multilateration location verification (WMLV), where multiple WiFi APs (also known as verifiers) verify the location information claimed by a prover. However, it is a high expenditure for any single location-based service provider to deploy densely covered WiFi facilities. Incentivizing independent WiFi owners to corporately verify location information is thus a feasible solution to this plight, yet none of the previous research has taken this into consideration. To this point, we design a double auction-based incentive mechanism for WMLV, which motivates the participation of both provers and verifiers. More importantly, we consider practical situations, where the provers have various verification precision requirements, and different number of verifiers are required by different provers. The proposed double auction mechanism achieves desirable economical properties, includingtruthfulness, individual rationality, computational efficiency, budget balance,andnonnegative social welfare.The desired properties are validated through both theoretical analysis and extensive simulations.
Wei Wang 0030, Jin Zhang 0001, Qian Zhang 0001
IEEE Internet Things J.2
2022 Editorial: Advances in Mobile, Edge and Cloud Computing
Xiaowen Chu 0001, Hongbo Jiang 0001, Bo Li 0001, Dan Wang 0002, Wei Wang 0030
Mob. Networks Appl.5
2022 Towards Dependency-Aware Cache Management for Data Analytics Applications
abstract
Memory caches are being used aggressively in today's data analytics systems such as Spark, Tez, and Piccolo. The significant performance impact of caches and their limited sizes call for efficient cache management in data analytics clusters. However, prevalent data analytics systems employ rather simple cache management policies—notably Least Recently Used (LRU) and Least Frequently Used (LFU)—that areobliviousto the application semantics of data dependency, expressed as directed acyclic graphs (DAGs). Without this knowledge, cache management can, at best, be performed by “guessing” the future data access patterns based on history, which frequently results in inefficient, erroneous caching with a low hit rate and a long response time. Worse still, the lack of data dependency knowledge makes it impossible to retain theall-or-nothingcache property of cluster applications, in that a compute task cannot be sped up unless all the dependent data has been kept in the main memory. In this paper, we propose a novel cache replacement policy, named Least Reference Count (LRC), which exploits the application's data dependency information to optimize the cache management. LRC keeps track of thereference countof each data block, defined as the number of dependent child blocks that have not been computed yet, and always evicts the block with the smallest reference count. Furthermore, we incorporate the all-or-nothing requirement into LRC by coordinately managing the reference counts of all the input data blocks for the same computation. We demonstrate the efficacy of LRC through both empirical analysis and cluster deployments against popular benchmarking workloads. Our Spark implementation shows that, the proposed policies well address the all-or-nothing requirement and significantly improve the cache performance. Compared with LRU and a recently proposed caching policy called MEMTUNE, LRC improves the caching performance of typical workloads in production clusters by 22 and 284 percent, respectively.
Yinghao Yu, Chengliang Zhang, Wei Wang 0030, Jun Zhang 0004, Khaled Ben Letaief
IEEE Trans. Cloud Comput.3
2022 Enabling Cost-Effective, SLO-Aware Machine Learning Inference Serving on Public Cloud
abstract
The remarkable advances of Machine Learning (ML) have spurred an increasing demand for ML- as-a-Service on public cloud: developers train and publish ML models as online services to provide low-latency inference for dynamic queries. The primary challenge of ML model serving is to meet the response-time Service-Level Objectives (SLOs) of inference workloads while minimizing serving cost. In this article, we proposes MArk (Model Ark), a general-purpose inference serving system, to tackle the dual challenge of SLO compliance and cost effectiveness. MArk employs three design choices tailored to inference workload. First, MArk dynamically batches requests and opportunistically serves them using expensive hardware accelerators (e.g., GPU) for improved performance-cost ratio. Second, instead of relying on feedback control scaling or over-provisioning to serve dynamic workload, which can be too slow or too expensive, MArk employs predictive autoscaling to hide the provisioning latency at low cost. Third, given the stateless nature of inference serving, MArk exploits the flexible, yet costly serverless instances to cover occasional load spikes that are hard to predict. We evaluated the performance of MArk using several state-of-the-art ML models trained in TensorFlow, MXNet, and Keras. Compared with the premier industrial ML serving platform SageMaker, MArk reduces the serving cost up to$7.8\times$while achieving even better latency performance.
Chengliang Zhang, Minchen Yu, Wei Wang 0030, Feng Yan 0001
IEEE Trans. Cloud Comput.3
2021 George: Learning to Place Long-Lived Containers in Large Clusters with Operation Constraints
abstract
Online cloud services are widely deployed as Long-Running Applications (LRAs) hosted in containers. Placing LRA containers turns out to be particularly challenging due to the complex interference between co-located containers and the operation constraints in production clusters such as fault tolerance, disaster avoidance and incremental deployment. Existing schedulers typically provide APIs for operators to manually specify the container scheduling requirements and offer only qualitative scheduling guidelines for container placement. Such schedulers, do not perform well in terms of both performance and scale, while also requiring manual intervention.
Suyi Li 0002, Wei Wang 0030, Yinghao Yu, Bo Li 0001
SoCC3
2021 Morphling: Fast, Near-Optimal Auto-Configuration for Cloud-Native Model Serving
abstract
Machine learning models are widely deployed in production cloud to provide online inference services. Efficiently deploying inference services requires careful tuning of hardware and runtime configurations (e.g., GPU type, GPU memory, batch size), which can significantly improve the model serving performance and reduce cost. However, existing autoconfiguration approaches for general workloads, such as Bayesian optimization and white-box prediction, are inefficient in navigating the high-dimensional configuration space of model serving, incurring high sampling cost.
Lingyun Yang, Yinghao Yu, Wei Wang 0030, Bo Li 0001, Xianchao Sun, Jian He 0004, Liping Zhang 0013
SoCC4
2021 Citadel: Protecting Data Privacy and Model Confidentiality for Collaborative Learning
abstract
Many organizations own data but have limited machine learning expertise (data owners). On the other hand, organizations that have expertise need data from diverse sources to train truly generalizable models (model owners). With the advancement of machine learning (ML) and its growing awareness, the data owners would like to pool their data and collaborate with model owners, such that both entities can benefit from the obtained models. In such a collaboration, the data owners want to protect the privacy of its training data, while the model owners desire the confidentiality of the model and the training method that may contain intellectual properties. Existing private ML solutions, such as federated learning and split learning, cannot simultaneously meet the privacy requirements of both data and model owners.
Chengliang Zhang, Junzhe Xia, Baichen Yang, Huancheng Puyang, Wei Wang 0030, Ruichuan Chen, Istemi Ekin Akkus, Paarijaat Aditya, Feng Yan 0001
SoCC5
2021 Communication-Efficient Federated Learning with Adaptive Parameter Freezing
abstract
Federated learning allows edge devices to collaboratively train a global model by synchronizing their local updates without sharing private data. Yet, with limited network bandwidth at the edge, communication often becomes a severe bottleneck. In this paper, we find that it is unnecessary to always synchronize the full model in the entire training process, because many parameters gradually stabilize prior to the ultimate model convergence, and can thus be excluded from being synchronized at an early stage. This allows us to reduce the communication overhead without compromising the model accuracy. However, challenges are that the local parameters excluded from global synchronization may diverge on different clients, and meanwhile some parameters may stabilize only temporally. To address these challenges, we propose a novel scheme called Adaptive Parameter Freezing (APF), which fixes (freezes) the non-synchronized stable parameters in intermittent periods. Specifically, the freezing periods are tentatively adjusted in an additively-increase and multiplicatively-decrease manner, depending on if the previously-frozen parameters remain stable in subsequent iterations. We implemented APF as a Python module in PyTorch. Our extensive array of experimental results show that APF can reduce data transfer by over 60%.
Chen Chen 0067, Hong Xu 0001, Wei Wang 0030, Baochun Li, Bo Li 0001, Li Chen 0008, Gong Zhang 0001
ICDCS3
2021 Gillis: Serving Large Neural Networks in Serverless Functions with Automatic Model Partitioning
abstract
The increased use of deep neural networks has stimulated the growing demand for cloud-based model serving platforms. Serverless computing offers a simplified solution: users deploy models as serverless functions and let the platform handle provisioning and scaling. However, serverless functions have constrained resources in CPU and memory, making them inefficient or infeasible to serve large neural networks-which have become increasingly popular. In this paper, we present Gillis, a serverless-based model serving system that automatically partitions a large model across multiple serverless functions for faster inference and reduced memory footprint per function. Gillis employs two novel model partitioning algorithms that respectively achieve latency-optimal serving and cost-optimal serving with SLO compliance. We have implemented Gillis on three serverless platforms-AWS Lambda, Google Cloud Functions, and KNIX-with MXNet as the serving backend. Experimental evaluations against popular models show that Gillis supports serving very large neural networks, reduces the inference latency substantially, and meets various SLOs with a low serving cost.
Minchen Yu, Zhifeng Jiang 0001, Hok Chun Ng, Wei Wang 0030, Ruichuan Chen, Bo Li 0001
ICDCS4
2021 Simplifying low-level GPU programming with GAS
abstract
Many low-level optimizations for NVIDIA GPU can only be implemented in native hardware assembly (SASS). However, programming in SASS is unproductive and not portable.
Da Yan 0002, Wei Wang 0030, Xiaowen Chu 0001
PPoPP2
2021 CrystalPerf: Learning to Characterize the Performance of Dataflow Computation through Code Analysis
Huangshi Tian, Minchen Yu, Wei Wang 0030
USENIX ATC3
2021 Toward Privacy-Preserving Task Assignment for Fully Distributed Spatial Crowdsourcing
abstract
With the proliferation of human-carried mobile devices, spatial crowdsourcing has emerged as a transformative system, where requesters outsource their spatiotemporal tasks to a set of workers who are willing to perform the tasks at the specified locations. However, in order to make efficient assignments, the existing spatial crowdsourcing system usually requires workers and/or tasks to expose their locations, which raises a significant concern of compromising location privacy. In addition, traditional spatial crowdsourcing systems employ a centralized server to manage the information of workers and tasks. Such a centralized design does not scale to a large number of workers/tasks, making the server easily a bottleneck. In this article, we present an online framework for assigning tasks to workers without compromising the location privacy in a fully distributed manner. Our system protects the location privacy of both workers and tasks through homomorphic encryption. We further propose a novel wait-and-decide mechanism and a proportional-backoff mechanism to increase the number of assigned tasks. Extensive experiments on real-world data sets illustrate that our proposed system achieves a large number of task assignments in an efficient and privacy-preserving manner.
Jingrou Wu, Wei Wang 0030, Jin Zhang 0001
IEEE Internet Things J.3
2020 Semi-dynamic load balancing: efficient distributed learning in non-dedicated environments
abstract
Machine learning (ML) models are increasingly trained in clusters with non-dedicated workers possessing heterogeneous resources. In such scenarios, model training efficiency can be negatively affected by stragglers---workers that run much slower than others. Efficient model training requires eliminating such stragglers, yet for modern ML workloads, existing load balancing strategies are inefficient and even infeasible. In this paper, we propose a novel strategy called semi-dynamic load balancing to eliminate stragglers of distributed ML workloads. The key insight is that ML workers shall be load-balanced at iteration boundaries, being non-intrusive to intra-iteration execution. We develop LB-BSP based on such an insight, which is an integrated worker coordination mechanism that adapts workers' load to their instantaneous processing capabilities by right-sizing the sample batches at the synchronization barriers. We have custom-designed the batch sizing algorithm respectively for CPU and GPU clusters based on their own characteristics. LB-BSP has been implemented as a Python module for ML frameworks like TensorFlow and PyTorch. Our EC2 deployment confirms that LB-BSP is practical, effective and light-weight, and is able to accelerating distributed training by up to 54%.
Chen Chen 0067, Qizhen Weng 0001, Wei Wang 0030, Baochun Li, Bo Li 0001
SoCC3
2020 RepBun: Load-Balanced, Shuffle-Free Cluster Caching for Structured Data
abstract
Cluster caching systems increasingly store structured data objects in the columnar format. However, these systems routinely face the imbalanced load that significantly impairs the I/O performance. Existing load-balancing solutions, while effective for reading unstructured data objects, fall short in handling columnar data. Unlike unstructured data that can only be read through a full-object scan, columnar data supports direct query of specific columns with two distinct access patterns: (1) columns have the heavily skewed popularity, and (2) hot columns are likely accessed together in a query job. Based on these two access patterns, we propose an effective load-balancing solution for structured data. Our solution, which we call RepBun, groups hot columns into a bundle. It then copies multiple replicas of the column bundle and stores them uniformly across servers. We show that RepBun achieves improved load balancing with reduced memory overhead, while avoiding data shuffling between cache servers. We implemented RepBun atop Alluxio, a popular in-memory distributed storage, and evaluate its performance through EC2 deployment against the TPC-H benchmark work-load. Experimental results show that RepBun outperforms the existing load-balancing solutions with significantly shorter read latency and faster query completion.
Minchen Yu, Yinghao Yu, Yunchuan Zheng, Baichen Yang, Wei Wang 0030
INFOCOM5
2020 Demystifying Tensor Cores to Optimize Half-Precision Matrix Multiply
abstract
Half-precision matrix multiply has played a key role in the training of deep learning models. The newly designed Nvidia Tensor Cores offer the native instructions for half-precision small matrix multiply, based on which Half-precision General Matrix Multiply (HGEMM) routines are developed and can be accessed through high-level APIs. In this paper, we, for the first time, demystify how Tensor Cores on NVIDIA Turing architecture work in great details, including the instructions used, the registers and data layout required, as well as the throughput and latency of Tensor Core operations. We further benchmark the memory system of Turing GPUs and conduct quantitative analysis of the performance. Our analysis shows that the bandwidth of DRAM, L2 cache and shared memory is the new bottleneck for HGEMM, whose performance is previously believed to be bound by computation. Based on our newly discovered features of Tensor Cores, we apply a series of optimization techniques on the Tensor Core-based HGEMM, including blocking size optimization, data layout redesign, data prefetching, and instruction scheduling. Extensive evaluation results show that our optimized HGEMM routine achieves an average of 1.73× and 1.46× speedup over the native implementation of cuBLAS 10.1 on NVIDIA Turing RTX2070 and T4 GPUs, respectively. The code of our implementation is written in native hardware assembly (SASS).
Da Yan 0002, Wei Wang 0030, Xiaowen Chu 0001
IPDPS2
2020 Not All Explorations Are Equal: Harnessing Heterogeneous Profiling Cost for Efficient MLaaS Training
abstract
Machine-Learning-as-a-Service (MLaaS) enables practitioners and AI service providers to train and deploy ML models in the cloud using diverse and scalable compute resources. A common problem for MLaaS users is to choose from a variety of training deployment options, notably scale-up (using more capable instances) and scale-out (using more instances), subject to the budget limits and/or time constraints. State-of-the-art (SOTA) approaches employ analytical modeling for finding the optimal deployment strategy. However, they have limited applicability as they must be tailored to specific ML model architectures, training framework, and hardware. To quickly adapt to the fast evolving design of ML models and hardware infrastructure, we propose a new Bayesian Optimization (BO) based method HeterBO for exploring the optimal deployment of training jobs. Unlike the existing BO approaches for general applications, we consider the heterogeneous exploration cost and machine learning specific prior to significantly improve the search efficiency. This paper culminates in a fully automated MLaaS training Cloud Deployment system (MLCD) driven by the highly efficient HeterBO search method. We have extensively evaluated MLCD in AWS EC2, and the experimental results show that MLCD outperforms two SOTA baselines, conventional BO and CherryPick, by 3.1× and 2.34×, respectively.
Chengliang Zhang, Wei Wang 0030, Cheng Li 0001, Feng Yan 0001
IPDPS3
2020 Optimizing batched Winograd convolution on GPUs
abstract
In this paper, we present an optimized implementation for single-precision Winograd convolution on NVIDIA Volta and Turing GPUs. Compared with the state-of-the-art Winograd convolution in cuDNN 7.6.1, our implementation achieves up to 2.13X speedup on Volta V100 and up to 2.65X speedup on Turing RTX2070. On both Volta and Turing GPUs, our implementation achieves up to 93% of device peak.
Da Yan 0002, Wei Wang 0030, Xiaowen Chu 0001
PPoPP2
2020 Metis: learning to schedule long-running applications in shared container clusters at scale
abstract
Online cloud services are increasingly deployed as long-running applications (LRAs) in containers. Placing LRA containers is known to be difficult as they often have sophisticated resource interferences and I/O dependencies. Existing schedulers rely on operators to manually express the container scheduling requirements as placement constraints and strive to satisfy as many constraints as possible. Such schedulers, however, fall short in performance as placement constraints only provide qualitative scheduling guidelines and minimizing constraint violations does not necessarily result in the optimal performance.In this work, we present Metis, a general-purpose scheduler that learns to optimally place LRA containers using deep reinforcement learning (RL) techniques. This eliminates the complex manual specification of placement constraints and offers, for the first time, concrete quantitative scheduling criteria. As directly training an RL agent does not scale, we develop a novel hierarchical learning technique that decomposes a complex container placement problem into a hierarchy of subproblems with significantly reduced state and action space. We show that many subproblems have similar structures and can hence be solved by training a unified RL agent offline. Large-scale EC2 deployment shows that compared with the traditional constraint-based schedulers, Metis improves the throughput by up to 61%, optimizes various performance metrics, and easily scales to a large cluster where 3K containers run on over 700 machines.
Qizhen Weng 0001, Wei Wang 0030, Chen Chen 0067, Bo Li 0001
SC3
2020 BatchCrypt: Efficient Homomorphic Encryption for Cross-Silo Federated Learning
Chengliang Zhang, Suyi Li 0002, Junzhe Xia, Wei Wang 0030, Feng Yan 0001, Yang Liu 0165
USENIX ATC4
2020 Achieving Load-Balanced, Redundancy-Free Cluster Caching with Selective Partition
abstract
Data-intensive clusters increasingly rely on in-memory storages to improve I/O performance. However, the routinely observed file popularity skew and load imbalance create hot spots, which significantly degrade the benefits of in- memory caching. Common approaches to tame load imbalance include copying multiple replicas of hot files and creating parity chunks using storage codes. Yet, these techniques either suffer from high memory overhead due to cache redundancy or incur non-trivial encoding/decoding complexity. In this paper, we propose an effective approach to achieve load balancing without cache redundancy or encoding/decoding overhead. Our solution, termed SP-Cache,selectively partitionsfiles based on the loads they contribute and evenly caches those partitions across the cluster. We develop an efficient algorithm to determine the optimal number of partitions for a hot file—too few partitions are incapable of mitigating hot spots, while too many are susceptible to stragglers. We have implemented SP-Cache atop Alluxio, a popular in-memory distributed storage system, and evaluated its performance through EC2 deployment and trace-driven simulations. SP-Cache can quickly react to the changing load by dynamically re-balancing cache servers. Compared to the state-of-the-art solution, SP-Cache reduces the file access latency by up to 40 percent in both the mean and the tail, using 40 percent less memory.
Yinghao Yu, Wei Wang 0030, Renfei Huang, Jun Zhang 0004, Khaled Ben Letaief
IEEE Trans. Parallel Distributed Syst.2
2019 Characterizing and Synthesizing Task Dependencies of Data-Parallel Jobs in Alibaba Cloud
abstract
Cluster schedulers routinely face data-parallel jobs with complex task dependencies expressed as DAGs (directed acyclic graphs). Understanding DAG structures and runtime characteristics in large production clusters hence plays a key role in scheduler design, which, however, remains an important missing piece in the literature. In this work, we present a comprehensive study of a recently released cluster trace in Alibaba. We examine the dependency structures of Alibaba jobs and find that their DAGs have sparsely connected vertices and can be approximately decomposed into multiple trees with bounded depth. We also characterize the runtime performance of DAGs and show that dependent tasks may have significant variability in resource usage and duration---even for recurring tasks. In both aspects, we compare the query jobs in the standard TPC benchmarks with the production DAGs and find the former inadequately representative. To better benchmark DAG schedulers at scale, we develop a workload generator that can faithfully synthesize task dependencies based on the production Alibaba trace. Extensive evaluations show that the synthesized DAGs have consistent statistical characteristics as the production DAGs, and the synthesized and real workloads yield similar scheduling results with various schedulers.
Huangshi Tian, Yunchuan Zheng, Wei Wang 0030
SoCC3
2019 CMFL: Mitigating Communication Overhead for Federated Learning
abstract
Federated Learning enables mobile users to collaboratively learn a global prediction model by aggregating their individual updates without sharing the privacy-sensitive data. As mobile devices usually have limited data plan and slow network connections to the central server where the global model is maintained, mitigating the communication overhead is of paramount importance. While existing works mainly focus on reducing the total bits transferred in each update via data compression, we study an orthogonal approach that identifies irrelevant updates made by clients and precludes them from being uploaded for reduced network footprint. Following this idea, we propose Communication-Mitigated Federated Learning (CMFL) in this paper. CMFL provides clients with feedback information regarding the global tendency of model updating. Each client checks if its update aligns with this global tendency and is relevant enough to model improvement. By avoiding uploading those irrelevant updates to the server, CMFL can substantially reduce the communication overhead while still guaranteeing the learning convergence. CMFL is shown to achieve general improvement in communication efficiency for almost all of the existing federated learning schemes. We evaluate CMFL through extensive simulations and EC2 emulations. Compared with vanilla Federated Learning, CMFL yields 13.97x communication efficiency in terms of the reduction of network footprint. When applied to Federated Multi-Task Learning, CMFL improves the communication efficiency by 5.7x with 4% higher prediction accuracy.
Wei Wang 0030, Bo Li 0001
ICDCS2
2019 LACS: Load-Aware Cache Sharing with Isolation Guarantee
abstract
Cluster caching has been increasingly deployed in front of cloud storage to improve I/O performance. In shared, multi-tenant environments such as cloud datacenters, cluster caches are constantly contended by many users. Enforcing performance isolation between users hence becomes imperative to cluster caching. A user's caching performance critically depends on two factors: (1) the amount of cache allocation and (2) the load of servers in which its files are cached. However, existing cache sharing policies only provide guarantees on the amount of cache allocation, while remaining agnostic to the load of cache servers. Consequently, "mice" users having files co-located with "elephants" contributing heavy data accesses may experience extremely long latency, hence receiving no isolation. In this paper, we propose a Load-Aware Cache Sharing scheme (LACS) to enforce isolation between users. LACS keeps track of the load contributed by each user and reins back the congestions caused by elephant users by throttling their cache usage and network bandwidth. We have implemented LACS atop Alluxio, a popular cluster caching system. EC2 deployment shows that LACS achieves performance isolation in the presence of elephants, while improving the mean read latency by up to 80.4% (25.3% on average) over the state-of-the-art load balancing technique.
Yinghao Yu, Wei Wang 0030, Jun Zhang 0004, Khaled Ben Letaief
ICDCS2
2019 Round-Robin Synchronization: Mitigating Communication Bottlenecks in Parameter Servers
abstract
Deep learning is usually performed in GPU clusters where each worker machine iteratively refines the model parameters by communicating the update with the Parameter Server (PS). More often than not, workers communicate in a synchronous manner, so as to avoid using out-of-dated parameters and make high-quality refinement in each iteration. However, as all workers synchronize with the PS simultaneously, the communication becomes a severe bottleneck. To address this problem, in this paper we propose the Round-Robin Synchronous Parallel (R2SP) scheme, which coordinates workers to make updates in an evenly-gapped, round-robin manner. This way, R2SP can minimize the network contention at a minimum cost of the refinement quality. We further extend R2SP to heterogeneous clusters by adaptively tuning the batch size of each worker based on its processing capability. We have implemented R2SP as a ready-to-use python library for status-quo deep learning frameworks. EC2 deployment in GPU clusters show that R2SP effectively mitigates the communication bottlenecks, accelerating the training of popular image classification models by up to 25%.
Chen Chen 0067, Wei Wang 0030, Bo Li 0001
INFOCOM2
2019 MArk: Exploiting Cloud Services for Cost-Effective, SLO-Aware Machine Learning Inference Serving
Chengliang Zhang, Minchen Yu, Wei Wang 0030, Feng Yan 0001
USENIX ATC3
2018 Fast Distributed Deep Learning via Worker-adaptive Batch Sizing
abstract
In heterogeneous or shared clusters, distributed learning processes are slowed down by straggling workers. In this work, we propose LB-BSP, a new synchronization scheme that eliminates stragglers by adapting each worker's training load (batch size) to its processing capability. For training in shared production clusters, a prerequisite for deciding the workers' batch sizes is to know their processing speeds before each iteration starts. To this end, we adopt NARX, an extended recurrent neural network that accounts for both the historical speeds and the driving factors such as CPU and memory in prediction.
Chen Chen 0067, Qizhen Weng 0001, Wei Wang 0030, Baochun Li, Bo Li 0001
SoCC3
2018 Continuum: A Platform for Cost-Aware, Low-Latency Continual Learning
abstract
Many machine learning applications operate in dynamic environments that change over time, in which models must be continually updated to capture the recent trend in data. However, most of today's learning frameworks perform training offline, without a system support for continual model updating.
Huangshi Tian, Minchen Yu, Wei Wang 0030
SoCC3
2018 Unraveling the RTT-fairness Problem for BBR: A Queueing Model
abstract
BBR is a congestion-based congestion control algorithm recently proposed by Google. It proactively measures the bottleneck bandwidth and round trip times (RTTs) of a connection pipe, based on which it governs its sending behaviors. Despite the significant throughput gains and latency reduction, some experimental studies reveal that BBR may result in a salient RTT-fairness problem, in that short-RTT flows can be starved of bandwidth allocation when comnetina with lons-R'I'T flows. In this paper, we study BBR's RTT-fairness problem from a theoretic perspective. We present a closed-form solution that characterizes the intrinsic dynamics of BBR flows and their interactions. Specifically, we model BBR's sending behaviors and bandwidth dynamics, based on which we establish an exponential relationship between the flows' bandwidth shares and their RTTs. We show that the degree of unfairness is dictated by the RTT ratio between two flows, irrespective of the other network parameters, such as the initial sending rates or link capacity. In particular, when the RTT ratio of the two flows is greater than 2, the short-RTT flow is starved of bandwidth allocation ( ≤ 0.1%), Our theoretical results are corroborated by simulations in a wide range of settings.
Yuechen Tao, Jingjie Jiang, Shiyao Ma, Wei Wang 0030, Bo Li 0001
GLOBECOM5
2018 Fair Coflow Scheduling without Prior Knowledge
abstract
Coflow scheduling improves the networking performance at the application level in datacenters. Ideally, a coflow scheduler should provide tenants with isolation guarantees to achieve predictable networking performance. Existing works in this regard (e.g., DRF [1] and HUG [2]) are limited to the clairvoyant scheduling, in that the complete knowledge of coflow sizes is assumed to be available before the communication starts. However, this assumption does not hold for many applications with pipelined computation, in which clairvoyant coflow schedulers become inapplicable. To bridge this gap, we develop a new non-clairvoyant coflow scheduler, called Non-Clairvoyant DRF (NC-DRF), which provides isolation guarantees between contending coflows without prior knowledge of coflow size. We show that NC-DRF achieves provable isolation guarantees in the long run. Cluster deployment and trace-driven simulations show that with NC-DRF, coflows are only delayed by 68% on average as compared with the clairvoyant, isolation-optimal DRF [1]. NC-DRF also outperforms existing alternatives (e.g., per-link fairness) by 1.7× in terms of the average coflow completion time.
Wei Wang 0030
ICDCS2
2018 OpuS: Fair and Efficient Cache Sharing for In-Memory Data Analytics
abstract
We study the fair cache allocation problem in shared cloud environments, where many users and applications contend for the main memory to cache shared datasets or files. Unlike other resources such as CPUs and networks, in-memory caches can be non-exclusively shared across many users, e.g., a cached columnar dataset queried by many Spark SQL jobs. This results in a unique challenge of the "free-riding" problem, where a user lies about its caching preferences to trick other users to cache files for it, using their allocated cache space. We show that existing cache allocation policies either suffer from such manipulations or result in poor efficiency. To address this problem, we propose a new cache allocation algorithm, termed OpuS, or Opportunistic Sharing for high efficiency. We show that OpuS provides performance isolation between users and is strategy-proof against "free-riding" manipulations. We have implemented OpuS as a pluggable cache manager in Alluxio, a popular memory-centric filesystem. Cluster deployment and trace-driven simulations demonstrate that OpuS allocates each user a fair share of caches while achieving near-optimal efficiency in cache utilization.
Yinghao Yu, Wei Wang 0030, Jun Zhang 0004, Qizhen Weng 0001, Khaled Ben Letaief
ICDCS2
2018 Stay Fresh: Speculative Synchronization for Fast Distributed Machine Learning
abstract
Large machine learning models are typically trained in parallel and distributed environments. The model parameters are iteratively refined by multiple worker nodes in parallel, each processing a subset of the training data. In practice, the training is usually conducted in an asynchronous parallel manner, where workers can proceed to the next iteration before receiving the latest model parameters. While this maximizes the rate of updates, the price paid is compromised training quality as the computation is usually performed using stale model parameters. To address this problem, we propose a new scheme, termed speculative synchronization. Our scheme allows workers to speculate about the recent parameter updates from others on the fly, and if necessary, the workers abort the ongoing computation, pull fresher parameters, and start over to improve the quality of training. We design an effective heuristic algorithm to judiciously determine when to restart training iterations with fresher parameters by quantifying the gain and loss. We implement our scheme in MXNet-a popular machine learning framework-and demonstrate its effectiveness through cluster deployment atop Amazon EC2. Experimental results show that speculative synchronization achieves up to 3× speedup over the asynchronous parallel scheme in many machine learning applications, with little additional communication overhead.
Chengliang Zhang, Huangshi Tian, Wei Wang 0030, Feng Yan 0001
ICDCS3
2018 Performance-Aware Fair Scheduling: Exploiting Demand Elasticity of Data Analytics Jobs
abstract
Efficient resource management is of paramount importance in today's production clusters. In this paper, we identify the demand elasticity of data-parallel jobs. Demand elasticity allows jobs to run with a significantly less amount of resources than they ideally need, at the expense of only a modest performance penalty. Our EC2 experiment using popular Spark benchmark suites confirms that running a job using 50% of demanded slots is sufficient to achieve at least 75% of the ideal performance. We show that such an elasticity is an intrinsic property of data-parallel jobs and can be exploited to speed up average job completion. In this regard, we propose Performance-Aware Fair (PAF) scheduler to identify the demand elasticity and use it to improve the average job performance, while still attaining near-optimal isolation guarantee close to fair sharing. PAF starts with a fair allocation and iteratively adjusts it by transferring resources from one job to another, improving the performance of resource-taker without penalizing resource-giver by a noticeable amount. We implemented PAF in Spark and evaluated its effectiveness through both EC2 experiments and large-scale simulations. Evaluation results show that compared with fair allocation, PAF improves the average job performance by 13%, while penalizing resource-givers by no more than 1%.
Chen Chen 0067, Wei Wang 0030, Bo Li 0001
INFOCOM2
2018 Utopia: Near-optimal Coflow Scheduling with Isolation Guarantee
abstract
Performance and service isolation come as two top objectives for coflow scheduling. However, the common wisdom is that these two objectives are often conflicting with each other and cannot be achieved simultaneously. Existing coflow scheduling frameworks either focus only on minimizing the average coflow completion time (CCT) (e.g., Varys), or providing optimal isolation between contending coflows by means of fair network sharing (e.g., HUG). In this paper, we make an attempt to achieve the best of both worlds through a novel coflow scheduler, Utopia, to attain near-optimal performance with provable isolation guarantee. This is particularly challenging given the correlation of bandwidth demands across multiple links from coflows. We show that Utopia is capable of reducing the average CCT dramatically, while still guaranteeing that no coflow will ever be delayed beyond a constant time than its CCT in a fair scheme. Both trace-driven simulation and EC2 deployment confirm that Utopia outperforms the fair sharing policy by 1.8 × in terms of average CCT, while producing no completion time delay for a single coflow. Even compared with performance-optimal Varys, Utopia speeds up average coflow completion by 9%.
Wei Wang 0030, Bo Li 0001
INFOCOM2
2018 SP-cache: load-balanced, redundancy-free cluster caching with selective partition
Yinghao Yu, Renfei Huang, Wei Wang 0030, Jun Zhang 0004, Khaled Ben Letaief
SC3
2017 Towards Online Checkpointing Mechanism for Cloud Transient Servers
abstract
Cloud providers such as Amazon EC2 and Google GCE typically price their idle compute resources at a significant discount in the form of transient servers. Unlike regular servers, cloud transient servers have no reliability guarantee and can be revoked anytime. In order to enjoy the price premium of transient servers while mitigating the impact of indefinite server revocations, cloud users continuously checkpoint current computational states to stable storage. To determine the optimal checkpointing interval, prior works have largely built upon predicting future revocations. However, we argue that existing prediction models may not work well if providers employ different server revocation algorithms in the future, and the required historical information can be unavailable in practice. In this work, we propose two online algorithms, one competitive and the other heuristic, that determine the checkpointing intervals without predicting future revocation based on history. Our competitive online algorithm achieves the best possible competitive ratio against the optimal, yet impractical, offline solution. Our heuristic algorithm adaptively reacts to server revocations, achieving better average performance than the competitive algorithm, especially for short-running tasks. Extensive trace-driven simulations confirm our analytical results and demonstrate the efficacy of the two online algorithms in the context of Amazon EC2 Spot Instances.
Wei Wang 0030, Bo Li 0001
GLOBECOM2
2017 LERC: Coordinated Cache Management for Data-Parallel Systems
abstract
Memory caches are being aggressively used in today's data- parallel frameworks such as Spark, Tez and Storm. By caching input and intermediate data in memory, compute tasks can witness speedup by orders of magnitude. To maximize the chance of in-memory data access, existing cache algorithms, be it recency- or frequency-based, settle on cache hit ratio as the optimization objective. However, unlike the conventional belief, we show in this paper that simply pursuing a higher cache hit ratio of individual data blocks does not necessarily translate into faster task completion in data-parallel environments. A data-parallel task typically depends on multiple input data blocks. Unless all of these blocks are cached in memory, no speedup will result. To capture this all-or-nothing property, we propose a more relevant metric, called effective cache hit ratio. Specifically, a cache hit of a data block is said to be effective if it can speed up a compute task. In order to optimize the effective cache hit ratio, we propose the Least Effective Reference Count (LERC) policy that persists the dependent blocks of a compute task as a whole in memory. We have implemented the LERC policy as a memory manager in Spark and evaluated its performance through Amazon EC2 deployment. Evaluation results demonstrate that LERC helps speed up data-parallel jobs by up to 37% compared with the widely employed least-recently-used (LRU) policy.
Yinghao Yu, Wei Wang 0030, Jun Zhang 0004, Khaled Ben Letaief
GLOBECOM2
2017 Speculative Slot Reservation: Enforcing Service Isolation for Dependent Data-Parallel Computations
abstract
Priority scheduling is a fundamental tool to provide service isolation for different jobs in shared clusters. Ideally, the performance of a high-priority job should not be dragged down by another with a lower priority. However, we show in this paper that simply assigning a high priority provides no isolation for jobs with dependent computations. A job, even receiving the highest priority, may give up compute slots to another before proceeding to the downstream computation, which is because of barrier, i.e., that the downstream computation cannot start until all the upstream tasks have completed. Such an interruption of execution inevitably results in a significant delay. In this paper, we propose speculative slot reservation that judiciously reserves slots for downstream computations, so as to retain service isolation for high-priority jobs. To mitigate the utilization loss due to slot reservation, we analyze the trade-off between utilization and isolation, and expose a tunable knob to navigate the trade-off. We also propose a complementary straggler mitigation strategy that uses the reserved slots to run extra copies of slow tasks. We have implemented speculative slot reservation in Spark. Evaluations based on both cluster deployment and trace-driven simulations show that our approach enforces strict service isolation for high-priority jobs, without slowing down the other jobs with a lower priority.
Chen Chen 0067, Wei Wang 0030, Bo Li 0001
ICDCS2
2017 Cluster fair queueing: Speeding up data-parallel jobs with delay guarantees
abstract
Cluster scheduler serves as a critical component to data-parallel systems in datacenters. Ideally, a scheduler should provide predictable performance with guarantees on the maximal job completion delay, while at the same time ensuring the minimal mean response time. Practically however, performance predictability and optimality are often conflicting with each other. The results often are a plethora of scheduling policies that either achieve predictable performance at the expense of long response times (e.g., max-min fairness), or run the risk of starving some jobs to obtain the minimal mean response time (e.g., Shortest Remaining Processing Time First). To address these problems, we develop a new scheduler, Cluster Fair Queueing (CFQ), which preferentially offers resources to jobs that complete the earliest under a fair sharing policy. We show that CFQ is able to minimize the mean response time while at the same time ensuring jobs to finish within a constant time after their completion under fair sharing. Our Spark deployment on a 100-node EC2 cluster demonstrates that compared to the built-in fair scheduler, CFQ can decrease the mean response time by 40%, which speeds up more than 40% of jobs by over 75% on average.
Chen Chen 0067, Wei Wang 0030, Shengkai Zhang, Bo Li 0001
INFOCOM2
2017 Coflex: Navigating the fairness-efficiency tradeoff for coflow scheduling
abstract
Fair and efficient coflow scheduling improves application-level networking performance in today's datacenters. Ideally, a coflow scheduler should provide isolation guarantees on the minimum coflow progress to achieve predictable networking performance. Network operators, on the other hand, strive to decrease the average coflow completion time (CCT). Unfortunately, optimal isolation guarantees and minimum average CCT are conflicting objectives and cannot be achieved at the same time. Existing coflow schedulers either optimize isolation guarantees at the expense of long CCTs (e.g., HUG [1]), or decrease the average CCT without performance isolation (e.g., Varys and Aalo [2], [3]). The lack of a smooth tradeoff in between poses a dilemma between low efficiency and no performance isolation. To bridge this gap, we develop a new coflow scheduler, Coflex, to navigate this tradeoff. Coflex allows network operators to specify the desired level of isolation guarantee using a tunable fairness knob, while at the same time decreasing the average CCT. Both our real-world deployments and trace-driven simulations have shown that Coflex offers a smooth tradeoff between fairness and efficiency. At an appropriate tradeoff level, Coflex outperforms fair schedulers by 2 × in minimizing the average CCT.
Wei Wang 0030, Shiyao Ma, Bo Li 0001, Baochun Li
INFOCOM1
2017 LRC: Dependency-aware cache management for data analytics clusters
abstract
Memory caches are being aggressively used in today's data-parallel systems such as Spark, Tez, and Piccolo. However, prevalent systems employ rather simple cache management policies — notably the Least Recently Used (LRU) policy — that are oblivious to the application semantics of data dependency, expressed as a directed acyclic graph (DAG). Without this knowledge, memory caching can at best be performed by “guessing” the future data access patterns based on historical information (e.g., the access recency and/or frequency), which frequently results in inefficient, erroneous caching with low hit ratio and a long response time. In this paper, we propose a novel cache replacement policy, Least Reference Count (LRC), which exploits the application-specific DAG information to optimize the cache management. LRC evicts the cached data blocks whose reference count is the smallest. The reference count is defined, for each data block, as the number of dependent child blocks that have not been computed yet. We demonstrate the efficacy of LRC through both empirical analysis and cluster deployments against popular benchmarking workloads. Our Spark implementation shows that, compared with LRU, LRC speeds up typical applications by 60%.
Yinghao Yu, Wei Wang 0030, Jun Zhang 0004, Khaled Ben Letaief
INFOCOM2
2016 Friends or foes: Revisiting strategy-proofness in cloud network sharing
abstract
Cloud networks consist of a large number of links, on which tenants have correlated and elastic bandwidth demands in the form of coflows. Ideally, a cloud network sharing policy should provide tenants with isolation guarantees on the minimum coflow progress, while at the same time attaining as high utilization as possible. Prior work shows that to achieve the optimal isolation guarantee, strategy-proofness is needed, in that tenants cannot lie about demands to obtain higher progresses. However, this requirement is derived under a simplified assumption that tenants are only interested in maximizing coflow progresses. We show in this work that a rational tenant should pursue more bandwidth allocation as a secondary objective after progress maximization. In this new model, enforcing strategy-proofness inevitably hurts the isolation guarantee. We propose a new network sharing policy to achieve the optimal isolation guarantee while attaining the highest possible utilization in spite of strategic, untruthful tenants. Trace-driven evaluations show that our policy outperforms existing alternatives with better isolation guarantee, higher utilization, and shorter coflow completion time (CCT).
Wei Wang 0030, A-Long Jin
ICNP1
2016 Multi-resource fair sharing for datacenter jobs with placement constraints
abstract
Providing quality-of-service guarantees by means of fair sharing has never been more challenging in datacenters. Due to the heterogeneity of machine configurations, datacenter jobs frequently specify placement constraints, restricting them to run on a particular class of machines meeting specific hardware/software requirements. In addition, jobs have diverse demands across multiple resource types, and may saturate any of the CPU, memory, or storage resources. Despite the rich body of recent work on datacenter scheduling, it remains unclear how multi-resource fair sharing is defined and achieved for jobs with placement constraints. In this paper, we propose a new sharing policy called Task Share Fairness (TSF). With TSF, jobs are better off sharing the datacenter, and are better off reporting demands and constraints truthfully. We have prototyped TSF on Apache Mesos and confirmed its service guarantees in a 50-node EC2 cluster. Trace-driven simulations have further revealed that TSF speeds up 60% of tasks over existing fair schedulers.
Wei Wang 0030, Baochun Li, Ben Liang 0001, Jun Li 0017
SC1
2016 Towards Multi-Resource Fair Allocation with Placement Constraints
abstract
Multi-resource fair schedulers have been widely implemented in compute clusters to provide service isolation guarantees. Existing multi-resource sharing policies, notably Dominant Resource Fairness (DRF) and its variants, are designed for unconstrained jobs that can run on all machines in a cluster. However, an increasing number of datacenter jobs specify placement constraints and can only run on a particular class of machines meeting specific hardware/software requirements (e.g., GPUs or a particular kernel version). We show that directly extending existing policies to constrained jobs either compromises isolation guarantees or allows users to gain more resources by deceiving the scheduler. It remains unclear how multi-resource fair sharing is defined and achieved in the presence of placement constraints. We address this open problem by a new sharing policy, called Task Share Fairness (TSF), that provides provable isolation guarantees and is strategy-proof against gaming the allocation policy. TSF is shown to be envy-free and Pareto optimal as well.
Wei Wang 0030, Baochun Li, Ben Liang 0001, Jun Li 0017
SIGMETRICS1
2015 TiSA: Time-dependent social network advertising
abstract
Much of today's online social network (OSN) system relies on advertising for financial support. To improve the effectiveness of advertising, online advertisers tend to leverage influential users to deliver ads. Most of existing efforts on online advertising have focused on single-shot scenarios or assume static OSN models, while they overlook the fact that actions of advertising affect users' behaviors. In this paper, we investigate the behaviors of Sina Weibo users over three months, and make the observation that advertising affects the behaviors of the user's followers, which in turn has an impact on the effectiveness of future advertising. Based on this observation, we propose TiSA, a time-dependent advertising framework, which considers the future impact of advertising. Under this framework, the advertiser and the user make their decisions based on their instant utilities as well as future utilities. We also devise a learning algorithm with provable convergence to obtain the optimal policies. Evaluations using three month traces of 975 Sina Weibo users have been conducted, and the results validate the effectiveness of the proposed framework by showing that the utilities of all entities are significantly improved compared with traditional systems.
Wei Wang 0030, Qing Liao 0001, Qian Zhang 0001
ICC1
2015 Multi-Resource Fair Allocation in Heterogeneous Cloud Computing Systems
abstract
We study the multi-resource allocation problem in cloud computing systems where the resource pool is constructed from a large number of heterogeneous servers, representing different points in the configuration space of resources such as processing, memory, and storage. We design a multi-resource allocation mechanism, called DRFH, that generalizes the notion of Dominant Resource Fairness (DRF) from a single server to multiple heterogeneous servers. DRFH provides a number of highly desirable properties. With DRFH, no user prefers the allocation of another user; no one can improve its allocation without decreasing that of the others; and more importantly, no coalition behavior of misreporting resource demands can benefit all its members. DRFH also ensures some level of service isolation among the users. As a direct application, we design a simple heuristic that implements DRFH in real-world systems. Large-scale simulations driven by Google cluster traces show that DRFH significantly outperforms the traditional slot-based scheduler, leading to much higher resource utilization with substantially shorter job completion times.
Wei Wang 0030, Ben Liang 0001, Baochun Li
IEEE Trans. Parallel Distributed Syst.1
2015 Optimal Online Multi-Instance Acquisition in IaaS Clouds
abstract
Infrastructure-as-a-service (IaaS) clouds offer diverse instance purchasing options. A user can either run instances on demand and pay only for what it uses, or it can prepay to reserve instances for a long period, during which a usage discount is entitled. An important problem facing a user is how these two instance options can be dynamically combined to serve time-varying demands at minimum cost. Existing strategies in the literature, however, require either exact knowledge or the distribution of demands in the long-term future, which significantly limits their use in practice. Unlike existing works, we propose two practical online algorithms, one deterministic and another randomized, that dynamically combine the two instance options online without any knowledge of the future. We show that the proposed deterministic (resp., randomized) algorithm incurs no more than 2 - α (resp., e/(e-1 + α)) times the minimum cost obtained by an optimal offline algorithm that knows the exact future a priori, where a is the entitled discount after reservation. Our online algorithms achieve the best possible competitive ratios in both the deterministic and randomized cases, and can be easily extended to cases when short-term predictions are reliable. Simulations driven by a large volume of real-world traces show that significant cost savings can be achieved with prevalent IaaS prices.
Wei Wang 0030, Ben Liang 0001, Baochun Li
IEEE Trans. Parallel Distributed Syst.1
2015 Dynamic Cloud Instance Acquisition via IaaS Cloud Brokerage
abstract
Infrastructure-as-a-Service clouds offer diverse pricing options, including on-demand and reserved instances with various discounts to attract different cloud users. A practical problem facing cloud users is how to minimize their costs by choosing among different pricing options based on their own demands. In this paper, we propose a new cloud brokerage service that reserves a large pool of instances from cloud providers and serves users with price discounts. The broker optimally exploits both pricing benefits of longterm instance reservations and multiplexing gains. We propose dynamic strategies for the broker to make instance reservations with the objective of minimizing its service cost. These strategies leverage dynamic programming and approximation algorithms to rapidly handle large volumes of demand. Our extensive simulations driven by large-scale Google cluster-usage traces have shown that significant price discounts can be realized via the broker.
Wei Wang 0030, Di Niu 0002, Ben Liang 0001, Baochun Li
IEEE Trans. Parallel Distributed Syst.1
2014 On the Fairness-Efficiency Tradeoff for Packet Processing with Multiple Resources
abstract
Middleboxes are widely deployed in today's networks. They apply a variety of complex network functions to transform, filter, and optimize incoming traffic based on the payload of packets. These functions require the support of multiple types of resources, such as CPU and link bandwidth, for processing incoming packets. Hence, a multi-resource packet scheduling algorithm is needed to allow flows to share these resources fairly and efficiently. However, unlike traditional fair queueing where bandwidth is the only concern, we show in this paper that fairness and efficiency are conflicting objectives that cannot be achieved simultaneously in the presence of multiple resources. Ideally, a scheduling algorithm should allow network operators to flexibly specify their fairness and efficiency requirements, so as to meet the Quality of Service demands while keeping the system at a high utilization level. Yet, existing multi-resource scheduling algorithms focus on fairness only, and may lead to poor resource utilization. In this paper, we propose a new scheduling algorithm to achieve a flexible tradeoff between fairness and efficiency for packet processing, consuming both CPU and link bandwidth. Experimental results based on both real-world implementation and trace-driven simulation suggest that trading off a modest level of fairness can potentially improve the efficiency to the point where the system capacity is almost saturated.
Wei Wang 0030, Chen Feng 0001, Baochun Li, Ben Liang 0001
CoNEXT1
2014 Dominant resource fairness in cloud computing systems with heterogeneous servers
abstract
We study the multi-resource allocation problem in cloud computing systems where the resource pool is constructed from a large number of heterogeneous servers, representing different points in the configuration space of resources such as processing, memory, and storage. We design a multi-resource allocation mechanism, called DRFH, that generalizes the notion of Dominant Resource Fairness (DRF) from a single server to multiple heterogeneous servers. DRFH provides a number of highly desirable properties. With DRFH, no user prefers the allocation of another user; no one can improve its allocation without decreasing that of the others; and more importantly, no user has an incentive to lie about its resource demand. As a direct application, we design a simple heuristic that implements DRFH in real-world systems. Large-scale simulations driven by Google cluster traces show that DRFH significantly outperforms the traditional slot-based scheduler, leading to much higher resource utilization with substantially shorter job completion times.
Wei Wang 0030, Baochun Li, Ben Liang 0001
INFOCOM1
2014 Low complexity multi-resource fair queueing with bounded delay
abstract
Middleboxes are ubiquitous in today's networks. They perform deep packet processing such as content-based filtering and transformation, which requires multiple categories of resources (e.g., CPU, memory bandwidth, and link bandwidth). Depending on the processing requirement of traffic, packet processing for different flows may consume vastly different amounts of resources. Multi-resource fair queueing allows flows to obtain a fair share of these resources, providing service isolation across flows. However, previous solutions for multi-resource fair queueing are either expensive to implement at high speeds, or incurring high scheduling delay for flows with uneven weights. In this paper, we present a new fair queueing algorithm, called Group Multi-Resource Round Robin (GMR3), that schedules packets in O(1) time, while achieving near-perfect fairness with a low scheduling delay bounded by a small constant. To our knowledge, it is the first provably fair, highly efficient multi-resource fair queueing algorithm with bounded delay.
Wei Wang 0030, Ben Liang 0001, Baochun Li
INFOCOM1
2014 Designing Truthful Spectrum Double Auctions with Local Markets
abstract
Market-driven spectrum auctions offer an efficient way to improve spectrum utilization by transferring unused or underused spectrum from its primary license holder to spectrum-deficient secondary users. Such a spectrum market exhibits strong locality in two aspects: 1) that spectrum is a local resource and can only be traded to users within the license area, and 2) that holders can partition the entire license areas and sell any pieces in the market. We design a spectrum double auction that incorporates such locality in spectrum markets, while keeping the auction economically robust and computationally efficient. Our designs are tailored to cases with and without the knowledge of bid distributions. Complementary simulation studies show that spectrum utilization can be significantly improved when distribution information is available. Therefore, an auctioneer can start from one design without any a priori information, and then switch to the other alternative after accumulating sufficient distribution knowledge. With minor modifications, our designs are also effective for a profit-driven auctioneer aiming to maximize the auction revenue.
Wei Wang 0030, Ben Liang 0001, Baochun Li
IEEE Trans. Mob. Comput.1
2013 Dynamic Cloud Resource Reservation via Cloud Brokerage
abstract
Infrastructure-as-a-Service clouds offer diverse pricing options, including on-demand and reserved instances with various discounts to attract different cloud users. A practical problem facing cloud users is how to minimize their costs by choosing among different pricing options based on their own demands. In this paper, we propose a new cloud brokerage service that reserves a large pool of instances from cloud providers and serves users with price discounts. The broker optimally exploits both pricing benefits of long-term instance reservations and multiplexing gains. We propose dynamic strategies for the broker to make instance reservations with the objective of minimizing its service cost. These strategies leverage dynamic programming and approximate algorithms to rapidly handle large volumes of demand. Our extensive simulations driven by large-scale Google cluster-usage traces have shown that significant price discounts can be realized via the broker.
Wei Wang 0030, Di Niu 0002, Baochun Li, Ben Liang 0001
ICDCS1
2013 Multi-Resource Round Robin: A low complexity packet scheduler with Dominant Resource Fairness
abstract
Middleboxes are widely deployed in today's enterprise networks. They perform a wide range of important network functions, including WAN optimizations, intrusion detection systems, network and application level firewalls, etc. Depending on the processing requirement of traffic, packet processing for different traffic flows may consume vastly different amounts of hardware resources (e.g., CPU and link bandwidth). Multi-resource fair queueing allows each traffic flow to receive a fair share of multiple middlebox resources. Previous schemes for multi-resource fair queueing, however, are expensive to implement at high speeds. Specifically, the time complexity to schedule a packet is O(log n), where n is the number of backlogged flows. In this paper, we design a new multi-resource fair queueing scheme that schedules packets in a way similar to Elastic Round Robin. Our scheme requires only O(1) work to schedule a packet and is simple enough to implement in practice. We show, both analytically and experimentally, that our queueing scheme achieves nearly perfect Dominant Resource Fairness.
Wei Wang 0030, Baochun Li, Ben Liang 0001
ICNP1
2013 Revenue maximization with dynamic auctions in IaaS cloud markets
abstract
Cloud service pricing plays a pivotal role towards the success of cloud computing. Existing pricing schemes, however, either provide no service guarantees (e.g., Spot Instances in Amazon EC2), or use static on-demand pricing in which the price cannot respond quickly to market dynamics (e.g., On-demand Instances in Amazon EC2). To overcome these problems, in this paper we design dynamic auctions where computing instances are periodically auctioned off to accommodate user demands over time. We address the two main challenges of revenue maximization and auction truthfulness. Our design encompasses a capacity allocation scheme, which determines the amount of instances to be auctioned off in each period, as well as the underlying auction mechanisms, based on dynamic payment schemes corresponding to the proposed capacity allocations over time. We show that our design is two-dimensionally truthful, and it is asymptotically optimal when demand is sufficiently high. Furthermore, by identifying certain optimization structures, we substantially reduce the computational complexity of our solution. Extensive simulations show that our design closely tracks market changes, while generating higher revenues than on-demand pricing.
Wei Wang 0030, Ben Liang 0001, Baochun Li
IWQoS1
2013 Multi-resource generalized processor sharing for packet processing
abstract
Middleboxes have found widespread adoption in today's networks. They perform a variety of network functions such as WAN optimization, intrusion detection, and network-level firewalls. Processing packets to serve these functions often require multiple middlebox resources, e.g., CPU and link band-width. Furthermore, different packet traffic flows may consume significantly different amounts of various resources, depending on the network functions that are applied. Multi-resource fair queueing is therefore needed to allow flows to share multiple middlebox resources in a fair manner. In this paper, we clarify the fairness requirements of a queueing scheme and present Dominant Resource Generalized Processor Sharing (DRGPS), a fluid flow-based fair queueing idealization that strictly realizes Dominant Resource Fairness (DRF) at all times. As a form of Generalized Processor Sharing (GPS) running on multiple resources, DRGPS serves as a benchmark that practical packet-by-packet fair queueing algorithm should follow. With DRGPS, techniques and insights that have been developed for traditional fair queueing can be leveraged to schedule multiple resources. As a case study, we extend Worst-case Fair Weighted Fair Queueing (WF2Q) to the multi-resource setting and analyze its performance.
Wei Wang 0030, Ben Liang 0001, Baochun Li
IWQoS1
2012 Towards Optimal Capacity Segmentation with Hybrid Cloud Pricing
abstract
Cloud resources are usually priced in multiple markets with different service guarantees. For example, Amazon EC2 prices virtual instances under three pricing schemes -- the subscription option (a.k.a., Reserved Instances), the pay-as-you-go offer (a.k.a., On-Demand Instances), and an auction-like spot market (a.k.a., Spot Instances) -- simultaneously. There arises a new problem of capacity segmentation: how can a provider allocate resources to different categories of pricing schemes, so that the total revenue is maximized? In this paper, we consider an EC2-like pricing scheme with traditional pay-as-you-go pricing augmented by an auction market, where bidders periodically bid for resources and can use the instances for as long as they wish, until the clearing price exceeds their bids. We show that optimal periodic auctions must follow the design of m+1-price auction with seller's reservation price. Theoretical analysis also suggests the connections between periodic auctions and EC2 spot market. Furthermore, we formulate the optimal capacity segmentation strategy as a Markov decision process over some demand prediction window. To mitigate the high computational complexity of the conventional dynamic programming solution, we develop a near-optimal solution that has significantly lower complexity and is shown to asymptotically approach the optimal revenue.
Wei Wang 0030, Baochun Li, Ben Liang 0001
ICDCS1
2011 District: Embracing local markets in truthful spectrum double auctions
abstract
Market-driven spectrum auctions offer an efficient way to improve spectrum utilization by transferring unused or under-used spectrum from its primary license holder to spectrum-deficient secondary users. Such a spectrum market exhibits strong locality in two aspects: 1) that spectrum is a local resource and can only be traded to users within the license area, and 2) that holders can partition the entire license areas and sell any pieces in the market. We design a spectrum double auction that incorporates such locality in spectrum markets, while keeping the auction economically robust and computationally efficient. Our designs in District are tailored to cases with and without knowledge of bid distributions. An auctioneer can start from one design without any a priori information, and then switch to the other alternative after accumulating sufficient distribution knowledge. Complementary simulation studies show that spectrum utilization can be significantly improved when distribution information is available.
Wei Wang 0030, Baochun Li, Ben Liang 0001
SECON1
2009 Sequential Greedy Localization in Wireless Sensor Networks With Inaccurate Anchor Positions
abstract
In this paper, we consider the range-based sensor network localization with inaccurate anchor position information. First, a novel optimization algorithm named sequential greedy optimization (SGO) algorithm is proposed, and then two distributed localization algorithms are obtained: the first is obtained by applying the SGO algorithm to a convex formulation of the localization problem, named CSGLA; while the second is obtained by applying the SGO algorithm to a nonconvex formulation of the localization problem, named NCSGLA. The CSGLA must converge globally while the NCSGLA may converge locally. Both algorithms are partially asynchronous and can be implemented in a distributed fashion in networks. We demonstrate the localization performance via simulations. Simulation results show that, 1) the CSGLA algorithm works faster than the synchronous algorithm with the same localization accuracy; 2) with a reasonably good initialization, the NCSGLA can work much better than the CSGLA.
Qingjiang Shi, Chen He 0001, Hongyang Chen 0001, Ling-ge Jiang, Wei Wang 0030
GLOBECOM5
2009 A Noncooperative Spectrum Sensing Game with Maximum Network Throughput
abstract
In this paper, we consider a noncooperative cognitive radio network with M selfish secondary users (SUs) opportunistically access N licensed channels. Every SU chooses one channel to sense and subsequently compete to access (based on the sensing outcome) to obtain the channel utility. Different channels may have different utilities. Each SU selfishly makes a sensing decision to maximize its obtained utility. The objective is to design an optimal sensing policy with maximum network throughput. This problem is formulated as a noncooperative game where a stable sensing policy reaches a Nash equilibrium (NE). A novel greedy algorithm with great efficiency is proposed to calculate all pure-strategy NE for a large class of utility functions. By slight modification, the algorithm is able to reach an optimal pure-strategy NE with the maximum network throughput. The algorithm can be practically implemented as a MAC protocol in a distributed way with negligible communication overhead.
Wei Wang 0030
GLOBECOM1