Demonstration venue · read-only. Every page can be browsed; the buttons that would change it are switched off. Create an account to run TaxoReview on your own data.

Karsten Schwan

dblp:s/KarstenSchwan · DBLP profile ↗
← Back
205ranked-venue papers
14as first author
0since 2021 · last 2017
—ORCID · none

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

Systems, architecture and hardware · 143 · 6 first-authorSoftware engineering, systems software and programming languages · 31 · 5 first-authorComputer networks · 10Applied, interdisciplinary, general and emerging computing · 10 · 2 first-authorSecurity and privacy · 6 · 1 first-authorGraphics, computer vision, multimedia, augmented reality and games · 4Artificial intelligence and machine learning · 2 · 1 first-authorHuman-computer interaction and ubiquitous computing · 2Databases, data management, data science and information retrieval · 1

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

Computer architecture, parallel and distributed computing, and storage systems
77 papers
Memory systems · 23% High-performance computing · 14% Storage systems · 13%
Software engineering, system software, and programming languages
18 papers
Operating systems · 98% Concurrent programming · 2% Debugging and program repair · 0%
Databases, data mining, and information retrieval
4 papers
Graph data management · 43% Transaction processing and concurrency control · 38% Data stream processing · 14%

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

TopicWeightPapersLastEvidence papers
Memory systems
non-volatile memory
0.742016
pVM: persistent virtual memory for efficient capacity scaling and object storage · EuroSys 2016
NVRAM-aware Logging in Transaction Systems · Proc. VLDB Endow. 2014
Reducing the cost of persistence for nonvolatile heaps in end user devices · HPCA 2014
Operating systems › resource management › memory management
virtual memory
0.632016
pVM: persistent virtual memory for efficient capacity scaling and object storage · EuroSys 2016
SpaceJMP: Programming with Multiple Virtual Address Spaces · ASPLOS 2016
Unified address translation for memory-mapped SSDs with FlashMap · ISCA 2015
Cloud and datacenter computing › resource management
datacenter memory management
0.312017
HeteroOS: OS Design for Heterogeneous Memory Management in Datacenter · ISCA 2017
Memory systems › memory management
heterogeneous memory management
0.312017
HeteroOS: OS Design for Heterogeneous Memory Management in Datacenter · ISCA 2017
Memory systems › non-volatile memory › persistent memory
byte-addressable persistent memory
0.322016
Reducing the cost of persistence for nonvolatile heaps in end user devices · HPCA 2014
pVM: persistent virtual memory for efficient capacity scaling and object storage · EuroSys 2016
Memory systems › non-volatile memory
persistent memory
0.322016
Reducing the cost of persistence for nonvolatile heaps in end user devices · HPCA 2014
SpaceJMP: Programming with Multiple Virtual Address Spaces · ASPLOS 2016
Operating systems › resource management › memory management
address space management
0.212016
SpaceJMP: Programming with Multiple Virtual Address Spaces · ASPLOS 2016
Operating systems › resource management
memory management
0.212016
An Evolutionary Study of Linux Memory Management for Fun and Profit · USENIX ATC 2016
Energy-efficient computing › power management
dynamic voltage and frequency scaling
0.212016
Exploiting variability for energy optimization of parallel programs · EuroSys 2016
Memory systems
hybrid memory
0.212016
Data tiering in heterogeneous memory systems · EuroSys 2016
Storage systems › storage hierarchy
tiered storage
0.212016
Data tiering in heterogeneous memory systems · EuroSys 2016
Parallel and multicore computing › load balancing › dynamic load balancing
work stealing
0.212016
Affinity-aware work-stealing for integrated CPU-GPU processors · PPoPP 2016
Cloud and datacenter computing
virtualization
0.232008
Protectit: trusted distributed services operating on sensitive data · EuroSys 2008
VirtualPower: coordinated power management in virtualized enterprise systems · SOSP 2007
High performance and scalable I/O virtualization via self-virtualized devices · HPDC 2007
Graph data management › graph analytics
large-scale graph analytics
0.212015
GraphReduce: processing large-scale graphs on accelerator-based systems · SC 2015
Memory systems › memory management › virtual memory
address translation
0.212015
Unified address translation for memory-mapped SSDs with FlashMap · ISCA 2015
Distributed systems
distributed graph processing
0.212015
Scaling iterative graph computations with GraphMap · SC 2015
Storage systems
flash and SSD
0.212015
Unified address translation for memory-mapped SSDs with FlashMap · ISCA 2015
Storage systems › flash and SSD › flash memory management
flash translation layer
0.212015
Unified address translation for memory-mapped SSDs with FlashMap · ISCA 2015
GPUs and heterogeneous computing
GPU graph processing
0.212015
GraphReduce: processing large-scale graphs on accelerator-based systems · SC 2015
Parallel and multicore computing › graph processing
iterative graph processing
0.212015
Scaling iterative graph computations with GraphMap · SC 2015
Storage systems › out-of-core computation
out-of-core graph processing
0.212015
GraphReduce: processing large-scale graphs on accelerator-based systems · SC 2015
Memory systems › memory management
virtual memory
0.212015
Unified address translation for memory-mapped SSDs with FlashMap · ISCA 2015
Transaction processing and concurrency control
logging and recovery
0.212014
NVRAM-aware Logging in Transaction Systems · Proc. VLDB Endow. 2014
Cloud and datacenter computing
cluster resource management and scheduling
0.212014
Scheduling Multi-tenant Cloud Workloads on Accelerator-Based Systems · SC 2014
GPUs and heterogeneous computing
GPU resource management
0.212014
Scheduling Multi-tenant Cloud Workloads on Accelerator-Based Systems · SC 2014
GPUs and heterogeneous computing
GPU scheduling
0.212014
Scheduling Multi-tenant Cloud Workloads on Accelerator-Based Systems · SC 2014
Distributed systems
stream processing
0.212014
ELF: Efficient Lightweight Fast Stream Processing at Scale · USENIX ATC 2014
High-performance computing › scientific data analysis
in-situ analysis
0.212013
GoldRush: resource efficient in situ scientific data analytics using fine-grained interference aware execution · SC 2013
Cloud and datacenter computing › resource management
resource harvesting
0.212013
GoldRush: resource efficient in situ scientific data analytics using fine-grained interference aware execution · SC 2013
High-performance computing
parallel i/o
0.222011
Six degrees of scientific data: reading patterns for extreme scale science IO · HPDC 2011
Efficient Wire Formats for High Performance Computing · SC 2000

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

virtualization · 0.7virtual memory management · 0.5OS design · 0.5page tables · 0.4performance modeling · 0.3offline and online analysis · 0.2empirical study · 0.2affinity-aware work-stealing · 0.2DVFS · 0.2two-level graph partitioning · 0.2locality-based optimization · 0.2gather-apply-scatter · 0.2asynchronous GPU streams · 0.2transcoder selection · 0.1dynamic protection rules · 0.1data filters · 0.1soft power states · 0.1server consolidation · 0.1
YearPublicationVenuePosition
2017 Fault-Scalable Virtualized Infrastructure Management
abstract
Large-scale virtualized datacenters require considerable automation in infrastructure management in order to operate efficiently. Automation is impaired, however, by the fact that deployments are prone to multiple types of subtle faults due to hardware failures, software bugs, misconfiguration, crashes, performance degraded hardware, etc. Existing Infrastructure-as-a-Service (IaaS) management stacks incorporate little to no resilience measures to shield end users from such cloud providerlevel failures and poor performance. This paper proposes and evaluates extensions to IaaS stacks that mask faults in a fault-agnostic manner while ensuring that the overheads can be proportional to observed failure rates. We also demonstrate that infrastructure automation services and end-user applications can use service-specific knowledge, together with our new interface, to achieve better outcomes.
Mukil Kesavan, Ada Gavrilovska, Karsten Schwan
ICDCS3
2017 HeteroOS: OS Design for Heterogeneous Memory Management in Datacenter
Sudarsun Kannan, Ada Gavrilovska, Vishal Gupta 0001, Karsten Schwan
ISCA4
2016 Energy Aware Persistence: Reducing Energy Overheads of Memory-based Persistence in NVMs
abstract
Next generation byte addressable nonvolatile memories (NVMs) such as PCM, Memristor, and 3D X-Point are attractive solutions for mobile and other end-user devices, as they offer memory scalability as well as fast persistent storage. However, NVM's limitations of slow writes and high write energy are magnified for applications that require atomic, consistent, isolated and durable (ACID) persistence. For maintaining ACID persistence guarantees, applications not only need to do extra writes to NVM but also need to execute a significant number of additional CPU instructions for performing NVM writes in a transactional manner. Our analysis shows that maintaining persistence with ACID guarantees increases CPU energy up to 7.3x and NVM energy up to 5.1x compared to a baseline with no ACID guarantees. For computing platforms such as mobile devices, where energy consumption is a critical factor, it is important that the energy cost of persistence is reduced.
Sudarsun Kannan, Moinuddin K. Qureshi, Ada Gavrilovska, Karsten Schwan
PACT4
2016 SpaceJMP: Programming with Multiple Virtual Address Spaces
abstract
Memory-centric computing demands careful organization of the virtual address space, but traditional methods for doing so are inflexible and inefficient. If an application wishes to address larger physical memory than virtual address bits allow, if it wishes to maintain pointer-based data structures beyond process lifetimes, or if it wishes to share large amounts of memory across simultaneously executing processes, legacy interfaces for managing the address space are cumbersome and often incur excessive overheads. We propose a new operating system design that promotes virtual address spaces to first-class citizens, enabling process threads to attach to, detach from, and switch between multiple virtual address spaces. Our work enables data-centric applications to utilize vast physical memory beyond the virtual range, represent persistent pointer-rich data structures without special pointer representations, and share large amounts of memory between processes efficiently.
Izzat El Hajj, Alex Merritt, Gerd Zellweger, Dejan S. Milojicic, Reto Achermann, Paolo Faraboschi, Wen-Mei W. Hwu, Timothy Roscoe, Karsten Schwan
ASPLOS9
2016 Landrush: Rethinking In-Situ Analysis for GPGPU Workflows
abstract
In-situ analysis on the output data of scientific simulations has been made necessary by ever-growing output data volumes and increasing costs of data movement as supercomputing is moving towards exascale. With hardware accelerators like GPUs becoming increasingly common in high end machines, new opportunities arise to co-locate scientific simulations and online analysis performed on the scientific data generated by the simulations. However, the asynchronous nature of GPGPU programming models and the limited context-switching capabilities on the GPU pose challenges to co-locating the scientific simulation and analysis on the same GPU. This paper dives deeper into these challenges to understand how best to co-locate analysis with scientific simulations on the GPUs in HPC clusters. Specifically, our 'Landrush' approach to GPU sharing proposes a solution that utilizes idle cycles on the GPU to provide an improved time-to-answer, that is, the total time to run the scientific simulation and analysis of the generated data. Landrush is demonstrated with experimental results obtained from leadership high-end applications on ORNL's Titan supercomputer, which show that (i) GPU-based scientific simulations have varying degrees of idle cycles to afford useful analysis task co-location, and (ii) the inability to context switch on the GPU at instruction granularity can be overcome by careful control of the analysis kernel launches and software-controlled early completion of analysis kernel executions. Results show that Landrush is superior in terms of time-to-answer compared to serially running simulations followed by analysis or by relying on the GPU driver and hardwired thread dispatcher to run analysis concurrently on a single GPU.
Anshuman Goswami, Yuan Tian 0004, Karsten Schwan, Fang Zheng 0003, Jeffrey Young 0001, Matthew Wolf, Greg Eisenhauer, Scott Klasky
CCGrid3
2016 FlashStager: Improving the Performance of SSD-Based Data Staging Systems via Write Redirection
abstract
When SSDs are used for in-situ execution of data-intensive scientific workflows, it is challenging to obtain consistently high I/O throughput because its I/O efficiency can be compromised for serving write and read requests simultaneously. This issue is so-called write-read interference. In this paper, we propose a novel scheme named FlashStager, which can isolate writes from reads using write redirection to improve data staging performance of SSDs by minimizing the write-read interference. Not only can it detect the interference, but also evaluate whether it is cost-effective to resolve it by executing the write redirection according to its correlation with write ratio and request size. Our experiments with both micro-benchmarks and real scientific applications show than FlashStager can improve I/O performance of staging by 40% on average.
Xuechen Zhang 0001, Fang Zheng 0003, Karsten Schwan, Matthew Wolf
CLUSTER3
2016 GraphIn: An Online High Performance Incremental Graph Processing Framework
Dipanjan Sengupta, Narayanan Sundaram, Theodore L. Willke, Jeffrey Young 0001, Matthew Wolf, Karsten Schwan
Euro-Par7
2016 Data tiering in heterogeneous memory systems
abstract
Memory-based data center applications require increasingly large memory capacities, but face the challenges posed by the inherent difficulties in scaling DRAM and also the cost of DRAM. Future systems are attempting to address these demands with heterogeneous memory architectures coupling DRAM with high capacity, low cost, but also lower performance, non-volatile memories (NVM) such as PCM and RRAM. A key usage model intended for NVM is as cheaper high capacity volatile memory. Data center operators are bound to ask whether this model for the usage of NVM to replace the majority of DRAM memory leads to a large slowdown in their applications? It is crucial to answer this question because a large performance impact will be an impediment to the adoption of such systems.
Subramanya Dulloor, Amitabha Roy 0002, Zheguang Zhao, Narayanan Sundaram, Nadathur Satish, Rajesh Sankaran, Jeff Jackson, Karsten Schwan
EuroSys8
2016 pVM: persistent virtual memory for efficient capacity scaling and object storage
abstract
Next-generation byte-addressable nonvolatile memories (NVMs), such as phase change memory (PCM) and Memristors, promise fast data storage, and more importantly, address DRAM scalability issues. State-of-the-art OS mechanisms for NVMs have focused on improving the block-based virtual file system (VFS) to manage both persistence and the memory capacity scaling needs of applications. However, using the VFS for capacity scaling has several limitations, such as the lack of automatic memory capacity scaling across DRAM and NVM, inefficient use of the processor cache and TLB, and high page access costs. These limitations reduce application performance and also impact applications that use NVM for persistent object storage with flat namespaces, such as photo stores, NoSQL databases, and others.
Sudarsun Kannan, Ada Gavrilovska, Karsten Schwan
EuroSys3
2016 Exploiting variability for energy optimization of parallel programs
abstract
In this paper we present optimizations that use DVFS mechanisms to reduce the total energy usage in scientific applications. Our main insight is that noise is intrinsic to large scale parallel executions and it appears whenever shared resources are contended. The presence of noise allows us to identify and manipulate any program regions amenable to DVFS. When compared to previous energy optimizations that make per core decisions using predictions of the running time, our scheme uses a qualitative approach to recognize the signature of executions amenable to DVFS. By recognizing the "shape of variability" we can optimize codes with highly dynamic behavior, which pose challenges to all existing DVFS techniques. We validate our approach using offline and online analyses for one-sided and two-sided communication paradigms. We have applied our methods to NWChem, and we show best case improvements in energy use of 12% at no loss in performance when using online optimizations running on 720 Haswell cores with one-sided communication. With NWChem on MPI two-sided and offline analysis, capturing the initialization, we find energy savings of up to 20%, with less than 1% performance cost.
Wim T. L. P. Lavrijsen, Costin Iancu, Wibe de Jong, Karsten Schwan
EuroSys5
2016 Phoenix: Memory Speed HPC I/O with NVM
abstract
In order to bridge the gap between the applications' I/O needs on future exascale platforms, and thecapabilities of conventional memory and storage technologies, HPC system designs started integrating components based onemerging non-volatile memory technologies. Non-volatile memory (NVRAM) provides persistent storage at close to memoryspeeds, with good capacity scaling, leading to opportunitiesto accelerate I/O in exascale machines. However, naive use ofNVRAM devices with current software stacks, exposes newbottlenecks due to the limited device bandwidth and slowerdevice access times compared to DRAM. To address this, we propose Phoenix (PHX), an NVRAM-bandwidth aware object store for persistent objects. PHXachieves efficiency through use of memory-centric objectinterfaces and device access stack specialized for NVRAM. Furthermore, PHX deals with the limited PCM bandwidththrough simultaneous use of NVRAM and DRAM devices, thus increasing the effective data movement bandwidth. Thisleads to reduction in the time length of the critical path I/Ooperations associated with the slow NVM device. To continueguaranteeing adequate reliability for the persistent objects, DRAM-resident object state is replicated across peer nodes'memory, accessible through high-bandwidth interconnects. Furthermore PHX minimizes the data movement overheads dueto additional data copies, by using a cost model that considersdevice bandwidths, remote storage distance and energy costs. Experimental analysis using real-world HPC applications onemulated NVRAM hardware shows that Phoenix's controlleduse of node-local and remote-node memory bandwidth, delivers up to ~ 1.2×, ~ 2× and ~ 12× speed-up for checkpoint I/Ofor the S3D, CM1 and GTC HPC applications, respectively. Furthermore PHX reduces total simulation checkpoint over-head of GTC up to ~ 18%.
Pradeep Fernando, Sudarsun Kannan, Ada Gavrilovska, Karsten Schwan
HiPC4
2016 Attribute-Based Partial Geo-Replication System
abstract
Existing partial geo-replication systems do not always provide optimal cost or latency, because their replication decisions are based on statically established data access popularity metrics, regardless of the application types. We demonstrate that additional reduction in cost and latency can be achieved by (1) using the right object attributes for making replication decisions for each type of application, (2) using multi-attribute-based replications, and (3) combining the popularity-based but reactive approach with the more random but proactive approach to data replication. Toward this end, we propose Acorn, an Attribute-based COntinuous partial geo-ReplicatioN system, and its prototype implementation based on Apache Cassandra. Experiments with two types of global-scale, data-sharing applications demonstrate up to 54% and 90% cost overhead reduction over existing systems or 38% and 91% latency overhead reduction.
Hobin Yoon, Ada Gavrilovska, Karsten Schwan
IC2E3
2016 Affinity-aware work-stealing for integrated CPU-GPU processors
abstract
Recent integrated CPU-GPU processors like Intel's Broadwell and AMD's Kaveri support hardware CPU-GPU shared virtual memory, atomic operations, and memory coherency. This enables fine-grained CPU-GPU work-stealing, but architectural differences between the CPU and GPU hurt the performance of traditionally-implemented work-stealing on such processors. These architectural differences include different clock frequencies, atomic operation costs, and cache and shared memory latencies. This paper describes a preliminary implementation of our work-stealing scheduler, Libra, which includes techniques to deal with these architectural differences in integrated CPU-GPU processors. Libra's affinity-aware techniques achieve significant performance gains over classically-implemented work-stealing. We show preliminary results using a diverse set of nine regular and irregular workloads running on an Intel Broadwell Core-M processor. Libra currently achieves up to a 2× performance improvement over classical work-stealing, with a 20% average improvement.
Naila Farooqui, Rajkishore Barik, Brian T. Lewis, Tatiana Shpeisman, Karsten Schwan
PPoPP5
2016 An Evolutionary Study of Linux Memory Management for Fun and Profit
Jian Huang 0006, Moinuddin K. Qureshi, Karsten Schwan
USENIX ATC3
2015 Understanding issue correlations: a case study of the Hadoop system
abstract
Over the last decade, Hadoop has evolved into a widely used platform for Big Data applications. Acknowledging its wide-spread use, we present a comprehensive analysis of the solved issues with applied patches in the Hadoop ecosystem. The analysis is conducted with a focus on Hadoop's two essential components: HDFS (storage) and MapReduce (computation), it involves a total of 4218 solved issues over the last six years, covering 2180 issues from HDFS and 2038 issues from MapReduce. Insights derived from the study concern system design and development, particularly with respect to correlated issues and correlations between root causes of issues and characteristics of the Hadoop subsystems. These findings shed light on the future development of Big Data systems, on their testing, and on bug-finding tools.
Jian Huang 0006, Xuechen Zhang 0001, Karsten Schwan
SoCC3
2015 ObsCon: Integrated Monitoring and Control for Parallel, Real-Time Applications
abstract
A large class of emerging compute-intensive applications demand real-time or near real-time processing guarantees on streaming data. Sensor processing in particular, has stringent latency requirements for carrying out its digital processing for rapidly incoming radar data streams. The consequent demands on the cluster middleware used to run such codes include (i) efficient online observation of current application performance, coupled with (ii) highly responsive controllers able to dynamically adjust the application's input-and data-dependent runtime behavior. We present the Obs(erver)Con(troller) software for online monitoring and control, which based on specifications of acceptable application states and tunable knobs within the execution environment, ensures that application performance falls within acceptable limits. ObsCon topologies are dynamic, making possible the runtime association of ObsCon methods with arbitrary DAG-structured, distributed/parallel stream processing applications running on high end cluster machines. This paper describes the ObsCon software and its 'grey box' use with a high performance cluster code that exports to ObsCon select 'hooks' for online monitoring and control -- Adaptive Digital Beamforming for a phase-array radar system.
Alan Nussbaum, Shwetha Mathangi Chandra Choodamani, Karsten Schwan
CLUSTER3
2015 SODA: Science-Driven Orchestration of Data Analytics
abstract
As scientific simulation applications evolve on the path towards exascale, a new model of scientific inquiry is required where concurrently with the running simulation, online analytics operate on the data it produces. By avoiding offline data storage except when absoluately necessary, it enables speeding up the scientific discovery process by providing rapid insights into the simulated science phenomena and affording more frequent, detailed data analytics than is possible with the traditional purely offline approach of using disk for intermediate data storage. However, a challenge for online analytics is to respond to behavior dynamics caused by changing simulation outputs and by unforeseen events on the underlying hardware/software platforms. This paper presents SODA, a set of run-time abstractions for online orchestration of data analytics, realized by embedding analytics tasks into workstations that monitor component behavior and enable responses to run-time changes in their resource demands and in the platform's resource availability. For high end simulations running on a leadership class machine, experimental evaluations show SODA can invoke efficient orchestration operations responding to a diverse set of run-time dynamics at different granularities to meet end-user and analysis specific requirements.
Jai Dayal, Jay F. Lofstead, Greg Eisenhauer, Karsten Schwan, Matthew Wolf, Hasan Abbasi, Scott Klasky
e-Science4
2015 Unified address translation for memory-mapped SSDs with FlashMap
abstract
Applications can map data on SSDs into virtual memory to transparently scale beyond DRAM capacity, permitting them to leverage high SSD capacities with few code changes. Obtaining good performance for memory-mapped SSD content, however, is hard because the virtual memory layer, the file system and the flash translation layer (FTL) perform address translations, sanity and permission checks independently from each other. We introduce FlashMap, an SSD interface that is optimized for memory-mapped SSD-files. FlashMap combines all the address translations into page tables that are used to index files and also to store the FTL-level mappings without altering the guarantees of the file system or the FTL. It uses the state in the OS memory manager and the page tables to perform sanity and permission checks respectively. By combining these layers, FlashMap reduces critical-path latency and improves DRAM caching efficiency. We find that this increases performance for applications by up to 3.32x compared to state-of-the-art SSD file-mapping mechanisms. Additionally, latency of SSD accesses reduces by up to 53.2%.
Jian Huang 0006, Anirudh Badam, Moinuddin K. Qureshi, Karsten Schwan
ISCA4
2015 Scaling iterative graph computations with GraphMap
abstract
In recent years, systems researchers have devoted considerable effort to the study of large-scale graph processing. Existing distributed graph processing systems such as Pregel, based solely on distributed memory for their computations, fail to provide seamless scalability when the graph data and their intermediate computational results no longer fit into the memory; and most distributed approaches for iterative graph computations do not consider utilizing secondary storage a viable solution. This paper presents GraphMap, a distributed iterative graph computation framework that maximizes access locality and speeds up distributed iterative graph computations by effectively utilizing secondary storage. GraphMap has three salient features: (1) It distinguishes data states that are mutable during iterative computations from those that are read-only in all iterations to maximize sequential access and minimize random access. (2) It entails a two-level graph partitioning algorithm that enables balanced workloads and locality-optimized data placement. (3) It contains a proposed suite of locality-based optimizations that improve computational efficiency. Extensive experiments on several real-world graphs show that GraphMap outperforms existing distributed memory-based systems for various iterative graph algorithms.
Kisung Lee, Ling Liu 0001, Karsten Schwan, Calton Pu, Qi Zhang 0009, Yang Zhou 0001, Emre Yigitoglu, Pingpeng Yuan
SC3
2015 GraphReduce: processing large-scale graphs on accelerator-based systems
abstract
Recent work on real-world graph analytics has sought to leverage the massive amount of parallelism offered by GPU devices, but challenges remain due to the inherent irregularity of graph algorithms and limitations in GPU-resident memory for storing large graphs. We present GraphReduce, a highly efficient and scalable GPU-based framework that operates on graphs that exceed the device's internal memory capacity. GraphReduce adopts a combination of edge- and vertex-centric implementations of the Gather-Apply-Scatter programming model and operates on multiple asynchronous GPU streams to fully exploit the high degrees of parallelism in GPUs with efficient graph data movement between the host and device. GraphReduce-based programming is performed via device functions that include gatherMap, gatherReduce, apply, and scatter, implemented by programmers for the graph algorithms they wish to realize. Extensive experimental evaluations for a wide variety of graph inputs and algorithms demonstrate that GraphReduce significantly outperforms other competing out-of-memory approaches.
Dipanjan Sengupta, Shuaiwen Song, Kapil Agarwal, Karsten Schwan
SC4
2015 HeteroVisor: Exploiting Resource Heterogeneity to Enhance the Elasticity of Cloud Platforms
abstract
This paper presents HeteroVisor, a heterogeneity-aware hypervisor, that exploits resource heterogeneity to enhance the elasticity of cloud systems. Introducing the notion of 'elasticity' (E) states, HeteroVisor permits applications to manage their changes in resource requirements as state transitions that implicitly move their execution among heterogeneous platform components. Masking the details of platform heterogeneity from virtual machines, the E-state abstraction allows applications to adapt their resource usage in a fine-grained manner via VM-specific 'elasticity drivers' encoding VM-desired policies. The approach is explored for the heterogeneous processor and memory subsystems evolving for modern server platforms, leading to mechanisms that can manage these heterogeneous resources dynamically and as required by the different VMs being run. HeteroVisor is implemented for the Xen hypervisor, with mechanisms that go beyond core scaling to also deal with memory resources, via the online detection of hot memory pages and transparent page migration. Evaluation on an emulated heterogeneous platform uses workload traces from real-world data, demonstrating the ability to provide high on-demand performance while also reducing resource usage for these workloads.
Vishal Gupta 0001, Min Lee, Karsten Schwan
VEE3
2015 A Framework for Emulating Non-Volatile Memory Systemswith Different Performance Characteristics
abstract
Exponential 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
ICPE7
2014 Flexpath: Type-Based Publish/Subscribe System for Large-Scale Science Analytics
abstract
As high-end systems move toward exascale sizes, a new model of scientific inquiry being developed is one in which online data analytics run concurrently with the high end simulations producing data outputs. Goals are to gain rapid insights into the ongoing scientific processes, assess their scientific validity, and/or initiate corrective or supplementary actions by launching additional computations when needed. The Flex path system presented in this paper addresses the fundamental problem of how to structure and efficiently implement the communications between high end simulations and concurrently running online data analytics, the latter comprised of componentized dynamic services and service pipelines. Using a type-based publish/subscribe approach, Flexpath encourages diversity by permitting analytics services to differ in their computational and scaling characteristics and even in their internal execution models. Flex path uses direct and MxN connections between interacting services to reduce data movements, to allow for runtime connectivity changes to accommodate component arrivals/departures, and to support the multiple underlying communication protocols used for analytics workflows in which simulation outputs are processed by analytics services residing on the same nodes where they are generated, on the same machine, and/or on attached or remote analytics engines. This paper describes the design and implementation of Flex path, and evaluates it with two widely used scientific applications and their associated data analytics methods.
Jai Dayal, Drew Bratcher, Greg Eisenhauer, Karsten Schwan, Matthew Wolf, Xuechen Zhang 0001, Hasan Abbasi, Scott Klasky, Norbert Podhorszki
CCGRID4
2014 Merlin: Application- and Platform-aware Resource Allocation in Consolidated Server Systems
abstract
Workload consolidation, whether via use of virtualization or with lightweight, container-based methods, is critically important for current and future datacenter and cloud computing systems. Yet such consolidation challenges the ability of current systems to meet application resource needs and isolate their resource shares, particularly for high core count or 'scaleup' servers. This paper presents the 'Merlin' approach to managing the resources of multicore platforms, which satisfies an application's resource requirements efficiently -- using low cost allocations -- and improves isolation -- measured as increased predictability of application execution. Merlin (i) creates a virtual platform (VP) as a system-level resource commitment to an application's resource shares, (ii) enforces its isolation, and (iii) operates with low runtime overhead. Further, Merlin's resource (re)-allocation and isolation methods operate by constructing online models that capture the resource 'sensitivities' of the currently running applications along all of their resource dimensions. Elevating isolation into a first-class management principle, these sensitivity- and cost-based allocation and sharing methods lead to efficient methods for shared resource use on scaleup server systems. Experimental evaluations on a large core-count machine demonstrate improved performance with reduced performance variation and increased system throughput and efficiency, for a wide range of popular datacenter workloads, compared with the methods used in prior work and with the state-of-art Xen hypervisor.
Priyanka Tembey, Ada Gavrilovska, Karsten Schwan
SoCC3
2014 HeteroCheckpoint: Efficient Checkpointing for Accelerator-Based Systems
abstract
Moving toward exascale, the number of GPUs in HPC machines is bound to increase, and applications will spend increasing amounts of time running on those GPU devices. While GPU usage has already led to substantial speedup for HPC codes, their failure rates due to overheating are at least 10 times higher than those seen for the CPUs now commonly used on HPC machines. This makes it increasingly important for GPUs to have robust checkpoint/restart mechanisms. This paper introduces a unified CPU-GPU checkpoint mechanism, which can efficiently checkpoint the combined GPU-CPU memory state resident on machine nodes. Efficiency is gained in part by addressing the end-to-end data movements required for check pointing - from GPU to storage - by introducing novel pre-copy and checksum methods. These methods reduce checkpoint data movement cost seen by HPC applications, with initial measurements using different benchmark applications showing up to 60% reduced checkpoint overhead. Additional exploration of the use of next-generation storage, like NVM, show further promises of reduced check pointing overheads.
Sudarsun Kannan, Naila Farooqui, Ada Gavrilovska, Karsten Schwan
DSN4
2014 Reducing the cost of persistence for nonvolatile heaps in end user devices
abstract
This paper explores the performance implications of using future byte addressable non-volatile memory (NVM) like PCM in end client devices. We explore how to obtain dual benefits - increased capacity and faster persistence - with low overhead and cost. Specifically, while increasing memory capacity can be gained by treating NVM as virtual memory, its use of persistent data storage incurs high consistency (frequent cache flushes) and durability (logging for failure) overheads, referred to as `persistence cost'. These not only affect the applications causing them, but also other applications relying on the same cache and/or memory hierarchy. This paper analyzes and quantifies in detail the performance overheads of persistence, which include (1) the aforementioned cache interference as well as (2) memory allocator overheads, and finally, (3) durability costs due to logging. Novel solutions to overcome such overheads include (1) a page contiguity algorithm that reduces interference-related cache misses, (2) a cache efficient NVM write aware memory allocator that reduces cache line flushes of allocator state by 8X, and (3) hybrid logging that reduces durability overheads substantially. With these solutions, experimental evaluations with different end user applications and SPEC2006 benchmarks show up to 12% reductions in cache misses, thereby reducing the total number of NVM writes.
Sudarsun Kannan, Ada Gavrilovska, Karsten Schwan
HPCA3
2014 Personal clouds: Sharing and integrating networked resources to enhance end user experiences
abstract
End user experiences on mobile devices with their rich sets of sensors are constrained by limited device battery lives and restricted form factors, as well as by the `scope' of the data available locally. The `Personal Cloud' distributed software abstractions address these issues by enhancing the capabilities of a mobile device via seamless use of both nearby and remote cloud resources. In contrast to vendor-specific, middleware-based cloud solutions, Personal Cloud instances are created at hypervisor-level, to create for each end user the federation of networked resources best suited for the current environment and use. Specifically, the Cirrostratus extensions of the Xen hypervisor can federate a user's networked resources to establish a personal execution environment, governed by policies that go beyond evaluating network connectivity to also consider device ownership and access rights, the latter managed in a secure fashion via standard Social Network Services. Experimental evaluations with both Linux- and Android-based devices, and using Facebook as the SNS, show the approach capable of substantially augmenting a device's innate capabilities, improving application performance and the effective functionality seen by end users.
Minsung Jang, Karsten Schwan, Ketan Bhardwaj, Ada Gavrilovska, Adhyas Avasthi
INFOCOM2
2014 Scibox: Online Sharing of Scientific Data via the Cloud
abstract
Collaborative science demands global sharing of scientific data. But it cannot leverage universally accessible cloud-based infrastructures like Drop Box, as those offer limited interfaces and inadequate levels of access bandwidth. We present the Scibox cloud facility for online sharing scientific data. It uses standard cloud storage solutions, but offers a usage model in which high end codes can write/read data to/from the cloud via the APIs they already use for their I/O actions. With Scibox, data upload/download volumes are controlled via Data Reduction-functions stated by end users and applied at the data source, before data is moved, with further gains in efficiency obtained by combining DR-functions to move exactly what is needed by current data consumers. We evaluate Scibox with science applications and their representative data analytics - the GTS fusion and the combustion image processing - demonstrating the potential for ubiquitous data access with substantial reductions in network traffic.
Jian Huang 0006, Xuechen Zhang 0001, Greg Eisenhauer, Karsten Schwan, Matthew Wolf, Stéphane Ethier, Scott Klasky
IPDPS4
2014 Scheduling Multi-tenant Cloud Workloads on Accelerator-Based Systems
abstract
Accelerator-based systems are making rapid inroads into becoming platforms of choice for high end cloud services. There is a need therefore, to move from the current model in which high performance applications explicitly and programmatically select the GPU devices on which to run, to a dynamic model where GPUs are treated as first class schedulable entities. The Strings scheduler realizes this vision by decomposing the GPU scheduling problem into a combination of load balancing and per-device scheduling. (i) Device-level scheduling efficiently uses all of a GPU's hardware resources, including its computational and data movement engines, and (ii) load balancing goes beyond obtaining high throughput, to ensure fairness through prioritizing GPU requests that have attained least service. With its methods, Strings achieves improvements in system throughput and fairness of up to 8.70× and 13%, respectively, compared to the CUDA runtime.
Dipanjan Sengupta, Anshuman Goswami, Karsten Schwan, Krishna Pallavi
SC3
2014 ELF: Efficient Lightweight Fast Stream Processing at Scale
Liting Hu, Karsten Schwan, Hrishikesh Amur
USENIX ATC2
2014 Hello ADIOS: the challenges and lessons of developing leadership class I/O frameworks
abstract
SUMMARY Applications running on leadership platforms are more and more bottlenecked by storage input/output (I/O). In an effort to combat the increasing disparity between I/O throughput and compute capability, we created Adaptable IO System (ADIOS) in 2005. Focusing on putting users first with a service oriented architecture, we combined cutting edge research into new I/O techniques with a design effort to create near optimal I/O methods. As a result, ADIOS provides the highest level of synchronous I/O performance for a number of mission critical applications at various Department of Energy Leadership Computing Facilities. Meanwhile ADIOS is leading the push for next generation techniques including staging and data processing pipelines. In this paper, we describe the startling observations we have made in the last half decade of I/O research and development, and elaborate the lessons we have learned along this journey. We also detail some of the challenges that remain as we look toward the coming Exascale era. Copyright © 2013 John Wiley & Sons, Ltd.
Qing Liu 0002, Jeremy Logan, Yuan Tian 0004, Hasan Abbasi, Norbert Podhorszki, Jong Choi 0001, Scott Klasky, Roselyne Tchoua, Jay F. Lofstead, Ron A. Oldfield, Manish Parashar, Nagiza F. Samatova, Karsten Schwan, Arie Shoshani, Matthew Wolf, Kesheng Wu, Weikuan Yu
Concurr. Comput. Pract. Exp.13
2014 Dynamic core affinity for high-performance file upload on Hadoop Distributed File System
Joong-Yeon Cho, Hyun-Wook Jin, Min Lee, Karsten Schwan
Parallel Comput.4
2014 NVRAM-aware Logging in Transaction Systems
abstract
Emerging byte-addressable, non-volatile memory technologies (NVRAM) like phase-change memory can increase the capacity of future memory systems by orders of magnitude. Compared to systems that rely on disk storage, NVRAM-based systems promise significant improvements in performance for key applications like online transaction processing (OLTP). Unfortunately, NVRAM systems suffer from two drawbacks: their asymmetric read-write performance and the notable higher cost of the new memory technologies compared to disk. This paper investigates the cost-effective use of NVRAM in transaction systems. It shows that using NVRAM only for the logging subsystem ( NV-Logging ) provides much higher transactions per dollar than simply replacing all disk storage with NVRAM. Specifically, for NV-Logging , we show that the software overheads associated with centralized log buffers cause performance bottlenecks and limit scaling. The per-transaction logging methods described in the paper help avoid these overheads, enabling concurrent logging for multiple transactions. Experimental results with a faithful emulation of future NVRAM-based servers using the TPCC, TATP, and TPCB benchmarks show that NV-Logging improves throughput by 1.42 - 2.72x over the costlier option of replacing all disk storage with NVRAM. Results also show that NV-Logging performs 1.21 - 6.71x better than when logs are placed into the PMFS NVRAM-optimized file system. Compared to state-of-the-art distributed logging, NV-Logging delivers 20.4% throughput improvements.
Jian Huang 0006, Karsten Schwan, Moinuddin K. Qureshi
Proc. VLDB Endow.2
2013 Memory-efficient groupby-aggregate using compressed buffer trees
abstract
The rapid growth of fast analytics systems, that require data processing in memory, makes memory capacity an increasingly-precious resource. This paper introduces a new compressed data structure called a Compressed Buffer Tree (CBT). Using a combination of techniques including buffering, compression, and serialization, CBTs improve the memory efficiency and performance of the GroupBy-Aggregate abstraction that forms the basis of not only batch-processing models like MapReduce, but recent fast analytics systems too. For streaming workloads, aggregation using the CBT uses 21--42% less memory than using Google SparseHash with up to 16% better throughput. The CBT is also compared to batch-mode aggregators in MapReduce runtimes such as Phoenix++ and Metis and consumes 4x and 5x less memory with 1.5--2x and 3--4x more performance respectively.
Hrishikesh Amur, Wolfgang Richter 0001, David G. Andersen, Michael Kaminsky, Karsten Schwan, Athula Balachandran, Erik Zawadzki
SoCC5
2013 FastMR: fast processing for large distributed data streams
abstract
FastMR is a graph-style framework for steam-oriented applications to realize near real-time streaming data record processing, and more importantly, complex coordinations between those applications. We introduces two components --- compressed buffer trees (CBTs) and shared reducer trees (SRTs) to assist with this task. CBTs address the problem of maintaining a significant amount of application-specific "accumulator" state in memory so that streaming data processing can combine current data with historical data. They do so by employing a novel, batch-oriented approach to updating the accumulator state. SRTs are basically P2P-based reducer trees that enable fine-grained queries (both one-shot and continual) to be efficiently rolled up concurrently. CBT's intermediate results are aggregated to the root of SRT via network aggregation. The roots of SRTs are analogous to vertices and anycast/multicast message transmission between the vertices (roots of SRTs) are analogous to edges in the graph-style computation model.
Liting Hu, Karsten Schwan, Hrishikesh Amur
SoCC2
2013 Distributed resource exchange: Virtualized resource management for SR-IOV InfiniBand clusters
abstract
The commoditization of high performance interconnects, like 40+ Gbps InfiniBand, and the emergence of low-overhead I/O virtualization solutions based on SR-IOV, is enabling the proliferation of such fabrics in virtualized datacenters and cloud computing platforms. As a result, such platforms are better equipped to execute workloads with diverse I/O requirements, ranging from throughput-intensive applications, such as `big data' analytics, to latency-sensitive applications, such as online applications with strict response-time guarantees. Improvements are also seen for the virtualization infrastructures used in data center settings, where high virtualized I/O performance supported by high-end fabrics enables more applications to be configured and deployed in multiple VMs - VM ensembles (VMEs) - distributed and communicating across multiple datacenter nodes. A challenge for I/O-intensive VM ensembles is the efficient management of the virtualized I/O and compute resources they share with other consolidated applications, particularly in lieu of VME-level SLA requirements like those pertaining to low or predictable end-to-end latencies for applications comprised of sets of interacting services. This paper addresses this challenge by presenting a management solution able to consider such SLA requirements, by supporting diverse SLA-aware policies, such as those maintaining bounded SLA guarantees for all VMEs, or those that minimize the impact of misbehaving VMEs. The management solution, termed Distributed Resource Exchange (DRX), borrows techniques from principles of microeconomics, and uses online resource pricing methods to provide mechanisms for such distributed and coordinated resource management. DRX and its mechanisms allow policies to be deployed on such a cluster in order to provide SLA guarantees to some applications by charging all the interfering VMEs `equally' or based on the `hurt', i.e. amount of I/O performed by the VMEs. While these mechanisms are general, our implementation is specifically for SR-IOV-based fabrics like InfiniBand and the KVM hypervisor. Our experimental evaluation consists of workloads representative of data-analytics, transactional and parallel benchmarks. The results demonstrate the feasibility of DRX and its utility to maintain SLA for transactional applications. We also show that the impact to the interfering workloads is also within acceptable bounds for certain policies.
Adit Ranadive, Ada Gavrilovska, Karsten Schwan
CLUSTER3
2013 A-Cache: Resolving cache interference for distributed storage with mixed workloads
abstract
Distributed key-value stores employ large main memory caches to mitigate the high costs of disk access. A challenge for such caches is that large scale distributed stores simultaneously face multiple workloads, often with drastically different characteristics. Interference between such competing workloads leads to performance degradation through inefficient use of the main memory cache. This paper diagnoses the cache interference seen for representative workloads and then develops A-Cache, an adaptive set of main memory caching methods for distributed key-value stores. Focused on read performance for common workload patterns, A-Cache leads to throughput improvements of up to 150% for competing data-intensive applications running on server class machines.
Bharath Ravi, Hrishikesh Amur, Karsten Schwan
CLUSTER3
2013 Oncilla: A GAS runtime for efficient resource allocation and data movement in accelerated clusters
abstract
Accelerated and in-core implementations of Big Data applications typically require large amounts of host and accelerator memory as well as efficient mechanisms for transferring data to and from accelerators in heterogeneous clusters. Scheduling for heterogeneous CPU and GPU clusters has been investigated in depth in the high-performance computing (HPC) and cloud computing arenas, but there has been less emphasis on the management of cluster resource that is required to schedule applications across multiple nodes and devices. Previous approaches to address this resource management problem have focused on either using low-performance software layers or on adapting complex data movement techniques from the HPC arena, which reduces performance and creates barriers for migrating applications to new heterogeneous cluster architectures. This work proposes a new system architecture for cluster resource allocation and data movement built around the concept of managed Global Address Spaces (GAS), or dynamically aggregated memory regions that span multiple nodes.We propose a software layer called Oncilla that uses a simple runtime and API to take advantage of non-coherent hardware support for GAS. The Oncilla runtime is evaluated using two different high-performance networks for microkernels representative of the TPC-H data warehousing benchmark, and this runtime enables a reduction in runtime of up to 81%, on average, when compared with standard disk-based data storage techniques. The use of the Oncilla API is also evaluated for a simple breadth-first search (BFS) benchmark to demonstrate how existing applications can incorporate support for managed GAS.
Jeffrey Young 0001, Se Hoon Shon, Sudhakar Yalamanchili, Alex Merritt, Karsten Schwan, Holger Fröning
CLUSTER5
2013 FlexQuery: An online query system for interactive remote visual data exploration at large scale
abstract
The remote visual exploration of live data generated by scientific simulations is useful for scientific discovery, performance monitoring, and online validation for the simulation results. Online visualization methods are challenged, however, by the continued growth in the volume of simulation output data that has to be transferred from its source - the simulation running on the high end machine - to where it is analyzed, visualized, and displayed. A specific challenge in this context is limits in the communication bandwidth between data source(s) and sinks. Previous work places queries `near' data sources, exploiting their data reduction capabilities, but such work does not address the common scenario in which scientists make multiple different queries on the data being produced. This paper considers the general case in which science users are interested in different (sub)sets of the data produced by a high end simulation. We offer the FlexQuery online data query system that can deploy and execute data queries `along' the I/O and analytics pipelines. FlexQuery carefully extends such analytics pipelines, using online performance monitoring and data location tracking, to realize data queries in ways that minimize additional data movement and offer low latency in data query execution. Using a real-world scientific application - the Maya astrophysics code and its analytics workflow - we demonstrate FlexQuery's ability to dynamically deploy queries for low-latency remote data visualization.
Hongbo Zou, Karsten Schwan, Magdalena Slawiñska, Matthew Wolf, Greg Eisenhauer, Fang Zheng 0003, Jai Dayal, Jeremy Logan, Qing Liu 0002, Scott Klasky, Tanja Bode, Michael Clark, Matthew Kinsey
CLUSTER2
2013 Cache Topology Aware Mapping of Stream Processing Applications onto CMPs
abstract
Data Stream Processing is an important class of data intensive applications in the "Big Data" era. Chip Multi-Processors (CMPs) are the standard hosting platforms in modern data centers. Gaining high performance for stream processing applications on CMPs is therefore of great interest. Since the performance of stream processing applications largely depends on their effective use of the complex cache structure present on CMPs, this paper proposes the StreamMap approach for tuning streaming applications' use of cache. Our major idea is to map application threads to CPU cores to facilitate data sharing AND mitigate memory resource contention among threads in a holistic manner. Applying StreamMap to the IBM's System S middleware leads to improvements of up to 1.8x in the performance of realistic applications over standard Linux OS scheduler on three different CMP platforms.
Fang Zheng 0003, Chitra Venkatramani, Rohit Wagle, Karsten Schwan
ICDCS4
2013 Optimizing Checkpoints Using NVM as Virtual Memory
abstract
Rapid checkpointing will remain key functionality for next generation high end machines. This paper explores the use of node-local nonvolatile memories (NVM) such as phase-change memory, to provide frequent, low overhead checkpoints. By adapting existing multi-level checkpoint techniques, we devise new methods, termed NVM-checkpoints, that efficiently store checkpoints on both local and remote node NVM. The checkpoint frequencies are guided by failure models that capture the expected accessibility of such data after failure. To lower overheads, NVM-checkpoints reduce the NVM and interconnect bandwidth used with a novel pre-copy mechanism, which incrementally moves checkpoint data from DRAM to NVM before a local checkpoint is started. This reduces local checkpoint cost by limiting the instantaneous data volume moved at checkpoint time, thereby freeing bandwidth for use by applications. In fact, the pre-copy method can reduce peak interconnect usage up to 46%. Since our approach treats NVM as memory rather than as 'Ramdisk', pre-copying can be generalized to directly move data to remote NVMs. This results in 40% faster application execution times compared to asynchronous approaches not using pre-copying.
Sudarsun Kannan, Ada Gavrilovska, Karsten Schwan, Dejan S. Milojicic
IPDPS3
2013 FlexIO: I/O Middleware for Location-Flexible Scientific Data Analytics
abstract
Increasingly severe I/O bottlenecks on High-End Computing machines are prompting scientists to process simulation output data online while simulations are running and before storing data on disk. There are several options to place data analytics along the I/O path: on compute nodes, on separate nodes dedicated to analytics, or after data is stored on persistent storage. Since different placements have different impact on performance and cost, there is a consequent need for flexibility in the location of data analytics. The FlexIO middleware described in this paper makes it easy for scientists to obtain such flexibility, by offering simple abstractions and diverse data movement methods to couple simulation with analytics. Various placement policies can be built on top of FlexIO to exploit the trade-offs in performing analytics at different levels of the I/O hierarchy. Experimental results demonstrate that FlexIO can support a variety of simulation and analytics workloads at large scale through flexible placement options, efficient data movement, and dynamic deployment of data manipulation functionalities.
Fang Zheng 0003, Hongbo Zou, Greg Eisenhauer, Karsten Schwan, Matthew Wolf, Jai Dayal, Jianting Cao, Hasan Abbasi, Scott Klasky, Norbert Podhorszki, Hongfeng Yu 0001
IPDPS4
2013 GoldRush: resource efficient in situ scientific data analytics using fine-grained interference aware execution
abstract
Severe I/O bottlenecks on High End Computing platforms call for running data analytics in situ. Demonstrating that there exist considerable resources in compute nodes un-used by typical high end scientific simulations, we leverage this fact by creating an agile runtime, termed GoldRush, that can harvest those otherwise wasted, idle resources to efficiently run in situ data analytics. GoldRush uses fine-grained scheduling to "steal" idle resources, in ways that minimize interference between the simulation and in situ analytics. This involves recognizing the potential causes of on-node resource contention and then using scheduling methods that prevent them. Experiments with representative science applications at large scales show that resources harvested on compute nodes can be leveraged to perform useful analytics, significantly improving resource efficiency, reducing data movement costs incurred by alternate solutions, and posing negligible impact on scientific simulations.
Fang Zheng 0003, Hongfeng Yu 0001, Can Hantas, Matthew Wolf, Greg Eisenhauer, Karsten Schwan, Hasan Abbasi, Scott Klasky
SC6
2013 Practical Compute Capacity Management for Virtualized Datacenters
abstract
We present CCM (Cloud Capacity Manager) - a prototype system and its methods for dynamically multiplexing the compute capacity of virtualized datacenters at scales of thousands of machines, for diverse workloads with variable demands. Extending prior studies primarily concerned with accurate capacity allocation and ensuring acceptable application performance, CCM also sheds light on the tradeoffs due to two unavoidable issues in large scale commodity datacenters: (i) maintaining low operational overhead given variable cost of performing management operations necessary to allocate resources, and (ii) coping with the increased incidences of these operations' failures. CCM is implemented in an industry-strength cloud infrastructure built on top of the VMware vSphere virtualization platform and is currently deployed in a 700 physical host datacenter. Its experimental evaluation uses production workload traces and a suite of representative cloud applications to generate dynamic scenarios. Results indicate that the pragmatic cloud-wide nature of CCM provides up to 25% more resources for workloads and improves datacenter utilization by up to 20%, compared to the common alternative approach of multiplexing capacity within multiple independent smaller datacenter partitions.
Mukil Kesavan, Irfan Ahmad 0005, Orran Krieger, Ravi Soundararajan, Ada Gavrilovska, Karsten Schwan
IEEE Trans. Cloud Comput.6
2012 Region scheduling: efficiently using the cache architectures via page-level affinity
abstract
The performance of modern many-core platforms strongly depends on the effectiveness of using their complex cache and memory structures. This indicates the need for a memory-centric approach to platform scheduling, in which it is the locations of memory blocks in caches rather than CPU idleness that determines where application processes are run. Using the term 'memory region' to denote the current set of physical memory pages actively used by an application, this paper presents and evaluates region-based scheduling methods for multicore platforms. This involves (i) continuously and at runtime identifying the memory regions used by executable entities, and their sizes, (ii) mapping these regions to caches to match performance goals, and (iii) maintaining region to cache mappings by ensuring that entities run on processors with direct access to the caches containing their regions. Region scheduling can implement policies that (i) offer improved performance to applications by 'unifying' the multiple caches present on the underlying physical machine and/or by 'balancing' cache usage to take maximum advantage of available cache space, (ii) better isolate applications from each other, particularly when their performance is strongly affected by cache availability, and also (iii) take advantage of standard scheduling and CPU-based load balancing when regioning is ineffective. The paper describes region scheduling and its system-level implementation and evaluates its performance with micro-benchmarks and representative multi-core applications. Single applications see performance improvements of up to 15% with region scheduling, and we observe 40% latency improvements when a platform is shared by multiple applications. Superior isolation is shown to be particularly important for cache-sensitive or real-time codes.
Min Lee, Karsten Schwan
ASPLOS2
2012 Interactive Use of Cloud Services: Amazon SQS and S3
abstract
Interactive use of cloud services is of keen interest to science end users, including for storing and accessing shared data sets. This paper evaluates the viability of interactively using two important cloud services offered by Amazon: SQS (Simple Queue Service) and S3 (Simple Storage Service). Specifically, we first measure the send-to-receive message latencies of SQS and then determine and devise rate controls to obtain suitable latencies and latency variations. Second, for S3, when transferring data into the cloud, we determine that increased parallelism in Transfer Manager can significantly improve upload performance, achieving up to 4 times improvements with careful elimination of upload bottlenecks.
Hobin Yoon, Ada Gavrilovska, Karsten Schwan, Jim Donahue
CCGRID3
2012 D2T: Doubly Distributed Transactions for High Performance and Distributed Computing
abstract
Current exascale computing projections suggest rather than a monolithic simulation running for the majority of the machine, a collection of components comprising the scientific discovery process will be employed in an online workflow. This move to an online workflow scenario requires knowledge that inter-step operations are completed and correct before the next phase begins. Further, dynamic load balancing or fault tolerance techniques may dynamically deploy or redeploy resources for optimal use of computing resources. These newly configured resources should only be used if they are successfully deployed. Our D2T system offers a mechanism to support these kinds of operations by providing database-like transactions with distributed servers and clients. Ultimately, with adequate hardware support, full ACID compliance is possible for the transactions. To prove the viability of this approach, we show that the D2T protocol has less than 1.2 seconds of overhead using 4096 clients and 32 servers with good scaling characteristics using this initial prototype implementation.
Jay F. Lofstead, Jai Dayal, Karsten Schwan, Ron A. Oldfield
CLUSTER3
2012 v-Bundle: Flexible Group Resource Offerings in Clouds
abstract
Traditional Infrastructure-as-a-Service offerings provide customers with large numbers of fixed-size virtual machine (VM) instances with resource allocations that are designed to meet application demands. With application demands varying over time, cloud providers gain efficiencies through resource consolidation and over-commitment. For cloud customers, however, this leads to inefficient use of the cloud resources they have purchased. To address cloud customers' dynamic application requirements, we present a new cloud resource offering, called v-Bundle, which makes flexible the exchange of resource capacity among multiple VM instances belonging to the same customer. Specifically targeting network resources, for each customer application, we first use DHT-based techniques to achieve an initial VM placement that minimizes its use of the data center network's bi-section bandwidth. When VMs' networking requirements change, the customer can then use v-Bundle to trade the networking resources allocated to her application. v-Bundle maintains information about network resources with any-cast tree-based methods implemented as extensions of the Pastry pub-sub core. Experimental evaluations show that the approach can scale well to thousands of hosts and VMs, and that v-Bundle can provide customers with better bandwidth utilization and improved application quality of service through borrowing extra bandwidth when needed, at no additional cost in terms of the total resources allocated to the customer.
Liting Hu, Kyung Dong Ryu, Dilma Da Silva, Karsten Schwan
ICDCS4
2012 Lynx: A dynamic instrumentation system for data-parallel applications on GPGPU architectures
abstract
As parallel execution platforms continue to proliferate, there is a growing need for real-time introspection tools to provide insight into platform behavior for performance debugging, correctness checks, and to drive effective resource management schemes. To address this need, we present the Lynx dynamic instrumentation system. Lynx provides the capability to write instrumentation routines that are (1) selective, instrumenting only what is needed, (2) transparent, without changes to the applications' source code, (3) customizable, and (4) efficient. Lynx is embedded into the broader GPU Ocelot system, which provides run-time code generation of CUDA programs for heterogeneous architectures. This paper describes (1) the Lynx framework and implementation, (2) its language constructs geared to the Single Instruction Multiple Data (SIMD) model of data-parallel programming used in current general-purpose GPU (GPGPU) based systems, and (3) useful performance metrics described via Lynx's instrumentation language that provide insights into the design of effective instrumentation routines for GPGPU systems. The paper concludes with a comparative analysis of Lynx with existing GPU profiling tools and a quantitative assessment of Lynx's instrumentation performance, providing insights into optimization opportunities for running instrumented GPU kernels.
Naila Farooqui, Andrew Kerr, Greg Eisenhauer, Karsten Schwan, Sudhakar Yalamanchili
ISPASS4
2012 VScope: Middleware for Troubleshooting Time-Sensitive Data Center Applications
Chengwei Wang, Infantdani Abel Rayan, Greg Eisenhauer, Karsten Schwan, Vanish Talwar, Matthew Wolf, Chad Huneycutt
Middleware4
2012 The Forgotten 'Uncore': On the Energy-Efficiency of Heterogeneous Cores
Vishal Gupta 0001, Paul Brett, David A. Koufaty, Dheeraj Reddy, Scott Hahn, Karsten Schwan, Ganapati Srinivasa
USENIX ATC6
2011 STRATUS: Assembling Virtual Platforms from Device Clouds
abstract
There is an increasing number of network-enabled computing devices in homes and offices, driven by continued improvements in device capabilities and network connectivity. By exploiting the virtualization technologies that have begun to pervade even the mobile domain, these devices -- hardware components, such as displays, input devices, disks, or processors, can be decoupled from the physical platforms on which they reside to form a resource pool or device cloud. By drawing on the composite resources of device clouds, applications can leverage the heterogeneity present in the cloud to exploit hardware/device differences in terms of power consumption, computational speeds, display sizes, or the presence of certain accelerators, and can take advantage of software diversity in terms of the different operating environments and applications that efficiently operate on individual devices. This paper implements and evaluates the concept of device clouds, in which virtual execution platforms dynamically composed from sets of devices are built for applications, using automated methods that are based on simple policies. Experimental results identify the basic overheads associated with device clouds and their use, and demonstrate the advantages of dynamically constructed virtual platforms rather than individual machines, both in terms of improvements in system properties like power usage and improvements in user experiences for media delivery and play out.
Minsung Jang, Karsten Schwan
IEEE CLOUD2
2011 HEaRS: A Hierarchical Energy-Aware Resource Scheduler for Virtualized Data Centers
abstract
With the increasing popularity of Internet-based cloud services, energy efficiency in large-scale Internet data centers has become important not only to curtail energy costs and alleviate environmental concern, but also because such systems can quickly reach the limits of power available to them. This paper investigates to what extent and how energy usage improvements through consolidation can benefit from taking into account the environmental influences and effects seen in data center systems. Toward that end, we present experimental results obtained in a fully instrumented, small scale data center and then use these results to propose a hierarchical energy-aware resource scheduler (HEaRS) for cluster workload placement and server provisioning, also considers the physical environment in which data center systems operate. Specifically, at the rack level, HEaRS tries to maintain a 'thermal balance' across the rack to avoid hot spots and reduce cooling costs. At the chassis level, HEaRS utilizes the proportional plus integral controller to achieve a balance in the levels of usage of electrical current between the two power domains in the chassis, which helps the chassis reach its most energy efficient state. Finally, at server level, HEaRS can employ known methods like dynamic voltage and frequency scaling or core idling to reduce power consumption. This results in a hierarchical set of controllers that jointly, implement holistic solutions to energy-aware resource scheduling for an entire rack, and this hierarchical solution can then be further extended to entire data centers. Our initial experiment result show opportunities for gains, with up to 16% in energy usage compared to methods that are not aware of the physical environment and up to 15% improvements in application performance.
Meina Song, Junde Song, Ada Gavrilovska, Karsten Schwan
CLUSTER5
2011 ResourceExchange: Latency-Aware Scheduling in Virtualized Environments with High Performance Fabrics
abstract
Virtualized infrastructures have seen strong acceptance in data center systems and applications, but have not yet seen adoptance for latency-sensitive codes which require I/O to arrive predictability, or response times to be generated within certain timeliness guarantees. Examples of such applications include certain classes of parallel HPC codes, server systems performing phonecall or multimedia delivery, or financial services in electronic trading platforms, like ICE and CME. In this paper, we argue that the use of high-performance, VMM-bypass capable devices can help create the virtualized infrastructures needed for the latency-sensitive applications listed above. However, to enable consolidation, problems to be solved go beyond efficient I/O virtualization, and include dealing with the shared use of I/O and compute resource, in ways that minimize or eliminate interference. Toward this end, we describe ResEx -- a resource management approach for virtualized RDMA-based platforms which incorporates concepts from supply-demand theory and congestion pricing to dynamically control the allocation of CPU and I/O resources of guest VMs. ResEx and its mechanisms and abstractions allow multiple 'pricing policies' to be deployed on these types of virtualized platforms, including such which reduce interference and enhance isolation by identifying and taxing VMs responsible for resource congestion. While the main ideas behind ResEx are more general, the design presented in this paper is specific for InfiniBand RDMA-based virtualized platforms due to the use of asynchronous monitoring needed to determine the VMs' I/O usage, and the methods to establish the trading rate for the underlying CPU and I/O resources. The latter is particularly necessary since the hypervisor's only mechanism to control I/O usage is by making appropriate adjustments in the VM's CPU resources. The experimental evaluation of our solution uses InfiniBand platforms virtualized with the open source Xen hyper visor, and an RDMA-based latency-sensitive benchmark, BenchEx, based on a model of a financial trading platform. The results demonstrate the utility of the ResEx approach in making RDMA-based virtualized platforms more manageable and better suited for hosting even latency-sensitive workloads. ResEx can reduce the latency interference by as much as 30% in some cases as shown.
Adit Ranadive, Ada Gavrilovska, Karsten Schwan
CLUSTER3
2011 Just in time: adding value to the IO pipelines of high performance applications with JITStaging
abstract
Large scale applications are generating a tsunami of data, with understanding driven by finding information hidden within this data. The ever-increasing sizes of output, however, are making it difficult for science users to inspect the data generated by their applications, understand its important properties, and/or organize it for subsequent analysis and visualization. This paper presents JITStager, a software infrastructure with which end users can dynamically customize and thus, add value to the output pipelines of their HEC applications. JITStager is able to customize data at scale, by leveraging the computational power of both compute nodes and of additional `data staging' nodes allocated by end users. Using existing, componentized I/O interfaces to decouple the compile-time specification of the program and the run-time customization of the data pipeline, JITStager employs efficient runtime methods for binary code generation and data movement to create custom pipelines for applications' output processes that provide end users with improved insights into the data being produced, without burdening the application's computational performance and without impeding output performance. This paper describes the JITStager architecture, evaluates its performance, and demonstrates the advantages derived from its use with representative HPC applications.
Hasan Abbasi, Greg Eisenhauer, Matthew Wolf, Karsten Schwan, Scott Klasky
HPDC4
2011 Six degrees of scientific data: reading patterns for extreme scale science IO
abstract
Petascale science simulations generate 10s of TBs of application data per day, much of it devoted to their checkpoint/restart fault tolerance mechanisms. Previous work demonstrated the importance of carefully managing such output to prevent application slowdown due to IO blocking, resource contention negatively impacting simulation performance and to fully exploit the IO bandwidth available to the petascale machine. This paper takes a further step in understanding and managing extreme-scale IO. Specifically, its evaluations seek to understand how to efficiently read data for subsequent data analysis, visualization, checkpoint restart after a failure, and other read-intensive operations. In their entirety, these actions support the 'end-to-end' needs of scientists enabling the scientific processes being undertaken. Contributions include the following. First, working with application scientists, we define 'read' benchmarks that capture the common read patterns used by analysis codes. Second, these read patterns are used to evaluate different IO techniques at scale to understand the effects of alternative data sizes and organizations in relation to the performance seen by end users. Third, defining the novel notion of a 'data district' to characterize how data is organized for reads, we experimentally compare the read performance seen with the ADIOS middleware's log-based BP format to that seen by the logically contiguous NetCDF or HDF5 formats commonly used by analysis tools. Measurements assess the performance seen across patterns and with different data sizes, organizations, and read process counts. Outcomes demonstrate that high end-to-end IO performance requires data organizations that offer flexibility in data layout and placement on parallel storage targets, including in ways that can make tradeoffs in the performance of data writes vs. reads.
Jay F. Lofstead, Milo Polte, Garth A. Gibson, Scott Klasky, Karsten Schwan, Ron A. Oldfield, Matthew Wolf, Qing Liu 0002
HPDC5
2011 Cloud4Home - Enhancing Data Services with @Home Clouds
abstract
Mobile devices, net books and laptops, and powerful home PCs are creating ever-growing computational capacity at the periphery of the Internet, and this capacity is supporting an increasingly rich set of services, including media-rich entertainment and social networks, gaming, home security applications, flexible data access and storage, and others. Such 'at the edge' capacity raises the question, however, about how to combine it with the capabilities present in the cloud computing infrastructures residing in data center systems and reachable via the Internet. The Cloud4Home project and approach presented in this paper addresses this topic, by enabling and exploring the aggregate use of @home and @datacenter computational and storage capabilities. Cloud4Home uses virtualization technologies to create content storage, access, and sharing services that are fungible both in terms of where stored objects are located and in terms of where they are manipulated. In this fashion, data services can provide low latency response to @home events as well as high throughput response when the higher and less predictable latencies of data center access can be tolerated. Cloud4Home is implemented with the Xen open source hypervisors for standard x86-based mobile to server platforms, and is evaluated using sample applications based on home security and video conversion services.
Sudarsun Kannan, Ada Gavrilovska, Karsten Schwan
ICDCS3
2011 Symbiotic Scheduling for Shared Caches in Multi-core Systems Using Memory Footprint Signature
abstract
As the trend of more cores sharing common resources on a single die and more systems crammed into enterprise computing space continue, optimizing the economies of scale for a given compute capacity is becoming more critical. One major challenge in performance scalability is the growing L2 cache contention caused by multiple contexts running on a multi-core processor either natively or under a virtual machine environment. Currently, an OS, at best, relies on history based affinity information to dispatch a process or thread onto a particular processor core. Unfortunately, this simple method can easily lead to destructive performance effect due to conflicts in common resources, thereby slowing down all processes. To ameliorate the allocation/management policy of a shared cache on a multi-core, in this paper, we propose Bloom filter signatures, a low-complexity architectural support to allow an OS or a Virtual Machine Monitor to infer cache footprint characteristics and interference of applications, and then perform job scheduling based on symbiosis. Our scheme integrates hardware-level counting Bloom filters in caches to efficiently summarize cache usage behavior on a per-core, per-process or per-VM basis. We then proposed and studied three resource allocation algorithms to determine the optimal process-to-core mapping to minimize interference in the L2. We executed applications using allocation generated by our new process to-core mapping algorithms on an Intel Core 2 Duo machine and showed an averaged 22% (up to 54%) improvement when applications run natively, and an averaged 9.5% improvement (up to 26%)when running inside VMs.
Mrinmoy Ghosh, Ripal Nathuji, Min Lee, Karsten Schwan, Hsien-Hsin S. Lee
ICPP4
2011 Statistical techniques for online anomaly detection in data centers
abstract
Online anomaly detection is an important step in data center management, requiring light-weight techniques that provide sufficient accuracy for subsequent diagnosis and management actions. This paper presents statistical techniques based on the Tukey and Relative Entropy statistics, and applies them to data collected from a production environment and to data captured from a testbed for multi-tier web applications running on server class machines. The proposed techniques are lightweight and improve over standard Gaussian assumptions in terms of performance.
Chengwei Wang, Krishnamurthy Viswanathan, Choudur Lakshminarayan, Vanish Talwar, Wade Satterfield, Karsten Schwan
Integrated Network Management6
2011 Pegasus: Coordinated Scheduling for Virtualized Accelerator-based Systems
Vishakha Gupta, Karsten Schwan, Niraj Tolia, Vanish Talwar, Parthasarathy Ranganathan
USENIX ATC2
2010 FaReS: Fair Resource Scheduling for VMM-Bypass InfiniBand Devices
abstract
In order to address the high performance I/O needs of HPC and enterprise applications, modern interconnection fabrics, such as InfiniBand and more recently, 10GigE, rely on network adapters with RDMA capabilities. In virtualized environments, these types of adapters are configured in a manner that bypasses the hypervisor and allows virtual machines (VMs) direct device access, so that they deliver near-native low-latency/high-bandwidth I/O. One challenge with the bypass approach is that it causes the hypervisor to lose control over VM-device interactions, including the ability to monitor such interactions and to ensure fair resource usage by VMs. Fairness violations, however, permit low-priority VMs to affect the I/O allocations of other higher priority VMs and more geerally, lack of supervision can lead to inefficiencies in the usage of platform resources. This paper describes the FaReS system-level mechanisms for monitoring VMs' usage of bypass I/O devices. Monitoring information acquired with FaReS is then used to adjust VMM-level scheduling in order to improve resource utilization and/or ensure fairness properties across the sets of VMs sharing platform resources. FaReS employs a memory introspection-based tool for asynchronously monitoring VMM-bypass devices, using InfiniBand HCAs as a concrete example. FaReS and its very low overhead (<;1%) monitoring methods are evaluated with microbenchmarks and representative HPC codes running on modern multicore platforms connected via InfiniBand, using the Xen hypervisor. For these codes fairness is achieved within 2% of the required value.
Adit Ranadive, Ada Gavrilovska, Karsten Schwan
CCGRID3
2010 Robust and flexible power-proportional storage
abstract
Power-proportional cluster-based storage is an important component of an overall cloud computing infrastructure. With it, substantial subsets of nodes in the storage cluster can be turned off to save power during periods of low utilization. Rabbit is a distributed file system that arranges its data-layout to provide ideal power-proportionality down to very low minimum number of powered-up nodes (enough to store a primary replica of available datasets). Rabbit addresses the node failure rates of large-scale clusters with data layouts that minimize the number of nodes that must be powered-up if a primary fails. Rabbit also allows different datasets to use different subsets of nodes as a building block for interference avoidance when the infrastructure is shared by multiple tenants. Experiments with a Rabbit prototype demonstrate its power-proportionality, and simulation experiments demonstrate its properties at scale.
Hrishikesh Amur, James Cipar, Varun Gupta 0004, Gregory R. Ganger, Michael A. Kozuch, Karsten Schwan
SoCC6
2010 Differential virtual time (DVT): rethinking I/O service differentiation for virtual machines
abstract
This paper investigates what it entails to provide I/O service differentiation and performance isolation for virtual machines on individual multicore nodes in cloud platforms. Sharing I/O between VMs is fundamentally different from sharing I/O between processes because guest VM operating systems use adaptive resource management mechanisms like TCP congestion avoidance, disk I/O schedulers, etc. The problem is that these mechanisms are generally sensitive to the magnitude and rate of change of service latencies, where failing to address these latency concerns while designing a service differentiation framework for I/O results in undue performance degradation and hence, insufficient isolation between VMs. This problem is addressed by the notion of Differential Virtual Time (DVT), which can provide service differentiation with performance isolation for VM guest OS resource management mechanisms. DVT is realized within a proportional share I/O scheduling framework for the Xen hypervisor, and its use requires no changes to guest OSs. DVT is applied to message-based I/O, but is also applicable to subsystems like disk I/O. Experimental results with DVT-based I/O scheduling for representative applications demonstrate the utility and effectiveness of the approach.
Mukil Kesavan, Ada Gavrilovska, Karsten Schwan
SoCC3
2010 vNUMA-mgr: Managing VM memory on NUMA platforms
abstract
Continuing improvements in the scale of many-core platforms are accompanied by increased asymmetry in their memory architectures. Such NUMA architectures, however, require systems software that understands this asymmetry to attain high levels of performance, leading to significant work in optimizing operating systems like Linux and Windows to increase locality of access to memory nodes and to consider differences in access latencies when accessing remote nodes. When running on today's virtualized hardware platforms, however, virtual machines (VMs) remain unaware of underlying memory architecture, leading to non-compliance with their operating system's (OS) policies for managing NUMA memory, and thereby, incurring undesirable performance overheads. This paper describes the vNUMA-mgr approach of dealing with the NUMA memory characteristics of future virtualized multicore architectures, which (1) properly manages VM memory at the hyper-visor (VMM) level through a set of VMM-level strategies, coupled with (2) 'enlightening' guest VM OSes, if possible, to aid their memory management, and also (3) provides mechanisms to maintain the distribution of VM memory (across physical nodes), even when memory is over-provisioned. The vNUMA-mgr implementation in the Xen VMM is evaluated with respect to each of its memory allocation strategies. The evaluation, which uses a set of memory-intensive benchmarks representative of many High-Performance Computing (HPC) applications, shows 30-50% performance improvement on our platform.
Subramanya Dulloor, Karsten Schwan
HiPC2
2010 PreDatA - preparatory data analytics on peta-scale machines
abstract
Peta-scale scientific applications running on High End Computing (HEC) platforms can generate large volumes of data. For high performance storage and in order to be useful to science end users, such data must be organized in its layout, indexed, sorted, and otherwise manipulated for subsequent data presentation, visualization, and detailed analysis. In addition, scientists desire to gain insights into selected data characteristics `hidden' or `latent' in these massive datasets while data is being produced by simulations. PreDatA, short for Preparatory Data Analytics, is an approach to preparing and characterizing data while it is being produced by the large scale simulations running on peta-scale machines. By dedicating additional compute nodes on the machine as `staging' nodes and by staging simulations' output data through these nodes, PreDatA can exploit their computational power to perform select data manipulations with lower latency than attainable by first moving data into file systems and storage. Such intransit manipulations are supported by the PreDatA middleware through asynchronous data movement to reduce write latency, application-specific operations on streaming data that are able to discover latent data characteristics, and appropriate data reorganization and metadata annotation to speed up subsequent data access. PreDatA enhances the scalability and flexibility of the current I/O stack on HEC platforms and is useful for data pre-processing, runtime data analysis and inspection, as well as for data exchange between concurrently running simulations.
Fang Zheng 0003, Hasan Abbasi, Ciprian Docan, Jay F. Lofstead, Qing Liu 0002, Scott Klasky, Manish Parashar, Norbert Podhorszki, Karsten Schwan, Matthew Wolf
IPDPS9
2010 Online detection of utility cloud anomalies using metric distributions
abstract
The online detection of anomalies is a vital element of operations in data centers and in utility clouds like Amazon EC2. Given ever-increasing data center sizes coupled with the complexities of systems software, applications, and workload patterns, such anomaly detection must operate automatically, at runtime, and without the need for prior knowledge about normal or anomalous behaviors. Further, detection should function for different levels of abstraction like hardware and software, and for the multiple metrics used in cloud computing systems. This paper proposes EbAT - Entropy-based Anomaly Testing - offering novel methods that detect anomalies by analyzing for arbitrary metrics their distributions rather than individual metric thresholds. Entropy is used as a measurement that captures the degree of dispersal or concentration of such distributions, aggregating raw metric data across the cloud stack to form entropy time series. For scalability, such time series can then be combined hierarchically and across multiple cloud subsystems. Experimental results on utility cloud scenarios demonstrate the viability of the approach. EbAT outperforms threshold-based methods with on average 57.4% improvement in accuracy of anomaly detection and also does better by 59.3% on average in false alarm rate with a `near-optimum' threshold-based method.
Chengwei Wang, Vanish Talwar, Karsten Schwan, Parthasarathy Ranganathan
NOMS3
2010 EFFIS: An End-to-end Framework for Fusion Integrated Simulation
abstract
The purpose of the Fusion Simulation Project is to develop a predictive capability for integrated modeling of magnetically confined burning plasmas. In support of this mission, the Center for Plasma Edge Simulation has developed an End-to-end Framework for Fusion Integrated Simulation (EFFIS) that combines critical computer science technologies in an effective manner to support leadership class computing and the coupling of complex plasma physics models. We describe here the main components of EFFIS and how they are being utilized to address our goal of integrated predictive plasma edge simulation.
Julian C. Cummings, Jay F. Lofstead, Karsten Schwan, Alex Sim, Arie Shoshani, Ciprian Docan, Manish Parashar, Scott Klasky, Norbert Podhorszki, Roselyne Tchoua
PDP3
2010 Managing Variability in the IO Performance of Petascale Storage Systems
abstract
Significant challenges exist for achieving peak or even consistent levels of performance when using IO systems at scale. They stem from sharing IO system resources across the processes of single largescale applications and/or multiple simultaneous programs causing internal and external interference, which in turn, causes substantial reductions in IO performance. This paper presents interference effects measurements for two different file systems at multiple supercomputing sites. These measurements motivate developing a 'managed' IO approach using adaptive algorithms varying the IO system workload based on current levels and use areas. An implementation of these methods deployed for the shared, general scratch storage system on Oak Ridge National Laboratory machines achieves higher overall performance and less variability in both a typical usage environment and with artificially introduced levels of 'noise'. The latter serving to clearly delineate and illustrate potential problems arising from shared system usage and the advantages derived from actively managing it.
Jay F. Lofstead, Fang Zheng 0003, Qing Liu 0002, Scott Klasky, Ron A. Oldfield, Todd Kordenbrock, Karsten Schwan, Matthew Wolf
SC7
2009 Extending I/O through high performance data services
abstract
The complexity of HPC systems has increased the burden on the developer as applications scale to hundreds of thousands of processing cores. Moreover, additional efforts are required to achieve acceptable I/O performance, where it is important how I/O is performed, which resources are used, and where I/O functionality is deployed. Specifically, by scheduling I/O data movement and by effectively placing operators affecting data volumes or information about the data, tremendous gains can be achieved both in the performance of simulation output and in the usability of output data. Previous studies have shown the value of using asynchronous I/O, of employing a staging area, and of performing select operations on data before it is written to disk. Leveraging such insights, this paper develops and experiments with higher level I/O abstractions, termed ldquodata servicesrdquo, that manage output data from `source to sink': where/when it is captured, transported towards storage, and filtered or manipulated by service functions to improve its information content. Useful services include data reduction, data indexing, and those that manage how I/O is performed, i.e., the control aspects of data movement. Our data services implementation distinguishes control aspects - the control plane - from data movement - the data plane, so that both may be changed separably. This results in runtime flexibility not only in which services to employ, but also in where to deploy them and how they use I/O resources. The outcome is consistently high levels of I/O performance at large scale, without requiring application change.
Hasan Abbasi, Jay F. Lofstead, Fang Zheng 0003, Karsten Schwan, Matthew Wolf, Scott Klasky
CLUSTER4
2009 DataStager: scalable data staging services for petascale applications
abstract
Known challenges for petascale machines are that (1) the costs of I/O for high performance applications can be substantial, especially for output tasks like checkpointing, and (2) noise from I/O actions can inject undesirable delays into the runtimes of such codes on individual compute nodes. This paper introduces the flexible 'DataStager' framework for data staging and alternative services within that jointly address (1) and (2). Data staging services moving output data from compute nodes to staging or I/O nodes prior to storage are used to reduce I/O overheads on applications' total processing times, and explicit management of data staging offers reduced perturbation when extracting output data from a petascale machine's compute partition. Experimental evaluations of DataStager on the Cray XT machine at Oak Ridge National Laboratory establish both the necessity of intelligent data staging and the high performance of our approach, using the GTC fusion modeling code and benchmarks running on 1000+ processors.
Hasan Abbasi, Matthew Wolf, Greg Eisenhauer, Scott Klasky, Karsten Schwan, Fang Zheng 0003
HPDC5
2009 Adaptable, metadata rich IO methods for portable high performance IO
abstract
Since IO performance on HPC machines strongly depends on machine characteristics and configuration, it is important to carefully tune IO libraries and make good use of appropriate library APIs. For instance, on current petascale machines, independent IO tends to outperform collective IO, in part due to bottlenecks at the metadata server. The problem is exacerbated by scaling issues, since each IO library scales differently on each machine, and typically, operates efficiently to different levels of scaling on different machines. With scientific codes being run on a variety of HPC resources, efficient code execution requires us to address three important issues: (1) end users should be able to select the most efficient IO methods for their codes, with minimal effort in terms of code updates or alterations; (2) such performance-driven choices should not prevent data from being stored in the desired file formats, since those are crucial for later data analysis; and (3) it is important to have efficient ways of identifying and selecting certain data for analysis, to help end users cope with the flood of data produced by high end codes. This paper employs ADIOS, the adaptable IO system, as an IO API to address (1)-(3) above. Concerning (1), ADIOS makes it possible to independently select the IO methods being used by each grouping of data in an application, so that end users can use those IO methods that exhibit best performance based on both IO patterns and the underlying hardware. In this paper, we also use this facility of ADIOS to experimentally evaluate on petascale machines alternative methods for high performance IO. Specific examples studied include methods that use strong file consistency vs. delayed parallel data consistency, as that provided by MPI-IO or POSIX IO. Concerning (2), to avoid linking IO methods to specific file formats and attain high IO performance, ADIOS introduces an efficient intermediate file format, termed BP, which can be converted, at small cost, to the standard file formats used by analysis tools, such as NetCDF and HDF-5. Concerning (3), associated with BP are efficient methods for data characterization, which compute attributes that can be used to identify data sets without having to inspect or analyze the entire data contents of large files.
Jay F. Lofstead, Fang Zheng 0003, Scott Klasky, Karsten Schwan
IPDPS4
2009 Isolation points: Creating performance-robust enterprise systems
abstract
This article explores a performance isolation-based approach to creating robust distributed applications. For each application, the approach is to understand the performance dependencies that pervade it and then impose constraints on the possible ‘spread’ of such dependencies through the application. The mechanisms used for this purpose, termed isolation points, are software abstractions inserted at key program locations: (1) in application interfaces, (2) in middleware implementations for making remote requests, and (3) in the system interfaces used by middleware and applications. This article demonstrates the utility of isolation points by using them to implement higher level abstractions that improve the performance-robustness of representative enterprise applications. The I-Queue abstraction uses isolation points to implement performance-robust messaging, targeting the message queues used in distributed enterprise codes. By appropriately orchestrating message dispatching, I-Queue can achieve an improvement of 16--32% in dispatched message locality based on traces obtained from the large-scale e-Pricing® search engine operated by Worldspan L.P.
Mohamed S. Mansour, Karsten Schwan, Sameh Abdelaziz
ACM Trans. Auton. Adapt. Syst.2
2008 Active CoordinaTion (ACT) - toward effectively managing virtualized multicore clouds
abstract
A key benefit of utility data centers and cloud computing infrastructure is the level of consolidation they can offer to arbitrary guest applications, and the substantial saving in operational costs and resources that can be derived in the process. However, significant challenges remain before it becomes possible to effectively and at low cost manage virtualized systems, particularly in the face of increasing complexity of individual many-core platforms, and given the dynamic behaviors and resource requirements exhibited by cloud guest VMs. This paper describes the active coordination (ACT) approach, aimed to address a specific issue in the management domain, which is the fact that management actions must (1) typically touch upon multiple resources in order to be effective, and (2) must be continuously refined in order to deal with the dynamism in the platform resource loads. ACT relies on the notion of class-of-service, associated with (sets of) guest VMs, based on which it maps VMs onto platform units, the latter encapsulating sets of platform resources of different types. Using these abstractions, ACT can perform active management in multiple ways, including a VM-specific approach and a black box approach that relies on continuous monitoring of the guest VMs' runtime behavior and on an adaptive resource allocation algorithm, termed Multiplicative Increase, Subtractive Decrease Algorithm with Wiggle Room. In addition, ACT permits explicit external events to trigger VM or application-specific resource allocations, e.g., leveraging emerging standards such as WSDM. The experimental analysis of the ACT prototype, built for Xen-based platforms, use industry-standard benchmarks, including RUBiS, Hadoop, and SPEC. They demonstrate ACT's ability to efficiently manage the aggregate platform resources according to the guest VMs' relative importance (class-of-service), for both the black-box and the VM-specific approach.
Mukil Kesavan, Adit Ranadive, Ada Gavrilovska, Karsten Schwan
CLUSTER4
2008 Exploiting Latent I/O Asynchrony in Petascale Science Applications
abstract
Current and emerging large-scale HPC applications face daunting I/O challenges. In existing codes, problems arise both from large data volumes and from the need to perform complex online data manipulations, including data staging, reorganization, and transformation. We describe three related techniques for enabling, encouraging, and exploiting latent I/O asynchrony in HPC applications: data taps, IOgraphs, and Metabots.
Mary Payne, Patrick M. Widener, Matthew Wolf, Hasan Abbasi, Scott McManus, Patrick G. Bridges, Karsten Schwan
eScience7
2008 Protectit: trusted distributed services operating on sensitive data
abstract
Protecting shared sensitive information is a key requirement for today's distributed applications. Our research uses virtualization technologies to create and maintain trusted data paths across distributed machines, for the services being run and their information exchanges. For trusted data paths, runtime protection methods control what data is visible to which distributed services operating on it, guided by online monitoring that determines the levels of trust inherent in the paths' machines, services, and service actions. This paper presents a key functional element of trusted data paths, which is the ProtectIT interception mechanism for controlling the data exchanges between the different virtual machines running trusted services. ProtectIT can be applied to any communication and/or I/O performed by virtual machines, and because ProtectIT does not require application, middleware, or operating system modifications, it can be used to construct trusted data paths without the knowledge or consent of such entities. Further, since ProtectIT operates in virtual machines isolated from those used by applications, it is not subject to the attacks faced by services exposed to the open Internet. ProtectIT's functionality consists of dynamic protection rules represented as data filters applied to virtual machines' communications. Examples presented in this paper include email services for which ProtectIT's filters control data visibility to mail servers and clients, and unsecured virtual machine communications morphed into secure ones via ProtectIT-based message interception.
Jiantao Kong, Karsten Schwan, Min Lee, Mustaque Ahamad
EuroSys2
2008 ShareStreams-V: A Virtualized QoS Packet Scheduling Accelerator
abstract
This paper introduces a virtualized FPGA-based accelerator for wire speed scheduling of packet streams under quality of service constraints. This work implements the dynamic window constrained scheduling algorithm and builds upon our previous custom accelerator by adding support for virtualization. This implementation is parametric, permitting tradeoffs between packet decision latency, decision throughput, and the number of virtual packet schedulers supported. When scheduling streams from multiple processes, ShareStreams-V 1 is able to schedule minimal size packets faster than one decision per 51.2 ns for up to 64 streams, the throughput required for 10 gbps Ethernet. The bottleneck currently is the host-accelerator HW/SW (PCIe) interface; this may be mitigated using high-speed interconnects/interfaces such as HyperTransport.
Kangtao Kendall Chuang, Sudhakar Yalamanchili, Ada Gavrilovska, Karsten Schwan
FCCM4
2008 Vpm tokens: virtual machine-aware power budgeting in datacenters
abstract
Power consumption and cooling overheads are becoming increasingly significant for large scale machines, affecting overall costs and the ability to extend resource capacities and performance capabilities. To help mitigate these issues, active power management technologies are being deployed aggressively, including power budgeting, which enables improved power provisioning and can address critical periods when power delivery or cooling capabilities are temporarily reduced. Given the use of virtualization to encapsulate application components into virtual machines (VMs), however, such power management capabilities must address the interplay between budgeting physical resources and the performance of the virtual machines used to run these applications. This paper proposes a set of cluster- and datacenter-level management components and abstractions for use by power budgeting policies. The key idea is to manage power from a VM-centric point of view, where the goal is to be aware of global utility tradeoffs between different virtual machines (and their applications) when maintaining power constraints for the physical hardware on which they run. Our approach to VM-aware power budgeting uses multiple distributed managers integrated into the VirtualPower Management (VPM) framework whose actions are coordinated via a new abstraction, termed VPM tokens. An implementation with the Xen hypervisor illustrates technical benefits of VPM tokens that include up to 43% improvements in global utility, highlighting the ability to dynamically improve cluster performance while still meeting power budgets.
Ripal Nathuji, Karsten Schwan
HPDC2
2008 Flexible Classification on Heterogenous Multicore Appliance Platforms
abstract
Emerging heterogeneous multicore systems are suitable platforms for efficient deployment of application- specific service components. Easily virtualizable and reprogrammable platforms, these 'appliances' for future information services make it possible to run familiar software stacks created with common development tools, in addition to offering acceleration capabilities for performance critical software components. Of particular importance are the efficient and scalable execution of the content-based services prevalent in future Internet applications, including flexible methods for data classification and forwarding. This paper provides a brief description of representative platform hardware and software components. It then describes in more detail the feasibility of supporting a range of flexible and reconfigurable application-specific classification operations necessary for future Internet applications.
Priyanka Tembey, Anish Bhatt, Subramanya Dulloor, Ada Gavrilovska, Karsten Schwan
ICCCN5
2008 WorldTravel: A Testbed for Service-Oriented Applications
Peter Budny, Srihari Govindharaj, Karsten Schwan
ICSOC3
2008 A state-space approach to SLA based management
abstract
Large complex systems (such as Enterprise systems) are often composed of several interacting, independent components. In many such systems, although the behavior of the constituent components is well characterized, the behavior that results from interaction between such components is more or less intractable; making it hard for the administrators to efficiently manage the system in conformance with the service level agreements or the SLAs. This paper presents an approach for deriving component-level objectives from system-level objectives or agreements, which if conformed to, imply conformance to the higher-level SLA. Our approach partitions the systempsilas state-space into homogeneous sub-spaces, creates micro-models for such subspaces, and then uses such micro-models to translate the higher-level objectives to component-level objectives. We have implemented a system, termed Pranaali, for evaluating our approach in realistic settings.
Vibhore Kumar, Karsten Schwan, Subu Iyer, Yuan Chen 0001, Akhil Sahai
NOMS2
2008 Netchannel: a VMM-level mechanism for continuous, transparentdevice access during VM migration
abstract
Efficient and seamless access to I/O devices during VM migration coupled with the ability to dynamically change the mappings of virtual to physical devices are required in multiple settings, including blade-servers, datacenters, and even in home-based personal computing environments. This paper develops a general solution for these problems, at a level of abstraction transparent to guest VMs and their device drivers. A key part of this solution is a novel VMM-level abstraction that transparently handles pending I/O transactions, termed Netchannel. Netchannel provides for (1) virtual device migration and device hot-swapping for networked as well as locally attached devices, and (2) remote access to devices not directly attached to networks via transparent device remoting, an example being a disk locally present on a bladeserver node. A Xen-based prototype of Netchannel demonstrates these capabilities for block and for USB devices, for both bulk and isochronous USB access methods. Within the same administrative domain, seamless access to these devices is maintained during live VM migration and during device hot-swapping. Experimental evaluations with micro-benchmarks and with representative server applications exhibit comparable performance for Netchannel-based remote vs. local devices.
Karsten Schwan
VEE2
2007 LIVE data workspace: A flexible, dynamic and extensible platform for petascale applications
abstract
The data needs of current and future PetaScale applications have increased over the last half decade to the extent that appropriate data management has become a crucial requirement. This concerns not only the storage of data produced by the new class of PetaScale applications, but also the data exchanges needed for coupling applications with concurrent analysis, online data visualization for validation, and others. To address such dynamic code coupling, we introduce the concept of an extensible, dynamic, and flexible data workspace, termed LIVE. In contrast to the data exchanges programmed with MPI, MPI-IO, or grid software, LIVE focuses on data exchanges carried out without a priori knowledge of potential data requirements. Examples include exchanges required by ad-hoc or dynamically determined methods for data validation, for general data analysis tasks, or for data visualization. Run on an execution environment comprised of integrated dynamic discovery and on-line management services, LIVE is used to create a dasiadata workspacepsila for a working molecular dynamics code base utilized by mechanical and materials engineers at Georgia Tech, for multi-scale materials modeling. Measurements of both this applicationpsilas data workspace and of the basic primitives in the LIVE framework demonstrate that the environmentpsilas substantial flexibility has minimal impact on overall performance, and in fact, that it improves performance in a number of usage scenarios. In particular, for a visualization pipeline example derived from our collaborators, we see a slight improvement over a solution based on MPI-IO, and a further improvement of up to 5% by utilizing LIVEpsilas ability to overlap communication with user-specified computation.
Hasan Abbasi, Matthew Wolf, Karsten Schwan
CLUSTER3
2007 E2EProf: Automated End-to-End Performance Management for Enterprise Systems
abstract
Distributed systems are becoming increasingly complex, caused by the prevalent use of Web services, multi-tier architectures, and grid computing, where dynamic sets of components interact with each other across distributed and heterogeneous computing infrastructures. For these applications to be able to predictably and efficiently deliver services to end users, it is therefore, critical to understand and control their runtime behavior. In a datacenter environment, for instance, understanding the end-to-end dynamic behavior of certain IT subsystems, from the time requests are made to when responses are generated and finally, received, is a key prerequisite for improving application response, to provide required levels of performance, or to meet service level agreements (SLAs). The E2EProf toolkit enables the efficient and nonintrusive capture and analysis of end-to-end program behavior for complex enterprise applications. E2EProf permits an enterprise to recognize and analyze performance problems when they occur - online, to take corrective actions as soon as possible and wherever necessary along the paths currently taken by user requests - end-to-end, and to do so without the need to instrument applications - nonintrusively. Online analysis exploits a novel signal analysis algorithm, termed pathmap, which dynamically detects the causal paths taken by client requests through application and backend servers and annotates these paths with end-to-end latencies and with the contributions to these latencies from different path components. Thus, with pathmap, it is possible to dynamically identify the bottlenecks present in selected servers or services and to detect the abnormal or unusual performance behaviors indicative of potential problems or overloads. Pathmap and the E2EProf toolkit successfully detect causal request paths and associated performance bottlenecks in the RUBiS ebay-like multi-tier Web application and in one of the datacenter of our industry partner, Delta Air Lines.
Sandip Agarwala, Fernando Alegre, Karsten Schwan, Jegannathan Mehalingham
DSN3
2007 LIVE: : a light-weight data workspace for computational science
abstract
We present the Lightweight Information Validation Environment, LIVE as asolution to the high complexity and data sizes of modern day computational science applications. LIVE is a data workspace that facilitates the creation of dynamic data processing overlays we call I/O graphs. We use LIVE as aplatform for dynamic extension of scientific applications using lightweight data extraction, runtime discovery and flexible data selection.
Hasan Abbasi, Matthew Wolf, Karsten Schwan
HPDC3
2007 High performance and scalable I/O virtualization via self-virtualized devices
abstract
While industry is making rapid advances in system virtualization, for server consolidation and for improving system maintenance and management, it has not yet become clear how virtualization can contribute to the performance of high end systems. In this context, this paper addresses a key issue in system virtualization - how to efficiently virtualize I/O subsystems and peripheral devices. We have developed a novel approach to I/O virtualization, termed self-virtualized devices, which improves I/O performance by off loading select virtualization functionality onto the device. This permits guest virtual machines to more efficiently (i.e., with less overhead and reduced latency) interact with the virtualized device. The concrete instance of such a device developed and evaluated in this paper is a self-virtualized network interface (SV-NIC), targeting the high end NICs used in thehigh performance domain. The SV-NIC (1) provides virtual interfaces (VIFs) to guest virtual machines for an underlying physical device, the network interface, (2) manages the wayin which the device's physical resources are used by guest operating systems, and (3) provides high performance, low overhead network access to guest domains. Experimental results are attained in a prototyping environment using an IXP 2400-based ethernet board as a programmable network device. The SV-NIC scales to large numbers of VIFs and guests, and offers VIFs with 77% higher throughput and 53% less latency compared to the current standard virtualized device implementations on hyper visor-based platforms.
Himanshu Raj, Karsten Schwan
HPDC2
2007 Automated Availability Management Driven by Business Policies
abstract
Policy-driven service management helps reduce IT management cost and it keeps the service management aligned with business objectives. While most of the previous research focuses on performance/resource managements, little has been researched in the area of availability management driven by business policies. This is a critically important task in enterprise IT, because a single failure in enterprise IT could cause huge business loss. It is still unclear how we can automate the availability management in a highly dynamic and complex system according to business level objectives for performance and risk attitude/preference. As a consequence, users can not manage the availability/performance ratio to match their risk tolerance. In this paper, we propose a policy-driven approach to automate run-time availability management in IT systems, according to high level availability and performance objectives. We further apply von Neumann-Morgenstern utility theory to deal with users' risk attitude policies, an important class of business policies specific for availability management, and the associated preference structures. Based on the proposed approach, we implement an automated decision engine for availability management. The initial evaluation of the solution illustrates the significance of the policy-driven approach and it demonstrates its applicability for availability management in complex IT environments. This way IT users can customize their availability to the risks tolerable by business objectives.
Zhongtang Cai, Yuan Chen 0001, Vibhore Kumar, Dejan S. Milojicic, Karsten Schwan
Integrated Network Management5
2007 iManage: Policy-Driven Self-management for Enterprise-Scale Systems
Vibhore Kumar, Brian F. Cooper, Greg Eisenhauer, Karsten Schwan
Middleware4
2007 Towards IQ-Appliances: Quality-awareness in Information Virtualization
abstract
Our research addresses "information appliances' used in modern large-scale distributed systems to: (1) virtualize their data flows by applying actions such as filtering, format translation, etc., and (2) separate such actions from enterprise applications' business logic, to make it easier for future service-oriented codes to inter-operate in diverse and dynamic environments. Our specific contribution is the enrichment of runtimes of these appliances with methods for QoS-awareness, thereby giving them the ability to deliver desired levels of QoS even under sudden requirement changes - IQ-appliances. For experimental evaluation, we prototype an IQ-appliance. Measurements demonstrate the feasibility and utility of the approach.
Radhika Niranjan, Ada Gavrilovska, Karsten Schwan, Priyanka Tembey
NCA3
2007 VirtualPower: coordinated power management in virtualized enterprise systems
abstract
Power management has become increasingly necessary in large-scale datacenters to address costs and limitations in cooling or power delivery. This paper explores how to integrate power management mechanisms and policies with the virtualization technologies being actively deployed in these environments. The goals of the proposed VirtualPower approach to online power management are (i) to support the isolated and independent operation assumed by guest virtual machines (VMs) running on virtualized platforms and (ii) to make it possible to control and globally coordinate the effects of the diverse power management policies applied by these VMs to virtualized resources. To attain these goals, VirtualPower extends to guest VMs `soft' versions of the hardware power states for which their policies are designed. The resulting technical challenge is to appropriately map VM-level updates made to soft power states to actual changes in the states or in the allocation of underlying virtualized hardware. An implementation of VirtualPower Management (VPM) for the Xen hypervisor addresses this challenge by provision of multiple system-level abstractions including VPM states, channels, mechanisms, and rules. Experimental evaluations on modern multicore platforms highlight resulting improvements in online power management capabilities, including minimization of power consumption with little or no performance penalties and the ability to throttle power consumption while still meeting application requirements. Finally, coordination of online methods for server consolidation with VPM management techniques in heterogeneous server systems is shown to provide up to 34% improvements in power consumption.
Ripal Nathuji, Karsten Schwan
SOSP2
2007 IQ-Paths: Predictably High Performance Data Streams Across Dynamic Network Overlays
Zhongtang Cai, Vibhore Kumar, Karsten Schwan
J. Grid Comput.3
2007 Advanced networking services for distributed multimedia streaming applications
Ada Gavrilovska, Srikanth Sundaragopalan, Karsten Schwan
Multim. Tools Appl.4
2006 IQ-Paths: Predictably High Performance Data Streams across Dynamic Network Overlays
abstract
Overlay networks are a key vehicle for delivering network and processing resources to high performance applications. For shared networks, however, to consistently deliver such resources at desired levels of performance, overlays must be managed at runtime, based on the continuous assessment and prediction of available distributed resources. Data-intensive applications, for example, must assess, predict, and judiciously use available network paths, and dynamically choose alternate or exploit concurrent paths. Otherwise, they cannot sustain the consistent levels of performance required by tasks like remote data visualization, online program steering, and remote access to high end devices. The multiplicity of data streams occurring in complex scientific workflows or in large-scale distributed collaborations exacerbate this problem, particularly when different streams have different performance requirements. This paper presents IQ-Paths, a set of techniques and their middleware realization that implement self-regulating overlay streams for data-intensive distributed applications. Self-regulation is based on (1) the dynamic and continuous assessment of the quality of each overlay path, (2) the use of online network monitoring and statistical analyses that provide probabilistic guarantees about available path bandwidth, loss rate, and RTT, and (3) self-management, via an efficient packet routing and scheduling algorithm that dynamically schedules data packets to different overlay paths in accordance with their available bandwidths. IQ-Paths offer probabilistic guarantees for application-level specifications of stream utility, based on statistical predictions of available network bandwidth. This affords applications with the ability, for instance, to send control or steering data across overlay paths that offer strong guarantees for future bandwidth vs. across less guaranteed paths. Experimental results presented in this paper use IQ-Paths to better handle the different kinds of data produced by two high performance applications: (1) a data-driven or interactive high performance code with user-defined utility requirements and (2) an adaptive overlay version of the popular Grid-FTP application
Zhongtang Cai, Vibhore Kumar, Karsten Schwan
HPDC3
2006 SysProf: Online Distributed Behavior Diagnosis through Fine-grain System Monitoring
abstract
Runtime monitoring is key to the effective management of enterprise and high performance applications. To deal with the complex behaviors of today’s multi-tier applications running across shared platforms, such monitoring must meet three criteria: (1) fine granularity, including being able to track the resource usage of specific application behaviors like individual client-server interactions, (2) real-time response, referring to the monitoring system’s ability to both capture and analyze currently needed monitoring information with the delays required for online management, and (3) enterprise-wide operation, which means that the monitoring information captured and analyzed must span across the entire software stack and set of machines involved in request generation, request forwarding, service provision, and return. This paper presents the SysProf system-level monitoring toolkit, which provides a flexible, low overhead framework for enterprise-wide monitoring. The toolkit permits the capture of monitoring information at different levels of granularity, ranging from tracking the system-level activities triggered by a single system call, to capturing the client-server interactions associated with certain request classes, to characterizing the server resources consumed by sets of clients or client behaviors. The paper demonstrates the efficacy of SysProf by using it to manage two different enterprise applications: (1) detecting performance bottlenecks in a high performance shared network file service, and (2) enforcing service level agreements in a multi-tier auctioning web site.
Sandip Agarwala, Karsten Schwan
ICDCS2
2006 AutoPower: Toward Energy-aware Software Systems for Distributed Mobile Robots
abstract
Autonomous robot systems have to manage their energy wisely in order to complete their missions. Typical approaches seek to conserve energy by energy-efficient motion or sensor planning. This paper puts forth a distributed systems approach to power management. Specifically, it develops and presents AutoPower, which is a model that characterizes robot software systems' computation and communication energy behaviors. With AutoPower, it is possible to make principled decisions about (1) where to deploy software components across the distributed computing resources of autonomous robotic systems, and (2) how the different systems involved should communicate to best meet overall mission objectives. We showcase AutoPower by using a multi-robot search-and-rescue mission as a guiding application. For this scenario, application of the model shows that there are counter-intuitive energy trade-offs in configuring such application software. Further, by using AutoPower to guide deployment and interconnects at runtime, for certain configurations, overall computing system lifetimes can be increased by up to 57% over a base-line configuration
Keith J. O'Hara, Ripal Nathuji, Himanshu Raj, Karsten Schwan, Tucker R. Balch
ICRA4
2006 I-Queue: Smart Queues for Service Management
Mohamed S. Mansour, Karsten Schwan, Sameh Abdelaziz
ICSOC2
2006 Utility-Driven Proactive Management of Availability in Enterprise-Scale Information Flows
Zhongtang Cai, Vibhore Kumar, Brian F. Cooper, Greg Eisenhauer, Karsten Schwan, Robert E. Strom
Middleware5
2006 Utilizing Network Processors in Distributed Enterprise Environments
abstract
The integration of legacy systems and enterprise applications with novel communications, Internet, or networking services creates the problems of mismatches, information integration and possibly, problems from the evolution of the systems themselves. Enterprise services to solve these problems are currently implemented via commodity server hardware or application specific integrated circuits (ASICs), which suffer from either poor performance or a lack of flexibility. This paper presents an alternative to implementing these services by using network processors (NPs) as fast, flexible, and cost-efficient network appliances. It characterizes the hardware and software requirements of NP-based enterprise services, and presents a case study of a real-world enterprise problem. We implement a solution service to the problem on the Intel IXP2400, and present evaluation results that generalize the strengths, weaknesses, and requirements of the NP's role as a platform for enterprise service deployment
Paul Royal, Mitch Halpin, Ada Gavrilovska, Karsten Schwan
NCA4
2006 Protected data paths: delivering sensitive data via untrusted proxies
abstract
No abstract available.
Jiantao Kong, Karsten Schwan, Patrick M. Widener
PST2
2006 IQ-Services: network-aware middleware for interactive large-data applications
abstract
Abstract IQ‐Services are application‐specific, resource‐aware code modules executed by data transport middleware. They constitute a ‘thin’ layer between application components and the underlying computational and communication resources. This layer implements the data manipulations necessary to permit wide‐area collaborations to proceed smoothly in the presence of dynamic resource variations. IQ‐Services interact with the application and resource layers via dynamic performance attributes, and end‐to‐end implementations of such attributes also permit clients to interact with data providers. The joint middleware/resource and provider/consumer interactions implemented with performance attributes may be used to realize effective methods for managing the data flows in the large‐data, distributed Grid applications targeted by our research. Experimental results in this paper demonstrate substantial performance improvements. These are attained by coordinating network‐level with service‐level adaptations of the data being transported and by permitting end users to dynamically deploy and use application‐specific services for manipulating data in ways suitable for their current needs. Copyright © 2005 John Wiley & Sons, Ltd.
Zhongtang Cai, Greg Eisenhauer, Qi He 0001, Vibhore Kumar, Karsten Schwan, Matthew Wolf
Concurr. Comput. Pract. Exp.5
2005 Addressing data compatibility on programmable network platforms
abstract
Large-scale applications require the efficient exchange of data across their distributed components, including data from heterogeneous sources and to widely varying clients. Inherent to such data exchanges are (1) discrepancies among the data representations used by sources, clients, or intermediate application components (e.g., due to natural mismatches or due to dynamic component evolution), and (2) requirements to route, combine, or otherwise manipulate data as it is being transferred. As a result, there is an ever growing need for data conversion services, handled by stubs in application servers, by middleware or messaging services, by the operating system, or by the network. This paper's goal is to demonstrate and evaluate the ability of modern network processors to efficiently address data compatibility issues, when data is 'in transit' between application-level services. Toward this end, we present the design and implementation of a network-level execution environment that permits systems to dynamically deploy and configure application-level data conversion services 'into' the network infrastructure. Experimental results obtained with a prototype implementation on Intel's IXP2400 network processors include measurements of XML-like data format conversions implemented with efficient binary data formats.
Ada Gavrilovska, Karsten Schwan
ANCS2
2005 Harnessing Shared Wide-area Clusters for Dynamic High End Services
abstract
Current trends in distributed computing have been moving towards the use of wide-area clusters that are managed by different entities. In this paper, we introduce middleware-level support to facilitate computational resource sharing with service guarantees using non-dedicated server systems in wide-area clusters. The aim is to ensure that sets of computational tasks submitted to such high end systems are completed reliably and in a timely fashion. Our approach develops methods that enhance basic job scheduling with information about the execution history and trust values for the computational nodes to which jobs are assigned. In essence, job scheduling is enriched with trust models constructed and maintained at runtime, and scheduling decisions are based on metrics that capture trust in remote server systems. An implementation of the approach is evaluated on Planetlab, with initial results demonstrating good success rates in completing jobs within their specific service level agreements, including under conditions of high system loads. Additional results are attained with a variant of the scheduling algorithm that uses redundancy to further improve the likelihood of meeting end user SLAs. A representative application considered in this paper is remote data visualization, where substantial computation must be applied to data before displaying it to end users. SLAs capture desired end-to-end delay, and distributed server or cluster systems are used to perform the required computations in a timely manner
R. Viswanath, Mustaque Ahamad, Karsten Schwan
CLUSTER3
2005 Service Augmentation for High End Interactive Data Services
abstract
Advances in computational science, combined with the increasingly interdisciplinary and geographically distributed research teams, have led to a need to support multi-tiered, data- and meta-data-rich collaboration infrastructures. Our research addresses the interactive, remote tasks undertaken in such collaborations, which require a flexible software infrastructure able to dynamically deploy services where and when needed, and to provide data to clients in the forms in which they require it with suitable levels of end-to-end performance. The concept of service augmentation advanced in this paper seeks to continuously adjust the differences or degrees of incompatibility between the data received and the data displayed or stored by clients. Difference adjustments occur anywhere on the paths between data providers and clients, and compatibility computations leverage all of the resources that may be brought to bear, including CPUs and GPUs on servers and additional data manipulations on server, overlay, and client nodes. A formal structure and experimental evaluations of this concept are performed with the SmartPointer scientific visualization and annotation framework, for which we show that data-driven SLAs provide improved client flexibility and the ability to maintain application-specific notions of quality of service
Matthew Wolf, Hasan Abbasi, Benjamin Collins, David Spain, Karsten Schwan
CLUSTER5
2005 Lightweight Morphing Support for Evolving Middleware Data Exchanges in Distributed Applications
abstract
Most systems must evolve as their missions or roles change and/or as they adapt to new execution environments. When evolving large distributed applications, it is particularly difficult to make changes to the data formats that underlie their components' communications, because such 'format evolution' can affect all or many application components. Prior approaches to the problem of implementing changes in the communications of a deployed system have relied upon ad-hoc solutions or on protocol negotiation to avoid message format mismatches. Unfortunately, such solutions tend to increase the complexity of application code. This paper presents a novel approach to the problem of data format evolution that combines meta-data about the data being exchanged with dynamic binary code generation to create a robust data exchange system that naturally supports application evolution. The idea is to specialize the communications of application components by dynamically generating the code that can automatically transform incoming data into forms that receiving components can understand. A realistic example in the context of publish/subscribe middleware is used to illustrate how this technique can be applied to enhance interoperability between different version of distributed applications
Sandip Agarwala, Greg Eisenhauer, Karsten Schwan
ICDCS3
2005 Resource-Aware Distributed Stream Management Using Dynamic Overlays
abstract
We consider distributed applications that continuously stream data across the network, where data needs to be aggregated and processed to produce a 'useful' stream of updates. Centralized approaches to performing data aggregation suffer from high communication overheads, lack of scalability, and unpredictably high processing workloads at central servers. This paper describes a scalable and efficient solution to distributed stream management based on (1) resource-awareness, which is middleware-level knowledge of underlying network and processing resources, (2) overlay-based in-network data aggregation, and (3) high-level programming constructs to describe data-flow graphs for composing useful streams. Technical contributions include a novel algorithm based on resource-aware network partitioning to support dynamic deployment of dataflow graph components across the network, where efficiency of the deployed overlay is maintained by making use of partition-level resource-awareness. Contributions also include efficient middleware-based support for component deployment, utilizing runtime code generation rather than interpretation techniques, thereby addressing both high performance and resource-constrained applications. Finally, simulation experiments and benchmarks attained with actual operational data corroborate this paper's claims.
Vibhore Kumar, Brian F. Cooper, Zhongtang Cai, Greg Eisenhauer, Karsten Schwan
ICDCS5
2005 Opportunistic Overlays: Efficient Content Delivery in Mobile Ad Hoc Networks
Yuan Chen 0001, Karsten Schwan
Middleware2
2005 I-RMI: Performance Isolation in Information Flow Applications
Mohamed S. Mansour, Karsten Schwan
Middleware2
2005 C-CORE: Using Communication Cores for High Performance Network Services
abstract
Recent hardware advances are creating multi-core systems with heterogeneous functionality. This paper explores how applications and middleware can utilize systems comprised of processors specialized for communication vs. computational tasks. The C-CORE execution environment enables applications, through middleware and underlying system functionality, to utilize both the computational capabilities of general purpose CPUs and the high performance communication hardware provided by specialized communication processors. Such future heterogeneous multi-core hardware is emulated by attaching a representative network processor - Intel's IXP2400 processor - to a general purpose CPU via a dedicated interconnect. For this platform, C-CORE provides abstractions to represent an application's communication actions, to efficiently couple such actions with application-level computations, and to dynamically create and configure the platform-resident 'chains' of computational and communication actions used by applications. C-CORE's functionality is evaluated with representative, communication-intensive applications. Measurements on our experimental platform establish the performance advantages afforded to applications by C-CORE
Ada Gavrilovska, Karsten Schwan, Srikanth Sundaragopalan
NCA3
2005 Platform Overlays: enabling in-network stream processing in large-scale distributed applications
abstract
The purpose of this research is to explore the capabilities of future, multi-core heterogeneous systems, with specialized communication support, to be used as efficient and flexible execution platforms in distributed streaming applications. On such platforms, we create overlays of hardware- and software-supported execution contexts -- platform overlays. Stream manipulations, represented via stream handlers, are deployed on top of such overlays, based on the ability of individual contexts to perform handler operations. As a result, stream processing is dynamically mapped to those platform resources best suited for it, and it can even be fully contained to the networking subsystems, thereby enabling in-network stream processing. Experimental results demonstrate the benefits of our approach towards meeting application-specific quality requirements.
Ada Gavrilovska, Srikanth Sundaragopalan, Karsten Schwan
NOSSDAV4
2005 KStreams: kernel support for efficient data streaming in proxy servers
abstract
Growth in broadband connectivity is making media streaming applications increasingly popular. For scalability, media is streamed across sets of proxy servers embedded in overlay networks, where the quality of delivered content depends both on available network capacities across overlay nodes and the capabilities of proxy servers. This paper addresses proxy server performance for media streaming and for the delivery of live media content. Our approach to efficient content delivery is to develop a set of kernel-level data streaming abstractions, termed KStreams. Compared to user-level solutions, KStreams (1) offers improved performance for the multiple data forwarding models commonly used in data distribution networks, (2) reduces per stream overheads by eliminating unnecessary system calls and memory copying, and (3) offers improved levels of predictability for the Quality of Service (QoS) experienced by media streams due to its use of non-preemptable kernel-level threads and its ability to directly interact with the CPU scheduler and other kernel-level resource managers. (4) Once initiated, KStreams operates without further involvement of and asynchronously to applications, permitting them to carry out other tasks.
Jiantao Kong, Karsten Schwan
NOSSDAV2
2005 Feedback-Based Dynamic Voltage and Frequency Scaling for Memory-Bound Real-Time Applications
abstract
Dynamic voltage and frequency scaling is increasingly being used to reduce the energy requirements of embedded and real-time applications by exploiting idle CPU resources, while still maintaining all application's real-time characteristics. Accurate predictions of task run-times are key to computing the frequencies and voltages that ensure that all tasks' real-time constraints are met. Past work has used feedback-based approaches, where applications' past CPU utilizations are used to predict future CPU requirements. Mispredictions in these approaches can lead to missed deadlines, suboptimal energy savings, or large overheads due to frequent changes to the chosen frequency or voltage. One shortcoming of previous approaches is that they ignore other 'indicators' of future CPU requirements, such as the frequency of I/O operations, memory accesses, or interrupts. This paper addresses the energy consumptions of memory-bound real-time applications via a feedback loop approach, based on measured task run-times and cache miss rates. Using cache miss rates as indicator for memory access rates introduces a more reliable predictor of future task run-times. Even in modern processor architectures, memory latencies can only be hidden partially, therefore, cache misses can be used to improve the run-time predictions by considering potential memory latencies. The results shown in this paper indicate improvements in both the number of deadlines met and the amount of energy saved.
Christian Poellabauer, Leo Singleton, Karsten Schwan
IEEE Real-Time and Embedded Technology and Applications Symposium3
2005 Combining Compiler and Operating System Support for Energy Efficient I/O on Embedded Platforms
abstract
Mobile and embedded platforms have experienced dramatic advances in capabilities, largely due to the development of associated peripheral devices for storage and communication. The incorporation of these I/O devices has increased the overall power envelope of these platforms. In fact, system-level power consumption of mobile platforms is often dominated by peripheral devices. Since battery technologies alone have been unable to provide the lifetimes required by many platforms, in order to conserve energy, most devices provide the ability to transition into low power states during idle periods. The resulting energy savings are heavily dependent upon the lengths and number of idle periods experienced by a device. This paper presents an infrastructure designed to take advantage of device low power states by increasing the burstiness of device accesses and idle periods to provide a reduced power profile, and thereby an improvement in battery life. Our approach combines compiler-based source modifications with operating system support to implement a dynamic solution for enhanced energy consumption. We evaluate our infrastructure on an XScale-based embedded platform with a Linux implementation.
Ripal Nathuji, Balasubramanian Seshasayee, Karsten Schwan
SCOPES3
2005 Flexible cross-domain event delivery for quality-managed multimedia applications
abstract
To meet end users' quality-of-service (QoS) requirements, online quality management for multimedia applications must include appropriate allocation of the underlying computing platform's resources. Previous work has developed novel operating system (OS) functionality for dynamic QoS management, including multimedia or real-time CPU schedulers and OS extensions for online performance monitoring and for adaptations, as well as QoS-aware applications that adapt their behavior to gain additional benefits from such functionality. This article describes a general OS mechanism that may be used to implement a wide variety of online quality management functions. ECalls is a communication mechanism that implements multiple cross-domain calling conventions that can be customized to the quality management needs of applications. The ECalls mechanism is based on the notions of events, event channels, and event handlers. Using events, applications can share relevant QoS attributes with OS services, and OS-level resource management services can efficiently provide monitoring data to target applications or application managers. Dynamically generated event handlers can be used to customize event delivery to meet diverse application needs, for example, to achieve high scalability for Web servers or small jitter for real-time data delivery.
Christian Poellabauer, Karsten Schwan
ACM Trans. Multim. Comput. Commun. Appl.2
2004 XChange: coupling parallel applications in a dynamic environment
abstract
Modern computational science applications are becoming increasingly multidisciplinary, involving widely distributed research teams and their underlying computational platforms. A common problem for the grid applications used in these environments is the necessity to couple multiple, parallel subsystems, with examples ranging from data exchanges between cooperating, linked parallel programs, to concurrent data streaming to distributed storage engines. This work presents the XChange/sub mxn/ middleware infrastructure for coupling componentized distributed applications. XChange/sub mxn/ implements the basic functionality of well-known services like the CCA Forum's MxN project, by providing efficient data redistribution across parallel application components. Beyond such basic functionality, however, XChange/sub mxn/ also addresses two of the problems faced by wide area scientific collaborations, which are (1) the need to deal with dynamic application/component behaviors, such as dynamic arrivals and departures due to the availability of additional resources, and (2) the need to 'match' data formats across disparate application components and research teams. In response to these needs, XChange/sub mxn/ uses an anonymous publish/subscribe model for linking interacting components, and the data being exchanged is dynamically specialized and transformed to match end point requirements. The pub/sub paradigm makes it easy to deal with dynamic component arrivals and departures. Dynamic data transformation enables the 'inflight' correction of data or needs mismatches for cooperating components. This work describes the design and implementation of XChange/sub mxn/, and it evaluates its implementation compared to those of less flexible transports like MPI. It also highlights the utility ofXChange/sub mxn/'s 'inflight' data specialization, by applying it to the SmartPointer parallel data visualization environment developed at our institution. Interestingly, using XChange/sub mxn/ did not significantly affect performance but led to a reduction in the size of the code base.
Hasan Abbasi, Matthew Wolf, Karsten Schwan, Greg Eisenhauer, A. Hilton
CLUSTER3
2004 ShareStreams: A Scalable Architecture and Hardware Support for High-Speed QoS Packet Schedulers
abstract
ShareStreams (scalable hardware architectures for stream schedulers) is a unified hardware architecture for realizing a range of wire-speed packet scheduling disciplines for output link scheduling. This paper presents opportunities to exploit parallelism, design issues, tradeoffs and evaluation of the FPGA hardware architecture for use in switch network interfaces. The architecture uses processor resources for queuing and data movement and FPGA hardware resources for accelerating decisions and priority updates. The hardware architecture stores state in register base blocks, stream service attributes are compared using single-cycle decision blocks arranged in a novel single-stage recirculating network. The architecture provides effective mechanisms to trade hardware complexity for lower execution-time in a predictable manner. The hardware realized in a Virtex-I and Virtex-II FPGA can meet the packet-time requirements of 10 Gbps links for 256 stream queues with window-constrained scheduling disciplines. The hardware can schedule 1536 stream queues with priority-class/fair-queuing scheduling disciplines using 16 service-classes to meet 10 Gbps packet-times.
Raj Krishnamurthy, Sudhakar Yalamanchili, Karsten Schwan, Richard West
FCCM3
2004 SOAP-binQ: High-Performance SOAP with Continuous Quality Management
abstract
There is substantial interest in using SOAP (simple object access protocol) in distributed applications' interprocess communications due to its promise of universal interoperability. The utility of SOAP is limited, however, by its inefficient implementation. Our research aims to make SOAP useful for high end or resource-constrained applications. The resulting SOAP-bin communication protocol exhibits substantially improved performance compared to regular SOAP communications. Gains are particularly evident when the same types of parameters are exchanged repeatedly, examples including transactional applications, remote graphics or visualization, and distributed scientific codes. A further improvement to SOAP-bin, termed SOAP-binQ, addresses resource-constrained applications like distributed media codes, where scarce communication bandwidth, for example, may prevent end users from interacting in real-time. SOAP-binQ offers quality management functions that permit SOAP to reduce parameter sizes dynamically, as and when needed. The methods used in size reduction are provided by end users and/or by applications, thereby enabling domain-specific tradeoffs in quality vs. performance. An adaptive use of SOAP-binQ's quality management techniques presented significantly reduces the jitter experienced in two sample applications like remote sensing and remote visualization.
Balasubramanian Seshasayee, Karsten Schwan, Patrick M. Widener
ICDCS2
2004 Efficient End to End Data Exchange Using Configurable Compression
abstract
We explore the use of compression methods to improve the middleware-based exchange of information in interactive or collaborative distributed applications. In such applications, good compression factors must be accompanied by compression speeds suitable for the data transfer rates sustainable across network links. Our approach combines methods that continuously monitor current network and processor resources and assess compression effectiveness, with techniques that automatically choose suitable compression techniques. The resulting network- and user-aware compression methods are evaluated experimentally across a range of network links and application data, the former ranging from low end links to homes, to wide-area Internet links, to high end links in intranets, the latter including both scientific (binary molecular dynamics data) and commercial (XML) data sets. Results attained demonstrate substantial improvements of this adaptive technique for data compression over non-adaptive approaches, where better compression methods are used when CPU loads are low and/or network links are slow, and where less effective and typically, faster compression techniques are used in high end network infrastructures.
Yair Wiseman, Karsten Schwan, Patrick M. Widener
ICDCS2
2004 StreamGen: A Workload Generation Tool for Distributed Information Flow Applications
abstract
This work presents the StreamGen load generator, which is targeted at distributed information flow applications. These include the event streaming services used in wide-area publish/subscribe systems or in operational information systems, the data streaming services used in remote visualization or collaboration, and the continuous data streams occurring in download services. Running across heterogeneous distributed platforms, these services are implemented by computational component that capture, manipulate, and produce information streams and are linked via overlay topologies. StreamGen can be used to produce the distributed computational and communication loads imposed by these applications. Dynamic application behaviors can be created with mathematical specifications or with behavior traces collected from application-level traces. An interesting set of traces presented in this paper is derived from long -term observations of the FTP download patterns observed at the Linux mirror site being run by the CERCS research center at the Georgia Institute of Technology. Two different flow-based applications are created and evaluated with StreamGen. The first emulates the data streaming behavior in a distributed scientific collaboration, where a scientific simulation (i.e., a molecular dynamics code) produces simulation data sent to and displayed for multiple, interactive remote users. The second emulates portions of the event-streaming behavior of an operational information system used by,a large U.S. corporation. Parametric studies with StreamGen's FTP traces applied to these applications are used to evaluate different load balancing strategies for the cluster machines manipulating these applications' data streams.
Mohamed S. Mansour, Matthew Wolf, Karsten Schwan
ICPP3
2004 IQ-Services: Resource-Aware Middleware for Heterogeneous Applications
abstract
Summary form only given. Heterogeneous computing platforms constitute a challenging execution environment for distributed applications. This article presents a 'systems' view of effective platform usage, by demonstrating the need for application software to be continuously 'aware' of the resources currently available on their underlying heterogeneous computing platforms. Our approach to the implementation of resource awareness is one that (1) provides a 'thin' middleware layer of resource aware services that permit applications to react to changes in resource availability and resources to be managed in accordance with application needs, and that (2) develops compiler- and application-level techniques for dynamic 'service morphing', the goal being to make it easy for application-level services to adjust to runtime changes in application needs or in platform resources. The specific results presented in this article are focused on large-data applications, for which the IQ-services "morphing" layer implements the data manipulations necessary to permit wide-area interactive or multimedia applications to proceed smoothly despite variations in underlying computing and network resources. Experimental results demonstrate substantial performance improvements attained by coordinating network-level with service-level adaptations of the data being transported and by permitting end users to dynamically deploy and use application-specific services for manipulating data in ways suitable for their current needs.
Zhongtang Cai, Greg Eisenhauer, Christian Poellabauer, Karsten Schwan, Matthew Wolf
IPDPS4
2004 Dynamic Data Access to the GT/CERCS Linux Mirror Site
abstract
Summary form only given. The purpose of this work is to better understand the dynamic data sharing behavior of certain classes of grid end users. Toward this end, we study a large-scale data repository and the end users' access behaviors to this repository. Interesting insights from the study include that (1) the use of parallel methods for data downloads, via download accelerators, is common, despite qualms expressed by the community about the impacts of such behavior on wide area data distribution networks, (2) high levels of burstiness exist for such data movements, as also observed for scientific or populist data repositories and Web sites (e.g., space imagery, sports events), and (3) a large number of remote data retrievals are by single clients for single files. We finally discuss the impact of our observations on grid applications in general.
Mohamed S. Mansour, Matthew Wolf, Karsten Schwan
IPDPS3
2004 Energy-Aware Media Transcoding in Wireless Systems
abstract
In distributed systems, transcoding techniques have been used to customize multimedia objects, utilizing trade-offs between the quality and sizes of these objects to provide differentiated services to clients. Our research uses transcoding techniques in wireless systems to customize video streams to the requirements of users, while minimizing the energy costs. We introduce an approach to dynamically determine which transcoders to execute and where to execute them (e.g., client or server). The goal is to select appropriate transcoders (a) to provide clients with the quality of service they desire while (b) minimizing the energy consumption of the end-hosts in accordance with application-specific global energy management directives. This paper investigates sample transcoder functions for video streaming on handheld devices and introduces a mechanism for selecting the most appropriate transcoders and transcoder parameters.
Christian Poellabauer, Karsten Schwan
PerCom2
2004 Energy-Aware Traffic Shaping for Wireless Real-Time Applications
abstract
Sleep modes of wireless network cards are used to switch these cards into low-power state when idle, but large timeout periods and frequent wake-ups can reduce the utility of this approach. Modern processors offer the ability to switch CPU voltages or clock frequencies and therefore reduce CPU energy consumption, however, that can reduce the sleep durations of a network device, adversely affecting the achievable energy savings. This paper describes an approach in which multiple resource managers cooperate to reduce a mobile device's energy consumption. This system-level approach is based on the integrated management of a real-time CPU scheduler, the frequency scaling capabilities of a modern processor, a QoS packet scheduler, and the low-power sleep mode of a wireless network card.
Christian Poellabauer, Karsten Schwan
IEEE Real-Time and Embedded Technology and Applications Symposium2
2004 Dynamic Window-Constrained Scheduling of Real-Time Streams in Media Servers
abstract
We describe an algorithm for scheduling packets in real-time multimedia data streams. Common to these classes of data streams are service constraints in terms of bandwidth and delay. However, it is typical for real-time multimedia streams to tolerate bounded delay variations and, in some cases, finite losses of packets. We have therefore developed a scheduling algorithm that assumes streams have window-constraints on groups of consecutive packet deadlines. A window-constraint defines the number of packet deadlines that can be missed (or, equivalently, 'must be met) in a window of deadlines for consecutive packets in a stream. Our algorithm, called dynamic window-constrained scheduling (DWCS), attempts to guarantee no more than re out of a window of y deadlines are missed for consecutive packets in real-time and multimedia streams. Using DWCS, the delay of service to real-time streams is bounded, even when the scheduler is overloaded. Moreover, DWCS is capable of ensuring independent delay bounds on streams, while, at the same time, guaranteeing minimum bandwidth utilizations over tunable and finite windows of time. We show the conditions under which the total demand for bandwidth by a set of window-constrained streams can exceed 100 percent and still ensure all window-constraints are met. In fact, we show how it is possible to strategically skip certain deadlines in overload conditions, yet fully utilize all available link capacity and guarantee worst-case per-stream bandwidth and delay constraints. Finally, we compare DWCS to the "distance-based" priority (DBP) algorithm, emphasizing the trade-offs of both approaches.
Richard West, Karsten Schwan, Christian Poellabauer
IEEE Trans. Computers3
2003 Differential Data Protection for Dynamic Distributed Application
abstract
We present a mechanism for providing differential data protection to publish/subscribe distributed systems, such as those used in peer-to-peer computing, grid environments, and others. This mechanism, termed "security overlays", incorporates credential-based communication channel creation, subscription and extension. We describe a conceptual model of publish/subscribe services that is made concrete by our mechanism. We also present an application, active video streams, whose reimplementation using security overlays allows it to react to high-level security policies specified in XML without significant performance loss or the necessity for embedding policy-specific code into the application.
Patrick M. Widener, Karsten Schwan, Fabián E. Bustamante
ACSAC2
2003 Resource-Aware Stream Management with the Customizable dproc Distributed Monitoring Mechanisms
abstract
Monitoring the resources of distributed systems is essential to the successful deployment and execution of grid applications, particularly when such applications have well-defined QoS requirements. The dproc system-level monitoring mechanisms implemented for standard Linux kernels have several key components. First, utilizing the familiar /proc filesystem, dproc extends this interface with resource information collected from both local and remote hosts. Second, to predictably capture and distribute monitoring information, dproc uses a kernel-level group communication facility, termed KECho, which is based on events and event channels. Third and the focus of this paper is dproc's run-time customizability for resource monitoring, which includes the generation and deployment of monitoring functionality within remote operating system kernels. Using dproc, we show that: (a) data streams can be customized according to a client's resource availabilities (dynamic stream management); (b) by dynamically varying distributed monitoring (dynamic filtering of monitoring information), appropriate balance can be maintained between monitoring overheads and application quality; and (c) by performing monitoring at kernel-level, the information captured enables decision making that takes into account the multiple resources used by applications.
Sandip Agarwala, Christian Poellabauer, Jiantao Kong, Karsten Schwan, Matthew Wolf
HPDC4
2003 Method Partitioning - Runtime Customization of Pervasive Programs without Design-time Application Knowledge
abstract
Method Partitioning is a dynamic technique for customizing performance-critical message-based interactions between program components, at runtime and without the need for design-time application knowledge. The technique partitions program units that implement message handling, with low costs and high levels of flexibility. It consists of (a) static analysis of a message handling method to produce candidate partitioning plans for the method, (b) cost models for evaluating the cost/benefits of different partitioning plans, (c) a Remote Continuation mechanism that "connects" the distributed parts of a partitioned method at runtime, and (d) Runtime Profiling and Reconfiguration which monitors actual costs of candidate plans and dynamically selects "best" plans from candidates. Experiments with prototypical implementation of Method Partitioning in the JECho distributed event system demonstrate significant performance improvements for both communication-bound and compute-intensive applications, with both applications having dynamic factors that are not predictable at design time.
Santosh Pande, Karsten Schwan
ICDCS3
2003 Opportunistic Channels: Mobility-Aware Event Delivery
Yuan Chen 0001, Karsten Schwan
Middleware2
2003 System-Level Resource Monitoring in High-Performance Computing Environments
Sandip Agarwala, Christian Poellabauer, Jiantao Kong, Karsten Schwan, Matthew Wolf
J. Grid Comput.4
2003 On Network CoProcessors for Scalable, Predictable Media Services
abstract
This paper presents the embedded realization and experimental evaluation of a media stream scheduler on network interface (NI) CoProcessor boards. When using media frames as scheduling units, the scheduler is able to operate in real-time on streams traversing the CoProcessor, resulting in its ability to stream video to remote clients at real-time rates. This paper presents a detailed evaluation of the effects of placing application or kernel-level functionality, like packet scheduling on NIs, rather than the host machines to which they are attached. The main benefits of such placement are: 1) that traffic is eliminated from the host bus and memory subsystem, thereby allowing increased host CPU utilization for other tasks, and 2) that NI-based scheduling is immune to host-CPU loading, unlike host-based media schedulers that are easily affected even by transient load conditions. An outcome of this work is a proposed cluster architecture for building scalable media servers by distributing schedulers and media stream producers across the multiple NIs used by a single server and by clustering a number of such servers using commodity network hardware and software.
Raj Krishnamurthy, Karsten Schwan, Richard West, Marcel-Catalin Rosu
IEEE Trans. Parallel Distributed Syst.2
2003 Dynamic Querying of Streaming Data with the dQUOB System
abstract
Data streaming has established itself as a viable communication abstraction in data-intensive parallel and distributed computations, occurring in applications such as scientific visualization, performance monitoring, and large-scale data transfer. A known problem in large-scale event communication is tailoring the data received at the consumer. It is the general problem of extracting data of interest from a data source, a problem that the database community has successfully addressed with SOL queries, a time tested, user-friendly way for noncomputer scientists to access data. By leveraging the efficiency of query processing provided by relational queries, the dQUOB system provides a conceptual relational data model and SOL query access over streaming data. Queries can be used to extract data, combine streams, and create new streams. The language augments queries with an action to enable more complex data transformations such as Fourier transforms. The dQUOB system has been applied to two large-scale distributed applications: a safety critical autonomous robotics simulation and scientific software visualization for global atmospheric transport modeling. In this paper, we present the dQUOB system and the results of performance evaluation undertaken to assess its applicability in data-intensive wide-area computations, where the benefit of portable data transformation must be evaluated against the cost of continuous query evaluation.
Beth Plale, Karsten Schwan
IEEE Trans. Parallel Distributed Syst.2
2002 A Case for Proactivity in Directory Services
abstract
In this paper, we argue that an exclusively inactive interface to directory services can hinder server scalability and indirectly restrict the behavior of potential applications. We propose to extend directory services' interfaces with a proactive mode by which clients can express their interest in (and be notified of) changes in the environment. These notification channels can be subsequently customized on a per-client basis through client-specified filters. Finally, in order to simplify the handling of client/server failures we adopt a leasing model for client registration to (and customization of) a notification channel. To validate our approach, we have designed and implemented the Proactive Directory Service (PDS).
Fabián E. Bustamante, Patrick M. Widener, Karsten Schwan
HPDC3
2002 IQ-RUDP: Coordinating Application Adaptation with Network Transport
abstract
Our research addresses the efficient transfer of large data across wide-area networks, focusing on applications like remote visualization and real-time collaboration. To attain high performance in the real-time exchange of data across collaborating machines and end users, we are developing and evaluating methods and techniques for coordinating application-level with network transport-level adaptations of data communication. Specifically, complementing previous work on TCP-friendly communication and on adaptive transport protocols, our approach is to strongly coordinate application-level with transport-level changes in communication behavior, so as to best meet application needs without violating fairness in network resource usage. The approach is evaluated with the IQ-ECho middleware, which implements the distribution of scientific data to remote collaborators. Using IQ-ECho, application-level adaptations like selective data down-sampling are triggered by transport-level information provided by the instrumented IQ-RUDP protocol underlying IQ-ECho's communications. The application- to network-layer exchange of information necessary for such coordinated adaptations is implemented with ECho attributes, which provide a lightweight way for an application to provide quality of service information and to describe its adaptation to the transport layer.
Qi He 0001, Karsten Schwan
HPDC2
2002 A Practical Approach for ?Zero? Downtime in an Operational Information System
abstract
An operational information system (OIS) supports a real-time view of an organization's information critical to its logistical business operations. A central component of an OIS is an engine that integrates data events captured from distributed, remote sources in order to derive meaningful real-time views of current operations. This event derivation engine (EDE) continuously updates these views and also publishes them to a potentially large number of remote subscribers. The paper first describes a sample OIS and EDE in the context of an airline's operations. It then defines the performance and availability requirements to be met by this system, specifically focusing on the EDE component. One particular requirement for the EDE is that subscribers to its output events should not experience downtime due to EDE failures, crashes or increased processing loads. Toward this end, we develop and evaluate a practical technique for masking failures and for hiding the costs of recovery from EDE subscribers. This technique utilizes redundant EDEs that coordinate view replicas with a relaxed synchronous fault tolerance protocol. A combination of pre- and post-buffering of replicas is used to attain a solution that offers low response times (i.e., 'zero' downtime) while also preventing system failures in the presence of deterministic faults like 'ill-formed' messages. Parallelism realized via a cluster machine and application-specific techniques for reducing synchronization across replicas are used to scale a 'zero' downtime EDE to support the large number of subscribers it must service.
Ada Gavrilovska, Karsten Schwan, Van Oleson
ICDCS2
2002 Cooperative run-time management of adaptive applications and distributed resources
abstract
This paper presents Q-fabric, which is a set of lightweight, kernel-level abstractions for cooperative, distributed resource management and system/application adaptation. The basis of Q-fabric is its kernel-level, anonymous, asynchronous event service. With this mechanism, (1) applications can monitor and manage the local and remote resources they are using, (2) system-level resource managers can customize their actions to meet the needs of individual applications, and (3) policies can be developed that combine application adaptation with distributed resource management. Results presented in this paper demonstrate the Q-fabric's ability to effectively adapt and manage the resources of a distributed multimedia application. In this application, media streams are adapted at application-level via data down-sampling, and their resources are managed at system-level (e.g., task scheduling) to cope with run-time variations in resource availability. The Q-fabric is implemented as kernel modules on standard Linux platforms.
Christian Poellabauer, Hasan Abbasi, Karsten Schwan
ACM Multimedia3
2002 Scalable directory services using proactivity
abstract
Common to computational grids and pervasive computing is the need for an expressive, efficient, and scalable directory service that provides information about objects in the environment. We argue that a directory interface that ‘pushes’ information to clients about changes to objects can significantly improve scalability. This paper describes the design, implementation, and evaluation of the Proactive Directory Service (PDS). PDS’ interface supports a customizable ‘proactive’ mode through which clients can subscribe to be notified about changes to their objects of interest. Clients can dynamically tune the detail and granularity of these notifications through filter functions instantiated at the server or at the object’s owner, and by remotely tuning the functionality of those filters. We compare PDS’ performance against off-the-shelf implementations of DNS and the Lightweight Directory Access Protocol. Our evaluation results confirm the expected performance advantages of this approach and demonstrate that customized notification through filter functions can reduce bandwidth utilization while improving the performance of both clients and directory servers.
Fabián E. Bustamante, Patrick M. Widener, Karsten Schwan
SC3
2002 SmartPointers: personalized scientific data portals in your hand
abstract
The SmartPointer system provides a paradigm for utilizing multiple light-weight client endpoints in a real-time scientific visualization infrastructure. Together, the client and server infrastructure form a new type of data portal for scientific computing. The clients can be used to personalize data for the needs of the individual scientist. This personalization of a shared dataset is designed to allow multiple scientists, each with their laptops or iPaqs to explore the dataset from different angles and with different personalized filters. As an example, iPaq clients can display 2D derived data functions which can be used to dynamically update and annotate the shared data space, which might be visualized separately on a large immersive display such as a CAVE. Measurements are presented for such a system, built upon the ECho middleware system developed at Georgia Tech.
Matthew Wolf, Zhongtang Cai, Weiyun Huang, Karsten Schwan
SC4
2002 Native Data Representation: An Efficient Wire Format for High-Performance Distributed Computing
abstract
New trends in high-performance software development such as tool- and component-based approaches have increased the need for flexible and high-performance communication systems. When trying to reap the well-known benefits of these approaches, the question of what communication infrastructure should be used to link the various components arises. In this context, flexibility and high-performance seem to be incompatible goals. Traditional HPC-style communication libraries, such as MPI, offer good performance, but are not intended for loosely-coupled systems. Object- and metadata-based approaches like XML offer the needed plug-and-play flexibility, but with significantly lower performance. We observe that the flexibility and baseline performance of data exchange systems are strongly determined by their wire formats, or by how they represent data for transmission in heterogeneous environments. After examining the performance implications of using a number of different wire formats, we propose an alternative approach for flexible high-performance data exchange, Native Data Representation, and evaluate its current implementation in the portable binary I/O library.
Greg Eisenhauer, Fabián E. Bustamante, Karsten Schwan
IEEE Trans. Parallel Distributed Syst.3
2001 Active Streams-An Approach to Adaptive Distributed Systems
abstract
Summary form only given. An increasing number of distributed applications aim to provide services to users by interacting with a correspondingly growing set of data-intensive network services. To support such requirements, we believe that new services need to be customizable, applications need to be dynamically extensible, and both applications and services need to be able to adapt to variations in resource availability and demand. A comprehensive approach to building new distributed applications can facilitate this by considering the contents of the information flowing across the application and its services and by adopting a component-based model to application/service programming. It should provide for dynamic adaptation at multiple levels and points in the underlying platform; and, since the mapping of components to resources in dynamic environment is too complicated, it should relieve programmers of this task. We propose Active Streams, a middleware approach and its associated framework for building distributed applications and services that exhibit these characteristics.
Fabián E. Bustamante, Greg Eisenhauer, Patrick M. Widener, Karsten Schwan, Calton Pu
HotOS4
2001 The Active Streams Approach to Adaptive Distrubuted Systems
abstract
The explosive growth of the Internet, with the emergence of new networking technologies and the increasing number of network-capable end devices, is paving the way for a number of novel distributed applications and services. Cooperative distributed systems have become a common computing model, and pervasive computing has caught the interest of academia and industry. To support future network applications, we believe that new services need to be customizable, applications need to be dynamically extensible, and both applications and services should be able to adapt to variations in resource availability and demand. propose Active Streams (F.E. Bustamante and K. Schwan, 1999), a middleware approach and its associated framework for building distributed applications and services that exhibit these characteristics.
Fabián E. Bustamante, Greg Eisenhauer, Karsten Schwan
HPDC3
2001 Adaptable Mirroring in Cluster Servers
abstract
This paper presents a software architecture for continuously mirroring streaming data received by one node of a cluster-based server to other cluster nodes. The intent is to distribute the load on the server generated by the data's processing and distribution to many clients. This is particularly important when the server not only processes streaming data, but also performs additional processing tasks that heavily depend on current application state. One such task is the preparation of suitable initialization state for thin clients, so that such clients can understand future data events being streamed to them. In particular, when large numbers of thin clients must be initialized at the same time, initialization must be performed without jeopardizing the quality of service offered to regular clients continuing to receive data streams. The mirroring framework presented and evaluated has several novel aspects. First, by performing mirroring at the middleware level, application semantics may be used to reduce mirroring traffic, including filtering events based on their content, by coalescing certain events, or by simply varying mirroring rates according to current application needs concerning the consistencies of mirrored vs. original data. Second, we present an adaptive algorithm that varies mirror consistency and thereby, mirroring overheads in response to changes in clients' request behavior. Third, our framework not only mirrors events, but it can also mirror the new states computed from incoming events, thus enabling dynamic tradeoffs in the communication vs. computation loads imposed on the server node receiving events and on its mirror nodes.
Ada Gavrilovska, Karsten Schwan, Van Oleson
HPDC2
2001 Open Metadata Formats: Efficient XML-Based Communication for High Performance Computing
abstract
High-performance computing faces considerable change as the Internet and the Grid mature. Applications that once were tightly-coupled and monolithic are now decentralized, with collaborating components spread across diverse computational elements. Such distributed systems most commonly communicate through the exchange of structured data. Definition and translation of metadata is incorporated in all systems that exchange structured data. We observe that the manipulation of this metadata can be decomposed into three separate steps: discovery, binding of program objects to the metadata, and marshaling of data to and from wire formats. We have designed a method of representing message formats in XML, using datatypes available in the XML Schema specification. We have implemented a tool, XMIT that uses such metadata and exploits this decomposition in order to provide flexible run-time metadata definition facilities for an efficient binary communication mechanism. We also demonstrate that the use of XMIT makes possible such flexibility at little performance cost.
Patrick M. Widener, Greg Eisenhauer, Karsten Schwan
HPDC3
2001 Open Metadata Formats: Efficient XML-Based Communication for Heterogeneous Distributed Systems
abstract
The definition and translation of metadata is incorporated in all systems that exchange structured data. We observe that the manipulation of this metadata can be decomposed into three separate steps: discovery, binding of program objects to the metadata, and marshaling of data to and from wire formats. We have designed a method of representing message formats in XML, using data types that are available in the XML schema specification. We have implemented a tool called xml2wire that uses such metadata and exploits this decomposition in order to provide flexible metadata definition facilities for an efficient binary communications mechanism. We also observe that the use of xml2wire makes possible such flexibility without intolerable performance effects.
Patrick M. Widener, Karsten Schwan, Greg Eisenhauer
ICDCS2
2001 Optimizations Enabled by Relational Data Model View to Querying Data Streams
abstract
We postulate that the popularity and efficiency of SQL for querying relational databases makes the language a viable solution to retrieving data from data streams. In response, we have developed a system, dQUOB, that uses SQL queries to extract data from streaming data in real time. The high performance needs of applications such as scientific visualization motivates our search for optimizations to improve query evaluation efficiency. The purpose of this paper is to discuss the unique optimizations we have realized by a database point of view to streaming data and to show that the enhanced conceptual model of viewing data streams as relations has reasonable overhead.
Beth Plale, Karsten Schwan
IPDPS2
2001 Taking the Step From Meta-Information to Communication Middleware in Computational Data Streams
abstract
It is our belief that network applications relying on globally distributed shared resources will increasingly adopt meta-level descriptions to describe the data streaming in the application at runtime. Our group has developed the notion of computational data streams to describe and act upon such data flows. Our work conceptualizes the data flows as database relations over which useful operations, such as querying, can be performed. This paper shows how one can make the step from a meta-level description of data flows to an actual implementation using CORBA-style event channels and binary I/O for data transport. 1
Beth Plale, Patrick M. Widener, Karsten Schwan
IPDPS3
2001 Eager Handlers: Communication Optimization in Java-based Distributed Applications with Fine-grained Code Migration
abstract
Java’s platform independence presents an opportunity for using function shipping to optimize communications in distributed applications. This paper introduces the concept of ‘eager handlers’, where handlers are functions that are tightly linked with certain communications performed by Java programs. Their ‘eagerness’ refers to fine-grained, transparent, and efficient function shipping performed for such handlers, with the intent of optimizing communications. The paper describes the design and implementation of eager handlers and also presents results evaluating this concept. It also includes a design for automatically generating eager handlers through static program analysis and online sampling and reconfiguration. Although we have been applying our work to stream-based peer-to-peer communication systems, it is also applicable to client-server applications.
Karsten Schwan
IPDPS2
2001 JECho: Supporting Distributed High Performance Applications with Java Event Channels
abstract
This paper presents JECho, a Java-based communication infrastructure for collaborative high performance applications. JECho implements a publish/subscribe communication paradigm, permitting distributed concurrent sets of components to provide interactive service to collaborating end users via event channels. JECho's eager handler concept allows individual event subscribers to dynamically tailor event flows to adapt to runtime changes in component behaviors and needs, and to changes in platform resources. Benchmark results suggest that JECho may be used for building large-scale, high-performance event delivery systems which can efficiently adapt to changes in user needs or the environment using eager handlers.
Karsten Schwan, Greg Eisenhauer, Yuan Chen 0001
IPDPS2
2001 Coordinated CPU and event scheduling for distributed multimedia applications
abstract
Distributed multimedia applications require support from the underlying operating system to achieve and maintain their desired Quality of Service (QoS). This has led to the creation of novel task and message schedulers and to the development of QoS mechanisms that allow applications to explicitly interact with relevant operating system services. However, the task scheduling techniques developed to date are not well equipped to take advantage of such interactions. As a result, important events such as position update messages in virtual environments may be ignored. If a CPU scheduler ignores these events, players will experience a lack of responsiveness or even inconsistencies in the virtual world. This paper argues that real-time and multimedia applications can benefit from coordinatedel event delivery mechanism, termed ECalls, that supports such coordination. We then show ECalls's ability to reduce variations in inter-frame times for media streams.
Christian Poellabauer, Karsten Schwan, Richard West
ACM Multimedia2
2001 Lightweight kernel/user communication for real-time and multimedia applications
abstract
Operating system enhancements to support real-time and multimedia appl ications often include specializations and extensions of kernel functionality, as with the kernel HTTP daemon (khttpd) in Linux, for instance. To enable efficient and flexible interactions of such extensions with user-level functionality, we have developed ECalls, a lightweight, bidirectional kernel/user event delivery facility, which not only supports the timely delivery of events, but it also reduces the cost and frequency of kernel/user boundary crossings. ECalls is a communication tool that allows (a) kernel extensions to register their offered services and (b) applications to register their interest in these services. Using ECalls, applications use lightweight system calls to generate events, while kernel extensions raise real-time signals or invoke handler functions (residing in either user or kernel space), or they may use kernel threads to handle events on behalf of applications. ECalls can also influence the CPU scheduler such that a process with pending events is given preference over other processes. To demonstrate its utility, this paper implements an I/O event delivery mechanism using ECalls. This mechanism is shown to improve the performance of two applications: a distributed video player and a web server.
Christian Poellabauer, Karsten Schwan, Richard West
NOSSDAV2
2001 CTK: Configurable Object Abstractions for Multiprocessors
abstract
The Configuration Toolkit (CTK) is a library for constructing configurable object based abstractions that are part of multiprocessor programs or operating systems. The library is unique in its exploration of runtime configuration for attaining performance improvements: 1) its programming model facilitates the expression and implementation of program configuration; and 2) its efficient runtime support enables performance improvements by the configuration of program components during their execution. Program configuration is attained without compromising the encapsulation or the reuse of software abstractions. CTK programs are configured using attributes associated with object classes, object instances, state variables, operations, and object invocations. At runtime, such attributes are interpreted by policy classes, which may be varied separately from the abstractions with which they are associated. Using policies and attributes, an object's runtime behavior may be varied by: 1) changing its performance or reliability while preserving the implementation of its functional behavior, or 2) changing the implementation of its internal computational strategy. CTK's multiprocessor implementation is layered on a Cthreads-compatible programming library, which results in its portability to a wide variety of uni- and multiprocessor machines, including a Kendall Square KSR-2 Supercomputer, SGI machines, various SUN workstations, and as a native kernel on the GP1000 BBN Butterfly multiprocessor. The platforms evaluated in the paper are the KSR and SGI machines.
Dilma Da Silva, Karsten Schwan, Greg Eisenhauer
IEEE Trans. Software Eng.2
2000 Event Services for High Performance Computing
abstract
The Internet and the Grid are changing the face of high-performance computing. Rather than tightly-coupled SPMD-style components running in a single cluster, on a parallel machine, or even on the Internet programmed in MPI, applications are evolving into sets of collaborating components scattered across diverse computational elements. These collaborating components may run on different operating systems and hardware platforms and may be written by different organizations in different languages. Complete "applications" are constructed by assembling these components in a plug-and-play fashion. This new vision for high-performance computing demands features and characteristics which are not easily provided by traditional high-performance communications middleware. In response to these needs, we have developed ECho, a high-performance event-delivery middleware that meets the new demands of the Grid environment. ECho provides efficient binary transmission of event data with unique features that support data-type discovery and enterprise-scale application evolution. We present measurements detailing ECho's performance to show that ECho significantly outperforms other systems intended to provide this functionality, and that it provides throughput and latency comparable to the most efficient middleware infrastructures available.
Greg Eisenhauer, Fabián E. Bustamante, Karsten Schwan
HPDC3
2000 dQUOB: Managing Large Data Flows using Dynamic Embedded Queries
abstract
The dQUOB system satisfies client need for specific information from high-volume data streams. The data streams we speak of are the flow of data existing during large-scale visualizations, video streaming to large numbers of distributed users, and high volume business transactions. We introduce the notion of conceptualizing a data stream as a set of relational database tables so that a scientist can request information with an SQL-like query. Transformation or computation that often needs to be performed on the data en-route can be conceptualized as computation performed on consecutive views of the data, with computation associated with each view. The dQUOB system moves the query code into the data stream as a quoblet; as compiled code. The relational database data model has the significant advantage of presenting opportunities for efficient reoptimizations of queries and sets of queries. Using examples from global atmospheric modeling, we illustrate the usefulness of the dQUOB system. We carry the examples through the experiments to establish the viability of the approach for high performance computing with a baseline benchmark. We define a cost-metric of end-to-end latency that can be used to determine realistic cases where optimization should be applied. Finally, we show that end-to-end latency can be controlled through a probability assigned to a query that a query will evaluate to true.
Beth Plale, Karsten Schwan
HPDC2
2000 A Network Co-Processor-Based Approach to Scalable Media Streaming in Servers
abstract
This paper presents the embedded construction and experimental results for a media scheduler on i960 RD equipped I20 Network Interfaces (NI) used for streaming. We utilize the Distributed Virtual Communication Machine (DVCM) infrastructure developed by us which allows run-time extensions to provide scheduling for streams that may require it. The scheduling overhead of such a scheduler is /spl ap/65 /spl mu/s with the ability to stream MPEG video to remote clients at requested rates. Moreover, placement of scheduler action 'close' to the network on the Network Interface (NI) allows tighter coupling of computation and communication, eliminating traffic from the host bus and memory subsystem, allowing increased host CPU utilization for other tasks without being affected by host-CPU loading. Architectures to build scalable media scheduling servers are explored-by distributing media schedulers and media stream producers among NIs within a server and clustering a number of such servers using commodity hardware and software.
Raj Krishnamurthy, Karsten Schwan, Richard West, Marcel-Catalin Rosu
ICPP2
2000 ACDS: Adapting Computational Data Streams for High Performance
abstract
Data-intensive, interactive applications are an important class of metacomputing (Grid) applications. They are characterized by large, time-varying data flows between data providers and consumers. The topic of this paper is the runtime adaptation of data streams, in response to changes in resource availability and/or in end user requirements, with the goal of continually providing to consumers data at the levels of quality they require. Our approach is one that associates computational objects with data streams. Runtime adaptation is achieved by adjusting objects' actions on streams, by splitting and merging objects, and by migrating them (and the streams on which they operate) across machines and network links. Adaptive streams also react to changes in resource availability detected by online monitoring.
Carsten Isert, Karsten Schwan
IPDPS2
2000 Support for Recoverable Memory in the Distributed Virtual Communication Machine
abstract
Distributed Virtual Communication Machine (DVCM) is a software communication architecture for clusters of workstations equipped with programmable network interfaces (Nls) for high-speed networks. DVCM is an extensible architecture, which promotes the transfer of application modules to the NI. By executing "closer" to the network, on the NI CoProcessor, these modules can communicate with significantly higher message rates and lower latencies than achievable at the CPU-level. This paper describes how DVCM modules can be used to enhance the performance of the Cluster Recoverable Memory system (CRMem), a transaction-processing kernel for memory-resident databases. By using the NI CoProcessor for CRMem's remote operations, our implementation achieves more than 3,000 trans/sec on a simplified TpcB benchmark.
Marcel-Catalin Rosu, Karsten Schwan
IPDPS2
2000 Efficient Wire Formats for High Performance Computing
abstract
High performance computing is being increasingly utilized in non-traditional circumstances where it must interoperate with other applications. For example, online visualization is being used to monitor the progress of applications, and real-world sensors are used as inputs to simulations. Whenever these situations arise, there is a question of what communications infrastructure should be used to link the different components. Traditional HPC-style communications systems such as MPI offer relatively high performance, but are poorly suited for developing these less tightly-coupled cooperating applications. Object-based systems and meta-data formats like XML offer substantial plug-and-play flexibility, but with substantially lower performance. We observe that the flexibility and baseline performance of all these systems is strongly determined by their `wire format', or how they represent data for transmission in a heterogeneous environment. We examine the performance implications of different wire formats and present an alternative with significant advantages in terms of both performance and flexibility.
Fabián E. Bustamante, Greg Eisenhauer, Karsten Schwan, Patrick M. Widener
SC3
1999 Min-Cut Methods for Mapping Dataflow Graphs
Volker Elling, Karsten Schwan
Euro-Par2
1999 Steering Data Streams in Distributed Computational Laboratories
abstract
This research supports the interactive access to large-scale scientific data by creation of active user interfaces (AUIs). An AUI continuously emits events describing its current information needs, based on which methods may be developed for controlling the potentially immense information streams directed at the interface. More precisely the purposes of stream control are twofold. First, stream control is performed to deal with heterogeneity in underlying systems, where low end displays may receive only small portions of the data shown at high end displays. Second, stream control is used to achieve scalability with respect to the size and complexity of data streams directed at a user interface, by filtering the data stream and by offloading certain computations from the AUI to the information generators or to information routing sites, by dynamically migrating such computations to appropriate locations, and by adapting these computations in order to effect tradeoffs in the amount of data moved across network links vs. the computations required.
Carsten Isert, Davis King 0001, Karsten Schwan, Beth Plale, Greg Eisenhauer
HPDC3
1999 Run-time Detection in Parallel and Distributed Systems: Application to Safety-Critical Systems
abstract
There is growing interest in run-time detection as parallel and distributed systems grow larger and more complex. This work targets run-time analysis of complex, interactive scientific applications for purposes of attaining scalability improvements with respect to the amount and complexity of the data transmitted, transformed, and shared among different application components. Such improvements are derived from using database techniques to manipulate data streams. Namely, by imposing a relational model on the data streams, constraints on the stream may be expressed as database queries evaluated against the data events comprising the stream. The application in the paper is to a safety-critical system. The paper also presents a tool, dQUOB, Dynamic QUery OBjects, which: (1) offers the means for dynamic creation of queries and for their application to large data streams; (2) permits implementation and runtime use of multiple "query optimization" techniques; and (3) supports dynamic reoptimization of queries based on streams' dynamic behavior.
Beth Plale, Karsten Schwan
ICDCS2
1999 FARACost: An Adaptation Cost Model Aware of Pending Constraints
abstract
The perturbations induced by adaptation and resource allocation decisions on the adapted applications may have the undesirable side effect of causing timing constraint failures. In order to benefit from available adaptation capabilities yet avoid critical timing failures, the dynamic resource allocation mechanism should be aware of the perturbation induced by its decisions. Therefore, the impact of adaptation on short-term performance should be considered a first-class decision criterion, along with traditional criteria such as long-term performance and application criticality. Towards this end, we propose the FARACost, an adaptation cost model that captures the impact of application-specific adaptation procedures and uses this information to evaluate adaptation choices. Experimental evaluations with two applications demonstrate that the use of models like FARACost reduces or prevents pending timing constraint failures, while leading to long-term performance improvements. The experiments are conducted in a cluster environment with a fully implemented infrastructure for adaptation and resource allocation based on the FARA framework.
Daniela Rosu 0001, Karsten Schwan
RTSS2
1999 Composing high-performance schedulers: a case study from real-time simulation
abstract
Dynamic, high-performance or real-time applications require scheduling latencies and throughput not typically offered by current kernel or user-level threads schedulers. Moreover, it is widely accepted that it is important to be able to specialize scheduling policies for specific target applications and their execution environments. This paper presents one solution to the construction of such high-performance, application-specific thread schedulers. Specifically, scheduler implementations are composed from modular components, where individual scheduler modules may be specialized to underlying hardware characteristics or implement precisely the mechanisms and policies desired by application programs. The resulting user-level schedulers' implementations can provide resource guarantees by interaction with kernel-level facilities which provide means of resource reservation. This paper demonstrates the concept of composable schedulers by construction of several compositions for highly dynamic target applications, where low scheduling latencies are critical to application performance. Claims about the importance and effectiveness of scheduler composition are validated experimentally on a shared-memory multiprocessor. Scheduler compositions are optimized to take advantage of different low-level hardware attributes and of knowledge about application requirements specific to certain applications, including a Time Warp-based real-time discrete event simulator. Experimental evaluations are based on synthetic workloads, on a real-time simulation blending simulated with implemented control system components, and on a dynamic robot control program. Measurements indicate that schedulers can be composed and specialized to offer performance similar to that of dedicated scheduling co-processors. Copyright © 1999 John Wiley & Sons, Ltd.
Richard M. Fujimoto, Karsten Schwan
Concurr. Pract. Exp.3
1998 Sender Coordination in the Distributed Virtual Communication Machine
abstract
The Distributed Virtual Communication Machine (DVCM) is an extensible communication architecture for tightly-coupled clusters of workstations (COWs) connected by high-speed networks. The DVCM is designed for off-the-shelf network interface cards equipped with communication coprocessors. Its main component is an active backplane implemented in firmware running on the coprocessors. This backplane can be extended with modules that implement application-specific-functionality and have access to some of the application's state. Consequently non-trivial collective computations can be implemented as DVCM extensions. We present a DVCM extension module that provides application-specific network flow control by coordinating the resource-competing components of a parallel application running on an ATM LAN. Our experiments show that this extension module helps eliminate message loss and achieve high link bandwidth utilization when there is significant link contention.
Marcel-Catalin Rosu, Karsten Schwan
HPDC2
1998 Techniques for Delayed Binding of Monitoring Mechanisms to Application-Specific Instrumentation Points
abstract
Online interaction with computer systems and applications allows developers to monitor, experiment with, and debug long-running resource-intensive applications at runtime. Traditionally, developers statically bind a monitoring mechanism to each application-specific instrumentation point. This approach has shortcomings for online, interactive monitoring. Namely, static binding limits portability among monitoring systems; it may mismatch monitoring mechanisms to interactive requests for monitoring data; and, predictions for the performance and execution paths of instrumentation for static bindings are left to the developer. To address these concerns, we have created a new technique called monitoring assertions that allows monitoring systems to delay binding of monitoring mechanisms to application-specific instrumentation points until runtime. Our empirical results show that we can alter the performance of both the application and the monitoring system by removing static binding requirement of application-specific monitoring systems.
Jeffrey S. Vetter, Karsten Schwan
ICPP2
1998 Falcon: On-line monitoring for steering parallel programs
abstract
Advances in high performance computing, communications and user interfaces enable developers to construct increasingly interactive high performance applications. The Falcon system presented in this paper supports such interactivity by providing runtime libraries, tools and user interfaces that permit the on-line monitoring and steering of large-scale parallel codes. The principal aspects of Falcon described in this paper are its abstractions and tools for capture and analysis of application-specific program information, performed on-line, with controlled latencies and scalable to parallel machines of substantial size. In addition, Falcon provides support for the on-line graphical display of monitoring information, and it allows programs to be steered during their execution, by human users or algorithmically. This paper presents our basic research motivation, outlines the Falcon system's functionality, and includes a detailed evaluation of its performance characteristics in light of its principal contributions. Falcon's functionality and performance evaluation are driven by our experiences with large-scale parallel applications being developed with end users in physics and in atmospheric sciences. The sample application highlighted in this paper is a molecular dynamics simulation program (MD) used by physicists to study the statistical mechanics of liquids. © 1998 John Wiley & Sons, Ltd.
Weiming Gu, Greg Eisenhauer, Karsten Schwan, Jeffrey S. Vetter
Concurr. Pract. Exp.3
1998 Indigo: user-level support for building distributed shared abstractions
abstract
Distributed systems that consist of workstations connected by high performance interconnects offer computational power comparable to moderate size parallel machines. Middleware like distributed shared memory (DSM) or distributed shared objects (DSO) attempts to improve the programmability of such hardware by presenting to application programmers interfaces similar to those offered by shared memory machines. This paper presents the portable Indigo data sharing library which provides a small set of primitives with which arbitrary shared abstractions are easily and efficiently implemented across distributed hardware platforms. Sample shared abstractions implemented with Indigo include DSM as well as fragmented objects, where the object state is split across different machines and where interfragment communications may be customized to application-specific consistency needs. The Indigo library's design and implementation are evaluated on two different target platforms: a workstation cluster and an IBM SP2 machine. As part of this evaluation, a novel DSM system and consistency protocol are implemented and evaluated with several high performance applications. Application performance attained with the DSM system is compared to the performance experienced when utilizing the underlying basic message-passing facilities or when employing Indigo to construct customized fragmented objects implementing the application's shared state. Such experimentation results in insights concerning the efficient implementation of DSM systems (e.g. how to deal with false sharing). It also leads to the conclusion that Indigo provides a sufficiently rich set of abstractions for efficient implementation of the next generation of parallel programming models for high performance machines. © 1998 John Wiley & Sons, Ltd.
Prince Kohli, Mustaque Ahamad, Karsten Schwan
Concurr. Pract. Exp.3
1998 DataExchange: High Performance Communications in Distributed Laboratories
Greg Eisenhauer, Beth Plale, Karsten Schwan
Parallel Comput.3
1997 Supporting Parallel Applications on Clusters of Workstations: The Intelligent Network Interface Approach
abstract
This paper presents a novel networking architecture designed for communication intensive parallel applications running on clusters of workstations (COWs) connected by high speed network. This architecture permits: (1) the transfer of selected communication-related functionality the host machine to the network interface coprocessor and (2) the exposure of this functionality directly to applications as instructions of a Virtual Communication Machine (VCM) implemented by the coprocessor. The user-level code interacts directly with the network coprocessor as the host kernel only 'connects' the application to the VCM and does not participate in the data transfers. The distinctive feature of our design is its flexibility: the integration of the network with the application can be varied to maximize performance. The resulting communication architecture is characterized by a very low overhead on the host processor by latency and bandwidth close to the hardware limits, and by an application interface which enables zero-copy messaging and eases the port of some shared-memory parallel applications to COWs. The architecture admits low cost implementations based only on off-the-shelf hardware components. Additionally, its current ATM-based implementation can be used to communicate with any ATM-enabled host.
Marcel-Catalin Rosu, Karsten Schwan, Richard M. Fujimoto
HPDC2
1997 Exploiting Temporal and Spatial Constraints on Distributed Shared Objects
abstract
Gigabit network technologies have made it possible to combine workstations into a distributed, massively-parallel computer system. Middleware, such as distributed shared objects (DSO), attempts to improve programmability of such systems, by providing globally accessible 'object' abstractions. Researchers have developed consistency protocols for replicated 'memory' objects. These protocols are well suited to scientific applications but less suited to multimedia or groupware applications. We address the state sharing needs of complex distributed applications with: high-frequency symmetric accesses to shared objects; unpredictable and limited locality of accesses; dynamically changing sharing behavior; and potential data races. We show that a DSO system exploiting application-level temporal and spatial constraints on shared objects can outperform shared object protocols which do not exploit application-level constraints. We compare our S(emantic) DSO against entry consistency using a sample application having the four properties mentioned above.
Richard West, Karsten Schwan, Ivan Tacic, Mustaque Ahamad
ICDCS2
1997 On adaptive resource allocation for complex real-time application
abstract
Resource allocation for high-performance real-time applications is challenging due to the applications' data-dependent nature, dynamic changes in their external environment, and limited resource availability in their target embedded system platforms. These challenges may be met by use of adaptive resource allocation (ARA) mechanisms that can promptly adjust resource allocation to changes in an application's resource needs, whenever there is a risk of failing to satisfy its timing constraints. By taking advantage of an application's adaptation capabilities, ARA eliminates the need for 'over-sizing' real-time systems to meet worst-case application needs. This paper proposes a model for describing an application's adaptation capabilities and the runtime variation of its resource needs. The paper also proposes a satisfiability-driven set of performance metrics for capturing the impact of ARA mechanisms on the performance of adaptable real-time applications. The relevance of the proposed set of metrics is demonstrated experimentally, using a synthetic application designed to represent time-critical applications in C31 systems.
Daniela Rosu 0001, Karsten Schwan, Sudhakar Yalamanchili, Rakesh Jha
RTSS2
1997 Software Approach to Hazard Detection Using On-line Analysis of Safety Constraints
abstract
Hazard situations in safety-critical systems are typically complex, so there is a need for means to detect complex hazards and react in a timely and meaningful way. This paper addresses the problem of hazard detection through the development of an online analysis tool. The approach allows the user to specify complex multi-source hazards using a query-like language, uses both synchronous and asynchronous online checking approaches to balance efficiency and expressiveness, accommodates dynamic applications through dynamic constraint addition, and supports distributed and parallel applications running in heterogeneous environments.
Beth A. Schroeder, Karsten Schwan, Sudhir Aggarwal
SRDS2
1996 Adaptive resource allocation for embedded parallel applications
abstract
Parallel and distributed computer architectures are increasingly being considered for application in a wide variety of computationally intensive embedded systems. Many such applications impose highly dynamic demands for resources (processors, memory, and communication network), because their computations are data-dependent, or because the applications must constantly interact with a rapidly changing physical environment, or because the applications themselves are adaptive. This paper presents a set of dynamic resource allocation techniques aimed at maintaining high levels of application performance in the presence of varying resource demands. It focuses on a class of applications structured as multiple pipelines of data-parallel stages, as this structure is common to many sensor-based applications. We discuss the issues involved in resource management for such applications, and present preliminary results from our implementations on Intel Paragon. Our approach uses feedback control-a real-time monitoring system is used to detect significant performance shortfalls, and resources are reallocated among the application components in an attempt to improve performance. The main contribution of this work is that it combines real-time monitoring of an application's performance with dynamic resource allocation, and focuses on practical implementations rather than simulation and analysis.
Rakesh Jha, Mustafa Muhammad, Sudhakar Yalamanchili, Karsten Schwan, Daniela Rosu 0001, Chris deCastro
HiPC4
1996 Improving Protocol Performance by Dynamic Control of Communication Resources
abstract
A problem frequently faced by complex distributed applications is to control the interaction of their communication and computational activities such that they jointly adhere to desired performance and timing requirements. This research concerns communication infrastructures able to cope with the varying processing and QoS/sup 1/ requirements imposed on them by application programs. Specifically, we describe and evaluate COMM/sup adapt/, a communication infrastructure enabling the on-line adaptation of a protocol's resource usage to currently available resources and application requirements. The key feature of COMM/sup adapt/ is its dynamic (auto-)configurability, which is its support of on-line configuration transparent to application programs. Such configuration is performed by a heuristic that accommodates changes in a connection's resource requirements by reallocating resources based on its knowledge of actual resource usage by other active connections. The heuristic's design and implementation are based on extensive investigations of the manner in which alternative assignments of protocol tasks to underlying processing resources can influence program-level latency and throughput requirements.
Daniela Rosu 0001, Karsten Schwan
ICECCS2
1996 A parallel spectral model for atmospheric transport processes
abstract
The paper describes a parallel implementation of a grand challenge problem: global atmospheric modeling. The novel contributions of our work include (1) a detailed investigation of opportunities for parallelism in atmospheric global modeling based on spectral solution methods, (2) the experimental evaluation of overheads arising from load imbalances and data movement for alternative parallelization methods, and (3) the development of a parallel code that can be monitored and steered interactively based on output data visualizations and animations of program functionality or performance. Code parallelization takes advantage of the relative independence of computations at different levels in the earth's atmosphere, resulting in parallelism of up to 40 processors, each independently performing computations for different atmospheric levels and requiring few communications between different levels across model time steps. Next, additional parallelism is attained within each level by taking advantage of the natural parallelism offered by the spectral computations being performed (e.g. taking advantage of independently computable terms in equations). Performance measurements are performed on a 64-node KSR2 supercomputer. However, the parallel code has been ported to several shared memory parallel machines, including SGI multiprocessors, and has also been ported to distributed memory platforms like the IBM SP-2.
Thomas Kindler, Karsten Schwan, Dilma Da Silva, Mary Trauner, Fred Alyea
Concurr. Pract. Exp.2
1996 Design and Analysis of a Parallel Molecular Dynamics Application
Greg Eisenhauer, Karsten Schwan
J. Parallel Distributed Comput.2
1996 Distributed Shared Abstractions (DSA) on Multiprocessor
abstract
Any parallel program has abstractions that are shared by the program's multiple processes. Such shared abstractions can considerably affect the performance of parallel programs, on both distributed and shared memory multiprocessors. As a result, their implementation must be efficient, and such efficiency should be achieved without unduly compromising program portability and maintainability. The primary contribution of the DSA library is its representation of shared abstractions as objects that may be internally distributed across different nodes of a parallel machine. Such distributed shared abstractions (DSA) are encapsulated so that their implementations are easily changed while maintaining program portability across parallel architectures. The principal results presented are: a demonstration that the fragmentation of object state across different nodes of a multiprocessor machine can significantly improve program performance; and that such object fragmentation can be achieved without compromising portability by changing object interfaces. These results are demonstrated using implementations of the DSA library on several medium scale multiprocessors, including the BBN Butterfly, Kendall Square Research, and SGI shared memory multiprocessors. The DSA library's evaluation uses synthetic workloads and a parallel implementation of a branch and bound algorithm for solving the traveling salesperson problem (TSP).
Christian Clémençon, Bodhisattwa Mukherjee, Karsten Schwan
IEEE Trans. Software Eng.3
1995 Indigo: User-Level Support for Building Distributed Shared Abstractions
abstract
Distributed systems that consist of workstations connected by high performance interconnects offer computational power comparable to moderate size parallel machines. It is desirable that such workstation clusters can also be programmed the same way as shared memory machines. We develop a portable, user-level library, called Indigo, that can be used to program a variety of state sharing techniques. In particular, Indigo can be used to program DSM protocols as well as distributed shared abstractions where objects can be fragmented/replicated and consistency actions are customized according to application needs. We present an evaluation of Indigo by using its calls to implement a distributed shared memory system as well as shared abstractions for a number of applications.
Prince Kohli, Mustaque Ahamad, Karsten Schwan
HPDC3
1995 Progress: A Toolkit for Interactive Program Steering
Jeffrey S. Vetter, Karsten Schwan
ICPP (2)2
1995 Data Interpretation and Experiment Planning in Performance Tools (Panel)
abstract
The parallel scientific computing community is placing increasing emphasis on portability and scalability of programs, languages, and architectures. This creates new challenges for developers of parallel performance analysis tools, who will have to deal with increasing volumes of performance data drawn from diverse platforms. One way to meet this challenge is to incorporate sophisticated facilities for data interpretation and experiment planning within the tools themselves, giving them increased flexibility and autonomy in gathering and selecting performance data. This panel discussion brings together four research groups that have made advances in this direction.
Allen D. Malony, B. Robert Helm, Jeffrey K. Hollingsworth, Barton P. Miller, Karsten Schwan
SIGMETRICS5
1995 An Integrated Approach for Steering, Visualization, and Analysis of Atmospheric Simulations
Yves Jean, Thomas Kindler, William Ribarsky, Weiming Gu, Greg Eisenhauer, Karsten Schwan, Fred Alyea
IEEE Visualization6
1994 Falcon - Toward Interactive Parallel Programs: The On-line Steering of a Molecular Dynamics Application
abstract
The paper focuses on the opportunities and costs of online steering as applied to a substantial parallel application. We demonstrate potential performance improvements through the use of the Falcon system, an experimental system for the online monitoring and steering of parallel programs. The visual presentation of program output along with animated displays of program performance information via Falcon's monitoring system enables the online capture, analysis, and display of program information required for program steering. Falcon also provides the mechanisms for the manipulations of program state that accomplish this online steering.>
Greg Eisenhauer, Karsten Schwan, Weiming Gu, Niru Mallavarupu
HPDC2
1993 Improving Performance by Use of Adaptive Objects: Experimentation with a Configurable Multiprocessor Thread Package
abstract
Since the mechanisms of an operating system can significantly affect the performance of parallel programs, it is important to customize operating system functionality for specific application programs. The authors first present a model for adaptive objects and the associated mechanisms, then they use this model to implement adaptive locks for multiprocessors which adapt themselves according to user-provided adaptation policies to suit changing application locking patterns. Using a parallel branch and bound program, they demonstrate the performance advantage of adaptive locks over existing locks.>
Bodhisattwa Mukherjee, Karsten Schwan
HPDC2
1993 Parallel and configurable protocols: experiences with a prototype and an architectural framework
abstract
The authors consider the use of parallelism and configurability to increase throughput and reduce protocol processing latencies. They obtain experimental results on parallel protocol performance using a prototype implemented on a shared memory multiprocessor. The results demonstrate the utility of parallel protocol processing, and they indicate the further research necessary for constructing viable communication protocols for large-scale parallel machines. Based on these experiences, the design of an object-oriented framework for parallel protocol programming which facilitates parallel protocol development and helps maximize protocol performance on a wide variety of multiprocessors is presented.>
Bert Lindgren, Mostafa H. Ammar, Bobby Krupczak, Karsten Schwan
ICNP4
1993 Experiments with Configurable Locks for Multiprocessors
abstract
Operating system kernels typically offer a fixed set of mechanisms and primitves. However, the attainment high of performance for a variety parallel application requires the availability of reconfigurable and extensible operattng system kernel primitives. In this paper, we present an implementation of multiprocessor locks that can be reconfigured dynamically.
Bodhisattwa Mukherjee, Karsten Schwan
ICPP (2)2
1993 CHAOS-arc: Kernel Support for Multiweight Objects, Invocations, and Atomicity in Real-Time Multiprocessor Applications
abstract
CHAOS arc is an object-based multiprocessor operating system kernel that provides primitives with which programmers may easily construct objects of differing types and object invocations of differing semantics, targeting multiprocessor systems, and real-time applications. The CHAOS arc can guarantee desired performance and functionality levels of selected computations in real-time applications. Such guarantees can be made despite possible uncertainty in execution environments by allowing programs to adapt in performance and functionality to varying operating conditions. This paper reviews the primitives offered by CHAOS arc and demonstrates how the required elements of the CHAOS arc real-time kernel are constructed with those primitives.
Ahmed Gheith, Karsten Schwan
ACM Trans. Comput. Syst.2
1993 Application-Dependent Dynamic Monitoring of Distributed and Parallel Systems
abstract
Achieving high performance for parallel or distributed programs often requires substantial amounts of information about the programs themselves, about the systems on which they are executing, and about specific program runs. The monitoring system that collects, analyzes, and makes application-dependent monitoring information available to the programmer and to the executing program is presented. The system may be used for off-line program analysis, for on-line debugging, and for making on-line, dynamic changes to parallel or distributed programs to enhance their performance. The authors use a high-level, uniform data model for the representation of program information and monitoring data. They show how this model may be used for the specification of program views and attributes for monitoring, and demonstrate how such specifications can be translated into efficient, program-specific monitoring code that uses alternative mechanisms for the distributed analysis and collection to be performed for the specified views. The model's utility has been demonstrated on a wide variety of parallel machines.>
David M. Ogle, Karsten Schwan, Richard T. Snodgrass
IEEE Trans. Parallel Distributed Syst.2
1992 Performance Effects of Information Sharing in a Distributed Multiprocessor Real-Time Scheduler
abstract
Two questions are examined, regarding real-time multiprocessor scheduling for large-scale nonuniform memory access (NUMA) architectures: how are the latency and the quality of scheduling affected by different degrees of completeness in the information shared among multiple potentially concurrent schedules, and how can scheduling information be represented so that it is efficiently and concurrently accessible? The authors present a real-time scheduling algorithm for multiprocessors that is scalable in the number of tasks performing scheduling and in the maximum amount of computation time consumed by those tasks. They also develop a flexible representation for shared information within the distributed scheduler that is easily varied regarding its degree of information completeness. It is then shown that the sharing of incomplete (vs. complete) information can lead to increased performance regarding scheduling latency with few or no losses in scheduling quality. In addition, it is shown that this holds for a variety of parallel machines, ranging from NUMA to distributed memory machines.>
Hongyi Zhou, Karsten Schwan, Ian F. Akyildiz
RTSS2
1992 Data base design for real-time adaptations
Prabha Gopinath, Rajiv Ramnath, Karsten Schwan
J. Syst. Softw.3
1992 Dynamic Scheduling of Hard Real-Time Tasks and Real-Time Threads
abstract
The authors investigate the dynamic scheduling of tasks with well-defined timing constraints. They present a dynamic uniprocessor scheduling algorithm with an O(n log n) worst-case complexity. The preemptive scheduling performed by the algorithm is shown to be of higher efficiency than that of other known algorithms. Furthermore, tasks may be related by precedence constraints, and they may have arbitrary deadlines and start times (which need not equal their arrival times). An experimental evaluation of the algorithm compares its average case behavior to the worst case. An analytic model used for explanation of the experimental results is validated with actual system measurements. The dynamic scheduling algorithm is the basis of a real-time multiprocessor operating system kernel developed in conjunction with this research. Specifically, this algorithm is used at the lowest, threads-based layer of the kernel whenever threads are created.>
Karsten Schwan, Hongyi Zhou
IEEE Trans. Software Eng.1
1991 Dynamic Adaption of Real-Time Software
abstract
In large, dynamic, real-time computer systems, it is frequently most cost effective to employ different software performance and reliability techniques at different levels of granularity, at different times, or within different subsystems. These techniques may include regulation of redundancy and resource allocation, multiversion and multipath execution, adjustments of program attributes such as time-out periods and others. The management of software in such systems is a difficult task. Software that may be adapted to meet varying performance and reliability requirements offers a solution. A REal-time Software Adaptation System (RESAS) includes a uniform model of adaptable software and provides the tool necessary for programmers to implement algorithms that choose and enact adaptations in real time. RESAS has been implemented on a testbed consisting of a multiprocessor and an attached workstation, and adaptation algorithms have been developed that address the problem of adapting software to achieve two goals: software execution within specified time constraints and software resiliency with respect to computer hardware failures.
Thomas E. Bihari, Karsten Schwan
ACM Trans. Comput. Syst.2
1991 Experimental Evaluation of a Real-Time Scheduler for a Multiprocessor System
abstract
A description is given of the design, implementation, and experimental evaluation of a multiprocessor scheduler used with robotics applications and other real-time programs. The scheduler makes decisions concerning both the assignment of processes and the scheduling of these processes on each processor such that a near-optimal numer of processor deadlines is satisfied. It assumes that process execution times, deadlines, and earliest possible start times are known.>
Ben A. Blake, Karsten Schwan
IEEE Trans. Software Eng.2
1990 From CHAOS^it base to CHAOS^it arc: A Family of Real-Time Kernels
abstract
The authors present a family of object-based real-time operating system kernels that address portability, extensibility, and customizability for low-level and subsystem-level operations. The family is extensible in that new abstractions and functionalities can be added easily and efficiently, thereby maintaining uniform kernel interfaces and permitting the implementation of domain or target machine specific features while preserving some given kernel interface for existing programs. It also provides an environment for experimenting with and prototyping of new operating system constructs and policies. The family is customizable in that existing kernel abstractions and functions can be modified easily, facilitating changes to an operating system for uses with different target architectures or application domains. The family is portable in that its implementation is based on the Mach C threads standard as a base layer for uniprocessors and parallel architectures.>
Karsten Schwan, Ahmed Gheith, Hongyi Zhou
RTSS1
1990 "Topologies" - Distributed Objects on Multicomputers
abstract
Application programs written for large-scale multicomputers with interconnection structures known to the programmer (e.g., hypercubes or meshes) use complex communication structures for connecting the applications' parallel tasks. Such structures implement a wide variety of functions, including the exchange of data or control information relevant to the task computations and/or the communications required for task synchronization, message forwarding/filtering under program control, and so on.Topologyis a programming and operating system construct that allows programmers to describe and efficiently implement such functionality as distributed objects with well-defined operational interfaces. As with abstract data types, topologies may be reused by any application desiring their functionality. However, in contrast to other research in parallel or distributed object-based operating systems, internally, a topology may be an entirely distributed implementation of the object's functionality, consisting of a communication graph and type-specific computations, which are triggered by messages traversing the graph. Sample computations may perform additions or minimizations of the values traversing a topology, thereby computing a global sum or minimum. Similarly, computations may concatenate or filter messages in order to implement program monitoring, I/O, file storage, or virtual terminal services. Topologies are implemented as an extension of the Intel iPSC hypercube's operating system kernel and have been used with several, large-scale parallel application programs.
Karsten Schwan, Win Bo
ACM Trans. Comput. Syst.1
1989 Object-Oriented Design of Real-Time Software
abstract
An examination is made of the ability of the object-oriented model to represent and encapsulate the temporal characteristics of adaptable, real-time software. The differences between object-oriented design and the more traditional function-oriented design of real-time systems are briefly reviewed. The various temporal characteristics present in a real-time, object-oriented system and techniques for the design of object-oriented real-time software are examined. The many open problems in the area are discussed. The object-oriented model and many of the concepts described have been used in the implementation of the control software for an experimental robotics project.>
Thomas E. Bihari, Prabha Gopinath, Karsten Schwan
RTSS3
1989 Global Data and Control in Multicomputers: Operating System Primitives and Experimentation with a Parallel Branch-and-Bound Algorithm
abstract
Abstract The static distribution of work among tasks is not possible in many parallel applications. Therefore, it is essential to implement convenient and efficient abstractions for ‘work sharing’ on multicomputers. This paper compares the utility of two operating system facilities for the implementation of such ‘work sharing’: (1) a system for the migration of processes from heavily to less loaded processors and (2) a more general OS construct for the implementation of arbitrary distributed objects. Both were implemented as extensions to the Intel iPSC/1 operating system on a 32‐node hypercube. Their experimental evaluation is based on a parallel implementation of a branch‐and‐bound algorithm. Two sets of results are attained. First, the necessity of the constructs for dynamic work sharing is demonstrated for applications with dynamic data domains, such as parallel branch‐and‐bound algorithms. This is followed by measurements that demonstrate the acceptable cost of process migration for a specific parallel branch‐and‐bound algorithm. These measurements are then compared with results attained with the construct for the implementation of distributed objects. Second, when using branch‐and‐bound to solve the Travelling Salesperson Problem (TSP), evaluation of the resulting parallel TSP program shows that some analytical and simulation results attained in past, published work may not hold.
Karsten Schwan, Ben A. Blake, Win Bo, J. Gawkowski
Concurr. Pract. Exp.1
1989 Automating resource allocation for multiprocessors
Karsten Schwan, Cheryl Gaimon
J. Syst. Softw.1
1988 A Comparison of Four Adaptation Algorithms for Increasing the Reliability of Real-Time Software
abstract
In a large, parallel, real-time computer system, it is frequently most cost-effective to use different software reliability techniques (e.g., retry, replication, and resource-allocation algorithms) at different levels of granularity or within different subsystems, dependent on the reliability requirements and fault models associated with each subsystem. The authors describe the results of applying four reliability algorithms, via RESAS (REal-time Software Adaptation System), to a sample real-time program executing on a multiprocessor.>
Thomas E. Bihari, Karsten Schwan
RTSS2
1988 A Language and System for the Construction and Tuning of Parallel Programs
abstract
The programming of efficient parallel software typically requires extensive experimentation with program prototypes. To facilitate such experimentation, any programming system that supports rapid prototyping of parallel programs should provide high-level language primitives with which programs can be explicitly, statically, or dynamically tuned with respect to performance and reliability. Such language primitives should be able to refer conveniently to the information about the executing program and the parallel hardware required for tuning. Such information may include monitoring data about the current or previous program or even hints regarding appropriate tuning decisions. Language primitives and an associated programming system for program tuning are presented. The primitives and system have been implemented, and have been tested with several parallel applications on a network of Unix workstations.>
Karsten Schwan, Rajiv Ramnath, Sridhar Vasudevan, David M. Ogle
IEEE Trans. Software Eng.1
1987 A System for Parallel Programming
Karsten Schwan, Rajiv Ramnath, Sridhar Vasudevan, David M. Ogle
ICSE1
1987 Adaptive, Reliable Software for Distributed and Parallel Real-Time Systems
Karsten Schwan, Thomas E. Bihari, Ben A. Blake
SRDS1
1987 CHAOS-Kernel Support for Objects in the Real-Time Domain
abstract
We describe our experiences with a real-time multiprocessor operating system, called GEM (Generalized Executive for Multiprocessor applications) and with an extension to GEM, called CHAOS (Concurrent Hierarchical Adaptable Object System). CHAOS offers kernel-level primitives that allow high-performance, large-scale, real-time software to be programmed as a system of interacting objects. This significantly improves modularity, reconfigurability, and maintainability.
Karsten Schwan, Prabha Gopinath, Win Bo
IEEE Trans. Computers1
1987 High-Performance Operating System Primitives for Robotics and Real-Time Control Systems
abstract
To increase speed and reliability of operation, multiple computers are replacing uniprocessors and wired-logic controllers in modern robots and industrial control systems. However, performance increases are not attained by such hardware alone. The operating software controlling the robots or control systems must exploit the possible parallelism of various control tasks in order to perform the necessary computations within given real-time and reliability constraints. Such software consists of both control programs written by application programmers and operating system software offering means of task scheduling, intertask communication, and device control. The Generalized Executive for real-time Multiprocessor applications (GEM) is an operating system that addresses several requirements of operating software. First, when using GEM, programmers can select one of two different types of tasks differing in size, called processes and microprocesses. Second, the scheduling calls offered by GEM permit the implementation of several models of task interaction. Third, GEM supports multiple models of communication with a parameterized communication mechanism. Fourth, GEM is closely coupled to prototype real-time programming environments that provide programming support for the models of computation offered by the operating system. GEM is being used on a multiprocessor with robotics application software of substantial size and complexity.
Karsten Schwan, Thomas E. Bihari, Bruce W. Weide, Gregor Taulbee
ACM Trans. Comput. Syst.1
1986 A High-Performance, Object-Based Operating System for Real-Time, Robotics Applications
Karsten Schwan, Win Bo, Prabha Gopinath
RTSS1
1986 Flexible Software Development for Multiple Computer Systems
abstract
The authors develop a model of concurrent software and the associated programming tools that jointly permit flexible software development for experimental programming on the Cm* multiprocessor. The model's implementation in the TASK and Bliss-11 programming languages is described using a sample concurrent program.
Karsten Schwan, Anita K. Jones
IEEE Trans. Software Eng.1
1985 Automating Resource Allocation for the Cm* Multiprocessor
Karsten Schwan, Cheryl Gaimon
ICDCS1
1985 GEM: Operating system primitives for robots and real-time control systems
abstract
To increase the speed and reliability of robots and of industrial control systems, multiple processing elements are used in their computing hardware. However, performance increases are not attained by hardware, alone. It is the hardware's operating software that must exploit the possible parallelism to gain the increases desired. Such software consists of control programs written by application programmers and operating system software offering means of task scheduling, inter-task communication, and hardware configuration control. The Generalized Executive for real-time Multiprocessor applications (GEM) is an operating system that addresses several problems arising due to the unique requirements of operating software, including: (1) GEM supports two different sizes of tasks and task scheduling, called processes and micro-processes, and offers a variety of real-time scheduling calls, and (2) GEM supports multiple models of communication.
Karsten Schwan, Thomas E. Bihari, Bruce W. Weide, Gregor Taulbee
ICRA1
1979 TASK Forces: Distributed Software for Solving Problems of Substantial Size
Anita K. Jones, Karsten Schwan
ICSE2
1979 StarOS, a Multiprocessor Operating System for the Support of Task Forces
abstract
StarOS is a message-based, object-oriented, multiprocessor operating system, specifically designed to support task forces, large collections of concurrently executing processes that cooperate to accomplish a single purpose. StarOS has been implemented at Carnegie-Mellon University for the 50 processor Cm* multi-microprocessor computer.
Anita K. Jones, Robert J. Chansler Jr., Ivor Durham, Karsten Schwan, Steven R. Vegdahl
SOSP4