Li Zhang 0002

dblp:89/5992-2 · DBLP profile ↗
← Back
66ranked-venue papers
2as first author
0since 2021 · last 2020
—ORCID · conflict

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

Systems, architecture and hardware · 34Computer networks · 13 · 1 first-authorSoftware engineering, systems software and programming languages · 8 · 1 first-authorApplied, interdisciplinary, general and emerging computing · 5Databases, data management, data science and information retrieval · 2Artificial intelligence and machine learning · 1Security and privacy · 1

Expertise — from the expertise taxonomy: the topics of the expert's papers under the CCF categories. A weight counts papers with recency: 1 for a paper about the topic, 0.3 when the topic is its context, halved every five years.

Computer architecture, parallel and distributed computing, and storage systems
21 papers
Cloud and datacenter computing · 43% Memory systems · 18% Distributed systems · 13%
Artificial intelligence
2 papers
Reinforcement learning · 84% Efficient and distributed learning · 16%
Computer networks
6 papers
Network performance modeling · 28% Routing and switching · 28% Internet architecture and protocols · 15%

Topics — the 30 heaviest of 58, each with the papers that count most for it

TopicWeightPapersLastEvidence papers
Cloud and datacenter computing
cluster resource management and scheduling
1.7102016
MapTask Scheduling in MapReduce With Data Locality: Throughput and Heavy-Traffic Optimality · IEEE/ACM Trans. Netw. 2016
Non-work-conserving effects in MapReduce: diffusion limit and criticality · SIGMETRICS 2014
MRONLINE: MapReduce online performance tuning · HPDC 2014
Cloud and datacenter computing › cluster resource management and scheduling › cluster scheduling
mapreduce scheduling
1.172014
Non-work-conserving effects in MapReduce: diffusion limit and criticality · SIGMETRICS 2014
DynMR: dynamic MapReduce with ReduceTask interleaving and MapTask backfilling · EuroSys 2014
Coupling task progress for MapReduce resource-aware scheduling · INFOCOM 2013
Memory systems
data locality
0.742016
MapTask Scheduling in MapReduce With Data Locality: Throughput and Heavy-Traffic Optimality · IEEE/ACM Trans. Netw. 2016
Coupling task progress for MapReduce resource-aware scheduling · INFOCOM 2013
Improving ReduceTask data locality for sequential MapReduce jobs · INFOCOM 2013
Distributed systems
fault tolerance
0.522017
GaDei: On Scale-Up Training as a Service for Deep Learning · ICDM 2017
HydraDB: a resilient RDMA-driven key-value middleware for in-memory cluster computing · SC 2015
Storage systems
key-value storage
0.522016
zExpander: a key-value cache with both high performance and fewer misses · EuroSys 2016
HydraDB: a resilient RDMA-driven key-value middleware for in-memory cluster computing · SC 2015
Distributed systems
clock synchronization
0.422015
Skewless Network Clock Synchronization Without Discontinuity: Convergence and Performance · IEEE/ACM Trans. Netw. 2015
Skewless network clock synchronization · ICNP 2013
Distributed systems › distributed machine learning
parameter server
0.312017
GaDei: On Scale-Up Training as a Service for Deep Learning · ICDM 2017
Memory systems › memory compression
cache compression
0.212016
zExpander: a key-value cache with both high performance and fewer misses · EuroSys 2016
Memory systems
cache management
0.212016
zExpander: a key-value cache with both high performance and fewer misses · EuroSys 2016
Memory systems › cache management
cache replacement
0.212016
zExpander: a key-value cache with both high performance and fewer misses · EuroSys 2016
Memory systems › cache
key-value cache
0.212016
zExpander: a key-value cache with both high performance and fewer misses · EuroSys 2016
Performance modeling and evaluation › queueing models › heavy-traffic analysis
heavy-traffic optimality
0.222016
Map task scheduling in MapReduce with data locality: Throughput and heavy-traffic optimality · INFOCOM 2013
MapTask Scheduling in MapReduce With Data Locality: Throughput and Heavy-Traffic Optimality · IEEE/ACM Trans. Netw. 2016
Performance modeling and evaluation
queueing models
0.222016
Map task scheduling in MapReduce with data locality: Throughput and heavy-traffic optimality · INFOCOM 2013
MapTask Scheduling in MapReduce With Data Locality: Throughput and Heavy-Traffic Optimality · IEEE/ACM Trans. Netw. 2016
Machine learning › Reinforcement learning › multi-agent reinforcement learning › multi-agent communication
information sharing
0.212015
Information sharing in distributed stochastic bandits · INFOCOM 2015
Machine learning › Reinforcement learning
multi-armed bandit
0.212015
Information sharing in distributed stochastic bandits · INFOCOM 2015
Storage systems › key-value storage
in-memory key-value store
0.212015
HydraDB: a resilient RDMA-driven key-value middleware for in-memory cluster computing · SC 2015
Performance modeling and evaluation
performance tuning
0.222014
MRONLINE: MapReduce online performance tuning · HPDC 2014
A smart hill-climbing algorithm for application server configuration · WWW 2004
Embedded and real-time systems › real-time scheduling
non-work-conserving scheduling
0.212014
Non-work-conserving effects in MapReduce: diffusion limit and criticality · SIGMETRICS 2014
Cloud and datacenter computing › configuration tuning
online tuning
0.212014
MRONLINE: MapReduce online performance tuning · HPDC 2014
Parallel and multicore computing › parallel scheduling
locality-aware scheduling
0.212013
Map task scheduling in MapReduce with data locality: Throughput and heavy-traffic optimality · INFOCOM 2013
Cloud and datacenter computing
resource management
0.212013
A Hierarchical Approach for the Resource Management of Very Large Cloud Platforms · IEEE Trans. Dependable Secur. Comput. 2013
Routing and switching › service disciplines
processor sharing
0.112012
Performance analysis of Coupling Scheduler for MapReduce/Hadoop · INFOCOM 2012
Network performance modeling
queueing analysis
0.112012
Performance analysis of Coupling Scheduler for MapReduce/Hadoop · INFOCOM 2012
Energy-efficient computing
energy-aware resource management
0.112012
Energy-Aware Autonomic Resource Allocation in Multitier Virtualized Environments · IEEE Trans. Serv. Comput. 2012
Cloud and datacenter computing
resource allocation
0.112012
Energy-Aware Autonomic Resource Allocation in Multitier Virtualized Environments · IEEE Trans. Serv. Comput. 2012
Electronic design automation › high-level synthesis
scheduling
0.112012
Coupling scheduler for MapReduce/Hadoop · HPDC 2012
Cloud and datacenter computing › virtualization
virtual machine
0.112012
Energy-Aware Autonomic Resource Allocation in Multitier Virtualized Environments · IEEE Trans. Serv. Comput. 2012
Cloud and datacenter computing › resource allocation
stochastic bin packing
0.112011
Consolidating virtual machines with dynamic bandwidth demand in data centers · INFOCOM 2011
Cloud and datacenter computing › virtualization › virtual machine management
virtual machine consolidation
0.112011
Consolidating virtual machines with dynamic bandwidth demand in data centers · INFOCOM 2011
Cloud and datacenter computing
datacenter network
0.112010
Improving the Scalability of Data Center Networks with Traffic-aware Virtual Machine Placement · INFOCOM 2010

Methods — techniques the papers use, named apart from their topics

convergence analysis · 0.8mini-batch size tuning · 0.6hyperparameter tuning · 0.6regret analysis · 0.4game theory · 0.4queueing analysis · 0.4maxweight policy · 0.4join the shortest queue · 0.4data compression · 0.2compact data organization · 0.2parameter optimization · 0.2multicore awareness · 0.2RDMA · 0.2stochastic optimization · 0.2receding horizon control · 0.2implementation · 0.2processor-sharing model · 0.1optimization · 0.1
YearPublicationVenuePosition
2020 Providing Performance Guarantees for Cloud-Deployed Applications
abstract
Applications with a dynamic workload demand need access to a flexible infrastructure to meet performance guarantees and minimize resource costs. While cloud computing provides the elasticity to scale the infrastructure on demand, cloud service providers lack control and visibility of user space applications, making it difficult to accurately scale the infrastructure. Thus, the burden of scaling falls on the user. That is, the user must determine when to trigger scaling and how much to scale. Scaling becomes even more challenging when applications exhibit dynamic changes in their behavior. In this paper, we propose a new cloud service, Dependable Compute Cloud (DC2), that automatically scales the infrastructure to meet the user-specified performance requirements, even when multiple user requests execute concurrently. DC2 employs Kalman filtering to automatically learn the (possibly changing) system parameters for each application, allowing it to proactively scale the infrastructure to meet performance guarantees. Importantly, DC2 is designed for the cloud - it is application-agnostic and does not require any offline application profiling or benchmarking, training data, or expert knowledge about the application. We evaluate DC2 via implementation on OpenStack using a multi-tier application under a range of workload mixes and arrival traces. Our experimental results demonstrate the robustness and superiority of DC2 over existing rule-based approaches with respect to avoiding SLA violations and minimizing resource consumption.
Anshul Gandhi, Parijat Dube, Alexei A. Karve, Andrzej Kochut, Li Zhang 0002
IEEE Trans. Cloud Comput.5
2019 Performance Prediction of GPU-based Deep Learning Applications
abstract
Recent years saw an increasing success in the application of deep learning methods across various domains and for tackling different problems, ranging from image recognition and classification to text processing and speech recognition. In this paper we propose and validate an approach to model the execution time for training convolutional neural networks (CNNs) deployed on GPGPUs. We demonstrate that our approach is generally applicable to a variety of CNN models and different types of G PG PU s with high accuracy, aiming at the preliminary design phases for system sizing.
Eugenio Gianniti, Li Zhang 0002, Danilo Ardagna
CLOSER2
2018 Performance Prediction of GPU-Based Deep Learning Applications
abstract
Recent years saw an increasing success in the application of deep learning methods across various domains and for tackling different problems, ranging from image recognition and classification to text processing and speech recognition. In this paper we propose and validate an approach to model the execution time for training convolutional neural networks (CNNs) deployed on GPGPUs. We demonstrate that our approach is generally applicable to a variety of CNN models and different types of G PG PU s with high accuracy, aiming at the preliminary design phases for system sizing.
Eugenio Gianniti, Li Zhang 0002, Danilo Ardagna
SBAC-PAD2
2018 Model-driven optimal resource scaling in cloud
Anshul Gandhi, Parijat Dube, Alexei A. Karve, Andrzej Kochut, Li Zhang 0002
Softw. Syst. Model.5
2017 GaDei: On Scale-Up Training as a Service for Deep Learning
abstract
Deep learning (DL) training-as-a-service (TaaS) is an important emerging industrial workload. TaaS must satisfy a wide range of customers who have no experience and/or resources to tune DL hyper-parameters (e.g., mini-batch size and learning rate), and meticulous tuning for each user's dataset is prohibitively expensive. Therefore, TaaS hyper-parameters must be fixed with values that are applicable to all users. Unfortunately, few research papers have studied how to design a system for TaaS workloads. By evaluating the IBM Watson Natural Language Classfier (NLC) workloads, the most popular IBM cognitive service used by thousands of enterprise-level clients globally, we provide empirical evidence that only the conservative hyper-parameter setup (e.g., small mini-batch size) can guarantee acceptable model accuracy for a wide range of customers. Unfortunately, smaller mini-batch size requires higher communication bandwidth in a parameter-server based DL training system. In this paper, we characterize the exceedingly high communication bandwidth requirement of TaaS using representative industrial deep learning workloads. We then present GaDei, a highly optimized shared-memory based scale-up parameter server design. We evaluate GaDei using both commercial benchmarks and public benchmarks and demonstrate that GaDei significantly outperforms the state-of-the-art parameter-server based implementation while maintaining the required accuracy. GaDei achieves near-best-possible runtime performance, constrained only by the hardware limitation. Furthermore, to the best of our knowledge, GaDei is the only scale-up DL system that provides fault-tolerance.
Wei Zhang 0057, Minwei Feng, Yunhui Zheng, Yufei Ren, Yandong Wang 0001, Peng Liu 0010, Bing Xiang, Li Zhang 0002, Bowen Zhou 0002, Fei Wang 0001
ICDM9
2017 Lightweight Replication Through Remote Backup Memory Sharing for In-memory Key-Value Stores
abstract
Memory price will continue dropping in the next few years according to Gartner. Such trend renders it affordable for in-memory key-value stores (IMKVs) to maintain redundant memory-resident copies of each key-value pair to provision enhanced reliability and high availability services. Though contemporary IMKVs have reached unprecedented performance, delivering single-digit microsecond-scale latency with up to tens of millions queries per second throughput, existing replication protocols are unable to keep pace with such an advancement of IMKVs, either incurring unbearable latency overhead or demanding intensive resource usage. Consequently, the adoption of those replication techniques always results in substantial performance degradation.In this paper, we propose MacR, a RDMA-based high-performance and lightweight replication protocol for IMKVs. The design of MacR centers around sharing the remote backup memory to enable RDMA-based replication protocol, and synthesizes a collection of optimizations, including memory allocator cooperative replication and adaptive bulk data synchronization to control the number of network operations and to enhance the recovery performance. Performance evaluations with a variety of YCSB workloads demonstrate that MacR can efficiently outperform alternative replication methods in terms of the throughput while preserving sufficiently low latency overhead. It can also efficiently speed up the recovery process.
Yandong Wang 0001, Li Zhang 0002, Michel Hack, Yufei Ren
MASCOTS2
2017 Nexus: Bringing Efficient and Scalable Training to Deep Learning Frameworks
abstract
Demand is mounting in the industry for scalable GPU-based deep learning systems. Unfortunately, existing training applications built atop popular deep learning frameworks, including Caffe, Theano, and Torch, etc, are incapable of conducting distributed GPU training over large-scale clusters. To remedy such a situation, this paper presents Nexus, a platform that allows existing deep learning frameworks to easily scale out to multiple machines without sacrificing model accuracy. Nexus leverages recently proposed distributed parameter management architecture to orchestrate distributed training by a large number of learners spread across the cluster. Through characterizing the run-time behavior of existing single-node based applications, Nexus is equipped with a suite of optimization schemes, including hierarchical and hybrid parameter aggregation, enhanced network and computation layer, and quality-guided communication adjustment, etc, to strengthen the communication channels and resource utilization. Empirical evaluations with a diverse set of deep learning applications demonstrate that Nexus is easy to integrate and can deliver efficient distributed training services to major deep learning frameworks. In addition, Nexus's optimization schemes are highly effective to shorten the training time with targeted accuracy bounds.
Yandong Wang 0001, Li Zhang 0002, Yufei Ren, Wei Zhang 0057
MASCOTS2
2016 Stage Aware Performance Modeling of DAG Based in Memory Analytic Platforms
abstract
Spark has grown both in popularity and complexity in recent years. In order to use available resources in an efficient way, users need to understand how the behavior of their applications is affected by the size of the datasets and various configuration settings. Indeed, Spark allows users to specify many configuration parameters and understanding the impact of these choices with respect to the application execution time is not easy. An accurate estimate of application execution time is important for cluster capacity planning and/or runtime scheduling. In this work we propose a gray-box approach to analyze the performance of Spark applications deployed in public cloud infrastructures. The approach is divided into two phases: during application profiling, the application is executed multiple times against different subsets of the input datasets to understand the effect of the data size on the execution time and its dependency on the main configuration parameters. Next, during the estimation phase, we use the data gathered in the first step to predict the execution time of the application, run against the entire dataset. The prediction approach builds several models in order to estimate separately the growth of the time required to execute each stage within the application. Finally, the DAG used by Spark to schedule the execution of stages is analyzed to aggregate the predictions of the stages execution times into the overall application execution time. Both phases are supported by our SLAP open source tool. Experimental results show that our model can effectively and accurately predict application execution time. The approach outperforms pure black-box polynomial regression methods obtaining 1-3% relative error.
Giovanni Paolo Gibilisco, Li Zhang 0002, Danilo Ardagna
CLOUD3
2016 zExpander: a key-value cache with both high performance and fewer misses
abstract
While key-value (KV) cache, such as memcached, dedicates a large volume of expensive memory to holding performance-critical data, it is important to improve memory efficiency, or to reduce cache miss ratio without adding more memory. As we find that optimizing replacement algorithms is of limited effect for this purpose, a promising approach is to use a compact data organization and data compression to increase effective cache size. However, this approach has the risk of degrading the cache's performance due to additional computation cost. A common perception is that a high-performance KV cache is not compatible with use of data compacting techniques.
Xingbo Wu, Li Zhang 0002, Yandong Wang 0001, Yufei Ren, Michel Hack, Song Jiang 0001
EuroSys2
2016 Autoscaling for Hadoop Clusters
abstract
Unforeseen events such as node failures and resource contention can have a severe impact on the performance of data processing frameworks, such as Hadoop, especially in cloud environments where such incidents are common. SLA compliance in the presence of such events requires the ability to quickly and dynamically resize infrastructure resources. Unfortunately, the distributed and stateful nature of data processing frameworks makes it challenging to accurately scale the system at run-time. In this paper, we present the design and implementation of a model-driven autoscaling solution for Hadoop clusters. We first develop novel gray-box performance models for Hadoop workloads that specifically relate job execution times to resource allocation and workload parameters. We then employ these models to dynamically determine the resources required to successfully complete the Hadoop jobs as per the user-specified SLA under various scenarios including node failures and multi-job executions. Our experimental results on three different Hadoop cloud clusters and across different workloads demonstrate the efficacy of our models and highlight their autoscaling capabilities.
Anshul Gandhi, Sidhartha Thota, Parijat Dube, Andrzej Kochut, Li Zhang 0002
IC2E5
2016 MEMTUNE: Dynamic Memory Management for In-Memory Data Analytic Platforms
abstract
Memory is a crucial resource for big data processing frameworks such as Spark and M3R, where the memory is used both for computation and for caching intermediate storage data. Consequently, optimizing memory is the key to extracting high performance. The extant approach is to statically split the memory for computation and caching based on workload profiling. This approach is unable to capture the varying workload characteristics and dynamic memory demands. Another factor that affects caching efficiency is the choice of data placement and eviction policy. The extant LRU policy is oblivious of task scheduling information from the analytic frameworks, and thus can lead to lost optimization opportunities. In this paper, we address the above issues by designing MEMTUNE, a dynamic memory manager for in-memory data analytics. MEMTUNE dynamically tunes computation/caching memory partitions at runtime based on workload memory demand and in-memory data cache needs. Moreover, if needed, the scheduling information from the analytic framework is leveraged to evict data that will not be needed in the near future. Finally, MEMTUNE also supports task-level data prefetching with a configurable window size to more effectively overlap computation with I/O. Our experiments show that MEMTUNE improves memory utilization, yields an overall performance gain of up to 46%, and achieves cache hit ratio of up to 41% compared to standard Spark.
Luna Xu, Li Zhang 0002, Ali Raza Butt, Yandong Wang 0001, Zane Zhenhua Hu
IPDPS3
2016 MapTask Scheduling in MapReduce With Data Locality: Throughput and Heavy-Traffic Optimality
abstract
MapReduce/Hadoop framework has been widely used to process large-scale datasets on computing clusters. Scheduling map tasks with data locality consideration is crucial to the performance of MapReduce. Many works have been devoted to increasing data locality for better efficiency. However, to the best of our knowledge, fundamental limits of MapReduce computing clusters with data locality, including the capacity region and theoretical bounds on the delay performance, have not been well studied. In this paper, we address these problems from a stochastic network perspective. Our focus is to strike the right balance between data locality and load balancing to simultaneously maximize throughput and minimize delay. We present a new queueing architecture and propose a map task scheduling algorithm constituted by the Join the Shortest Queue policy together with the MaxWeight policy. We identify an outer bound on the capacity region, and then prove that the proposed algorithm can stabilize any arrival rate vector strictly within this outer bound. It shows that the outer bound coincides with the actual capacity region, and the proposed algorithm is throughput-optimal. Furthermore, we study the number of backlogged tasks under the proposed algorithm, which is directly related to the delay performance based on Little's law. We prove that the proposed algorithm is heavy-traffic optimal, i.e., it asymptotically minimizes the number of backlogged tasks as the arrival rate vector approaches the boundary of the capacity region. Therefore, the proposed algorithm is also delay-optimal in the heavy-traffic regime. The proofs in this paper deal with random processing times with heterogeneous parameters and nonpreemptive task execution, which differentiate our work from many existing works on MaxWeight-type algorithms, so the proof techniques themselves for the stability analysis and the heavy-traffic analysis are also novel contributions.
Weina Wang 0001, Kai Zhu 0002, Lei Ying 0001, Jian Tan 0001, Li Zhang 0002
IEEE/ACM Trans. Netw.5
2015 Information sharing in distributed stochastic bandits
abstract
Information sharing is an important issue for stochastic bandit problems in a distributed setting. Consider N players dealing with the same multi-armed bandit problem. All players receive requests simultaneously and must choose one of M actions for each request. Sharing information among these N players can decrease the regret for each of them but also incurs cooperation and communication overhead. In this setting, we study how cooperation and communication can impact the system performance measured by regret and communication cost. For both scenarios, we establish a uniform lower bound to the regret for the entire system as a function of time and network size. Concerning cooperation, we study the problem from a game-theoretic perspective. When each player's actions and payoffs are immediately visible to all others, we identify strategies for all players under which co-operative exploration is ensured. Regarding the communication cost, we consider incomplete information sharing such that a player's payoffs and actions are not entirely available to others. The players communicate observations to each other to reduce their regret, however with a cost. We show that a logarithmic communication cost is necessary to achieve the optimal regret. For Bernoulli arrivals, we specify a policy that achieves the optimal regret with a logarithmic communication cost. Our work opens a novel direction towards understanding information sharing for active learning in a distributed environment.
Swapna Buccapatnam, Jian Tan 0001, Li Zhang 0002
INFOCOM3
2015 HydraDB: a resilient RDMA-driven key-value middleware for in-memory cluster computing
abstract
In this paper, we describe our experiences and lessons learned from building a general-purpose in-memory key-value middleware, called HydraDB. HydraDB synthesizes a collection of state-of-the-art techniques, including continuous fault-tolerance, Remote Direct Memory Access (RDMA), as well as awareness for multicore systems, etc, to deliver a high-throughput, low-latency access service in a reliable manner for cluster computing applications.
Yandong Wang 0001, Li Zhang 0002, Jian Tan 0001, Xavier Guerin, Xiaoqiao Meng, Shicong Meng
SC2
2015 Skewless Network Clock Synchronization Without Discontinuity: Convergence and Performance
abstract
This paper examines synchronization of computer clocks connected via a data network and proposes a skewless algorithm to synchronize them. Unlike existing solutions, which either estimate and compensate the frequency difference (skew) among clocks or introduce offset corrections that can generate jitter and possibly even backward jumps, our solution achieves synchronization without these problems. We first analyze the convergence property of the algorithm and provide explicit necessary and sufficient conditions on the parameters to guarantee synchronization. We then study the effect of noisy measurements (jitter) and frequency drift (wander) on the offsets and synchronization frequency, and further optimize the parameter values to minimize their variance. Our study reveals a few insights, for example, we show that our algorithm can converge even in the presence of timing loops and noise, provided that there is a well-defined leader. This marks a clear contrast with current standards such as NTP and PTP, where timing loops are specifically avoided. Furthermore, timing loops can even be beneficial in our scheme as it is demonstrated that highly connected subnetworks can collectively outperform individual clients when the time source has large jitter. The results are supported by experiments running on a cluster of IBM BladeCenter servers with Linux.
Enrique Mallada, Xiaoqiao Meng, Michel Hack, Li Zhang 0002, Ao Tang
IEEE/ACM Trans. Netw.4
2014 C-Hint: An Effective and Reliable Cache Management for RDMA-Accelerated Key-Value Stores
abstract
Recently, many in-memory key-value stores have started using a High-Performance network protocol, Remote Direct Memory Access (RDMA), to provision ultra-low latency access services. Among various solutions, previous studies have recognized that leveraging RDMA Read to optimize GET operations and continuing using message passing for other requests can offer tremendous performance improvement while avoiding read-write races. However, although such a design can utilize the power of RDMA when there is sufficient memory space, it has also raised new challenges on the cache management that do not exist in traditional key-value stores. First, RDMA Read deprives servers of the awareness of the read operations. Therefore, how to track popular items and make replacement decisions at the server side becomes a critical issue. Second, without the access knowledge from the clients, new approaches are needed for servers to efficiently and reliably reclaim the resources. Lastly, the remote pointers hold by clients to conduct RDMA are highly susceptible to the evictions made by remote servers. Thus, any replacement algorithm that solely considers the server-side hit ratio is insufficient and can cause severe underutilization of RDMA.
Yandong Wang 0001, Xiaoqiao Meng, Li Zhang 0002, Jian Tan 0001
SoCC3
2014 DynMR: dynamic MapReduce with ReduceTask interleaving and MapTask backfilling
abstract
In order to improve the performance of MapReduce, we design DynMR. It addresses the following problems that persist in the existing implementations: 1) difficulty in selecting optimal performance parameters for a single job in a fixed, dedicated environment, and lack of capability to configure parameters that can perform optimally in a dynamic, multi-job cluster; 2) long job execution resulting from a task long-tail effect, often caused by ReduceTask data skew or heterogeneous computing nodes; 3) inefficient use of hardware resources, since ReduceTasks bundle several functional phases together and may idle during certain phases.
Jian Tan 0001, Alicia Chin, Zane Zhenhua Hu, Yonggang Hu, Shicong Meng, Xiaoqiao Meng, Li Zhang 0002
EuroSys7
2014 MRONLINE: MapReduce online performance tuning
abstract
MapReduce job parameter tuning is a daunting and time consuming task. The parameter configuration space is huge; there are more than 70 parameters that impact job performance. It is also difficult for users to determine suitable values for the parameters without first having a good understanding of the MapReduce application characteristics. Thus, it is a challenge to systematically explore the parameter space and select a near-optimal configuration. Extant offline tuning approaches are slow and inefficient as they entail multiple test runs and significant human effort.
Liangzhao Zeng, Shicong Meng, Jian Tan 0001, Li Zhang 0002, Ali Raza Butt, Nicholas C. Fuller
HPDC5
2014 Modeling the Impact of Workload on Cloud Resource Scaling
abstract
Cloud computing offers the flexibility to dynamically size the infrastructure in response to changes in workload demand. While both horizontal and vertical scaling of infrastructure is supported by major cloud providers, these scaling options differ significantly in terms of their cost, provisioning time, and their impact on workload performance. Importantly, the efficacy of horizontal and vertical scaling critically depends on the workload characteristics, such as the workload's parallelizability and its core scalability. In today's cloud systems, the scaling decision is left to the users, requiring them to fully understand the tradeoffs associated with the different scaling options. In this paper, we present our solution for optimizing the resource scaling of cloud deployments via implementation in OpenStack. The key component of our solution is the modelling engine that characterizes the workload and then quantitatively evaluates different scaling options for that workload. Our modelling engine leverages Amdahl's Law to model service time scaling in scaleup environments and queueing-theoretic concepts to model performance scaling in scale-out environments. We further employ Kalman filtering to account for inaccuracies in the model-based methodology, and to dynamically track changes in the workload and cloud environment.
Anshul Gandhi, Parijat Dube, Alexei A. Karve, Andrzej Kochut, Li Zhang 0002
SBAC-PAD5
2014 Non-work-conserving effects in MapReduce: diffusion limit and criticality
abstract
Sequentially arriving jobs share a MapReduce cluster, each desiring a fair allocation of computing resources to serve its associated map and reduce tasks. The model of such a system consists of a processor sharing queue for the MapTasks and a multi-server queue for the ReduceTasks. These two queues are dependent through a constraint that the input data of each ReduceTask are fetched from the intermediate data generated by the MapTasks belonging to the same job. A more generalized form of MapReduce queueing model can capture the essence of other distributed data processing systems that contain interdependent processor sharing queues and multi-server queues.
Jian Tan 0001, Yandong Wang 0001, Weikuan Yu, Li Zhang 0002
SIGMETRICS4
2014 Estimating life-time distribution by observing population continuously
Hanhua Feng, Parijat Dube, Li Zhang 0002
Perform. Evaluation3
2013 Improving Multi-job MapReduce Scheduling in an Opportunistic Environment
abstract
As a state-of-the-art programming model for big data analytics, MapReduce is well suited for parallel processing of large data sets in opportunistic environments. Existing research on MapReduce in opportunistic environment has focused on improving single job performance, the issue of fairness that is critical in the more dominant scenario of multiple concurrent jobs remains unexplored. We address this problem by proposing an opportunistic fair scheduling algorithm, which extends the broadly adopted Fair Scheduler to an environment where nodes are intermittently available with possibly different availability patterns. The proposed scheduler maintains statistics specific to the opportunistic environment, e.g., node availability rates and pairwise availability correlations, and utilizes this information in scheduling decisions to improve fairness. Using a Hadoop-based implementation, we compare our scheduler with the current Hadoop Fair Scheduler on representative benchmarks. Our experiments verify that our scheduler can significantly reduce the variability in job completion times.
Yuting Ji, Lang Tong 0001, Ting He 0001, Jian Tan 0001, Kang-Won Lee 0002, Li Zhang 0002
IEEE CLOUD6
2013 Skewless network clock synchronization
abstract
This paper examines synchronization of computer clocks connected via a data network and proposes a skewless algorithm to synchronize them. Unlike existing solutions, which either estimate and compensate the frequency difference (skew) among clocks or introduce offset corrections that can generate jitter and possibly even backward jumps, our algorithm achieves synchronization without these problems. We first analyze the convergence property of the algorithm and provide necessary and sufficient conditions on the parameters to guarantee synchronization. We then implement our solution on a cluster of IBM BladeCenter servers running Linux and study its performance. In particular, both analytically and experimentally, we show that our algorithm can converge in the presence of timing loops. This marks a clear contrast with current standards such as NTP and PTP, where timing loops are specifically avoided. Furthermore, timing loops can even be beneficial in our scheme. For example, it is demonstrated that highly connected subnetworks can collectively outperform individual clients when the time source has large jitter. It is also experimentally demonstrated that our algorithm outperforms other well-established software-based solutions such as the NTPv4 and IBM Coordinated Cluster Time (IBM CCT).
Enrique Mallada, Xiaoqiao Meng, Michel Hack, Li Zhang 0002, Ao Tang
ICNP4
2013 Improving ReduceTask data locality for sequential MapReduce jobs
abstract
Improving data locality for MapReduce jobs is critical for the performance of large-scale Hadoop clusters, embodying the principle of moving computation close to data for big data platforms. Scheduling tasks in the vicinity of stored data can significantly diminish network traffic, which is crucial for system stability and efficiency. Though the issue on data locality has been investigated extensively for MapTasks, most of the existing schedulers ignore data locality for ReduceTasks when fetching the intermediate data, causing performance degradation. This problem of reducing the fetching cost for ReduceTasks has been identified recently. However, the proposed solutions are exclusively based on a greedy approach, relying on the intuition to place ReduceTasks to the slots that are closest to the majority of the already generated intermediate data. The consequence is that, in presence of job arrivals and departures, assigning the ReduceTasks of the current job to the nodes with the lowest fetching cost can prevent a subsequent job with even better match of data locality from being launched on the already taken slots. To this end, we formulate a stochastic optimization framework to improve the data locality for ReduceTasks, with the optimal placement policy exhibiting a threshold-based structure. In order to ease the implementation, we further propose a receding horizon control policy based on the optimal solution under restricted conditions. The improved performance is further validated through simulation experiments and real performance tests on our testbed.
Jian Tan 0001, Shicong Meng, Xiaoqiao Meng, Li Zhang 0002
INFOCOM4
2013 Coupling task progress for MapReduce resource-aware scheduling
abstract
Schedulers are critical in enhancing the performance of MapReduce/Hadoop in presence of multiple jobs with different characteristics and performance goals. Though current schedulers for Hadoop are quite successful, they still have room for improvement: map tasks (MapTasks) and reduce tasks (ReduceTasks) are not jointly optimized, albeit there is a strong dependence between them. This can cause job starvation and unfavorable data locality. In this paper, we design and implement a resource-aware scheduler for Hadoop. It couples the progresses of MapTasks and ReduceTasks, utilizing Wait Scheduling for ReduceTasks and Random Peeking Scheduling for MapTasks to jointly optimize the task placement. This mitigates the starvation problem and improves the overall data locality. Our extensive experiments demonstrate significant improvements in job response times.
Jian Tan 0001, Xiaoqiao Meng, Li Zhang 0002
INFOCOM3
2013 Map task scheduling in MapReduce with data locality: Throughput and heavy-traffic optimality
abstract
Scheduling map tasks to improve data locality is crucial to the performance of MapReduce. Many works have been devoted to increasing data locality for better efficiency. However, to the best of our knowledge, fundamental limits of MapReduce computing clusters with data locality, including the capacity region and theoretical bounds on the delay performance, have not been studied. In this paper, we address these problems from a stochastic network perspective. Our focus is to strike the right balance between data-locality and load-balancing to simultaneously maximize throughput and minimize delay. We present a new queueing architecture and propose a map task scheduling algorithm constituted by the Join the Shortest Queue policy together with the MaxWeight policy. We identify an outer bound on the capacity region, and then prove that the proposed algorithm stabilizes any arrival rate vector strictly within this outer bound. It shows that the algorithm is throughput optimal and the outer bound coincides with the actual capacity region. Further, we study the number of backlogged tasks under the proposed algorithm, which is directly related to the delay performance based on Little's law. We prove that the proposed algorithm is heavy-traffic optimal, i.e., it asymptotically minimizes the number of backlogged tasks as the arrival rate vector approaches the boundary of the capacity region. Therefore, the proposed algorithm is also delay optimal in the heavy-traffic regime.
Weina Wang 0001, Kai Zhu 0002, Lei Ying 0001, Jian Tan 0001, Li Zhang 0002
INFOCOM5
2013 Joint optimization of overlapping phases in MapReduce
Minghong Lin, Li Zhang 0002, Adam Wierman, Jian Tan 0001
Perform. Evaluation2
2013 A Hierarchical Approach for the Resource Management of Very Large Cloud Platforms
abstract
Worldwide interest in the delivery of computing and storage capacity as a service continues to grow at a rapid pace. The complexities of such cloud computing centers require advanced resource management solutions that are capable of dynamically adapting the cloud platform while providing continuous service and performance guarantees. The goal of this paper is to devise resource allocation policies for virtualized cloud environments that satisfy performance and availability guarantees and minimize energy costs in very large cloud service centers. We present a scalable distributed hierarchical framework based on a mixed-integer nonlinear optimization of resource management acting at multiple timescales. Extensive experiments across a wide variety of configurations demonstrate the efficiency and effectiveness of our approach.
Bernardetta Addis, Danilo Ardagna, Barbara Panicucci, Mark S. Squillante, Li Zhang 0002
IEEE Trans. Dependable Secur. Comput.5
2012 Coupling scheduler for MapReduce/Hadoop
abstract
Current schedulers of MapReduce/Hadoop are quite successful in providing good performance. However improving spaces still exist: map and reduce tasks are not jointly optimized for scheduling, albeit there is a strong dependence between them. This can cause job starvation and bad data locality. We design a resource-aware scheduler for Hadoop, which couples the progresses of mappers and reducers, and jointly optimize the placements for both of them. This mitigates the starvation problem and improves the overall data locality. Our experiments demonstrate improvements to job response times by up to an order of magnitude.
Jian Tan 0001, Xiaoqiao Meng, Li Zhang 0002
HPDC3
2012 Performance analysis of Coupling Scheduler for MapReduce/Hadoop
abstract
For MapReduce/Hadoop, map and reduce phases exhibit fundamentally distinguishing characteristics. Additionally, these two phases admit complicated and tight dependency on each other, causing the repeatedly observed starvation problem with the widely used Fair Scheduler. To mitigate this problem, we design Coupling Scheduler, which, among other new features, jointly schedules map and reduce tasks by coupling their progresses, different from existing ones that treat them separately. This design is based on the intuition that allocating excess resources to reduce tasks without balancing with the map task progress of the same job is likely to result in resource underutilization since a job is deemed done only when both phases complete. In order to analytically understand the performance of this design, we propose a model that captures the fundamental scheduling characteristics for MapReduce. Specifically, the map phase is modeled by a processor sharing queue, and the reduce phase by a “sticky processor sharing” queue. Along with the important dependence between these two types of tasks, we show that, for a class of jobs with regularly varying map service times, the job processing time distribution under Coupling Scheduler can be one order better than Fair Scheduler. These theoretical results are validated through simulations and the improved performance is further illustrated through real experiments on our testbed.
Jian Tan 0001, Xiaoqiao Meng, Li Zhang 0002
INFOCOM3
2012 Performance modeling and characterization of large last level caches
abstract
Different workloads exhibit different memory footprint and have different dependency on the size and configuration of the memory subsystem. To quantify the performance implication of various memory system architectures and attributes one needs to understand the program behavior of the workload and its use of the memory subsystem. We developed the large cache simulator (LCS) to study cache performance of different applications and its sensitivity to different architecture parameters. The LCS is a multi-processor system that runs applications with a coherently attached FPGA which can emulate different cache configurations for long periods of time. The LCS measures different statistics associated with cache performance which are then used to develop cache performance models. Our goal is to explicitly characterize the performance of large last level cache for different workloads and model its dependency on cache configuration parameters.
Parijat Dube, Michael Tsao, Li Zhang 0002, Alan Bivens
ISPASS3
2012 Performance Modeling and Characterization of Large Last Level Caches
abstract
Different workloads exhibit different memory footprint and have different dependency on the size and configuration of the memory hierarchy. To quantify the performance implication of various memory system architectures and attributes one needs to understand the program behavior of an application/workload and its use of the memory subsystem. One such architecture studied in this paper is where multiple memory technologies are integrated into the memory system with one (typically the faster, more expensive technology) acting as a large cache for the other (typically a slower, cheaper technology). We develop a large cache prototype to study the performance of different applications and their sensitivity to different architecture parameters. The prototype measures different metrics associated with a cache performance which are in turn used to characterize the performance implications of such a memory architecture on different workloads and the dependency on different configuration parameters.
Parijat Dube, Michael Tsao, Li Zhang 0002, Alan Bivens
MASCOTS3
2012 Delay tails in MapReduce scheduling
abstract
MapReduce/Hadoop production clusters exhibit heavy-tailed characteristics for job processing times. These phenomena are resultant of the workload features and the adopted scheduling algorithms. Analytically understanding the delays under different schedulers for MapReduce can facilitate the design and deployment of large Hadoop clusters. The map and reduce tasks of a MapReduce job have fundamental difference and tight dependence between them, complicating the analysis. This also leads to an interesting starvation problem with the widely used Fair Scheduler due to its greedy approach to launching reduce tasks. To address this issue, we design and implement Coupling Scheduler, which gradually launches reduce tasks depending on map task progresses. Real experiments demonstrate improvements to job response times by up to an order of magnitude.
Jian Tan 0001, Xiaoqiao Meng, Li Zhang 0002
SIGMETRICS3
2012 Experiences in building and scaling an enterprise application on multicore systems
abstract
SUMMARY Even though Java is the de facto programming language for enterprise applications, there exist only a limited number of Java‐based benchmarks to understand the performance on emerging multicore systems. To bridge this gap, this paper presents a report generation benchmark that is developed on top of Open Source Apache Geronimo's DayTrader benchmark. Report generation and rendering is at the heart of many enterprise business analytics and business intelligence software products, and it is used by many enterprise applications. We evaluate the performance scalability of this benchmark on a state‐of‐the‐art Power7 multicore system with 8 Power7 cores and 32 hardware threads. The benchmark throughput scales linearly up to eight hardware threads, but beyond that point, the throughput falls sharply. Significant locking in the Java class libraries for non‐shared objects results in this performance drop. Splitting the locks on these shared classes results in near linear scaling from eight to 32 threads and improved the throughput by 80%. We also show that the Linux operating system load balancing could result in a degraded application performance in hardware multithreaded systems and simultaneous‐multithreads‐aware task scheduling results in uniform core‐resource utilization as well as improved application performance. Copyright © 2011 John Wiley & Sons, Ltd.
Seetharami Seelam, Parijat Dube, Megumi Ito, Deniz Binay, Michael Dawson 0001, Pramod Nagaraja, Graeme Johnson, Liana L. Fong, Michel Hack, Xiaoqiao Meng, Li Zhang 0002
Concurr. Comput. Pract. Exp.13
2012 Energy-Aware Autonomic Resource Allocation in Multitier Virtualized Environments
abstract
With the increase of energy consumption associated with IT infrastructures, energy management is becoming a priority in the design and operation of complex service-based systems. At the same time, service providers need to comply with Service Level Agreement (SLA) contracts which determine the revenues and penalties on the basis of the achieved performance level. This paper focuses on the resource allocation problem in multitier virtualized systems with the goal of maximizing the SLAs revenue while minimizing energy costs. The main novelty of our approach is to address—in a unifying framework—service centers resource management by exploiting as actuation mechanisms allocation of virtual machines (VMs) to servers, load balancing, capacity allocation, server power state tuning, and dynamic voltage/frequency scaling. Resource management is modeled as an NP-hard mixed integer nonlinear programming problem, and solved by a local search procedure. To validate its effectiveness, the proposed model is compared to top-performing state-of-the-art techniques. The evaluation is based on simulation and on real experiments performed in a prototype environment. Synthetic as well as realistic workloads and a number of different scenarios of interest are considered. Results show that we are able to yield significant revenue gains for the provider when compared to alternative methods (up to 45 percent). Moreover, solutions are robust to service time and workload variations.
Danilo Ardagna, Barbara Panicucci, Marco Trubian, Li Zhang 0002
IEEE Trans. Serv. Comput.4
2011 Consolidating virtual machines with dynamic bandwidth demand in data centers
abstract
Recent advances in virtualization technology have made it a common practice to consolidate virtual machines(VMs) into a fewer number of servers. An efficient consolidation scheme requires that VMs are packed tightly, yet receive resources commensurate with their demands. However, measurements from production data centers show that the network bandwidth demands of VMs are dynamic, making it difficult to characterize the demands by a fixed value and to apply traditional consolidation schemes. In this work, we formulate the VM consolidation into a Stochastic Bin Packing problem and propose an online packing algorithm by which the number of servers required is within (1+∈)(√2+1) of the optimum for any ∈ >; 0. The result can be improved to within (√2+1) of the optimum in a special case. In addition, we use numerical experiments to evaluate the proposed consolidation algorithm and observe 30% server reduction compared to several benchmark algorithms.
Xiaoqiao Meng, Li Zhang 0002
INFOCOM3
2010 Autonomic Management of Cloud Service Centers with Availability Guarantees
abstract
Modern cloud infrastructures live in an open world, characterized by continuous changes in the environment and in the requirements they have to meet. Continuous changes occur autonomously and unpredictably, and they are out of control of the cloud provider. Therefore, advanced solutions have to be developed able to dynamically adapt the cloud infrastructure, while providing continuous service and performance guarantees. A number of autonomic computing solutions have been developed such that resources are dynamically allocated among running applications on the basis of short-term demand estimates. However, only performance and energy trade-off have been considered so far with a lower emphasis on the infrastructure dependability/availability which has been demonstrated to be the weakest link in the chain for early cloud providers. The aim of this paper is to fill this literature gap devising resource allocation policies for cloud virtualized environments able to identify performance and energy trade-offs, providing a priori availability guarantees for cloud end-users.
Bernardetta Addis, Danilo Ardagna, Barbara Panicucci, Li Zhang 0002
IEEE CLOUD4
2010 Improving the Scalability of Data Center Networks with Traffic-aware Virtual Machine Placement
abstract
The scalability of modern data centers has become a practical concern and has attracted significant attention in recent years. In contrast to existing solutions that require changes in the network architecture and the routing protocols, this paper proposes using traffic-aware virtual machine (VM) placement to improve the network scalability. By optimizing the placement of VMs on host machines, traffic patterns among VMs can be better aligned with the communication distance between them, e.g. VMs with large mutual bandwidth usage are assigned to host machines in close proximity. We formulate the VM placement as an optimization problem and prove its hardness. We design a two-tier approximate algorithm that efficiently solves the VM placement problem for very large problem sizes. Given the significant difference in the traffic patterns seen in current data centers and the structural differences of the recently proposed data center architectures, we further conduct a comparative analysis on the impact of the traffic patterns and the network architectures on the potential performance gain of traffic-aware VM placement. We use traffic traces collected from production data centers to evaluate our proposed VM placement algorithm, and we show a significant performance improvement compared to existing general methods that do not take advantage of traffic patterns and data center network characteristics.
Xiaoqiao Meng, Vasileios Pappas, Li Zhang 0002
INFOCOM3
2010 Program behavior characterization in large memory systems
abstract
We introduce models to characterize large cache performance in terms of various statistics related to sojourn time of a line in the cache. These statistics themselves depend on cache configuration parameters and we are currently working to isolate this dependency using LCS data and models. This will then help us in obtaining explicit relation between cache performance and its configuration parameters which will be helpful in identifying an optimal set of configuration parameters during early design phase of large memory systems.
Parijat Dube, Michael Tsao, Dan E. Poff, Li Zhang 0002, Alan Bivens
ISPASS4
2010 Linear-speed interior-path algorithms for distributed control of information networks
Hanhua Feng, Cathy H. Xia, Zhen Liu 0001, Li Zhang 0002
Perform. Evaluation4
2009 Connection and performance model driven optimization of pageview response time
abstract
Managing client perceived pageview response time for multiple classes of service is essential in today's highly competitive, e-commerce environment. We present Connection and Performance Model Driven Optimization (CP-MDO), a novel approach for providing optimal QoS as defined by a cost objective based on client perceived pageview response time and pageview drop rate. Our approach combines two vital models: (1) a latency model for connection establishment that captures the interactions between web browsers and web servers across network protocol layers and (2) a server performance model based on queueing theory that models performance across all tiers of a server complex. An algorithm capable of enforcing the optimal admission control based on the inter-arrival time between pageview admissions is given. Our approach has been implemented and evaluated in an experimental setting, demonstrating how CP-MDO achieves the minimal cost while providing minimal pageview response times under minimal drop rates across multiple classes of service.
David P. Olshefski, Li Zhang 0002
MASCOTS3
2009 Real-time performance modeling for adaptive software systems with multi-class workload
abstract
Modern, adaptive software systems must often adjust or reconfigure their architecture in order to respond to continuous changes in their execution environment. Efficient autonomic control in such systems is highly dependent on the accuracy of their representative performance model. In this paper, we are concerned with real-time estimation of a performance model for adaptive software systems that process multiple classes of transactional workload. Based on an open queueing network model and an Extended Kalman Filter (EKF), experiments in this work show that: (1) the model parameter estimates converge to the actual value very slowly when the variation in incoming workload is very low, (2) the estimates fail to converge quickly to the new value when there is a step-change caused by adaptive reconfiguration of the actual software parameters. We therefore propose a modified EKF design in which the measurement model is augmented with a set of constraints based on past measurement values. Experiments demonstrate the effectiveness of our approach that leads to significant improvement in convergence in the two cases.
Asser N. Tantawi, Li Zhang 0002
MASCOTS3
2008 Model Identification for Energy-Aware Management of Web Service Systems
Mara Tanelli, Danilo Ardagna, Marco Lovera, Li Zhang 0002
ICSOC4
2007 Scalability of the Nutch search engine
abstract
Nutch is an open source search engine that is gaining increasing popularity in the commercial world. The Nutch architecture leads itself to a wide range of parallelization techniques. Multiple backend servers can be used to both partition the corpus of search data, thus increasing the rate of queries serviced, and to increase the size of the search data while preserving the service rate. Alternatively, multiple search engines can operate in parallel, further increasing the query rate. In this paper, we analyze the performance and scalability of various configurations of Nutch. The configurations were implemented as part of the Commercial Scale Out project at IBM Research, and were used to investigate the applicability of scale-out architectures in commercial environments. We conclude that Nutch is highly scalable, with the different configurations behaving differently from a performance perspective.
José E. Moreira, Maged M. Michael, Dilma Da Silva, Doron Shiloach, Parijat Dube, Li Zhang 0002
ICS6
2007 Almost Peer-to-Peer Clock Synchronization
abstract
In this paper, an almost peer-to-peer (AP2P) clock synchronization protocol is proposed. AP2P is almost peer-to-peer in the sense that it provides the desirable features of a purely hierarchical (client/server) clock synchronization protocol while avoiding the undesirable consequences of a purely peer-to-peer one. In AP2P, a unique node is elected as a leader in a distributed manner. Each non-leader node adjusts its clock rate based on message exchanges with its neighbors, taking into consideration that neighbors that are closer to the leader have more effect on the adjustment than the neighbors that are further away from the leader. We compare the performance of AP2P with that of the server time protocol (STP), which is a purely hierarchical clock synchronization protocol. Simulation results, which have been conducted on several network topologies, have shown that AP2P can provide a clock synchronization accuracy that is indistinguishable from that of STP. Furthermore, AP2P is more fault-tolerant because it can recover from certain types of failures that STP cannot recover from.
Ahmed Sobeih, Michel Hack, Zhen Liu 0001, Li Zhang 0002
IPDPS4
2007 Performance Studies of a WebSphere Application, Trade, in Scale-out and Scale-up Environments
abstract
Scale-out approach, in contrast to scale-up approach (exploring increasing performance by utilizing more powerful shared-memory servers), refers to deployment of applications on a large number of small, inexpensive, but tightly packaged and tightly interconnected servers. Recently, there has been an increasing interest in scale-out approach. The purpose of this study is to discover advantages or disadvantages of scale-out systems with a typical enterprise workload, IBM Trade Performance Benchmark Sample for Websphere application server (a.k.a. Trade6). In this work, through cross system performance comparison, we show that for such workload, scale-out approach has better performance/cost effect. In term of scalability, we show that Websphere application server packages for distributed environment scale well while the possible bottleneck of the application deployment is the database tier. We present preliminary results to show that both database partitioning feature (DPF) and federated database server approaches are not exactly suitable for providing scale-out solution for the database tier of workloads similar to Trade (small tables and short transactions). In addition, we discuss our on-going effort on further performance study: (1) studies of performance/scalability for larger deployments by adopting the IBM AMBIENCE queuing network modeling tool, (2) performance breakdowns utilizing IBM ACTC hardware counter library.
Hao Yu 0008, José E. Moreira, Parijat Dube, I-Hsin Chung, Li Zhang 0002
IPDPS5
2007 Performance Evaluation of a Commercial Application, Trade, in Scale-out Environments
abstract
Scale-out approach, in contrast to scale-up approach (exploring increasing performance by utilizing more powerful shared-memory servers), refers to deployment of applications on a large number of small, inexpensive, but tightly packaged and tightly interconnected servers. The purpose of this study is to understand the performance of scale-out architectures with a typical enterprise workload, IBM Trade Performance Benchmark Sample for WebSphere Application Server (a.k.a. Trade). We describe a performance evaluation methodology that gives accurate predictions of application performance and system utilization by utilizing experimental data driven model development. Through experiments and extrapolation from the derived model, we show that for such workload, WebSphere Application Server packages for distributed environments scale well while the possible bottleneck of the application deployment is the database tier.
Parijat Dube, Hao Yu 0008, Li Zhang 0002, José E. Moreira
MASCOTS3
2007 SLA based resource allocation policies in autonomic environments
Danilo Ardagna, Marco Trubian, Li Zhang 0002
J. Parallel Distributed Comput.3
2007 Load shedding and distributed resource control of stream processing networks
Hanhua Feng, Zhen Liu 0001, Cathy H. Xia, Li Zhang 0002
Perform. Evaluation4
2006 Distributed Resource Allocation for Stream Data Processing
Ao Tang, Zhen Liu 0001, Cathy H. Xia, Li Zhang 0002
HPCC4
2006 Distributed Resource Allocation in Stream Processing Systems
Cathy H. Xia, James Broberg, Zhen Liu 0001, Li Zhang 0002
DISC4
2005 SLA Based Profit Optimization in Multi-tier Systems
abstract
Nowadays, large service centers provide computational capacity to many customers by sharing a pool of IT resources. The service providers and their customers negotiate utility based service level agreement (SLA) to determine the costs and penalties on the base of the achieved performance level. The system is often based on a multi-tier architecture to service requests. The service provider would like to maximize the SLA revenues, while minimizing its operating costs. The system we consider is based on a centralized network dispatcher which controls the allocation of applications to servers, the request volumes at various servers and the scheduling policy at each server. The dispatcher can also decide to turn ON or OFF servers depending on the system load. This paper designs a resource allocation scheduler for such multi-tier environments so as to maximize the profits associated with multiple class SLAs. The overall problem is NP-hard. We develop heuristic solutions by implementing a local-search algorithm. Results are presented to demonstrate the benefits of our approach
Danilo Ardagna, Marco Trubian, Li Zhang 0002
NCA3
2005 Optimal capacity allocation for Web systems with end-to-end delay guarantees
Wuqin Lin, Zhen Liu 0001, Cathy H. Xia, Li Zhang 0002
Perform. Evaluation4
2005 Web traffic modeling at finer time scales and performance implications
Cathy H. Xia, Zhen Liu 0001, Mark S. Squillante, Li Zhang 0002, Naceur Malouch
Perform. Evaluation4
2004 Overlay Multicast Trees of Minimal Delay
abstract
Overlay multicast (or application-level multicast) has become an increasingly popular alternative to IP-supported multicast. End nodes participating in overlay multicast can form a directed tree rooted at the source using existing unicast links. For each receiving node there is always only one incoming link. Very often, nodes can support no more than a fixed number of outgoing links due to bandwidth constraints. Here, we describe an algorithm for constructing a multicast tree with the objective of minimizing the maximum communication delay (i.e. the longest path in the tree), while satisfying degree constraints at nodes. We show that the algorithm is a constant-factor approximation algorithm. We further prove that the algorithm is asymptotically optimal if the communicating nodes can be mapped into Euclidean space such that the nodes are uniformly distributed in a convex region. We evaluate the performance of the algorithm using randomly generated configurations of up to 5,000,000 nodes.
Anton Riabov, Zhen Liu 0001, Li Zhang 0002
ICDCS3
2004 SLA based profit optimization in autonomic computing systems
abstract
With the development of the Service Oriented Architecture (SOA), organizations are able to compose complex applications from distributed services supported by third party providers. Under this scenario, large data centers provide services to many customers by sharing available IT resources. This leads to the efficient use of resources and the reduction of operating costs. Service providers and their customers often negotiate utility based Service Level Agreements (SLAs) to determine costs and penalties based on the achieved performance levels. Data centers often employ an autonomic computing infrastructure and use a centralized dispatch and control component (a dispatcher) to distribute the user requests to backend servers, and to set the scheduling policies at each server. This dispatcher can also decide to turn ON or OFF servers depending on the system load. This paper designs a set of dispatching and control policies for the dispatcher in such service oriented environments. The objective is to maximize the provider's profits associated with multiple class of SLAs. We show that the overall problem is NP-hard, and develop meta-heuristic solutions based on the tabu-search algorithm. Experimental results are presented to show the benefits of our approach.
Li Zhang 0002, Danilo Ardagna
ICSOC1
2004 A smart hill-climbing algorithm for application server configuration
abstract
The overwhelming success of the Web as a mechanism for facilitating information retrieval and for conducting business transactions has ledto an increase in the deployment of complex enterprise applications. These applications typically run on Web Application Servers, which assume the burden of managing many tasks, such as concurrency, memory management, database access, etc., required by these applications. The performance of an Application Server depends heavily on appropriate configuration. Configuration is a difficult and error-prone task dueto the large number of configuration parameters and complex interactions between them. We formulate the problem of finding an optimal configuration for a given application as a black-box optimization problem. We propose a smart hill-climbing algorithm using ideas of importance sampling and Latin Hypercube Sampling (LHS). The algorithm is efficient in both searching and random sampling. It consists of estimating a local function, and then, hill-climbing in the steepest descent direction. The algorithm also learns from past searches and restarts in a smart and selective fashion using the idea of importance sampling. We have carried out extensive experiments with an on-line brokerage application running in a WebSphere environment. Empirical results demonstrate that our algorithm is more efficient than and superior to traditional heuristic methods.
Bowei Xi, Zhen Liu 0001, Mukund Raghavachari, Cathy H. Xia, Li Zhang 0002
WWW5
2004 Efficiently serving dynamic data at highly accessed web sites
abstract
We present architectures and algorithms for efficiently serving dynamic data at highly accessed Web sites together with the results of an analysis motivating our design and quantifying its performance benefits. This includes algorithms for keeping cached data consistent so that dynamic pages can be cached at the Web server and dynamic content can be served at the performance level of static content. We show that our system design is able to achieve cache hit ratios close to 100% for cached data which is almost never obsolete by more than a few seconds, if at all. Our architectures and algorithms provide more than an order of magnitude improvement in performance using an order of magnitude fewer servers over that obtained under conventional methods.
Jim Challenger, Paul Dantzig, Arun Iyengar, Mark S. Squillante, Li Zhang 0002
IEEE/ACM Trans. Netw.5
2003 New Algorithms for Content-Based Publication-Subscription Systems
abstract
This paper introduces new algorithms specifically designed for content-based publication-subscription systems. These algorithms can be used to determine multicast groups with as much commonality as possible, based on the totality of subscribers' interests. The algorithms are based oil concepts borrowed from the literature on spatial databases and clustering. These algorithms perform well in the context of highly heterogeneous subscriptions, and they also scale well. Based on concepts borrowed from the spatial database literature, we develop an algorithm to match publications to subscribers in real-time. We also investigate the benefits of dynamically determining whether to unicast, multicast or broadcast information about the events over the network to the matched subscribers. We call this the distribution method problem. Some of these same concepts can be applied to match publications to subscribers in real-time, and also to determine dynamically whether to unicast, multicast or broadcast information about the events over the network to the matched subscribers. We demonstrate the quality of our algorithms via a number of realistic simulation experiments.
Anton Riabov, Zhen Liu 0001, Joel L. Wolf, Philip S. Yu, Li Zhang 0002
ICDCS5
2002 Analysis of measurement data from sporting event Web sites
abstract
With the growing popularity of Web applications, there is a considerable increase in the importance of managing Web sites to deliver high levels of performance and scalability to accommodate future growth and evolution. One of the key issues in this regard concerns a better understanding of the traffic patterns at multiple levels, such as the levels of requests, pages and sessions. This paper presents a detailed analysis of measurement data from various sources pertaining to a specific multi-tiered, geographically distributed architecture that has been used to host the Web sites for a number of recent, popular sporting events. Our analysis of the request-level and page-level patterns demonstrate differences among the Web sites depending upon the type of event and the breadth of interests of the user community. Some of these patterns are consistent with commercial Web sites, while others are significantly different These results further illustrate geographical differences in the request-level and page-level patterns. Our analysis also investigates in detail session-level characteristics. This includes an analysis of the session durations, the think time distributions, the dependence structure of the session arrival process, and the page views comprising each session.
Zhen Liu 0001, Mark S. Squillante, Cathy H. Xia, S.-Z. Yu, Li Zhang 0002, Naceur Malouch, Paul Dantzig
GLOBECOM5
2002 Clustering Algorithms for Content-Based Publication-Subscription Systems
abstract
We consider efficient communication schemes based on both network-supported and application-level multicast techniques for content-based publication-subscription systems. We show that the communication costs depend heavily on the network configurations, distribution of publications and subscriptions. We devise new algorithms and adapt existing partitional data clustering algorithms. These algorithms can be used to determine multicast groups with as much commonality as possible, based on the totality of subscribers' interests. They perform well in the context of highly heterogeneous subscriptions, and they also scale well. An efficiency of 60% to 80% with respect to the ideal solution can be achieved with a small number of multicast groups (less than 100 in our experiments). Some of these same concepts can be applied to match publications to subscribers in real-time, and also to determine dynamically whether to unicast, multicast or broadcast information about the events over the network to the matched subscribers. We demonstrate the quality of our algorithms via simulation experiments.
Anton Riabov, Zhen Liu 0001, Joel L. Wolf, Philip S. Yu, Li Zhang 0002
ICDCS5
2002 Clock Synchronization Algorithms for Network Measurements
abstract
Packet delay traces are important measurements for analyzing end-to-end performance and for designing traffic control algorithms in computer networks. Due to the fact that the clocks at the end systems are usually not synchronized and running at different speeds, these measurements can be quite inaccurate. We propose several algorithms to estimate and remove the relative clock skews from delay measurements based on the computation of convex hulls. Compared with existing techniques, such as linear regression and linear programming, the convex-hull approach provides better insight and allows us to handle more error metrics. We obtain algorithms which are linear in the number of measurement points for the case with no clock resets. For the more challenging case with clock resets, i.e., the clocks are reset to some reference times during the measurement period, we develop linear algorithms to identity the clock resets, and derive the best clock skew lines. We extend this analysis to environments in which at least one of the clocks is controlled by NTP (network time protocol). These algorithms can greatly improve the accuracy of the measurements, and can be used both online and offline. They can also be extended for active clock synchronization, to replace or further improve NTP. Numerical experiments are presented to demonstrate the robustness of the algorithms.
Li Zhang 0002, Zhen Liu 0001, Cathy H. Xia
INFOCOM1
2002 Optimal scheduling in queuing network models of high-volume commercial web sites
Mark S. Squillante, Cathy H. Xia, Li Zhang 0002
Perform. Evaluation3
1999 Analysis of Job Arrival Patterns and Parallel Scheduling Performance
Mark S. Squillante, David D. Yao, Li Zhang 0002
Perform. Evaluation3
1999 Analysis and Characterization of Large-Scale Web Server Access Patterns and Performance
Arun Iyengar, Mark S. Squillante, Li Zhang 0002
World Wide Web3
1998 A General Methodology for Characterizing Access Patterns and Analyzing Web Server Performance
abstract
We develop a general methodology for characterizing Web server access patterns based on a spectral analysis of finite collections of observed data from real systems. Our approach is used together with the access logs from the IBM Web site for the 1996 Olympic Games to demonstrate some of its advantages over previous methods and to analyze certain aspects of large-scale Web server performance.
Arun Iyengar, Edward A. MacNair, Mark S. Squillante, Li Zhang 0002
MASCOTS4