Chuan Wu 0001

dblp:34/3772-1 · DBLP profile ↗
← Back
205ranked-venue papers
17as first author
65since 2021 · last 2026
0000-0002-3144-4398ORCID · verified

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

Computer networks · 115 · 12 first-author · 21 since 2021Systems, architecture and hardware · 51 · 3 first-author · 26 since 2021Graphics, computer vision, multimedia, augmented reality and games · 16 · 5 first-author · 4 since 2021Artificial intelligence and machine learning · 14 · 12 since 2021Applied, interdisciplinary, general and emerging computing · 9 · 1 since 2021Software engineering, systems software and programming languages · 8 · 2 since 2021Databases, data management, data science and information retrieval · 2 · 2 since 2021Theory of computation · 1
YearPublicationVenuePosition
2026 HetAuto: Cross-Cluster Auto-Parallelism for Heterogeneous Distributed Training
abstract
As large neural network models (e.g., LLMs) grow in scale, single-cluster resources become insufficient, making cross-cluster distributed training essential. Cross-cluster training is challenging: hardware heterogeneity complicates load balancing and parallelization strategy and introduces hardware compatibility issues in implementation; cross-cluster communication bottlenecks severely impact training throughput. We present HetAuto, an automatic parallelization system for efficient cross-cluster heterogeneous large model training. HetAuto contributes three key innovations: (1) a principle-guided MCTS algorithm with a random forest-enhanced cost model that efficiently searches parallelization strategies and quickly evaluates their performance under heterogeneous configurations; (2) cross-cluster communication optimizations including Virtual-1F1B scheduling that overlaps communication with computation and an optimized resharding strategy for inter-stage communication; and (3) a unified API enabling seamless integration of diverse accelerators. We evaluate HetAuto across 4 different clusters with up to 736 heterogeneous devices. The evaluation results show that HetAuto achieves up to 1.57× training throughput improvement over representative baselines, and strikes an efficient balance between solution quality and search overhead.
Guicheng Qi, Junwei Su, Liqi Yang, Tingwen Xie, Yerui Sun, Chuan Wu 0001
EuroSys8
2026 Laminar: A Scalable Asynchronous RL Post-Training Framework
abstract
Reinforcement learning (RL) post-training for Large Language Models (LLMs) is now scaling to large clusters and running for extended durations to enhance model reasoning performance. However, the scalability of existing RL frameworks is limited, as extreme long-tail skewness in RL trajectory generation causes severe GPU underutilization. Current asynchronous RL systems attempt to mitigate this, but they rely on global weight synchronization between the actor and all rollouts, which creates a rigid model update schedule. This global synchronization is ill-suited for the highly skewed and evolving distribution of trajectory generation latency in RL training, crippling training efficiency. Our key insight is that efficient scaling requires breaking this lockstep through trajectory-level asynchrony, which generates and consumes each trajectory independently. We propose Laminar, a scalable and robust RL post-training system built on a fully decoupled architecture. First, we replace global updates with a tier of relay workers acting as a distributed parameter service. This enables asynchronous and fine-grained weight synchronization, allowing rollouts to pull the latest weight anytime without stalling the actor's training loop. Second, a dynamic repack mechanism consolidates long-tail trajectories onto a few dedicated rollouts, maximizing generation throughput. The fully decoupled design also isolates failures, ensuring robustness for long-running jobs. Our evaluation on a 1024-GPU cluster shows that Laminar achieves up to 5.48$\times$ training throughput speedup over state-of-the-art systems, while reducing model convergence time.
Guangming Sheng, Yuxuan Tong, Borui Wan, Wang Zhang 0017, Chaobo Jia, Xibin Wu, Xiang Li 0067, Chi Zhang 0022, Yanghua Peng, Haibin Lin, Xin Liu 0086, Chuan Wu 0001
EuroSys13
2026 MegaScale-Data: Scaling DataLoader for Multisource Large Foundation Model Training
abstract
Modern frameworks for training large foundation models (LFMs) employ dataloaders in a data-parallel manner, with each loader processing a disjoint subset of training data. When preparing data for LFM training that originates from multiple, distinct sources, two fundamental challenges arise. First, due to the quadratic computational complexity of the attention operator, the non-uniform sample distribution over data-parallel ranks leads to significant workload imbalance among dataloaders, degrading the training efficiency. Second, supporting diverse data sources requires per-dataset file access states that are redundantly replicated across parallel loaders, consuming excessive memory. This also hinders dynamic data mixing (e.g., curriculum learning) and causes redundant access/memory overhead in hybrid parallelism.
Juntao Zhao 0002, Borui Wan, Lei Zuo 0004, Junda Feng, Jianyu Jiang, Yangrui Chen, Shuaishuai Cao, Jialing He, Kaihua Jiang, Shibiao Nong, Yanghua Peng, Haibin Lin, Chuan Wu 0001
EuroSys16
2026 BROS: Efficient LLM Serving on Hybrid Real-time and Best-effort Requests
Borui Wan, Juntao Zhao 0002, Chenyu Jiang 0002, Chuanxiong Guo, Chuan Wu 0001
INFOCOM5
2026 Joint Bitrate and Resource Adaptation for Super-Resolution Video Streaming in Multi-Cluster Edge Networks: A New Online Learning Approach
abstract
Today's video streaming service providers have exploited cloud-edge collaborative networks for video delivery across geo-distributed edge clusters and end users. The existing content delivery network (CDN) scheduling and adaptive bitrate algorithms may not fully utilize edge resources or lack a global control to optimize resource sharing. The emerging super-resolution (SR) approach can unleash the potential of leveraging computation resources to compensate for bandwidth consumption, by producing high-quality videos from low-resolution contents. Yet the uncertain SR resource sensitivity and its interplay with bitrate adaptation are under-explored. In this work, we proposeRosevin, the first resource scheduler that jointly decides the bitrates and fine-grained resource allocation to perform SR at the edge, which can learn to optimize the long-term QoE for distributed end users. To handle the time-varying and complex space of decisions as well as a non-smooth objective function,Rosevinrealizes a novel online combinatorial learning algorithm, which nicely integrates convex optimization theories and online learning techniques, addressing the switching cost issues. In addition to theoretically analyzing its performance, we implement an SR-assisted video streaming prototype ofRosevinand demonstrate its advantages over several video delivery benchmarks.
Xiaoxi Zhang 0001, Longhao Zou, Jingpu Duan, Chuan Wu 0001, Yali Xue, Zuozhou Chen, Chaoqi Zhou, Xu Chen 0004
IEEE Trans. Mob. Comput.5
2025 DyOrc: Efficient Serving of Dynamic Machine Learning Workflows
abstract
The landscape of machine learning applications has shifted from monolithic end-to-end models to compositions of pretrained large foundation models. For instance, multi-modal chatbots are often built by composition of a large language model and modality-specific encoder models. Such applications often feature dynamic workflows, with models conditionally evoked according to different inputs and intermediate processing results. Conditional model execution prevents conventional request batching and hinders efficient hardware utilization, due to dynamic, diverging execution paths across requests. Separately deploying models as dedicated services and invoking them on the go during dynamic workflow executions can potentially allow service-wise request batching, boosting resource efficiency. However, generic workflow orchestrators are proven inefficient for machine learning applications, due to schedulers that do not exploit batching, communication methods that are suboptimal for GPU tensors, and the considerable cold-start delays associated with model deployment.
Shiwei Zhang 0002, Lansong Diao, Zisheng Meng, Siyu Wang 0006, Wei Lin 0016, Chuan Wu 0001
SoCC6
2025 SplitQuant: Resource-Efficient LLM Offline Serving on Heterogeneous GPUs via Phase-Aware Model Partition and Adaptive Quantization
abstract
Modern large language models (LLMs) serving systems address distributed deployment challenges through two key techniques: distributed model partitioning for parallel computation across accelerators and quantization for reducing parameter size. While existing systems assume homogeneous GPU environments, we reveal significant untapped potential in heterogeneous systems with mixed-capacity accelerators where two critical limitations persist: (1) uniform partitioning and quantization strategies fail to adapt to hardware heterogeneity, exacerbating resource imbalance, and (2) decoupled optimization of partitioning and quantization overlooks critical performance synergies between these techniques. We present SplitQuant, a phase-aware distributed serving system that co-optimizes mixedprecision quantization, phase-aware model partitioning, and micro-batch sizing for heterogeneous environments. Our approach combines analytical modeling of quality-runtime tradeoffs with a lightweight planning algorithm to maximize throughput while preserving user-specified model quality targets. Evaluations across 10 production clusters show SplitQuant achieves up to$2.34 \times(1.61 \times$mean) higher throughput than state-of-theart approaches without violating accuracy targets. Our results underscore the value of co-designing quantization and model partitioning strategies for heterogeneous environments.
Juntao Zhao 0002, Borui Wan, Yanghua Peng, Haibin Lin, Chuan Wu 0001
CLUSTER5
2025 HybridFlow: A Flexible and Efficient RLHF Framework
abstract
Reinforcement Learning from Human Feedback (RLHF) is widely used in Large Language Model (LLM) alignment. Traditional RL can be modeled as a dataflow, where each node represents computation of a neural network (NN) and each edge denotes data dependencies between the NNs. RLHF complicates the dataflow by expanding each node into a distributed LLM training or generation program, and each edge into a many-to-many multicast. Traditional RL frameworks execute the dataflow using a single controller to instruct both intra-node computation and inter-node communication, which can be inefficient in RLHF due to large control dispatch overhead for distributed intra-node computation. Existing RLHF systems adopt a multi-controller paradigm, which can be inflexible due to nesting distributed computation and data communication. We propose HybridFlow, which combines single-controller and multi-controller paradigms in a hybrid manner to enable flexible representation and efficient execution of the RLHF data flow. We carefully design a set of hierarchical APIs that decouple and encapsulate computation and data dependencies in the complex RLHF dataflow, allowing efficient operation orchestration to implement RLHF algorithms and flexible mapping of the computation onto various devices. We further design a 3D-HybridEngine for efficient actor model resharding between training and generation phases, with zero memory redundancy and significantly reduced communication overhead. Our experimental results demonstrate 1.53x~20.57× throughput improvement when running various RLHF algorithms using HybridFlow, as compared with state-of-the-art baselines. HybridFlow source code is available at https://github.com/volcengine/verl
Guangming Sheng, Chi Zhang 0022, Zilingfeng Ye, Xibin Wu, Wang Zhang 0017, Ru Zhang 0006, Yanghua Peng, Haibin Lin, Chuan Wu 0001
EuroSys9
2025 ECCheck: Enhancing In-Memory Checkpoint with Erasure Coding in Distributed DNN Training
abstract
Distributed large model training is intensively time and resource consuming. Failures during the long training period are often inevitable, and can incur substantial recovery costs. Checkpointing has been the standard fault tolerance approach, which periodically stores the latest model states at remote persistent storage. This process can be time-consuming due to limited network bandwidth, and adversely affects training throughput. In-memory checkpointing addresses this issue by saving checkpoint data into host memory instead of remote storage. However, host memory is non-persistent, and may not provide sufficient resilience in case of machine failure. We propose ECCheck, a novel in-memory checkpoint system that employs erasure coding to enhance fault tolerance in distributed deep neural network training. ECCheck advocates serialization-free encoding and decoding in model checkpointing. Several techniques are proposed to minimize computation and communication overhead incurred by erasure coding. Extensive experiments demonstrate that ECCheck achieves superior fault tolerance compared to state-of-the-art solutions, while maintaining high checkpointing frequency, low checkpointing stalls, and fast recovery from failures.
Guicheng Qi, Zongpeng Li, Chuan Wu 0001, Zhuwei Peng, Yi Zheng 0007
ICDCS3
2025 On the Interplay between Graph Structure and Learning Algorithms in Graph Neural Networks
abstract
This paper studies the interplay between learning algorithms and graph structure for graph neural networks (GNNs). Existing theoretical studies on the learning dynamics of GNNs primarily focus on the convergence rates of learning algorithms under the interpolation regime (noise-free) and offer only a crude connection between these dynamics and the actual graph structure (e.g., maximum degree). This paper aims to bridge this gap by investigating the excessive risk (generalization performance) of learning algorithms in GNNs within the generalization regime (with noise). Specifically, we extend the conventional settings from the learning theory literature to the context of GNNs and examine how graph structure influences the performance of learning algorithms such as stochastic gradient descent (SGD) and Ridge regression. Our study makes several key contributions toward understanding the interplay between graph structure and learning in GNNs. First, we derive the excess risk profiles of SGD and Ridge regression in GNNs and connect these profiles to the graph structure through spectral graph theory. With this established framework, we further explore how different graph structures (regular vs. power-law) impact the performance of these algorithms through comparative analysis. Additionally, we extend our analysis to multi-layer linear GNNs, revealing an increasing non-isotropic effect on the excess risk profile, thereby offering new insights into the over-smoothing issue in GNNs from the perspective of learning algorithms. Our empirical results align with our theoretical predictions, collectively showcasing a coupling relation among graph structure, GNNs and learning algorithms, and providing insights on GNN algorithm design and selection in practice.
Junwei Su, Chuan Wu 0001
ICML2
2025 A Non-Asymptotic Convergent Analysis for Scored-Based Graph Generative Model via a System of Stochastic Differential Equations
abstract
This paper investigates the convergence behavior of score-based graph generative models (SGGMs). Unlike common score-based generative models (SGMs) that are governed by a single stochastic differential equation (SDE), SGGMs utilize a system of dependent SDEs, where the graph structure and node features are modeled separately, while accounting for their inherent dependencies. This distinction makes existing convergence analyses from SGMs inapplicable for SGGMs. In this work, we present the first convergence analysis for SGGMs, focusing on the convergence bound (the risk of generative error) across three key graph generation paradigms: (1) feature generation with a fixed graph structure, (2) graph structure generation with fixed node features, and (3) joint generation of both graph structure and node features. Our analysis reveals several unique factors specific to SGGMs (e.g., the topological properties of the graph structure) which significantly affect the convergence bound. Additionally, we offer theoretical insights into the selection of hyperparameters (e.g., sampling steps and diffusion length) and advocate for techniques like normalization to improve convergence. To validate our theoretical findings, we conduct a controlled empirical study using a synthetic graph model. The results in this paper contribute to a deeper theoretical understanding of SGGMs and offer practical guidance for designing more efficient and effective SGGMs.
Junwei Su, Chuan Wu 0001
ICML2
2025 PASTA: Training Acceleration for Vertical Federated Learning via Adaptive Pipeline Parallelism
abstract
Vertical federated learning (VFL) enables collaborative model training among geo-distributed participants, each with different features of the same samples, but only one party possesses the labels. Communication delays between active and passive parties in VFL significantly hinder its training efficiency. Existing VFL methods adopt asynchronous schemes or multiple local updates per communication round, but they either introduce heavy computation overhead or fail to adapt to dynamic network conditions. This work proposes PASTA, a novel framework employing Adaptive Pipeline Parallelism with Staleness Control for VFL, designed to mitigate these delays and balance training efficiency and model performance. PASTA enables concurrent communication and computation, maximizing resource utilization and minimizing idle time by strategically using stale gradients. Each passive party can send one or more batches of embeddings per communication and conduct stale local training, so that computation times can overlap with communication latency. Since staleness impedes model accuracy despite its benefits in reducing time, a dynamic feedback-based mechanism is proposed to adjust the numbers of embeddings sent and local training iterations based on system heterogeneity. Extensive experiments across various datasets demonstrate that PASTA significantly enhances convergence speed by$1.8 \times$to$4.6 \times$compared to leading VFL systems, without compromising final accuracy. The source code is available at https://github.com/PointerA/PASTA.
Ziwei Zhan, Jingpu Duan, Chuan Wu 0001, Jinhang Zuo, Xu Chen 0004, Xiaoxi Zhang 0001
IWQoS6
2025 SAILS: A Synchronous Accessible Immersive Online Learning System for Young Learners
abstract
Online distance learning emerged as a prominent means of education during the pandemic and is expected to continue as a long-lasting trend, enabling access to top education resources without the constraints of physical distance. However, conferencing software predominantly used for synchronous online teaching and learning is not suitable or sufficient for young students and teachers. The learning experiences of these students heavily rely on hands-on activities and interactions with peers and teachers that are challenging to replicate. Additionally, teachers face difficulties in effectively monitoring student progress. To address these challenges, we propose a Synchronous Accessible Immersive Online Learning System (SAILS). It offers an immersive and interactive platform for young learners by integrating real-life school activities with a virtual learning environment. Using the system, teachers can easily organize classes, assess students' work. We evaluated the effectiveness of SAILS through a user study involving 40 young learners from a kindergarten and several primary schools, along with their teachers. The results demonstrated a strong preference for our system among the participants, highlighting its ability to provide a more immersive and engaging online learning experience.
Yuran Sun, Zhuoying Zhang, Zhenxiao Luo, Alan William Dougherty, Man Ho Yip, Yi-King Choi, Chuan Wu 0001
MMSys7
2025 ByteCheckpoint: A Unified Checkpointing System for Large Foundation Model Development
Borui Wan, Mingji Han, Yiyao Sheng, Yanghua Peng, Haibin Lin, Mofan Zhang, Zhichao Lai, Menghan Yu, Junda Zhang, Zuquan Song, Xin Liu 0086, Chuan Wu 0001
NSDI12
2025 DCP: Addressing Input Dynamism In Long-Context Training via Dynamic Context Parallelism
abstract
Context parallelism has emerged as a key technique to support long-context training, a growing trend in generative AI for modern large models. However, existing context parallel methods rely on static parallelization configurations that overlook the dynamic nature of training data, specifically, the variability in sequence lengths and token relationships (i.e., attention patterns) across samples. As a result, these methods often suffer from unnecessary communication overhead and imbalanced computation. In this paper, we present DCP, a dynamic context parallel training framework that introduces fine-grained blockwise partitioning of both data and computation. By enabling flexible mapping of data and computation blocks to devices, DCP can adapt to varying sequence characteristics, effectively reducing communication and improving memory and computation balance. Micro-benchmarks demonstrate that DCP accelerates attention by 1.19x~2.45x under causal masks and 2.15x~3.77x under sparse attention patterns. Additionally, we observe up to 0.94x~1.16x end-to-end training speed-up for causal masks, and 1.00x~1.46x for sparse masks.
Chenyu Jiang 0002, Zhenkun Cai, Zhen Jia 0001, Yida Wang 0003, Chuan Wu 0001
SOSP6
2025 Robust LLM Training Infrastructure at ByteDance
abstract
The training scale of large language models (LLMs) has reached tens of thousands of GPUs and is still continuously expanding, enabling faster learning of larger models. Accompanying the expansion of the resource scale is the prevalence of failures (CUDA error, NaN values, job hang, etc.), which poses significant challenges to training stability. Any large-scale LLM training infrastructure should strive for minimal training interruption, efficient fault diagnosis, and effective failure tolerance to enable highly efficient continuous training. This paper presents ByteRobust, a large-scale GPU infrastructure management system tailored for robust and stable training of LLMs. It exploits the uniqueness of LLM training process and gives top priorities to detecting and recovering failures in a routine manner. Leveraging parallelisms and characteristics of LLM training, ByteRobust enables high-capacity fault tolerance, prompt fault demarcation, and localization with an effective data-driven approach, comprehensively ensuring continuous and efficient training of LLM tasks. ByteRobust is deployed on a production GPU platform with over 200,000 GPUs and advances the state of the art in training robustness by achieving 97% ETTR for a three-month training job on 9,600 GPUs.
Borui Wan, Gaohong Liu, Zuquan Song, Jun Wang 0039, Guangming Sheng, Shuguang Wang, Houmin Wei, Weiqiang Lou, Mofan Zhang, Kaihua Jiang, Cheng Ren, Xiaoyun Zhi, Menghan Yu, Zhe Nan, Zhuolin Zheng, Baoquan Zhong, Qinlong Wang, Jinxin Chi, Wang Zhang 0017, Zixian Du, Sida Zhao, Jingzhe Tang, Zherui Liu, Chuan Wu 0001, Yanghua Peng, Haibin Lin, Wencong Xiao, Xin Liu 0086
SOSP30
2025 Heta: Distributed Training of Heterogeneous Graph Neural Networks
abstract
Heterogeneous Graphs (HetGs) that capture relationships among different types of nodes are ubiquitous in real-world applications such as academic networks and e-commerce. Although Heterogeneous Graph Neural Networks (HGNNs) have demonstrated superior performance in learning from these complex structures, distributed training of HGNNs on large-scale graphs with billions of edges faces substantial communication overhead. This challenge is exacerbated by heterogeneous characteristics such as varying feature dimensions across node types and featureless nodes requiring learnable parameters. Existing systems and communication reduction techniques designed for homogeneous graphs become suboptimal or even inapplicable for HetGs and HGNNs by overlooking both these heterogeneous characteristics and the inherent computational structure of HGNNs. We present Heta , a framework designed to address the communication bottleneck in distributed HGNN training. Heta leverages the key insight that HGNN aggregation is order-invariant and decomposable into relation-specific computations. Built on this insight, we introduce three key innovations: (1) a Relation-Aggregation-First (RAF) paradigm that conducts relation-specific aggregations within partitions and exchanges only partial aggregations across machines, proven to reduce communication complexity; (2) a meta-partitioning strategy that divides a HetG based on its graph schema and HGNN computation dependency while minimizing cross-partition communication and maintaining computation and storage balance; and (3) a heterogeneity-aware GPU cache system that accounts for varying miss-penalty ratios across node types. Through extensive evaluation of billion-edge heterogeneous graphs, we demonstrate that Heta achieves up to 5.3X and 4.4X speedup over state-of-the-art systems DGL and GraphLearn while maintaining model accuracy.
Yuchen Zhong, Junwei Su, Chuan Wu 0001
Proc. VLDB Endow.3
2025 TransXNet: Learning Both Global and Local Dynamics With a Dual Dynamic Token Mixer for Visual Recognition
abstract
Recent studies have integrated convolutions into transformers to introduce inductive bias and improve generalization performance. However, the static nature of conventional convolution prevents it from dynamically adapting to input variations, resulting in a representation discrepancy between convolution and self-attention as self-attention calculates attention matrices dynamically. Furthermore, when stacking token mixers that consist of convolution and self-attention to form a deep network, the static nature of convolution hinders the fusion of features previously generated by self-attention into convolution kernels. These two limitations result in a suboptimal representation capacity of the constructed networks. To find a solution, we propose a lightweight dual dynamic token mixer (D-Mixer) to simultaneously learn global and local dynamics, that is, mechanisms that compute weights for aggregating global contexts and local details in an input-dependent manner. D-Mixer works by applying an efficient global attention module and an input-dependent depthwise convolution separately on evenly split feature segments, endowing the network with strong inductive bias and an enlarged effective receptive field. We use D-Mixer as the basic building block to design TransXNet, a novel hybrid CNN-transformer vision backbone network that delivers compelling performance. In the ImageNet-1K image classification task, TransXNet-T surpasses Swin-T by 0.3% in top-1 accuracy while requiring less than half of the computational cost. Furthermore, TransXNet-S and TransXNet-B exhibit excellent model scalability, achieving top-1 accuracy of 83.8% and 84.6%, respectively, with reasonable computational costs. In addition, our proposed network architecture demonstrates strong generalization capabilities in various dense prediction tasks, outperforming other state-of-the-art networks while having lower computational costs. Code is publicly available at https://github.com/LMMMEng/TransXNet.
Meng Lou, Shu Zhang 0001, Sibei Yang, Chuan Wu 0001, Yizhou Yu
IEEE Trans. Neural Networks Learn. Syst.5
2024 FaPES: Enabling Efficient Elastic Scaling for Serverless Machine Learning Platforms
abstract
Serverless computing platforms have become increasingly popular for running machine learning (ML) tasks due to their user-friendliness and decoupling from underlying infrastructure. However, auto-scaling to efficiently serve incoming requests still remains a challenge, especially for distributed ML training or inference jobs in a serverless GPU cluster. Distributed training and inference jobs are highly sensitive to resource configurations, and demand high model efficiency throughout their lifecycle. We propose FaPES, a FaaS-oriented Performance-aware Elastic Scaling system to enable efficient resource allocation in serverless platforms for ML jobs. FaPES enables flexible resource loaning between virtual clusters for running training and inference jobs. For running inference jobs, servers are reclaimed on demand with minimal preemption overhead to guarantee service level objective (SLO); for training jobs, optimal GPU allocation and model hyperparameters are jointly adapted based on an ML-based performance model and a resource usage prediction board, alleviating users from model tuning and resource specification. Evaluation on a 128-GPU testbed demonstrates up to 24.8% job completion time reduction and ×1.8 Goodput improvement, as compared to representative elastic scaling schemes.
Xiao-Yang Zhao 0005, Siran Yang, Jiamang Wang, Lansong Diao, Lin Qu, Chuan Wu 0001
SoCC6
2024 On the Topology Awareness and Generalization Performance of Graph Neural Networks
Junwei Su, Chuan Wu 0001
ECCV (84)2
2024 CDMPP: A Device-Model Agnostic Framework for Latency Prediction of Tensor Programs
abstract
Deep Neural Networks (DNNs) have shown excellent performance in a wide range of machine learning applications. Knowing the latency of running a DNN model or tensor program on a specific device is useful in various tasks, such as DNN graph- or tensor-level optimization and device selection. Considering the large space of DNN models and devices that impedes direct profiling of all combinations, recent efforts focus on building a predictor to model the performance of DNN models on different devices. However, none of the existing attempts have achieved a cost model that can accurately predict the performance of various tensor programs while supporting both training and inference accelerators. We propose CDMPP, an efficient tensor program latency prediction framework for both cross-model and cross-device prediction. We design an informative but efficient representation of tensor programs, called compact ASTs, and a pre-order-based positional encoding method, to capture the internal structure of tensor programs. We develop a domain-adaption-inspired method to learn domain-invariant representations and devise a KMeans-based sampling algorithm, for the predictor to learn from different domains (i.e., different DNN operators and devices). Our extensive experiments on a diverse range of DNN models and devices demonstrate that CDMPP significantly outperforms state-of-the-art baselines with 14.03% and 10.85% prediction error for cross-model and cross-device prediction, respectively, and one order of magnitude higher training efficiency. The implementation and the expanded dataset are available at https://github.com/joapolarbear/cdmpp.
Hanpeng Hu, Junwei Su, Juntao Zhao 0002, Yanghua Peng, Yibo Zhu 0001, Haibin Lin, Chuan Wu 0001
EuroSys7
2024 DynaPipe: Optimizing Multi-task Training through Dynamic Pipelines
abstract
Multi-task model training has been adopted to enable a single deep neural network model (often a large language model) to handle multiple tasks (e.g., question answering and text summarization). Multi-task training commonly receives input sequences of highly different lengths due to the diverse contexts of different tasks. Padding (to the same sequence length) or packing (short examples into long sequences of the same length) is usually adopted to prepare input samples for model training, which is nonetheless not space or computation efficient. This paper proposes a dynamic micro-batching approach to tackle sequence length variation and enable efficient multi-task model training. We advocate pipelineparallel training of the large model with variable-length micro-batches, each of which potentially comprises a different number of samples. We optimize micro-batch construction using a dynamic programming-based approach, and handle micro-batch execution time variation through dynamic pipeline and communication scheduling, enabling highly efficient pipeline training. Extensive evaluation on the FLANv2 dataset demonstrates up to 4.39x higher training throughput when training T5, and 3.25x when training GPT, as compared with packing-based baselines. DynaPipe's source code is publicly available at https://github.com/awslabs/optimizing-multitask-training-through-dynamic-pipelines.
Chenyu Jiang 0002, Zhen Jia 0001, Shuai Zheng 0004, Yida Wang 0003, Chuan Wu 0001
EuroSys5
2024 HAP: SPMD DNN Training on Heterogeneous GPU Clusters with Automated Program Synthesis
abstract
Single-Program-Multiple-Data (SPMD) parallelism has recently been adopted to train large deep neural networks (DNNs). Few studies have explored its applicability on heterogeneous clusters, to fully exploit available resources for large model learning. This paper presents HAP, an automated system designed to expedite SPMD DNN training on heterogeneous clusters. HAP jointly optimizes the tensor sharding strategy, sharding ratios across heterogeneous devices and the communication methods for tensor exchanges for optimized distributed training with SPMD parallelism. We novelly formulate model partitioning as a program synthesis problem, in which we generate a distributed program from scratch on a distributed instruction set that semantically resembles the program designed for a single device, and systematically explore the solution space with an A-based search algorithm. We derive the optimal tensor sharding ratios by formulating it as a linear programming problem. Additionally, HAP explores tensor communication optimization in a heterogeneous cluster and integrates it as part of the program synthesis process, for automatically choosing optimal collective communication primitives and applying sufficient factor broadcasting technique. Extensive experiments on representative workloads demonstrate that HAP achieves up to 2.41x speed-up on heterogeneous clusters.
Shiwei Zhang 0002, Lansong Diao, Chuan Wu 0001, Zongyan Cao, Siyu Wang 0006, Wei Lin 0016
EuroSys3
2024 AdapCC: Making Collective Communication in Distributed Machine Learning Adaptive
abstract
As deep learning (DL) models continue to grow in size, there is a pressing need for distributed model learning using a large number of devices (e.g., G PU s) and servers. Collective communication among devices/servers (for gradient synchronization, intermediate data exchange, etc.) introduces significant overheads, rendering major performance bottlenecks in distributed learning. A number of communication libraries, such as NCCL, Gloo and MPI, have been developed to optimize collective communication. Predefined communication strategies (e.g., ring or tree) are largely adopted, which may not be efficient or adaptive enough for inter-machine communication, especially in cloud-based scenarios where instance configurations and network performance can vary substantially. We propose AdapCC, a novel communication library that dynamically adapts to resource heterogeneity and network variability for optimized communication and training performance. AdapCC generates communication strategies based on run-time profiling, mitigates resource waste in waiting for computation stragglers, and executes efficient data transfers among DL workers. Experimental results under various settings demonstrate 2x communication speed-up and 31 % training throughput improvement with AdapCC, as compared to NCCL and other representative communication backends.
Xiao-Yang Zhao 0005, Chuan Wu 0001
ICDCS3
2024 PRES: Toward Scalable Memory-Based Dynamic Graph Neural Networks
abstract
Memory-based Dynamic Graph Neural Networks (MDGNNs) are a family of dynamic graph neural networks that leverage a memory module to extract, distill, and memorize long-term temporal dependencies, leading to superior performance compared to memory-less counterparts. However, training MDGNNs faces the challenge of handling entangled temporal and structural dependencies, requiring sequential and chronological processing of data sequences to capture accurate temporal patterns. During the batch training, the temporal data points within the same batch will be processed in parallel, while their temporal dependencies are neglected. This issue is referred to as temporal discontinuity and restricts the effective temporal batch size, limiting data parallelism and reducing MDGNNs' flexibility in industrial applications. This paper studies the efficient training of MDGNNs at scale, focusing on the temporal discontinuity in training MDGNNs with large temporal batch sizes. We first conduct a theoretical study on the impact of temporal batch size on the convergence of MDGNN training. Based on the analysis, we propose PRES, an iterative prediction-correction scheme combined with a memory coherence learning objective to mitigate the effect of temporal discontinuity, enabling MDGNNs to be trained with significantly larger temporal batches without sacrificing generalization performance. Experimental results demonstrate that our approach enables up to a 4 $\times$ larger temporal batch (3.4$\times$ speed-up) during MDGNN training.
Junwei Su, Difan Zou, Chuan Wu 0001
ICLR3
2024 Expediting Distributed GNN Training with Feature-only Partition and Optimized Communication Planning
abstract
Feature-only partition of large graph data in distributed Graph Neural Network (GNN) training offers advantages over commonly adopted graph structure partition, such as minimal graph preprocessing cost and elimination of cross-worker subgraph sampling burdens. Nonetheless, performance bottleneck of GNN training with feature-only partitions still largely lies in the substantial communication overhead due to cross-worker feature fetching. To reduce the communication overhead and expedite distributed training, we first investigate and answer two key questions on convergence behaviors of GNN model in feature-partition based distribute GNN training: 1) As no worker holds a complete copy of each feature, can gradient exchange among workers compensate for the information loss due to incomplete local features? 2) If the answer to the first question is negative, is feature fetching in every training iteration of the GNN model necessary to ensure model convergence? Based on our theoretical findings on these questions, we derive an optimal communication plan that decides the frequency for feature fetching during the training process, taking into account bandwidth levels among workers and striking a balance between model loss and training time. Extensive evaluation demonstrates consistent results with our theoretical analysis, and the effectiveness of our proposed design.
Bingqian Du, Jun Liu 0002, Ziyue Luo, Chuan Wu 0001, Qiankun Zhang 0001, Hai Jin 0001
INFOCOM4
2024 Rosevin: Employing Resource- and Rate-Adaptive Edge Super-Resolution for Video Streaming
abstract
Today’s video streaming service providers have exploited cloud-edge collaborative networks for geo-distributed video delivery. The existing content delivery network (CDN) scheduling and adaptive bitrate algorithms may not fully utilize edge resources or lack a global control to optimize resource sharing. The emerging super-resolution (SR) approach can unleash the potential of leveraging computation resources to compensate for bandwidth consumption, by producing high-quality videos from low-resolution contents. Yet the uncertain SR resource sensitivity and its interplay with bitrate adaptation are underexplored. In this work, we propose Rosevin, the first resource scheduler that jointly decides the bitrates and fine-grained resource allocation to perform SR at the edge, which can learn to optimize the long-term QoE for distributed end users. To handle the time-varying and complex space of decisions as well as a non-smooth objective function, Rosevin realizes a novel online combinatorial learning algorithm, which nicely integrates convex optimization theories and online learning techniques. In addition to theoretically analyzing its performance, we implement an SR-assisted video streaming prototype of Rosevin and demonstrate its advantages over several video delivery benchmarks.
Xiaoxi Zhang 0001, Longhao Zou, Jingpu Duan, Chuan Wu 0001, Yali Xue, Zuozhou Chen, Xu Chen 0004
INFOCOM5
2024 QSync: Quantization-Minimized Synchronous Distributed Training Across Hybrid Devices
abstract
A number of production deep learning clusters have attempted to explore inference hardware for DNN training, at the off-peak serving hours with many inference GPUs idling. Conducting DNN training with a combination of heterogeneous training and inference GPUs, known as hybrid device training, presents considerable challenges due to disparities in compute capability and significant differences in memory capacity. We propose QSync, a training system that enables efficient synchronous data-parallel DNN training over hybrid devices by strategically exploiting quantized operators. According to each device’s available resource capacity, QSync selects a quantization-minimized setting for operators in the distributed DNN training graph, minimizing model accuracy degradation but keeping the training efficiency brought by quantization. We carefully design a predictor with a bi-directional mixed-precision indicator to reflect the sensitivity of DNN layers on fixed-point and floating-point low-precision operators, a replayer with a neighborhood-aware cost mapper to accurately estimate the latency of distributed hybrid mixed-precision training, and then an allocator that efficiently synchronizes workers with minimized model accuracy degradation. QSync bridges the computational graph on PyTorch to an optimized backend for quantization kernel performance and flexible support for various GPU architectures. Extensive experiments show that QSync’s predictor can accurately simulate distributed mixed-precision training with < 5% error, with a consistent 0.27 − 1.03% accuracy improvement over the from-scratch training tasks compared to uniform precision.
Juntao Zhao 0002, Borui Wan, Yanghua Peng, Haibin Lin, Yibo Zhu 0001, Chuan Wu 0001
IPDPS6
2024 MSPipe: Efficient Temporal GNN Training via Staleness-Aware Pipeline
abstract
Memory-based Temporal Graph Neural Networks (MTGNNs) are a class of temporal graph neural networks that utilize a node memory module to capture and retain long-term temporal dependencies, leading to superior performance compared to memory-less counterparts. However, the iterative reading and updating process of the memory module in MTGNNs to obtain up-to-date information needs to follow the temporal dependencies. This introduces significant overhead and limits training throughput. Existing optimizations for static GNNs are not directly applicable to MTGNNs due to differences in training paradigm, model architecture, and the absence of a memory module. Moreover, these optimizations do not effectively address the challenges posed by temporal dependencies, making them ineffective for MTGNN training. In this paper, we propose MSPipe, a general and efficient framework for memory-based TGNNs that maximizes training throughput while maintaining model accuracy. Our design specifically addresses the unique challenges associated with fetching and updating node memory states in MTGNNs by integrating staleness into the memory module. However, simply introducing a predefined staleness bound in the memory module to break temporal dependencies may lead to suboptimal performance and lack of generalizability across different models and datasets. To overcome this, we introduce an online pipeline scheduling algorithm in MSPipe that strategically breaks temporal dependencies with minimal staleness and delays memory fetching to obtain fresher memory states. This is achieved without stalling the MTGNN training stage or causing resource contention. Additionally, we design a staleness mitigation mechanism to enhance training convergence and model accuracy. Furthermore, we provide convergence analysis and demonstrate that MSPipe maintains the same convergence rate as vanilla sampling-based GNN training. Experimental results show that MSPipe achieves up to 2.45× speed-up without sacrificing accuracy, making it a promising solution for efficient MTGNN training. The implementation of our paper can be found at the following link: https://github.com/PeterSH6/MSPipe.
Guangming Sheng, Junwei Su, Chao Huang 0001, Chuan Wu 0001
KDD4
2024 FedMoE-DA: Federated Mixture of Experts via Domain Aware Fine-Grained Aggregation
abstract
Federated learning (FL) is a collaborative machine learning approach that enables multiple clients to train models without sharing their private data. With the rise of deep learning, large-scale models have garnered significant attention due to their exceptional performance. However, a key challenge in FL is the limitation imposed by clients with constrained computational and communication resources, which hampers the deployment of these large models. The Mixture of Experts (MoE) architecture addresses this challenge with its sparse activation property, which reduces computational workload and communication demands during inference and updates. Additionally, MoE facilitates better personalization by allowing each expert to specialize in different subsets of the data distribution. To alleviate the communication burdens between the server and clients, we propose FedMoE-DA, a new FL model training framework that leverages the MoE architecture and incorporates a novel domain-aware, fine-grained aggregation strategy to enhance the robustness, personalizability, and communication efficiency simultaneously. Specifically, the correlation between both intra-client expert models and inter-client data heterogeneity is exploited. Moreover, we utilize peer-to-peer (P2P) communication between clients for selective expert model synchronization, thus significantly reducing the server-client transmissions. Experiments demonstrate that our FedMoE-DA achieves excellent performance while reducing the communication pressure on the server.
Ziwei Zhan, Wenkuan Zhao, Xiaoxi Zhang 0001, Chee-Wei Tan 0001, Chuan Wu 0001, Deke Guo, Xu Chen 0004
MSN7
2024 MPVSched: Multipath Transmissions and Video Frame Scheduling for Content Delivery Networks
abstract
With the widespread adoption of video streaming applications, effective video delivery solutions are crucial for providing seamless user experiences. Recent studies have revealed that multipath transmissions are beneficial to video streaming applications, given their potential of better load balancing and fault tolerance, relative to single path settings. However, the necessity of cross-layer co-design of multipath routing and video frame scheduling is overlooked. This work identifies that preset or path-oblivious frame scheduling used in existing works cannot adapt to network dynamics and fail to enhance the quality of experiences (QoE) in multipath transmissions. Therefore, we propose MPVSched, a novel framework that unifies the design of multipath routing and application-layer frame scheduling, with a particular focus on improving the rebuffer rate for short video delivery. At the network layer, we propose to use network-assisted routing that selects the optimal paths for each video transmission, with per-hop per-frame latency prediction. We implement an end-to-end QUIC-based video streaming system by integrating our routing strategy and application-layer frame scheduler, which effectively improves streaming efficiency and prevents user-side freezes. Our testbed experiments with real-world short video request traces demonstrate that MPVSched can achieve reductions of up to 28.58% in rebuffer ratio, compared to representative baseline methods.
Xiaoxi Zhang 0001, Jingpu Duan, Chuan Wu 0001, Jinhang Zuo, Xuan Zeng 0002, Yubing Qiu, Xu Chen 0004
NAS4
2024 POSTER: LLM-PQ: Serving LLM on Heterogeneous Clusters with Phase-Aware Partition and Adaptive Quantization
abstract
The immense sizes of Large-scale language models (LLMs) have led to high resource demand and cost for running the models. Though the models are largely served using uniform high-caliber GPUs nowadays, utilizing a heterogeneous cluster with a mix of available high- and low-capacity GPUs can potentially substantially reduce the serving cost. This paper proposes LLM-PQ, a system that advocates adaptive model quantization and phase-aware partition to improve LLM serving efficiency on heterogeneous GPU clusters. Extensive experiments on production inference workloads demonstrate throughput improvement in inference, showing great advantages over state-of-the-art works.
Juntao Zhao 0002, Borui Wan, Chuan Wu 0001, Yanghua Peng, Haibin Lin
PPoPP3
2024 Dynamic Flow Scheduling for DNN Training Workloads in Data Centers
abstract
Distributed deep learning (DL) training constitutes a significant portion of workloads in modern data centers that are equipped with high computational capacities, such as GPU servers. However, frequent tensor exchanges among workers during distributed deep neural network (DNN) training can result in heavy traffic in the data center network, leading to congestion at server NICs and in the switching network. Unfortunately, none of the existing DL communication libraries support active flow control to optimize tensor transmission performance, instead relying on passive adjustments to the congestion window or sending rate based on packet loss or delay. To address this issue, we propose a flow scheduler per host that dynamically tunes the sending rates of outgoing tensor flows from each server, maximizing network bandwidth utilization and expediting job training progress. Our scheduler comprises two main components: a monitoring module that interacts with state-of-the-art communication libraries supporting parameter server and all-reduce paradigms to track the training progress of DNN jobs, and a congestion control protocol that receives in-network feedback from traversing switches and computes optimized flow sending rates. For data centers where switches are not programmable, we provide a software solution that emulates switch behavior and interacts with the scheduler on servers. Experiments with real-world GPU testbed and trace-driven simulation demonstrate that our scheduler outperforms common rate control protocols and representative learning-based schemes in various settings.
Xiao-Yang Zhao 0005, Chuan Wu 0001
IEEE Trans. Netw. Serv. Manag.2
2024 Optimizing Task Placement and Online Scheduling for Distributed GNN Training Acceleration in Heterogeneous Systems
abstract
Training Graph Neural Networks (GNNs) on large graphs is resource-intensive and time-consuming, mainly due to the large graph data that cannot be fit into the memory of a single machine, but have to be fetched from distributed graph storage and processed on the go. Unlike distributed deep neural network (DNN) training, the bottleneck in distributed GNN training lies largely in large graph data transmission for constructing mini-batches of training samples. Existing solutions often advocate data-computation colocation, and do not work well with limited resources and heterogeneous training devices in heterogeneous clusters. The potentials of strategical task placement and optimal scheduling of data transmission and task execution have not been well explored. This paper designs an efficient algorithm framework for task placement and execution scheduling of distributed GNN training in heterogeneous systems, to better resource utilization, improve execution pipelining, and expedite training completion. Our framework consists of two modules: (i) an online scheduling algorithm that schedules the execution of training tasks, and the data transmission plan; and (ii) an exploratory task placement scheme that decides the placement of each training task. We conduct thorough theoretical analysis, testbed experiments and simulation studies, and observe up to 48% training speed-up with our algorithm as compared to representative baselines in our testbed settings.
Ziyue Luo, Yixin Bao, Chuan Wu 0001
IEEE/ACM Trans. Netw.3
2024 An Offline-Transfer-Online Framework for Cloud-Edge Collaborative Distributed Reinforcement Learning
abstract
Recent advances in deep reinforcement learning (DRL) have made it possible to train various powerful agents to perform complex tasks in real-time environments. With the next-generation communication technologies, making cloud-edge collaborative artificial intelligence service with evolved DRL agents can be a significant scenario. However, agents with different algorithms and architectures in the same DRL scenario may not be compatible, and training them is either time-consuming or resource-demanding. In this paper, we design a novel cloud-edge collaborative DRL training framework, named Offline-Transfer-Online, which is a new approach that can speed up the convergence of online DRL agents at the edge by interacting with offline agents in the cloud, with the minimum data interchanged and without relying on high-quality offline datasets. Therein, we propose a novel algorithm-independent knowledge distillation algorithm for online RL agents, by leveraging pre-trained models and the interface between agents and the environment to transfer distilled knowledge among multiple heterogeneous agents efficiently. Extensive experiments show that our algorithm can accelerate the convergence of various online agents in a double to decuple speed, with comparable reward achieved in different environments.
Tianyu Zeng, Xiaoxi Zhang 0001, Jingpu Duan, Chao Yu 0004, Chuan Wu 0001, Xu Chen 0004
IEEE Trans. Parallel Distributed Syst.5
2024 Swift: Expedited Failure Recovery for Large-Scale DNN Training
abstract
As the size of deep learning models gets larger and larger, training takes longer time and more resources, making fault tolerance more and more critical. Existing state-of-the-art methods like CheckFreq and Elastic Horovod need to back up a copy of the model state (i.e., parameters and optimizer states) in memory, which is costly for large models and leads to non-trivial overhead. This article presentsSwift, a novel recovery design for distributed deep neural network training that significantly reduces the failure recovery overhead without affecting training throughput and model accuracy. Instead of making an additional copy of the model state,Swiftresolves the inconsistencies of the model state caused by the failure and exploits the replicas of the model state in data parallelism for failure recovery. We propose a logging-based approach when replicas are unavailable, which records intermediate data and replays the computation to recover the lost state upon a failure. The re-computation is distributed across multiple machines to accelerate failure recovery further. We also log intermediate data selectively, exploring the trade-off between recovery time and intermediate data storage overhead. Evaluations show thatSwiftsignificantly reduces the failure recovery time and achieves similar or better training throughput during failure-free execution compared to state-of-the-art methods without degrading final model accuracy.Swiftcan also achieve up to 1.16x speedup in total training time compared to state-of-the-art methods.
Yuchen Zhong, Guangming Sheng, Jinhui Yuan, Chuan Wu 0001
IEEE Trans. Parallel Distributed Syst.5
2023 MixSynthFormer: A Transformer Encoder-like Structure with Mixed Synthetic Self-attention for Efficient Human Pose Estimation
abstract
Human pose estimation in videos has wide-ranging practical applications across various fields, many of which require fast inference on resource-scarce devices, necessitating the development of efficient and accurate algorithms. Previous works have demonstrated the feasibility of exploiting motion continuity to conduct pose estimation using sparsely sampled frames with transformer-based models. However, these methods only consider the temporal relation while neglecting spatial attention, and the complexity of dot product self-attention calculations in transformers are quadratically proportional to the embedding size. To address these limitations, we propose MixSynthFormer, a transformer encoder-like model with MLP-based mixed synthetic attention. By mixing synthesized spatial and temporal attentions, our model incorporates inter-joint and inter-frame importance and can accurately estimate human poses in an entire video sequence from sparsely sampled frames. Additionally, the flexible design of our model makes it versatile for other motion synthesis tasks. Our extensive experiments on 2D/3D pose estimation, body mesh recovery, and motion prediction validate the effectiveness and efficiency of MixSynthFormer. The code is available at https://github.com/ireneesun/MixSynthFormer.git
Yuran Sun, Alan William Dougherty, Zhuoying Zhang, Yi-King Choi, Chuan Wu 0001
ICCV5
2023 Towards Robust Graph Incremental Learning on Evolving Graphs
abstract
Incremental learning is a machine learning approach that involves training a model on a sequence of tasks, rather than all tasks at once. This ability to learn incrementally from a stream of tasks is crucial for many real-world applications. However, incremental learning is a challenging problem on graph-structured data, as many graph-related problems involve prediction tasks for each individual node, known as Node-wise Graph Incremental Learning (NGIL). This introduces non-independent and non-identically distributed characteristics in the sample data generation process, making it difficult to maintain the performance of the model as new tasks are added. In this paper, we focus on the inductive NGIL problem, which accounts for the evolution of graph structure (structural shift) induced by emerging tasks. We provide a formal formulation and analysis of the problem, and propose a novel regularization-based technique called Structural-Shift-Risk-Mitigation (SSRM) to mitigate the impact of the structural shift on catastrophic forgetting of the inductive NGIL problem. We show that the structural shift can lead to a shift in the input distribution for the existing tasks, and further lead to an increased risk of catastrophic forgetting. Through comprehensive empirical studies with several benchmark datasets, we demonstrate that our proposed method, Structural-Shift-Risk-Mitigation (SSRM), is flexible and easy to adapt to improve the performance of state-of-the-art GNN incremental learning frameworks in the inductive setting.
Junwei Su, Difan Zou, Chuan Wu 0001
ICML4
2023 HES: Edge Sampling for Heterogeneous Graphs
abstract
In light of the success of graph neural networks (GNNs), recent years have seen significant developments in modeling graphstructured data. Heterogeneous graphs have been widely adopted to model complex systems for various ML tasks. However, although some researchers have proposed methods for heterogeneous graphs, they merely focus on node features while neglecting the effectiveness of edge features. Besides, some research projects attempted incorporating edges into GNNs, but most regarded edge features as shared weights between node pairs. In this paper, we propose a two-stage method HES to learn the correlations among edge neighbors: (1) Graph Trans-formation: We convert the original heterogeneous graph into an undirected graph while preserving the orientation information, (2) Group Edge Sampling: To reduce the computation cost for the edge sampling in a heterogeneous graph, we propose to sample the most important edges over a group of edge neighbors instead of the whole graph, which leverages the edge features based on the one-hot encodings to describe the mutual influences between any adjacent edges. Finally, the experimental results on multiple public datasets show that HES outperforms existing state-of-the-art (SOTA) graph sampling methods. We further apply our approach to some existing GNN models as a pre-training process, demonstrating that HES can augment GNN-based models effectively.
Le Fang 0003, Chuan Wu 0001
IJCNN2
2023 TapFinger: Task Placement and Fine-Grained Resource Allocation for Edge Machine Learning
abstract
Machine learning (ML) tasks are one of the major workloads in today's edge computing networks. Existing edge-cloud schedulers allocate the requested amounts of resources to each task, falling short of best utilizing the limited edge resources flexibly for ML task performance optimization. This paper proposes TapFinger, a distributed scheduler that minimizes the total completion time of ML tasks in a multi-cluster edge network, through co-optimizing task placement and fine-grained multi-resource allocation. To learn the tasks' uncertain resource sensitivity and enable distributed online scheduling, we adopt multi-agent reinforcement learning (MARL), and propose several techniques to make it efficient for our ML-task resource allocation. First, TapFinger uses a heterogeneous graph attention network as the MARL backbone to abstract inter-related state features into more learnable environmental patterns. Second, the actor network is augmented through a tailored task selection phase, which decomposes the actions and encodes the optimization constraints. Third, to mitigate decision conflicts among agents, we novelly combine Bayes' theorem and masking schemes to facilitate our MARL model training. Extensive experiments using synthetic and test-bed ML task traces show that TapFinger can achieve up to 28.6% reduction in the average task completion time and improve resource efficiency as compared to state-of-the- art resource schedulers.
Tianyu Zeng, Xiaoxi Zhang 0001, Jingpu Duan, Chuan Wu 0001
INFOCOM5
2023 Two-level Graph Caching for Expediting Distributed GNN Training
abstract
Graph Neural Networks (GNNs) are increasingly popular due to excellent performance on learning graph-structured data in various domains. With fast expanding graph sizes and feature dimensions, distributed GNN training has been adopted, with multiple concurrent workers learning on different portions of a large graph. It has been observed that a main bottleneck in distributed GNN training lies in graph feature fetching across servers, which dominates the training time of each training iteration at each worker. This paper studies efficient feature caching on each worker to minimize feature fetching overhead, in order to expedite distributed GNN training. Current distributed GNN training systems largely adopt static caching of fixed neighbor nodes. We propose a novel two-level dynamic cache design exploiting both GPU memory and host memory at each worker, and design efficient two-level dynamic caching algorithms based on online optimization and a lookahead batching mechanism. Our dynamic caching algorithms consider node requesting probabilities and heterogeneous feature fetching costs from different servers, achieving an O(log3k) competitive ratio in terms of overall feature-fetching communication cost (where k is the cache capacity). We evaluate practical performance of our caching design with testbed experiments, and show that our design achieves up to 5.4x convergence speed-up.
Ziyue Luo, Chuan Wu 0001
INFOCOM3
2023 BGL: GPU-Efficient GNN Training by Optimizing Graph Data I/O and Preprocessing
Tianfeng Liu, Yangrui Chen, Dan Li 0001, Chuan Wu 0001, Yibo Zhu 0001, Yanghua Peng, Hongzheng Chen, Chuanxiong Guo
NSDI4
2023 Swift: Expedited Failure Recovery for Large-Scale DNN Training
abstract
As the size of deep learning models gets larger and larger, training takes longer time and more resources, making fault tolerance critical. Existing state-of-the-art methods like Check-Freq and Elastic Horovod need to back up a copy of the model state in memory, which is costly for large models and leads to non-trivial overhead. This paper presents Swift, a novel failure recovery design for distributed deep neural network training that significantly reduces the failure recovery overhead without affecting training throughput and model accuracy. Instead of making an additional copy of the model state, Swift resolves the inconsistencies of the model state caused by the failure and exploits replicas of the model state in data parallelism for failure recovery. We propose a logging-based approach when replicas are unavailable, which records intermediate data and replays the computation to recover the lost state upon a failure. Evaluations show that Swift significantly reduces the failure recovery time and achieves similar or better training throughput during failure-free execution compared to state-of-the-art methods without degrading final model accuracy.
Yuchen Zhong, Guangming Sheng, Jinhui Yuan, Chuan Wu 0001
PPoPP5
2023 SP-GNN: Learning structure and position information from graphs
Yangrui Chen, Jiaxuan You, Yanghua Peng, Chuan Wu 0001, Yibo Zhu 0001
Neural Networks6
2023 Online Scheduling Algorithm for Heterogeneous Distributed Machine Learning Jobs
abstract
Distributed machine learning (ML) has played a key role in today's proliferation of AI services. A typical model of distributed ML is to partition training datasets over multiple worker nodes to update model parameters in parallel, adopting aparameter serverorAllReducearchitecture. ML training jobs are typically resource elastic, completed using various time lengths with different resource configurations. A fundamental problem in a distributed ML cluster is how to explore the demand elasticity of ML jobs and schedule them with different resource configurations, such that the utilization of resources is maximized and average job completion time is minimized. To address it, we propose an online scheduling algorithm to decide the execution time window, the number and the type of concurrent workers and parameter servers for each job upon its arrival, with a goal of minimizing the weighted average completion time. Our online algorithm consists of (i) an online scheduling framework that groups unprocessed ML training jobs into a batch iteratively, and (ii) a batch scheduling algorithm that configures each ML job to maximize the total weight of scheduled jobs in the current iteration. Our online algorithm guarantees a good parameterized competitive ratio with polynomial time complexity. Extensive evaluations using real-world data demonstrate that it outperforms state-of-the-art schedulers in today's AI cloud systems.
Ruiting Zhou, Jinlong Pang, Chuan Wu 0001, Lei Jiao 0002, Zongpeng Li
IEEE Trans. Cloud Comput.4
2023 Caching in Dynamic Environments: A Near-Optimal Online Learning Approach
abstract
The rapid growth of rich multimedia data in today’s Internet, especially video traffic, has challenged the content delivery networks (CDNs). Caching serves as an important means to reduce user access latency so as to enable faster content downloads. Motivated by the dynamic nature of the real-world edge traces, this paper introduces aprovably wellonline caching policy in dynamic environments where: 1) the popularity is highly dynamic; 2) no regular stochastic pattern can model this dynamic evaluation process. First, we design an online optimization framework, which aims to minimize thedynamic regretthat finds the distance between an online caching policy and the best dynamic policy in hindsight. Second, we propose a dynamic online learning method to solve the non-stationary caching problem formulated in the previous framework. Compared to the linear dynamic regret of previous methods, our proposal is proved to achieve asublinear dynamic regret, from which it is guaranteed to be nearly optimal. We verify the design using both synthetic and real-world traces: the proposed policy achieves the best performance in the synthetic traces with different levels of dynamicity, which verifies the dynamic adaptation; our proposal consistently achieves at least 9.4% improvement than the baselines, including LRU, LFU, Static Online Learning based replacement, and Deep Reinforcement Learning based replacement, in random edge areas from real-world traces (from iQIYI), further verifying the effectiveness and robustness on the edge.
Shiji Zhou, Zhi Wang 0001, Chenghao Hu, Yinan Mao, Haopeng Yan, Shanghang Zhang, Chuan Wu 0001, Wenwu Zhu 0001
IEEE Trans. Multim.7
2023 Deep Learning-Based Job Placement in Distributed Machine Learning Clusters With Heterogeneous Workloads
abstract
Nowadays, most leading IT companies host a variety of distributed machine learning (ML) workloads in ML clusters to support AI-driven services, such as speech recognition, machine translation, and image processing. While multiple jobs are executed concurrently in a shared cluster to improve resource utilization, interference among co-located ML jobs can lead to significant performance downgrade. Existing cluster schedulers, such as YARN and Mesos, are interference-agnostic in their job placement, leading to suboptimal resource efficiency and usage. Some literature has studied interference-aware job placement policy, but relies on detailed workload profiling and interference modeling, which is not a general solution. In this work, we present Harmony, a deep learning-driven ML cluster scheduler that places heterogeneous training jobs (either with parameter server architecture or all-reduce architecture) in a manner that minimizes interference and maximizes performance (i.e., training completion time minimization). The design of Harmony is based on a carefully designed deep reinforcement learning (DRL) framework enhanced with reward modeling. The DRL integrates a dynamic sequence-to-sequence model with the state-of-the-art techniques to stabilize training and improve convergence, including actor-critic algorithm, job-aware action space exploration, multi-head attention, and experience replay. In view of a common lack of reward samples corresponding to different placement decisions, we build an auxiliary sequence-to-sequence reward prediction model, which is trained with historical samples and used for producing reward for unseen placement. Experiments using real ML workloads in a Kubernetes cluster of 6 GPU servers show that Harmony outperforms representative schedulers by 16%–42% in terms of average job completion time.
Yixin Bao, Yanghua Peng, Chuan Wu 0001
IEEE/ACM Trans. Netw.3
2023 Task Placement and Resource Allocation for Edge Machine Learning: A GNN-Based Multi-Agent Reinforcement Learning Paradigm
abstract
Machine learning (ML) tasks are one of the major workloads in today's edge computing networks. Existing edge-cloud schedulers allocate the requested amounts of resources to each task, falling short of best utilizing the limited edge resources for ML tasks. This paper proposesTapFinger, a distributed scheduler for edge clusters that minimizes the total completion time of ML tasks through co-optimizing task placement and fine-grained multi-resource allocation. To learn the tasks’ uncertain resource sensitivity and enable distributed scheduling, we adopt multi-agent reinforcement learning (MARL) and propose several techniques to make it efficient, including a heterogeneous graph attention network as the MARL backbone, a tailored task selection phase in the actor network, and the integration of Bayes’ theorem and masking schemes. We first implement asingle-task schedulingversion, which schedules at most one task each time. Then we generalize to themulti-task schedulingcase, in which a sequence of tasks is scheduled simultaneously. Our design can mitigate the expanded decision space and yield fast convergence to optimal scheduling solutions. Extensive experiments using synthetic and test-bed ML task traces show thatTapFingercan achieve up to 54.9% reduction in the average task completion time and improve resource efficiency as compared to state-of-the-art schedulers.
Xiaoxi Zhang 0001, Tianyu Zeng, Jingpu Duan, Chuan Wu 0001, Di Wu 0001, Xu Chen 0004
IEEE Trans. Parallel Distributed Syst.5
2023 Expediting Distributed DNN Training With Device Topology-Aware Graph Deployment
abstract
This paper presents TAG, an automatic system to derive optimized DNN training graph and its deployment onto any device topology, for expedited training in device- and topology- heterogeneous ML clusters. We novelly combine both the DNN computation graph and the device topology graph as input to a graph neural network (GNN), and join the GNN with a search-based method to quickly identify optimized distributed training strategies. To reduce communication in a heterogeneous cluster, we further explore a lossless gradient compression technique and solve a combinatorial optimization problem to automatically apply the technique for training time minimization. We evaluate TAG with various representative DNN models and device topologies, showing that it can achieve up to 4.56x training speed-up as compared to existing schemes. TAG can produce efficient deployment strategies for both unseen DNN models and unseen device topologies, without heavy fine-tuning.
Shiwei Zhang 0002, Xiaodong Yi 0001, Lansong Diao, Chuan Wu 0001, Siyu Wang 0006, Wei Lin 0016
IEEE Trans. Parallel Distributed Syst.4
2022 Accelerating large-scale distributed neural network training with SPMD parallelism
abstract
Deep neural networks (DNNs) with trillions of parameters have emerged, e.g., Mixture-of-Experts (MoE) models. Training models of this scale requires sophisticated parallelization strategies like the newly proposed SPMD parallelism, that shards each tensor along different dimensions. A common problem using SPMD is that computation stalls during communication due to data dependencies, resulting in low GPU utilization and long training time. We present a general technique to accelerate SPMD-based DNN training by maximizing computation-communication overlap and automatic SPMD strategy search. The key idea is to duplicate the DNN model into two copies that have no dependency, and interleave their execution such that computation of one copy overlaps with communication of the other. We propose a dynamic programming algorithm to automatically identify optimized sharding strategies that minimize model training time by maximally enabling computation-communication overlap. Experiments show that our designs achieve up to 61% training speed-up as compared to existing frameworks.
Shiwei Zhang 0002, Lansong Diao, Chuan Wu 0001, Siyu Wang 0006, Wei Lin 0016
SoCC3
2022 Optimizing Task Placement and Online Scheduling for Distributed GNN Training Acceleration
abstract
Training Graph Neural Networks (GNN) on large graphs is resource-intensive and time-consuming, mainly due to the large graph data that cannot be fit into the memory of a single machine, but have to be fetched from distributed graph storage and processed on the go. Unlike distributed deep neural network (DNN) training, the bottleneck in distributed GNN training lies largely in large graph data transmission for constructing mini-batches of training samples. Existing solutions often advocate data-computation colocation, and do not work well with limited resources where the colocation is infeasible. The potentials of strategical task placement and optimal scheduling of data transmission and task execution have not been well explored. This paper designs an efficient algorithm framework for task placement and execution scheduling of distributed GNN training, to better resource utilization, improve execution pipelining, and expediting training completion. Our framework consists of two modules: (i) an online scheduling algorithm that schedules the execution of training tasks, and the data transmission plan; and (ii) an exploratory task placement scheme that decides the placement of each training task. We conduct thorough theoretical analysis, testbed experiments and simulation studies, and observe up to 67% training speed-up with our algorithm as compared to representative baselines.
Ziyue Luo, Yixin Bao, Chuan Wu 0001
INFOCOM3
2022 Efficient Pipeline Planning for Expedited Distributed DNN Training
abstract
To train modern large DNN models, pipeline parallelism has recently emerged, which distributes the model across GPUs and enables different devices to process different microbatches in pipeline. Earlier pipeline designs allow multiple versions of model parameters to co-exist (similar to asynchronous training), and cannot ensure the same model convergence and accuracy performance as without pipelining. Synchronous pipelining has recently been proposed which ensures model performance by enforcing a synchronization barrier between training iterations. Nonetheless, the synchronization barrier requires waiting for gradient aggregation from all microbatches and thus delays the training progress. Optimized pipeline planning is needed to minimize such wait and hence the training time, which has not been well studied in the literature. This paper designs efficient, near-optimal algorithms for expediting synchronous pipeline-parallel training of modern large DNNs over arbitrary inter-GPU connectivity. Our algorithm framework comprises two components: a pipeline partition and device mapping algorithm, and a pipeline scheduler that decides processing order of microbatches over the partitions, which together minimize the per-iteration training time. We conduct thorough theoretical analysis, extensive testbed experiments and trace-driven simulation, and demonstrate our scheme can accelerate training up to 157% compared with state-of-the-art designs.
Ziyue Luo, Xiaodong Yi 0001, Guoping Long, Shiqing Fan, Chuan Wu 0001, Jun Yang 0052, Wei Lin 0016
INFOCOM5
2022 GADGET: Online Resource Optimization for Scheduling Ring-All-Reduce Learning Jobs
abstract
Fueled by advances in distributed deep learning (DDL), recent years have witnessed a rapidly growing demand for resource-intensive distributed/parallel computing to process DDL computing jobs. To resolve network communication bottleneck and load balancing issues in distributed computing, the so-called "ring-all-reduce" decentralized architecture has been increasingly adopted to remove the need for dedicated parameter servers. To date, however, there remains a lack of theoretical understanding on how to design resource optimization algorithms for efficiently scheduling ring-all-reduce DDL jobs in computing clusters. This motivates us to fill this gap by proposing a series of new resource scheduling designs for ring-all-reduce DDL jobs. Our contributions in this paper are threefold: i) We propose a new resource scheduling analytical model for ring-all-reduce deep learning, which covers a wide range of objectives in DDL performance optimization (e.g., excessive training avoidance, energy efficiency, fairness); ii) Based on the proposed performance analytical model, we develop an efficient resource scheduling algorithm called GADGET (greedy ring-all-reduce distributed graph embedding technique), which enjoys a provable strong performance guarantee; iii) We conduct extensive trace-driven experiments to demonstrate the effectiveness of the GADGET approach and its superiority over the state of the art.
Menglu Yu, Bo Ji 0001, Chuan Wu 0001, Hridesh Rajan, Jia Liu 0002
INFOCOM4
2022 Federated Graph Learning with Periodic Neighbour Sampling
abstract
Graph Convolutional Networks (GCN) proposed recently have achieved promising results on various graph learning tasks. Federated learning (FL) for GCN training is needed when learning from geo-distributed graph datasets. Existing FL paradigms are inefficient for geo-distributed GCN training since neighbour sampling across geo-locations will soon dominate the whole training process and consume large WAN bandwidth. We derive a practical federated graph learning algorithm, carefully striking the trade-off among GCN convergence error, wall-clock runtime, and neighbour sampling interval. Our analysis is divided into two cases according to the budget for neighbour sampling. In the unconstrained case, we obtain the optimal neighbour sampling interval, that achieves the best trade-off between convergence and runtime; in the constrained case, we show that determining the optimal sampling interval is actually an online problem and we propose a novel online algorithm with bounded competitive ratio to solve it. Combining the two cases, we propose a unified algorithm to decide the neighbour sampling interval in federated graph learning, and demonstrate its effectiveness with extensive simulation over graph datasets from real applications.
Bingqian Du, Chuan Wu 0001
IWQoS2
2022 SAPipe: Staleness-Aware Pipeline for Data Parallel DNN Training
abstract
Data parallelism across multiple machines is widely adopted for accelerating distributed deep learning, but it is hard to achieve linear speedup due to the heavy communication. In this paper, we propose SAPipe, a performant system that pushes the training speed of data parallelism to its fullest extent. By introducing partial staleness, the communication overlaps the computation with minimal staleness in SAPipe. To mitigate additional problems incurred by staleness, SAPipe adopts staleness compensation techniques including weight prediction and delay compensation with provably lower error bounds. Additionally, SAPipe presents an algorithm-system co-design with runtime optimization to minimize system overhead for the staleness training pipeline and staleness compensation. We have implemented SAPipe in the BytePS framework, compatible to both TensorFlow and PyTorch. Our experiments show that SAPipe achieves up to 157% speedups over BytePS (non-stale), and outperforms PipeSGD in accuracy by up to 13.7%.
Yangrui Chen, Juncheng Gu, Yanghua Peng, Haibin Lin, Chuan Wu 0001, Yibo Zhu 0001
NeurIPS7
2022 Differentially private recommender system with variational autoencoders
Le Fang 0003, Bingqian Du, Chuan Wu 0001
Knowl. Based Syst.3
2022 Large-Scale Machine Learning Cluster Scheduling via Multi-Agent Graph Reinforcement Learning
abstract
Efficient scheduling of distributed deep learning (DL) jobs in large GPU clusters is crucial for resource efficiency and job performance. While server sharing among jobs improves resource utilization, interference among co-located DL jobs occurs due to resource contention. Interference-aware job placement has been studied, with white-box approaches based on explicit interference modeling and black-box schedulers with reinforcement learning. In today’s clusters containing thousands of GPU servers, running a single scheduler to manage all arrival jobs in a timely and effective manner is challenging, due to the large workload scale. We adopt multiple schedulers in a large-scale cluster/data center, and propose a multi-agent reinforcement learning (MARL) scheduling framework to cooperatively learn fine-grained job placement policies, towards the objective of minimizing job completion time (JCT). To achieve topology-aware placements, our proposed framework uses hierarchical graph neural networks to encode the data center topology and server architecture. In view of a common lack of precise reward samples corresponding to different placements, a job interference model is further devised to predict interference levels in face of various co-locations, for training of the MARL schedulers. Testbed and trace-driven evaluations show that our scheduler framework outperforms representative scheduling schemes by more than 20% in terms of average JCT, and is adaptive to various machine learning cluster topologies.
Xiao-Yang Zhao 0005, Chuan Wu 0001
IEEE Trans. Netw. Serv. Manag.2
2022 Optimizing DNN Compilation for Distributed Training With Joint OP and Tensor Fusion
abstract
This article proposesDisCo, an automatic deep learning compilation module for data-parallel distributed training. Unlike most deep learning compilers that focus on training or inference on a single device,DisCooptimizes a DNN model for distributed training over multiple GPU machines. Existing single-device compilation strategies do not work well in distributed training, due mainly to communication inefficiency that they incur.DisCogenerates optimized, joint computation operator and communication tensor fusion strategies to enable highly efficient distributed training. A GNN-based simulator is built to effectively estimate per-iteration training time achieved by operator/tensor fusion candidates. A backtracking search algorithm is driven by the simulator, navigating efficiently in the large strategy space to identify good operator/tensor fusion strategies that minimize distributed training time. We compareDisCowith existing DL fusion schemes and show that it achieves good training speed-up close to the ideal, full computation-communication overlap case.
Xiaodong Yi 0001, Shiwei Zhang 0002, Lansong Diao, Chuan Wu 0001, Zhen Zheng, Shiqing Fan, Siyu Wang 0006, Jun Yang 0052, Wei Lin 0016
IEEE Trans. Parallel Distributed Syst.4
2021 A Sum-of-Ratios Multi-Dimensional-Knapsack Decomposition for DNN Resource Scheduling
abstract
In recent years, to sustain the resource-intensive computational needs for training deep neural networks (DNNs), it is widely accepted that exploiting the parallelism in large-scale computing clusters is critical for the efficient deployments of DNN training jobs. However, existing resource schedulers for traditional computing clusters are not well suited for DNN training, which results in unsatisfactory job completion time performance. The limitations of these resource scheduling schemes motivate us to propose a new computing cluster resource scheduling framework that is able to leverage the special layered structure of DNN jobs and significantly improve their job completion times. Our contributions in this paper are three-fold: i) We develop a new resource scheduling analytical model by considering DNN's layered structure, which enables us to analytically formulate the resource scheduling optimization problem for DNN training in computing clusters; ii) Based on the proposed performance analytical model, we then develop an efficient resource scheduling algorithm based on the widely adopted parameter-server architecture using a sum-of-ratios multi-dimensional-knapsack decomposition (SMD) method to offer strong performance guarantee; iii) We conduct extensive numerical experiments to demonstrate the effectiveness of the proposed schedule algorithm and its superior performance over the state of the art.
Menglu Yu, Chuan Wu 0001, Bo Ji 0001, Jia Liu 0002
INFOCOM2
2021 Near-Optimal Topology-adaptive Parameter Synchronization in Distributed DNN Training
abstract
Distributed machine learning with multiple concurrent workers has been widely adopted to train large deep neural networks (DNNs). Parameter synchronization is a key component in each iteration of distributed training, where workers exchange locally computed gradients through an AllReduce operation or parameter servers, for global parameter updates. Parameter synchronization often constitutes a significant portion of the training time; minimizing the communication time contributes substantially to DNN training speed-up. Standard ring-based AllReduce or PS architecture work efficiently mostly with homogeneous inter-worker connectivity. However, available bandwidth among workers in real-world clusters is often heterogeneous, due to different hardware configurations, switching topologies, and contention with concurrent jobs. This work investigates the best parameter synchronization topology and schedule among workers for most expedited communication in distributed DNN training. We show that the optimal parameter synchronization topology should be comprised of trees with different workers as roots, each for aggregating or broadcasting a partition of gradients/parameters. We identify near-optimal forest packing to maximally utilize available bandwidth and overlap aggregation and broadcast stages to minimize communication time. We provide theoretical analysis of the performance bound, and show that our scheme outperforms state-of-the-art parameter synchronization schemes by up to 18.3 times with extensive evaluation under various settings.
Chuan Wu 0001, Zongpeng Li
INFOCOM2
2021 DAPPLE: a pipelined data parallel approach for training large models
abstract
It is a challenging task to train large DNN models on sophisticated GPU platforms with diversified interconnect capabilities. Recently, pipelined training has been proposed as an effective approach for improving device utilization. However, there are still several tricky issues to address: improving computing efficiency while ensuring convergence, and reducing memory usage without incurring additional computing costs. We propose DAPPLE, a synchronous training framework which combines data parallelism and pipeline parallelism for large DNN models. It features a novel parallelization strategy planner to solve the partition and placement problems, and explores the optimal hybrid strategies of data and pipeline parallelism. We also propose a new runtime scheduling algorithm to reduce device memory usage, which is orthogonal to re-computation approach and does not come at the expense of training throughput. Experiments show that DAPPLE planner consistently outperforms strategies generated by PipeDream's planner by up to 3.23× speedup under synchronous training scenarios, and DAPPLE runtime outperforms GPipe by 1.6× speedup of training throughput and saves 12% of memory consumption at the same time.
Shiqing Fan, Zongyan Cao, Siyu Wang 0006, Zhen Zheng, Chuan Wu 0001, Guoping Long, Jun Yang 0052, Lixue Xia, Lansong Diao, Wei Lin 0016
PPoPP7
2021 Joint Model and Data Adaptation for Cloud Inference Serving
abstract
Real-time deep learning inference serving systems often require prohibitive resources and diverse user requirements. The existing design of inference serving systems mainly focusing on computation resource efficiency, largely ignoring the trade-off between computation and bandwidth resources in need. Sub-optimal resource utilization usually leads to huge serving cost waste. In this paper, we tackle the dual challenge of computation-bandwidth trade-off and cost-effectiveness by proposing A2, an efficient joint Adaptive model, and Adaptive data deep learning serving solution across the geo-datacenters. Inspired by the insight that a trade-off between computational cost and bandwidth cost in achieving the same accuracy, we design a real-time inference serving framework, which selectively places different "versions" of the deep learning models at different geo-locations, and schedules different data sample versions to be sent to those model versions for inference. The goal is to minimize the total serving cost while meeting latency and accuracy demand for the serving requests. We formulate a joint placement and serving problem and propose an efficient approximation algorithm to solve it with a theoretical performance guarantee. We deploy A2on Amazon EC2 for experiments, which shows that A2achieves 30%-50% serving cost reduction under the same required latency and accuracy as compared to baselines.
Jingyan Jiang, Ziyue Luo, Chenghao Hu, Zhaoliang He, Zhi Wang 0001, Shutao Xia, Chuan Wu 0001
RTSS7
2021 A Nonintrusive Elderly Home Monitoring System
abstract
Home anomaly monitoring is crucial for the elderly who live alone. A number of IoT-based home monitoring systems have been available, but most rely on privacy-intrusive cameras. With more and more concerns on privacy and security of human data, anomaly detection based on nonintrusive IoT devices becomes more desirable. Considering the elderly consumers, a low-cost system with good detection accuracy is further critical for the system's acceptability by elderly users. We propose a smart home monitoring system for living-alone senior citizens, relying on carefully designed, low-cost infrared sensor devices, as well as a cloud-based data processing and anomaly detection platform. Our PIR sensor device is effective in continuous monitoring of motion data in a user's apartment, and an open-hardware software platform is devised to support sensors manufactured by various vendors in the IoT system, all for cost reduction purpose. For privacy preservation, we encrypt collected data and store data indices in a blockchain system, to achieve efficient data access control and auditing. For motion anomaly detection, we propose a simple but effective environment adaptation method to work with the one-class support vector machine (OCSVM) method. Experiments driven by real-world traces show good reliability, accuracy, and efficiency of our system.
Le Fang 0003, Yu Wu 0010, Chuan Wu 0001, Yizhou Yu
IEEE Internet Things J.3
2021 Dynamic VM Scaling: Provisioning and Pricing through an Online Auction
abstract
Today's IaaS clouds allow dynamic scaling of VMs allocated to a user, according to real-time demand of the user. There are two types of scaling: horizontal scaling (scale-out) by allocating more VM instances to the user, and vertical scaling (scale-up) by boosting resources of VMs owned by the user. It has been a daunting issue how to efficiently allocate the resources on physical servers to meet the scaling demand of users on the go, which achieves the best server utilization and user utility. An accompanying critical challenge is how to effectively charge the incremental resources, such that the economic benefits of both the cloud provider and cloud users are guaranteed. There has been online auction design dealing with dynamic VM provisioning, where the resource bids are not related to each other, failing to handle VM scaling where later bids may rely on earlier bids of the same user. As the first in the literature, this paper designs an efficient, truthful online auction for resource provisioning and pricing in the practical cases of dynamic VM scaling, where: (i) users bid for customized VMs to use in future durations, and can bid again in the following time to increase resources, indicating both scale-up and scale-out options; (ii) the cloud provider packs the demanded VMs on heterogeneous servers for energy cost minimization on the go. We carefully design resource prices maintained for each type of resource on each server to achieve threshold-based online allocation and charging, as well as a novel competitive analysis technique based on submodularity of the offline objective, to show a good competitive ratio is achieved. The efficacy of the online auction is validated through solid theoretical analysis and trace-driven simulations.
Xiaoxi Zhang 0001, Zhiyi Huang 0002, Chuan Wu 0001, Zongpeng Li, Francis C. M. Lau 0001
IEEE Trans. Cloud Comput.3
2021 DL2: A Deep Learning-Driven Scheduler for Deep Learning Clusters
abstract
Efficient resource scheduling is essential for maximal utilization of expensive deep learning (DL) clusters. Existing cluster schedulers either are agnostic to machine learning (ML) workload characteristics, or use scheduling heuristics based on operators' understanding of particular ML framework and workload, which are less efficient or not general enough. In this article, we show that DL techniques can be adopted to design a generic and efficient scheduler. Specifically, we propose DL2, a DL-driven scheduler for DL clusters, targeting global training job expedition by dynamically resizing resources allocated to jobs. DL2 advocates a joint supervised learning and reinforcement learning approach: a neural network is warmed up via offline supervised learning based on job traces produced by the existing cluster scheduler; then the neural network is plugged into the live DL cluster, fine-tuned by reinforcement learning carried out throughout the training progress of the DL jobs, and used for deciding job resource allocation in an online fashion. We implement DL2 on Kubernetes and enable dynamic resource scaling in DL jobs on MXNet. Extensive evaluation shows that DL2 outperforms fairness scheduler (i.e., DRF) by 44.1 percent and expert heuristic scheduler (i.e., Optimus) by 17.5 percent in terms of average job completion time.
Yanghua Peng, Yixin Bao, Yangrui Chen, Chuan Wu 0001, Wei Lin 0016
IEEE Trans. Parallel Distributed Syst.4
2020 Distributed Machine Learning through Heterogeneous Edge Systems
abstract
Many emerging AI applications request distributed machine learning (ML) among edge systems (e.g., IoT devices and PCs at the edge of the Internet), where data cannot be uploaded to a central venue for model training, due to their large volumes and/or security/privacy concerns. Edge devices are intrinsically heterogeneous in computing capacity, posing significant challenges to parameter synchronization for parallel training with the parameter server (PS) architecture. This paper proposes ADSP, a parameter synchronization model for distributed machine learning (ML) with heterogeneous edge systems. Eliminating the significant waiting time occurring with existing parameter synchronization models, the core idea of ADSP is to let faster edge devices continue training, while committing their model updates at strategically decided intervals. We design algorithms that decide time points for each worker to commit its model update, and ensure not only global model convergence but also faster convergence. Our testbed implementation and experiments show that ADSP outperforms existing parameter synchronization models significantly in terms of ML model convergence time, scalability and adaptability to large heterogeneity.
Hanpeng Hu, Chuan Wu 0001
AAAI3
2020 Elastic parameter server load distribution in deep learning clusters
abstract
In distributed DNN training, parameter servers (PS) can become performance bottlenecks due to PS stragglers, caused by imbalanced parameter distribution, bandwidth contention, or computation interference. Few existing studies have investigated efficient parameter (aka load) distribution among PSs. We observe significant training inefficiency with the current parameter assignment in representative machine learning frameworks (e.g., MXNet, TensorFlow), and big potential for training acceleration with better PS load distribution. We design PSLD, a dynamic parameter server load distribution scheme, to mitigate PS straggler issues and accelerate distributed model training in the PS architecture. An exploitation-exploration method is carefully designed to scale in and out parameter servers and adjust parameter distribution among PSs on the go. We also design an elastic PS scaling module to carry out our scheme with little interruption to the training process. We implement our module on top of open-source PS architectures, including MXNet and BytePS. Testbed experiments show up to 2.86x speed-up in model training with PSLD, for different ML models under various straggler settings.
Yangrui Chen, Yanghua Peng, Yixin Bao, Chuan Wu 0001, Yibo Zhu 0001, Chuanxiong Guo
SoCC4
2020 Optimizing distributed training deployment in heterogeneous GPU clusters
abstract
This paper proposes HeteroG, an automatic module to accelerate deep neural network training in heterogeneous GPU clusters. To train a deep learning model with large amounts of data, distributed training using data or model parallelism has been widely adopted, mostly over homogeneous devices (GPUs, network bandwidth). Heterogeneous training environments may often exist in shared clusters with GPUs of different models purchased in different batches and network connections of different bandwidth availability (e.g., due to contention). Classic data parallelism does not work well in a heterogeneous cluster, while model-parallel training is hard to plan. HeteroG enables highly-efficient distributed training over heterogeneous devices, by automatically converting a single-GPU training model to a distributed one according to the deep learning graph and available resources. HeteroG embraces operation-level hybrid parallelism, communication architecture selection and execution scheduling, based on a carefully designed strategy framework exploiting both GNN-based learning and combinatorial optimization. We compare HeteroG with existing parallelism schemes and show that it achieves up-to 222% training speed-up. HeteroG also enables efficient training of large models over a set of heterogeneous devices where simple parallelism is infeasible.
Xiaodong Yi 0001, Shiwei Zhang 0002, Ziyue Luo, Guoping Long, Lansong Diao, Chuan Wu 0001, Zhen Zheng, Jun Yang 0052, Wei Lin 0016
CoNEXT6
2020 Preemptive All-reduce Scheduling for Expediting Distributed DNN Training
abstract
Data-parallel training is widely used for scaling DNN training over large datasets, using the parameter server or all-reduce architecture. Communication scheduling has been promising to accelerate distributed DNN training, which aims to overlap communication with computation by scheduling the order of communication operations. We identify two limitations of previous communication scheduling work. First, layer-wise computation graph has been a common assumption, while modern machine learning frameworks (e.g., TensorFlow) use a sophisticated directed acyclic graph (DAG) representation as the execution model. Second, the default sizes of tensors are often less than optimal for transmission scheduling and bandwidth utilization. We propose PACE, a communication scheduler that preemptively schedules (potentially fused) all-reduce tensors based on the DAG of DNN training, guaranteeing maximal overlapping of communication with computation and high bandwidth utilization. The scheduler contains two integrated modules: given a DAG, we identify the best tensor-preemptive communication schedule that minimizes the training time; exploiting the optimal communication scheduling as an oracle, a dynamic programming approach is developed for generating a good DAG, which merges small communication tensors for efficient bandwidth utilization. Experiments in a GPU testbed show that PACE accelerates training with representative system configurations, achieving up to 36% speed-up compared with state-of-the-art solutions.
Yixin Bao, Yanghua Peng, Yangrui Chen, Chuan Wu 0001
INFOCOM4
2020 Fast Training of Deep Learning Models over Multiple GPUs
abstract
This paper proposes FastT, a transparent module to work with the TensorFlow framework for automatically identifying a satisfying deployment and execution order of operations in DNN models over multiple GPUs, for expedited model training. We propose white-box algorithms to compute the strategies with small computing resource consumption in a short time. Recently, similar studies have been done to optimize device placement using reinforcement learning. Compared to those works which learn to optimize device placement of operations in several hours using large amounts of computing resources, our approach can find excellent device placement and execution order within minutes using the same computing node as for training. We design a list of scheduling algorithms to compute the device placement and execution order for each operation and also design an algorithm to split operations in the critical path to support fine-grained (mixed) data and model parallelism to further improve the training speed in each iteration. We compare FastT with representative strategies and obtain insights on the best strategies for training different types of DNN models based on extensive testbed experiments.
Xiaodong Yi 0001, Ziyue Luo, Mengdi Wang 0001, Guoping Long, Chuan Wu 0001, Jun Yang 0052, Wei Lin 0016
Middleware6
2020 Online scheduling of heterogeneous distributed machine learning jobs
abstract
Distributed machine learning (ML) has played a key role in today's proliferation of AI services. A typical model of distributed ML is to partition training datasets over multiple worker nodes to update model parameters in parallel, adopting a parameter server architecture. ML training jobs are typically resource elastic, completed using various time lengths with different resource configurations. A fundamental problem in a distributed ML cluster is how to explore the demand elasticity of ML jobs and schedule them with different resource configurations, such that the utilization of resources is maximized and average job completion time is minimized. To address it, we propose an online scheduling algorithm to decide the execution time window, the number and the type of concurrent workers and parameter servers for each job upon its arrival, with a goal of minimizing the weighted average completion time. Our online algorithm consists of (i) an online scheduling framework that groups unprocessed ML training jobs into a batch iteratively, and (ii) a batch scheduling algorithm that configures each ML job to maximize the total weight of scheduled jobs in the current iteration. Our online algorithm guarantees a good parameterized competitive ratio with polynomial time complexity. Extensive evaluations using real-world data demonstrate that it outperforms state-of-the-art schedulers in today's AI cloud systems.
Ruiting Zhou, Chuan Wu 0001, Lei Jiao 0002, Zongpeng Li
MobiHoc3
2020 Fair Online Power Capping for Emergency Handling in Multi-Tenant Cloud Data Centers
abstract
In view of the high capital expense for scaling up power capacity to meet the escalating demand, maximizing the utilization of built capacity has become a top priority for multi-tenant data center operators, where many cloud providers house their physical servers. The traditional power provisioning guarantees a high availability, but is very costly and results in a significant capacity under-utilization. On the other hand, power oversubscription (i.e., deploying more servers than what the capacity allows) improves utilization but offers no availability guarantees due to the necessity of power reduction to handle the resulting power emergencies. Given these limitations, we propose a novel hybrid power provisioning approach, called HyPP, which provides a combination of two different power availabilities to tenants: capacity with a very high availability (100 percent or nearly 100 percent), plus additional capacity with a medium availability that may be unavailable for up to a certain amount during each billing period. For HyPP, we design an online algorithm for the operator to coordinate tenants' power reduction at runtime when the tenants' aggregate power demand exceeds the power capacities. Our algorithm aims at achieving long-term fairness in tenants' power reduction (defined as the ratio of total actual power reduction by a tenant to its contracted reduction budget over a billing period). We analyze the theoretical performance of our online algorithm and derive a good competitive ratio in terms of fairness compared to the offline optimum. We also validate our algorithm through simulations under realistic settings.
Shaolei Ren, Chuan Wu 0001
IEEE Trans. Cloud Comput.3
2020 A Truthful $(1-\epsilon)$(1-ε)-Optimal Mechanism for On-Demand Cloud Resource Provisioning
abstract
On-demand resource provisioning in cloud computing provides tailor-made resource packages (typically in the form of VMs) to meet users' demands. Public clouds nowadays provide elaborated types of VMs, but have yet to offer the most flexible dynamic VM assembly, which is partly due to the lack of a mature mechanism for pricing tailor-made VMs. This work proposes an efficient randomized auction mechanism based on a novel application of smoothed analysis and randomized reduction, for dynamic VM provisioning and pricing in geo-distributed cloud data centers. To the best of our knowledge, it is the first one in literature that achieves (i) truthfulness in expectation, (ii) polynomial running time in expectation, and (iii) (1 - ε)-optimal social welfare in expectation for resource allocation, where ε can be arbitrarily close to0. Our mechanism consists of three modules: (1) an exact algorithm to solve the NP-hard social welfare maximization problem, which has polynomial run-time in expectation, (2) a perturbation-based randomized resource allocation scheme which produces an allocation solution that is (1 - ε)-optimal and (3) an auction mechanism prices the customized VMs using a randomized VCG payment, with a guarantee in truthfulness in expectation. We validate the efficacy of the mechanism through theoretical analysis and trace-driven simulations.
Xiaoxi Zhang 0001, Chuan Wu 0001, Zongpeng Li, Francis C. M. Lau 0001
IEEE Trans. Cloud Comput.2
2020 An Online Algorithm for VNF Service Chain Scaling in Datacenters
abstract
Built on top of virtualization technologies, network function virtualization (NFV) provides flexible and scalable software implementation of various network functions. Virtual network functions (VNFs), which are network functions implemented as virtual machines, are chained together to provide network services. Dynamic deployment of VNFs while satisfying incoming network traffic demand is the key to cost optimization of an NFV system. Besides considering server resource capacity and incoming traffic rates, an optimal scaling policy needs to strike a balance between VNF's operational costs, the costs for maintaining VNF instances, and VNF deployment costs, additional costs when setting up new VNF instances on a server. This paper targets dynamic scaling of VNF instances in a cloud data center where multiple VNF chains are running. We propose an online scaling algorithm to adjust the deployment of VNF instances according to time-varying traffic demand, ensuring a good competitive ratio. Through theoretical analysis and trace-driven simulation, we demonstrate effectiveness of the proposed online VNF scaling algorithm.
Ziyue Luo, Chuan Wu 0001
IEEE/ACM Trans. Netw.2
2020 Online Placement and Scaling of Geo-Distributed Machine Learning Jobs via Volume-Discounting Brokerage
abstract
Geo-distributed machine learning (ML) often uses large geo-dispersed data collections produced over time to train global models, without consolidating the data to a central site. In the parameter server architecture, “workers” and “parameter servers” for a geo-distributed ML job should be strategically deployed and adjusted on the fly, to allow easy access to the datasets and fast exchange of the model parameters at anytime. Despite many cloud platforms now provide volume discounts to encourage the usage of their ML resources, different geo-distributed ML jobs that run in the clouds often rent cloud resources separately and respectively, thus rarely enjoying the benefit of discounts. We study an ML broker service that aggregates geo-distributed ML jobs into cloud data centers for volume discounts via dynamic online placement and scaling of workers and parameter servers in individual jobs for long-term cost minimization. To decide the number and the placement of workers and parameter servers, we propose an efficient online algorithm which first decomposes the online problem into a series of one-shot optimization problems solvable at each individual time slot by the technique of regularization, and afterwards round the fractional decisions to the integer ones via a carefully-designed dependent rounding method. We prove a parameterized-constant competitive ratio for our online algorithm as the theoretical performance analysis, and also conduct extensive simulation studies to exhibit its close-to-offline-optimum practical performance in realistic settings.
Ruiting Zhou, Lei Jiao 0002, Chuan Wu 0001, Yuhang Deng, Zongpeng Li
IEEE Trans. Parallel Distributed Syst.4
2019 Learning Resource Allocation and Pricing for Cloud Profit Maximization
abstract
Cloud computing has been widely adopted to support various computation services. A fundamental problem faced by cloud providers is how to efficiently allocate resources upon user requests and price the resource usage, in order to maximize resource efficiency and hence provider profit. Existing studies establish detailed performance models of cloud resource usage, and propose offline or online algorithms to decide allocation and pricing. Differently, we adopt a blackbox approach, and leverage model-free Deep Reinforcement Learning (DRL) to capture dynamics of cloud users and better characterize inherent connections between an optimal allocation/pricing policy and the states of the dynamic cloud system. The goal is to learn a policy that maximizes net profit of the cloud provider through trial and error, which is better than decisions made on explicit performance models. We combine long short-term memory (LSTM) units with fully-connected neural networks in our DRL to deal with online user arrivals, and adjust the output and update methods of basic DRL algorithms to address both resource allocation and pricing. Evaluation based on real-world datasets shows that our DRL approach outperforms basic DRL algorithms and state-of-theart white-box online cloud resource allocation/pricing algorithms significantly, in terms of both profit and the number of accepted users.
Bingqian Du, Chuan Wu 0001, Zhiyi Huang 0002
AAAI2
2019 A Provably-Efficient Online Algorithm for Re-Utilizing Unused VM Resources for Edge Providers
abstract
In recent years, as a result of the rapidly growing volume of generated data, computation has been increasingly migrating from megascale data centers to Internet edges (a.k.a edge computing), for avoiding high latencies and overwhelmed bandwidths. Unlike centralized clouds, edge computing processes workloads generated by users nearby. Thus, due to the lack of statistical multiplexing from a large group of users, the resource demand at an edge data center exhibits more fluctuations, resulting in time-varying unused computation resources. In this paper, we propose our UNusEd spAred VM Re-uTilizing mecHanism, UNEARTH, to utilize different types of unused resources, such as storage, CPU, GPU, and so on, offered by an edge computing provider. Notably, the exact amount of unused VM resources is unknown before selling them. We evaluate the performance of our algorithms under realistic settings, showing that our proposed VM bundle allocation algorithm can achieve (1+ Ω/Ω-1 ε (e)M1/Ω-1-1)-approximation in the worst case compared with optimums; and overall, our algorithms outperform the existing and heuristic algorithms.
Shaolei Ren, Chuan Wu 0001
ICC3
2019 FlowShader: a Generalized Framework for GPU-accelerated VNF Flow Processing
abstract
GPU acceleration has been widely investigated for packet processing in virtual network functions (NFs), but not for L7 flow-processing NFs. In L7 NFs, reassembled TCP messages of the same flow should be processed in order in the same processing thread, and the uneven sizes among flows pose a major challenge for full realization of GPU's parallel computation power. To exploit GPUs for L7 NF processing, this paper presents FlowShader, a GPU acceleration framework to achieve both high generality and throughput even under skewed flow size distributions. We carefully design an efficient scheduling algorithm that fully exploits available GPU and CPU capacities; in particular, we dispatch large flows which seriously break up the size balance to CPU and the rest of flows to GPU. Furthermore, FlowShader allows similar NF logic (as CPU-based NFs) to run on individual threads in a GPU, which is more generalized and easy to take on as compared to redesigning an NF for operation parallelism on GPU. We implemented a number of L7 flow processing NFs based on FlowShader. Evaluations are conducted under both synthetic and real-world traffic traces and results show that the throughput achieved by FlowShader is up to 6x that of the CPU-only baseline and 3x of the GPU-only design.
Xiaodong Yi 0001, Jingpu Duan, Wei Bai 0001, Chuan Wu 0001, Yongqiang Xiong, Dongsu Han
ICNP5
2019 Deep Learning-based Job Placement in Distributed Machine Learning Clusters
abstract
Production machine learning (ML) clusters commonly host a variety of distributed ML workloads, e.g., speech recognition, machine translation. While server sharing among jobs improves resource utilization, interference among co-located ML jobs can lead to significant performance downgrade. Existing cluster schedulers (e.g., Mesos) are interference-oblivious in their job placement, causing suboptimal resource efficiency. Interference-aware job placement has been studied in the literature, but was treated using detailed workload profiling and interference modeling, which is not a general solution. This paper presents Harmony, a deep learning-driven ML cluster scheduler that places training jobs in a manner that minimizes interference and maximizes performance (i.e., training completion time). Harmony is based on a carefully designed deep reinforcement learning (DRL) framework augmented with reward modeling. The DRL employs state-of-the-art techniques to stabilize training and improve convergence, including actor-critic algorithm, job-aware action space exploration and experience replay. In view of a common lack of reward samples corresponding to different placement decisions, we build an auxiliary reward prediction model, which is trained using historical samples and used for producing reward for unseen placement. Experiments using real ML workloads in a Kubernetes cluster of 6 GPU servers show that Harmony outperforms representative schedulers by 25% in terms of average job completion time.
Yixin Bao, Yanghua Peng, Chuan Wu 0001
INFOCOM3
2019 A generic communication scheduler for distributed DNN training acceleration
abstract
We present ByteScheduler, a generic communication scheduler for distributed DNN training acceleration. ByteScheduler is based on our principled analysis that partitioning and rearranging the tensor transmissions can result in optimal results in theory and good performance in real-world even with scheduling overhead. To make ByteScheduler work generally for various DNN training frameworks, we introduce a unified abstraction and a Dependency Proxy mechanism to enable communication scheduling without breaking the original dependencies in framework engines. We further introduce a Bayesian Optimization approach to auto-tune tensor partition size and other parameters for different training models under various networking conditions. ByteScheduler now supports TensorFlow, PyTorch, and MXNet without modifying their source code, and works well with both Parameter Server (PS) and all-reduce architectures for gradient synchronization, using either TCP or RDMA. Our experiments show that ByteScheduler accelerates training with all experimented system configurations and DNN models, by up to 196% (or 2.96X of original speed).
Yanghua Peng, Yibo Zhu 0001, Yangrui Chen, Yixin Bao, Bairen Yi, Chang Lan, Chuan Wu 0001, Chuanxiong Guo
SOSP7
2019 NetStar: A Future/Promise Framework for Asynchronous Network Functions
abstract
Network functions (NFs) are more than simple packet processors that apply various transformations to the packet content. Modern NFs often resort to various external services to achieve their purposes, e.g., storing flow states in an external storage or looking up a DNS. Working with external services is usually implemented using callback-based asynchronous programming, which is complex and error-prone. This paper proposes NetStar, a new NF programming framework that brings the future/promise abstraction to the NF dataplane for flow processing. NetStar simplifies asynchronous NF programming via a carefully designed async-flow interface that exploits the future/promise paradigm by chaining multiple continuation functions for asynchronous operations handling. The programs implemented using the NetStar framework mimic simple synchronous programming but are able to achieve full flow processing asynchrony. We have used NetStar to implement a number of representative NFs. Our experience and evaluation results show that NetStar can effectively simplify asynchronous NF programming by substantially reducing the lines of code, while still approaching line-rate packet processing speeds.
Jingpu Duan, Xiaodong Yi 0001, Chuan Wu 0001, Franck Le
IEEE J. Sel. Areas Commun.4
2019 NFVactor: A Resilient NFV System Using the Distributed Actor Model
abstract
Resilience functionality, including failure resilience and flow migration, is of pivotal importance in practical network function virtualization (NFV) systems. However, existing failure recovery procedures incur high packet processing delay due to heavyweight process checkpointing, while flow migration has poor performance due to centralized control. This paper proposes NFVactor, a novel NFV system that aims to provide lightweight failure resilience and high-performance flow migration. NFVactorenables these by using actor model to provide a per-flow execution environment, so that each flow can replicate and migrate itself with improved parallelism, while the efficiency of the actor model is guaranteed by a carefully designed runtime system. Moreover, NFVactorachieves transparent resilience: once a new network function (NF) is implemented for NFVactor, the NF automatically acquires resilience support. Our evaluation result shows that NFVactorachieves 10-Gbps packet processing, flow migration completion time that is 144 times faster than the existing system, and packet processing delay stabilized at around 20 μs during replication.
Jingpu Duan, Xiaodong Yi 0001, Shixiong Zhao, Chuan Wu 0001, Heming Cui, Franck Le
IEEE J. Sel. Areas Commun.4
2019 Scaling Geo-Distributed Network Function Chains: A Prediction and Learning Framework
abstract
Geo-distributed virtual network function (VNF) chaining has been useful, such as in network slicing in 5G networks and for network traffic processing in the WAN. Agile scaling of the VNF chains according to real-time traffic rates is the key in network function virtualization. Designing efficient scaling algorithms is challenging, especially for geo-distributed chains, where bandwidth costs and latencies incurred by the WAN traffic are important but difficult to handle in making scaling decisions. Existing studies have largely resorted to optimization algorithms in scaling design. Aiming at better decisions empowered by in-depth learning from experiences, this paper proposes a deep learning-based framework for scaling of the geo-distributed VNF chains, exploring inherent pattern of traffic variation and good deployment strategies over time. We novelly combine a recurrent neural network as the traffic model for predicting upcoming flow rates and a deep reinforcement learning (DRL) agent for making chain placement decisions. We adopt the experience replay technique based on the actor-critic DRL algorithm to optimize the learning results. Trace-driven simulation shows that with limited offline training, our learning framework adapts quickly to traffic dynamics online and achieves lower system costs, compared to the existing representative algorithms.
Ziyue Luo, Chuan Wu 0001, Zongpeng Li
IEEE J. Sel. Areas Commun.2
2019 An Efficient Online Placement Scheme for Cloud Container Clusters
abstract
Containers represent an agile alternative to virtual machines (VMs), for providing cloud computing services. Containers are more flexible and lightweight, and can be easily instrumented. Enterprise users often create clusters of inter-connected containers to provision complex services. Compared to traditional cloud services, key challenges in container cluster (CC) provisioning lie in the optimal placement of containers while considering inter-container traffic in a CC. The challenge further escalates, when CCs are provisioned in an online fashion. We propose an online algorithm to address the above challenges, aiming to maximize the aggregate value of all served clusters. We first study a one-shot CC placement problem. Leveraging techniques of exhaustive sampling and ST rounding, we design an efficient one-shot algorithm to determine the placement scheme of a given CC. We then propose a primal-dual online placement scheme that employs the one-shot algorithm as a building block to make decisions upon the arrival of each CC request. Through both theoretical analysis and trace-driven simulations, we verify that the online placement algorithm is computationally efficient and achieves a good competitive ratio.
Ruiting Zhou, Zongpeng Li, Chuan Wu 0001
IEEE J. Sel. Areas Commun.3
2018 Optimus: an efficient dynamic resource scheduler for deep learning clusters
abstract
Deep learning workloads are common in today's production clusters due to the proliferation of deep learning driven AI services (e.g., speech recognition, machine translation). A deep learning training job is resource-intensive and time-consuming. Efficient resource scheduling is the key to the maximal performance of a deep learning cluster. Existing cluster schedulers are largely not tailored to deep learning jobs, and typically specifying a fixed amount of resources for each job, prohibiting high resource efficiency and job performance. This paper proposes Optimus, a customized job scheduler for deep learning clusters, which minimizes job training time based on online resource-performance models. Optimus uses online fitting to predict model convergence during training, and sets up performance models to accurately estimate training speed as a function of allocated resources in each job. Based on the models, a simple yet effective method is designed and used for dynamically allocating resources and placing deep learning tasks to minimize job completion time. We implement Optimus on top of Kubernetes, a cluster manager for container orchestration, and experiment on a deep learning cluster with 7 CPU servers and 6 GPU servers, running 9 training jobs using the MXNet framework. Results show that Optimus outperforms representative cluster schedulers by about 139% and 63% in terms of job completion time and makespan, respectively.
Yanghua Peng, Yixin Bao, Yangrui Chen, Chuan Wu 0001, Chuanxiong Guo
EuroSys4
2018 Online Cloud Resource Allocation and Pricing with Server Speed Scaling
abstract
The provisioning of cloud computing services typically incurs huge electricity costs. Utilization maximization of the cloud resources and efficient resource pricing have been key factors determining a cloud provider's revenue. On the other hand, dynamic CPU speed scaling has been widely supported by modern operating systems and hypervisors as an efficient technique for CPU energy saving, potentially useful for cutting down provider's electricity bill. In this paper, we propose an online mechanism for resource allocation and pricing on a cloud platform, which enables dynamic CPU speed scaling for achieving the best job execution efficiency. Using a novel compact infinite optimization technique and the primal-dual online algorithm design framework, our online mechanism achieves computational efficiency, truthfulness, and near-optimal social welfare during the long run of the cloud system. Trace-driven simulation studies further demonstrate good performance of our mechanism in realistic settings.
Ziyue Luo, Zongpeng Li, Chuan Wu 0001
ICC3
2018 Online Job Scheduling in Distributed Machine Learning Clusters
abstract
Nowadays large-scale distributed machine learning systems have been deployed to support various analytics and intelligence services in IT firms. To train a large dataset and derive the prediction/inference model, e.g., a deep neural network, multiple workers are run in parallel to train partitions of the input dataset, and update shared model parameters. In a shared cluster handling multiple training jobs, a fundamental issue is how to efficiently schedule jobs and set the number of concurrent workers to run for each job, such that server resources are maximally utilized and model training can be completed in time. Targeting a distributed machine learning system using the parameter server framework, w e design an online algorithm for scheduling the arriving jobs and deciding the adjusted numbers of concurrent workers and parameter servers for each job over its course, to maximize overall utility of all jobs, contingent on their completion times. Our online algorithm design utilizes a primal-dual framework coupled with efficient dual subroutines, achieving good long-term performance guarantees with polynomial time complexity. Practical effectiveness of the online algorithm is evaluated using trace-driven simulation and testbed experiments, which demonstrate its outperformance as compared to commonly adopted scheduling algorithms in today's cloud systems.
Yixin Bao, Yanghua Peng, Chuan Wu 0001, Zongpeng Li
INFOCOM3
2018 Occupation-Oblivious Pricing of Cloud Jobs via Online Learning
abstract
State-of-the-art cloud platforms adopt pay-as-you-go pricing, where users pay for the resources on demand according to occupation time. Simple and intuitive as it is, such a pricing scheme is a mismatch for new workloads today such as large-scale machine learning, whose completion time is hard to estimate beforehand. To supplement existing cloud pricing schemes, we propose an occupation-oblivious online pricing mechanism for cloud jobs without pre-specified time duration and for users who prefer a pre-determined cost for job execution. Our strategy posts unit resource prices upon user arrival and decides a fixed charge for completing the user's job, without the need to know how long the job is to occupy the requested resources. At the core of our design is a novel multi-armed bandit based online learning algorithm for estimating unknown input by exploration and exploitation of past resource sales, and deciding resource prices to maximize profit of the cloud provider in an online setting. Our online learning algorithm achieves a low regret sublinear with the time horizon, in terms of overall provider profit, compared with an omniscient benchmark. We also conduct trace-driven simulations to verify efficacy of the algorithm in real-world settings.
Xiaoxi Zhang 0001, Chuan Wu 0001, Zhiyi Huang 0002, Zongpeng Li
INFOCOM2
2018 A Shapley-Value Mechanism for Bandwidth On Demand between Datacenters
abstract
Recent studies in cloud resource allocation and pricing have focused on computing and storage resources but not network bandwidth. Cloud users nowadays customarily deploy services across multiple geo-distributed datacenters, with significant inter-datacenter traffic generated, paid by cloud providers to ISPs. An effective bandwidth allocation and charging mechanism is needed between the cloud provider and the cloud users. Existing volume based static charging schemes lack market efficiency. This work presents the first dynamic pricing mechanism for inter-data-center on-demand bandwidth, via a Shapley value based auction. Our auction is expressive enough to accept bids as a flat bandwidth rate plus a time duration, or a data volume with a transfer deadline. We start with an offline auction, design an optimal end-to-end traffic scheduling approach, and exploit the Shapley value in computing payments. Our auction is truthful, individual rational, budget balanced and approximately efficient in social welfare. An online version of the auction follows, where decisions are made instantly upon the arrival of each user's realtime transmission demand. We propose an efficient online traffic scheduling algorithm, and approximate the offline Shapley value based payments on the fly. We validate our mechanism design with solid theoretical analysis, as well as trace-driven simulation studies.
Chuan Wu 0001, Zongpeng Li
IEEE Trans. Cloud Comput.2
2018 A Truthful Online Mechanism for Location-Aware Tasks in Mobile Crowd Sensing
abstract
Effective incentive mechanisms are invaluable in mobile crowd sensing, for stimulating participation of smartphone users. Online auction mechanisms represent a natural solution for such sensing task allocation. Departing from existing studies that focus on an isolated system round, we optimize social cost across the system lifespan, while considering location constraints and capacity constraints when assigning sensing tasks to users. The winner determination problem (WDP) at each round is NP-hard even without inter-round coupling imposed by user capacity constraints. We first propose a truthful one-round auction, comprising of an approximation algorithm for solving the one-round WDP and a payment scheme for computing remuneration to winners. We then propose an online algorithm framework that employs the one-round auction as a building block towards a flexible mechanism that makes on-spot decisions upon dynamically arriving bids. Through both theoretical analysis and trace-driven simulations, we demonstrate that our online auction is truthful, individually rational, computationally efficient, and achieves a good competitive ratio.
Ruiting Zhou, Zongpeng Li, Chuan Wu 0001
IEEE Trans. Mob. Comput.3
2018 Online Scaling of NFV Service Chains Across Geo-Distributed Datacenters
Yongzheng Jia, Chuan Wu 0001, Zongpeng Li, Franck Le, Alex X. Liu
IEEE/ACM Trans. Netw.2
2018 Scheduling Frameworks for Cloud Container Services
abstract
Compared with traditional virtual machines, cloud containers are more flexible and lightweight, emerging as the new norm of cloud resource provisioning. We exploit this new algorithm design space, and propose scheduling frameworks for cloud container services. Our offline and online schedulers permit partial execution, and allow a job to specify its job deadline, desired cloud containers, and inter-container dependence relations. We leverage the following classic and new techniques in our scheduling algorithm design. First, we apply the compact-exponential technique to express and handle nonconventional scheduling constraints. Second, we adopt the primal-dual framework that determines the primal solution based on its dual constraints in both the offline and online algorithms. The offline scheduling algorithm includes a new separation oracle to separate violated dual constraints, and works in concert with the randomized rounding technique to provide a near-optimal solution. The online scheduling algorithm leverages the online primal-dual framework with a learning-based scheme for obtaining dual solutions. Both theoretical analysis and trace-driven simulations validate that our scheduling frameworks are computationally efficient and achieve close-to-optimal aggregate job valuation.
Ruiting Zhou, Zongpeng Li, Chuan Wu 0001
IEEE/ACM Trans. Netw.3
2017 Online Learning-Assisted VNF Service Chain Scaling with Network Uncertainties
abstract
Network function virtualization has emerged as a promising technology to enable rapid network service composition/innovation, energy conservation and cost minimization for network operators. To optimally operate a virtualized network service, it is of key importance to optimally deploy a VNF (virtualized network function) service chain within the provisioning infrastructure (e.g., servers and the network within a cloud datacenter), and dynamically scale it in response to flow traffic changes. Most of the existing work on VNF scaling assume access to precise network bandwidth information for placement decisions, while in reality, network bandwidth typically fluctuates following an unknown pattern and an effective way to adapt to it is to do trials. In this paper, we address dynamic VNF service chain deployment and scaling by a novel combination of an online provisioning algorithm and a multi-armed bandit optimization framework, which exploits online learning of the available bandwidths to enable optimal deployment of a scaled service chain. Specifically, we adopt the online algorithm to minimize the cost for provisioning VNF instances on the go, and a bandit-based online learning algorithm to place the VNF instances which minimizes the congestion in a datacenter network. We demonstrate effectiveness of our algorithms using solid theoretical analysis and trace-driven evaluation.
Chuan Wu 0001, Franck Le, Francis C. M. Lau 0001
CLOUD2
2017 GPUNFV: a GPU-Accelerated NFV System
abstract
This paper presents GPUNFV, a high-performance NFV system providing flow-level micro services for stateful service chains with Graphics Processing Unit (GPU) acceleration. GPUNFV exploits the massively-parallel processing power of GPU to maximize the throughput of the NFV system. Combined with the customized flow handler, GPUNFV achieves a much better throughput than the existing NFV systems. With a carefully designed GPU-based virtualized network function framework, GPUNFV is able to efficiently support both stateful and stateless network functions. We have implemented a number of GPU-based network functions and a preliminary GPUNFV system to demonstrate the lexibility and potential of our design.
Xiaodong Yi 0001, Jingpu Duan, Chuan Wu 0001
APNet3
2017 Double Auction for Resource Allocation in Cloud Computing
Fei Chen 0013, T.-H. Hubert Chan, Chuan Wu 0001
CLOSER4
2017 A Scalable and Distributed Approach for NFV Service Chain Cost Minimization
abstract
Network function virtualization (NFV) represents the latest technology advancement in network service provisioning. Traditional hardware middleboxes are replaced by software programs running on industry standard servers and virtual machines, for service agility, flexibility, and cost reduction. NFV users are provisioned with service chains composed of virtual network functions (VNFs). A fundamental problem in NFV service chain provisioning is to satisfy user demands with minimum system-wide cost. We jointly consider two types of cost in this work: nodal resource cost and link delay cost, and formulate the service chain provisioning problem using nonlinear optimization. Through the method of auxiliary variables, we transform the optimization problem into its separable form, and then apply the alternating direction method of multipliers (ADMM) to design scalable and fully distributed solutions. Through simulation studies, we verify the convergence and efficacy of our distributed algorithm design.
Zongpeng Li, Chuan Wu 0001, Chuanhe Huang
ICDCS3
2017 Virtualized Network Coding Functions on the Internet
abstract
Network coding is a fundamental tool that enables higher network capacity and lower complexity in routing algorithms, by encouraging the mixing of information flows in the middle of a network. Implementing network coding in the core Internet is subject to practical concerns, since Internet routers are often overwhelmed by packet forwarding tasks, leaving little processing capacity for coding operations. Inspired by the recent paradigm of network function virtualization, we propose implementing network coding as a new network function, and deploying such coding functions in geo-distributed cloud data centers, to practically enable network coding on the Internet. We target multicast sessions (including unicast flows as special cases), strategically deploy relay nodes (network coding functions) in selected data centers between senders and receivers, and embrace high bandwidth efficiency brought by network coding with dynamic coding function deployment. We design and implement the network coding function on typical virtual machines, featuring efficient packet processing. We propose an efficient algorithm for coding function deployment, scaling in and out, in the presence of system dynamics. Real-world implementation on Amazon EC2 and Linode demonstrates significant throughput improvement and higher robustness of multicast via coding functions as well as efficiency of the dynamic deployment and scaling algorithm.
Linquan Zhang, Shangqi Lai, Chuan Wu 0001, Zongpeng Li, Chuanxiong Guo
ICDCS3
2017 Proactive VNF provisioning with multi-timescale cloud resources: Fusing online learning and online optimization
abstract
Network Function Virtualization (NFV) represents a new paradigm of network service provisioning. NFV providers acquire cloud resources, install virtual network functions (VNFs), assemble VNF service chains for customer usage, and dynamically scale VNF deployment against input traffic fluctuations. While existing literature on VNF scaling mostly adopts a reactive approach, we target a proactive approach that is more practical given the time overhead for VNF deployment. We aim to effectively estimate upcoming traffic rates and adjust VNF deployment a priori, for flow service quality assurance and resource cost minimization. We adapt online learning techniques for predicting future service chain workloads. We further combine the online learning method with a multi-timescale online optimization algorithm for VNF scaling, through minimization of the regret due to inaccurate demand prediction and minimization of the cost incurred by sub-optimal online decisions in a joint online optimization framework. The resulting proactive online VNF provisioning algorithm achieves a good performance guarantee, as shown by both theoretical analysis and simulation under realistic settings.
Xiaoxi Zhang 0001, Chuan Wu 0001, Zongpeng Li, Francis C. M. Lau 0001
INFOCOM2
2017 deTector: a Topology-aware Monitoring System for Data Center Networks
Yanghua Peng, Chuan Wu 0001, Chuanxiong Guo, Chengchen Hu, Zongpeng Li
USENIX ATC3
2017 Virtualized resource sharing in cloud radio access networks: An auction approach
Ruiting Zhou, Xunrui Yin, Zongpeng Li, Chuan Wu 0001
Comput. Commun.4
2017 Dynamic Scaling of Virtualized, Distributed Service Chains: A Case Study of IMS
abstract
The emerging paradigm of network function virtualization advocates deploying virtualized network functions (VNFs) on standard virtualization platforms for significant cost reduction and management flexibility. There have been system designs for managing dynamic deployment and scaling of VNF service chains within one cloud datacenter. Many real-world network services involve geo-distributed service chains, with prominent examples of mobile core networks and IP multimedia subsystems (IMSs)). Virtualizing these service chains requires efficient coordination of dynamic VNF deployment across geo-distributed data centers, calling for a new management system. This paper designs a dynamic scaling system for geo-distributed VNF service chains, using the case of an IMS. IMSs are widely used subsystems for delivering multimedia services among mobile users in a 3G/4G network, whose virtualization has been broadly advocated in the industry for reducing cost, improving network usage efficiency and enabling dynamic network topology reconfiguration for performance optimization. Our scaling system design caters to key control-plane and data-plane service chains in an IMS, combining proactive and reactive approaches for timely, cost-effective scaling of the service chains. The design principles are applicable to scaling of other systems with multiple related service chains. We evaluate our system using real-world experiments on both an emulation platform and a geo-distributed public cloud.
Jingpu Duan, Chuan Wu 0001, Franck Le, Alex X. Liu, Yanghua Peng
IEEE J. Sel. Areas Commun.2
2017 Online Stochastic Buy-Sell Mechanism for VNF Chains in the NFV Market
abstract
With the recent advent of network functions virtualization (NFV), enterprises and businesses are looking into network service provisioning through the service chains of virtual network functions (VNFs), instead of relying on dedicated hardware middleboxes. Accompanying this trend, an NFV market is emerging, where NFV service providers create VNF instances, assemble VNF service chains, and sell them for the use of customers, using resources (computing, bandwidth) that they own or rent from other resource suppliers. Efficient service chain provisioning and pricing mechanisms are still missing, to charge assembled service chains according to demand and the supply of resources at any time. We propose an online stochastic auction mechanism for on-demand service chain provisioning and pricing at an NFV provider. Our auction takes in buy bids for service chains from multiple customers and sell bids from various resource suppliers to supplement the NFV provider's geo-distributed resource pool, with resource occupation/contribution durations. We extend online primal-dual optimization framework for handling both buyers and sellers, with a new competitive analysis. The online mechanism maximizes the expected social welfare of the NFV ecosystem (the NFV provider, customers and resource suppliers) with a good competitive ratio as compared with the expected offline optimal social welfare, while guaranteeing truthfulness in bidding, individual rationality for both buyers and sellers, and polynomial time for computation. We evaluate our mechanism through trace-driven simulation studies, and demonstrate a close-to-offline-optimal performance in expected social welfare under realistic settings.
Xiaoxi Zhang 0001, Zhiyi Huang 0002, Chuan Wu 0001, Zongpeng Li, Francis C. M. Lau 0001
IEEE J. Sel. Areas Commun.3
2017 Orchestrating Bulk Data Transfers across Geo-Distributed Datacenters
abstract
As it has become the norm for cloud providers to host multiple datacenters around the globe, significant demands exist for inter-datacenter data transfers in large volumes, e.g., migration of big data. A challenge arises on how to schedule the bulk data transfers at different urgency levels, in order to fully utilize the available inter-datacenter bandwidth. The Software Defined Networking (SDN) paradigm has emerged recently which decouples the control plane from the data paths, enabling potential global optimization of data routing in a network. This paper aims to design a dynamic, highly efficient bulk data transfer service in a geo-distributed datacenter system, and engineer its design and solution algorithms closely within an SDN architecture. We model data transfer demands as delay tolerant migration requests with different finishing deadlines. Thanks to the flexibility provided by SDN, we enable dynamic, optimal routing of distinct chunks within each bulk data transfer (instead of treating each transfer as an infinite flow), which can be temporarily stored at intermediate datacenters to mitigate bandwidth contention with more urgent transfers. An optimal chunk routing optimization model is formulated to solve for the best chunk transfer schedules over time. To derive the optimal schedules in an online fashion, three algorithms are discussed, namely a bandwidth-reserving algorithm, a dynamically-adjusting algorithm, and a future-demand-friendly algorithm, targeting at different levels of optimality and scalability. We build an SDN system based on the Beacon platform and OpenFlow APIs, and carefully engineer our bulk data transfer algorithms in the system. Extensive real-world experiments are carried out to compare the three algorithms as well as those from the existing literature, in terms of routing optimality, computational delay and overhead.
Yu Wu 0010, Chuan Wu 0001, Chuanxiong Guo, Zongpeng Li, Francis C. M. Lau 0001
IEEE Trans. Cloud Comput.3
2017 Virtualized Resource Sharing in Cloud Radio Access Networks Through Truthful Mechanisms
abstract
In the recent paradigm of cloud radio access networks (C-RAN), signal processing functions at the base stations (BSs) are virtualized and migrated into a mobile cloud that maintains a pool of virtual BS (VBS) instances. Remote radio heads and antennae at the BSs are connected to the VBS pool by fronthaul fiber links. Mobile operators may lease resources from the tower company who owns the C-RAN infrastructure. We study auction mechanisms for efficiently sharing C-RAN resources among mobile operators. Leveraging randomized rounding, we design an offline C-RAN auction mechanism that can achieve truthfulness and near-optimal social welfare. For the more realistic setting of online bid arrival, we design an online algorithm that executes in polynomial time and achieves a competitive ratio of $(1-\epsilon )$ . A tailored fractional Vickrey-Clarke-Groves mechanism works in concert with the online algorithm to elicit truthful bids. Extensive simulation studies verify the efficacy of our C-RAN auction mechanisms.
Sijia Gu, Zongpeng Li, Chuan Wu 0001, Huyin Zhang
IEEE Trans. Commun.3
2017 Cost-Effective Low-Delay Design for Multiparty Cloud Video Conferencing
abstract
Multiparty cloud video conferencing architecture has been recently advocated to exploit rich computing and bandwidth resources in the cloud to effectively improve video conferencing performance. As a typical design in this architecture, multiple agents, i.e., virtual machines, are deployed in different cloud sites, and users are assigned to the agents. Then, the users communicate through the agents, and the agents might transcode the recorded videos given the heterogeneities among devices in terms of hardware specification and connectivity. In this architecture, two critical and nontrivial challenges are: 1) assigning users to agents to reduce the operational cost and the user-to-user conferencing delay and 2) identifying best agents to perform transcoding tasks, taking into account the heterogeneous bandwidth and processing availabilities. To address these challenges, we cast a joint problem of user-to-agent assignment and transcoding-agent selection. The ultimate objective is to simultaneously minimize the cost of the service provider and the conferencing delay. The problem is combinatorial in nature, which belongs to the NP-hard node assignment problems. We leverage the Markov approximation framework and devise an adaptive parallel algorithm that finds a close-to-optimal solution to our problem with a bounded performance guarantee. To evaluate the performance of our solution, we implement a prototype video conferencing system and carry out trace-driven experiments. In a set of largescale experiments using PlanetLab traces, our solution decreases the operational cost by 77% and simultaneously yields lower conferencing delay compared with an existing alternative.
Mohammad Hajiesmaili, Lok To Mak, Zhi Wang 0001, Chuan Wu 0001, Minghua Chen 0001, Ahmad Khonsari
IEEE Trans. Multim.4
2017 Online Auctions in IaaS Clouds: Welfare and Profit Maximization With Server Costs
abstract
Auction design has recently been studied for dynamic resource bundling and virtual machine (VM) provisioning in IaaS clouds, but is mostly restricted to one-shot or offline setting. This paper targets a more realistic case of online VM auction design, where: 1) cloud users bid for resources into the future to assemble customized VMs with desired occupation durations, possibly located in different data centers; 2) the cloud provider dynamically packs multiple types of resources on heterogeneous physical machines (servers) into the requested VMs; 3) the operational costs of servers are considered in resource allocation; and 4) both social welfare and the cloud provider's net profit are to be maximized over the system running span. We design truthful, polynomial time auctions to achieve social welfare maximization and/or the provider's profit maximization with good competitive ratios. Our mechanisms consist of two main modules: 1) an online primal-dual optimization framework for VM allocation to maximize the social welfare with server costs, and for revealing the payments through the dual variables to guarantee truthfulness and 2) a randomized reduction algorithm to convert the social welfare maximizing auctions to ones that provide a maximal expected profit for the provider, with competitive ratios comparable to those for social welfare. We adopt a new application of Fenchel duality in our primal-dual framework, which provides richer structures for convex programs than the commonly used Lagrangian duality, and our optimization framework is general and expressive enough to handle various convex server cost functions. The efficacy of the online auctions is validated through careful theoretical analysis and trace-driven simulation studies.
Xiaoxi Zhang 0001, Zhiyi Huang 0002, Chuan Wu 0001, Zongpeng Li, Francis C. M. Lau 0001
IEEE/ACM Trans. Netw.3
2017 An Efficient Cloud Market Mechanism for Computing Jobs With Soft Deadlines
abstract
This paper studies the cloud market for computing jobs with completion deadlines, and designs efficient online auctions for cloud resource provisioning. A cloud user bids for future cloud resources to execute its job. Each bid includes: 1) a utility, reflecting the amount that the user is willing to pay for executing its job and 2) a soft deadline, specifying the preferred finish time of the job, as well as a penalty function that characterizes the cost of violating the deadline. We target cloud job auctions that executes in an online fashion, runs in polynomial time, provides truthfulness guarantee, and achieves optimal social welfare for the cloud ecosystem. Towards these goals, we leverage the following classic and new auction design techniques. First, we adapt the posted pricing auction framework for eliciting truthful online bids. Second, we address the challenge posed by soft deadline constraints through a new technique of compact exponential-size LPs coupled with dual separation oracles. Third, we develop efficient social welfare approximation algorithms using the classic primal-dual framework based on both LP duals and Fenchel duals. Empirical studies driven by real-world traces verify the efficacy of our online auction design.
Ruiting Zhou, Zongpeng Li, Chuan Wu 0001, Zhiyi Huang 0002
IEEE/ACM Trans. Netw.3
2017 An Online Auction Mechanism for Dynamic Virtual Cluster Provisioning in Geo-Distributed Clouds
abstract
It is common for cloud users to require clusters of inter-connected virtual machines (VMs) in a geo-distributed IaaS cloud, to run their services. Compared to isolated VMs, key challenges on dynamic virtual cluster (VC) provisioning (computation + communication resources) lie in two folds: (1) optimal placement of VCs and inter-VM traffic routing involve NP-hard problems, which are non-trivial to solve offline, not to mention if an online efficient algorithm is sought; (2) an efficient pricing mechanism is missing, which charges a market-driven price for each VC as a whole upon request, while maximizing system efficiency or provider revenue over the entire span. This paper proposes efficient online auction mechanisms to address the above challenges. We first design SWMOA, a novel online algorithm for dynamic VC provisioning and pricing, achieving truthfulness, individual rationality, computation efficiency, and (1 + 2 log μ)-competitiveness in social welfare, where m is related to the problem size. Next, applying a randomized reduction technique, we convert the social welfare maximizing auction into a revenue maximizing online auction, PRMOA, achieving O(log μ)-competitiveness in provider revenue, as well as truthfulness, individual rationality and computation efficiency. We investigate auction design in different cases of resource cost functions in the system. We validate the efficacy of the mechanisms through solid theoretical analysis and trace-driven simulations.
Chuan Wu 0001, Zongpeng Li
IEEE Trans. Parallel Distributed Syst.2
2016 Online VNF Scaling in Datacenters
abstract
Network Function Virtualization (NFV) is a promising technology that promises to significantly reduce the operational costs of network services by deploying virtualized network functions (VNFs) to commodity servers in place of dedicated hardware middleboxes. The VNFs are typically running on virtual machine instances in a cloud infrastructure, where the virtualization technology enables dynamic provisioning of VNF instances, to process the fluctuating traffic that needs to go through the network functions in a network service. In this paper, we target dynamic provisioning of enterprise network services - expressed as one or multiple service chains - in cloud datacenters, and design efficient online algorithms without requiring any information on future traffic rates. The key is to decide the number of instances of each VNF type to provision at each time, taking into consideration the server resource capacities and traffic rates between adjacent VNFs in a service chain. In the case of a single service chain, we discover an elegant structure of the problem and design an efficient randomized algorithm achieving a e/(e-1) competitive ratio. For multiple concurrent service chains, an online heuristic algorithm is proposed, which is O(1)-competitive. We demonstrate the effectiveness of our algorithms using solid theoretical analysis and trace-driven simulations.
Chuan Wu 0001, Franck Le, Alex X. Liu, Zongpeng Li, Francis C. M. Lau 0001
CLOUD2
2016 An efficient auction mechanism for service chains in the NFV market
abstract
Network Function Virtualization (NFV) is emerging as a new paradigm for providing elastic network functions through flexible virtual network function (VNF) instances executed on virtualized computing platforms exemplified by cloud datacenters. In the new NFV market, well defined VNF instances each realize an atomic function that can be chained to meet user demands in practice. This work studies the dynamic market mechanism design for the transaction of VNF service chains in the NFV market, to help relinquish the full power of NFV. Combining the techniques of primal-dual approximation algorithm design with Myerson's characterization of truthful mechanisms, we design a VNF chain auction that runs efficiently in polynomial time, guarantees truthfulness, and achieves near-optimal social welfare in the NFV eco-system. Extensive simulation studies verify the efficacy of our auction mechanism.
Sijia Gu, Zongpeng Li, Chuan Wu 0001, Chuanhe Huang
INFOCOM3
2016 An online mechanism for dynamic virtual cluster provisioning in geo-distributed clouds
abstract
It is common for cloud users to require clusters of inter-connected virtual machines (VMs) in a geo-distributed IaaS cloud, to run their services. Compared to isolated VMs, key challenges on dynamic virtual cluster (VC) provisioning (computation + communication resources) lie in two folds: (1) optimal placement of VCs and inter-VM traffic routing involve NP-hard problems, which are non-trivial to solve offline, not to mention if an online efficient algorithm is sought; (2) an efficient pricing mechanism is missing, which charges a market-driven price for each VC as a whole upon request, while maximizing system efficiency or provider revenue over the entire span. This paper proposes efficient online auction mechanisms to address the above challenges. We first design SWMOA, a novel online algorithm for dynamic VC provisioning and pricing, achieving truthfulness, individual rationality, computation efficiency, and (1 + 2 log μ)-competitiveness in social welfare, where μ is related to the problem size. Next, applying a randomized reduction technique, we convert the social welfare maximizing auction into a revenue maximizing online auction, PRMOA, achieving O(log μ)-competitiveness in provider revenue, as well as truthfulness, individual rationality and computation efficiency. We validate the efficacy of the mechanisms through solid theoretical analysis and trace-driven simulations.
Chuan Wu 0001, Zongpeng Li
INFOCOM2
2016 Online influence maximization in non-stationary Social Networks
abstract
Social networks have been popular platforms for information propagation. An important use case is viral marketing: given a promotion budget, an advertiser can choose some influential users as the seed set and provide them free or discounted sample products; in this way, the advertiser hopes to increase the popularity of the product in the users' friend circles by the world-of-mouth effect, and thus maximizes the number of users that information of the production can reach. There has been a body of literature studying the influence maximization problem. Nevertheless, the existing studies mostly investigate the problem on a one-off basis, assuming fixed known influence probabilities among users, or the knowledge of the exact social network topology. In practice, the social network topology and the influence probabilities are typically unknown to the advertiser, which can be varying over time, i.e., in cases of newly established, strengthened or weakened social ties. In this paper, we focus on a dynamic non-stationary social network and design a randomized algorithm, RSB, based on multi-armed bandit optimization, to maximize influence propagation over time. The algorithm produces a sequence of online decisions and calibrates its explore-exploit strategy utilizing outcomes of previous decisions. It is rigorously proven to achieve an upper-bounded regret in reward and applicable to large-scale social networks. Practical effectiveness of the algorithm is evaluated using real-world datasets, which demonstrates that our algorithm outperforms previous stationary methods under non-stationary conditions.
Yixin Bao, Zhi Wang 0001, Chuan Wu 0001, Francis C. M. Lau 0001
IWQoS4
2016 More is Better? Measurement of MPTCP Based Cellular Bandwidth Aggregation in the Wild
abstract
4G/3G Networks have been widely deployed around the world to provide high wireless bandwidth for mobile users. However, the achievable 3G/4G bandwidth is still much lower than their theoretic maximum. Signal strengths and available backhaul capacities may vary significantly at different locations and times, often leading to unsatisfactory performance. Band-width aggregation, which uses multiple interfaces concurrently for data transfer, is a readily deployable solution. Specifically, Multi-Path TCP (MPTCP) has been advocated as a promising approach for leveraging multiple source-destination paths simultaneously in the transport layer. In this paper, we investigate the efficiency of an MPTCP-based bandwidth aggregation frame-work based on extensive measurements. In particular, we evaluate the gain for bandwidth aggregation across up to 4 cellular operators' networks, with respect to factors such as time, user location, data size, aggregation proxy location and congestion control algorithm. Our measurement studies reveal that (1) bandwidth aggregation in general improves the cellular network bandwidth experienced by mobile users, but the performance gain is significant only for bandwidth-intensive delay-tolerant flows, (2) the effectiveness of aggregation depends on many network factors, including QoS of individual cellular interfaces and the location of aggregation proxy, (3) contextual factors, including the time of day and the mobility of a user, also affect the aggregation performance.
Zhixiong Niu, Zhi Wang 0001, Hong Xu 0001, Chuan Wu 0001, Francis C. M. Lau 0001
MASS4
2016 Colocation Demand Response: Joint Online Mechanisms for Individual Utility and Social Welfare Maximization
abstract
Data centers with high yet elastic energy demand are ideal candidates for participation in demand response programs. This paper studies emergency demand response (EDR) at multi-tenant colocation data centers (colocations). While the colocation has no direct control over tenants' servers, we design online mechanisms to incentivize and coordinate tenants' energy reduction. Our mechanism is online in nature, aiming to maximize not only social welfare but also tenant utility. Our main proposal is a truthful incentive auction that provides tenants with monetary remuneration for EDR energy reduction, minimizing social cost, which combines seamlessly with an online primal-dual framework for each tenant to schedule their delay-tolerant workloads. The online optimization at each tenant targets its utility maximization, concurrently reporting valuation functions for the tenant to participate in the auction. Our online algorithms achieve long-term performance guarantees in both tenants' utility and social welfare maximization, while fulfilling the EDR requirement with minimal diesel generation. We validate the efficiency of our algorithms through both the theoretical analysis and real-world trace-driven simulations.
Chuan Wu 0001, Zongpeng Li, Shaolei Ren
IEEE J. Sel. Areas Commun.2
2016 Virtual Machine Trading in a Federation of Clouds: Individual Profit and Social Welfare Maximization
abstract
By sharing resources among different cloud providers, the paradigm of federated clouds exploits temporal availability of resources and geographical diversity of operational costs for efficient job service. While interoperability issues across different cloud platforms in a cloud federation have been extensively studied, fundamental questions on cloud economics remain: When and how should a cloud trade resources (e.g., virtual machines) with others, such that its net profit is maximized over the long run, while a close-to-optimal social welfare in the entire federation can also be guaranteed? To answer this question, a number of important, interrelated decisions, including job scheduling, server provisioning, and resource pricing, should be dynamically and jointly made, while the long-term profit optimality is pursued. In this work, we design efficient algorithms for intercloud virtual machine (VM) trading and scheduling in a cloud federation. For VM transactions among clouds, we design a double-auction-based mechanism that is strategy-proof, individual-rational, ex-post budget-balanced, and efficient to execute over time. Closely combined with the auction mechanism is a dynamic VM trading and scheduling algorithm, which carefully decides the true valuations of VMs in the auction, optimally schedules stochastic job arrivals with different service level agreements (SLAs) onto the VMs, and judiciously turns on and off servers based on the current electricity prices. Through rigorous analysis, we show that each individual cloud, by carrying out the dynamic algorithm in the online double auction, can achieve a time-averaged profit arbitrarily close to the offline optimum. Asymptotic optimality in social welfare is also achieved under homogeneous cloud settings. We carry out simulations to verify the effectiveness of our algorithms, and examine the achievable social welfare under heterogeneous cloud settings, as driven by the real-world Google cluster usage traces.
Hongxing Li 0002, Chuan Wu 0001, Zongpeng Li, Francis C. M. Lau 0001
IEEE/ACM Trans. Netw.2
2016 An Online Auction Framework for Dynamic Resource Provisioning in Cloud Computing
abstract
Auction mechanisms have recently attracted substantial attention as an efficient approach to pricing and allocating resources in cloud computing. This work, to the authors' knowledge, represents the first online combinatorial auction designed for the cloud computing paradigm, which is general and expressive enough to both: 1) optimize system efficiency across the temporal domain instead of at an isolated time point; and 2) model dynamic provisioning of heterogeneous virtual machine (VM) types in practice. The final result is an online auction framework that is truthful, computationally efficient, and guarantees a competitive ratio ≈ 3.30 in social welfare in typical scenarios. The framework consists of three main steps: 1) a tailored primal-dual algorithm that decomposes the long-term optimization into a series of independent one-shot optimization problems, with a small additive loss in competitive ratio; 2) a randomized subframework that applies primal-dual optimization for translating a centralized cooperative social welfare approximation algorithm into an auction mechanism, retaining the competitive ratio while adding truthfulness; and 3) a primal-dual algorithm for approximating the one-shot optimization with a ratio close to e. We also propose two extensions: 1) a binary search algorithm that improves the average-case performance; 2) an improvement to the online auction framework when a minimum budget spending fraction is guaranteed, which produces a better competitive ratio. The efficacy of the online auction framework is validated through theoretical analysis and trace-driven simulation studies. We are also in the hope that the framework can be instructive in auction design for other related problems.
Linquan Zhang, Chuan Wu 0001, Zongpeng Li, Francis C. M. Lau 0001
IEEE/ACM Trans. Netw.3
2016 The Streaming Capacity of Sparsely Connected P2P Systems With Distributed Control
abstract
Peer-to-peer (P2P) streaming technologies can take advantage of the upload capacity of clients, and hence can scale to large content distribution networks with lower cost. A fundamental question for P2P streaming systems is the maximum streaming rate that all users can sustain. Prior works have studied the optimal streaming rate for a complete network, where every peer is assumed to be able to communicate with all other peers. This is, however, an impractical assumption in real systems. In this paper, we are interested in the achievable streaming rate when each peer can only connect to a small number of neighbors. We show that even with a random peer-selection algorithm and uniform rate allocation, as long as each peer maintains Ω(logN) downstream neighbors, where N is the total number of peers in the system, the system can asymptotically achieve a streaming rate that is close to the optimal streaming rate of a complete network. These results reveal a number of important insights into the dynamics of the system, based on which we then design simple improved algorithms that can reduce the constant factor in front of the Ω(logN) term, yet can achieve the same level of performance guarantee. Simulation results are provided to verify our analysis.
Can Zhao 0006, Xiaojun Lin 0001, Chuan Wu 0001
IEEE/ACM Trans. Netw.3
2016 Capacity of P2P On-Demand Streaming With Simple, Robust, and Decentralized Control
abstract
The performance of large-scale peer-to-peer (P2P) video-on-demand (VoD) streaming systems can be very challenging to analyze due to sparse connectivity and complex, random dynamics. Specifically, in practical P2P VoD systems, each peer only interacts with a small number of other peers/neighbors. Furthermore, its upload capacity, downloading position, and content availability change dynamically and randomly. In this paper, we rigorously study large-scale P2P VoD systems with sparse connectivity among peers and investigate simple and decentralized P2P control strategies that can provably achieve close-to-optimal streaming capacity. We first focus on a single streaming channel. Using a simple algorithm that assigns each peer a random set of Θ(logN) neighbors and allocates upload capacity uniformly, we show that a close-to-optimal streaming rate can be asymptotically achieved for all peers with high probability as the number of peers N increases. Furthermore, the tracker does not need to obtain detailed knowledge of which chunks each peer caches, and hence incurs low overhead. We then study multiple streaming channels where peers watching one channel may help peers in another channel with insufficient upload bandwidth. We propose a simple random cache-placement strategy and show that a close-to-optimal streaming capacity region for all channels can be attained with high probability, again with only Θ(logN) per-peer neighbors. These results provide important insights into the dynamics of large-scale P2P VoD systems, which will be useful for guiding the design of improved P2P control protocols.
Can Zhao 0006, Jian Zhao 0008, Xiaojun Lin 0001, Chuan Wu 0001
IEEE/ACM Trans. Netw.4
2015 Cost-Minimizing Online VM Purchasing for Application Service Providers with Arbitrary Demands
abstract
Recent years witness the proliferation of Infrastructure-as-a-Service (IaaS) cloud services, which provide on-demand resources (CPU, RAM, disk) in the form of virtual machines (VMs) for hosting applications/services of third parties. Given the state-of-the-art IaaS offerings, it is still a problem of fundamental importance how the Application Service Providers (ASPs) should rent VMs from the clouds to serve their application needs, in order to minimize the cost while meeting their job demands over a long run. Cloud providers offer different pricing options to meet computing requirements of a variety of applications. However, the challenge facing an ASP is how these pricing options can be dynamically combined to serve arbitrary demands at the optimal cost. In this paper, we propose an online VM purchasing algorithm based on the Lyapunov optimization technique, for minimizing the long-term-averaged VM rental cost of an ASP with time-varying and delay-tolerant workloads, while bounding the maximum response delay of its jobs. In stark contrast with the existing studies, the proposed algorithm enables an ASP to optimally decide the amount of reserved, on-demand and spot instances to purchase simultaneously. Rigorous analysis shows that our algorithm can achieve a time-averaged resource cost close to the offline optimum. Trace-driven simulations further verify the efficacy of our algorithm.
Shengkai Shi, Chuan Wu 0001, Zongpeng Li
CLOUD2
2015 Hierarchical Virtual Machine Placement in Modular Data Centers
abstract
This work studies how to minimize communication cost for placing Virtual Machines (VMs) in a modular data center. We consider a number of cooperative VMs implementing the same job, with known inter-VM communication patterns. The modular data center has a two-layer network structure, where computing pods constitute basic building blocks and are connected by a core network. At the core network layer, we design spectral clustering algorithms to partition VMs into computing pods, minimizing inter-pod communication cost. We then further apply an SDP relaxation approach to decide the VM placement within each computing pod, targeting both load balancing among physical servers and inter-server communication cost minimization. Extensive simulations are conducted to validate the efficacy of the proposed hierarchical VM placement scheme.
Linquan Zhang, Xunrui Yin, Zongpeng Li, Chuan Wu 0001
CLOUD4
2015 Responsive multipath TCP in SDN-based datacenters
abstract
A basic need in datacenter networks is to provide high throughput for large flows such as the massive shuffle traffic flows in a MapReduce application. Multipath TCP (MPTCP) has been investigated as an effective approach toward this goal, by spreading one TCP flow onto multiple paths. However, the current MPTCP implementation has two major limitations: (1) a fixed number of subflows are used without reacting to the actual traffic condition; (2) the routing of subflows of a multipath TCP connection relies heavily on the ECMP-based random hashing. The former may lead to a waste of both the server and network resources, while the latter can cause throughput degradation when multiple subflows collide on the same path. This paper proposes a responsive MPTCP system to resolve the two limitations simultaneously. Our system employs a centralized controller for intelligent subflow route calculation and a monitor running on each server for actively adjusting the number of subflows. Working in synergy, the two modules enable MPTCP flows to respond to the traffic conditions and pursue high throughput on the fly, at very low computation and messaging overhead. NS3-based experiments show that our system achieves satisfactory throughput with less resource overhead, or better throughput at similar amounts of overhead, as compared to common alternatives.
Jingpu Duan, Zhi Wang 0001, Chuan Wu 0001
ICC3
2015 Cost-Effective Low-Delay Cloud Video Conferencing
abstract
The cloud computing paradigm has been advocated in recent video conferencing system design, which exploits the rich on-demand resources spanning multiple geographic regions of a distributed cloud, for better conferencing experience. A typical architectural design in cloud environment is to create video conferencing agents, i.e., Virtual machines, in each cloud site, assign users to the agents, and enable inter-user communication through the agents. Given the diversity of devices and network connectivities of the users, the agents may also transcode the conferencing streams to the best formats and bitrates. In this architecture, two key issues exist on how to effectively assign users to agents and how to identify the best agent to perform a Transco ding task, which are nontrivial due to the following: (1) the existing proximity-based assignment may not be optimal in terms of inter-user delay, which fails to consider the whereabouts of the other users in a conferencing session, (2) the agents may have heterogeneous bandwidth and processing availability, such that the best Transco ding agents should be carefully identified, for cost minimization while best serving all the users requiring the transcoded streams. To address these challenges, we formulate the user-to-agent assignment and Transco ding-agent selection problems, which targets at minimizing the operational cost of the conferencing provider while keeping the conferencing delay low. The optimization problem is combinatorial in nature and difficult to solve. Using Markov approximation framework, we design a decentralized algorithm that provably converges to a bounded neighborhood of the optimal solution. An agent ranking scheme is also proposed to properly initialize our algorithm so as to improve its convergence. The results from a prototype system implementation show that our design in a set of Internet-scale scenarios reduces the operational cost by 77% as compared to a commonly-adopted alternative, while simultaneously yielding lower conferencing delays.
Mohammad Hajiesmaili, Lok To Mak, Zhi Wang 0001, Chuan Wu 0001, Minghua Chen 0001, Ahmad Khonsari
ICDCS4
2015 Socially-optimal online spectrum auctions for secondary wireless communication
abstract
Spectrum auctions are efficient mechanisms for licensed users to relinquish their under-utilized spectrum to secondary links for monetary remuneration. Truthfulness and social welfare maximization are two natural goals in such auctions, but cannot be achieved simultaneously with polynomial-time complexity by existing methods, even in a static network with fixed parameters. The challenge escalates in practical systems with QoS requirements and volatile traffic demands for secondary communication. Online, dynamic decisions are required for rate control, channel evaluation/bidding, and packet dropping at each secondary link, as well as for winner determination and pricing at the primary user. This work proposes an online spectrum auction framework with cross-layer decision making and randomized winner determination on the fly. The framework is truthful-in-expectation, and achieves close-to-offline-optimal time-averaged social welfare and individual utilities with polynomial time complexity. A new method is introduced for online channel evaluation in a stochastic setting. Simulation studies further verify the efficacy of the proposed auction in practical scenarios.
Hongxing Li 0002, Chuan Wu 0001, Zongpeng Li
INFOCOM2
2015 A truthful incentive mechanism for emergency demand response in colocation data centers
abstract
Data centers are key participants in demand response programs, including emergency demand response (EDR), where the grid coordinates large electricity consumers for demand reduction in emergency situations to prevent major economic losses. While existing literature concentrates on owner-operated data centers, this work studies EDR in multi-tenant colocation data centers where servers are owned and managed by individual tenants. EDR in colocation data centers is significantly more challenging, due to lack of incentives to reduce energy consumption by tenants who control their servers and are typically on fixed power contracts with the colocation operator. Consequently, to achieve demand reduction goals set by the EDR program, the operator has to rely on the highly expensive and/or environmentally-unfriendly on-site energy backup/generation. To reduce cost and environmental impact, an efficient incentive mechanism is therefore in need, motivating tenants' voluntary energy reduction in case of EDR. This work proposes a novel incentive mechanism, Truth-DR, which leverages a reverse auction to provide monetary remuneration to tenants according to their agreed energy reduction. Truth-DR is computationally efficient, truthful, and achieves 2-approximation in colocation-wide social cost. Trace-driven simulations verify the efficacy of the proposed auction mechanism.
Linquan Zhang, Shaolei Ren, Chuan Wu 0001, Zongpeng Li
INFOCOM3
2015 A truthful (1-ε)-optimal mechanism for on-demand cloud resource provisioning
abstract
On-demand resource provisioning in cloud computing provides tailor-made resource packages (typically in the form of VMs) to meet users' demands. Public clouds nowadays provide more and more elaborated types of VMs, but have yet to offer the most flexible dynamic VM assembly, which is partly due to the lack of a mature mechanism for pricing tailor-made VMs on the spot. This work proposes an efficient randomized auction mechanism based on a novel application of smoothed analysis and randomized reduction, for dynamic VM provisioning and pricing in geo-distributed cloud data centers. This auction, to the best of our knowledge, is the first one in literature that achieves (i) truthfulness in expectation, (ii) polynomial running time in expectation, and (iii) (1 − ε)-optimal social welfare in expectation for resource allocation, where e can be arbitrarily close to 0. Our mechanism consists of three modules: (1) an exact algorithm to solve the NP-hard social welfare maximization problem, which runs in polynomial time in expectation, (2) a perturbation-based randomized resource allocation scheme which produces a VM provisioning solution that is (1 − ε)-optimal and (3) an auction mechanism that applies the perturbation-based scheme for dynamic VM provisioning and prices the customized VMs using a randomized VCG payment, with a guarantee in truthfulness in expectation. We validate the efficacy of the mechanism through careful theoretical analysis and trace-driven simulations.1
Xiaoxi Zhang 0001, Chuan Wu 0001, Zongpeng Li, Francis C. M. Lau 0001
INFOCOM2
2015 An online procurement auction for power demand response in storage-assisted smart grids
abstract
The quintessential problem in a smart grid is the matching between power supply and demand - a perfect balance across the temporal domain, for the stable operation of the power network. Recent studies have revealed the critical role of electricity storage devices, as exemplified by rechargeable batteries and plug-in electric vehicles (PEVs), in helping achieve the balance through power arbitrage. Such potential from batteries and PEVs can not be fully realized without an appropriate economic mechanism that incentivizes energy discharging at times when supply is tight. This work aims at a systematic study of such demand response problem in storage-assisted smart grids through a well-designed online procurement auction mechanism. The long-term social welfare maximization problem is naturally formulated into a linear integer program. We first apply a primal-dual optimization algorithm to decompose the online auction design problem into a series of one-round auction design problems, achieving a small loss in competitive ratio. For the one round auction, we show that social welfare maximization is still NP-hard, and design a primal-dual approximation algorithm that works in concert with the decomposition algorithm. The end result is a truthful power procurement auction that is online, truthful, and 2-competitive in typical scenarios.
Ruiting Zhou, Zongpeng Li, Chuan Wu 0001
INFOCOM3
2015 Fair rewarding in colocation data centers: Truthful mechanism for emergency demand response
abstract
Reducing servers' power usage in data centers upon utility's request has been emerging as a valuable demand response resource for enhancing power grid's efficiency and reliability, especially during emergency events (e.g., extreme weather) that result in electricity production shortage and put the grid in jeopardy. Nonetheless, for demand response in multi-tenant colocation data centers, operators may have to leverage expensive and environmentally-unfriendly diesel generation, because individual tenants manage their own servers' power usage without coordination and are typically charged by data center operators based on fixed power contracts that provide no incentives for demand response. This paper focuses on emergency demand response (EDR) and proposes an auction-based incentive mechanism, called FairDR, that incentivizes and coordinates tenants' energy reduction through financial rewards for enabling cost-effective and low-carbon EDR in colocation data center. FairDR decides tenants' energy reduction online without knowing a priori the future energy reduction requirements. It is proved that FairDR ensures tenants' truthfulness in the auction process, attains a bounded overall cost saving compared to the offline optimum which knows all the demands, and guarantees fairness (i.e., similar rewards are offered if tenants reduce the same amount of energy) that is largely absent in the existing auction mechanisms. Finally, trace-driven simulations are performed to validate our analysis and demonstrate that FairDR outperforms the existing mechanisms by improving fairness and achieving a good cost saving that is comparable to the offline optimum.
Chuan Wu 0001, Shaolei Ren, Zongpeng Li
IWQoS2
2015 Online cost minimization for operating geo-distributed cloud CDNs
abstract
Cloud-based content delivery networks (Cloud CDN) cache and deliver contents from geo-distributed cloud data centers to end users across the globe, exploiting "infinite" on-demand cloud resources to address volatile user demands. It is critically important to efficiently manage cloud resources in different locations over time, for minimization of the operational cost of the CDN provider, while delivering short response delay to user requests. Although many have studied cost-aware replica placement and request redirection in CDN systems, most are restricted to an offline or one-time setting, or resort to greedy heuristics for online operation. This work proposes an efficient online algorithm for dynamic content replication and request dispatching in cloud CDNs operating over a long time span, targeting overall cost minimization with performance guarantees. Our online algorithm consists of two main modules: (1) a regularization method from the online learning literature to convert the offline cost-minimization optimization problem into a sequence of regularized problems, each to be efficiently solvable in one time slot; (2) a randomized approach to convert the optimal fractional solutions from the regularized problems to integer solutions of the original problem, achieving a good competitive ratio. The effectiveness of our online algorithm is validated through solid theoretical analysis and trace-driven simulations.
Xiaoxi Zhang 0001, Chuan Wu 0001, Zongpeng Li, Francis C. M. Lau 0001
IWQoS2
2015 Software Defined Mobile Multicast
abstract
Mobile multicast has been deployed in telecommunication networks for information dissemination applications such as IPTV and video conferencing. Recent studies of mobile multicast focused on fast handover protocols, and algorithms for multicast tree management have witnessed little improvement over the years. Shortest path trees represent the status quo of multicast topology in real-world systems. Steiner trees were investigated extensively in the theory community and are known to be bandwidth efficient, but come with an associated complexity. Recent developments in the Software Defined Networking (SDN) paradigm have shed light on implementing more sophisticated protocols for better routing performance. We propose an SDN-based design to combat the complexity vs. Performance dilemma in mobile multicast. We construct low-cost Steiner trees for multicastin a mobile network, employing an SDN controller for coordinating tree construction and morphing. Highlights of our design include a set of efficient online algorithms for tree adjustment when nodes arrive and depart on the fly, and an SDN rule update framework based on constraints expressed by boolean logic to ensure loop free rule updates. The algorithms are proven to achieve a constant competitive ratio against the offline optimal Steiner tree, with an amortized constant number of edge swaps per adjustment. Mininet-based implementation and evaluation further validate the efficacy of our design.
Shunyi Xu, Chuan Wu 0001, Zongpeng Li
MASS2
2015 Online Auctions in IaaS Clouds: Welfare and Profit Maximization with Server Costs
abstract
Auction design has recently been studied for dynamic resource bundling and VM provisioning in IaaS clouds, but is mostly restricted to the one-shot or offline setting. This work targets a more realistic case of online VM auction design, where: (i) cloud users bid for resources into the future to assemble customized VMs with desired occupation durations; (ii) the cloud provider dynamically packs multiple types of resources on heterogeneous physical machines (servers) into the requested VMs; (iii) the operational costs of servers are considered in resource allocation; (iv) both social welfare and the cloud provider's net profit are to be maximized over the system running span. We design truthful, polynomial time auctions to achieve social welfare maximization and/or the provider's profit maximization with good competitive ratios. Our mechanisms consist of two main modules: (1) an online primal-dual optimization framework for VM allocation to maximize the social welfare with server costs, and for revealing the payments through the dual variables to guarantee truthfulness; and (2) a randomized reduction algorithm to convert the social welfare maximizing auctions to ones that provide a maximal expected profit for the provider, with competitive ratios comparable to those for social welfare. We adopt a new application of Fenchel duality in our primal-dual framework, which provides richer structures for convex programs than the commonly used Lagrangian duality, and our optimization framework is general and expressive enough to handle various convex server cost functions. The efficacy of the online auctions is validated through careful theoretical analysis and trace-driven simulation studies.
Xiaoxi Zhang 0001, Zhiyi Huang 0002, Chuan Wu 0001, Zongpeng Li, Francis C. M. Lau 0001
SIGMETRICS3
2015 Online Electricity Cost Saving Algorithms for Co-Location Data Centers
abstract
This work studies the online electricity cost minimization problem at a co-location data center. A co-location data center serves multiple tenants who rent the physical infrastructure within the data center to run their respective cloud computing services. Consequently, the co-location operator has no direct control over power consumption of its tenants, and an efficient mechanism is desired for eliciting desirable consumption patterns from the co-location tenants. Electricity billing faced by a data center is nowadays based on both the total volume consumed and the peak consumption rate. This leads to an interesting new combinatorial optimization structure on the electricity cost optimization problem, which also exhibits an online nature due to the definition of peak consumption. We model and solve the problem through two approaches: the pricing approach and the auction approach. For the former, we design an offline 2-approximation algorithm as well as an online algorithm with a small competitive ratio in most practical settings. For the latter, we design an efficient (2+c)-competitive online algorithm, where c is a system dependent parameter close to 1.49, and then convert it into an efficient mechanism that executes in an online fashion, runs in polynomial time, and guarantees truthful bidding and (2+2c)-competitive in social cost.
Linquan Zhang, Zongpeng Li, Chuan Wu 0001, Shaolei Ren
SIGMETRICS3
2015 Online Electricity Cost Saving Algorithms for Co-Location Data Centers
abstract
This work studies the online electricity cost minimization problem at a co-location data center, which serves multiple tenants who rent the physical infrastructure within the data center to run their respective cloud computing services. The co-location operator has no direct control over power consumption of its tenants, and an efficient mechanism is desired for eliciting desirable consumption patterns from the tenants. Electricity billing faced by a data center is nowadays based on both the total volume consumed and the peak consumption rate. This leads to an interesting new combinatorial optimization structure on the electricity cost optimization problem, which also exhibits an online nature due to the definition of peak consumption. We model and solve the problem through two approaches: the pricing approach and the auction approach, and design online algorithms with small competitive ratios.
Linquan Zhang, Zongpeng Li, Chuan Wu 0001, Shaolei Ren
IEEE J. Sel. Areas Commun.3
2015 Demand Response in Smart Grids: A Randomized Auction Approach
abstract
The smart grid is a modern power grid that achieves high efficiency and robustness through sophisticated information and communications technology. Demand response has great potential in helping balance demand and supply in a smart grid, cutting generation cost and carbon footprint, and improving system stability. Auctions represent a natural and efficient approach for carrying out demand response between the power grid and large electricity users, microgrids, and electricity storage devices. This work explores the modeling and design space of demand response auctions, targeting expressive power, truthful information revelation, computational efficiency, and economic efficiency. We present a randomized auction that explores the underlying problem structure of demand response, and prove that it is truthful, runs in polynomial time, and achieves (1 + ϵ)-optimal social cost for an arbitrarily small constant ϵ. The key technique lies in the marriage of smoothed analysis and randomized reduction, which makes its debut in this work among literature on mechanism design, and can be applied to problems where social welfare optimization is NP-hard but admits a smoothed polynomial-time algorithm.
Ruiting Zhou, Zongpeng Li, Chuan Wu 0001, Minghua Chen 0001
IEEE J. Sel. Areas Commun.3
2015 Locality-aware streaming in hybrid P2P-cloud CDN systems
Jian Zhao 0008, Chuan Wu 0001, Xiaojun Lin 0001
Peer-to-Peer Netw. Appl.2
2015 Introduction to the Special Section on Visual Computing in the Cloud: Fundamentals and Applications
abstract
Cloud computing involves a large number of terminals connected through a real-time high-speed network (such as the Internet). The adoption rates for private and hybrid cloud services increased to 40% in 2013, with computing shifting from on-premise infrastructure to the cloud. To keep pace with the ever-accelerating rate of innovation, companies are moving to the cloud. However, visual computing in the cloud brings great challenges, such as how to measure and then improve the quality of experience in cloud computing. This Special Section provides the image/video community a forum to present new academic research and industrial development in running visual computing services in the cloud. This Special Section aims to address fundamental and practical aspects of visual computing in the cloud, such as how to build cloud platforms that can cope with seemingly unlimited supply of content coming from traditional media sources as well as new media uploaded to the Internet (YouTube, Facebook, etc.); how to leverage cloud technology to build high-quality image/video browsing and delivery experiences for a global audience; how to ingest, encode, process, adapt, as well as protect contents and privacy of users; how to provide both on-demand and live-streaming capabilities; how to tag image/video and allow consumers to access the image/video contents with high availability; how to support image/video services in mobile devices; and how to perform real-time image/video analytics in the cloud, to mention a few among a diverse range of challenges.
Jiangchuan Liu, Wenwu Zhu 0001, Touradj Ebrahimi, John G. Apostolopoulos, Xian-Sheng Hua 0001, Chuan Wu 0001
IEEE Trans. Circuits Syst. Video Technol.6
2015 A Joint Online Transcoding and Delivery Approach for Dynamic Adaptive Streaming
abstract
Dynamic adaptive streaming has emerged as a popular approach for video services in today's Internet. To date, the two important components in dynamic adaptive streaming, video transcoding that generates the adaptive bitrates of a video and video delivery that streams the videos to users, have been separately studied, resulting in a huge waste of computation and storage resource due to producing and caching different versions of videos regardless of their demands. We conduct extensive measurement studies of video sharing systems, including an IPTV service which streams regular, professionally made videos and an instant video clip sharing service which provides extremely short user-generated videos, as well as the availability of computation resource in conventional content delivery networks (CDNs). Based on the measurement insights, we propose an online joint transcoding and delivery approach for adaptive video streaming. We formulate optimization problems to enable high streaming quality for the users, and low computation and replication costs for the system. In particular, our strategy connects video transcoding and video delivery based on users' preferences of CDN regions and regional preferences of video versions. We analyze hardness of these problems and design distributed solutions. Extensive trace-driven experiments further demonstrate the superiority of our design.
Zhi Wang 0001, Lifeng Sun, Chuan Wu 0001, Wenwu Zhu 0001, Qidong Zhuang, Shiqiang Yang
IEEE Trans. Multim.3
2015 Scaling Social Media Applications Into Geo-Distributed Clouds
abstract
Federation of geo-distributed cloud services is a trend in cloud computing that, by spanning multiple data centers at different geographical locations, can provide a cloud platform with much larger capacities. Such a geo-distributed cloud is ideal for supporting large-scale social media applications with dynamic contents and demands. Although promising, its realization presents challenges on how to efficiently store and migrate contents among different cloud sites and how to distribute user requests to the appropriate sites for timely responses at modest costs. These challenges escalate when we consider the persistently increasing contents and volatile user behaviors in a social media application. By exploiting social influences among users, this paper proposes efficient proactive algorithms for dynamic, optimal scaling of a social media application in a geo-distributed cloud. Our key contribution is an online content migration and request distribution algorithm with the following features: 1) future demand prediction by novelly characterizing social influences among the users in a simple but effective epidemic model; 2) one-shot optimal content migration and request distribution based on efficient optimization algorithms to address the predicted demand; and 3) a Δ(t)-step look-ahead mechanism to adjust the one-shot optimization results toward the offline optimum. We verify the effectiveness of our online algorithm by solid theoretical analysis, as well as thorough comparisons to ready algorithms including the ideal offline optimum, using large-scale experiments with dynamic realistic settings on Amazon Elastic Compute Cloud (EC2).
Yu Wu 0010, Chuan Wu 0001, Bo Li 0001, Linquan Zhang, Zongpeng Li, Francis C. M. Lau 0001
IEEE/ACM Trans. Netw.2
2015 Cost-Minimizing Dynamic Migration of Content Distribution Services into Hybrid Clouds
abstract
With the recent advent of cloud computing technologies, a growing number of content distribution applications are contemplating a switch to cloud-based services, for better scalability and lower cost. Two key tasks are involved for such a move: to migrate the contents to cloud storage, and to distribute the Web service load to cloud-based Web services. The main issue is to best utilize the cloud as well as the application provider's existing private cloud, to serve volatile requests with service response time guarantee at all times, while incurring the minimum operational cost. While it may not be too difficult to design a simple heuristic, proposing one with guaranteed cost optimality over a long run of the system constitutes an intimidating challenge. Employing Lyapunov optimization techniques, we design a dynamic control algorithm to optimally place contents and dispatch requests in a hybrid cloud infrastructure spanning geo-distributed data centers, which minimizes overall operational cost overtime, subject to service response time constraints. Rigorous analysis shows that the algorithm nicely bounds the response times within the preset QoS target, and guarantees that the overall cost is within a small constant gap from the optimum achieved by a T-slot lookahead mechanism with known future information. We verify the performance of our dynamic algorithm with prototype-based evaluation.
Xuanjia Qiu, Hongxing Li 0002, Chuan Wu 0001, Zongpeng Li, Francis C. M. Lau 0001
IEEE Trans. Parallel Distributed Syst.3
2015 Enhancing Internet-Scale Video Service Deployment Using Microblog-Based Prediction
abstract
Online microblogging has been very popular in today's Internet, where users follow other people they are interested in and exchange information between themselves. Among these exchanges, video links are a representative type on a microblogging site. The impact is fundamental-not only are viewers in a video service directly coming from the microblog sharing and recommendation, but also are the users in the microblogging site representing a promising sample to all the viewers. It is intriguing to study a proactive service deployment for such videos, using the propagation patterns of microblogs. Based on extensive traces from Youku and Tencent Weibo, a popular video sharing site and a favored microblogging system, we explore how video propagation patterns in the microblogging system are correlated with video popularity on the video sharing site. Using influential factors summarized from the measurement studies, we further design a neural network-based learning framework to predict the number of potential viewers and their geographic distribution. We then design proactive video deployment algorithms based on the prediction framework, which not only determines the upload capacities of servers in different regions, but also strategically replicates videos to these regions to serve users. Our PlanetLab-based experiments verify the effectiveness of our design.
Zhi Wang 0001, Lifeng Sun, Chuan Wu 0001, Shiqiang Yang
IEEE Trans. Parallel Distributed Syst.3
2014 Core-Selecting Auctions for Dynamically Allocating Heterogeneous VMs in Cloud Computing
abstract
In a cloud market, the cloud provider provisions heterogeneous virtual machine (VM) instances from its resource pool, for allocation to cloud users. Auction-based allocations are efficient in assigning VMs to users who value them the most. Existing auction design often overlooks the heterogeneity of VMs, and does not consider dynamic, demand-driven VM provisioning. Moreover, the classic VCG auction leads to unsatisfactory seller revenues and vulnerability to a strategic bidding behavior known as shill bidding. This work presents a new type of core-selecting VM auctions, which are combinatorial auctions that always select bidder charges from the core of the price vector space, with guaranteed economic efficiency under truthful bidding. These auctions represent a comprehensive three-phase mechanism that instructs the cloud provider to judiciously assemble, allocate, and price VM bundles. They are proof against shills, can improve seller revenue over existing auction mechanisms, and can be tailored to maximize truthfulness.
Haoming Fu, Zongpeng Li, Chuan Wu 0001, Xiaowen Chu 0001
IEEE CLOUD3
2014 Federated Private Clouds via Broker's Marketplace: A Stackelberg-Game Perspective
abstract
More and more enterprises have set up their own private clouds by applying virtualization to their data centers, the benefit is flexible resource supply to different internal demands. Aiming to meet the peak demand in their resource provisioning, private clouds are often under-utilized. A new paradigm has emerged that advocates leasing the spare resources to external users, when and if adequate rental prices are offered. A broker is typically employed which pools the spare resources of multiple private clouds together and leases them to serve external users' jobs. Good mechanisms have yet to be derived for the broker to set the offered prices to buy spare resources from the private clouds, and to schedule jobs on the available resources, such that the economic benefits of both the broker and the private clouds are maximized. The design of the mechanism is especially challenging when we consider the dynamic arrival of users' jobs and volatile availability of spare resources at the private clouds, while aiming at long-term profit optimality. In this paper, we model the interaction between the broker and the private clouds as a two-stage Stackelberg game. As the leader in the game, the broker decides and offers prices for renting VMs of different types from each private cloud. As a follower, each private cloud responds with the number of VMs of each type that it is willing to lease. By combining with the Stackelberg game model we design online algorithms for the broker to set the prices and schedule jobs on the private clouds, and for the private cloud to decide the numbers of VMs to lease, based on the Lyapunov optimization theory. We prove that the broker achieves a time-averaged profit that is close to the offline optimum with complete information on future job arrivals and resource availability, while each private cloud makes their best earning. The proposed online algorithm is carefully evaluated based on usage traces of Google cluster and Amazon EC2.
Xuanjia Qiu, Chuan Wu 0001, Hongxing Li 0002, Zongpeng Li, Francis C. M. Lau 0001
IEEE CLOUD2
2014 Characterizing cascade dynamics in a microblogging system
abstract
Online microblogging sites have become increasingly important platforms for information diffusion in today's world, where users post short messages and follow various messages posted by people that they are interested in. It is intriguing to qualitatively study the temporal dynamics of an information cascade in a microblogging system, in terms of the number of users influenced at any given time, which may provide valuable input to facilitate emerging applications such as online advertising and content distribution. In this paper, we model information diffusion in a microblogging network as an age-dependent branching process, based on practical observations from Tencent Weibo, a popular microblogging site in China. This model enables careful characterization of the diffusion topology, the different delays for users to respond to new information, and the evolution of the size of the information cascade over time. We derive the expected cascade size at any time. We validate our model based on Tencent Weibo traces, and demonstrate its effectiveness in capturing information diffusion dynamics in the real world.
Shengkai Shi, Zhi Wang 0001, Chuan Wu 0001, Xiaojun Lin 0001
ICC3
2014 Joint online transcoding and geo-distributed delivery for dynamic adaptive streaming
abstract
Dynamic adaptive video streaming has emerged as a popular approach for video streaming in today's Internet. To date the two important components in dynamic adaptive streaming, video transcoding which generates the adaptive bitrates of a video and video delivery which streams the videos to users, have been separately studied, resulting in a huge waste of computation and storage resource due to transcoding useless videos and suboptimal streaming quality due to homogeneous video replication. In this paper, we propose to jointly perform video transcoding and video delivery for adaptive streaming in an online manner. We conduct extensive measurement studies of a video sharing system and a CDN to motivate our design. We formulate and solve optimization problems to enable high streaming quality for the users, and low computation and replication costs for the system. In particular, our design connects video transcoding and video delivery based on users' preferences of CDN regions and regional preferences of video versions. Extensive trace-driven experiments further confirm the superiority of our design.
Zhi Wang 0001, Lifeng Sun, Chuan Wu 0001, Wenwu Zhu 0001, Shiqiang Yang
INFOCOM3
2014 Dynamic resource provisioning in cloud computing: A randomized auction approach
abstract
This work studies resource allocation in a cloud market through the auction of Virtual Machine (VM) instances. It generalizes the existing literature by introducing combinatorial auctions of heterogeneous VMs, and models dynamic VM provisioning. Social welfare maximization under dynamic resource provisioning is proven NP-hard, and modeled with a linear integer program. An efficient α-approximation algorithm is designed, with α ~ 2.72 in typical scenarios. We then employ this algorithm as a building block for designing a randomized combinatorial auction that is computationally efficient, truthful in expectation, and guarantees the same social welfare approximation factor α. A key technique in the design is to utilize a pair of tailored primal and dual LPs for exploiting the underlying packing structure of the social welfare maximization problem, to decompose its fractional solution into a convex combination of integral solutions. Empirical studies driven by Google Cluster traces verify the efficacy of the randomized auction.
Linquan Zhang, Zongpeng Li, Chuan Wu 0001
INFOCOM3
2014 Online algorithms for uploading deferrable big data to the cloud
abstract
This work studies how to minimize the bandwidth cost for uploading deferral big data to a cloud computing platform, for processing by a MapReduce framework, assuming the Internet service provider (ISP) adopts the MAX contract pricing scheme. We first analyze the single ISP case and then generalize to the MapReduce framework over a cloud platform. In the former, we design a Heuristic Smoothing algorithm whose worst-case competitive ratio is proved to fall between 2−1/(D+1) and 2(1 − 1/e), where D is the maximum tolerable delay. In the latter, we employ the Heuristic Smoothing algorithm as a building block, and design an efficient distributed randomized online algorithm, achieving a constant expected competitive ratio. The Heuristic Smoothing algorithm is shown to outperform the best known algorithm in the literature through both theoretical analysis and empirical studies. The efficacy of the randomized online algorithm is also verified through simulation studies.
Linquan Zhang, Zongpeng Li, Chuan Wu 0001, Minghua Chen 0001
INFOCOM3
2014 Dynamic pricing and profit maximization for the cloud with geo-distributed data centers
abstract
Cloud providers often choose to operate datacenters over a large geographic span, in order that users may be served by resources in their proximity. Due to time and spatial diversities in utility prices and operational costs, different datacenters typically have disparate charges for the same services. Cloud users are free to choose the datacenters to run their jobs, based on a joint consideration of monetary charges and quality of service. A fundamental problem with significant economic implications is how the cloud should price its datacenter resources at different locations, such that its overall profit is maximized. The challenge escalates when dynamic resource pricing is allowed and long-term profit maximization is pursued. We design an efficient online algorithm for dynamic pricing of VM resources across datacenters in a geo-distributed cloud, together with job scheduling and server provisioning in each datacenter, to maximize the profit of the cloud provider over a long run. Theoretical analysis shows that our algorithm can schedule jobs within their respective deadlines, while achieving a time-average overall profit closely approaching the offline maximum, which is computed by assuming that perfect information on future job arrivals are freely available. Empirical studies further verify the efficacy of our online profit maximizing algorithm.
Jian Zhao 0008, Hongxing Li 0002, Chuan Wu 0001, Zongpeng Li, Francis C. M. Lau 0001
INFOCOM3
2014 RSMOA: A revenue and social welfare maximizing online auction for dynamic cloud resource provisioning
abstract
We study online cloud resource auctions where users can arrive anytime and bid for heterogeneous types of virtual machines (VMs) assembled and provisioned on the fly. The proposed auction mechanism RSMOA, to the authors' knowledge, represents the first truthful online mechanism that timely responds to incoming users' demands and makes dynamic resource provisioning and allocation decisions, while guaranteeing efficiency in both the provider's revenue and system social welfare. RSMOA consists of two components: (1) an online mechanism that computes resource allocation and users' payments based on a global, non-decreasing pricing curve, and guarantees truthfulness; (2) a judiciously designed pricing curve, which is derived from a threat-based strategy and guarantees a competitive ratio O(ln(p)) in both system social welfare and the provider's revenue, as compared to the celebrated offline Vickrey-Clarke-Groves (VCG) auction. Here p is the ratio between the upper and lower bounds of users' marginal valuation of a type of resource. The efficacy of RSMOA is validated through extensive theoretical analysis and trace-driven simulation studies.
Chuan Wu 0001, Zongpeng Li
IWQoS2
2014 An online auction framework for dynamic resource provisioning in cloud computing
abstract
Auction mechanisms have recently attracted substantial attention as an efficient approach to pricing and resource allocation in cloud computing. This work, to the authors' knowledge, represents the first online combinatorial auction designed in the cloud computing paradigm, which is general and expressive enough to both (a) optimize system efficiency across the temporal domain instead of at an isolated time point, and (b) model dynamic provisioning of heterogeneous Virtual Machine (VM) types in practice. The final result is an online auction framework that is truthful, computationally efficient, and guarantees a competitive ratio ~ e+ 1 over e-1 ~ 3.30 in social welfare in typical scenarios. The framework consists of three main steps: (1) a tailored primal-dual algorithm that decomposes the long-term optimization into a series of independent one-shot optimization problems, with an additive loss of 1 over e-1 in competitive ratio, (2) a randomized auction sub-framework that applies primal-dual optimization for translating a centralized co-operative social welfare approximation algorithm into an auction mechanism, retaining a similar approximation ratio while adding truthfulness, and (3) a primal-dual update plus dual fitting algorithm for approximating the one-shot optimization with a ratio λ close to e. The efficacy of the online auction framework is validated through theoretical analysis and trace-driven simulation studies. We are also in the hope that the framework, as well as its three independent modules, can be instructive in auction design for other related problems.
Linquan Zhang, Chuan Wu 0001, Zongpeng Li, Francis C. M. Lau 0001
SIGMETRICS3
2014 Randomized auction design for electricity markets between grids and microgrids
abstract
This work studies electricity markets between power grids and microgrids, an emerging paradigm of electric power generation and supply. It is among the first that addresses the economic challenges arising from such grid integration, and represents the first power auction mechanism design that explicitly handles the Unit Commitment Problem (UCP), a key challenge in power grid optimization previously investigated only for centralized cooperative algorithms. The proposed solution leverages a recent result in theoretical computer science that can decompose an optimal fractional (infeasible) solution to NP-hard problems into a convex combination of integral (feasible) solutions. The end result includes randomized power auctions that are (approximately) truthful and computationally efficient, and achieve small approximation ratios for grid-wide social welfare under UCP constraints and temporal demand correlations. Both power markets with grid-to-microgrid and microgrid-to-grid energy sales are studied, with an auction designed for each, under the same randomized power auction framework. Trace driven simulations are conducted to verify the efficacy of the two proposed inter-grid power auctions.
Linquan Zhang, Zongpeng Li, Chuan Wu 0001
SIGMETRICS3
2014 Latency-minimizing data aggregation in wireless sensor networks under physical interference model
Hongxing Li 0002, Chuan Wu 0001, Qiang-Sheng Hua, Francis C. M. Lau 0001
Ad Hoc Networks2
2014 The performance and locality tradeoff in bittorrent-like file sharing systems
Wei Huang 0027, Chuan Wu 0001, Zongpeng Li, Francis C. M. Lau 0001
Peer-to-Peer Netw. Appl.2
2014 A Geometric Perspective to Multiple-Unicast Network Coding
abstract
The multiple-unicast network coding conjecture states that for multiple unicast sessions in an undirected network, network coding is equivalent to routing. Simple and intuitive as it appears, the conjecture has remained open since its proposal in 2004, and is now a well-known unsolved problem in the field of network coding. Based on a recently proposed tool of space information flow, we present a geometric framework for analyzing the multiple-unicast conjecture. The framework consists of four major steps, in which the conjecture is transformed from its throughput version to cost version, from the graph domain to the space domain, and then from high dimension to 1-D, where it is to be eventually proved. We apply the geometric framework to derive unified proofs to known results of the conjecture, as well as new results previously unknown. A possible proof to the conjecture based on this framework is outlined.
Tang Xiahou, Zongpeng Li, Chuan Wu 0001, Jiaqing Huang
IEEE Trans. Inf. Theory3
2013 An Auction and League Championship Algorithm Based Resource Allocation Mechanism for Distributed Cloud
Jiajia Sun, Xingwei Wang 0001, Keqin Li 0001, Chuan Wu 0001, Min Huang 0001
APPT4
2013 SmartDPSS: Cost-Minimizing Multi-source Power Supply for Datacenters with Arbitrary Demand
abstract
To tackle soaring power costs, significant carbon emission and unexpected power outage, Cloud Service Providers (CSPs) typically equip their Datacenters with a Power Supply System (DPSS) nurtured by multiple sources: (1) smart grid with time-varying electricity prices, (2) uninterrupted power supply (UPS), and (3) renewable energy with intermittent and uncertain supply. It remains a significant challenge how to operate among multiple power supply sources in a complementary manner, to deliver reliable energy to datacenter users with arbitrary demand over time, while minimizing a CSP's operation cost over the long run. This paper proposes an efficient, online control algorithm for DPSS, SmartDPSS, based on the two-timescale Lyapunov optimization techniques. Without requiring a priori knowledge of system statistics, SmartDPSS allows CSPs to make online decisions on how much power demand, including delay-sensitive demand and delay-tolerant demand, to serve at each time, the amount of power to purchase from the long-term-ahead and realtime grid markets, and charging and discharging of UPS over time, in order to fully leverage the available renewable energy and time-varying prices from the grid markets, for minimum operational cost. We thoroughly analyze the performance of our online control algorithm with rigorous theoretical analysis. We also demonstrate its optimality in terms of operational cost, demand service delay, datacenter availability, system robustness and scalability, using extensive simulations based on one-month worth of traces from live power systems.
Fangming Liu, Hai Jin 0001, Chuan Wu 0001
ICDCS4
2013 Profit-maximizing virtual machine trading in a federation of selfish clouds
abstract
The emerging federated cloud paradigm advocates sharing of resources among cloud providers, to exploit temporal availability of resources and diversity of operational costs for job serving. While extensive studies exist on enabling interoperability across different cloud platforms, a fundamental question on cloud economics remains unanswered: When and how should a cloud trade VMs with others, such that its net profit is maximized over the long run? In order to answer this question by the federation, a number of important, correlated decisions, including job scheduling, server provisioning and resource pricing, need to be dynamically made, with long-term profit optimality being a goal. In this work, we design efficient algorithms for inter-cloud resource trading and scheduling in a federation of geo-distributed clouds. For VM trading among clouds, we apply a double auction-based mechanism that is strategy proof, individual rational, and ex-post budget balanced. Coupling with the auction mechanism is an efficient, dynamic resource trading and scheduling algorithm, which carefully decides the true valuations of VMs in the auction, optimally schedules stochastic job arrivals with different SLAs onto the VMs, and judiciously turns on and off servers based on the current electricity prices. Through rigorous analysis, we show that each individual cloud, by carrying out our dynamic algorithm, can achieve a time-averaged profit arbitrarily close to the offline optimum.
Hongxing Li 0002, Chuan Wu 0001, Zongpeng Li, Francis C. M. Lau 0001
INFOCOM2
2013 Socially-optimal multi-hop secondary communication under arbitrary primary user mechanisms
abstract
In a cognitive radio system, licensed primary users can lease idle spectrum to secondary users for monetary remuneration. Secondary users acquire available spectrum for their data delivery needs, with the goal of achieving high throughput and low spectrum charges. Maximizing such a net utility (throughput utility minus spectrum cost) is a central problem faced by a multihop secondary network. Optimal decision making is challenging, since it involves multiple data flows, cross-layer coordination, and economic constraints (budgets of sources). The picture is further complicated by the inter-play between secondary data communication and primary spectrum leasing mechanisms. This work is the first to investigate the full spectrum of socially optimal secondary user communication. We design a social welfare maximization framework for multi-session multi-hop secondary data dissemination based on Lyapunov optimization techniques. A salient feature of the framework is that it takes any given primary user mechanism as input, and produces correspondingly a dynamic, distributed rate control, routing, and spectrum allocation and pricing protocol that can achieve longterm maximization of the overall system utility. Through rigorous theoretical analysis, we prove that our online protocol can achieve a social welfare that is arbitrarily close to the offline optimum, with only finite buffer space requirement at each secondary user, and guarantee of no buffer overflow. Empirical studies are conducted to examine the performance of the protocol.
Hongxing Li 0002, Chuan Wu 0001, Zongpeng Li, Francis C. M. Lau 0001
INFOCOM2
2013 Moving big data to the cloud
abstract
Cloud computing, rapidly emerging as a new computation paradigm, provides agile and scalable resource access in a utility-like fashion, especially for the processing of big data. An important open issue here is how to efficiently move the data, from different geographical locations over time, into a cloud for effective processing. The de facto approach of hard drive shipping is not flexible, nor secure. This work studies timely, cost-minimizing upload of massive, dynamically-generated, geodispersed data into the cloud, for processing using a MapReducelike framework. Targeting at a cloud encompassing disparate data centers, we model a cost-minimizing data migration problem, and propose two online algorithms, for optimizing at any given time the choice of the data center for data aggregation and processing, as well as the routes for transmitting data there. The first is an online lazy migration (OLM) algorithm achieving a competitive ratio of as low as 2.55, under typical system settings. The second is a randomized fixed horizon control (RFHC) algorithm achieving a competitive ratio of 1+ 1/l+λ κ/λ with a lookahead window of l, where κ and λ are system parameters of similar magnitude.
Linquan Zhang, Chuan Wu 0001, Zongpeng Li, Chuanxiong Guo, Minghua Chen 0001, Francis C. M. Lau 0001
INFOCOM2
2013 Capacity of P2P on-demand streaming with simple, robust and decentralized control
abstract
The performance of large-scaled peer-to-peer (P2P) video-on-demand (VoD) streaming systems can be very challenging to analyze. In practical P2P VoD systems, each peer only interacts with a small number of other peers/neighbors. Further, its upload capacity may vary randomly, and both its downloading position and content availability change dynamically. In this paper, we rigorously study the achievable streaming capacity of large-scale P2P VoD systems with sparse connectivity among peers, and investigate simple and decentralized P2P control strategies that can provably achieve close-to-optimal streaming capacity. We first focus on a single streaming channel. We show that a close-to-optimal streaming rate can be asymptotically achieved for all peers with high probability as the number of peers N increases, by assigning each peer a random set of Θ(log N) neighbors and using a uniform rate-allocation algorithm. Further, the tracker does not need to obtain detailed knowledge of which chunks each peer caches, and hence incurs low overhead. We then study multiple streaming channels where peers watching one channel may help in another channel with insufficient upload bandwidth. We propose a simple random cache-placement strategy, and show that a close-to-optimal streaming capacity region for all channels can be attained with high probability, again with only Θ(log N) per-peer neighbors. These results provide important insights into the dynamics of large-scale P2P VoD systems, which will be useful for guiding the design of improved P2P control protocols.
Can Zhao 0006, Jian Zhao 0008, Xiaojun Lin 0001, Chuan Wu 0001
INFOCOM4
2013 On uniform matroidal networks
abstract
Matroidal networks play a fundamental role in proving theoretical results on the limits of network coding. This can be explained by the underlying connections between network coding and matroid theory, both of which build upon the fundamental concept of independence. Two existing methods are known in the network coding literature for constructing networks from a matroid. The method due to Dougherty et al. [5] is high in time complexity but can create relatively simple network structures from a given matroid. Another method due to El Rouayheb et al. [3] is low in time complexity, but results in rather complex network structures. This work studies the design of matroidal networks from uniform matroids, targetting both low time complexity and minimum network sizes. Our construction is based on the new technique of dependence deduction, which may serve as a promising direction for constructing general matroidal networks. Some of our constructions lead to new networks for understanding network coding in terms of base field requirement.
Zongpeng Li, Chuan Wu 0001, Xunrui Yin
ISIT3
2013 Cost-minimizing preemptive scheduling of mapreduce workloads on hybrid clouds
abstract
MapReduce has become the dominant programming model for processing massive amounts of data on cloud platforms. More and more enterprises are now utilizing hybrid clouds, consisting of private infrastructure owned by themselves and public clouds such as Amazon EC2, to process their spiky MapReduce workloads, which fully utilize their own on-premise resources while outsourcing the tasks only when needed. With disparate workloads of different MapReduce tasks, an efficient scheduling mechanism is in need to enable efficient utilization of the on-premise resources and to minimize the task outsourcing cost, while meeting the task completion time requirements as well. In this paper, a fine-grained model is described to characterize the scheduling of heterogeneous MapReduce workloads, and an online algorithm is proposed for joint task admission control into the private cloud, task outsourcing to the public cloud, and VM allocation to execute the admitted tasks on the private cloud, such that the time-averaged task outsourcing cost is minimized over the long run. The online algorithm features preemptive scheduling of the tasks, where a task executed partially on the on-premise infrastructure can be paused and scheduled to run later. It also achieves desirable properties such as meeting a pre-set task admission ratio and bounding the worst-case task completion time, as proven by our rigorous theoretical analysis.
Xuanjia Qiu, Wai-Leong Yeow, Chuan Wu 0001, Francis C. M. Lau 0001
IWQoS3
2013 Moving Big Data to The Cloud: An Online Cost-Minimizing Approach
abstract
Cloud computing, rapidly emerging as a new computation paradigm, provides agile and scalable resource access in a utility-like fashion, especially for the processing of big data. An important open issue here is to efficiently move the data, from different geographical locations over time, into a cloud for effective processing. The de facto approach of hard drive shipping is not flexible or secure. This work studies timely, cost-minimizing upload of massive, dynamically-generated, geo-dispersed data into the cloud, for processing using a MapReduce-like framework. Targeting at a cloud encompassing disparate data centers, we model a cost-minimizing data migration problem, and propose two online algorithms: an online lazy migration (OLM) algorithm and a randomized fixed horizon control (RFHC) algorithm , for optimizing at any given time the choice of the data center for data aggregation and processing, as well as the routes for transmitting data there. Careful comparisons among these online and offline algorithms in realistic settings are conducted through extensive experiments, which demonstrate close-to-offline-optimum performance of the online algorithms.
Linquan Zhang, Chuan Wu 0001, Zongpeng Li, Chuanxiong Guo, Minghua Chen 0001, Francis C. M. Lau 0001
IEEE J. Sel. Areas Commun.2
2013 CloudMoV: Cloud-Based Mobile Social TV
abstract
The rapidly increasing power of personal mobile devices (smartphones, tablets, etc.) is providing much richer contents and social interactions to users on the move. This trend however is throttled by the limited battery lifetime of mobile devices and unstable wireless connectivity, making the highest possible quality of service experienced by mobile users not feasible. The recent cloud computing technology, with its rich resources to compensate for the limitations of mobile devices and connections, can potentially provide an ideal platform to support the desired mobile services. Tough challenges arise on how to effectively exploit cloud resources to facilitate mobile services, especially those with stringent interaction delay requirements. In this paper, we propose the design of a Cloud-based, novel Mobile sOcial tV system (CloudMoV). The system effectively utilizes both PaaS (Platform-as-a-Service) and IaaS (Infrastructure-as-a-Service) cloud services to offer the living-room experience of video watching to a group of disparate mobile users who can interact socially while sharing the video. To guarantee good streaming quality as experienced by the mobile users with time-varying wireless connectivity, we employ a surrogate for each user in the IaaS cloud for video downloading and social exchanges on behalf of the user. The surrogate performs efficient stream transcoding that matches the current connectivity quality of the mobile user. Given the battery life as a key performance bottleneck, we advocate the use of burst transmission from the surrogates to the mobile users, and carefully decide the burst size which can lead to high energy efficiency and streaming quality. Social interactions among the users, in terms of spontaneous textual exchanges, are effectively achieved by efficient designs of data storage with BigTable and dynamic handling of large volumes of concurrent messages in a typical PaaS cloud. These various designs for flexible transcoding capabilities, battery efficiency of mobile devices and spontaneous social interactivity together provide an ideal platform for mobile social TV services. We have implemented CloudMoV on Amazon EC2 and Google App Engine and verified its superior performance based on real-world experiments.
Yu Wu 0010, Chuan Wu 0001, Zongpeng Li, Francis C. M. Lau 0001
IEEE Trans. Multim.3
2013 Peer-Assisted Social Media Streaming with Social Reciprocity
abstract
Online video sharing and social networking are cross-pollinating rapidly in today's Internet: Online social network users are sharing more and more media contents among each other, while online video sharing sites are leveraging social connections among users to promote their videos. An intriguing development as it is, the operational challenge in previous video sharing systems persists, em i.e., the large server cost demanded for scaling of the systems. Peer-to-peer video sharing could be a rescue, only if the video viewers' mutual resource contribution has been fully incentivized and efficiently scheduled. Exploring the unique advantages of a social network based video sharing system, we advocate to utilize social reciprocities among peers with social relationships for efficient contribution incentivization and scheduling, so as to enable high-quality video streaming with low server cost. We exploit social reciprocity with two give-and-take ratios at each peer: (1) peer contribution ratio (em PCR), which evaluates the reciprocity level between a pair of social friends, and (2) system contribution ratio (em SCR), which records the give-and-take level of the user to and from the entire system. We design efficient peer-to-peer mechanisms for video streaming using the two ratios, where each user optimally decides which other users to seek relay help from and help in relaying video streams, respectively, based on combined evaluations of their social relationship and historical reciprocity levels. Our design achieves effective incentives for resource contribution, load balancing among relay peers, as well as efficient social-aware resource scheduling. We also discuss practical implementation and implement our design in a prototype social media sharing system. Our extensive evaluations based on PlanetLab experiments verify that high-quality large-scale social media sharing can be achieved with conservative server costs.
Zhi Wang 0001, Chuan Wu 0001, Lifeng Sun, Shiqiang Yang
IEEE Trans. Netw. Serv. Manag.2
2013 Aggregation Latency-Energy Tradeoff in Wireless Sensor Networks with Successive Interference Cancellation
abstract
Minimizing latency and energy consumption is the prime objective of the design of data aggregation in battery-powered wireless networks. A tradeoff exists between the aggregation latency and the energy consumption, which has been widely studied under the protocol interference model. There has been, however, no investigation of the tradeoff under the physical interference model that is known to capture more accurately the characteristics of wireless interferences. When coupled with the technique of successive interference cancellation, by which a receiver may recover signals from multiple simultaneous senders, the model can lead to much reduced latency but increased energy usage. In this paper, we investigate the latency-energy tradeoff for data aggregation in wireless sensor networks under the physical interference model and using successive interference cancellation. We present theoretical lower bounds on both latency and energy as well as their tradeoff, and give an efficient approximation algorithm that can achieve the asymptotical optimum in both aggregation latency and latency-energy tradeoff. We show that our algorithm can significantly reduce the aggregation latency, for which the energy consumption is kept at its lowest possible level.
Hongxing Li 0002, Chuan Wu 0001, Dongxiao Yu, Qiang-Sheng Hua, Francis C. M. Lau 0001
IEEE Trans. Parallel Distributed Syst.2
2013 Signal Alignment: Enabling Physical Layer Network Coding for MIMO Networking
abstract
We apply signal alignment (SA), a wireless communication technique that enables physical layer network coding (PNC) in multi-input multi-output (MIMO) wireless networks. Through calculated precoding, SA contracts the perceived signal space at a node to match its receive capability, and hence facilitates the demodulation of linearly combined data packets. PNC coupled with SA (PNC-SA) has the potential of fully exploiting the precoding space at the senders, and can better utilize the spatial diversity of a MIMO network for higher system degrees-of-freedom (DoF). PNC-SA adopts the idea of `demodulating a linear combination' from PNC. The design of PNC-SA is also inspired by recent advances in IA, though SA aligns signals not interferences. We study the optimal precoding and power allocation problem of PNC-SA, for SNR (singal-to-noise-ratio) maximization at the receiver. The mapping from SNR to BER is then analyzed, revealing that the DoF gain of PNC-SA does not come with a sacrifice in BER. We then design a general PNC-SA algorithm in larger systems, and demonstrate general applications of PNC-SA, and show via network level simulations that it can substantially increase the throughput of unicast and multicast sessions, by opening previously unexplored solution spaces in multi-hop MIMO routing.
Ruiting Zhou, Zongpeng Li, Chuan Wu 0001, Carey L. Williamson
IEEE Trans. Wirel. Commun.3
2012 Epidemic forwarding in mobile social networks
abstract
Recent years have witnessed the prosperity of mobile social networks, where various information is shared among mobile users through their opportunistic contacts. To investigate efficiency of information dissemination in wireless networks, epidemic models have been employed to study message forwarding delays, presuming message delivery whenever an opportunistic contact occurs. A practical concern is typically neglected, that one mobile user may only be willing to pass information onto others with social ties, rather than anyone upon contact. Under such a constraint, information dissemination may behave differently, according to the pattern of social ties that exist in the network. In this paper, we model social-aware epidemic forwarding in mobile social networks using mean-field equations, and carefully study the end-to-end unicast message propagation delays under different levels of social ties among users. Both cases of limited and unlimited message validity are considered in our models, i.e., whether relay nodes may delete a message after carrying it for some finite time T or never. Through careful theoretical analysis and empirical studies, we made a number of intriguing observations: First, the topology of social relation graphs significantly influences message forwarding delays, i.e., the more skewed the social relationship distribution is, the larger delay it results in. Second, the average delivery delay remains fairly stable with the growth of system scale, presenting a sharp contrast with the case without social awareness. Third, we observe that with a moderate choice of T, message delivery can achieve a successful ratio of almost 100% with an expected delay very close to the case of unlimited validity, signifying that a good tradeoff can be achieved between end-to-end message delivery efficiency and energy/storage overhead at the relay nodes in a network. All these provide useful guidance for efficient information dissemination protocol design in practical mobile social networks.
Hongxian Sun, Chuan Wu 0001
ICC2
2012 Buddy Routing: A Routing Paradigm for NanoNets Based on Physical Layer Network Coding
abstract
NanoNets are networks of nanomachines at extremely small dimensions, on the order of nanometers or micrometers. Recent advances in physics and engineering have made basic computing and communication feasible on nanomachines, and NanoNets are envisioned as an important emerging technology with broad future applications. Traditional networking solutions require significant modifications for application in NanoNets. In this paper, we focus on routing algorithm design in NanoNets. Based on the salient features of a NanoNet, including low node cost and very low available power, we propose a new routing paradigm for multi-hop data transmission in NanoNets. Our design, termed {\em Buddy Routing (BR)}, is enabled by latest advancements in physical layer network coding, and argues for pair-to-pair data forwarding in place of traditional node-to-node data forwarding. Through both analysis and simulations, we compare BR with point-to-point routing, in terms of raw throughput, error rate, energy efficiency, and protocol overhead, and show the advantages of BR in NanoNets.
Ruiting Zhou, Zongpeng Li, Chuan Wu 0001, Carey L. Williamson
ICCCN3
2012 Stochastic optimal multirate multicast in socially selfish wireless networks
abstract
Multicast supporting non-uniform receiving rates is an effective means of data dissemination to receivers with diversified bandwidth availability. Designing efficient rate control, routing and capacity allocation to achieve optimal multirate multicast has been a difficult problem in fixed wireline networks, let alone wireless networks with random channel fading and volatile node mobility. The challenge escalates if we consider also the selfishness of users who prefer to relay data for others with strong social ties. Such social selfishness of users is a new constraint in network protocol design. Its impact on efficient multicast in wireless networks has yet to be explored especially when multiple receiving rates are allowed. In this paper, we design an efficient, social-aware multirate multicast scheme that can maximize the overall utility of socially selfish users in a wireless network, and its distributed implementation. We model social preferences of users as differentiated costs for packet relay, which are weighted by the strength of social tie between the relay and the destination. Stochastic Lyapunov optimization techniques are utilized to design optimal scheduling of multicast transmissions, which are combined with multi-resolution coding and random linear network coding. With rigorous theoretical analysis, we study the optimality, stability, and complexity of our algorithm, as well as the impact of social preferences. Empirical studies further confirm the superiority of our algorithm under different social selfishness patterns.
Hongxing Li 0002, Chuan Wu 0001, Zongpeng Li, Wei Huang 0027, Francis C. M. Lau 0001
INFOCOM2
2012 Cost-minimizing dynamic migration of content distribution services into hybrid clouds
abstract
The recent advent of cloud computing technologies has enabled agile and scalable resource access for a variety of applications. Content distribution services are a major category of popular Internet applications. A growing number of content providers are contemplating a switch to cloud-based services, for better scalability and lower cost. Two key tasks are involved for such a move: to migrate their contents to cloud storage, and to distribute their web service load to cloud-based web services. The main challenge is to make the best use of the cloud as well as their existing on-premise server infrastructure, to serve volatile content requests with service response time guarantee at all times, while incurring the minimum operational cost. Employing Lyapunov optimization techniques, we present an optimization framework for dynamic, cost-minimizing migration of content distribution services into a hybrid cloud infrastructure that spans geographically distributed data centers. A dynamic control algorithm is designed, which optimally places contents and dispatches requests in different data centers to minimize overall operational cost over time, subject to service response time constraints. Rigorous analysis shows that the algorithm nicely bounds the response times within the preset QoS target in cases of arbitrary request arrival patterns, and guarantees that the overall cost is within a small constant gap from the optimum achieved by a T-slot lookahead mechanism with known information into the future.
Xuanjia Qiu, Hongxing Li 0002, Chuan Wu 0001, Zongpeng Li, Francis C. M. Lau 0001
INFOCOM3
2012 Guiding internet-scale video service deployment using microblog-based prediction
abstract
Online microblogging has been very popular in today's Internet, where users exchange short messages and follow various contents shared by people that they are interested in. Among the variety of exchanges, video links are a representative type on a microblogging site. More and more viewers of an Internet video service are coming from microblog recommendations. It is intriguing research to explore the connections between the patterns of microblog exchanges and the popularity of videos, in order to potentially use the propagation patterns of microblogs to guide proactive service deployment of a video sharing system. Based on extensive traces from Youku and Tencent Weibo, a popular video sharing site and a favored microblogging system in China, we explore how patterns of video link propagation in the microblogging system are correlated with video popularity on the video sharing site, at different times and in different geographic regions. Using influential factors summarized from the measurement studies, we further design neural network-based learning frameworks to predict the number of potential viewers of different videos and the geographic distribution of viewers. Experiments show that our neural network-based frameworks achieve better prediction accuracy, as compared to a classical approach that relies on historical numbers of views. We also briefly discuss how proactive video service deployment can be effectively enabled by our prediction frameworks.
Zhi Wang 0001, Lifeng Sun, Chuan Wu 0001, Shiqiang Yang
INFOCOM3
2012 Scaling social media applications into geo-distributed clouds
abstract
Federation of geo-distributed cloud services is a trend in cloud computing which, by spanning multiple data centers at different geographical locations, can provide a cloud platform with much larger capacities. Such a geo-distributed cloud is ideal for supporting large-scale social media streaming applications (e.g., YouTube-like sites) with dynamic contents and demands, owing to its abundant on-demand storage/bandwidth capacities and geographical proximity to different groups of users. Although promising, its realization presents challenges on how to efficiently store and migrate contents among different cloud sites (i.e. data centers), and to distribute user requests to the appropriate sites for timely responses at modest costs. These challenges escalate when we consider the persistently increasing contents and volatile user behaviors in a social media application. By exploiting social influences among users, this paper proposes efficient proactive algorithms for dynamic, optimal scaling of a social media application in a geo-distributed cloud. Our key contribution is an online content migration and request distribution algorithm with the following features: (1) future demand prediction by novelly characterizing social influences among the users in a simple but effective epidemic model; (2) oneshot optimal content migration and request distribution based on efficient optimization algorithms to address the predicted demand, and (3) a Δ(t)-step look-ahead mechanism to adjust the one-shot optimization results towards the offline optimum. We verify the effectiveness of our algorithm using solid theoretical analysis, as well as large-scale experiments under dynamic realistic settings on a home-built cloud platform.
Yu Wu 0010, Chuan Wu 0001, Bo Li 0001, Linquan Zhang, Zongpeng Li, Francis C. M. Lau 0001
INFOCOM2
2012 Space information flow: Multiple unicast
abstract
The multiple unicast network coding conjecture states that for multiple unicast in an undirected network, network coding is equivalent to routing. Simple and intuitive as it appears, the conjecture has remained open since its proposal in 2004 [1], [2], and is now a well-known unsolved problem in the field of network coding. In this work, we provide a proof to the conjecture in its space/geometric version. Space information flow is a new paradigm being proposed [3], [4]. It studies the transmission of information in a geometric space, where information flows are free to propagate along any trajectories, and may be encoded wherever they meet. The goal is to minimize a natural bandwidth-distance sum-product (network volume), while sustaining end-to-end unicast and multicast communication demands among terminals at known coordinates. The conjecture is true in networks only if it is true in space. Our main result is that network coding is indeed equivalent to routing in the space model. Besides its own merit, this partially verifies the original conjecture, and further leads to a geometric framework [5] for a hopeful proof to the conjecture.
Zongpeng Li, Chuan Wu 0001
ISIT2
2012 Auction-based P2P VoD streaming: Incentives and optimal scheduling
abstract
Real-world large-scale Peer-to-Peer (P2P) Video-on-Demand (VoD) streaming applications face more design challenges as compared to P2P live streaming, due to higher peer dynamics and less buffer overlap. The situation is further complicated when we consider the selfish nature of peers, who in general wish to download more and upload less, unless otherwise motivated. Taking a new perspective of distributed dynamic auctions, we design efficient P2P VoD streaming algorithms with simultaneous consideration of peer incentives and streaming optimality. In our solution, media block exchanges among peers are carried out through local auctions, in which budget-constrained peers bid for desired blocks from their neighbors, which in turn deliver blocks to the winning bidders and collect revenue. With strategic design of a discriminative second price auction with seller reservation, a supplying peer has full incentive to maximally contribute its bandwidth to increase its budget; requesting peers are also motivated to bid in such a way that optimal media block scheduling is achieved effectively in a fully decentralized fashion. Applying techniques from convex optimization and mechanism design, we prove (a) the incentive compatibility at the selling and buying peers, and (b) the optimality of the induced media block scheduling in terms of social welfare maximization. Large-scale empirical studies are conducted to investigate the behavior of the proposed auction mechanisms in dynamic P2P VoD systems based on real-world settings.
Chuan Wu 0001, Zongpeng Li, Xuanjia Qiu, Francis C. M. Lau 0001
ACM Trans. Multim. Comput. Commun. Appl.1
2012 Diagnosing network-wide P2P live streaming inefficiencies
abstract
Large-scale live peer-to-peer (P2P) streaming applications have been successfully deployed in today's Internet. While they can accommodate hundreds of thousands of users simultaneously with hundreds of channels of programming, there still commonly exist channels and times where and when the streaming quality is unsatisfactory. In this paper, based on more than two terabytes and one year worth of live traces from UUSee, a large-scale commercial P2P live streaming system, we show an in-depth network-wide diagnosis of streaming inefficiencies, commonly present in typical mesh-based P2P live streaming systems. As the first highlight of our work, we identify an evolutionary pattern of low streaming quality in the system, and the distribution of streaming inefficiencies across various streaming channels and in different geographical regions. We then carry out an extensive investigation to explore the causes to such streaming inefficiencies over different times and across different channels/regions at specific times, by investigating the impact of factors such as the number of peers, peer upload bandwidth, inter-peer bandwidth availability, server bandwidth consumption, and many more. The original discoveries we have brought forward include the two-sided effects of peer population on the streaming quality in a streaming channel, the significant impact of inter-peer bandwidth bottlenecks at peak times, and the inefficient utilization of server capacities across concurrent channels. Based on these insights, we identify problems within the existing P2P live streaming design and discuss a number of suggestions to improve real-world streaming protocols operating at a large scale.
Chuan Wu 0001, Baochun Li, Shuqiao Zhao
ACM Trans. Multim. Comput. Commun. Appl.1
2011 CloudMedia: When Cloud on Demand Meets Video on Demand
abstract
Internet-based cloud computing is a new computing paradigm aiming to provide agile and scalable resource access in a utility-like fashion. Other than being an ideal platform for computation-intensive tasks, clouds are believed to be also suitable to support large-scale applications with periods of flash crowds by providing elastic amounts of bandwidth and other resources on the fly. The fundamental question is how to configure the cloud utility to meet the highly dynamic demands of such applications at a modest cost. In this paper, we address this practical issue with solid theoretical analysis and efficient algorithm design using Video on Demand (VoD) as the example application. Having intensive bandwidth and storage demands in real time, VoD applications are purportedly ideal candidates to be supported on a cloud platform, where the on-demand resource supply of the cloud meets the dynamic demands of the VoD applications. We introduce a queueing network based model to characterize the viewing behaviors of users in a multichannel VoD application, and derive the server capacities needed to support smooth playback in the channels for two popular streaming models: client-server and P2P. We then propose a dynamic cloud resource provisioning algorithm which, using the derived capacities and instantaneous network statistics as inputs, can effectively support VoD streaming with low cloud utilization cost. Our analysis and algorithm design are verified and extensively evaluated using large-scale experiments under dynamic realistic settings on a home-built cloud platform.
Yu Wu 0010, Chuan Wu 0001, Bo Li 0001, Xuanjia Qiu, Francis C. M. Lau 0001
ICDCS2
2011 Strategyproof auctions for balancing social welfare and fairness in secondary spectrum markets
abstract
Secondary spectrum access is emerging as a promising approach for mitigating the spectrum scarcity in wireless networks. Coordinated spectrum access for secondary users can be achieved using periodic spectrum auctions. Recent studies on such auction design mostly neglect the repeating nature of such auctions, and focus on greedily maximizing social welfare. Such auctions can cause subsets of users to experience starvation in the long run, reducing their incentive to continue participating in the auction. It is desirable to increase the diversity of users allocated spectrum in each auction round, so that a trade-off between social welfare and fairness is maintained. We study truthful mechanisms towards this objective, for both local and global fairness criteria. For local fairness, we introduce randomization into the auction design, such that each user is guaranteed a minimum probability of being assigned spectrum. Computing an optimal, interference-free spectrum allocation is NP-Hard; we present an approximate solution, and tailor a payment scheme to guarantee truthful bidding is a dominant strategy for all secondary users. For global fairness, we adopt the classic max-min fairness criterion. We tailor another auction by applying linear programming techniques for striking the balance between social welfare and max-min fairness, and for finding feasible channel allocations. In particular, a pair of primal and dual linear programs are utilized to guide the probabilistic selection of feasible allocations towards a desired tradeoff in expectation.
Ajay Gopinathan, Zongpeng Li, Chuan Wu 0001
INFOCOM3
2011 Anonymous communication with network coding against traffic analysis attack
abstract
Flow untraceability is one critical requirement for anonymous communication with network coding, which prevents malicious attackers with wiretapping and traffic analysis abilities from relating the senders to the receivers, using linear dependency of the received packets. There have recently been proposals advocating encryptions on the Global Encoding Vectors (GEV) of network coding to thwart such attacks [1], [2]. Nevertheless, there has been no exploration of the capability of networking coding itself, to constitute more efficient and effective algorithms which guarantee anonymity. In this paper, we design a novel, simple, and effective linear network coding mechanism (ALNCode) to achieve flow untraceability in a communication network with multiple unicast flows. With solid theoretical analysis, we first show that linear network coding (LNC) can be applied to thwart traffic analysis attacks without the need of encrypting GEVs. Our key idea is to mix multiple flows at their intersection nodes by generating downstream GEVs from the common basis of upstream GEVs belonging to multiple flows, in order to hide the correlation of upstream and downstream GEVs in each flow. We then design a deterministic LNC scheme to implement our idea, by which the downstream GEVs produced are guaranteed to obfuscate their correlation with the corresponding upstream GEVs. We also give extensive theoretical analysis on the intersection probability of GEV bases and the influential factors to the effectiveness of our scheme, as well as the algorithm complexity to support its efficiency.
Jin Wang 0009, Jianping Wang 0001, Chuan Wu 0001, Kejie Lu, Naijie Gu
INFOCOM3
2011 The streaming capacity of sparsely-connected P2P systems with distributed control
abstract
Peer-to-Peer (P2P) streaming technologies can take advantage of the upload capacity of clients, and hence can scale to large content distribution networks with lower cost. A fundamental question for P2P streaming systems is the maximum streaming rate that all users can sustain. Prior works have studied the optimal streaming rate for a complete network, where every peer is assumed to communicate with all other peers. This is however an impractical assumption in real systems. In this paper, we are interested in the achievable streaming rate when each peer can only connect to a small number of neighbors. We show that even with a random peer selection algorithm and uniform rate allocation, as long as each peer maintains Ω(log N) downstream neighbors, where N is the total number of peers in the system, the system can asymptotically achieve a streaming rate that is close to the optimal streaming rate of a complete network.We then extend our analysis to multi-channel P2P networks, and we study the scenario where “helpers” from channels with excessive upload capacity can help peers in channels with insufficient upload capacity. We show that by letting each peer select Ω(log N) neighbors randomly from either the peers in the same channel or from the helpers, we can achieve a close-to-optimal streaming capacity region. Simulation results are provided to verify our analysis.
Can Zhao 0006, Xiaojun Lin 0001, Chuan Wu 0001
INFOCOM3
2011 Peer-assisted online games with social reciprocity
abstract
Online games and social networks are cross-pollinating rapidly in today's Internet: Online social network sites are deploying more and more games in their systems, while online game providers are leveraging social networks to power their games. An intriguing development as it is, the operational challenge in the previous game persists, i.e., the large server operational cost remains a non-negligible obstacle for deploying high-quality multi-player games. Peer-to-peer based game network design could be a rescue, only if the game players' mutual resource contribution has been fully incentivized and efficiently scheduled. Exploring the unique advantage of social network based games (social games), we advocate to utilize social reciprocities among peers with social relationships for efficient contribution incentivization and scheduling, so as to power a high-quality online game with low server cost. In this paper, social reciprocity is exploited with two give-and-take ratios at each peer: (1) peer contribution ratio (PCR), which evaluates the reciprocity level between a pair of social friends, and (2) system contribution ratio (SCR), which records the give-and-take level of the player to and from the entire network. We design efficient peer-to-peer mechanisms for game state distribution using the two ratios, where each player optimally decides which other players to seek relay help from and help in relaying game states, respectively, based on combined evaluations of their social relationship and historical reciprocity levels. Our design achieves effective incentives for resource contribution, load balancing among relay peers, as well as efficient social-aware resource scheduling. We also discuss practical implementation concerns and implement our design in a prototype online social game. Our extensive evaluations based on experiments on PlanetLab verify that high-quality large-scale social games can be achieved with conservative server costs.
Zhi Wang 0001, Chuan Wu 0001, Lifeng Sun, Shiqiang Yang
IWQoS2
2011 Utility-Maximizing Data Dissemination in Socially Selfish Cognitive Radio Networks
abstract
In cognitive radio networks, the occupation patterns of the primary users can be very dynamic, which makes optimization (e.g., utility maximization) of data dissemination among secondary users difficult. Even under the assumption that all secondary users are fully collaborative, the optimization requires cross-layer decision making which is challenging. The challenge escalates if users are socially selfish, who prefer to relay data only to those other users with whom there are social ties. Such social selfishness of users translates into new constraints on network protocol design. There has been no study so far on the impact of social selfishness on data dissemination in cognitive radio networks. In this paper, we consider social selfishness of secondary users, and propose the design of a joint end-to-end rate control, routing, and channel allocation protocol which can maximize the overall throughput utility of multi-session unicast in cognitive radio networks. We give a distributed implementation of the protocol. Based on a Lyapunov optimization framework, we address social preferences of users using differentiated buffer sizes and relay rates for different data sessions, and apply back-pressure based transmission scheduling to achieve guaranteed utility optimality. A unique contribution of our Lyapunov optimization is that only a finite-sized buffer is required at each user node, which sets our design apart from other designs in existing literature where they assume infinite buffers. We investigate the the optimality of our protocol and the impact of user social selfishness using both theoretical analysis and extensive simulations.
Hongxing Li 0002, Wei Huang 0027, Chuan Wu 0001, Zongpeng Li, Francis C. M. Lau 0001
MASS3
2011 Physical Layer Network Coding with Signal Alignment for MIMO Wireless Networks
abstract
We propose signal alignment (SA), a new wireless communication technique that enables physical layer network coding (PNC) in multi-input multi-output (MIMO) wireless networks. Through calculated preceding, SA contracts the perceived signal space at a node to match its receive diversity, and hence facilitates the demodulation of linearly combined data packets. PNC coupled with SA (PNC-SA) has the potential of fully exploiting the preceding space at the senders, and can better utilize the spatial diversity of a MIMO network for higher transmission rates, outperforming existing techniques including MIMO or PNC alone, interference alignment (IA) and interference alignment and cancellation (IAC). PNC-SA adopts the seminal idea of 'demodulate a linear combination' from PNC. The design of PNC-SA is also inspired by recent advances in IA, though SA aligns signals not interferences. We study the optimal preceding and power allocation problem of PNC-SA, for SNR maximization at the receiver. The mapping from SNR to BER is then analyzed, revealing that the throughput gain of PNC-SA does not come with a sacrifice in BER. We finally demonstrate general applications of PNC-SA, and show via network level simulations that it can substantially increase the throughput of unicast and multicast sessions, by opening previously unexplored solution spaces in multi-hop MIMO routing.
Ruiting Zhou, Zongpeng Li, Chuan Wu 0001, Carey L. Williamson
MASS3
2011 SMS: Collaborative S\mathcal{S}treaming in M\mathcal{M}obile S\mathcal{S}ocial Networks
Chenguang Kong, Chuan Wu 0001, Victor O. K. Li
Networking (2)2
2011 On Dynamic Server Provisioning in Multichannel P2P Live Streaming
abstract
To guarantee the streaming quality in live peer-to-peer (P2P) streaming channels, it is preferable to provision adequate levels of upload capacities at dedicated streaming servers, compensating for peer instability and time-varying peer upload bandwidth availability. Most commercial P2P streaming systems have resorted to the practice of overprovisioning a fixed amount of upload capacity on streaming servers. In this paper, we have performed a detailed analysis on 10 months of run-time traces from UUSee, a commercial P2P streaming system, and observed that available server capacities are not able to keep up with the increasing demand by hundreds of channels. We propose a novel online server capacity provisioning algorithm that proactively adjusts server capacities available to each of the concurrent channels, such that the supply of server bandwidth in each channel dynamically adapts to the forecasted demand, taking into account the number of peers, the streaming quality, and the channel priority. The algorithm is able to learn over time, has full Internet service provider (ISP) awareness to maximally constrain P2P traffic within ISP boundaries, and can provide differentiated streaming qualities to different channels by manipulating their priorities. To evaluate its effectiveness, our experiments are based on an implementation of the algorithm, which replays real-world traces.
Chuan Wu 0001, Baochun Li, Shuqiao Zhao
IEEE/ACM Trans. Netw.1
2010 Strategies of Collaboration in Multi-Channel P2P VoD Streaming
abstract
As compared to live peer-to-peer (P2P) streaming, modern P2P video-on-demand (VoD) systems have brought much larger volumes of videos and more interactive controls to the Internet users. Nevertheless, the larger number of available videos and the flexibility of allowing users to jump back and forth in a video, have led to much fewer numbers of concurrent peers watching at a similar pace, that reduces the chance for collaborative chunk supply among peers and thus significantly increases the server bandwidth cost. Towards the ultimate goal of maximizing peer resource utilization, in this paper, we design effective strategies for both cross-channel and intra-channel collaborations in multi- channel P2P VoD systems, such that individual peer's resources, including download/upload bandwidths and the cache capacity, are effectively utilized to maximize the streaming qualities in all the channels. In particular, each peer actively and strategically determines the supply-and-demand imbalance in different channels, as well as that among different chunks within each video, makes use of its surplus download capacity to fetch chunks with the most need, and then serves those chunks using its idle upload bandwidth, all without impairing its own streaming quality. Our extensive trace-driven simulations show the effectiveness of our strategies in reducing the server cost while guaranteeing high streaming qualities in the entire system, even during extreme scenarios such as unexpected flash crowds.
Zhi Wang 0001, Chuan Wu 0001, Lifeng Sun, Shiqiang Yang
GLOBECOM2
2010 The Performance and Locality Tradeoff in BitTorrent-Like P2P File-Sharing Systems
abstract
The recent surge of large-scale peer-to-peer (P2P) applications has brought huge amounts of P2P traffic, which significantly changes the Internet traffic pattern and increases the traffic-relay cost at the Internet Service Providers (ISPs). To alleviate the stress on networks, localized peer selection has been proposed that advocates neighbor selection within the same network (AS or ISP) to reduce the cross-ISP traffic. Nevertheless, localized peer selection may potentially lead to the downgrade of downloading speed at the peers, rendering a non-negligible tradeoff between the downloading performance and traffic localization in the P2P system. Aiming at effective peer selection strategies that achieve any desired Pareto optimum in face of the tradeoff, in this paper, we characterize the performance and locality tradeoff as a multi-objective b-matching optimization problem. In particular, we first present a generic maximum weight b-matching model that characterizes the tit-for-tat in BitTorrent-like peer selection. We then introduce multiple optimization objectives into the model, which effectively characterize the performance and locality tradeoff using simultaneous objectives to optimize. We also design fully distributed peer selection algorithms that can effectively achieve any desired Pareto optimum of the global multi-objective optimization, that represents a desired tradeoff point between performance and locality in the entire system. Our models and algorithms are supported by rigorous analysis and extensive simulations.
Wei Huang 0027, Chuan Wu 0001, Francis C. M. Lau 0001
ICC2
2010 UUSee: Large-Scale Operational On-Demand Streaming with Random Network Coding
abstract
Since the inception of network coding in information theory, we have witnessed a sharp increase of research interest in its applications in communications and networking, where the focus has been on more practical aspects. However, thus far, network coding has not been deployed in real-world commercial systems in operation at a large scale, and in a production setting. In this paper, we present the objectives, rationale, and design in the first production deployment of random network coding, where it has been used in the past year as the cornerstone of a large-scale production on-demand streaming system, operated by UUSee Inc., delivering thousands of on-demand video channels to millions of unique visitors each month. To achieve a thorough understanding of the performance of network coding, we have collected 200 Gigabytes worth of real-world traces throughout the 17-day Summer Olympic Games in August 2008, and present our lessons learned after an in-depth trace-driven analysis.
Zimu Liu, Chuan Wu 0001, Baochun Li, Shuqiao Zhao
INFOCOM2
2010 Minimum-latency aggregation scheduling in wireless sensor networks under physical interference model
abstract
Minimum-Latency Aggregation Scheduling (MLAS) is a problem of fundamental importance in wireless sensor networks. There however has been very little effort spent on designing algorithms to achieve sufficiently fast data aggregation under the physical interference model which is a more realistic model than traditional protocol interference model. In particular, a distributed solution to the problem under the physical interference model is challenging because of the need for global-scale information to compute the cumulative interference at any individual node. In this paper, we propose a distributed algorithm that solves the MLAS problem under the physical interference model in networks of arbitrary topology in O(K) time slots, where K is the logarithm of the ratio between the lengths of the longest and shortest links in the network. We also give a centralized algorithm to serve as a benchmark for comparison purposes, which aggregates data from all sources in O(log3n) time slots (where n is the total number of nodes). This is the current best algorithm for the problem in the literature. The distributed algorithm partitions the network into cells according to the value K, thus obviating the need for global information. The centralized algorithm strategically combines our aggregation tree construction algorithm with the non-linear power assignment strategy in [9]. We prove the correctness and efficiency of our algorithms, and conduct empirical studies under realistic settings to validate our analytical results.
Hongxing Li 0002, Qiang-Sheng Hua, Chuan Wu 0001, Francis C. M. Lau 0001
MSWiM3
2010 InstantLeap: an architecture for fast neighbor discovery in large-scale P2P VoD streaming
Xuanjia Qiu, Wei Huang 0027, Chuan Wu 0001, Francis C. M. Lau 0001, Xiaola Lin
Multim. Syst.3
2009 Distilling Superior Peers in Large-Scale P2P Streaming Systems
abstract
In large-scale peer-to-peer (P2P) live streaming systems with a limited supply of server bandwidth, increasing the amount of upload bandwidth supplied by peers becomes critically important to the "well being" of streaming sessions in live channels. Intuitively, two types of peers are preferred to be kept up in a live session: peers that contribute a higher percentage of their upload capacities, and peers that are stable for a long period of time. The fundamental challenge is to identify, and satisfy the needs of, these types of "superior" peers in a live session, and to achieve this goal with minimum disruption to the traditional pull-based protocols that real-world live streaming protocols use. In this paper, we conduct a comprehensive and in-depth statistical analysis based on more than 130 GB worth of runtime traces from hundreds of streaming channels in a large- scale real-world live streaming system, UUSee (among the top three commercial systems in popularity in mainland China). Our objective is to discover critical factors that may influence the longevity and bandwidth contribution ratio of peers, using survival analysis techniques such as the Cox proportional hazards model and the Mantel-Haenszel test. Once these influential factors are found, they can be used to form a superiority index to distill superior peers from the general peer population. The index can be used in any way to favor superior peers, and we simulate the use of a simple ranking mechanism in a natural selection algorithm to show the effectiveness of the index, based on a replay of real-world traces from UUSee.
Zimu Liu, Chuan Wu 0001, Baochun Li, Shuqiao Zhao
INFOCOM2
2009 Diagnosing Network-Wide P2P Live Streaming Inefficiencies
abstract
Large-scale live peer-to-peer (P2P) streaming applications have been successfully deployed in today's Internet. While they can accommodate millions of users simultaneously with hundreds of channels of programming, there still commonly exist channels and times where and when the streaming quality is unsatisfactory. In this paper, based on more than two terabytes and one year worth of live traces from UUSee, a large-scale commercial P2P live streaming system, we show an in-depth network-wide diagnosis of streaming inefficiencies, commonly present in mesh-based P2P streaming systems. We first identify an evolutionary pattern of low streaming quality in the system and the distribution of streaming inefficiencies across various streaming channels. We then carry out an extensive investigation to explore the causes to such streaming inefficiencies over different times and across different channels at specific times. The original discoveries we have brought forward include the two-sided effects of peer population on the streaming quality in a channel, the significant impact of inter-peer bandwidth bottlenecks at peak times, and the inefficient utilization of server capacities across concurrent channels. We conclude with a number of suggestions to improve real-world large-scale P2P streaming.
Chuan Wu 0001, Baochun Li, Shuqiao Zhao
INFOCOM1
2009 Why Are Peers Less Stable in Unpopular P2P Streaming Channels?
Zimu Liu, Chuan Wu 0001, Baochun Li, Shuqiao Zhao
Networking2
2009 InstantLeap: fast neighbor discovery in P2P VoD streaming
abstract
A fundamental challenge in peer-to-peer (P2P) Video-on-Demand (VoD) streaming is to quickly locate new supplying peers whenever a VCR command is issued, in order to achieve smooth viewing experiences. For most existing commercial systems which resort to tracking servers for such neighbor discovery, the increasing scale of P2P VoD systems has brought heavy load onto the dedicated servers. To avoid overloading the servers and achieve instant neighbor discovery over the self-organizing P2P overlay, we design a novel method of organizing peers watching the same video, that constitutes a light-weighted indexing structure to support efficient streaming and fast neighbor discovery at the same time. InstantLeap achieves an O(1) neighbor discovery efficiency upon any playback "leaps" across the media stream in streaming overlays of any sizes, with a low messaging cost for the overlay maintenance. We support our design with rigorous analysis and extensive simulations.
Xuanjia Qiu, Chuan Wu 0001, Xiaola Lin, Francis C. M. Lau 0001
NOSSDAV2
2008 Multi-Channel Live P2P Streaming: Refocusing on Servers
abstract
Due to peer instability and time-varying peer upload bandwidth availability in live peer-to-peer (P2P) streaming channels, it is preferable to provision adequate levels of stable upload capacities at dedicated streaming servers, in order to guarantee the streaming quality in all channels. Most commercial P2P streaming systems have resorted to the practice of over-provisioning upload capacities on streaming servers. We have performed a detailed analysis on 400 GB and 7 months of run-time traces from UUSee, a commercial P2P streaming system, and observed that available capacities on streaming servers are not able to keep up with the increasing demand imposed by hundreds of channels. We propose a novel online server capacity provisioning algorithm that proactively adjusts the server capacities available to each of the concurrent channels, such that the supply of server bandwidth in each channel dynamically adapts to the forecasted demand, taking into account the number of peers, the streaming quality, and the priorities of channels. The algorithm is able to learn over time, and has full ISP awareness to maximally constrain P2P traffic within ISP boundaries. To evaluate the effectiveness of our solution, our experimental studies are based on an implementation of the algorithm with actual channels of P2P streaming traffic, with real-world traces replayed within a server cluster.
Chuan Wu 0001, Baochun Li, Shuqiao Zhao
INFOCOM1
2008 Exploring large-scale peer-to-peer live streaming topologies
abstract
Real-world live peer-to-peer (P2P) streaming applications have been successfully deployed in the Internet, delivering live multimedia content to millions of users at any given time. With relative simplicity in design with respect to peer selection and topology construction protocols and without much algorithmic sophistication, current-generation live P2P streaming applications are able to provide users with adequately satisfying viewing experiences. That said, little existing research has provided sufficient insights on the time-varying internal characteristics of peer-to-peer topologies in live streaming. This article presents Magellan , our collaborative work with UUSee Inc., Beijing, China, for exploring and charting graph theoretical properties of practical P2P streaming topologies, gaining important insights in their topological dynamics over a long period of time. With more than 120 GB worth of traces starting September 2006 from a commercially deployed P2P live streaming system that represents UUSee's core product, we have completed a thorough and in-depth investigation of the topological properties in large-scale live P2P streaming, as well as their evolutionary behavior over time, for example, at different times of the day and in flash crowd scenarios. We seek to explore real-world P2P streaming topologies with respect to their graph theoretical metrics, such as the degree, clustering coefficient, and reciprocity. In addition, we compare our findings with results from existing studies on topological properties of P2P file sharing applications, and present new and unique observations specific to streaming. We have observed that live P2P streaming sessions demonstrate excellent scalability, a high level of reciprocity, a clustering phenomenon in each ISP, and a degree distribution that does not follow the power-law distribution.
Chuan Wu 0001, Baochun Li, Shuqiao Zhao
ACM Trans. Multim. Comput. Commun. Appl.1
2008 rStream: Resilient and Optimal Peer-to-Peer Streaming with Rateless Codes
abstract
Due to the lack of stability and reliability in peer-topeer networks, multimedia streaming over peer-to-peer networks represents several fundamental engineering challenges. First, multimedia streaming sessions need to be resilient to volatile network dynamics and node departures that are characteristic in peer-to-peer networks. Second, they need to take full advantage of the existing bandwidth capacities, by minimizing the delivery of redundant content and the need for content reconciliation among peers during streaming. Finally, streaming peers need to be optimally selected to construct high-quality streaming topologies, so that end-to-end latencies are taken into consideration. The original contributions of this paper are two-fold. First, we propose to use a recent coding technique, referred to as rateless codes, to code the multimedia bitstreams before they are transmitted over peer-to-peer links. The use of rateless codes eliminates the requirements of content reconciliation, as well as the risks of delivering redundant content over the network. Rateless codes also help the streaming sessions to adapt to volatile network dynamics. Second, we minimize end-to-end latencies in streaming sessions by optimizing towards a latency-related objective in a linear optimization problem, the solution to which can be efficiently derived in a decentralized and iterative fashion. The validity and effectiveness of our new contributions are demonstrated in extensive experiments in emulated realistic peer-to-peer environments with our rStream implementation.
Chuan Wu 0001, Baochun Li
IEEE Trans. Parallel Distributed Syst.1
2008 Dynamic Bandwidth Auctions in Multioverlay P2P Streaming with Network Coding
abstract
In peer-to-peer (P2P) live streaming applications such as IPTV, it is natural to accommodate multiple coexisting streaming overlays, corresponding to channels of programming. In the case of multiple overlays, it is a challenging task to design an appropriate bandwidth allocation protocol, such that these overlays efficiently share the available upload bandwidth on peers, media content is efficiently distributed to achieve the required streaming rate, as well as the streaming costs are minimized. In this paper, we seek to design simple, effective, and decentralized strategies to resolve conflicts among coexisting streaming overlays in their bandwidth competition and combine such strategies with network-coding-based media distribution to achieve efficient multioverlay streaming. Since such strategies of conflict are game theoretic in nature, we characterize them as a decentralized collection of dynamic auction games, in which downstream peers bid for upload bandwidth at the upstream peers for the delivery of coded media blocks. With extensive theoretical analysis and performance evaluation, we show that these local games converge to an optimal topology for each overlay in realistic asynchronous environments. Together with network-coding-based media dissemination, these streaming overlays adapt to peer dynamics, fairly share peer upload bandwidth to achieve satisfactory streaming rates, and can be prioritized.
Chuan Wu 0001, Baochun Li, Zongpeng Li
IEEE Trans. Parallel Distributed Syst.1
2007 Magellan: Charting Large-Scale Peer-to-Peer Live Streaming Topologies
abstract
Live peer-to-peer (P2P) streaming applications have been successfully deployed in the Internet. With relatively simple peer selection protocol design, modern live P2P streaming applications are able to provide millions of concurrent users adequately satisfying viewing experiences. That said, few existing research has provided sufficient insights on the time-varying internal characteristics of P2P topologies in live streaming. With 120 GB worth of traces in late 2006 from a commercial P2P live streaming system of UUSee Inc. in Beijing, this paper represents the first attempt in the research community to explore topological properties in practical P2P streaming, and how they behave over time. Starting from classical graph metrics, such as degree, clustering coefficient, and reciprocity, we explore and extend them in specific perspectives of streaming applications. We also compare our findings with existing insights from topological studies of P2P file sharing applications, which shed new and unique insights specific to streaming. Our characterization reveals the scalability of the commercial P2P streaming application even in case of large flash crowds, the clustering phenomenon of peers in each ISP, as well as the reciprocal behavior among peers, all of which play important roles in achieving its current success.
Chuan Wu 0001, Baochun Li, Shuqiao Zhao
ICDCS1
2007 Strategies of Conflict in Coexisting Streaming Overlays
abstract
In multimedia applications such as IPTV, it is natural to accommodate multiple coexisting peer-to-peer streaming overlays, corresponding to channels of programming. With coexisting streaming overlays, one wonders how these overlays may efficiently share the available upload bandwidth on peers. In order to satisfy the required streaming rate in each overlay, as well as to minimize streaming costs. In this paper, we seek to design simple, effective and decentralized strategies to resolve conflicts among coexisting streaming overlays. Since such strategies of conflict are game theoretic in nature, we characterize them as a decentralized collection of dynamic auction games, in which downstream peers submit bids for bandwidth at the upstream peers. With extensive theoretical analysis and performance evaluation, we show that the outcome of these local games is an optimal topology for each overlay that minimizes streaming costs. These overlay topologies evolve and adapt to peer dynamics, fairly share peer upload bandwidth, and can be prioritized.
Chuan Wu 0001, Baochun Li
INFOCOM1
2007 Optimal Rate Allocation in Overlay Content Distribution
Chuan Wu 0001, Baochun Li
Networking1
2007 Outburst: Efficient Overlay Content Distribution with Rateless Codes
Chuan Wu 0001, Baochun Li
Networking1
2007 Diverse: application-layer service differentiation in peer-to-peer communications
abstract
The peer-to-peer communication paradigm, when used to disseminate bulk content or to stream real-time multimedia, has enjoyed the distinct advantage of scalability when compared to the client-server model, since it takes advantage of available upload bandwidth at participating peers to alleviate server load. As multiple concurrent peer-to-peer sessions co-exist in the Internet, it is natural to demand differentiated services in different sessions, with respect to Quality of Service metrics such as bit rates and latencies. The problem of service differentiation across sessions, however, has never been addressed in the literature at the application layer. In this paper, we open a new direction of research that treats different peer-to-peer sessions with different priorities, and present Diverse, a novel application-layer approach to achieve service differentiation across different sessions. An extensive evaluation of our implementation of Diverse in an emulated peer-to-peer environment has demonstrated its effectiveness in achieving our design objectives.
Chuan Wu 0001, Baochun Li
IEEE J. Sel. Areas Commun.1
2007 Characterizing Peer-to-Peer Streaming Flows
abstract
The fundamental advantage of peer-to-peer (P2P) multimedia streaming applications is to leverage peer upload capacities to minimize bandwidth costs on dedicated streaming servers. The available bandwidth among peers is of pivotal importance to P2P streaming applications, especially as the number of peers in the streaming session reaches a very large scale. In this paper, we utilize more than 230 GB of traces collected from a commercial P2P streaming system, UUSee, over a four-month period of time. With such traces, we seek to thoroughly understand and characterize the achievable bandwidth of streaming flows among peers in large-scale real-world P2P live streaming sessions, in order to derive useful insights towards the improvement of current-generation P2P streaming protocols, such as peer selection. Using continuous traces over a long period of time, we explore evolutionary properties of inter-peer bandwidth. Focusing on representative snapshots of the entire topology at specific times, we investigate distributions of inter-peer bandwidth in various peer ISP/area/type categories, and statistically test and model the deciding factors that cause the variance of such inter-peer bandwidth. Our original discoveries in this study include: (1) The ISPs that peers belong to are more correlated to inter-peer bandwidth than their geographic locations; (2) There exist excellent linear correlations between peer last-mile bandwidth availability and inter-peer bandwidth within the same ISP, and between a subset of ISPs as well; and (3) The evolution of inter-peer bandwidth between two ISPs exhibits daily variation patterns. Based on these insights, we design a throughput expectation index that facilitates high-bandwidth peer selection without performing any measurements.
Chuan Wu 0001, Baochun Li, Shuqiao Zhao
IEEE J. Sel. Areas Commun.1
2006 Echelon: Peer-to-Peer Network Diagnosis with Network Coding
abstract
It is critical to monitor the performance and "health" of large-scale peer-to-peer applications. As an example, operators of peer-to-peer live streaming applications may be interested in observing performance bottlenecks, peer failures, and network topologies. In most cases, such observations are used to diagnose potential problems in the protocol design, to troubleshoot network outage, or to improve the Quality of Service of the peer-to-peer network in general. They are not time sensitive in nature, as delayed observations up to minutes or even hours are still valuable. However, such historical and delay-tolerant observations should include measurements of peers that have already failed or departed, as peer dynamics significantly affect the health of peer-to-peer applications. Such a delay-tolerant observation of peer-to-peer applications over a historical period of time is referred to as a diagnosis. In this paper, we present Echelon, a time-insensitive way to construct the diagnosis of a large-scale peer-to-peer application. Replacing the traditional wisdom of logging servers, we leverage the power of network coding to collect application-specific measurements on each peer, and disseminate them to other peers in a coded form. Over time, measurements of departed peers can still be recovered, simply by probing a small subset of peers in the network. Simulation studies have shown that Echelon is highly configurable, bandwidth efficient, and extremely tolerant of peer dynamics, thanks to the advantages of randomized network coding
Chuan Wu 0001, Baochun Li
IWQoS1
2005 rStream: resilient peer-to-peer streaming with rateless codes
abstract
The inherent instability and unreliability of peer-to-peer networks introduce several fundamental engineering challenges to multimedia streaming over peer-to-peer networks. First, multimedia streaming sessions need to be resilient to the volatile network dynamics in peer-to-peer networks. Second, they need to take full advantage of the existing bandwidth capacities, by minimizing the delivery of redundant content during streaming. In this paper, we propose to use a recent coding technique, referred to as rateless codes, to code the multimedia bitstreams before they are transmitted over peer-to-peer links. The use of rateless codes eliminates the requirements of content reconciliation, as well as the risks of delivering redundant content over the network. It also helps the streaming sessions to adapt to volatile network dynamics. Our preliminary simulation results demonstrate the validity and effectiveness of our new contribution, as compared to traditional solutions with or without erasure codes.
Chuan Wu 0001, Baochun Li
ACM Multimedia1
2002 Events recognition by semantic inference for sports video
abstract
In this paper, we propose a knowledge-based semantic inference scheme for events recognition in sports video. The framework includes three layers. At the bottom layer, low-level features are extracted at the frame level and semantic clips are segmented. Then we map the semantic clips to semantic concepts by a neural network and decision-tree at the second layer. Finally, semantic inference toward events recognition is performed on the predefined finite-state machine models at the top layer. The effectiveness and efficiency of our approach are demonstrated by the experimental results on events recognition in track and field videos.
Chuan Wu 0001, Yufei Ma 0006, HongJiang Zhang, Yuzhuo Zhong
ICME (1)1