Michael Stumm

dblp:82/5150 · DBLP profile ↗
← Back
61ranked-venue papers
0as first author
9since 2021 · last 2025
0000-0002-9377-2493ORCID · corroborated

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

Systems, architecture and hardware · 40 · 8 since 2021Software engineering, systems software and programming languages · 17 · 1 since 2021Databases, data management, data science and information retrieval · 3 · 3 since 2021Security and privacy · 2Computer networks · 1Graphics, computer vision, multimedia, augmented reality and games · 1Applied, interdisciplinary, general and emerging computing · 1
YearPublicationVenuePosition
2025 PaperCache: In-Memory Caching with Dynamic Eviction Policies
abstract
In-memory caches play a critical role in storage environments by reducing data access latencies and loads on backend data stores. A cache's eviction policy significantly impacts its attained miss ratio, and recent modeling techniques allow for efficient evaluation of different eviction policies at runtime. However, modern in-memory caches lack the ability to switch between eviction policies at runtime, except for Redis that can only switch between LRU and LFU. We present PaperCache, an in-memory cache capable of switching between multiple different eviction policies at runtime. Our evaluation shows that immediately after an eviction policy switch, PaperCache's behavior closely mirrors that of a cache implementing the target policy exactly (with a miss ratio typically within 1%) for a short period of time, after which PaperCache's behavior is fully inline with an exact policy implementation. Further, PaperCache is able to periodically and automatically switch to the policy exhibiting the lowest miss ratio, reducing the overall miss ratio by up to 48.5%.
Kia Shakiba, Michael Stumm
HotStorage2
2024 TTLs Matter: Efficient Cache Sizing with TTL-Aware Miss Ratio Curves and Working Set Sizes
abstract
In-memory caches play a pivotal role in optimizing distributed systems by significantly reducing query response times. Correctly sizing these caches is critical, especially considering that prominent organizations use terabytes and even petabytes of DRAM for these caches. The Miss Ratio Curve (MRC) and Working Set Size (WSS) are the most widely used tools for sizing these caches.
Sari Sultan, Kia Shakiba, Paul Chen, Michael Stumm
EuroSys5
2024 Kosmo: Efficient Online Miss Ratio Curve Generation for Eviction Policy Evaluation
Kia Shakiba, Sari Sultan, Michael Stumm
FAST3
2022 ctFS: Replacing File Indexing with Hardware Memory Translation through Contiguous File Allocation for Persistent Memory
Ruibin Li, Xiang Ren 0003, Xu Zhao 0004, Siwei He, Michael Stumm, Ding Yuan 0004
FAST5
2022 Investigating Managed Language Runtime Performance: Why JavaScript and Python are 8x and 29x slower than C++, yet Java and Go can be Faster?
David Lion, Adrian Chiu, Michael Stumm, Ding Yuan 0004
USENIX ATC3
2022 ctFS: Replacing File Indexing with Hardware Memory Translation through Contiguous File Allocation for Persistent Memory
abstract
Persistent byte-addressable memory (PM) is poised to become prevalent in future computer systems. PMs are significantly faster than disk storage, and accesses to PMs are governed by the Memory Management Unit (MMU) just as accesses with volatile RAM. These unique characteristics shift the bottleneck from I/O to operations such as block address lookup—for example, in write workloads, up to 45% of the overhead in ext4-DAX is due to building and searching extent trees to translate file offsets to addresses on persistent memory. We propose a novel contiguous file system, ctFS, that eliminates most of the overhead associated with indexing structures such as extent trees in the file system. ctFS represents each file as a contiguous region of virtual memory, hence a lookup from the file offset to the address is simply an offset operation, which can be efficiently performed by the hardware MMU at a fraction of the cost of software-maintained indexes. Evaluating ctFS on real-world workloads such as LevelDB shows it outperforms ext4-DAX and SplitFS by 3.6× and 1.8×, respectively.
Ruibin Li, Xiang Ren 0003, Xu Zhao 0004, Siwei He, Michael Stumm, Ding Yuan 0004
ACM Trans. Storage5
2021 Evolution of Development Priorities in Key-value Stores Serving Large-scale Applications: The RocksDB Experience
Siying Dong, Andrew Kryczka, Yanqin Jin, Michael Stumm
FAST4
2021 Bladerunner: Stream Processing at Scale for a Live View of Backend Data Mutations at the Edge
abstract
Consider a social media platform with hundreds of millions of online users at any time, utilizing a social graph that has many billions of nodes and edges. The problem this paper addresses is how to provide each user a continuously fresh, up-to-date view of the parts of the social graph they are currently interested in, so as to provide a positive interactive user experience. The problem is challenging because the social graph mutates at a high rate, users change their focus of interest frequently, and some mutations are of interest to many online users.
Jeff Barber, Ximing Yu, Laney Kuenzel Zamore, Zhiyuan Jerry Lin, Vahid Jazayeri, Shie Erlich, Tony Savor, Michael Stumm
SOSP8
2021 RocksDB: Evolution of Development Priorities in a Key-value Store Serving Large-scale Applications
abstract
This article is an eight-year retrospective on development priorities for RocksDB, a key-value store developed at Facebook that targets large-scale distributed systems and that is optimized for Solid State Drives (SSDs). We describe how the priorities evolved over time as a result of hardware trends and extensive experiences running RocksDB at scale in production at a number of organizations: from optimizing write amplification, to space amplification, to CPU utilization. We describe lessons from running large-scale applications, including that resource allocation needs to be managed across different RocksDB instances, that data formats need to remain backward- and forward-compatible to allow incremental software rollouts, and that appropriate support for database replication and backups are needed. Lessons from failure handling taught us that data corruption errors needed to be detected earlier and that data integrity protection mechanisms are needed at every layer of the system. We describe improvements to the key-value interface. We describe a number of efforts that in retrospect proved to be misguided. Finally, we describe a number of open problems that could benefit from future research.
Siying Dong, Andrew Kryczka, Yanqin Jin, Michael Stumm
ACM Trans. Storage4
2019 An analysis of performance evolution of Linux's core operations
abstract
This paper presents an analysis of how Linux's performance has evolved over the past seven years. Unlike recent works that focus on OS performance in terms of scalability or service of a particular workload, this study goes back to basics: the latency of core kernel operations (e.g., system calls, context switching, etc.). To our surprise, the study shows that the performance of many core operations has worsened or fluctuated significantly over the years. For example, the select system call is 100% slower than it was just two years ago. An in-depth analysis shows that over the past seven years, core kernel subsystems have been forced to accommodate an increasing number of security enhancements and new features. These additions steadily add overhead to core kernel operations but also frequently introduce extreme slowdowns of more than 100%. In addition, simple misconfigurations have also severely impacted kernel performance. Overall, we find most of the slowdowns can be attributed to 11 changes.
Xiang Ren 0003, Kirk Rodrigues, Luyuan Chen, Juan Camilo Vega, Michael Stumm, Ding Yuan 0004
SOSP5
2019 The inflection point hypothesis: a principled debugging approach for locating the root cause of a failure
abstract
The end goal of failure diagnosis is to locate the root cause. Prior root cause localization approaches almost all rely on statistical analysis. This paper proposes taking a different approach based on the observation that if we model an execution as a totally ordered sequence of instructions, then the root cause can be identified by the first instruction where the failure execution deviates from the non-failure execution that has the longest instruction sequence prefix in common with that of the failure execution. Thus, root cause analysis is transformed into a principled search problem to identify the non-failure execution with the longest common prefix. We present Kairux, a tool that does just that. It is, in most cases, capable of pinpointing the root cause of a failure in a distributed system, in a fully automated way. Kairux uses tests from the system's rich unit test suite as building blocks to construct the non-failure execution that has the longest common prefix with the failure execution in order to locate the root cause. By evaluating Kairux on some of the most complex, real-world failures from HBase, HDFS, and ZooKeeper, we show that Kairux can accurately pinpoint each failure's respective root cause.
Yongle Zhang 0007, Kirk Rodrigues, Yu Luo 0006, Michael Stumm, Ding Yuan 0004
SOSP4
2018 Sharding the Shards: Managing Datastore Locality at Scale with Akkio
Muthukaruppan Annamalai, Kaushik Ravichandran 0003, Harish Srinivas, Igor Zinkovsky, Luning Pan, Tony Savor, David Nagle, Michael Stumm
OSDI8
2017 The Game of Twenty Questions: Do You Know Where to Log?
abstract
A production system's printed logs are often the only source of runtime information available for postmortem debugging, performance analysis and profiling, security auditing, and user behavior analytics. Therefore, the quality of this data is critically important. Recent work has attempted to enhance log quality by recording additional variable values, but logging statement placement, i.e., where to place a logging statement, which is the most challenging and fundamental problem for improving log quality, has not been adequately addressed so far. This position paper proposes we automate the placement of logging statements by measuring how much uncertainty, i.e., the expected number of possible execution code paths taken by the software, can be removed by adding a logging statement to a basic block. Guided by ideas from information theory, we describe a simple approach that automates logging statement placement. Preliminary results suggest that our algorithm can effectively cover, and further improve, the existing logging statement placements selected by developers. It can compute an optimal logging statement placement that disambiguates the entire function call path with only 0.218% of slowdown.
Xu Zhao 0004, Kirk Rodrigues, Yu Luo 0006, Michael Stumm, Ding Yuan 0004, Yuanyuan Zhou 0001
HotOS4
2017 The SEPO Model of Computation to Enable Larger-Than-Memory Hash Tables for GPU-Accelerated Big Data Analytics
abstract
The massive parallelism and high memory bandwidth of GPU's are particularly well matched with the exigencies of Big Data analytics applications, for which many independent computations and high data throughput are prevalent. These applications often produce (intermediary or final) results in the form of key-value (KV) pairs, and hash tables are particularly well-suited for storing these KV pairs in memory. How such hash tables are implemented on GPUs, however, has a large impact on performance. Unfortunately, all hash table solutions designed for GPUs to date have limitations that prevent acceleration for Big Data analytics applications. In this paper, we present the design and implementation of a GPU-based hash table for efficiently storing the KV pairs of Big Data analytics applications. The hash table is able to grow beyond the size of available GPU memory without excessive performance penalties. Central to our hash table design is the SEPO model of computation, where the processing of individual tasks is selectively postponed when processing is expected to be inefficient. A performance evaluation on seven GPU-based Big Data analytics applications, each processing several Gigabytes of input data, shows that our hash table allows the applications to achieve, on average, a speedup of 3.5 over their CPU-based multi-threaded implementations. This gain is realized despite having hash tables that grow up to four times larger than the size of available GPU memory.
Reza Mokhtari, Michael Stumm
IPDPS2
2017 Log20: Fully Automated Optimal Placement of Log Printing Statements under Specified Overhead Threshold
abstract
When systems fail in production environments, log data is often the only information available to programmers for postmortem debugging. Consequently, programmers' decision on where to place a log printing statement is of crucial importance, as it directly affects how effective and efficient postmortem debugging can be. This paper presents Log20, a tool that determines a near optimal placement of log printing statements under the constraint of adding less than a specified amount of performance overhead. Log20 does this in an automated way without any human involvement. Guided by information theory, the core of our algorithm measures how effective each log printing statement is in disambiguating code paths. To do so, it uses the frequencies of different execution paths that are collected from a production environment by a low-overhead tracing library. We evaluated Log20 on HDFS, HBase, Cassandra, and ZooKeeper, and observed that Log20 is substantially more efficient in code path disambiguation compared to the developers' manually placed log printing statements. Log20 can also output a curve showing the trade-off between the informativeness of the logs and the performance slowdown, so that a developer can choose the right balance.
Xu Zhao 0004, Kirk Rodrigues, Yu Luo 0006, Michael Stumm, Ding Yuan 0004, Yuanyuan Zhou 0001
SOSP4
2016 Social Hash: An Assignment Framework for Optimizing Distributed Systems Operations on Social Networks
Alon Shalita, Brian Karrer, Igor Kabiljo, Alessandro Presta, Aaron Adcock, Herald Kllapi, Michael Stumm
NSDI8
2016 Non-Intrusive Performance Profiling for Entire Software Stacks Based on the Flow Reconstruction Principle
Xu Zhao 0004, Kirk Rodrigues, Yu Luo 0006, Ding Yuan 0004, Michael Stumm
OSDI5
2016 Continuous deployment of mobile software at facebook (showcase)
abstract
Continuous deployment is the practice of releasing software updates to production as soon as it is ready, which is receiving increased adoption in industry. The frequency of updates of mobile software has traditionally lagged the state of practice for cloud-based services for a number of reasons. Mobile versions can only be released periodically. Users can choose when and if to upgrade, which means that several different releases coexist in production. There are hundreds of Android hardware variants, which increases the risk of having errors in the software being deployed.
Chuck Rossi, Elisa Shibley, Shi Su, Kent L. Beck, Tony Savor, Michael Stumm
SIGSOFT FSE6
2014 BigKernel - High Performance CPU-GPU Communication Pipelining for Big Data-Style Applications
abstract
GPUs offer an order of magnitude higher compute power and memory bandwidth than CPUs. GPUs therefore might appear to be well suited to accelerate computations that operate on voluminous data sets in independent ways, e.g., for transformations, filtering, aggregation, partitioning or other "Big Data" style processing. Yet experience indicates that it is difficult, and often error-prone, to write GPGPU programs which efficiently process data that does not fit in GPU memory, partly because of the intricacies of GPU hardware architecture and programming models, and partly because of the limited bandwidth available between GPUs and CPUs. In this paper, we propose Big Kernel, a scheme that provides pseudo-virtual memory to GPU applications and is implemented using a 4-stage pipeline with automated prefetching to (i) optimize CPU-GPU communication and (ii) optimize GPU memory accesses. Big Kernel simplifies the programming model by allowing programmers to write kernels using arbitrarily large data structures that can be partitioned into segments where each segment is operated on independently, these kernels are transformed into Big Kernel using straight-forward compiler transformations. Our evaluation on six data-intensive benchmarks shows that Big Kernel achieves an average speedup of 1.7 over state-of-the-art double-buffering techniques and an average speedup of 3.0 over corresponding multi-threaded CPU implementations.
Reza Mokhtari, Michael Stumm
IPDPS2
2014 User-Guided Device Driver Synthesis
Leonid Ryzhyk, John Keys, Alexander Legg, Arun Raghunath, Michael Stumm, Mona Vij
OSDI6
2014 Simple Testing Can Prevent Most Critical Failures: An Analysis of Production Failures in Distributed Data-Intensive Systems
Ding Yuan 0004, Yu Luo 0006, Xin Zhuang, Guilherme Renna Rodrigues, Xu Zhao 0004, Yongle Zhang 0007, Pranay Jain, Michael Stumm
OSDI8
2014 lprof: A Non-intrusive Request Flow Profiler for Distributed Systems
Xu Zhao 0004, Yongle Zhang 0007, David Lion, Muhammad Faizan Ullah, Yu Luo 0006, Ding Yuan 0004, Michael Stumm
OSDI7
2011 Exception-Less System Calls for Event-Driven Servers
Livio B. Soares, Michael Stumm
USENIX ATC2
2010 Otherworld: giving applications a chance to survive OS kernel crashes
abstract
The default behavior of all commodity operating systems today is to restart the system when a critical error is encountered in the kernel. This terminates all running applications with an attendant loss of "work in progress" that is nonpersistent.
Alex Depoutovitch, Michael Stumm
EuroSys2
2010 FlexSC: Flexible System Call Scheduling with Exception-Less System Calls
Livio B. Soares, Michael Stumm
OSDI2
2009 RapidMRC: approximating L2 miss rate curves on commodity systems for online optimizations
abstract
Miss rate curves (MRCs) are useful in a number of contexts. In our research, online L2 cache MRCs enable us to dynamically identify optimal cache sizes when cache-partitioning a shared-cache multicore processor. Obtaining L2 MRCs has generally been assumed to be expensive when done in software and consequently, their usage for online optimizations has been limited. To address these problems and opportunities, we have developed a low-overhead software technique to obtain L2 MRCs online on current processors, exploiting features available in their performance monitoring units so that no changes to the application source code or binaries are required. Our technique, called RapidMRC, requires a single probing period of roughly 221 million processor cycles (147 ms), and subsequently 124 million cycles (83 ms) to process the data. We demonstrate its accuracy by comparing the obtained MRCs to the actual L2 MRCs of 30 applications taken from SPECcpu2006, SPECcpu2000, and SPECjbb2000. We show that RapidMRC can be applied to sizing cache partitions, helping to achieve performance improvements of up to 27%.
David K. Tam, Livio B. Soares, Michael Stumm
ASPLOS4
2008 Reducing the harmful effects of last-level cache polluters with an OS-level, software-only pollute buffer
abstract
It is well recognized that LRU cache-line replacement can be ineffective for applications with large working sets or non-localized memory access patterns. Specifically, in last-level processor caches, LRU can cause cache pollution by inserting non-reuseable elements into the cache while evicting reusable ones. The work presented in this paper addresses last-level cache pollution through a dynamic operating system mechanism, called ROCS, requiring no change to underlying hardware and no change to applications. ROCS employs hardware performance counters on a commodity processor to characterize application cache behavior at run-time. Using this online profiling, cache unfriendly pages are dynamically mapped to a pollute buffer in the cache, eliminating competition between reusable and non-reusable cache lines. The operating system implements the pollute buffer through a page-coloring based technique, by dedicating a small slice of the last-level cache to store non-reusable pages. Measurements show that ROCS, implemented in the Linux 2.6.24 kernel and running on a 2.3 GHz PowerPC 970FX, improves performance of memory-intensive SPEC CPU 2000 and NAS benchmarks by up to 34%, and 16% on average.
Livio B. Soares, David K. Tam, Michael Stumm
MICRO3
2007 Thread clustering: sharing-aware scheduling on SMP-CMP-SMT multiprocessors
abstract
The major chip manufacturers have all introduced chip multiprocessing (CMP) and simultaneous multithreading (SMT) technology into their processing units. As a result, even low-end computing systems and game consoles have become shared memory multiprocessors with L1 and L2 cache sharing within a chip. Mid- and large-scale systems will have multiple processing chips and hence consist of an SMP-CMP-SMT configuration with non-uniform data sharing overheads. Current operating system schedulers are not aware of these new cache organizations, and as a result, distribute threads across processors in a way that causes many unnecessary, long-latency cross-chip cache accesses.
David K. Tam, Michael Stumm
EuroSys3
2007 Path: page access tracking to improve memory management
abstract
Traditionally, operating systems use a coarse approximation of memory accesses to implement memory management algorithms by monitoring page faults or scanning page table entries. With finer-grained memory access information, however, the operating system can manage memory muchmore effectively. Previous work has proposed the use of a software mechanism based on virtual page protection and soft faults to track page accesses at finer granularity. In this paper, we show that while this approach is effective for some applications, for many others it results in an unacceptably high overhead. We propose simple Page Access Tracking Hardware (PATH)to provide accurate page access information to the operating system. The suggested hardware support is generic andcan be used by various memory management algorithms. In this paper, we show how the information generated by PATH can be used to implement (i) adaptive page replacement policies, (ii) smart process memory allocation to improve performance or to provide isolation and better process prioritization, and (iii) effectively prefetch virtual memory pages when applications have non-trivial memory access patterns. Our simulation results show that these algorithms can dramatically improve performance (up to 500%) with PATH-provided information, especially when the system is under memory pressure. We show that the software overhead of processing PATH information is less than 6% acrossthe applications we examined (less than 3% in all but two applications), which is at least an order of magni.
Livio B. Soares, Michael Stumm, Thomas Walsh 0002, Angela Demke Brown
ISMM3
2007 Experience distributing objects in an SMMP OS
abstract
Designing and implementing system software so that it scales well on shared-memory multiprocessors (SMMPs) has proven to be surprisingly challenging. To improve scalability, most designers to date have focused on concurrency by iteratively eliminating the need for locks and reducing lock contention. However, our experience indicates that locality is just as, if not more, important and that focusing on locality ultimately leads to a more scalable system. In this paper, we describe a methodology and a framework for constructing system software structured for locality, exploiting techniques similar to those used in distributed systems. Specifically, we found two techniques to be effective in improving scalability of SMMP operating systems: (i) an object-oriented structure that minimizes sharing by providing a natural mapping from independent requests to independent code paths and data structures, and (ii) the selective partitioning, distribution, and replication of object implementations in order to improve locality. We describe concrete examples of distributed objects and our experience implementing them. We demonstrate that the distributed implementations improve the scalability of operating-system-intensive parallel workloads.
Jonathan Appavoo, Dilma Da Silva, Orran Krieger, Marc A. Auslander, Michal Ostrowski, Bryan S. Rosenburg, Amos Waterland, Robert W. Wisniewski, Jimi Xenidis, Michael Stumm, Livio B. Soares
ACM Trans. Comput. Syst.10
2005 Online performance analysis by statistical sampling of microprocessor performance counters
abstract
Hardware performance counters (HPCs) are increasingly being used to analyze performance and identify the causes of performance bottlenecks. However, HPCs are difficult to use for several reasons. Microprocessors do not provide enough counters to simultaneously monitor the many different types of events needed to form an over-all understanding of performance. Moreover, HPCs primarily count low-level micro-architectural events from which it is difficult to extract high-level insight required for identifying causes of performance problems.We describe two techniques that help overcome these difficulties, allowing HPCs to be used in dynamic real-time optimizers. First, statistical sampling is used to dynamically multiplex HPCs and make a larger set of logical HPCs available. Using real programs, we show experimentally that it is possible through this sampling to obtain counts of hardware events that are statistically similar (within 15%) to complete non-sampled counts, thus allowing us to provide a much larger set of logical HPCs. Second, we observe that stall cycles are a primary source of inefficiencies, and hence they should be major targets for software optimization. Based on this observation, we build a simple model in real-time that speculatively associates each stall cycle to a processor component that likely caused the stall. The information needed to produce this model is obtained using our HPC multiplexing facility to monitor a large number of hardware components simultaneously. Our analysis shows that even in an out-of-order superscalar micro-processor such a speculative approach yields a fairly accurate model with run-time overhead for collection and computation of under 2%.These results demonstrate that we can effective analyze on-line performance of application and system code running at full speed. The stall analysis shows where performance is being lost on a given processor.
Michael Stumm, Robert W. Wisniewski
ICS2
2005 Shared-buffer smoothing of variable bit-rate streams
Stergios V. Anastasiadis, Kenneth C. Sevcik, Michael Stumm
Perform. Evaluation3
2005 Scalable and fault-tolerant support for variable bit-rate data in the exedra streaming server
abstract
We describe the design and implementation of the Exedra continuous media server, and experimentally evaluate alternative resource management policies using a prototype system that we built. Exedra has been designed to provide scalable and efficient support for variable bit-rate media streams whose compression efficiency leads to reduced storage space and bandwidth requirements in comparison to constant bit-rate streams of equivalent quality. We examine alternative disk striping policies, and quantify the benefits of innovative techniques for storage space allocation, buffer management, and resource reservation, which we developed to achieve both predictability and high-performance in handling disk and network data transfers of variable size. Additionally, we investigate the differences between diverse data replication schemes over disk arrays, and compare methods for disk access time reservation that enable tolerance of disk failures at minimal cost. Overall, we demonstrate the feasibility of building network media servers that exploit the latest advances in media compression technology towards reducing the cost of wide-scale streaming services for stored data.
Stergios V. Anastasiadis, Kenneth C. Sevcik, Michael Stumm
ACM Trans. Storage3
2003 System Support for Online Reconfiguration
Craig A. N. Soules, Jonathan Appavoo, Kevin Hui, Robert W. Wisniewski, Dilma Da Silva, Gregory R. Ganger, Orran Krieger, Michael Stumm, Marc A. Auslander, Michal Ostrowski, Bryan S. Rosenburg, Jimi Xenidis
USENIX ATC, General Track8
2002 Maximizing Throughput in Replicated Disk Striping of Variable Bit-Rate Streams
Stergios V. Anastasiadis, Kenneth C. Sevcik, Michael Stumm
USENIX ATC, General Track3
2001 Supporting Hot-Swappable Components for System Software
abstract
Summary form only given. A hot-swappable component is one that can be replaced with a new or different implementation while the system is running and actively using the component. For example, a component of a TCP/IP protocol stack, when hot-swappable, can be replaced (perhaps to handle new denial-of-service attacks or improve performance), without disturbing existing network connections. The capability to swap components offers a number of potential advantages such as: online upgrades for high availability systems, improved performance due to dynamic adaptability and simplified software structures by allowing distinct policy and implementation options to be implemented in separate components (rather than as a single monolithic component) and dynamically swapped as needed. In order to hot-swap a component, it is necessary to (i) instantiate a replacement component; (ii) establish a quiescent state in which the component is temporarily idle; (iii) transfer state from the old component to the new component; (iv) swap the new component for the old; and (v) deallocate the old component.
Kevin Hui, Jonathan Appavoo, Robert W. Wisniewski, Marc A. Auslander, David Edelsohn, Benjamin Gamsa, Orran Krieger, Bryan S. Rosenburg, Michael Stumm
HotOS9
2001 Server-based smoothing of variable bit-rate streams
abstract
We introduce an algorithm that uses buffer space available at the server for smoothing disk transfers of variable bit-rate streams. Previous smoothing techniques prefetched stream data into the client buffer space, instead. However, emergence of personal computing devices with widely different hardware configurations means that we should not always assume abundance of resources at the client side. The new algorithm is shown to have optimal smoothing effect under the specified constraints. We incorporate it into a prototype server, and demonstrate significant increase in the number of streams concurrently supported at different system scales. We also extend our algorithm for striping variable bit-rate streams on heterogeneous disks. High bandwidth utilization is achieved across all the different disks, which leads to server throughput improved by several factors at high loads.
Stergios V. Anastasiadis, Kenneth C. Sevcik, Michael Stumm
ACM Multimedia3
2000 The NUMAchine Multiprocessor
abstract
Small-scale multiprocessors are becoming increasingly economical and common, whereas larger multiprocessors continue to have higher per-node costs. The NUMAchine multiprocessor project seeks to make large-scale multiprocessors more economical while maintaining high performance by exploring architectural and hardware features for low-cost, modular multiprocessors. To demonstrate our approach, we have implemented a prototype system that is scalable to 128 processors. An efficient directory-based cache coherence protocol exploits our hierarchical ring-based interconnect and supports sequential consistency. This paper documents the design choices and the resulting performance of the system using both simulation results and measurements on the prototype hardware.
R. Grindley, Tarek S. Abdelrahman, Stephen Brown 0003, S. Caranci, D. DeVries, Benjamin Gamsa, A. Grbic, M. Gusat, R. Ho, Orran Krieger, Guy Lemieux, K. Loveless, Naraig Manjikian, P. McHardy, Sinisa Srbljic, Michael Stumm, Zvonko G. Vranesic, Zeljko Zilic
ICPP16
1999 Tornado: Maximizing Locality and Concurrency in a Shared Memory Multiprocessor Operating System
Benjamin Gamsa, Orran Krieger, Jonathan Appavoo, Michael Stumm
OSDI4
1998 Design and Implementation of the NUMAchine Multiprocessor
abstract
This paper describes the design and implementation of the NUMAchine multiprocessor. As the market for CC-NUMA multiprocessors expands, this research project provides a timely architectural design and cost-effective prototype. The key to the successful implementation of our 48-processor prototype is the use of off-the-shelf components and programmable logic devices. Since this machine will serve as a research vehicle for parallel software development, a number of hardware features to enhance experimentation have been included in the design.
A. Grbic, Stephen Brown 0003, S. Caranci, R. Grindley, M. Gusat, Guy Lemieux, K. Loveless, Naraig Manjikian, Sinisa Srbljic, Michael Stumm, Zvonko G. Vranesic, Zeljko Zilic
DAC10
1998 On topology and bisection bandwidth of hierarchical-ring networks for shared-memory multiprocessors
abstract
Hierarchical-ring based multiprocessors are interesting alternatives to the more popular two-dimensional direct networks. They allow for simple router designs and wider communication paths than their direct network counterparts. There are several ways hierarchical-ring networks can be configured for a given number of processors. Feasible topologies range from tall, lean networks to short, wide networks, but only a few of these possess high throughput and low latency. We present the results of a simulation study: to determine how large hierarchical-ring networks can become before their performance deteriorates due to their bisection bandwidth constraints; and to derive topologies with high throughput and low latency for a given number of processors. We show that a system with a maximum of 120 processors and three levels of hierarchy can sustain most memory access behaviours, but that larger systems can be sustained, only if their bisection bandwidth is increased.
Govindan Ravindran, Michael Stumm
HiPC2
1998 Prioritized Multiprocessor Networks: Design and Performance
abstract
This paper proposes and evaluates prioritized direct shared-memory multiprocessor networks. We use three components to implement prioritized networks, namely, priority-based link arbitration, priority inheritance, and dynamic virtual channels. The two major results from our study are: (i) adding priorities to direct shared-memory multiprocessor networks can lead to reduced average transaction latencies and increased system throughput when running traditional parallel applications, and (ii) a prioritized multiprocessor network can be used to reduce the worst-case latencies of time-constrained traffic when it co-exists with best-effort traffic, without penalizing the average performance of best-effort traffic.
Govindan Ravindran, Michael Stumm
MASCOTS2
1997 A Performance Comparison of Hierarchical Ring- and Mesh-Connected Multiprocessor Networks
abstract
This paper compares the performance of hierarchical ring- and mesh-connected wormhole routed shared memory multiprocessor networks in a simulation study. Hierarchical rings are interesting alternatives to meshes since (i) they can be clocked at faster rates, (ii) they can have wider data paths and hence shorter message sates, (iii) they allow addition and removal of processing nodes at arbitrary locations, (iv) their topology allows natural exploitation in the spatial locality of application memory access patterns, and (v) their topology allows efficient implementation of broadcasts. Our study shows that for workloads with little locality, meshes scale better than ring networks because ring-based systems have limited bisection bandwidth. However, for workloads with some memory access locality hierarchical rings outperform meshes by 20-40% for system sizes of up to 128 processors. Even with poor access locality, hierarchical rings will outperform meshes for these system sizes if the mesh router buffers are only 1-flit large, and they will outperform meshes an systems with less than 36 processors regardless of mesh router buffer size.
Govindan Ravindran, Michael Stumm
HPCA2
1997 Linear and Extended Linear Transformations for Shared-Memory Multiprocessors
abstract
Advances in program transformation frameworks have significantly advanced compiler technology over the past few years. Program transformation frameworks provide mathematical abstractions of loop and data structures and formal methods for manipulating these structures. It is these frameworks that have allowed the development of algorithms capable of automatically tailoring an application for a target architecture. In this paper, we focus on the utility of these frameworks in improving the performance of mainly parallel applications on shared-memory multiprocessors. Data locality-oriented program optimizations are a key to good performance on shared-memory multiprocessors, since these optimizations can often improve performance by a factor of 10 or more. In this paper, we show the effectiveness of three key loop and data transformation frameworks in optimizing parallel programs on shared-memory multiprocessors. In particular, we describe our computation decomposition and alignment (CDA) framework, which can modify both the composition and the execution order of the re-composed iterations. We show how fine-grain transformations within the CDA framework enable new optimizations such as local optimizations that are otherwise achieved by global data transformations.
Dattatraya Kulkarni, Michael Stumm
Comput. J.2
1997 Analytical Prediction of Performance for Cache Coherence Protocols
abstract
In this paper, we introduce new analytical models for predicting the performance of parallel applications under various cache coherence protocol assumptions. The purpose of these models is to determine which protocols are to be used for which data blocks, and, in the case of dynamic protocols, also to determine when to change protocols. Although we focus on tightly-coupled multiprocessor systems, similar models can be derived for loosely-coupled distributed systems, such as networks of workstations. Our models are unique in that they lie between a large body of theoretical models that assume independence and a uniform distribution of memory accesses across processors, and a large body of address-trace oriented models that assume the availability of a precise characterization of interleaving behavior of memory accesses. The former are not very realistic, and the latter are not suitable for compile-time and run-time usage. In contrast, our models enable us to choose different input parameters depending on how the models will be used and depending on the needed accuracy in performance prediction. We present the models and show how the required parameters can be obtained. We assess the accuracy of our models on 15 parallel applications. For these applications, our most complete model predicts performance within a 10 percent margin when compared to a simulation of a sequentially consistent multiprocessor system. As part of this study, we also show the potential advantage of using dynamic hybrid protocols.
Sinisa Srbljic, Zvonko G. Vranesic, Michael Stumm, Leo Budin
IEEE Trans. Computers3
1997 HFS: A Performance-Oriented Flexible File System Based on Building-Block Compositions
abstract
The Hurricane File System (HFS) is designed for (potentially large-scale) shared-memory multiprocessors. Its architecture is based on the principle that, in order to maximize performance for applications with diverse requirements, a file system must support a wide variety of file structures, file system policies, and I/O interfaces. Files in HFS are implemented using simple building blocks composed in potentially complex ways. This approach yields great flexibility, allowing an application to customize the structure and policies of a file to exactly meet its requirements. As an extreme example, HFS allows a file's structure to be optimized for concurrent random-access write-only operations by 10 threads, something no other file system can do. Similarly, the prefetching, locking, and file cache management policies can all be chosen to match an application's access pattern. In contrast, most parallel file systems support a single file structure and a small set of policies. We have implemented HFS as part of the Hurricane operating system running on the Hector shared-memory multiprocessor. We demonstrate that the flexibility of HFS comes with little processing or I/O overhead. We also show that for a number of file access patterns, HFS is able to deliver to the applications the full I/O bandwidth of the disks on our system.
Orran Krieger, Michael Stumm
ACM Trans. Comput. Syst.2
1995 Implementing Flexible Computation Rules with Subexpression-level Loop Transformation
Dattatraya Kulkarni, Michael Stumm, Ronald C. Unrau
Euro-Par2
1995 On the Scalability of Demand-Driven Parallel Systems
Ronald C. Unrau, Michael Stumm, Orran Krieger
Euro-Par2
1995 Hierarchical Ring Topologies and the Effect of Their Bisection Bandwidth Constraints
Govindan Ravindran, Michael Stumm
ICPP (1)2
1995 Scalable cache consistency for hierarchically structured multiprocessors
Keith I. Farkas, Zvonko G. Vranesic, Michael Stumm
J. Supercomput.3
1995 Hierarchical clustering: A structure for scalable multiprocessor operating system design
Ronald C. Unrau, Orran Krieger, Benjamin Gamsa, Michael Stumm
J. Supercomput.4
1994 Optimizing IPC Performance for Shared-Memory Multiprocessors
abstract
We assert that in order to perform well, a shared-memory multiprocessor inter-process communication (IPC) facility must avoid a) accessing any shared data, and b) acquiring any locks. In addition, such a multiprocessor IPC facility must preserve the locality and concurrency of the applications themselves so that the high performance of the IPC facility can be fully exploited. In this paper we describe the design and implementation of a new shared-memory multiprocessor IPC facility that in the common case internally requires no accesses to shared data and no locking. In addition, the model of IPC we support and our implementation ensure that local resources are made available to the server to allow it to exploit any locality and concurrency available in the service. To the best of our knowledge, this is the first IPC subsystem with these attributes. The performance data we present demonstrates that the end-to- end performance of our multiprocessor IPC facility is competitive with the fastest uniprocessor IPC times.
Benjamin Gamsa, Orran Krieger, Michael Stumm
ICPP (1)3
1994 Experiences with Locking in a NUMA Multiprocessor Operating System Kernel
Ronald C. Unrau, Orran Krieger, Benjamin Gamsa, Michael Stumm
OSDI4
1994 Performance Evaluation of Hierarchical Ring-Based Shared Memory Multiprocessors
abstract
Investigates the performance of word-packet, slotted unidirectional ring-based hierarchical direct networks in the context of large-scale shared memory multiprocessors. Slotted unidirectional rings are attractive because their electrical characteristics and simple interfaces allow for fast cycle times and large bandwidths. For large-scale systems, it is necessary to use multiple rings for increased aggregate bandwidth. Hierarchies are attractive because the topology ensures unique paths between nodes, simple node interfaces and simple inter-ring connections. To ensure that a realistic region of the design space is examined, the architecture of the network used in the Hector prototype is adopted as the initial design point. A simulator of that architecture has been developed and validated with measurements from the prototype. The system and workload parameterization reflects conditions expected in the near future. The results of this study shows the importance of system balance on performance.>
Mark A. Holliday, Michael Stumm
IEEE Trans. Computers2
1993 A Fair Fast Scalable Reader-Writer Lock
abstract
A reader-writer (RW) lock allows either multiple readers to inspect shared data or a single writer exclusive access for modifying that data. On shared memory multiprocessors,cost of acquiring and releasing these locks can have a large impact on the performance of parallel applications. A major problem with naive implementations of these locks, where processors spin on a global lock variable waiting for the lock to become available, is that the memory containing the lock and the interconnection network to that memory will also become contended when the lock is contended.
Orran Krieger, Michael Stumm, Ronald C. Unrau, Jonathan Hanna
ICPP (2)2
1993 Locality and Loop Scheduling on NUMA Multiprocessors
abstract
An improtant issue in the parallel execution of loops is how to partition and schedule the loops onto the available processors. While most existing dynamic scheduling algorithms manage to load imbalance well, they fail to take locality into account and therefore perform poorly on parallel systems with non-uniform memory access times.
Sudarsan Tandri, Michael Stumm, Kenneth C. Sevcik
ICPP (2)3
1992 Cache Consistency in Hierarchical-Ring-Based Multiprocessors
abstract
A cache consistency scheme is presented for a class of multiprocessors based on a hierarchy of rings. By taking advantage of the natural broadcast and ordering properties of rings, cache consistency is achieved via a simple, selective-broadcast based protocol requiring no complex hardware. Using address-trace driven simulations of the Hector shared-memory multiprocessor, it is shown that the scheme performs well.
Keith I. Farkas, Zvonko G. Vranesic, Michael Stumm
SC3
1992 Heterogeneous Distributed Shared Memory
abstract
The design, implementation, and performance of heterogeneous distributed shared memory (HDSM) are studied. A prototype HDSM system that integrates very different types of hosts has been developed, and a number of applications of this system are reported. Experience shows that despite a number of difficulties in data conversion, HDSM is implementable with minimal loss in functional and performance transparency when compared to homogeneous DSM systems.>
Songnian Zhou, Michael Stumm, Kai Li 0001, David B. Wortman
IEEE Trans. Parallel Distributed Syst.2
1990 An Architecture for a Trusted Network
E. Stewart Lee, Brian W. Thomson, Peter I. P. Boulton, Michael Stumm
ESORICS4
1990 Using Deducibility in Secure Network Modelling
Brian W. Thomson, E. Stewart Lee, Peter I. P. Boulton, Michael Stumm, David M. Lewis
ESORICS4
1990 Extending Distributed Shared Memory to Heterogeneous Environments
abstract
The problems of building a distributed shared memory system on a network of heterogeneous machines are discussed. An existing algorithm (due to K. Li, 1986) that implements distributed shared memory is extended to a heterogeneous environment. An implementation that runs on Sun and DEC Firefly multiprocessor workstations connected by Ethernet is described. Related implementation and performance issues are discussed. On the basis of measurements of the applications ported to the system, it is concluded that heterogeneous distributed shared memory is not only feasible but can also be compared in performance to its homogeneous counterpart.>
Songnian Zhou, Michael Stumm, Tim McInerney
ICDCS2