VLDB 2026 Research / reviewers in the wild / expert
Kathryn Mohror
dblp:05/772 · also Kathryn M. Mohror
· DBLP profile ↗
56ranked-venue papers
6as first author
17since 2021 · last 2025
0000-0002-1366-1655ORCID · verified
Domains — the database's venue-derived domains; a paper can count in several
Systems, architecture and hardware · 47 · 6 first-author · 14 since 2021Graphics, computer vision, multimedia, augmented reality and games · 2 · 2 since 2021Applied, interdisciplinary, general and emerging computing · 1
| Year | Publication | Venue | Position |
|---|---|---|---|
| 2025 | WisIO: Automated I/O Bottleneck Detection with Multi-Perspective Views for HPC WorkflowsabstractWhy I/O Bottlenecks Matter in HPC• Modern HPC workloads (AI, simulations) involve massive data transfers that are crucial for enabling scientific discoveries• The large volume of these data transfers often lead to workloads spending significant amount of time performing I/O• Recent studies show that it is between 25-40% of total runtime • As a result, tuning the performance of data transfers via I/O analysis has become a routine task for application developers 6/27/2025 High-Level Execution Flow of WisIO 6/27/2025 WisIO 7 • Transform raw trace data into multi-perspective views • File, process, timeline, or user-defined High-Level Execution Flow of WisIO 6/27/2025 WisIO 8 • Transform raw trace data into multi-perspective views • File, process, timeline, or user-defined High-Level Execution Flow of WisIO 6/27/2025 WisIO 9 • Severity-based classification using I/O metrics • Quantifies how "bad" an I/O behavior is via a relative severity angle High-Level Execution Flow of WisIO 6/27/2025 WisIO 10 • Severity-based classification using I/O metrics • Quantifies how "bad" an I/O behavior is via a relative severity angle High-Level Execution Flow of WisIO 6/27/2025 WisIO 11 • Explains bottlenecks via rule-based reasoning • Identifies one or more causes (e.g., small reads, metadata overhead) High-Level Execution Flow of WisIO 6/27/2025 WisIO 12 • Explains bottlenecks via rule-based reasoning • Identifies one or more causes (e.g., small reads, metadata overhead) Implementation & API 6/27/2025 WisIO 13 Implemented in Python, for versions 3.8 and above • Parallel and distributed via Dask Works out-of-the-box with trace data from common I/O monitoring tools • Darshan, DFTracer, Recorder Two user-facing interfaces: • CLI: Installable via pip, highly configurable • Python API: Allows interactive analysis Multiple output types: Izzet Yildirim, Hariharan Devarajan, Antonios Kougkas, Xian-He Sun, Kathryn Mohror |
ICS | 5 |
| 2025 | VerifyIO: Verifying Adherence to Parallel I/O Consistency SemanticsabstractHigh-performance computing (HPC) applications generate and consume substantial amounts of data, typically managed by parallel file systems. These applications access file systems either through the POSIX interface or by using highlevel I/O libraries. While the POSIX consistency model remains dominant in HPC, emerging file systems and popular I/O libraries increasingly adopt alternative consistency models that relax semantics in various ways, creating significant challenges for correctness and portability. This paper addresses these challenges by proposing a trace-driven I/O consistency verification workflow, implemented in our open-source tool, VerifyIO, which collects execution traces, detects data conflicts, and verifies proper synchronization against specified consistency models. Our extensive evaluation of 91 test case executions across three widely used I/O libraries with four I/O consistency models reveals critical consistency issues at both application and implementation levels. Chen Wang 0004, Zhaobin Zhu, Kathryn Mohror, Sarah Neuwirth, Marc Snir |
IPDPS | 3 |
| 2025 | H5Intent: Autotuning HDF5 With User IntentabstractThe complexity of data management in HPC systems stems from the diversity in I/O behavior exhibited by new workloads, multistage workflows, and multitiered storage systems. The HDF5 library is a popular interface to interact with storage systems in HPC workloads. The library manages the complexity of diverse I/O behaviors by providing user-level configurations to optimize the I/O for HPC workloads. The HDF5 library exposes hundreds of configuration properties that can be set to alter how HDF5 manages I/O requests for better performance. However, determining which properties to set is quite challenging for users who lack expertise in HDF5 library internals. We propose a paradigm change through our H5Intent software, where users specify the intent of I/O operations and the software can set various HDF5 properties automatically to optimize the I/O behavior. This work demonstrates several use cases where mapping user-defined intents to HDF5 properties can be exploited to optimize I/O. In this study, we make three observations. First, I/O intents can accurately define HDF5 properties while managing conflicts between various properties and improving the I/O performance of microbenchmarks by up to 22×. Second, I/O intents can be efficiently passed to HDF5 with a small footprint of 6.74MB per node for thousands of intents per process. Third, an H5Intent VOL connector can dynamically map I/O intents to HDF5 properties for various I/O behaviors exhibited by our microbenchmark and improve I/O performance by up to 8.8×. Overall, H5Intent software improves the I/O performance of complex large-scale workloads we studied by up to 11×. Hariharan Devarajan, Gerd Heber, Kathryn Mohror |
IEEE Trans. Parallel Distributed Syst. | 3 |
| 2024 | ML-based Modeling to Predict I/O Performance on Different Storage Sub-systemsabstractParallel applications can spend a significant amount of time performing I/O on large-scale supercomputers. Fast near-compute storage accelerators called burst buffers can reduce the time a processor spends performing I/O and mitigate I/O bottlenecks. However, determining if a given application could be accelerated using burst buffers is not straightforward even for storage experts. The relationship between an application's I/O characteristics (such as I/O volume, processes involved, etc.) and the best storage sub-system for it can be complicated. As a result, adapting parallel applications to use burst buffers efficiently is a trial-and-error process. In this work, we present a Python-based tool called PrismIO that enables programmatic analysis of I/O traces. Using PrismIO, we identify performance bottlenecks when using burst buffers and parallel file systems, and explain why certain I/O patterns perform poorly. Further, we use machine learning to model the relationship between I/O characteristics and file system selections. We use IOR, an I/O benchmark to gather performance data for training the machine learning model. Our model can predict the better performing storage system for unseen IOR scenarios with an accuracy of 94.47% and for four real applications with an accuracy of 95.86%. Yiheng Xu, Pranav Sivaraman, Hariharan Devarajan, Kathryn Mohror, Abhinav Bhatele |
HiPC | 4 |
| 2024 | TailorFS: An Adaptive File System to Support Dynamic I/O requirements of HPC WorkloadsabstractHigh-Performance Computing (HPC) systems typically provide a global storage system to support large-scale workload’s I/O requirements. However, these workloads have diversified from traditional simulation to include big data analytics and AI workloads. HPC systems include hardware accelerators and storage software with configurations to support this workload diversification. However, selecting the correct combination of these configurations for each workload is a complex task even for I/O experts. To address this problem, we designed TailorFS, a software abstraction using FSView that transparently selects the appropriate file system characteristics including software abstractions, hardware accelerators, and optimizations to accelerate I/O for a given workload. Given a workload’s behavior, TailorFS chooses an appropriate software existing in the system and a corresponding configuration to dynamically optimize the workload. With TailorFS, users can utilize any interface to perform I/O on the file system and get the correct software selected for them based on the described behavior. Finally, TailorFS can accelerate I/O for different workload types by up to 46 × by utilizing workload characteristics. In conclusion, TailorFS can accelerate complex HPC workflows by up to 6.5 × better I/O performance on the Lassen supercomputer by using dynamically created workload-aware FSViews. Hariharan Devarajan, Kathryn Mohror |
SBAC-PAD | 2 |
| 2024 | DFTracer: An Analysis-Friendly Data Flow Tracer for AI-Driven WorkflowsabstractModern HPC workflows involve intricate coupling of simulation, data analytics, and artificial intelligence (AI) applications to improve time to scientific insight. These workflows require a cohesive set of performance analysis tools to provide a comprehensive understanding of data exchange patterns in HPC systems. However, current tools are not designed to work with an AI-based I/O software stack that requires tracing at multiple levels of the application. To this end, we developed a data flow tracer called DFTracer to capture data-centric events from workflows and the I/O stack to build a detailed understanding of the data exchange within AI-driven workflows. DFTracer has the following three novel features, including a unified interface to capture trace data from different layers in the software stack, a trace format that is analysis-friendly and optimized to support efficiently loading multi-million events in a few seconds, and the capability to tag events with workflow-specific context to perform domain-centric data flow analysis for workflows. Additionally, we demonstrate that DFTracer has a $1.44 x$ smaller runtime overhead and 1.3-7.1x smaller trace size than state-of-the-art tracing tools such as Score-P, Recorder, and Darshan. Moreover, with AI-driven workflows, Score-P, Recorder, and Darshan cannot find I/O accesses from dynamically spawned processes, and their load performance of 100 M events is three orders of magnitude slower than DFTracer. In conclusion, we demonstrate that DFTracer can capture multi-level performance data, including contextual event tagging with a low overhead of 1-5% from AI-driven workflows such as MuMMI and Microsoft’s Megatron Deepspeed running on large-scale HPC systems. Hariharan Devarajan, Loïc Pottier, Kaushik Velusamy, Huihuo Zheng, Izzet Yildirim, Olga Kogiou, Weikuan Yu, Antonios Kougkas, Xian-He Sun, Jae-Seung Yeom, Kathryn Mohror |
SC | 11 |
| 2024 | Formal Definitions and Performance Comparison of Consistency Models for Parallel File SystemsabstractThe semantics of HPC storage systems are defined by the consistency models to which they abide. Storage consistency models have been less studied than their counterparts in memory systems, with the exception of the POSIX standard and its strict consistency model. The use of POSIX consistency imposes a performance penalty that becomes more significant as the scale of parallel file systems increases and the access time to storage devices, such as node-local solid storage devices, decreases. While some efforts have been made to adopt relaxed storage consistency models, these models are often defined informally and ambiguously as by-products of a particular implementation. In this work, we establish a connection between memory consistency models and storage consistency models and revisit the key design choices of storage consistency models from a high-level perspective. Further, we propose a formal and unified framework for defining storage consistency models and a layered implementation that can be used to easily evaluate their relative performance for different I/O workloads. Finally, we conduct a comprehensive performance comparison of two relaxed consistency models on a range of commonly seen parallel I/O workloads, such as checkpoint/restart of scientific applications and random reads of deep learning applications. We demonstrate that for certain I/O scenarios, a weaker consistency model can significantly improve the I/O performance. For instance, in small random reads that are typically found in deep learning applications, session consistency achieved a 5x improvement in I/O bandwidth compared to commit consistency, even at small scales. Chen Wang 0004, Kathryn Mohror, Marc Snir |
IEEE Trans. Parallel Distributed Syst. | 2 |
| 2024 | A Visual Comparison of Silent Error PropagationabstractHigh-performance computing (HPC) systems play a critical role in facilitating scientific discoveries. Their scale and complexity (e.g., the number of computational units and software stack) continue to grow as new systems are expected to process increasingly more data and reduce computing time. However, with more processing elements, the probability that these systems will experience a random bit-flip error that corrupts a program's output also increases, which is often recognized as silent data corruption. Analyzing the resiliency of HPC applications in extreme-scale computing to silent data corruption is crucial but difficult. An HPC application often contains a large number of computation units that need to be tested, and error propagation caused by error corruption is complex and difficult to interpret. To accommodate this challenge, we propose an interactive visualization system that helps HPC researchers understand the resiliency of HPC applications and compare their error propagation. Our system models an application's error propagation to study a program's resiliency by constructing and visualizing its fault tolerance boundary. Coordinating with multiple interactive designs, our system enables domain experts to efficiently explore the complicated spatial and temporal correlation between error propagations. At the end, the system integrated a nonmonotonic error propagation analysis with an adjustable graph propagation visualization to help domain experts examine the details of error propagation and answer such questions as why an error is mitigated or amplified by program execution. Harshitha Menon, Kathryn Mohror, Shusen Liu 0001, Luanzheng Guo, Peer-Timo Bremer, Valerio Pascucci |
IEEE Trans. Vis. Comput. Graph. | 3 |
| 2023 | UnifyFS: A User-level Shared File System for Unified Access to Distributed Local StorageabstractWe introduce UnifyFS, a user-level file system that aggregates node-local storage tiers available on high performance computing (HPC) systems and makes them available to HPC applications under a unified namespace. UnifyFS employs transparent I/O interception, so it does not require changes to application code and is compatible with commonly used HPC I/O libraries. The design of UnifyFS supports the predominant HPC I/O workloads and is optimized for bulk-synchronous I/O patterns. Furthermore, UnifyFS provides customizable file system semantics to flexibly adapt its behavior for diverse I/O workloads and storage devices. In this paper, we discuss the unique design goals and architecture of UnifyFS and evaluate its performance on a leadership-class HPC system. In our experimental results, we demonstrate that UnifyFS exhibits excellent scaling performance for write operations and can improve the performance of application checkpoint operations by as much as 3× versus a tuned configuration. Michael J. Brim, Adam Moody, Seung-Hwan Lim, Ross G. Miller, Swen Böhm, Cameron Stanavige, Kathryn Mohror, Sarp Oral |
IPDPS | 7 |
| 2023 | Mimir: Extending I/O Interfaces to Express User Intent for Complex Workloads in HPCabstractThe complexity of data management in HPC systems stems from the diversity in I/O behavior exhibited by new workloads, multistage workflows, and the presence of multitiered storage systems. This complexity is managed by the storage systems, which provide user-level configurations to allow the tuning of workload I/O within the system. However, these configurations are difficult to set by users who lack expertise in I/O subsystems. We propose a paradigm change in which users specify the intent of I/O operations and storage systems automatically set various configurations based on the supplied intent. To this end, we developed the Mimir infrastructure to assist users in passing I/O intent to the underlying storage system. We demonstrate several use cases that map user-defined intents to storage configurations that lead to optimized I/O. In this study, we make three observations. First, I/O intents should be applied to each level of the I/O storage stack, from HDF5 to MPI-IO to POSIX, and integrated using lightweight adaptors in the existing stack. Second, the Mimir infrastructure supports up to 400M Ops/sec throughput of intents in the system, with a low memory overhead of 6.85KB per node. Third, intents assist in configuring a hierarchical cache to preload I/O, buffer in a node-local device, and store data in a global cache to optimize I/O workloads by 2.33×, 4×, and 2.1×, respectively. Our Mimir infrastructure optimizes complex large-scale workflows by up to 4× better I/O performance on the Lassen supercomputer by using automatically derived I/O intents. Hariharan Devarajan, Kathryn Mohror |
IPDPS | 2 |
| 2022 | Extracting and characterizing I/O behavior of HPC workloadsabstractSystem administrators set default storage-system configuration parameters with the goal of providing high per-formance for their system's I/O workloads. However, this gener-alized configuration can lead to suboptimal I/O performance for individual workloads. Users can provide parameter settings to the storage system to obtain better performance for individual applications, but it can be very challenging to determine which parameters to set and to what values. This problem is further ex-acerbated by the increased complexity of modern storage systems. In this work, we move towards solving this problem by providing a systematic categorization of workload-related information that users or middleware libraries can pass to the storage system to optimize I/O performance for specific workloads. We study applications and workflows from different scientific domains to cover a broad range of HPC use cases. Through our categorization, we find that a) workload features differ based on the hardware, software, and data components involved in the execution of workloads and b) multiple workload features together drive I/O optimizations. The methodology proposed in this work optimizes complex scientific workloads by 2.2 x−8 x, using workload-aware I/Ooptimizations. Using the proposed methodology, users can pragmatically characterize their workload, and this characterization can assist the storage system in configuring itself to optimize I/Operformance for individual workloads in HPC systems. Hariharan Devarajan, Kathryn Mohror |
CLUSTER | 2 |
| 2022 | DFMan: A Graph-based Optimization of Dataflow Scheduling on High-Performance Computing SystemsabstractScientific research and development campaigns are materialized by workflows of applications executing on high-performance computing (HPC) systems. These applications con-sist of tasks that can have inter- or intra-application flows of data to achieve the research goals successfully. These dataflows create dependencies among the tasks and cause resource con-tention on shared storage systems, thus limiting the aggregated I/O bandwidth achieved by the workflow. However, these I/O performance issues are often solved by tedious and manual efforts that demand holistic knowledge about the data dependencies in the workflow and the information about the infrastructure being utilized. Taking this into consideration, we design DFMan, a graph-based dataflow management and optimization framework for maximizing I/O bandwidth by leveraging the powerful storage stack on HPC systems to manage data sharing optimally among the tasks in the workflows. In particular, we devise a graph-based optimization algorithm that can leverage an intuitive graph representation of dataflow- and system-related information, and automatically carry out co-scheduling of task and data placement. According to our experiments, DFMan optimizes a wide variety of scientific workflows such as Hurricane 3D on Cloud Model 1 (CM1), Montage Carina Nebula (NGC3372), and an emulated dataflow kernel of the Multiscale Machine-learned Modeling Infrastructure (MuMMI I/O) on the Lassen supercomputer, and improves their aggregated I/O bandwidth by up to 5.42 x, 2.12 x and 1.29 x, respectively, compared to the baseline bandwidth. Fahim Chowdhury, Francesco Di Natale, Adam Moody, Kathryn Mohror, Weikuan Yu |
IPDPS | 4 |
| 2021 | O(1) Communication for Distributed SGD through Two-Level Gradient AveragingabstractLarge neural network models present a hefty communication challenge to distributed Stochastic Gradient Descent (SGD), with a per-iteration communication complexity of $\mathcal{O}(n)$ per worker for a model of n parameters. Many sparsification and quantization techniques have been proposed to compress the gradients, some reducing the per-iteration communication complexity to $\mathcal{O}(k)$, where $k\ll n$. In this paper, we introduce a strategy called two-level gradient averaging (A2SGD) to consolidate all gradients down to merely two local averages per worker before the computation of two global averages for an updated model. A2SGD also retains local errors to maintain the variance for fast convergence. Our analysis shows that A2SGD converges similar to the default distributed SGD algorithm. Our evaluation validates the conclusion and demonstrates that A2SGD significantly reduces the communication traffic per worker, and improves the overall training time of LSTM-PTB by $3.2\times$ and $23.2\times$, compared to Top-K and QSGD, respectively. We evaluate the effectiveness of our approach using two kinds of optimizers, SGD and Adam. Also, our evaluation with various communication options demonstrates the strength of our approach both in terms of communication reduction and convergence. To the best of our knowledge, A2SGD is the first to achieve $\mathcal{O}$ (1) communication complexity per worker without incurring a significant accuracy degradation of DNN models while communicating only two scalars representing gradients per worker for distributed SGD. Subhadeep Bhattacharya, Weikuan Yu, Fahim Chowdhury, Kathryn Mohror |
CLUSTER | 4 |
| 2021 | File System Semantics Requirements of HPC ApplicationsabstractMost widely-deployed parallel file systems (PFSs) implement POSIX semantics, which implies sequential consistency for reads and writes. Strict adherence to POSIX semantics is known to impede performance and thus several new PFSs with relaxed consistency semantics and better performance have been introduced. Such PFSs are useful provided that applications can run correctly on a PFS with weaker semantics. While it is widely assumed that HPC applications do not require strict POSIX semantics, to our knowledge there has not been systematic work to support this assumption. In this paper, we address this gap with a categorization of the consistency semantics guarantees of PFSs and develop an algorithm to determine the consistency semantics requirements of a variety of HPC applications. We captured the I/O activity of 17 representative HPC applications and benchmarks as they performed I/O through POSIX or I/O libraries and examined the metadata operations used and their file access patterns. From this analysis, we find that 16 of the 17 applications can utilize PFSs with weaker semantics. Chen Wang 0004, Kathryn Mohror, Marc Snir |
HPDC | 2 |
| 2021 | Understanding a program's resiliency through error propagationabstractAggressive technology scaling trends have worsened the transient fault problem in high-performance computing (HPC) systems. Some faults are benign, but others can lead to silent data corruption (SDC), which represents a serious problem; a fault introducing an error that is not readily detected nto an HPC simulation. Due to the insidious nature of SDCs, researchers have worked to understand their impact on applications. Previous studies have relied on expensive fault injection campaigns with uniform sampling to provide overall SDC rates, but this solution does not provide any feedback on the code regions without samples. Harshitha Menon, Kathryn Mohror, Peer-Timo Bremer, Yarden Livnat, Valerio Pascucci |
PPoPP | 3 |
| 2021 | Understanding the use of message passing interface in exascale proxy applicationsabstractSummary The Exascale Computing Project (ECP) focuses on the development of future exascale‐capable applications. Most ECP applications use the message passing interface (MPI) as their parallel programming model with mini‐apps serving as proxies. This paper explores the explicit usage of MPI in such ECP proxy applications. We empirically analyze 14 proxy applications from the ECP Proxy Apps Suite. We use the MPI profiling interface (PMPI) to collect MPI usage patterns in ECP proxy apps. Our analysis shows that a small subset of features from MPI is commonly used in the proxies of exascale‐capable applications, even when they reference third‐party libraries. This study is intended to provide a better understanding of the use of MPI in current exascale applications. The findings can help focus software investments made for exascale systems in the MPI middleware including optimization, fault‐tolerance, tuning, and hardware‐offload. Nawrin Sultana, Martin Ruefenacht, Anthony Skjellum, Purushotham V. Bangalore, Ignacio Laguna, Kathryn Mohror |
Concurr. Comput. Pract. Exp. | 6 |
| 2021 | SpotSDC: Revealing the Silent Data Corruption Propagation in High-Performance Computing SystemsabstractThe trend of rapid technology scaling is expected to make the hardware of high-performance computing (HPC) systems more susceptible to computational errors due to random bit flips. Some bit flips may cause a program to crash or have a minimal effect on the output, but others may lead to silent data corruption (SDC), i.e., undetected yet significant output errors. Classical fault injection analysis methods employ uniform sampling of random bit flips during program execution to derive a statistical resiliency profile. However, summarizing such fault injection result with sufficient detail is difficult, and understanding the behavior of the fault-corrupted program is still a challenge. In this article, we introduce SpotSDC, a visualization system to facilitate the analysis of a program's resilience to SDC. SpotSDC provides multiple perspectives at various levels of detail of the impact on the output relative to where in the source code the flipped bit occurs, which bit is flipped, and when during the execution it happens. SpotSDC also enables users to study the code protection and provide new insights to understand the behavior of a fault-injected program. Based on lessons learned, we demonstrate how what we found can improve the fault injection campaign method. Harshitha Menon, Dan Maljovec, Yarden Livnat, Shusen Liu 0001, Kathryn Mohror, Peer-Timo Bremer, Valerio Pascucci |
IEEE Trans. Vis. Comput. Graph. | 6 |
| 2020 | Understanding HPC Application I/O Behavior Using System Level StatisticsabstractThe processor performance of high performance computing (HPC) systems is increasing at a much higher rate than storage performance. This imbalance leads to I/O performance bottlenecks in massively parallel HPC applications. Therefore, there is a need for improvements in storage and file system designs to meet the ever-growing I/O needs of HPC applications. Storage and file system designers require a deep understanding of how HPC application I/O behavior affects current storage system installations in order to improve them. In this work, we contribute to this understanding using application-agnostic file system statistics gathered on compute nodes as well as metadata and object storage file system servers. We analyze file system statistics of more than 4 million jobs over a period of three years on two systems at Lawrence Livermore National Laboratory that include a 15 PiB Lustre file system for storage. The results of our study add to the state-of-the-art in I/O understanding by providing insight into how general HPC workloads affect the performance of large-scale storage systems. Some key observations in our study show that reads and writes are evenly distributed across the storage system; applications which perform I/O, spread that I/O across ~78% of the minutes of their runtime on average; less than 22% of HPC users who submit write-intensive jobs perform efficient writes to the file system; and I/O contention seriously impacts I/O performance. Arnab Kumar Paul, Olaf Faaland, Adam Moody, Elsa Gonsiorowski, Kathryn Mohror, Ali Raza Butt |
HiPC | 5 |
| 2020 | EReinit: Scalable and efficient fault-tolerance for bulk-synchronous MPI applicationsabstractSummary Scientists from many different fields have been developing Bulk‐Synchronous MPI applications to simulate and study a wide variety of scientific phenomena. Since failure rates are expected to increase in larger‐scale future HPC systems, providing efficient fault‐tolerance mechanisms for this class of applications is paramount. The global‐restart model has been proposed to decrease the time of failure recovery in Bulk‐Synchronous applications by allowing a fast reinitialization of MPI. However, the current implementations of this model have several drawbacks: they lack efficiency; their scalability have not been shown; and they require the use of the MPI profiling interface, which precludes the use of tools. In this paper, we present EReinit, an implementation of the global‐restart model that addresses these problems. Our key idea and optimization is the co‐design of basic fault‐tolerance mechanisms such as failure detection, notification, and recovery between MPI and the resource manager in contrast to current approaches on which these mechanisms are implemented in MPI only. We demonstrate EReinit in three HPC programs and show that it is up to four times more efficient than existing solutions at 4,096 processes. Sourav Chakraborty 0003, Ignacio Laguna, Murali Emani, Kathryn Mohror, Dhabaleswar K. Panda 0001, Martin Schulz 0001, Hari Subramoni |
Concurr. Comput. Pract. Exp. | 4 |
| 2020 | Ad Hoc File Systems for High-Performance Computing
André Brinkmann, Kathryn Mohror, Weikuan Yu, Philip H. Carns, Toni Cortes, Scott Klasky, Alberto Miranda, Franz-Josef Pfreundt, Robert B. Ross, Marc-Andre Vef |
J. Comput. Sci. Technol. | 2 |
| 2020 | QMPI: A next generation MPI profiling interface for modern HPC platforms
Bengisu Elis, Dai Yang, Olga Pearce, Kathryn Mohror, Martin Schulz 0001 |
Parallel Comput. | 4 |
| 2019 | Mitigating Inter-Job Interference via Process-Level Quality-of-ServiceabstractJobs on most high-performance computing (HPC) systems share the network with other concurrently executing jobs. This sharing creates contention that can severely degrade performance. We investigate the use of Quality of Service (QoS) mechanisms to reduce the negative impacts of network contention. Our results show that careful use of QoS reduces the impact of contention for specific jobs, resulting in up to a 27% performance improvement. In some cases the impact of contention is completely eliminated. These improvements are achieved with limited negative impact to other jobs; any job that experiences performance loss typically degrades less than 5%, often much less. Our approach can help ensure that HPC machines maintain high throughput as per-node compute power continues to increase faster than network bandwidth. Lee Savoie, David K. Lowenthal, Bronis R. de Supinski, Kathryn Mohror |
CLUSTER | 4 |
| 2019 | Efficient User-Level Storage Disaggregation for Deep LearningabstractOn large-scale high performance computing (HPC) systems, applications are provisioned with aggregated resources to meet their peak demands for brief periods. This results in resource underutilization because application requirements vary a lot during execution. This problem is particularly pronounced for deep learning applications that are running on leadership HPC systems with a large pool of burst buffers in the form of flash or non-volatile memory (NVM) devices. In this paper, we examine the I/O patterns of deep neural networks and reveal their critical need of loading many small samples randomly for successful training. We have designed a specialized Deep Learning File System (DLFS) that provides a thin set of APIs. Particularly, we design the metadata management of DLFS through an in-memory tree-based sample directory and its file services through the user-level SPDK protocol that can disaggregate the capabilities of NVM Express (NVMe) devices to parallel training tasks. Our experimental results show that DLFS can dramatically improve the throughput of training for deep neural networks on NVMe over Fabric, compared with the kernel-based Ext4 file system. Furthermore, DLFS achieves efficient user-level storage disaggregation with very little CPU utilization. Yue Zhu 0002, Weikuan Yu, Bing Jiao, Kathryn Mohror, Adam Moody, Fahim Chowdhury |
CLUSTER | 4 |
| 2019 | I/O Characterization and Performance Evaluation of BeeGFS for Deep LearningabstractParallel File Systems (PFSs) are frequently deployed on leadership High Performance Computing (HPC) systems to ensure efficient I/O, persistent storage and scalable performance. Emerging Deep Learning (DL) applications incur new I/O and storage requirements to HPC systems with batched input of small random files. This mandates PFSs to have commensurate features that can meet the needs of DL applications. BeeGFS is a recently emerging PFS that has grabbed the attention of the research and industry world because of its performance, scalability and ease of use. While emphasizing a systematic performance analysis of BeeGFS, in this paper, we present the architectural and system features of BeeGFS, and perform an experimental evaluation using cutting-edge I/O, Metadata and DL application benchmarks. Particularly, we have utilized AlexNet and ResNet-50 models for the classification of ImageNet dataset using the Livermore Big Artificial Neural Network Toolkit (LBANN), and ImageNet data reader pipeline atop TensorFlow and Horovod. Through extensive performance characterization of BeeGFS, our study provides a useful documentation on how to leverage BeeGFS for the emerging DL applications. Fahim Chowdhury, Yue Zhu 0002, Todd Heer, Saul Paredes, Adam Moody, Robin Goldstone, Kathryn Mohror, Weikuan Yu |
ICPP | 7 |
| 2019 | VeloC: Towards High Performance Adaptive Asynchronous Checkpointing at Large ScaleabstractGlobal checkpointing to external storage (e.g., a parallel file system) is a common I/O pattern of many HPC applications. However, given the limited I/O throughput of external storage, global checkpointing can often lead to I/O bottlenecks. To address this issue, a shift from synchronous checkpointing (i.e., blocking until writes have finished) to asynchronous checkpointing (i.e., writing to faster local storage and flushing to external storage in the background) is increasingly being adopted. However, with rising core count per node and heterogeneity of both local and external storage, it is non-trivial to design efficient asynchronous checkpointing mechanisms due to the complex interplay between high concurrency and I/O performance variability at both the node-local and global levels. This problem is not well understood but highly important for modern supercomputing infrastructures. This paper proposes a versatile asynchronous checkpointing solution that addresses this problem. To this end, we introduce a concurrency-optimized technique that combines performance modeling with lightweight monitoring to make informed decisions about what local storage devices to use in order to dynamically adapt to background flushes and reduce the checkpointing overhead. We illustrate this technique using the VeloC prototype. Extensive experiments on a pre-Exascale supercomputing system show significant benefits. Bogdan Nicolae, Adam Moody, Elsa Gonsiorowski, Kathryn Mohror, Franck Cappello |
IPDPS | 4 |
| 2019 | A large-scale study of MPI usage in open-source HPC applicationsabstractUnderstanding the state-of-the-practice in MPI usage is paramount for many aspects of supercomputing, including optimizing the communication of HPC applications and informing standardization bodies and HPC systems procurements regarding the most important MPI features. Unfortunately, no previous study has characterized the use of MPI on applications at a significant scale; previous surveys focus either on small data samples or on MPI jobs of specific HPC centers. This paper presents the first comprehensive study of MPI usage in applications. We survey more than one hundred distinct MPI programs covering a significantly large space of the population of MPI applications. We focus on understanding the characteristics of MPI usage with respect to the most used features, code complexity, and programming models and languages. Our study corroborates certain findings previously reported on smaller data samples and presents a number of interesting, previously un-reported insights. Ignacio Laguna, Ryan J. Marshall, Kathryn Mohror, Martin Ruefenacht, Anthony Skjellum, Nawrin Sultana |
SC | 3 |
| 2019 | The MPI_T events interface: An early evaluation and overview of the interface
Marc-André Hermanns, Nathan T. Hjelm, Michael Knobloch, Kathryn Mohror, Martin Schulz 0001 |
Parallel Comput. | 4 |
| 2019 | Failure recovery for bulk synchronous applications with MPI stages
Nawrin Sultana, Martin Ruefenacht, Anthony Skjellum, Ignacio Laguna, Kathryn Mohror |
Parallel Comput. | 5 |
| 2018 | Entropy-Aware I/O Pipelining for Large-Scale Deep Learning on HPC SystemsabstractDeep neural networks have recently gained tremendous interest due to their capabilities in a wide variety of application areas such as computer vision and speech recognition. Thus it is important to exploit the unprecedented power of leadership High-Performance Computing (HPC) systems for greater potential of deep learning. While much attention has been paid to leverage the latest processors and accelerators, I/O support also needs to keep up with the growth of computing power for deep neural networks. In this research, we introduce an entropy-aware I/O framework called DeepIO for large-scale deep learning on HPC systems. Its overarching goal is to coordinate the use of memory, communication, and I/O resources for efficient training of datasets. DeepIO features an I/O pipeline that utilizes several novel optimizations: RDMA (Remote Direct Memory Access)-assisted in-situ shuffling, input pipelining, and entropy-aware opportunistic ordering. In addition, we design a portable storage interface to support efficient I/O on any underlying storage system. We have implemented DeepIO as a prototype for the popular TensorFlow framework and evaluated it on a variety of different storage systems. Our evaluation shows that DeepIO delivers significantly better performance than existing memory-based storage systems. Yue Zhu 0002, Fahim Chowdhury, Huansong Fu, Adam Moody, Kathryn Mohror, Kento Sato, Weikuan Yu |
MASCOTS | 5 |
| 2018 | DisCVar: discovering critical variables using algorithmic differentiation for transient faultsabstractAggressive technology scaling trends have made the hardware of high performance computing (HPC) systems more susceptible to faults. Some of these faults can lead to silent data corruption (SDC), and represent a serious problem because they alter the HPC simulation results. In this paper, we present a full-coverage, systematic methodology called DisCVar to identify critical variables in HPC applications for protection against SDC. DisCVar uses automatic differentiation (AD) to determine the sensitivity of the simulation output to errors in program variables. We empirically validate our approach in identifying vulnerable variables by comparing the results against a full-coverage code-level fault injection campaign. We find that our DisCVar correctly identifies the variables that are critical to ensure application SDC resilience with a high degree of accuracy compared to the results of the fault injection campaign. Additionally, DisCVar requires only two executions of the target program to generate results, whereas in our experiments we needed to perform millions of executions to get the same information from a fault injection campaign. Harshitha Menon, Kathryn Mohror |
PPoPP | 2 |
| 2018 | Enabling callback-driven runtime introspection via MPI_TabstractUnderstanding the behavior of parallel applications that use the Message Passing Interface (MPI) is critical for optimizing communication performance. Performance tools for MPI currently rely on the PMPI Profiling Interface or the MPI Tools Information Interface, MPI_T, for portably collecting information for performance measurement and analysis. While tools using these interfaces have proven to be extremely valuable for performance tuning, these interfaces only provide synchronous information, i.e., when an MPI or an MPI_T function is called. There is currently no option for collecting information about asynchronous events from within the MPI library. In this work we propose a callback-driven interface for event notification from MPI implementations. Our approach is integrated in the existing MPI_T interface and provides a portable API for tools to discover and register for events of interest. We demonstrate the functionality and usability of the interface with a prototype implementation in Open MPI, a small logging tool (MEL) and the measurement infrastructure Score-P. Marc-André Hermanns, Nathan T. Hjelm, Michael Knobloch, Kathryn Mohror, Martin Schulz 0001 |
EuroMPI | 4 |
| 2018 | MPI Stages: Checkpointing MPI State for Bulk Synchronous ApplicationsabstractWhen an MPI program experiences a failure, the most common recovery approach is to restart all processes from a previous checkpoint and to re-queue the entire job. A disadvantage of this method is that, although the failure occurred within the main application loop, live processes must start again from the beginning of the program, along with new replacement processes---this incurs unnecessary overhead for live processes. To avoid such overheads and concomitant delays, we introduce the concept of "MPI Stages." MPI Stages saves internal MPI state in a separate checkpoint in conjunction with application state. Upon failure, both MPI and application state are recovered, respectively, from their last synchronous checkpoints and continue without restarting the overall MPI job. Live processes roll back only a few iterations within the main loop instead of rolling back to the beginning of the program, while a replacement of failed process restarts and reintegrates, thereby achieving faster failure recovery. This approach integrates well with large-scale, bulk synchronous applications and checkpoint/restart. Nawrin Sultana, Anthony Skjellum, Ignacio Laguna, Matthew Shane Farmer, Kathryn Mohror, Murali Emani |
EuroMPI | 5 |
| 2018 | ADAPT: algorithmic differentiation applied to floating-point precision tuning
Harshitha Menon, Michael O. Lam, Daniel Osei-Kuffuor, Markus Schordan, Kathryn Mohror, Jeffrey A. F. Hittinger |
SC | 6 |
| 2017 | MetaKV: A Key-Value Store for Metadata Management of Distributed Burst BuffersabstractDistributed burst buffers are a promising storage architecture for handling I/O workloads for exascale computing. Their aggregate storage bandwidth grows linearly with system node count. However, although scientific applications can achieve scalable write bandwidth by having each process write to its node-local burst buffer, metadata challenges remain formidable, especially for files shared across many processes. This is due to the need to track and organize file segments across the distributed burst buffers in a global index. Because this global index can be accessed concurrently by thousands or more processes in a scientific application, the scalability of metadata management is a severe performance-limiting factor. In this paper, we propose MetaKV: a key-value store that provides fast and scalable metadata management for HPC metadata workloads on distributed burst buffers. MetaKV complements the functionality of an existing key-value store with specialized metadata services that efficiently handle bursty and concurrent metadata workloads: compressed storage management, supervised block clustering, and log-ring based collective message reduction. Our experiments demonstrate that MetaKV outperforms the state-of-the-art key-value stores by a significant margin. It improves put and get metadata operations by as much as 2.66× and 6.29×, respectively, and the benefits of MetaKV increase with increasing metadata workload demand. Teng Wang 0001, Adam Moody, Yue Zhu 0002, Kathryn Mohror, Kento Sato, Tanzima Z. Islam, Weikuan Yu |
IPDPS | 4 |
| 2016 | Managing I/O Interference in a Shared Burst Buffer SystemabstractIn this work, we investigate the problem of inter-application interference in a shared Burst Buffer (BB) system. A BB is a new storage technology for HPC architectures that acts as an intermediate layer between performance-hungry HPC applications and the slow parallel file system. While the BB is meant to alleviate the problem of slow I/O in HPC systems, it is itself prone to performance degradation under interference. We observe that the magnitude of interference effects can reach a level that matters to the HPC system and the jobs that run on it. We investigate I/O scheduling techniques as a mechanism to mitigate BB I/O interference. With our results, we show that scheduling techniques tuned to BBs can control interference and significant performance benefits can be achieved. Sagar Thapaliya, Purushotham V. Bangalore, Jay F. Lofstead, Kathryn Mohror, Adam Moody |
ICPP | 4 |
| 2016 | I/O Aware Power ShiftingabstractPower limits on future high-performance computing (HPC) systems will constrain applications. However, HPC applications do not consume constant power over their lifetimes. Thus, applications assigned a fixed power bound may be forced to slow down during high-power computation phases, but may not consume their full power allocation during low-power I/O phases. This paper explores algorithms that leverage application semantics -- phase frequency, duration and power needs -- to shift unused power from applications in I/O phases to applications in computation phases, thus improving system-wide performance. We design novel techniques that include explicit staggering of applications to improve power shifting. Compared to executing without power shifting, our algorithms can improve average performance by up to 8% or improve performance of a single, high-priority application by up to 32%. Lee Savoie, David K. Lowenthal, Bronis R. de Supinski, Tanzima Z. Islam, Kathryn Mohror, Barry Rountree, Martin Schulz 0001 |
IPDPS | 5 |
| 2016 | Structural Clustering: A New Approach to Support Performance Analysis at ScaleabstractThe increasing complexity of high performance computing systems creates high demands on performance tools and human analysts due to an unmanageable volume of data gathered for performance analysis. A promising approach for reducing data volume is classification of data from multiple processes into groups of similar behavior to aid in analyzing application performance and identifying hot spots. However, existing approaches for structural and temporal classification of performance data suffer from lack of scalability or produce misleading results. To address this problem, we present a novel and effective structural similarity measure to efficiently classify data from parallel processes and introduce a method for efficient storage of the classified data. Using four examples, we show how existing performance analysis techniques benefit from our structural classification. Finally, we present a case study with 15 applications on up to 65,536 parallel processes that demonstrates the generality and scalability of our classification approach. Matthias Weber 0002, Ronny Brendel, Tobias Hilbrich, Kathryn Mohror, Martin Schulz 0001, Holger Brunst |
IPDPS | 4 |
| 2016 | MPI Sessions: Leveraging Runtime Infrastructure to Increase Scalability of Applications at ExascaleabstractMPI includes all processes in MPI_COMM_WORLD; this is untenable for reasons of scale, resiliency, and overhead. This paper offers a new approach, extending MPI with a new concept called Sessions, which makes two key contributions: a tighter integration with the underlying runtime system; and a scalable route to communication groups. This is a fundamental change in how we organise and address MPI processes that removes well-known scalability barriers by no longer requiring the global communicator MPI_COMM_WORLD. Daniel J. Holmes, Kathryn Mohror, Ryan E. Grant, Anthony Skjellum, Martin Schulz 0001, Wesley Bland, Jeffrey M. Squyres |
EuroMPI | 2 |
| 2016 | Allowing MPI tools builders to forget about FortranabstractC tool writers are forced to deal with a number of Fortran and C interoperability issues when intercepting MPI routines and completing them with PMPI. The C based tool has to intercept the Fortran MPI routines and marshal arguments between C and Fortran, which is not always easily done from C. Further, there is a subset of MPI routines that need to call PMPI from the original language they were called from, forcing the C tool to go back to a Fortran layer. Combined, these issues make writing tools that apply to C and Fortran applications both error-prone and time consuming. In this paper, we present WMPI, a wrapper generator that solves these issues by generating multiple lightweight wrappers to handle the marshalling, correct language specific reentry and other incompatibilities. Søren Rasmussen, Martin Schulz 0001, Kathryn Mohror |
EuroMPI | 3 |
| 2016 | An ephemeral burst-buffer file system for scientific applicationsabstractBurst buffers are becoming an indispensable hardware resource on large-scale supercomputers to buffer the bursty I/O from scientific applications. However, there is a lack of software support for burst buffers to be efficiently shared by applications within a batch-submitted job and recycled across different batch jobs. In addition, burst buffers need to cope with a variety of challenging I/O patterns from data-intensive scientific applications. In this study, we have designed an ephemeral Burst Buffer File System (BurstFS) that supports scalable and efficient aggregation of I/O bandwidth from burst buffers while having the same life cycle as a batch-submitted job. BurstFS features several techniques including scalable metadata indexing, co-located I/O delegation, and server-side read clustering and pipelining. Through extensive tuning and analysis, we have validated that BurstFS has accomplished our design objectives, with linear scalability in terms of aggregated I/O bandwidth for parallel writes and reads. Teng Wang 0001, Kathryn Mohror, Adam Moody, Kento Sato, Weikuan Yu |
SC | 2 |
| 2015 | The Role of Container Technology in Reproducible Computer Systems ResearchabstractEvaluating experimental results in the field of computer systems is a challenging task, mainly due to the many changes in software and hardware that computational environments go through. In this position paper, we analyze salient features of container technology that, if leveraged correctly, can help reduce the complexity of reproducing experiments in systems research. We present a use case in the area of distributed storage systems to illustrate the extensions that we envision, mainly in terms of container management infrastructure. We also discuss the benefits and limitations of using containers as a way of reproducing research in other areas of experimental systems research. Ivo Jimenez, Carlos Maltzahn, Adam Moody, Kathryn Mohror, Jay F. Lofstead, Remzi H. Arpaci-Dusseau, Andrea C. Arpaci-Dusseau |
IC2E | 4 |
| 2014 | A User-Level InfiniBand-Based File System and Checkpoint Strategy for Burst BuffersabstractCheckpoint/Restart is an indispensable fault tolerance technique commonly used by high-performance computing applications that run continuously for hours or days at a time. However, even with state-of-the-art checkpoint/restart techniques, high failure rates at large scale will limit application efficiency. To alleviate the problem, we consider using burst buffers. Burst buffers are dedicated storage resources positioned between the compute nodes and the parallel file system, and this new tier within the storage hierarchy fills the performance gap between node-local storage and parallel file systems. With burst buffers, an application can quickly store checkpoints with increased reliability. In this work, we explore how burst buffers can improve efficiency compared to using only node-local storage. To fully exploit the bandwidth of burst buffers, we develop a user-level Infini Band-based file system (IBIO). We also develop performance models for coordinated and uncoordinated checkpoint/restart strategies, and we apply those models to investigate the best checkpoint strategy using burst buffers on future large-scale systems. Kento Sato, Kathryn Mohror, Adam Moody, Todd Gamblin, Bronis R. de Supinski, Naoya Maruyama, Satoshi Matsuoka |
CCGRID | 2 |
| 2014 | FMI: Fault Tolerant Messaging Interface for Fast and Transparent RecoveryabstractFuture supercomputers built with more components will enable larger, higher-fidelity simulations, but at the cost of higher failure rates. Traditional approaches to mitigating failures, such as checkpoint/restart (C/R) to a parallel file system incur large overheads. On future, extreme-scale systems, it is unlikely that traditional C/R will recover a failed application before the next failure occurs. To address this problem, we present the Fault Tolerant Messaging Interface (FMI), which enables extremely low-latency recovery. FMI accomplishes this using a survivable communication runtime coupled with fast, in-memory C/R, and dynamic node allocation. FMI provides message-passing semantics similar to MPI, but applications written using FMI can run through failures. The FMI runtime software handles fault tolerance, including check pointing application state, restarting failed processes, and allocating additional nodes when needed. Our tests show that FMI runs with similar failure-free performance as MPI, but FMI incurs only a 28% overhead with a very high mean time between failures of 1 minute. Kento Sato, Adam Moody, Kathryn Mohror, Todd Gamblin, Bronis R. de Supinski, Naoya Maruyama, Satoshi Matsuoka |
IPDPS | 3 |
| 2014 | Detailed Modeling and Evaluation of a Scalable Multilevel Checkpointing SystemabstractHigh-performance computing (HPC) systems are growing more powerful by utilizing more components. As the system mean time before failure correspondingly drops, applications must checkpoint frequently to make progress. However, at scale, the cost of checkpointing becomes prohibitive. A solution to this problem is multilevel checkpointing, which employs multiple types of checkpoints in a single run. Lightweight checkpoints can handle the most common failure modes, while more expensive checkpoints can handle severe failures. We designed a multilevel checkpointing library, the Scalable Checkpoint/Restart (SCR) library, that writes lightweight checkpoints to node-local storage in addition to the parallel file system. We present probabilistic Markov models of SCR's performance. We show that on future large-scale systems, SCR can lead to a gain in machine efficiency of up to 35 percent, and reduce the load on the parallel file system by a factor of two. Additionally, we predict that checkpoint scavenging, or only writing checkpoints to the parallel file system on application termination, can reduce the load on the parallel file system by 20 × on today's systems and still maintain high application efficiency. Kathryn Mohror, Adam Moody, Greg Bronevetsky, Bronis R. de Supinski |
IEEE Trans. Parallel Distributed Syst. | 1 |
| 2013 | Alignment-Based Metrics for Trace Comparison
Matthias Weber 0002, Kathryn Mohror, Martin Schulz 0001, Bronis R. de Supinski, Holger Brunst, Wolfgang E. Nagel |
Euro-Par | 2 |
| 2013 | A 1 PB/s file system to checkpoint three million MPI tasks
Raghunath Rajachandrasekar, Adam Moody, Kathryn Mohror, Dhabaleswar K. Panda 0001 |
HPDC | 3 |
| 2013 | There goes the neighborhood: performance degradation due to nearby jobsabstractPredictable performance is important for understanding and alleviating application performance issues; quantifying the effects of source code, compiler, or system software changes; estimating the time required for batch jobs; and determining the allocation requests for proposals. Our experiments show that on a Cray XE system, the execution time of a communication-heavy parallel application ranges from 28% faster to 41% slower than the average observed performance. Blue Gene systems, on the other hand, demonstrate no noticeable run-to-run variability. In this paper, we focus on Cray machines and investigate potential causes for performance variability such as OS jitter, shape of the allocated partition, and interference from other jobs sharing the same network links. Reducing such variability could improve overall throughput at a computer center and save energy costs. Abhinav Bhatele, Kathryn Mohror, Steve H. Langer, Katherine E. Isaacs |
SC | 2 |
| 2012 | McrEngine: a scalable checkpointing system using data-aware aggregation and compressionabstractHigh performance computing (HPC) systems use checkpoint-restart to tolerate failures. Typically, applications store their states in checkpoints on a parallel file system (PFS). As applications scale up, checkpoint-restart incurs high overheads due to contention for PFS resources. The high overheads force large-scale applications to reduce checkpoint frequency, which means more compute time is lost in the event of failure. We alleviate this problem through a scalable checkpointrestart system, MCRENGINE. MCRENGINE aggregates checkpoints from multiple application processes with knowledge of the data semantics available through widely-used I/O libraries, e.g., HDF5 and netCDF, and compresses them. Our novel scheme improves compressibility of checkpoints up to 115% over simple concatenation and compression. Our evaluation with large-scale application checkpoints show that MCRENGINE reduces checkpointing overhead by up to 87% and restart overhead by up to 62% over a baseline with no aggregation or compression. Tanzima Z. Islam, Kathryn Mohror, Saurabh Bagchi, Adam Moody, Bronis R. de Supinski, Rudolf Eigenmann |
SC | 2 |
| 2012 | Design and modeling of a non-blocking checkpointing systemabstractAs the capability and component count of systems increase, the MTBF decreases. Typically, applications tolerate failures with checkpoint/restart to a parallel file system (PFS). While simple, this approach can suffer from contention for PFS resources. Multi-level checkpointing is a promising solution. However, while multi-level checkpointing is successful on today's machines, it is not expected to be sufficient for exascale class machines, which are predicted to have orders of magnitude larger memory sizes and failure rates. Our solution combines the benefits of non-blocking and multi-level checkpointing. In this paper, we present the design of our system and model its performance. Our experiments show that our system can improve efficiency by 1.1 to 2.0x on future machines. Additionally, applications using our checkpointing system can achieve high efficiency even when using a PFS with lower bandwidth. Kento Sato, Naoya Maruyama, Kathryn Mohror, Adam Moody, Todd Gamblin, Bronis R. de Supinski, Satoshi Matsuoka |
SC | 3 |
| 2012 | Trace profiling: Scalable event tracing on high-end parallel systems
Kathryn Mohror, Karen L. Karavanic |
Parallel Comput. | 1 |
| 2010 | Design, Modeling, and Evaluation of a Scalable Multi-level Checkpointing SystemabstractHigh-performance computing (HPC) systems are growing more powerful by utilizing more hardware components. As the system mean-time-before-failure correspondingly drops, applications must checkpoint more frequently to make progress. However, as the system memory sizes grow faster than the bandwidth to the parallel file system, the cost of checkpointing begins to dominate application run times. Multi-level checkpointing potentially solves this problem through multiple types of checkpoints with different costs and different levels of resiliency in a single run. This solution employs lightweight checkpoints to handle the most common failure modes and relies on more expensive checkpoints for less common, but more severe failures. This theoretically promising approach has not been fully evaluated in a large- scale, production system context. We have designed the Scalable Checkpoint/Restart (SCR) library, a multi-level checkpoint system that writes checkpoints to RAM, Flash, or disk on the compute nodes in addition to the parallel file system. We present the performance and reliability properties of SCR as well as a probabilistic Markov model that predicts its performance on current and future systems. We show that multi-level checkpointing improves efficiency on existing large-scale systems and that this benefit increases as the system size grows. In particular, we developed low-cost checkpoint schemes that are 100x-1000x faster than the parallel file system and effective against 85% of our system failures. This leads to a gain in machine efficiency of up to 35%, and it reduces the the load on the parallel file system by a factor of two on current and future systems. Adam Moody, Greg Bronevetsky, Kathryn Mohror, Bronis R. de Supinski |
SC | 3 |
| 2009 | Evaluating similarity-based trace reduction techniques for scalable performance analysisabstractEvent traces are required to correctly diagnose a number of performance problems that arise on today's highly parallel systems. Unfortunately, the collection of event traces can produce a large volume of data that is difficult, or even impossible, to store and analyze. One approach for compressing a trace is to identify repeating trace patterns and retain only one representative of each pattern. However, determining the similarity of sections of traces, i.e., identifying patterns, is not straightforward. In this paper, we investigate pattern-based methods for reducing traces that will be used for performance analysis. We evaluate the different methods against several criteria, including size reduction, introduced error, and retention of performance trends, using both benchmarks with carefully chosen performance behaviors, and a real application. Kathryn Mohror, Karen L. Karavanic |
SC | 1 |
| 2007 | Towards Scalable Event Tracing for High End Systems
Kathryn Mohror, Karen L. Karavanic |
HPCC | 1 |
| 2007 | A study of tracing overhead on a high-performance linux clusterabstractOur goal in this work was to identify and quantify the overheads of tracing parallel applications. We investigate several different sources of overhead related to tracing: trace instrumentation, periodic writing of trace files to disk, differing trace buffer sizes, system changes, and increasing numbers of processors in the target application. We encountered overheads as large as 26.7% for writing the trace file to disk. We found that buffer sizes can make a difference in the overheads, and that differences in system software can also contribute to the level of the perturbation. Our results show that the overhead of instrumentation correlates strongly with the number of events, while the overhead of writing the trace buffer increases with increasing numbers of processors. Kathryn Mohror, Karen L. Karavanic |
PPoPP | 1 |
| 2005 | Integrating Database Technology with Comparison-based Parallel Performance Diagnosis: The PerfTrack Performance Experiment Management ToolabstractPerfTrack is a data store and interface for managing performance data from large-scale parallel applications. Data collected in different locations and formats can be compared and viewed in a single performance analysis session. The underlying data store used in PerfTrack is implemented with a database management system (DBMS). PerfTrack includes interfaces to the data store and scripts for automatically collecting data describing each experiment, such as build and platform details. We have implemented a prototype of PerfTrack that can use Oracle or PostgreSQL for the data store. We demonstrate the prototype's functionality with three case studies: one is a comparative study of an ASC purple benchmark on high-end Linux and AIX platforms; the second is a parameter study conducted at Lawrence Livermore National Laboratory (LLNL) on two high end platforms, a 128 node cluster of IBM Power 4 processors and BlueGene/L; the third demonstrates incorporating performance data from the Paradyn Parallel Performance Tool into an existing PerfTrack data store. Karen L. Karavanic, John May, Kathryn Mohror, Brian Miller 0001, Kevin A. Huck, Rashawn L. Knapp, Brian Pugh |
SC | 3 |
| 2004 | Performance Tool Support for MPI-2 on LinuxabstractProgrammers of message-passing codes for clusters of workstations face a daunting challenge in understanding the performance bottlenecks of their applications. This is largely due to the vast amount of performance data that is collected, and the time and expertise necessary to use traditional parallel performance tools to analyze that data. This paper reports on our recent efforts developing a performance tool for MPI applications on Linux clusters. Our target MPI implementations were LAM/MPI and MPICH2, both of which support portions of the MPI-2 Standard. We started with an existing performance tool and added support for non-shared file systems, MPI-2 one-sided communications, dynamic process creation, and MPI Object naming. We present results using the enhanced version of the tool to examine the performance of several applications. We describe a new performance tool benchmark suite we have developed, PPerfMark, and present results for the benchmark using the enhanced tool. Kathryn Mohror, Karen L. Karavanic |
SC | 1 |