EDBT 2026 Demo / reviewers in the wild / expert
Ludmila Cherkasova
dblp:88/4130 · also Lucy Cherkasova
· DBLP profile ↗
96ranked-venue papers
31as first author
3since 2021 · last 2023
0000-0002-9333-4901ORCID · verified
Domains — the database's venue-derived domains; a paper can count in several
Systems, architecture and hardware · 32 · 12 first-authorComputer networks · 21 · 8 first-author · 2 since 2021Software engineering, systems software and programming languages · 18 · 1 first-authorTheory of computation · 8 · 5 first-authorSecurity and privacy · 7 · 3 first-authorDatabases, data management, data science and information retrieval · 3 · 1 first-authorGraphics, computer vision, multimedia, augmented reality and games · 3 · 3 first-authorApplied, interdisciplinary, general and emerging computing · 3Artificial intelligence and machine learning · 1 · 1 first-author
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
14 papers |
Embedded and real-time systems · 34% Performance modeling and evaluation · 21% Cloud and datacenter computing · 21% | |
| Computer networks
7 papers |
Network management and operations · 49% Internet of things and sensor networks · 42% Content delivery and video streaming · 5% | |
| Databases, data mining, and information retrieval
1 paper |
Information retrieval · 77% Data integration and cleaning · 23% | |
| Artificial intelligence
1 paper |
Time series and sequential data · 100% |
Topics — the 30 heaviest of 40, each with the papers that count most for it
| Topic | Weight | Papers | Last | Evidence papers |
|---|---|---|---|---|
Network management and operations › network robustness › fault tolerance
proactive fault tolerance |
0.7 | 1 | 2023 | DeepFT: Fault-Tolerant Edge Computing using a Self-Supervised Deep Surrogate Model · INFOCOM 2023 |
Embedded and real-time systems
real-time scheduling |
0.7 | 1 | 2023 | DeepFT: Fault-Tolerant Edge Computing using a Self-Supervised Deep Surrogate Model · INFOCOM 2023 |
Internet of things and sensor networks › wireless sensor network
sensor deployment |
0.4 | 1 | 2020 | Optimizing Sensor Deployment and Maintenance Costs for Large-Scale Environmental Monitoring · IEEE Trans. Comput. Aided Des. Integr. Circuits Syst. 2020 |
Embedded and real-time systems
embedded software |
0.4 | 1 | 2020 | eWASM: Practical Software Fault Isolation for Reliable Embedded Devices · IEEE Trans. Comput. Aided Des. Integr. Circuits Syst. 2020 |
Cloud and datacenter computing
cluster resource management and scheduling |
0.3 | 3 | 2013 | Orchestrating an Ensemble of MapReduce Jobs for Minimizing Their Makespan · IEEE Trans. Dependable Secur. Comput. 2013 Efficient resource allocation and power saving in multi-tiered systems · WWW 2010 Capacity planning tool for streaming media services · ACM Multimedia 2003 |
Performance modeling and evaluation
workload characterization |
0.3 | 5 | 2012 | Dealing with Burstiness in Multi-Tier Applications: Models and Their Parameterization · IEEE Trans. Software Eng. 2012 Analysis of enterprise media server workloads: access patterns, locality, content evolution, and rates of change · IEEE/ACM Trans. Netw. 2004 Capacity planning tool for streaming media services · ACM Multimedia 2003 |
Machine learning › Time series and sequential data
fault detection and diagnosis |
0.2 | 1 | 2023 | DeepFT: Fault-Tolerant Edge Computing using a Self-Supervised Deep Surrogate Model · INFOCOM 2023 |
Electronic design automation › high-level synthesis › scheduling
makespan minimization |
0.2 | 1 | 2013 | Orchestrating an Ensemble of MapReduce Jobs for Minimizing Their Makespan · IEEE Trans. Dependable Secur. Comput. 2013 |
Cloud and datacenter computing › cluster resource management and scheduling › cluster scheduling
mapreduce scheduling |
0.2 | 1 | 2013 | Orchestrating an Ensemble of MapReduce Jobs for Minimizing Their Makespan · IEEE Trans. Dependable Secur. Comput. 2013 |
Performance modeling and evaluation
performance prediction |
0.1 | 1 | 2012 | Dealing with Burstiness in Multi-Tier Applications: Models and Their Parameterization · IEEE Trans. Software Eng. 2012 |
Performance modeling and evaluation
queueing models |
0.1 | 1 | 2012 | Dealing with Burstiness in Multi-Tier Applications: Models and Their Parameterization · IEEE Trans. Software Eng. 2012 |
Internet of things and sensor networks › wireless sensor network
environmental monitoring |
0.1 | 1 | 2020 | Optimizing Sensor Deployment and Maintenance Costs for Large-Scale Environmental Monitoring · IEEE Trans. Comput. Aided Des. Integr. Circuits Syst. 2020 |
Energy-efficient computing
datacenter power management |
0.1 | 1 | 2010 | Lightning: self-adaptive, energy-conserving, multi-zoned, commodity green cloud storage system · HPDC 2010 |
Cloud and datacenter computing › resource provisioning
dynamic resource provisioning |
0.1 | 1 | 2010 | Efficient resource allocation and power saving in multi-tiered systems · WWW 2010 |
Storage systems
energy-efficient storage |
0.1 | 1 | 2010 | Lightning: self-adaptive, energy-conserving, multi-zoned, commodity green cloud storage system · HPDC 2010 |
Energy-efficient computing
power management |
0.1 | 1 | 2010 | Efficient resource allocation and power saving in multi-tiered systems · WWW 2010 |
Energy-efficient computing › low-power design
power optimization |
0.1 | 1 | 2010 | Efficient resource allocation and power saving in multi-tiered systems · WWW 2010 |
Information retrieval › text matching
document matching |
0.1 | 1 | 2009 | Applying syntactic similarity algorithms for enterprise information management · KDD 2009 |
Information retrieval › similarity measure
document similarity |
0.1 | 1 | 2009 | Applying syntactic similarity algorithms for enterprise information management · KDD 2009 |
Data integration and cleaning
entity resolution |
0.1 | 1 | 2009 | Applying syntactic similarity algorithms for enterprise information management · KDD 2009 |
Information retrieval › similarity search
near-duplicate detection |
0.1 | 1 | 2009 | Applying syntactic similarity algorithms for enterprise information management · KDD 2009 |
Distributed systems
anomaly detection |
0.1 | 1 | 2009 | Automated anomaly detection and performance modeling of enterprise applications · ACM Trans. Comput. Syst. 2009 |
Performance modeling and evaluation
capacity planning |
0.1 | 2 | 2009 | Capacity planning tool for streaming media services · ACM Multimedia 2003 Automated anomaly detection and performance modeling of enterprise applications · ACM Trans. Comput. Syst. 2009 |
Cloud and datacenter computing
virtualization |
0.1 | 1 | 2005 | Measuring CPU Overhead for I/O Processing in the Xen Virtual Machine Monitor · USENIX ATC, General Track 2005 |
Performance modeling and evaluation
benchmarking |
0.1 | 2 | 2005 | Measuring the capacity of a streaming media server in a Utility Data Center environment · ACM Multimedia 2002 Measuring CPU Overhead for I/O Processing in the Xen Virtual Machine Monitor · USENIX ATC, General Track 2005 |
Distributed systems › distributed system architecture
multi-tier application |
0.0 | 1 | 2012 | Dealing with Burstiness in Multi-Tier Applications: Models and Their Parameterization · IEEE Trans. Software Eng. 2012 |
Network measurement and analytics › network performance measurement
internet performance monitoring |
0.0 | 1 | 2002 | EtE: Passive End-to-End Internet Service Performance Monitoring · USENIX ATC, General Track 2002 |
Embedded and real-time systems › real-time scheduling
admission control |
0.0 | 1 | 2002 | Session-Based Admission Control: A Mechanism for Peak Load Management of Commercial Web Sites · IEEE Trans. Computers 2002 |
Cloud and datacenter computing
overload control |
0.0 | 1 | 2002 | Session-Based Admission Control: A Mechanism for Peak Load Management of Commercial Web Sites · IEEE Trans. Computers 2002 |
Storage systems › storage hierarchy
tiered storage |
0.0 | 1 | 2010 | Lightning: self-adaptive, energy-conserving, multi-zoned, commodity green cloud storage system · HPDC 2010 |
Methods — techniques the papers use, named apart from their topics
self-supervised learning · 2.0deep surrogate model · 2.0co-simulation · 2.0webassembly compilation · 0.9software sandboxing · 0.9sparse nonlinear optimizer · 0.4particle swarm optimization · 0.4mutual information · 0.4artificial bee colony · 0.4simulation · 0.3johnson algorithm · 0.2balancedpools heuristic · 0.2markov-modulated process · 0.1index of dispersion · 0.1shingling · 0.1fingerprint sampling · 0.1content-based chunking · 0.1stochastic petri nets · 0.0
| Year | Publication | Venue | Position |
|---|---|---|---|
| 2023 | DeepFT: Fault-Tolerant Edge Computing using a Self-Supervised Deep Surrogate ModelabstractThe emergence of latency-critical AI applications has been supported by the evolution of the edge computing paradigm. However, edge solutions are typically resource-constrained, posing reliability challenges due to heightened contention for compute capacities and faulty application behavior in the presence of overload conditions. Although a large amount of generated log data can be mined for fault prediction, labeling this data for training is a manual process and thus a limiting factor for automation. Due to this, many companies resort to unsupervised fault-tolerance models. Yet, failure models of this kind can incur a loss of accuracy when they need to adapt to non-stationary workloads and diverse host characteristics. Thus, we propose a novel modeling approach, DeepFT, to proactively avoid system overloads and their adverse effects by optimizing the task scheduling decisions. DeepFT uses a deep-surrogate model to accurately predict and diagnose faults in the system and co-simulation based self-supervised learning to dynamically adapt the model in volatile settings. Experimentation on an edge cluster shows that DeepFT can outperform state-of-the-art methods in fault-detection and QoS metrics. Specifically, DeepFT gives the highest F1 scores for fault-detection, reducing service deadline violations by up to 37% while also improving response time by up to 9%. Shreshth Tuli, Giuliano Casale, Ludmila Cherkasova, Nicholas R. Jennings |
INFOCOM | 3 |
| 2023 | Automating and Optimizing Reliability-Driven Deployment in Energy-Harvesting IoT NetworksabstractRecent years have witnessed a significant expansion in Internet-of-Things (IoT) applications. Although the battery energy availability can be improved with energy harvesting, the overall device reliability management has been overlooked in the existing literature. State-of-the-art reliability models of solar panels, electronics and rechargeable batteries show exponential dependence of failures on temperature. This work is the first to develop a comprehensive reliability deployment framework for energy-harvesting IoT networks, reflecting the non-negligible thermal stresses on each hardware component. Our framework improves the reliability on both pre-deployment and post-deployment stages. Prior to deployment, given the historical temperature and solar radiation of the region, we formulate a Mixed Integer Linear Program (MILP) to place the minimum number of nodes, while ensuring (i) full target coverage, (ii) complete connectivity, (iii) energy-neutral operation, and (iv) reliability constraints at each deployed node. We propose a polynomial-time heuristic, R-TSH, to approximate the optimal placement in large-scale deployments. While R-TSH optimizes long-term reliability, the prompt temperature or link quality differences from the historical patterns can significantly degrade device reliability after deployment. The post-deployment section of our design consists of a reliability-driven routing algorithm, AODV-Rel, that adapts to real-time environmental and link quality changes. Extensive analysis is done using a real-world dataset from the National Solar Radiation Database. Simulations in ns-3 show that R-TSH meets all reliability constraints even after 5 years of deployment as compared to the state of the art. In addition, it is 2000x faster than the optimal solution, while placing only 28% more nodes. AODV-Rel further extends the minimal operational lifetime by 1.5 and 2.8 months under temperature deviation and wireless interference. Xiaofan Yu 0001, Kazim Ergun, Xueyang Song, Ludmila Cherkasova, Tajana Rosing |
IEEE Trans. Netw. Serv. Manag. | 4 |
| 2021 | Automating Reliable and Fault-Tolerant Design of LoRa-based IoT NetworksabstractLow-Power Wide-Area Networks (LPWAN) has recently been scaling rapidly, targeting at large-scale and low-power applications. LoRa and LoRaWAN have been adopted in many practical deployments. While the advantages of LoRa have been well-demonstrated, challenges in scalability and reliability impede LoRa networks from further expansion. Traditional deployment strategies for reliability and fault tolerance, which ensure that multiple networking paths are available, are not directly applicable because of LoRa's single-hop and Aloha medium-access design. In this paper, we study how to design LoRa networks in large regions so that their transmission reliability and fault tolerance against gateway failures and interference are met for LoRaWAN technology. We first introduce m-gateway connectivity to guarantee fault tolerance due to LoRa's unique properties. Next, we leverage state-of-the-art transmission reliability model based on estimated path loss from satellite maps. Combining the above two contributions, we formulate an Integer Nonlinear Program (INLP) that minimizes the number of gateways through strategic gateway placement and resource allocation. Constraints are imposed to achieve (i) fault tolerance, (ii) reliable transmission (i.e., satisfactory QoS), (iii) sufficiently long lifetime. Due to high complexity of INLP, we design a greedy heuristic, RFT-LoRa, to acquire a high-quality solution for larger size problems. Comprehensive evaluation is performed with ns-3 simulator using real-world datasets. The results demonstrate that RFT-LoRa enhances average packet delivery ratio by 10% - 54% over the existing heuristic under gateway failures and interference. Xiaofan Yu 0001, Ludmila Cherkasova, Tajana Rosing |
CNSM | 3 |
| 2020 | Reliability-Driven Deployment in Energy-Harvesting Sensor NetworksabstractRecent years have witnessed a significant expansion in Internet-of-Things (IoT) applications, especially in environmental monitoring, which aims at providing full coverage over potential targets. With energy harvesting ability, sensor devices can be replenished by external energy sources, and thus their lifetime is prolonged. While existing literature focuses on minimizing deployment costs, the reliability management is overlooked. Previous research has addressed that a higher temperature exponentially accelerates hardware failure rates. The versatile outdoor environments impose a non-negligible thermal stress on the hardware and consequently reduce the reliability of devices. In this paper, we are the first to propose a reliability-driven sensor deployment approach to achieve minimum nodes, while satisfying (i) full target coverage, (ii) complete connectivity, (iii) energy-neutral operation, and (iv) reliability constraints. Given external temperature distribution, we propose an algorithm to convert reliability constraints to a single-value power threshold for each location. A Mixed Integer Linear Programming (MILP) model is formulated and solved with CPLEX. Due to the complex nature of MILP, we propose a heuristic, named Reliability-driven TwoStage Heuristic (R-TSH), to approximate the optimal solution for large-scale problems. Extensive simulations are performed on a real-world dataset from the National Solar Radiation Database. Our results indicate that R-TSH meets all reliability constraints with only 20% more sensors than the optimal solution, while executing more than 1500x faster. Compared to state-of-the-art heuristics, R-TSH avoids 20 - 80% of reliability violations with a comparable number of nodes and execution time. Xiaofan Yu 0001, Xueyang Song, Ludmila Cherkasova, Tajana Rosing |
CNSM | 3 |
| 2020 | Sledge: a Serverless-first, Light-weight Wasm Runtime for the EdgeabstractEmerging IoT applications with real-time latency constraints require new data processing systems operating at the Edge. Serverless computing offers a new compelling paradigm, where a user can execute a small application without handling the operational issues of server provisioning and resource management. Despite a variety of existing commercial and open source serverless platforms (utilizing VMs and containers), these solutions are too heavy-weight for a resource-constrained Edge systems (due to large memory footprint and high invocation time). Moreover, serverless workloads that focus on per-client, short-running computations are not an ideal fit for existing general purpose computing systems. Phani Kishore Gadepalli, Sean McBride, Gregor Peach, Ludmila Cherkasova, Gabriel Parmer |
Middleware | 4 |
| 2020 | eWASM: Practical Software Fault Isolation for Reliable Embedded DevicesabstractAs we connect more microcontrollers to the Internet and employ them to control the physical world around us, their reliability and security are increasingly important. Many microcontrollers provide limited facilities for hardware isolation, and real-time OSes offer custom APIs, that require coupling applications into the ecosystem and abstractions of that specific OS to leverage isolation. This article investigates the use of software sandboxing of applications to support isolation for resource-constrained devices. Toward this, we detail the design of eWASM, a processes abstraction that adapts a popular sandbox, Wasm, for microcontrollers. eWASM provides a runtime to constrain memory accesses and control flow, enabled by our aWsm Wasm compiler. We discuss and evaluate its multiple implementations that effectively trade time and space, optimizing for the constraints of embedded systems. This enables popular languages (e.g., C) to be effectively sandboxed by software. We demonstrate performance within 40% of native C on Polybench. We believe this is a practical and compelling result for many IoT domains, and it represents the first compiled sandboxing environment for microcontrollers. We show that restrictions of the current Wasm specification lead to significant memory consumption and provide suggestions for the creation of an embedded-specific Wasm variant. Gregor Peach, Runyu Pan, Zhuoyi Wu, Gabriel Parmer, Christopher Haster, Ludmila Cherkasova |
IEEE Trans. Comput. Aided Des. Integr. Circuits Syst. | 6 |
| 2020 | Optimizing Sensor Deployment and Maintenance Costs for Large-Scale Environmental MonitoringabstractRecent advances in low-power long-range communication schemes such as LoRa have opened up new potentials in large-scale Internet-of-Things (IoT) applications, especially environmental monitoring. However, the versatile environment and the long traveling distance have imposed significant challenges to maintenance. Previous research has shown that higher temperature exponentially accelerates electronics failure rates. The maintenance cost can take as much as 80% of the total deployment expenses if not managed carefully. In this article, we formulate a sensor deployment problem to preventively minimize maintenance costs while ensuring tolerable sensing quality and complete connectivity. We are the first to derive a maintenance cost model for IoT networks considering thermal degradation and battery depletion. To assess the spatial phenomena of interest, we adopt the sensing quality metric based on mutual information. While the proposed problem is nonconvex, we bring up a relaxed form and solve it with a sparse nonlinear optimizer. We further apply two population-based metaheuristics, i.e., particle swarm optimization (PSO) and artificial bee colony (ABC) algorithm, to approximate the optimal solution. Extensive simulations are performed on two real-world datasets of the Southern California region in the U.S. Our metaheuristics save up to 40% of maintenance cost compared with the existing greedy heuristics under the same acceptable sensing quality. Xiaofan Yu 0001, Kazim Ergun, Ludmila Cherkasova, Tajana Rosing |
IEEE Trans. Comput. Aided Des. Integr. Circuits Syst. | 3 |
| 2019 | Analysis and Demand Forecasting of Residential Energy Consumption at Multiple Time Scales
Poojitha Amin, Ludmila Cherkasova, Robert C. Aitken, Vikas Kache |
IM | 2 |
| 2019 | Challenges and Opportunities for Efficient Serverless Computing at the EdgeabstractServerless computing frameworks allow users to execute a small application (dedicated to a specific task) without handling operational issues such as server provisioning, resource management, and resource scaling for the increased load. Serverless computing originally emerged as a Cloud computing framework, but might be a perfect match for IoT data processing at the Edge. However, the existing serverless solutions, based on VMs and containers, are too heavy-weight (large memory footprint and high function invocation time) for operating efficiency and elastic scaling at the Edge. Moreover, many novel IoT applications require low-latency data processing and near real-time responses, which makes the current cloud-based serverless solutions unsuitable. Recently, WebAssembly (Wasm) has been proposed as an alternative method for running serverless applications at near-native speeds, while having a small memory footprint and optimized invocation time. In this paper, we discuss some existing serverless solutions, their design details, and unresolved performance challenges for an efficient serverless management at the Edge. We outline our serverless framework, called aWsm, based on the WebAssembly approach, and discuss the opportunities enabled by the aWsm design, including function profiling and SLO-driven performance management of users' functions. Finally, we present an initial assessment of aWsm performance featuring average startup time (12μs to 30μs) and an economical memory footprint (ranging from 10s to 100s of kB) for a subset of MiBench microbenchmarks used as functions. Phani Kishore Gadepalli, Gregor Peach, Ludmila Cherkasova, Robert C. Aitken, Gabriel Parmer |
SRDS | 3 |
| 2018 | ProfDP: A Lightweight Profiler to Guide Data Placement in Heterogeneous Memory SystemsabstractNew memory technologies, such as non-volatile memory and stacked memory, have reformed the memory hierarchies in modern and emerging computer architectures. It becomes common to see memories of different types integrated into the same system, as known as heterogeneous memory. Typically, a heterogeneous memory system consists of a small fast component and a large slow component. This encourages new style of data processing and exposes developers with a new problem: given two memory types, how shall we redesign applications to benefit from this memory arrangement and decide on the efficient data placement? Existing methods perform detailed memory access pattern analysis to guide data placement. However, these methods are heavyweight and ignore the interactions between software and hardware. Shasha Wen, Ludmila Cherkasova, Felix Xiaozhu Lin, Xu Liu 0001 |
ICS | 2 |
| 2018 | Evaluating Scalability and Performance of a Security Management Solution in Large Virtualized EnvironmentsabstractVirtualized infrastructure is a key capability of modern enterprise data centers and cloud computing, enabling a more agile and dynamic IT infrastructure with fast IT provisioning, simplified, automated management, and flexible resource allocation to handle a broad set of workloads. However, at the same time, virtualization introduces new challenges, since securing virtual servers is more difficult than physical machines. HyTrust Inc. has developed an innovative security solution, called HyTrust Cloud Control (HTCC), to mitigate risks associated with virtualization and cloud technologies. HTCC is a virtual appliance deployed as a transparent proxy in front of a VMware-based virtualized environment. Since HTCC serves as a gateway to a customer virtualized environment, it is important to carefully assess its performance and scalability as well as provide its accurate resource sizing. In this work, we introduce a novel approach for accomplishing this goal. First, we describe a special framework, based on a nested virtualization technique, which enables the creation and deployment of a large scale virtualized environment (with 30,000 VMs) using a limited number of physical servers (4 servers in our experiments). Second, we introduce a design and implementation of a novel, extensible benchmark, called HT-vmbench, that allows to mimic the session-based activities of different system administrators and users in virtualized environments. The benchmark is implemented using VMware Web Service SDK. By executing HT-vmbench in the emulated large-scale virtualized environments, we can support an efficient performance assessment of management and security solutions (such as HTCC), their overhead, and provide capacity planning rules and resource sizing recommendations. Lishan Yang 0001, Ludmila Cherkasova, Rajeev Badgujar, Jack Blancaflor, Rahul Konde, Jason Mills, Evgenia Smirni |
ICPE | 2 |
| 2017 | Sparkle: optimizing spark for large memory machines and analyticsabstractGiven the growing availability of affordable scale-up servers, our goal is to bring the performance benefits of in-memory processing on scale-up servers to an increasingly common class of data analytics applications that process small to medium size datasets (up to a few 100GBs) that can easily fit in the memory of a typical scale-up server To achieve this goal, we leverage Spark, an existing memory-centric data analytics framework with wide-spread adoption among data scientists. Bringing Spark's data analytic capabilities to a scale-up system requires rethinking the original design assumptions, which, although effective for a scale-out system, are a poor match to a scale-up system resulting in unnecessary communication and memory inefficiencies. Mijung Kim, Jun Li 0008, Haris Volos 0001, Manish Marwah, Alexander Ulanov, Kimberly Keeton, Joseph A. Tucek, Ludmila Cherkasova, Pradeep Fernando |
SoCC | 8 |
| 2017 | Predictive modeling and scalability analysis for large graph analyticsabstractMany HPC and modern large graph processing applications belong to a class of scale-out applications, where the application dataset is partitioned and processed by a cluster of machines. Assessing the application scalability is one of the primary goals during such application implementation. Typically, in the design phase, programmers are limited by a small size cluster available for their experiments. Therefore, predictive modeling is required for the analysis of the application scalability and its performance in a larger cluster. While in an increased size cluster, each node will process a smaller portion of the original dataset, a higher communication volume between a larger number of nodes may cripple the application scalability and provide diminishing performance benefits. One of the main challenges is the analysis of bandwidth demands due to an increased communication volume in a larger size cluster. In this paper1, we introduce a novel regression-based approach to assess the scalability and performance of a distributed memory program for execution in a large-scale cluster. Our solution involves 1) a limited set of traditional experiments performed in a small size cluster and 2) an additional set of similar experiments performed with an “interconnect bandwidth throttling” tool, which exposes the bandwidth impact on the application performance. These measurements are used in creating an ensemble of analytical models for performance and scalability analysis. Using a linear regression approach, step by step, we incorporate into the model the following important parameters: i) the number of cluster nodes and application processes, ii) the dataset size, and iii) interconnect bandwidth. We demonstrate our solution, its power, and accuracy using a popular Graph500 benchmark, which implements a Breadth First Search algorithm on large, synthetically generated graphs. By utilizing measurements collected in a 32-node cluster, we are able to project the program performance in a large size cluster with hundreds of nodes. The proposed approach and derived models help to provide an early feedback to programmers on the scalability and efficiency of their solution. Sourav Medya, Ludmila Cherkasova, Ambuj K. Singh |
IM | 2 |
| 2017 | DyScale: A MapReduce Job Scheduler for Heterogeneous Multicore ProcessorsabstractThe functionality of modern multi-core processors is often driven by a given power budget that requires designers to evaluate different decision trade-offs, e.g., to choose between many slow, power-efficient cores, or fewer faster, power-hungry cores, or a combination of them. Here, we prototype and evaluate a new Hadoop scheduler, called DyScale, that exploits capabilities offered by heterogeneous cores within a single multi-core processor for achieving a variety of performance objectives. A typical MapReduce workload contains jobs with different performance goals: large, batch jobs that are throughput oriented, and smaller interactive jobs that are response time sensitive. Heterogeneous multi-core processors enable creating virtual resource pools based on "slow" and "fast" cores for multi-class priority scheduling. Since the same data can be accessed with either "slow" or "fast" slots, spare resources (slots) can be shared between different resource pools. Using measurements on an actual experimental setting and via simulation, we argue in favor of heterogeneous multi-core processors as they achieve "faster" (up to 40 percent) processing of small, interactive MapReduce jobs, while offering improved throughput (up to 40 percent) for large, batch jobs. We evaluate the performance benefits of DyScale versus the FIFO and Capacity job schedulers that are broadly used in the Hadoop community. Feng Yan 0001, Ludmila Cherkasova, Zhuoyao Zhang, Evgenia Smirni |
IEEE Trans. Cloud Comput. | 2 |
| 2016 | Parallel Graph Processing on Modern Multi-core Servers: New Findings and Remaining ChallengesabstractBig Data analytics and new problems in social networks, computational biology, and web connectivity led to a renewed research interest in graph processing. Due to "irregularity" of graph computations, efficient parallel graph processing faces a set of software and hardware challenges debated in literature. In this paper, by utilizing hardware performance counters, we characterize system bottlenecks, resource usage, and the efficiency of popular graph applications on the modern commodity hardware. We analyze selected graph applications (implemented in the Galois framework) on a variety of graph datasets: both scale-free graphs and meshes. Our profiling shows that with an increased number of cores the analyzed graph applications achieve a good speedup, which is highly correlated with utilized memory bandwidth. Contrary to traditional past stereotypes, we find that graph applications significantly benefit from hardware prefetchers. Moreover, the use of transparent huge pages (THP) exhibits a "double win" impact: 1) THP significantly decrease the TLB misses and page walk durations, and 2) THP boost the hardware prefetchers' performance. These insights shed light to understand the performance of emerging systems with large memories. Our profiling framework reports hardware counter values over time. It reveals the danger of using averages for a bottleneck and resource usage analysis: many applications have a time-varying behavior and stretch the usage of system resources to their peak. We discuss the new insights and remaining challenges for guiding the design of future hardware and software components for efficient graph processing. Assaf Eisenman, Ludmila Cherkasova, Guilherme Magalhaes, Qiong Cai, Sachin Katti |
MASCOTS | 2 |
| 2016 | Parallel Graph Processing: Prejudice and State of the ArtabstractLarge graph processing has attracted much renewed attention due to its increased importance for a social network analysis. The efficient parallel graph processing faces a set of software and hardware issues, discussed in literature. The main cause of these challenges is the "irregularity" of graph computations and related difficulties in efficient parallelization of graph processing. Unbalanced computations, caused by uneven data partitioning, can affect application scalability. Moreover, the issue of poor data locality is another major concern, that makes the graph processing applications memory-bound. In this paper, we aim to profile how large, parallel graph applications (based on Galois framework) utilize modern systems, in particular, memory subsystem. We found that modern graph processing frameworks executed on the latest Intel multi-core systems (a single node server) exhibit a good data locality and achieve a good speedup with an increased number of cores, contrary to traditional past stereotypes. The application processing speedup is highly correlated with utilized memory bandwidth. At the same time, our measurements show that the memory bandwidth is not a bottleneck, and the analyzed graph applications are memory-latency bound. These new insights can help us in matching the resource demands of the graph processing applications to future system design parameters. Assaf Eisenman, Ludmila Cherkasova, Guilherme Magalhaes, Qiong Cai, Paolo Faraboschi, Sachin Katti |
ICPE | 2 |
| 2016 | Towards Performance and Scalability Analysis of Distributed Memory Programs on Large-Scale ClustersabstractMany HPC and modern Big Data processing applications belong to a class of so-called scale-out applications, where the application dataset is partitioned and processed by a cluster of machines. Understanding and assessing the scalability of the designed application is one of the primary goals during the application implementation. Typically, in the design and implementation phase, the programmer is bound to a limited size cluster for debugging and performing profiling experiments. The challenge is to assess the scalability of the designed program for its execution on a larger cluster. While in an increased size cluster, each node needs to process a smaller fraction of the original dataset, the communication volume and communication time might be significantly increased, which could become detrimental and provide diminishing performance benefits. The distributed memory applications exhibit complex behavior: they tend to interleave computations and communications, use bursty transfers, and utilize global synchronization primitives. Therefore, one of the main challenges is the analysis of bandwidth demands due to increased communication volume as a function of a cluster size. In this paper, we introduce a novel approach to assess the scalability and performance of a distributed memory program for execution on a large-scale cluster. Our solution involves 1) a limited set of traditional experiments performed in a medium size cluster and 2) an additional set of similar experiments performed with an "interconnect bandwidth throttling" tool, which enables the assessment of the communication demands with respect to available bandwidth. This approach enables a prediction of a cluster size, where a communication cost becomes a dominant component, at which point the performance benefits of the increased cluster lead to a diminishing return. We demonstrate the proposed approach using a popular Graph500 benchmark. Sourav Medya, Ludmila Cherkasova, Guilherme Magalhaes, Kivanc M. Ozonat, Chaitra Padmanabha, Jiban Sarma, Imran Sheikh 0002 |
ICPE | 2 |
| 2016 | Interconnect Emulator for Aiding Performance Analysis of Distributed Memory ApplicationsabstractMany modern large graph and Big Data processing applications operate on datasets that do not fit into DRAM of a single machine. This leads to a design of scale-out applications, where the application dataset is partitioned and processed by a cluster of machines. Distributed memory applications exhibit complex behavior: they tend to interleave computations and communications, use bursty transfers, and utilize global synchronization primitives. This makes it difficult to analyze the impact of communication layer on the application performance and answer the questions: how interconnect latency or bandwidth characteristics may change the application performance will the application performance scale when processed by a larger system? In this work, we introduce a novel emulation framework, called InterSense, which is implemented on top of existing high-speed interconnect, such as InfiniBand, and which provides two performance knobs for changing the (today's) interconnect bandwidth and latency. This approach offers an easy-to-use framework for a sensitivity analysis of complex distributed applications to communication layer performance instead of creating customized and time-consuming application models to answer the same questions. We evaluate the emulator accuracy with popular OSU MPI benchmark suite and two clusters with different generation InfiniBand interconnects (DDR and FDR): InterSense emulates the specified andwidth and latency values with less than 2% error between the expected and measured values. To demonstrate the InterSense's ease of use, we present a case study, where we apply InterSense for sensitivity analysis of four applications and benchmarks for getting non-trivial insights. Ludmila Cherkasova, Jun Li 0008, Haris Volos 0001 |
ICPE | 2 |
| 2015 | InterSense: Interconnect Performance Emulator for Future Scale-out Distributed Memory ApplicationsabstractA common approach for improving application performance is to process its working set from memory. For datasets that do not fit into DRAM of a single machine this leads to a design of scale-out applications, where the application dataset is partitioned and processed by a cluster of machines. Performance of distributed memory applications, implemented using MPI (Message Passing Interface), inherently depends on performance of communication layer, which is largely defined by performance characteristics of underlying interconnect. During last couple years, many Big Data applications, e.g., Hadoop, Spark, Memcached, were re-written to take advantage of Remote Direct Memory Access (RDMA) technology and RDMA-capable interconnects which provide fast and high-bandwidth communications. The application analysis of potential performance improvements due to faster and higher bandwidth interconnects is a challenging task. Does the existing application implementation take a full advantage of the underlying interconnect or not? Will the application performance get worse if the interconnect has X% increased latency or Y% lower bandwidth? In this work, we introduce a novel emulation framework, called InterSense, which is implemented on top of existing high-speed interconnect, such as InfiniBand, and which provides two performance knobs for changing the interconnect bandwidth and latency. This approach offers an easy-to-use framework for a sensitivity analysis of complex distributed applications to communication layer performance instead of creating customized and time-consuming application models to answer the same questions. We evaluate the emulator accuracy with popular OSU MPI benchmarks: InterSense emulates the specified bandwidth and latency values with less than 2% error between the expected and measured values. We apply InterSense for sensitivity analysis of two new benchmarks, such as GUPS and Graph 500 to demonstrate the emulator's ease of use in getting non-trivial insights. Ludmila Cherkasova, Jun Li 0008, Haris Volos 0001 |
MASCOTS | 2 |
| 2015 | Quartz: A Lightweight Performance Emulator for Persistent Memory SoftwareabstractNext-generation non-volatile memory (NVM) technologies, such as phase-change memory and memristors, can enable computer systems infrastructure to continue keeping up with the voracious appetite of data-centric applications for large, cheap, and fast storage. Persistent memory has emerged as a promising approach to accessing emerging byte-addressable non-volatile memory through processor load/store instructions. Due to lack of commercially available NVM, system software researchers have mainly relied on emulation to model persistent memory performance. However, existing emulation approaches are either too simplistic, or too slow to emulate large-scale workloads, or require special hardware. To fill this gap and encourage wider adoption of persistent memory, we developed a performance emulator for persistent memory, called Quartz. Quartz enables an efficient emulation of a wide range of NVM latencies and bandwidth characteristics for performance evaluation of emerging byte-addressable NVMs and their impact on applications performance (without modifying or instrumenting their source code) by leveraging features available in commodity hardware. Our emulator is implemented on three latest Intel Xeon-based processor architectures: Sandy Bridge, Ivy Bridge, and Haswell. To assist researchers and engineers in evaluating design decisions with emerging NVMs, we extend Quartz for emulating the application execution on future systems with two types of memory: fast, regular volatile DRAM and slower persistent memory. We evaluate the effectiveness of our approach by using a set of specially designed memory-intensive benchmarks and real applications. The accuracy of the proposed approach is validated by running these programs both on our emulation platform and a multisocket (NUMA) machine that can support a range of memory latencies. We show that Quartz can emulate a range of performance characteristics with low overhead and good accuracy (with emulation errors 0.2% - 9%). Haris Volos 0001, Guilherme Magalhaes, Ludmila Cherkasova, Jun Li 0008 |
Middleware | 3 |
| 2015 | A Framework for Emulating Non-Volatile Memory Systemswith Different Performance CharacteristicsabstractExponential increase of online data and a corresponding growth of data-centric applications (Big Data analytics) forces system architects to revisit assumptions and requirements of the future system design. New non-volatile memory (NVM) technologies, such as Phase-Change Memory (PCM) and HP Memristor offer significantly improved latency and power efficiency compared to flash and hard drives. Many future systems are expected to have both DRAM and NVM. This can radically change system and software design, and enable new style of Big Data processing applications. However, the commercial unavailability of new NVMs technologies and uncertainty of their performance characteristics make it difficult to assess new system software stacks and to study their performance impact on future workloads. To bridge this gap and encourage an early design phase, we are building a DRAM-based performance emulation platform, called NVMpro, that leverages features available in commodity hardware, to emulate different latency and bandwidth characteristics of future NVM technologies. NVMpro enables an efficient and accurate emulation of a wide range of NVM latencies and bandwidth characteristics for performance evaluation of emerging byte-addressable NVMs and their impact on applications performance without modifying or instrumenting their source code. Dipanjan Sengupta, Haris Volos 0001, Ludmila Cherkasova, Jun Li 0008, Guilherme Magalhaes, Karsten Schwan |
ICPE | 4 |
| 2014 | Optimizing Power and Performance Trade-offs of MapReduce Job Processing with Heterogeneous Multi-core ProcessorsabstractModern processors are often constrained by a given power budget that forces designers to consider different trade-offs, e.g., to choose between either many slow, power-efficient cores, or fewer faster, power-hungry cores, or to select a combination of them. In this work, we design and evaluate a new Hadoop scheduler, called DyScale, that exploits capabilities offered by heterogeneous cores within a single multi-core processor for achieving a variety of performance objectives. A typical MapReduce workload contains jobs with different performance goals: large, batch jobs that are throughput oriented, and smaller interactive jobs that are response-time sensitive. Heterogeneous multi-core processors enable creating virtual resource pools based on the different core types for multi-class priority scheduling. These virtual Hadoop clusters, based on "slow" cores versus "fast" cores can effectively support different performance objectives that cannot be achieved in a Hadoop cluster with homogeneous processors. Using detailed measurements and extensive simulation study we argue in favor of heterogeneous multi-core processors as they provide performance means for "faster" processing of the small, interactive MapReduce jobs (up to 40% faster), while at the same time offer an improved throughput (up to 40% higher) for large, batch job processing. Feng Yan 0001, Ludmila Cherkasova, Zhuoyao Zhang, Evgenia Smirni |
IEEE CLOUD | 2 |
| 2014 | Heterogeneous cores for MapReduce processing: Opportunity or challenge?abstractTo offer diverse computing capabilities, the emergent modern system on a chip (SoC) might include heterogeneous multi-core processors. The current SoC design is often constrained by a given power budget that forces designers to consider different decision trade-offs, e.g., to choose between many slow cores, fewer faster cores, or to select a combination of them. In this work, we design a new Hadoop scheduler, called DyScale, that exploits capabilities offered by heterogeneous cores for achieving a variety of performance objectives. Our preliminary performance evaluation results confirm potential benefits of heterogeneous multi-core processors for “faster” processing of the small, interactive MapReduce jobs, while at the same time offering an improved throughput and performance for large, batch job processing. Feng Yan 0001, Ludmila Cherkasova, Zhuoyao Zhang, Evgenia Smirni |
NOMS | 2 |
| 2014 | Optimizing cost and performance trade-offs for MapReduce job processing in the cloudabstractCloud computing offers a new, attractive option to customers for provisioning a suitable size Hadoop cluster, consuming resources as a service, executing the MapReduce workload, and paying for the time these resources were used. One of the open questions in such environments is the choice and the amount of resources that a user should lease from the service provider. In this work1, we offer a framework for evaluating and selecting the right underlying platform (e.g., small, medium, or large EC2 instances) and achieving the desirable Service Level Objectives (SLOs). A user can define a set of different SLOs: i) achieving a given completion time for a set of MapReduce jobs while minimizing the cost (budget), or ii) for a given budget select the type and the number of instances that optimize the MapReduce workload performance (i.e., the completion time of the jobs). We demonstrate that the application performance of a customer workload may vary significantly on different platforms. This makes a selection of the best cost/performance platform for a given workload being a challenging problem. Our evaluation study and experiments with Amazon EC2 platform reveal that for different workload mixes the optimized platform choice may result in 37-70% cost savings for achieving the same performance objectives when using different (but seemingly equivalent) choices. The results of our simulation study are validated through experiments with Hadoop clusters deployed on different Amazon EC2 instances. Zhuoyao Zhang, Ludmila Cherkasova, Boon Thau Loo |
NOMS | 2 |
| 2014 | Parameterizable benchmarking framework for designing a MapReduce performance modelabstractSUMMARY In MapReduce environments, many applications have to achieve different performance goals for producing time relevant results. One of typical user questions is how to estimate the completion time of a MapReduce program as a function of varying input dataset sizes and given cluster resources. In this work, we offer a novel performance evaluation framework for answering this question. We analyze the MapReduce processing pipeline and utilize the fact that the execution of map (reduce) tasks consists of specific, well‐defined data processing phases. Only map and reduce functions are custom, and their executions are user‐defined for different MapReduce jobs. The executions of the remaining phases aregeneric(i.e., defined by the MapReduce framework code) and depend on the amount of data processed by the phase and the performance of the underlying Hadoop cluster. First, we designa set of parameterizable microbenchmarksto profile the execution of generic phases and to derivea platform performance modelof a given Hadoop cluster. Then, using the job past executions, we summarize job's properties and performance of its custom map/reduce functions in a compact job profile. Finally, by combining the knowledge of the job profile and the derived platform performance model, we introducea MapReduce performance modelthat estimates the program completion time for processing a new dataset. The proposed benchmarking approach derives an accurate performance model of Hadoop's generic execution phases (once), and then, this model isreusedfor predicting the performance of different applications. The evaluation study justifies our approach and the proposed framework: We use a diverse suite of 12 MapReduce applications to validate the proposed model. The predicted completion times for most experiments are within 10% of the measured ones (with a worst case resulting in 17% of error) on our 66‐node Hadoop cluster. Copyright © 2014 John Wiley & Sons, Ltd Zhuoyao Zhang, Ludmila Cherkasova, Boon Thau Loo |
Concurr. Comput. Pract. Exp. | 2 |
| 2014 | Special Issue on "Quantitative Evaluation of SysTems" (QEST 2012)
Giuliano Casale, Ludmila Cherkasova, Holger Hermanns |
Perform. Evaluation | 2 |
| 2014 | Profiling and evaluating hardware choices for MapReduce environments: An application-aware approach
Ludmila Cherkasova, Roy H. Campbell |
Perform. Evaluation | 2 |
| 2013 | Performance Modeling of MapReduce Jobs in Heterogeneous Cloud EnvironmentsabstractMany companies start using Hadoop for advanced data analytics over large datasets. While a traditional Hadoop cluster deployment assumes a homogeneous cluster, many enterprise clusters are grown incrementally over time, and might have a variety of different servers in the cluster. The nodes' heterogeneity represents an additional challenge for efficient cluster and job management. Due to resource heterogeneity, it is often unclear which resources introduce inefficiency and bottlenecks, and how such a Hadoop cluster should be configured and optimized. In this work1, we explore the efficiency and performance accuracy of the bounds-based performance model for predicting the MapReduce job completion times in heterogeneous Hadoop clusters. We validate the accuracy of the proposed performance model using a diverse set of 13 realistic applications and two different heterogeneous clusters. Since one of the Hadoop clusters is formed by different capacity VM instances in Amazon EC2 environment, we additionally explore and discuss factors that impact the MapReduce job performance in the Cloud. Zhuoyao Zhang, Ludmila Cherkasova, Boon Thau Loo |
IEEE CLOUD | 2 |
| 2013 | Workload analysis and demand prediction for the HP ePrint Service
Vipul Garg, Ludmila Cherkasova, Swaminathan Packirisami, Jerome A. Rolia |
IM | 2 |
| 2013 | ACE: Automated capacity evaluation for HP ePrint
Vipul Garg, Ludmila Cherkasova, Swaminathan Packirisami, Jerome A. Rolia |
IM | 2 |
| 2013 | Getting more for less in optimized MapReduce workflows
Zhuoyao Zhang, Ludmila Cherkasova, Boon Thau Loo |
IM | 2 |
| 2013 | Benchmarking approach for designing a mapreduce performance modelabstractIn MapReduce environments, many of the programs are reused for processing a regularly incoming new data. A typical user question is how to estimate the completion time of these programs as a function of a new dataset and the cluster resources. In this work1 , we offer a novel performance evaluation framework for answering this question. We observe that the execution of each map (reduce) tasks consists of specific, well-defined data processing phases. Only map and reduce functions are custom and their executions are user-defined for different MapReduce jobs. The executions of the remaining phases are generic and depend on the amount of data processed by the phase and the performance of underlying Hadoop cluster. First, we design a set of parameterizable microbenchmarks to measure generic phases and to derive a platform performance model of a given Hadoop cluster. Then using the job past executions, we summarize job's properties and performance of its custom map/reduce functions in a compact job profile. Finally, by combining the knowledge of the job profile and the derived platform performance model, we offer a MapReduce performance model that estimates the program completion time for processing a new dataset. The evaluation study justifies our approach and the proposed framework: we are able to accurately predict performance of the diverse set of twelve MapReduce applications. The predicted completion times for most experiments are within 10% of the measured ones (with a worst case resulting in 17% of error) on our 66-node Hadoop cluster. Zhuoyao Zhang, Ludmila Cherkasova, Boon Thau Loo |
ICPE | 2 |
| 2013 | Performance Modeling and Optimization of Deadline-Driven Pig ProgramsabstractMany applications associated with live business intelligence are written as complex data analysis programs defined by directed acyclic graphs of MapReduce jobs, for example, using Pig, Hive, or Scope frameworks. An increasing number of these applications have additional requirements for completion time guarantees. In this article, we consider the popular Pig framework that provides a high-level SQL-like abstraction on top of MapReduce engine for processing large data sets. There is a lack of performance models and analysis tools for automated performance management of such MapReduce jobs. We offer a performance modeling environment for Pig programs that automatically profiles jobs from the past runs and aims to solve the following inter-related problems: (i) estimating the completion time of a Pig program as a function of allocated resources; (ii) estimating the amount of resources (a number of map and reduce slots) required for completing a Pig program with a given (soft) deadline. First, we design a basic performance model that accurately predicts completion time and required resource allocation for a Pig program that is defined as a sequence of MapReduce jobs: predicted completion times are within 10% of the measured ones. Second, we optimize a Pig program execution by enforcing the optimal schedule of its concurrent jobs. For DAGs with concurrent jobs, this optimization helps reducing the program completion time: 10%--27% in our experiments. Moreover, it eliminates possible nondeterminism of concurrent jobs’ execution in the Pig program, and therefore, enables a more accurate performance model for Pig programs. Third, based on these optimizations, we propose a refined performance model for Pig programs with concurrent jobs. The proposed approach leads to significant resource savings (20%--60% in our experiments) compared with the original, unoptimized solution. We validate our solution using a 66-node Hadoop cluster and a diverse set of workloads: PigMix benchmark, TPC-H queries, and customized queries mining a collection of HP Labs’ web proxy logs. Zhuoyao Zhang, Ludmila Cherkasova, Boon Thau Loo |
ACM Trans. Auton. Adapt. Syst. | 2 |
| 2013 | Orchestrating an Ensemble of MapReduce Jobs for Minimizing Their MakespanabstractCloud computing offers an attractive option for businesses to rent a suitable size MapReduce cluster, consume resources as a service, and pay only for resources that were consumed. A key challenge in such environments is to increase the utilization of MapReduce clusters to minimize their cost. One way of achieving this goal is to optimize the execution of Mapreduce jobs on the cluster. For a set of production jobs that are executed periodically on new data, we can perform an offline analysis for evaluating performance benefits of different optimization techniques. In this work, we consider a subset of production workloads that consists of MapReduce jobs with no dependencies. We observe that the order in which these jobs are executed can have a significant impact on their overall completion time and the cluster resource utilization. Our goal is to automate the design of a job schedule that minimizes the completion time (makespan) of such a set of MapReduce jobs. We introduce a simple abstraction where each MapReduce job is represented as a pair of map and reduce stage durations. This representation enables us to apply the classic Johnson algorithm that was designed for building an optimal two-stage job schedule. We evaluate the performance benefits of the constructed schedule through an extensive set of simulations over a variety of realistic workloads. The results are workload and cluster-size dependent, but it is typical to achieve up to 10-25 percent of makespan improvements by simply processing the jobs in the right order. However, in some cases, the simplified abstraction assumed by Johnson's algorithm may lead to a suboptimal job schedule. We design a novel heuristic, called BalancedPools, that significantly improves Johnson's schedule results (up to 15-38 percent), exactly in the situations when it produces suboptimal makespan. Overall, we observe up to 50 percent in the makespan improvements with the new BalancedPools algorithm. The results of our simulation study are validated through experiments on a 66-node Hadoop cluster. Ludmila Cherkasova, Roy H. Campbell |
IEEE Trans. Dependable Secur. Comput. | 2 |
| 2012 | Selling T-shirts and Time Shares in the CloudabstractCloud computing has emerged as a new and alternative approach for providing computing services. Customers acquire and release resources by requesting and returning virtual machines to the cloud. Different service models and pricing schemes are offered by cloud service providers. This can make it difficult for customers to compare cloud services and select an appropriate solution. Cloud Infrastructure-as-a-Service vendors offer a t-shirt approach for Virtual Machines (VMs) on demand. Customers can select from a set of fixed size VMs and vary the number of VMs as their demands change. Private clouds often offer another alternative, called time-sharing, where the capacity of each VM is permitted to change dynamically. With this approach each virtual machine is allocated a dynamic amount of CPU and memory resources over time to better utilize available resources. We present a tool that can help customers make informed decisions about which approach works most efficiently for their workloads in aggregate and for each workload separately. A case study using data from an enterprise customer with 312 workloads demonstrates the use of the tool. It shows that for the given set of workloads the t-shirt model requires almost twice the number of physical servers as the time share model. The costs for such infrastructure must ultimately be passed on to the customer in terms of monetary costs or performance risks. We conclude that private and public clouds should consider offering both resource sharing models to meet the needs of customers. Daniel Gmach, Jerome A. Rolia, Ludmila Cherkasova |
CCGRID | 3 |
| 2012 | Optimizing Completion Time and Resource Provisioning of Pig ProgramsabstractAs cloud computing continues to mature, IT managers have started concentrating on the support of additional performance requirements: quality of service and tailored resource allocation for achieving service performance goals. In this paper, we consider the popular Pig framework that provides a high-level SQL-like abstraction on top of MapReduce engine for processing large data sets. Programs written in such frameworks are compiled into directed acyclic graphs (DAGs) of MapReduce jobs. Often, data processing applications have to produce results by a certain time deadline. We design a performance modeling framework for Pig programs that solves two inter-related problems: (i) estimating the completion time of a Pig program as a function of allocated resources, (ii) estimating the amount of resources (a number of map and reduce slots) required for completing a Pig program with a given (soft) deadline. To achieve these goals, we first, optimize a Pig program execution by enforcing the optimal schedule of its concurrent jobs. This optimization reduces a program completion time (10%-27% in our experiments), and moreover, it eliminates possible non-determinism in the DAGs execution. Based on our optimization, we propose an accurate performance model for Pig programs. This approach leads to significant resource savings (20%-60% in our experiments) compared with the original, unoptimized solution. We validate our approach in a 66-node Hadoop cluster using two workload sets: TPC-H queries and a set of customized queries mining a collection of HP Labs' web proxy logs. Zhuoyao Zhang, Ludmila Cherkasova, Boon Thau Loo |
CCGRID | 2 |
| 2012 | Two Sides of a Coin: Optimizing the Schedule of MapReduce Jobs to Minimize Their Makespan and Improve Cluster PerformanceabstractLarge-scale MapReduce clusters that routinely process petabytes of unstructured and semi-structured data represent a new entity in the changing landscape of clouds. A key challenge is to increase the utilization of these MapReduce clusters. In this work, we consider a subset of the production workload that consists of MapReduce jobs with no dependencies. We observe that the order in which these jobs are executed can have a significant impact on their overall completion time and the cluster resource utilization. Our goal is to automate the design of a job schedule that minimizes the completion time (makespan) of such a set of MapReduce jobs. We offer a novel abstraction framework and a heuristic, called BalancedPools, that efficiently utilizes performance properties of MapReduce jobs in a given workload for constructing an optimized job schedule. Simulations performed over a realistic workload demonstrate that 15%-38% makespan improvements are achievable by simply processing the jobs in the right order. Ludmila Cherkasova, Roy H. Campbell |
MASCOTS | 2 |
| 2012 | Comparing efficiency and costs of cloud computing modelsabstractPublic and private clouds are being adopted as a cost-effective approach for sharing IT resources. Customers acquire and release resources by requesting and returning virtual machines to the cloud. Different service models are proposed for virtual machine resource management. Some public cloud providers follow a t-shirt model for VM resource sizing. A second approach for resource management is based on a time share model. This paper compares the two approaches from the perspective of resource usage for both the service provider and workload owner. Using data from 312 customer applications, we show that the t-shirt model requires 40% more infrastructure than when a finer degree of resource sharing based on time varying resource shares is permitted. Daniel Gmach, Jerome A. Rolia, Ludmila Cherkasova |
NOMS | 3 |
| 2012 | Deadline-based workload management for MapReduce environments: Pieces of the performance puzzleabstractHadoop and the associated MapReduce paradigm, has become the de facto platform for cost-effective analytics over “Big Data”. There is an increasing number of MapReduce applications associated with live business intelligence that require completion time guarantees. In this work, we introduce and analyze a set of complementary mechanisms that enhance workload management decisions for processing MapReduce jobs with deadlines. The three mechanisms we consider are the following: 1) a policy for job ordering in the processing queue; 2) a mechanism for allocating a tailored number of map and reduce slots to each job with a completion time requirement; 3) a mechanism for allocating and deallocating (if necessary) spare resources in the system among the active jobs. We analyze the functionality and performance benefits of each mechanism via an extensive set of simulations over diverse workload sets. The proposed mechanisms form the integral pieces in the performance puzzle of automated workload management in MapReduce environments. Ludmila Cherkasova, Vijay S. Kumar, Roy H. Campbell |
NOMS | 2 |
| 2012 | Dealing with Burstiness in Multi-Tier Applications: Models and Their ParameterizationabstractWorkloads and resource usage patterns in enterprise applications often show burstiness resulting in large degradation of the perceived user performance. In this paper, we propose a methodology for detecting burstiness symptoms in multi-tier applications but, rather than identifying the root cause of burstiness, we incorporate this information into models for performance prediction. The modeling methodology is based on the index of dispersion of the service process at a server, which is inferred by observing the number of completions within the concatenated busy times of that server. The index of dispersion is used to derive a Markov-modulated process that captures burstiness and variability of the service process at each resource well and that allows us to define queueing network models for performance prediction. Experimental results and performance model predictions are in excellent agreement and argue for the effectiveness of the proposed methodology under both bursty and nonbursty workloads. Furthermore, we show that the methodology extends to modeling flash crowds that create burstiness in the stream of requests incoming to the application. Giuliano Casale, Ningfang Mi, Ludmila Cherkasova, Evgenia Smirni |
IEEE Trans. Software Eng. | 3 |
| 2011 | Play It Again, SimMR!abstractA typical MapReduce cluster is shared among different users and multiple applications. A challenging problem in such shared environments is the ability to efficiently control resource allocations among the running and submitted jobs for achieving users' performance goals. To ease the task of evaluating and comparing different provisioning and scheduling approaches in MapReduce environments, we designed and implemented a simulation environment Sim MR which is comprised of three inter-related components: i) Trace Generator that creates a replayable MapReduce workload, ii) Simulator Engine that accurately emulates the job master functionality in Hadoop, and iii) a pluggable scheduling policy that dictates the scheduler decisions on job ordering and the amount of resources allocated to different jobs over time. We validate the accuracy of Sim MR environment by, first, executing a set of realistic MapReduce applications in a 66-node Hadoop cluster and then by replaying the collected job execution traces in SimMR. Our simulator accurately reproduces the original job processing: the completion times of the simulated jobs are within 5% of the original ones. SimMR can process over one million events per second. This allows users to simulate complex workloads in a few seconds instead of multi-hour executions in the real test bed. Finally, by using SimMR we analyze and compare performance of two novel deadline-driven schedulers over a diverse set of real and synthetic workloads. Ludmila Cherkasova, Roy H. Campbell |
CLUSTER | 2 |
| 2011 | Resource and virtualization costs up in the cloud: Models and design choicesabstractVirtualization offers the potential for cost-effective service provisioning. For service providers who make significant investments in new virtualized data centers in support of private or public clouds, one of the serious challenges is the problem of recovering costs for new server hardware, software, network, storage, management, etc. Gaining visibility and accurately determining the cost of shared resources used by collocated services is essential for implementing a proper chargeback approach in cloud environments. We introduce and compare three different models for apportioning cost and champion the one that is least sensitive to workload placement decisions and provides the most robust and repeatable cost estimates. A detailed study involving 312 workloads from an HP customer environment demonstrates the result. Finally, we employ the cost model in a case study that evaluates the impact on the cost of exploiting different virtualization platform alternatives for the 312 workloads. For example, some workloads may cost more to host using certain virtualization platforms than on others or on standalone hosts. We demonstrate different decision points with potential cost savings of nearly 20% by “right-virtualizing” the workloads. Daniel Gmach, Jerome A. Rolia, Ludmila Cherkasova |
DSN | 3 |
| 2011 | Run-time performance optimization and job management in a data protection solutionabstractThe amount of stored data in enterprise Data Centers quadruples every 18 months. This trend presents a serious challenge for backup management and sets new requirements for performance efficiency of traditional backup and archival tools. In this work, we discuss potential performance shortcomings of the existing backup solutions. During a backup session a predefined set of objects (client filesystems) should be backed up. Traditionally, no information on the expected duration and throughput requirements of different backup jobs is provided. This may lead to an inefficient job schedule and the increased backup session time. We analyze historic data on backup processing from eight backup servers in HP Labs, and introduce two additional metrics associated with each backup job, called job duration and job throughput. Our goal is to use this additional information for automated design of a backup schedule that minimizes the overall completion time for a given set of backup jobs. This problem can be formulated as a resource constrained scheduling problem which is known to be NP-complete. Instead, we propose an efficient heuristics for building an optimized job schedule, called FlexLBF. The new job schedule provides a significant reduction in the backup time (up to 50%) and reduced resource usage (up to 2-3 times). Moreover, we design a simulation-based tool that aims to automate parameter tuning for avoiding manual configuration by system administrators while helping them to achieve nearly optimal performance. Ludmila Cherkasova, Roger Lau, Harald Burose, Subramaniam Venkata Kalambur, Bernhard Kappler, Kuttiraja Veeranan |
Integrated Network Management | 1 |
| 2011 | Chargeback model for resource pools in the cloudabstractThis paper presents three methods for apportioning server costs among workloads in shared resource environments such as computing clouds. We consider a fine sharing of resources, the impact of time varying resource usage, large ratios for peak to mean workload demands, and the influence of random choices for the co-placement of workloads on shared servers. These features can affect the quantity of servers needed to support workloads as well as the robustness of the cost values assigned to each workload. We compare the three methods for apportioning costs and recommend the method that assigns costs in the most repeatable manner. A detailed study involving 312 workloads from an HP customer environment demonstrates the result. Daniel Gmach, Jerome A. Rolia, Ludmila Cherkasova |
Integrated Network Management | 3 |
| 2011 | Resource Provisioning Framework for MapReduce Jobs with Performance Goals
Ludmila Cherkasova, Roy H. Campbell |
Middleware | 2 |
| 2011 | Performance modeling in mapreduce environments: challenges and opportunitiesabstractUnstructured data is the largest and fastest growing portion of most enterprise's assets, often representing 70% to 80% of online data. These steep increase in volume of information being produced often exceeds the capabilities of existing commercial databases. MapReduce and its open-source implementation Hadoop represent an economically compelling alternative that offers an efficient distributed computing platform for handling large volumes of data and mining petabytes of unstructured information. It is increasingly being used across the enterprise for advanced data analytics, business intelligence, and enabling new applications associated with data retention, regulatory compliance, e-discovery and litigation issues. Ludmila Cherkasova |
ICPE | 1 |
| 2010 | DP+IP = design of efficient backup schedulingabstractMany industries experience an explosion in digital content. This explosion of electronic documents, along with new regulations and document retention rules, sets new requirements for performance efficiency of traditional data protection and archival tools. During a backup session a predefined set of objects (client filesystems) should be backed up. Traditionally, no information on the expected duration and throughput requirements of different backup jobs is provided. This may lead to a suboptimal job schedule that results in the increased backup session time. In this work, we characterize each backup job via two metrics, called job duration and job throughput. These metrics are derived from collected historic information about backup jobs during previous backup sessions. Our goal is to automate the design of a backup schedule that minimizes the overall completion time for a given set of backup jobs. This problem can be formulated as a resource constrained scheduling problem where a set of n jobs should be scheduled on m machines with given capacities. We provide an integer programming (IP) formulation of this problem and use available IP-solvers for finding an optimized schedule, called binpacking schedule. Performance benefits of the new bin-packing schedule are evaluated via a broad variety of realistic experiments using backup processing data from six backup servers in HP Labs. The new bin-packing job schedule significantly optimizes the backup session time (20%-60% of backup time reduction). HP Data Protector (DP) is HP's enterprise backup offering and it can directly benefit from the designed technique. Moreover, significantly reduced backup session times guarantee an improved resource/power usage of the overall backup solution. Ludmila Cherkasova, Alex Zhang |
CNSM | 1 |
| 2010 | Lightning: self-adaptive, energy-conserving, multi-zoned, commodity green cloud storage systemabstractThe objective of this research is to present an energy-conserving, self-adaptive Commodity Green Cloud Storage, called Lightning. Lightning's File System dynamically configures the servers in the Cloud Storage into logical Hot and Cold Zones. Lightning uses data-classification driven data placement to realize guaranteed, substantially long, periods (several days) of idleness in a significant subset of servers designated as the Cold Zone, in the commodity datacenter backing the Cloud Storage. These servers are then transitioned to inactive power modes and the resulting energy savings substantially reduce the operating costs of the datacenter. Furthermore, the energy savings allow Lightning to improve the data access performance by incorporation of high-performance, though high-cost Solid State Drives (SSD) without exceeding the total cost of ownership (TCO) of the datacenter. Analytical cost model analysis of Lightning suggests savings in the upwards of $24 million in the TCO of a 20,000 server datacenter. The simulation results show that Lightning can achieve 46% energy costs reduction even when the datacenter is at 80% capacity utilization. Rini T. Kaushik, Ludmila Cherkasova, Roy H. Campbell, Klara Nahrstedt |
HPDC | 2 |
| 2010 | AWAIT: Efficient Overload Management for Busy Multi-tier Web Services under Bursty Workloads
Ludmila Cherkasova, Vittoria de Nitto Persone, Ningfang Mi, Evgenia Smirni |
ICWE | 2 |
| 2010 | Efficient resource allocation and power saving in multi-tiered systemsabstractIn this paper, we present Fastrack, a parameter-free algorithm for dynamic resource provisioning that uses simple statistics to promptly distill information about changes in workload burstiness. This information, coupled with the application's end-to-end response times and system bottleneck characteristics, guide resource allocation that shows to be very effective under a broad variety of burstiness profiles and bottleneck scenarios. Andrew Caniff, Ningfang Mi, Ludmila Cherkasova, Evgenia Smirni |
WWW | 4 |
| 2009 | Applying syntactic similarity algorithms for enterprise information managementabstractFor implementing content management solutions and enabling new applications associated with data retention, regulatory compliance, and litigation issues, enterprises need to develop advanced analytics to uncover relationships among the documents, e.g., content similarity, provenance, and clustering. In this paper, we evaluate the performance of four syntactic similarity algorithms. Three algorithms are based on Broder's "shingling" technique while the fourth algorithm employs a more recent approach, "content-based chunking". For our experiments, we use a specially designed corpus of documents that includes a set of "similar" documents with a controlled number of modifications. Our performance study reveals that the similarity metric of all four algorithms is highly sensitive to settings of the algorithms' parameters: sliding window size and fingerprint sampling frequency. We identify a useful range of these parameters for achieving good practical results, and compare the performance of the four algorithms in a controlled environment. We validate our results by applying these algorithms to finding near-duplicates in two large collections of HP technical support documents. Ludmila Cherkasova, Kave Eshghi, Charles B. Morrey III, Joseph A. Tucek, Alistair C. Veitch |
KDD | 1 |
| 2009 | Enhancing and optimizing a data protection solutionabstractAnalyzing and managing large amounts of unstructured information is a high priority task for many companies. For implementing content management solutions, companies need a comprehensive view of their unstructured data. In order to provide a new level of intelligence and control over data resident within the enterprise, one needs to build a chain of tools and automated processes that enable the evaluation, analysis, and visibility into information assets and their dynamics during the information life-cycle. We propose a novel framework to utilize the existing backup infrastructure by integrating additional content analysis routines and extracting already available filesystem metadata over time. This is used to perform data analysis and trending to add performance optimization and self-management capabilities to backup and information management tasks. Backup management faces serious challenges on its own: processing ever increasing amount of data while meeting the timing constraints of backup windows could require adaptive changes in backup scheduling routines. We revisit a traditional backup job scheduling and demonstrate that random job scheduling may lead to inefficient backup processing and an increased backup time. In this work, we use a historic information about the object backup processing time and suggest an additional job scheduling, and automated parameter tuning which may significantly optimize the overall backup time. Under this scheduling, called LBF, the longest backups (the objects with longest backup time) are scheduled first. We evaluate the performance benefits of the introduced scheduling using a realistic workload collected from the seven backup servers at HP Labs. Significant reduction of the backup time (up to 30%) and improved quality of service can be achieved under the proposed job assignment policy. Ludmila Cherkasova, Roger Lau, Harald Burose, Bernhard Kappler |
MASCOTS | 1 |
| 2009 | Resource pool management: Reactive versus proactive or let's be friends
Daniel Gmach, Jerome A. Rolia, Ludmila Cherkasova, Alfons Kemper |
Comput. Networks | 3 |
| 2009 | Automated anomaly detection and performance modeling of enterprise applicationsabstractAutomated tools for understanding application behavior and its changes during the application lifecycle are essential for many performance analysis and debugging tasks. Application performance issues have an immediate impact on customer experience and satisfaction. A sudden slowdown of enterprise-wide application can effect a large population of customers, lead to delayed projects, and ultimately can result in company financial loss. Significantly shortened time between new software releases further exacerbates the problem of thoroughly evaluating the performance of an updated application. Our thesis is that online performance modeling should be a part of routine application monitoring. Early, informative warnings on significant changes in application performance should help service providers to timely identify and prevent performance problems and their negative impact on the service. We propose a novel framework for automated anomaly detection and application change analysis. It is based on integration of two complementary techniques: (i) a regression-based transaction model that reflects a resource consumption model of the application, and (ii) an application performance signature that provides a compact model of runtime behavior of the application. The proposed integrated framework provides a simple and powerful solution for anomaly detection and analysis of essential performance changes in application behavior. An additional benefit of the proposed approach is its simplicity: It is not intrusive and is based on monitoring data that is typically available in enterprise production environments. The introduced solution further enables the automation of capacity planning and resource provisioning tasks of multitier applications in rapidly evolving IT environments. Ludmila Cherkasova, Kivanc M. Ozonat, Ningfang Mi, Julie Symons, Evgenia Smirni |
ACM Trans. Comput. Syst. | 1 |
| 2008 | Anomaly? application change? or workload change? towards automated detection of application performance anomaly and changeabstractAutomated tools for understanding application behavior and its changes during the application life-cycle are essential for many performance analysis and debugging tasks. Application performance issues have an immediate impact on customer experience and satisfaction. A sudden slowdown of enterprise-wide application can effect a large population of customers, lead to delayed projects and ultimately can result in company financial loss. We believe that online performance modeling should be a part of routine application monitoring. Early, informative warnings on significant changes in application performance should help service providers to timely identify and prevent performance problems and their negative impact on the service. We propose a novel framework for automated anomaly detection and application change analysis. It is based on integration of two complementary techniques: i) a regression-based transaction model that reflects a resource consumption model of the application, and ii) an application performance signature that provides a compact model of run-time behavior of the application. The proposed integrated framework provides a simple and powerful solution for anomaly detection and analysis of essential performance changes in application behavior. An additional benefit of the proposed approach is its simplicity: it is not intrusive and is based on monitoring data that is typically available in enterprise production environments. Ludmila Cherkasova, Kivanc M. Ozonat, Ningfang Mi, Julie Symons, Evgenia Smirni |
DSN | 1 |
| 2008 | An integrated approach to resource pool management: Policies, efficiency and quality metricsabstractThe consolidation of multiple servers and their workloads aims to minimize the number of servers needed thereby enabling the efficient use of server and power resources. At the same time, applications participating in consolidation scenarios often have specific quality of service requirements that need to be supported. To evaluate which workloads can be consolidated to which servers we employ a trace-based approach that determines a near optimal workload placement that provides specific qualities of service. However, the chosen workload placement is based on past demands that may not perfectly predict future demands. To further improve efficiency and application quality of service we apply the trace-based technique repeatedly, as a workload placement controller. We integrate the workload placement controller with a reactive controller that observes current behavior to i) migrate workloads off of overloaded servers and ii) free and shut down lightly-loaded servers. To evaluate the effectiveness of the approach, we developed a new host load emulation environment that simulates different management policies in a time effective manner. A case study involving three months of data for 138 SAP applications compares our integrated controller approach with the use of each controller separately. The study considers trade-offs between i) required capacity and power usage, ii) resource access quality of service for CPU and memory resources, and iii) the number of migrations. We consider two typical enterprise environments: blade and server based resource pool infrastructures. The results show that the integrated controller approach outperforms the use of either controller separately for the enterprise application workloads in our study. We show the influence of the blade and server pool infrastructures on the effectiveness of the management policies. Daniel Gmach, Jerome A. Rolia, Ludmila Cherkasova, Guillaume Belrose, Tom Turicchi, Alfons Kemper |
DSN | 3 |
| 2008 | Burstiness in Multi-tier Applications: Symptoms, Causes, and New Models
Ningfang Mi, Giuliano Casale, Ludmila Cherkasova, Evgenia Smirni |
Middleware | 3 |
| 2008 | Profiling and Modeling Resource Usage of Virtualized Applications
Timothy Wood 0001, Ludmila Cherkasova, Kivanc M. Ozonat, Prashant J. Shenoy |
Middleware | 2 |
| 2008 | Analysis of application performance and its change via representative application signaturesabstractApplication servers are a core component of a multitier architecture that has become the industry standard for building scalable client-server applications. A client communicates with a service deployed as a multi-tier application via request-reply transactions. A typical server reply consists of the web page dynamically generated by the application server. The application server may issue multiple database calls while preparing the reply. Understanding the cascading effects of the various tasks that are sprung by a single request-reply transaction is a challenging task. Furthermore, significantly shortened time between new software releases further exacerbates the problem of thoroughly evaluating the performance of an updated application. We address the problem of efficiently diagnosing essential performance changes in application behavior in order to provide timely feedback to application designers and service providers. In this work, we propose a new approach based on an application signature that enables a quick performance comparison of the new application signature against the old one, while the application continues its execution in the production environment. The application signature is built based on new concepts that are introduced here, namely the transaction latency profiles and transaction signatures. These become instrumental for creating an application signature that accurately reflects important performance characteristics. We show that such an application signature is representative and stable under different workload characteristics. We also show that application signatures are robust as they effectively capture changes in transaction times that result from software updates. Application signatures provide a simple and powerful solution that can further be used for efficient capacity planning, anomaly detection, and provisioning of multi-tier applications in rapidly evolving IT environments. Ningfang Mi, Ludmila Cherkasova, Kivanc M. Ozonat, Julie Symons, Evgenia Smirni |
NOMS | 2 |
| 2007 | Capacity Management and Demand Prediction for Next Generation Data CentersabstractAdvances in server, network, and storage virtualization are enabling the creation of resource pools of servers that permit multiple application workloads to share each server in the pool. This paper proposes and evaluates aspects of a capacity management process for automating the efficient use of such pools when hosting large numbers of services. We use a trace based approach to capacity management that relies on i) a definition for required capacity, ii) the characterization of workload demand patterns, iii) the generation of synthetic workloads that predict future demands based on the patterns, and iv) a workload placement recommendation service. A case study with 6 months of data representing the resource usage of 139 workloads in an enterprise data center demonstrates the effectiveness of the proposed capacity management process. Our results show that when consolidating to 8 processor systems, we predicted future per-server required capacity to within one processor 95% of the time. The approach enabled a 35% reduction in processor usage as compared to today's current best practice for workload placement. Daniel Gmach, Jerome A. Rolia, Ludmila Cherkasova, Alfons Kemper |
ICWS | 3 |
| 2007 | A Capacity Planning Framework for Multi-tier Enterprise Services with Real WorkloadsabstractWith complexity of systems increasing and customer requirements for QoS growing, new methods and modeling techniques that explain large-systems' behavior and help predict their future performance are required to effectively tackle the emerging performance issues. To accurately answer capacity planning questions for an existing production system with a real workload mix, we propose a new capacity planning framework that is based on the following three components: i) a Workload Profiler that dynamically builds the workload profile; ii) a Regression-based Solver that is used for deriving the CPU demand of client transactions on a given hardware; and iii) an Analytical model that is based on a network of queues representing the different tiers. To validate our approach, we conduct a detailed case study using the access logs from two heterogeneous production servers that represent customized client accesses to a popular and actively used HP Open View Service Desk application. Qi Zhang 0012, Ludmila Cherkasova, Guy Mathews, Wayne Greene, Evgenia Smirni |
Integrated Network Management | 2 |
| 2007 | R-Capriccio: A Capacity Planning and Anomaly Detection Tool for Enterprise Services with Live Workloads
Qi Zhang 0012, Ludmila Cherkasova, Guy Mathews, Wayne Greene, Evgenia Smirni |
Middleware | 2 |
| 2007 | Modeling and generating realistic streaming media server workloads
Wenting Tang, Yun Fu 0003, Ludmila Cherkasova, Amin Vahdat |
Comput. Networks | 3 |
| 2006 | R-Opus: A Composite Framework for Application Performability and QoS in Shared Resource PoolsabstractWe consider shared resource pool management taking into account per-application quality of service (QoS) requirements and server failures. Application QoS requirements are defined by complementary specifications for acceptable and time-limited degraded performance. Furthermore, a requirement specification is provided for both the normal case and for the case where an application server fails in the pool. Independently, the resource pool operator provides a resource access QoS commitment for two classes of service (CoS). These govern statistical multiplexing within the pool. A QoS translation automatically maps application demands onto the resource pool's CoS to best enable sharing. A workload placement service consolidates applications to a small number of servers while satisfying application QoS requirements. The service reports whether a spare server is needed or how applications affected by a single failure can operate according to failure QoS constraints using remaining servers until the failure can be repaired. A case study demonstrates the approach Ludmila Cherkasova, Jerome A. Rolia |
DSN | 1 |
| 2006 | Enforcing Performance Isolation Across Virtual Machines in Xen
Diwaker Gupta, Ludmila Cherkasova, Rob Gardner, Amin Vahdat |
Middleware | 2 |
| 2006 | Configuring Workload Manager Control Parameters for Resource PoolsabstractResource pools are computing environments that offer virtualized access to shared resources. When used effectively they can align the use of capacity with business needs (flexibility), lower infrastructure costs (via resource sharing), and lower operating costs (via automation). Using resources effectively can rely on a combination of workload placement and workload management technologies. Workload placement decides which workloads will share resources. Workload management governs short term access to resource capacity. It provides performance isolation within resource pools to ensure resource sharing even under high loads. A workload manager can have a direct impact both on an application's overall resource access quality of service and on the number of workloads that can be assigned to a pool. In this paper we take a detailed look at an application workload's demands. We explore tradeoffs in resource access quality of service received by the application and the minimum allocation of resources for the workload. We show that by careful selection of workload scheduling parameters along with a proposed fast allocation policy we can sometimes more than triple the number of workloads that can be assigned to a pool without sacrificing application workload quality of service or the efficiency of the resource pool. Jerome A. Rolia, Ludmila Cherkasova, Clifford McCarthy |
NOMS | 2 |
| 2006 | A unified benchmarking and model-based framework for building QoS-aware streaming media services
Ludmila Cherkasova, Wenting Tang, Amin Vahdat |
Multim. Syst. | 1 |
| 2005 | Optimizing the Reliable Distribution of Large Files within CDNsabstractContent delivery networks (CDNs) provide an efficient support for serving http and streaming media content white minimizing the network impact of content delivery as well as overcoming the server overload problem. For serving the large documents and media files, there is an additional problem of the original content distribution across the CDN edge servers. We propose an algorithm, called ALM-fastreplica, for optimizing replication of large files across the edge servers in CDNs. The original file is partitioned into k subfiles, and each subfile is replicated via a correspondingly constructed multicast tree. Nodes from the different multicast trees use additional cross-nodes connections to exchange their corresponding subfiles such that each node eventually receives an entire file. This new replication method significantly reduces file replication time, up to 5-15 times compared to the traditional unicast (or point-to-point) schema. Since a single node failure in the multicast tree during the file distribution may impact the file delivery to a significant number of nodes, it is important to design an algorithm which is able to deal with node failures. We augment ALM-FastReplica with an efficient reliability mechanism, that can deal with node failures by making local repair decisions within a particular replication group of nodes. Under the proposed algorithm, the load of the failed node is shared among the nodes of the corresponding replication group, making the performance degradation gradual. Ludmila Cherkasova |
ISCC | 1 |
| 2005 | Measuring CPU Overhead for I/O Processing in the Xen Virtual Machine Monitor
Ludmila Cherkasova, Rob Gardner |
USENIX ATC, General Track | 1 |
| 2004 | Sizing the streaming media cluster solution for a given workloadabstractThe goal of the proposed benchmarking and capacity planning framework is to evaluate the amount of resources needed for processing a given workload while meeting the specified performance requirements. There are two essential components in our capacity planning framework: (i) the capacity measurements of different hardware and software solutions using a specially designed set of media benchmarks; and (ii) a media service workload profiler, called MediaProf, which extracts a set of quantitative and qualitative parameters characterizing the service demand. The capacity planning tool matches the requirements of the media service workload profile, SLA and configuration constraints to produce the best available cost/performance solution. In case of a multi-node configuration, the capacity planner performs a cluster sizing evaluation by taking into account the choice of load balancing solution. Ludmila Cherkasova, Wenting Tang |
CCGRID | 1 |
| 2004 | An SLA-Oriented Capacity Planning Tool for Streaming Media ServicesabstractThis paper addresses the problem of mapping the requirements of a known media service workload into the corresponding system resource requirements and accurately sizing a media server cluster to handle the workload. In this paper, we propose a new capacity planning framework for evaluating the resources needed for processing a given streaming media workload with specified performance requirements. The performance requirements are specified in a service level agreement (SLA) containing: i) basic capacity requirements that define the percentage of time the configuration is capable of processing the workload without performance degradation while satisfying bounds on system utilization; and ii) performability requirements that define the acceptable degradation of service performance during the remaining, non-compliant time and in case of node failures. Using a set of specially benchmarked media server configurations, the capacity planning tool matches the overall capacity requirements of the media service workload profile with the specified SLAs to identify the number of nodes necessary to support the required service performance. Ludmila Cherkasova, Wenting Tang, Sharad Singhal |
DSN | 1 |
| 2004 | Analysis of enterprise media server workloads: access patterns, locality, content evolution, and rates of changeabstractUnderstanding the nature of media server workloads is crucial to properly designing and provisioning current and future media services. The main issue we address in this paper is the workload analysis of today's enterprise media servers. This analysis aims to establish a set of properties specific to the enterprise media server workloads and to compare them to well-known related observations about the web server workloads. We partition the media workload properties in two groups: static and temporal. While the static properties provide more traditional and general characteristics of the underlying media fileset and quantitative properties of client accesses to those files (independent of the access time), the temporal properties reflect the dynamics and evolution of accesses to the media content over time. We propose two new metrics characterizing the temporal properties: 1) the new files impact metric characterizing the site evolution due to new content and 2) the life span metric reflecting the rates of change in accesses to the newly introduced files. We illustrate these new metrics with the analysis of two different enterprise media server workloads collected over a significant period of time. Ludmila Cherkasova |
IEEE/ACM Trans. Netw. | 1 |
| 2003 | Building a Performance Model of Streaming Media Applications in Utility Data Center EnvironmentabstractUtility Data Center (UDC) provides a flexible, cost-effective infrastructure to support the hosting of applications for Internet services. In order to enable the design of a "utility-aware" streaming media service which automatically requests the necessary resources from UDC infrastructure, we introduce a set of benchmarks for measuring the basic capacities of streaming media systems. The benchmarks allow one to derive the scaling rules of server capacity for delivering media files which are: i) encoded at different bit rates, ii) streamed from memory vs disk. Using an experimental testbed, we show that these scaling rules are non-trivial. In this paper, we develop a workload-aware, media server performance model which is based on a cost function derived from the set of basic benchmark measurements. We validate this performance model by comparing the predicted and measured media server capacities for a set of synthetic workloads. Ludmila Cherkasova, Loren Staley |
CCGRID | 1 |
| 2003 | Capacity planning tool for streaming media servicesabstractThe goal of the proposed capacity planning tool is to provide the best cost/performance configuration for support of a known media service workload. There are two essential components in our capacity planning tool: i) the capacity measurements of different h/w and s/w solutions using a specially designed set of media benchmarks and ii) a media service workload profiler, called MediaProf, which extracts a set of quantitative and qualitative parameters characterizing the service demand. The capacity planning tool matches the requirements of the media service workload profile, SLAs and configuration constraints to produce the best available cost/performance solution. Ludmila Cherkasova, Wenting Tang |
ACM Multimedia | 1 |
| 2003 | MediSyn: a synthetic streaming media service workload generatorabstractCurrently, Internet hosting centers and content distribution networks leverage statistical multiplexing to meet the performance requirements of a number of competing hosted network services. Developing efficient resource allocation mechanisms for such services requires an understanding of both the short-term and long-term behavior of client access patterns to these competing services. At the same time, streaming media services are becoming increasingly popular, presenting new challenges for designers of shared hosting services. These new challenges result from fundamentally new characteristics of streaming media relative to traditional web objects, principally different client access patterns and significantly larger computational and bandwidth overhead associated with a streaming request. To understand the characteristics of these new workloads we use two long-term traces of streaming media services to develop MediSyn, a publicly available streaming media workload generator. In summary, this paper makes the following contributions: i) we model the long-term behavior of network services capturing the process of file introduction and changing file popularity, ii) we present a novel generalized Zipf-like distribution that captures recently-observed popularity of both web objects and streaming media not captured by existing Zipf-like distributions, and iii) we capture a number of characteristics unique to streaming media services, including file duration, encoding bit rate, session duration and non-stationary popularity of media accesses. Wenting Tang, Yun Fu 0003, Ludmila Cherkasova, Amin Vahdat |
NOSSDAV | 3 |
| 2003 | Measuring and characterizing end-to-end Internet service performanceabstractFundamental to the design of reliable, high-performance network services is an understanding of the performance characteristics of the service as perceived by the client population as a whole. Understanding and measuring such end-to-end service performance is a challenging task. Current techniques include periodic sampling of service characteristics from strategic locations in the network and instrumenting Web pages with code that reports client-perceived latency back to a performance server. Limitations to these approaches include potentially nonrepresentative access patterns in the first case and determining the location of a performance bottleneck in the second.This paper presents EtE monitor, a novel approach to measuring Web site performance. Our system passively collects packet traces from a server site to determine service performance characteristics. We introduce a two-pass heuristic and a statistical filtering mechanism to accurately reconstruct different client page accesses and to measure performance characteristics integrated across all client accesses. Relative to existing approaches, EtE monitor offers the following benefits: i) a latency breakdown between the network and server overhead of retrieving a Web page, ii) longitudinal information for all client accesses, not just the subset probed by a third party, iii) characteristics of accesses that are aborted by clients, iv) an understanding of the performance breakdown of accesses to dynamic, multitiered services, and v) quantification of the benefits of network and browser caches on server performance. Our initial implementation and performance analysis across three different commercial Web sites confirm the utility of our approach. Ludmila Cherkasova, Yun Fu 0003, Wenting Tang, Amin Vahdat |
ACM Trans. Internet Techn. | 1 |
| 2002 | Measuring the capacity of a streaming media server in a Utility Data Center environmentabstractAbstract In order to design a "utility-aware" streaming media service which automatically requests the necessary resources from Utility Data Center infrastructure, several classic performance questions should be answered: how to measure the basic capacity of a streaming media server? what is the set of basic benchmarks exposing the performance limits and main bottlenecks of a media server? In this paper, we propose a set of benchmarks for measuring the basic capacities of streaming media systems for different expected workloads, and demonstrate the results using an experimental testbed. Ludmila Cherkasova, Loren Staley |
ACM Multimedia | 1 |
| 2002 | Characterizing locality, evolution, and life span of accesses in enterprise media server workloadsabstractThe main issue we address in this paper is the workload analysis of today's enterprise media servers. This analysis aims to establish a set of properties specific for enterprise media server workloads and to compare them with well known related observations about web server workloads. We propose two new metrics to characterize the dynamics and evolution of the accesses, and the rate of change in the site access pattern, and illustrate them with the analysis of two different enterprise media server workloads collected over a significant period of time. Another goal of our workload analysis study is to develop a media server log analysis tool, called MediaMetrics, that produces a media server traffic access profile and its system resource usage in a way useful to service providers. Ludmila Cherkasova |
NOSSDAV | 1 |
| 2002 | EtE: Passive End-to-End Internet Service Performance Monitoring
Yun Fu 0003, Amin Vahdat, Ludmila Cherkasova, Wenting Tang |
USENIX ATC, General Track | 3 |
| 2002 | Session-Based Admission Control: A Mechanism for Peak Load Management of Commercial Web SitesabstractWe consider a new, session-based workload for measuring web server performance. We define a session as a sequence of client's individual requests. Using a simulation model, we show that an overloaded web server can experience a severe loss of throughput measured as a number of completed sessions compared against the server throughput measured in requests per second. Moreover, statistical analysis of completed sessions reveals that the overloaded web server discriminates against longer sessions. For e-commerce retail sites, longer sessions are typically the ones that would result in purchases, so they are precisely the ones for which the companies want to guarantee completion. To improve Web QoS for commercial Web servers, we introduce a session-based admission control (SBAC) to prevent a web server from becoming overloaded and to ensure that longer sessions can be completed. We show that a Web server augmented with the admission control mechanism is able to provide a fair guarantee of completion, for any accepted session, independent of a session length. This provides a predictable and controllable platform for web applications and is a critical requirement for any e-business. Additionally, we propose two new adaptive admission control strategies, hybrid and predictive, aiming to optimize the performance of SBAC mechanism. These new adaptive strategies are based on a self-tunable admission control function, which adjusts itself accordingly to variations in traffic loads. Ludmila Cherkasova, Peter Phaal |
IEEE Trans. Computers | 1 |
| 2001 | Dynamics and Evolution of Web Sites: Analysis, Metrics and Design IssuesabstractOur goal is to develop a Web server log analysis tool that produces a Web site profile and its system resource usage in a way useful to service providers. Understanding the nature of traffic to the Web site is crucial in properly designing site support infrastructure, especially for large, busy sites. The main questions we address are the new access patterns of today's WWW, how to characterize dynamics or evolution of Web sites, and how to measure the rate of changes. We propose a set of new metrics to characterize the site dynamics, and we illustrate them with analysis of three different Web sites. Ludmila Cherkasova, Magnus Karlsson 0002 |
ISCC | 1 |
| 2000 | Characterizing temporal locality and its impact on web server performanceabstractThe presence of temporal locality in web traces has long been recognized. However, the close proximity of requests for the same file in a trace can be attributed to two orthogonal reasons: long-term popularity and short-term correlation. The former reflects the fact that requests for a popular document appear "frequently" thus they are likely to be "close" in an absolute sense. The latter reflects the fact that requests for a given document might concentrate around particular points in the trace due to a variety of reasons, such as deadlines or surges in user interests, hence it focuses on "relative" closeness. We introduce a new measure of temporal locality, the scaled stack distance, which is insensitive to popularity and captures instead the impact of short-term correlation, and use it to parameterize a synthetic trace generator. Then, we validate the appropriateness of this quantity by comparing the file and byte miss ratios corresponding to either the original or the synthetic traces. Ludmila Cherkasova, Gianfranco Ciardo |
ICCCN | 1 |
| 2000 | Detecting Timed-Out Client Requests for Avoiding Livelock and Improving Web Server PerformanceabstractA Web server's listen queue capacity is often set to a large value to accommodate bursts of traffic and to accept as many client requests as possible. We show that, under certain overload conditions, this results in a significant loss of server performance due to the processing of so-called "dead requests": timed-out client requests whose associated connection has been closed from the client side. In some pathological cases, these overload conditions can lead to server livelock, where the server is busily processing only dead requests and is not doing any useful work. We propose a method of detecting these dead requests and avoiding unnecessary overhead related to their processing. This provides a predictable and controllable platform for Web applications, thus improving their overall performance. DTOC strategy is implemented as a part of WebQoS product for HP9000 servers [HP-WebQoS]. Richard Carter, Ludmila Cherkasova |
ISCC | 2 |
| 2000 | FLEX: Load Balancing and Management Strategy for Scalable Web Hosting ServiceabstractFLEX is a new scalable "locality aware" solution for achieving both load balancing and efficient memory usage on a cluster of machines hosting several Web sites. FLEX allocates the sites to different machines in the cluster based on their traffic characteristics. This aims to avoid the unnecessary document replication to improve the overall performance of the system. The desirable routing can be done by submitting the corresponding configuration files to the DNS server since each hosted Web site has a unique domain name. FLEX can be easily implemented on top of the current infrastructure used by Web hosting service providers. Using a simulation model and a synthetic trace generator, we compare the Round-Robin based solutions and FLEX over the range of different workloads. For generated traces, FLEX outperforms Round-Robin based solutions 2-5 times. Ludmila Cherkasova |
ISCC | 1 |
| 2000 | Optimizing a "Content-Aware" Load Balancing Strategy for Shared Web Hosting ServiceabstractFLEX is a new scalable "locality-aware" solution for achieving both load balancing and efficient memory usage on a cluster of machines hosting several Web sites. FLEX allocates the sites to different machines in the cluster based on their traffic characteristics. We propose a set of new methods and algorithms (Simple+, Advanced and Advanced+) to improve the allocation of the Web sites to different machines. New methods show additional performance benefits as compared to the original Simple strategy. Experiments show that FLEX outperforms traditional load-balancing solutions by 50% to 100% (in throughput), even for a four-node cluster. The miss ratio is improved 2-3 times. Ludmila Cherkasova, Shankar Ponnekanti |
MASCOTS | 1 |
| 2000 | Predictive Admission Control Strategy for Overloaded Commercial Web ServerabstractUses a session-based workload to measure a Web server's performance. An overloaded Web server can experience a severe loss of throughput when measured as the number of completed sessions. Moreover, the overloaded Web server discriminates against longer sessions. Session-based admission control (SBAC) prevents a Web server from becoming overloaded and ensures that longer sessions can be completed. In this paper, we propose a new "predictive" strategy as a robust and flexible admission control mechanism, aiming to optimize the SBAC performance. Ludmila Cherkasova, Peter Phaal |
MASCOTS | 1 |
| 1996 | Components of Congestion ControlabstractThis paper presents three original and complementary ideas on flow control mechanisms for packet switched interconnects, backpressure flow control, alpha message scheduling and balanced injection.These three components were adopted for the I'ed-.Y fabric.Each of the three addresses performance problems caused by a particular characteristic of realistic network workloads. Ludmila Cherkasova, Al Davis, Robin Hodgson, Vadim E. Kotov, Ian N. Robinson, Tomas Rokicki |
SPAA | 1 |
| 1995 | Modeling A Fibre Channel Switch with Stochastic Petri NetsabstractNo abstract available. Gianfranco Ciardo, Ludmila Cherkasova, Vadim E. Kotov, Tomas Rokicki |
SIGMETRICS | 2 |
| 1995 | Bounded Self-Stabilizing Petri Nets
Ludmila Cherkasova, Rodney R. Howell, Louis E. Rosier |
Acta Informatica | 1 |
| 1992 | Compositional Generation of Home States in Free Choice NetsabstractAbstract Free choice nets are a class of Petri nets which allow the modelling of concurrency and nondeterministic choice, but with the restriction that choices cannot be influenced externally. Home states are initial states which lead to a strongly connected state graph, that is, a home state can be reached from any of its successor states. The main result of this paper characterises the home states of a well-formed free choice net compositionally by recourse to its decomposition into T-components. Eike Best, Ludmila Cherkasova, Jörg Desel |
Formal Aspects Comput. | 2 |
| 1991 | Compositional Generation of Home States in Free Choice Systems
Eike Best, Ludmila Cherkasova, Jörg Desel |
STACS | 2 |
| 1991 | An Algebra of Concurrent Non-Deterministic Processes
Ludmila Cherkasova, Vadim E. Kotov |
Theor. Comput. Sci. | 1 |
| 1989 | Concurrent Nondeterministic Processes: Adequacy of Structure and Behaviour
Ludmila Cherkasova, Vadim E. Kotov |
MFCS | 1 |
| 1988 | On Models and Algebras for Concurrent Processes
Ludmila Cherkasova |
MFCS | 1 |
| 1987 | On Generalized Process Logic
Vadim E. Kotov, Ludmila Cherkasova |
FCT | 2 |
| 1981 | Structured Nets
Ludmila Cherkasova, Vadim E. Kotov |
MFCS | 1 |