VLDB 2026 Research / reviewers in the wild / expert
William Gropp
dblp:g/WilliamGropp · also Bill Gropp, William D. Gropp
· DBLP profile ↗
124ranked-venue papers
21as first author
8since 2021 · last 2025
0000-0003-2905-3029ORCID · verified
Domains — the database's venue-derived domains; a paper can count in several
Systems, architecture and hardware · 98 · 16 first-author · 7 since 2021Computer networks · 3Software engineering, systems software and programming languages · 2Security and privacy · 1Applied, interdisciplinary, general and emerging computing · 1 · 1 since 2021
| Year | Publication | Venue | Position |
|---|---|---|---|
| 2025 | HiCCL: A Hierarchical Collective Communication LibraryabstractHiCCL (Hierarchical Collective Communication Library) addresses the growing complexity and diversity in highperformance network architectures. As GPU systems have evolved into networks of GPUs with different multilevel communication hierarchies, optimizing each collective function for a specific system has become a challenging task. Consequently, many collective libraries struggle to adapt to different hardware and software, especially across systems from different vendors. HiCCL's library design decouples the collective communication logic from network-specific optimizations through a compositional API. The communication logic is composed using multicast, reduction, and fence primitives, which are then factorized for a specified network hieararchy using only point-to-point operations within a level. Finally, striping and pipelining optimizations streamline execution. Performance evaluation of HiCCL across four different machines-two with Nvidia GPUs, one with AMD GPUs, and one with Intel GPUs—demonstrates an average$17 \times$higher throughput than the collectives of highly specialized GPU-aware MPI implementations, and competitive throughput with those of vendor-specific libraries (NCCL, RCCL and OneCCL), while providing portability across all four machines. Mert Hidayetoglu, Simon Garcia de Gonzalo, Elliott Slaughter, Pinku Surana, Wen-Mei W. Hwu, William Gropp, Alex Aiken |
IPDPS | 6 |
| 2024 | UniNet: Accelerating the Container Network Data Plane in IaaS CloudsabstractKubernetes ($K$8s) is a container orchestration plat-form for cloud-based IaaS environments. While it operates on either bare-metal servers or VMs, users prefer VMs for cost savings and agility reasons despite the added network overhead. This overhead, stemming from dual network tunneling at the VM and container levels, degrades performance. To address this, we present UniNet, a SmartNIC-based solution that offloads container-level network tunneling. We designed UniNet to be compatible with leading Container Network Interfaces (CNIs). This approach involves three key elements: (1) transforming VF-based NICs into a container network gateway, (2) offloading the critical path of the data plane functionalities to SmartNICs for enhanced performance and reduced latency, and (3) instituting an isolated control plane that separates VM- and container-level rule insertions, making it tenant-accessible. UniNet boosts CNI throughput by 7.08 x on average, cuts tail latency by 41.6 %, and reduces CPU usage by up to 5.6 x for the receiver and 4.02 x for the sender, respectively. William Gropp, Hubertus Franke, Bharat Sukhwani, Sameh W. Asaad, Jinjun Xiong, Volodymyr V. Kindratenko, Deming Chen |
CLOUD | 3 |
| 2024 | CommBench: Micro-Benchmarking Hierarchical Networks with Multi-GPU, Multi-NIC NodesabstractModern high-performance computing systems have multiple GPUs and network interface cards (NICs) per node. The resulting network architectures have multilevel hierarchies of subnetworks with different interconnect and software technologies. These systems offer multiple vendor-provided communication capabilities and library implementations (IPC, MPI, NCCL, RCCL, OneCCL) with APIs providing varying levels of performance across the different levels. Understanding this performance is currently difficult because of the wide range of architectures and programming models (CUDA, HIP, OneAPI). Mert Hidayetoglu, Simon Garcia de Gonzalo, Elliott Slaughter, Yu Li 0041, Christopher Zimmer 0001, Tekin Bicer, Bin Ren 0002, William Gropp, Wen-Mei W. Hwu, Alex Aiken |
ICS | 8 |
| 2024 | Quantum-centric supercomputing for materials science: A perspective on challenges and future directions
Yuri Alexeev, Maximilian Amsler, Marco Antonio Barroca, Sanzio Bassini, Torey Battelle, Daan Camps, David Casanova, Young Jay Choi, Fred Chong, Charles Chung, Christopher Codella, Antonio D. Córcoles, James Cruise, Alberto Di Meglio, Ivan Duran, Thomas Eckl, Sophia E. Economou, Stephan J. Eidenbenz, Bruce Elmegreen, Clyde Fare, Ismael Faro, Cristina Sanz Fernández, Rodrigo Neumann Barros Ferreira, Keisuke Fuji, Bryce Fuller, Laura Gagliardi, Giulia Galli, Jennifer R. Glick, Isacco Gobbi, Pranav Gokhale, Salvador de la Puente Gonzalez, Johannes Greiner, William Gropp, Michele Grossi, Emanuel Gull, Burns Healy, Matthew R. Hermes, Benchen Huang, Travis S. Humble, Nobuyasu Ito, Artur F. Izmaylov, Ali Javadi-Abhari, Douglas M. Jennewein, Shantenu Jha, Bert de Jong, Petar Jurcevic, William M. Kirby, Stefan Kister, Masahiro Kitagawa, Joel Klassen, Katherine Klymko, Kwangwon Koh, Masaaki Kondo, Doga Murat Kürkçüoglu, Krzysztof Kurowski, Teodoro Laino, Ryan Landfield, Matthew L. Leininger, Vicente Leyton-Ortega, Ang Li 0006, Meifeng Lin, Junyu Liu, Nicolás Lorente, André Luckow, Simon Martiel, Francisco Martín-Fernández, Margaret Martonosi, Claire Marvinney, Arcesio Castañeda Medina, Dirk Merten, Antonio Mezzacapo, Kristel Michielsen, Abhishek Mitra, Tushar Mittal, Kyungsun Moon, Joel Moore, Sarah Mostame, Mario Motta, Young-Hye Na, Yunseong Nam, Prineha Narang, Yu-ya Ohnishi, Daniele Ottaviani, Matthew Otten, Scott Pakin, Vincent R. Pascuzzi, Edwin Pednault, Tomasz Piontek, Jed W. Pitera, Patrick Rall, Gokul Subramanian Ravi, Niall Robertson, Matteo A. C. Rossi, Piotr Rydlichowski, Hoon Ryu, Georgy Samsonidze, Mitsuhisa Sato, Nishant Saurabh, Kunal Sharma, Soyoung Shin, George Slessman, Mathias Steiner, Iskandar Sitdikov, In-Saeng Suh, Eric D. Switzer, Joel Thompson, Synge Todo, Minh C. Tran, Dimitar Trenev, Christian Trott, Huan-Hsin Tseng, Norm M. Tubman, Esin Tureci, David García Valiñas, Sofia Vallecorsa, Christopher Wever, Konrad W. Wojciechowski, Xiaodi Wu 0001, Shinjae Yoo, Nobuyuki Yoshioka, Victor Wen-zhe Yu, Seiji Yunoki, Sergiy Zhuk, Dmitry Zubarev |
Future Gener. Comput. Syst. | 33 |
| 2023 | Fine-grained Policy-driven I/O Sharing for Burst BuffersabstractA burst buffer is a common method to bridge the performance gap between the I/O needs of modern supercomputing applications and the performance of the shared file system on large-scale supercomputers. However, existing I/O sharing methods require resource isolation, offline profiling, or repeated execution that significantly limit the utilization and applicability of these systems. Here we present ThemisIO, a policy-driven I/O sharing framework for a remote-shared burst buffer: a dedicated group of I/O nodes, each with a local storage device. ThemisIO preserves high utilization by implementing opportunity fairness so that it can reallocate unused I/O resources to other applications. ThemisIO accurately and efficiently allocates I/O cycles among applications, purely based on real-time I/O behavior without requiring user-supplied information or offline-profiled application characteristics. ThemisIO supports a variety of fair sharing policies, such as user-fair, size-fair, as well as composite policies, e.g., group-then-user-fair. All these features are enabled by its statistical token design. ThemisIO can alter the execution order of incoming I/O requests based on assigned tokens to precisely balance I/O cycles between applications via time slicing, thereby enforcing processing isolation. Experiments using I/O benchmarks show that ThemisIO sustains 13.5--13.7% higher I/O throughput and 19.5--40.4% lower performance variation than existing algorithms. For real applications, ThemisIO significantly reduces the slowdown by 59.1--99.8% caused by I/O interference. Ed Karrels, Lei Huang 0019, Yuhong Kan, Ishank Arora, Yinzhi Wang, Daniel S. Katz, William Gropp, Zhao Zhang 0007 |
SC | 7 |
| 2023 | Characterizing the performance of node-aware strategies for irregular point-to-point communication on heterogeneous architectures
Shelby Lockhart, Amanda Bienz, William Gropp, Luke N. Olson |
Parallel Comput. | 3 |
| 2022 | Exploring Spatial Indexing for Accelerated Feature Retrieval in HPCabstractDespite the critical role that range queries play in analysis and visualization for HPC applications, there has been no comprehensive analysis of indices that are designed to accelerate range queries and the extent to which they are viable in HPC. In this paper we present the first such evaluation, examining 20 open-source C and C++ libraries that support range queries. Contributions of this paper include answering the following questions: which of the implementations are viable in HPC, how do these libraries compare in terms of build time, query time, memory usage, and scalability, what are other trade-offs between these implementations, is there a single overall best solution, and when does a brute force solution offer the best performance? We also share key insights learned during this process that can assist both HPC application scientists and spatial index developers. While we find that there is no single best solution, three libraries, Boost, CGAL and R-tree, offer some of the best performance, scalability, memory overheads, and support for different mesh types. We find several areas where the spatial indices could be substantially improved: better performance when there are a large number of query matches, reduced memory overheads, and better support for GPUs or other accelerators. Margaret Lawson, William Gropp, Jay F. Lofstead |
CCGRID | 2 |
| 2022 | EMPRESS: Accelerating Scientific Discovery through Descriptive Metadata ManagementabstractHigh-performance computing scientists are producing unprecedented volumes of data that take a long time to load for analysis. However, many analyses only require loading in the data containing particular features of interest and scientists have many approaches for identifying these features. Therefore, if scientists store information (descriptive metadata) about these identified features, then for subsequent analyses they can use this information to only read in the data containing these features. This can greatly reduce the amount of data that scientists have to read in, thereby accelerating analysis. Despite the potential benefits of descriptive metadata management, no prior work has created a descriptive metadata system that can help scientists working with a wide range of applications and analyses to restrict their reads to data containing features of interest. In this article, we present EMPRESS, the first such solution. EMPRESS offers all of the features needed to help accelerate discovery: It can accelerate analysis by up to 300 ×, supports a wide range of applications and analyses, is high-performing, is highly scalable, and requires minimal storage space. In addition, EMPRESS offers features required for a production-oriented system: scalable metadata consistency techniques, flexible system configurations, fault tolerance as a service, and portability. Margaret Lawson, William Gropp, Jay F. Lofstead |
ACM Trans. Storage | 2 |
| 2020 | FFT, FMM, and multigrid on the road to exascale: Performance challenges and opportunities
Huda Ibeid, Luke N. Olson, William Gropp |
J. Parallel Distributed Comput. | 3 |
| 2019 | Locus: A System and a Language for Program OptimizationabstractWe discuss the design and the implementation of Locus, a system and a language to orchestrate the optimization of applications. The increasing complexity of machines and the large space of program variants, produced by the many transformations available, conspire to make compilers deliver unsatisfactory performance. As a result, optimization experts must intervene to manually explore the space of program variants seeking the best version for each target machine. This intervention is unproductive, and maintaining and managing sequences of transformations as new architectures are adopted and new application features are incorporated is challenging.Locus allows collections of program transformation sequences to be specified separately from the application code. The language is able to represent in a clear notation complex collections of transformations that are applied to code regions selected by the programmer. The system integrates multiple optimization modules as well as search modules that facilitate the efficient traversal of the space of program variants. Locus is intended to help experts in the optimization process, specially for complex, long-lived applications that are to be executed on different environments. Four examples are presented to illustrate the power and simplicity of the language. Although not the primary focus of this paper, the examples also show that exploring the space of variants typically leads to better performing codes than those produced by conventional compiler optimizations that are based on heuristics. Thiago S. F. X. Teixeira, Corinne Ancourt, David A. Padua, William Gropp |
CGO | 4 |
| 2019 | Using performance models to understand scalable Krylov solver performance at scale for structured grid problemsabstractKrylov solvers are key kernels in many large-scale science and engineering applications for solving sparse linear systems. Applications running at scale can experience significant slowdown due to factors such as network congestion, off-node congestion, network distance, and performance variation across processes. Performance models can help us better understand factors limiting performance, however simple models fail to capture slowdowns often occurring at scale and performance variation across multiple runs of the same code. This work develops performance models that capture behavior found at scale and uses these models to guide optimizations for Krylov solvers and related kernels using both blocking and non-blocking communication for structured grid problems at scale. We use detailed performance analysis with network performance counters to show how network behavior relates to observed performance and guide the development of performance models that capture the runtime impact of network congestion, network distance, communication and computation overlap, and process mappings. These models guide us to optimize kernels using MPI protocol changes, node-aware communication, and topology-aware communication. The resulting tools and analysis provide us with a better understanding of how to improve performance at scale that can benefit a wider range of applications. Paul R. Eller, Torsten Hoefler, William Gropp |
ICS | 3 |
| 2019 | Node aware sparse matrix-vector multiplication
Amanda Bienz, William Gropp, Luke N. Olson |
J. Parallel Distributed Comput. | 2 |
| 2019 | Using node and socket information to implement MPI Cartesian topologies
William Gropp |
Parallel Comput. | 1 |
| 2019 | Guest editor's introduction: Special issue on best papers from EuroMPI/USA 2017
William Gropp, Rajeev Thakur |
Parallel Comput. | 1 |
| 2018 | Improving Performance Models for Irregular Point-to-Point CommunicationabstractParallel applications are often unable to take full advantage of emerging parallel architectures due to scaling limitations, which arise due to inter-process communication. Performance models are used to analyze the sources of communication costs. However, traditional models for point-to-point communication fail to capture the full cost of many irregular operations, such as sparse matrix methods. In this paper, a node-aware based model is presented. Furthermore, the model is extended to include communication queue search time as well as an additional parameter estimating network contention. The resulting model is applied to a variety of irregular communication patterns throughout matrix operations, displaying improved accuracy over traditional models. Amanda Bienz, William Gropp, Luke N. Olson |
EuroMPI | 2 |
| 2018 | Using Node Information to Implement MPI Cartesian TopologiesabstractThe MPI API provides support for Cartesian process topologies, including the option to reorder the processes to achieve better communication performance. But MPI implementations rarely provide anything useful for the reorder option, typically ignoring it. One argument made is that modern interconnects are fast enough that applications are less sensitive to the exact layout of processes onto the system. However, intranode communication performance is much greater than internode communication performance. In this paper, we show a simple approach that takes into account only information about which MPI processes are on the same node to provide a fast and effective implementation of the MPI Cartesian topology. While not optimal, this approach provides a significant improvement over all tested MPI implementations and provides an implementation that may be used as the default in any MPI implementation of MPI_Cart_create. William Gropp |
EuroMPI | 1 |
| 2017 | A DSL for Performance OrchestrationabstractThe complexity and diversity of today's computer architectures are requiring more attention from the software developers in order to harness all the computing power available. Furthermore, each different modern architecture requires a potentially non-overlapping set of optimizations to attain a higher fraction of its nominal peak speed. This leads to challenges about performance portability and code maintainability, in particular, how to manage different optimized versions of the same code tailored to different architectures and how to keep them up to date as new algorithmic features are added. This increasing complexity of the architectures and the extension of the optimization space tends to make compilers deliver unsatisfactory performance, and the gap between the performance of hand-tuned and compiler-generated code has grown dramatically. Even the use of advanced optimization flags is not enough to narrow this gap. On the other hand, optimizing applications manually is very time-consuming, and the developer needs to understand and interact with many different hardware features for each architecture. Successful research has been developed to assist the programmer in this painful and error-prone process of implementing, optimizing and porting applications to different architectures. Nonetheless, the adoption of these works has been mostly restricted to specific domains, such as dense linear algebra, Fourier transforms, and signal processing. We have developed the framework ICE that decouples the performance expert role from the application expert role (separation of concerns). It allows the use of architecture-specific optimizations while keeping the code maintainable on the long term. It is responsible to orchestrate the use of multiple optimization tools to application's baseline version and perform an empirical search to find the best sequence of optimizations and their parameters. The baseline version is regarded as not having any architecture- or compiler-specific optimizations. The optimizations and the empirical search are directed by a domain-specific language (DSL) in an external file. Application's code are often dramatically altered by adding multiple optimization cases for each architecture used. This DSL allows the performance expert to apply optimizations without disarrange the original code. The DSL has constructs to expose the options of the optimizations and generates a search space that can be traversed by different search tools. For instance, it has conditional statements that can be used to specify which optimizations should be carried out for each compiler. The DSL is not only the input of the empirical search, but also the output. It can be used so save the best sequence of transformations found in previous searches. The application's code is annotated with unique identifiers that are referenced in the DSL. Currently, source-to-source loop optimizations, algorithm and pragmas selection are accepted. The framework interface is flexible to integrate new optimization and search tools. And in case of any failure it falls back to the baseline version. We have applied the framework to linear algebra problems, stencil computations and to a production code for the simulation of plasma-coupled combustion~xpacc achieving up to 3x speedup. Other works have tried to solve the problem of facilitating optimizing applications, but they lack of important features comprised by ICE. CHiLL, Orio, and X Language simplifies the generation of optimized code. CHiLL is the only one among these that the instructions to carry out the optimizations are given using an external file, but it references loops by their position on the source and modifications in the source require modifications in the external file, restricting its use in large production codes. Only Orio empirically evaluates variants of the annotated code. Summarizing, the contributions of the framework are: the separation of concerns, incremental adoption, a DSL to specify the optimization space, interface to plug-in and compare different optimization and search tools, combination of empirical search with expert knowledge. Thiago S. F. X. Teixeira, David A. Padua, William Gropp |
PACT | 3 |
| 2017 | Towards a More Complete Understanding of SDC PropagationabstractWith the rate of errors that can silently effect an application's state/output expected to increase on future HPC machines, numerous application-level detection and recovery schemes have been proposed. Recovery is more efficient when errors are contained and affect only part of the computation's state. Containment is usually achieved by verifying all information leaking out of a statically defined containment domain, which is an expensive procedure. Alternatively, error propagation can be analyzed to bound the domain that is affected by a detected error. This paper investigates how silent data corruption (SDC) due to soft errors propagates through three HPC applications: HPCCG, Jacobi, and CoMD. To allow for more detailed view of error propagation, the paper tracks propagation at the instruction and application variable level. The impact of detection latency on error propagation is shown along with an application's ability to recover. Finally, the impact of compiler optimizations are explored along with the impact of local problem size on error propagation. Jon Calhoun 0001, Marc Snir, Luke N. Olson, William Gropp |
HPDC | 4 |
| 2017 | Eliminating contention bottlenecks in multithreaded MPI
Hoang-Vu Dang, Marc Snir, William Gropp |
Parallel Comput. | 3 |
| 2016 | Scalability Challenges in Current MPI One-Sided ImplementationsabstractMPI one-sided or remote memory access (RMA) communication provides a different execution model from traditional two-sided or group communication and is better suited for some classes of applications. However, current implementations of MPI RMA are notorious for their inability to scale to large systems or problem sizes. In this paper, we present a study of the RMA infrastructure in popular open-source MPI implementations. Our objective is to identify critical scalability limitations with respect to memory usage in these implementations. We then perform a thorough evaluation on two cluster computers to demonstrate those scalability limitations, and we provide suggestions on how they can be alleviated. Pavan Balaji, William Gropp |
ISPDC | 3 |
| 2016 | Towards millions of communicating threadsabstractWe explore in this paper the advantages that accrue from avoiding the use of wildcards in MPI. We show that, with this change, one can efficiently support millions of concurrently communicating light-weight threads using send-receive communication. Hoang-Vu Dang, Marc Snir, William Gropp |
EuroMPI | 3 |
| 2016 | Modeling MPI Communication Performance on SMP Nodes: Is it Time to Retire the Ping Pong TestabstractThe "postal" model of communication [3, 8] T = α + βn, for sending n bytes of data between two processes with latency α and bandwidth 1/β, is perhaps the most commonly used communication performance model in parallel computing. This performance model is often used in developing and evaluating parallel algorithms in high-performance computing, and was an effective model when it was first proposed. Consequently, numerous tests of "ping pong" communication have been developed in order to measure these parameters in the model. However, with the advent of multicore nodes connected to a single (or a few) network interfaces, the model has become a poor match to modern hardware. In this paper, we show a simple three-parameter model that better captures the behavior of current parallel computing systems, and demonstrate its accuracy on several systems. In support of this model, which we call the max-rate model, we have developed an open source benchmark1 that can be used to determine the model parameters. William Gropp, Luke N. Olson, Philipp Samfass |
EuroMPI | 1 |
| 2016 | Scalable non-blocking preconditioned conjugate gradient methodsabstractThe preconditioned conjugate gradient method (PCG) is a popular method for solving linear systems at scale. PCG requires frequent blocking allreduce collective operations that can limit performance at scale. We investigate PCG variations designed to reduce communication costs by decreasing the number of allreduces and by overlapping communication with computation using a non-blocking allreduce. These variations include two methods we have developed, non-blocking PCG and 2-step pipelined PCG, and pipelined PCG from Ghysels and Vanroose. Performance modeling for communication and computation costs shows the expected performance of these methods. Weak and strong scaling experiments on up to 128k cores show that scalable PCG methods can outperform standard PCG at scale. We observe that the fastest method varies depending on the work per core, suggesting we need a suite of scalable solvers to obtain the best performance. Experiments with multiple preconditioners and linear systems show the robustness of these methods. Paul R. Eller, William Gropp |
SC | 2 |
| 2016 | An implementation and evaluation of the MPI 3.0 one-sided communication interfaceabstractSummary The Message Passing Interface (MPI) 3.0 standard includes a significant revision to MPI's remote memory access (RMA) interface, which provides support for one‐sided communication. MPI‐3 RMA is expected to greatly enhance the usability and performance of MPI RMA. We present the first complete implementation of MPI‐3 RMA and document implementation techniques and performance optimization opportunities enabled by the new interface. Our implementation targets messaging‐based networks and is publicly available in the latest release of the MPICH MPI implementation. Using this implementation, we explore the performance impact of new MPI‐3 functionality and semantics. Results indicate that the MPI‐3 RMA interface provides significant advantages over the MPI‐2 interface by enabling increased communication concurrency through relaxed semantics in the interface and additional routines that provide new window types, synchronization modes, and atomic operations. Copyright © 2016 John Wiley & Sons, Ltd. James Dinan, Pavan Balaji, Darius Buntinas, David Goodell, William Gropp, Rajeev Thakur |
Concurr. Comput. Pract. Exp. | 5 |
| 2015 | Runtime Support for Irregular Computation in MPI-Based ApplicationsabstractIn recent years more and more applications have been using irregular computation models in various domains such as bioinformatics and social network analysis. Traditional data movement approaches are not well suited for such applications because of the irregular communication patterns, sparse data structures, fast growth rate of data movement as system size or problem size rises, and so forth. Active Messages (AM) is an alternative programming paradigm that is more suitable for irregular computations. It allows small pieces of data to be dynamically moved to the remote process and certain computation to be triggered, and the remote process does not need to explicitly receive the data. In this paper, an outline of the first author's Ph.D. thesis, focusing on runtime support for irregular computation, is presented. In the first part, we combine the capability of AM with traditional MPI data movement patterns, and we propose a generalized MPI-interoperable AM framework (MPI-AM). In the second part, we extend the MPI-AM framework to provide a model of dynamic task parallelism for data-driven computation. In each part we describe critical issues, demonstrate the current status of the work and performance gain, and discuss remaining challenges to be solved. Pavan Balaji, William Gropp |
CCGRID | 3 |
| 2015 | A Multiplatform Study of I/O Behavior on Petascale SupercomputersabstractWe examine the I/O behavior of thousands of supercomputing applications "in the wild," by analyzing the Darshan logs of over a million jobs representing a combined total of six years of I/O behavior across three leading high-performance computing platforms. We mined these logs to analyze the I/O behavior of applications across all their runs on a platform; the evolution of an application's I/O behavior across time, and across platforms; and the I/O behavior of a platform's entire workload. Our analysis techniques can help developers and platform owners improve I/O performance and I/O system utilization, by quickly identifying underperforming applications and offering early intervention to save system resources. We summarize our observations regarding how jobs perform I/O and the throughput they attain in practice. Huong Luu 0002, Marianne Winslett, William Gropp, Robert B. Ross, Philip H. Carns, Kevin Harms, Prabhat, Surendra Byna, Yushu Yao |
HPDC | 3 |
| 2015 | DAME: A Runtime-Compiled Engine for Derived DatatypesabstractIn order to achieve high performance on modern and future machines, applications need to make effective use of the complex, hierarchical memory system. Writing performance-portable code continues to be challenging since each architecture has unique memory access characteristics. In addition, some optimization decisions can only reasonably be made at runtime. This suggests that a two-pronged approach to address the challenge is required. First, provide the programmer with a means to express memory operations declaratively which will allow a runtime system to transparently access the memory in the best way and second, exploit runtime information. MPI's derived datatypes accomplish the former although their performance in current MPI implementations shows scope for improvement. JIT-compilation can be used for the latter. In this work, we present DAME --- a language and interpreter that is used as the backend for MPI's derived datatypes. We also present DAME-L and DAME-X, two JIT-enabled implementations of DAME. All three implementations have been integrated into MPICH. We evaluate the performance of our implementations using DDTBench and two mini-applications written with MPI derived datatypes and obtain communication speedups of up to 20x and mini-application speedup of 3x. Tarun Prabhu, William Gropp |
EuroMPI | 2 |
| 2014 | Nonblocking Epochs in MPI One-Sided CommunicationabstractThe synchronization model of the MPI one-sided communication paradigm can lead to serialization and latency propagation. For instance, a process can propagate non-RMA communication-related latencies to remote peers waiting in their respective epoch-closing routines in matching epochs. In this work, we discuss six latency issues that were documented for MPI-2.0 and show how they evolved in MPI-3.0. Then, we propose entirely nonblocking RMA synchronizations that allow processes to avoid waiting even in epoch-closing routines. The proposal provides contention avoidance in communication patterns that require back to back RMA epochs. It also fixes the latency propagation issues. Moreover, it allows the MPI progress engine to orchestrate aggressive schedulings to cut down the overall completion time of sets of epochs without introducing memory consistency hazards. Our test results show noticeable performance improvements for a lower-upper matrix decomposition as well as an application pattern that performs massive atomic updates. Judicael A. Zounmevo, Pavan Balaji, William Gropp, Ahmad Afsahi |
SC | 4 |
| 2013 | Toward Asynchronous and MPI-Interoperable Active MessagesabstractMany new large-scale applications have emerged recently and become important in areas such as bioinformatics and social networks. These applications are often data-intensive and involve irregular communication patterns and complex operations on remote processes. Active messages have proven effective for parallelizing such nontraditional applications. However, most current active messages frameworks are low-level and system specific, do not efficiently support asynchronous progress, and are not interoperable with two-sided and collective communications. In this paper, we present the design and implementation of an active messages framework inside MPI to provide portability and programmability, and we explore challenges when asynchronously handling active messages and other messages from the network as well as from shared memory. We test our implementation with a set of comprehensive benchmarks. Evaluation results show that our framework has the advantages of overlapping and interoperability, while introducing only a modest overhead. Darius Buntinas, Judicael A. Zounmevo, James Dinan, David Goodell, Pavan Balaji, Rajeev Thakur, Ahmad Afsahi, William Gropp |
CCGRID | 9 |
| 2013 | Runtime system design of decoupled execution paradigm for data-intensive high-end computingabstractHigh performance computing are widely used for scientific discoveries by running scientific computation programs. Many of these applications are getting more and more data intensive [1]. They generate or access huge amount of data during some execution phases. However, traditional supercomputers are designed for computing-intensive tasks. They usually have highdensity clusters of processing cores and their storage systems are placed remotely and connected to the computing clusters with networks. This separation of the computing system and the storage system causes the data Input/Output performance bottleneck, especially for the data-intensive phases of HPC applications. This bottleneck degrades the HPC system's efficiency. Yanlong Yin, Hassan Eslami, Xian-He Sun, Yong Chen 0001, Rajeev Thakur, William Gropp |
CLUSTER | 8 |
| 2013 | Optimization Strategies for MPI-Interoperable Active MessagesabstractData-intensive applications, such as those in bioinformatics and social network analysis, differ from traditional scientific applications in that they often involve data-driven and irregular computation/communication patterns, making them ill-suited for traditional data movement approaches. Active Messages (AM) is an alternative programming model that allows dynamically moving computation closer to data, rather than moving the data to the local process. In our previous work, we proposed an MPI-interoperable AM framework that allows existing MPI applications to incrementally take advantage of AM capabilities. While that work presented a baseline implementation of how AMs semantically interact with the rest of the MPI infrastructure, it had several performance shortcomings. In this paper, we analyze these performance shortcomings and propose three optimization strategies: one implicitly derived by the MPI implementation and two explicitly hinted to by the application user. In addition to the detailed description of these optimization strategies, the paper presents a thorough performance evaluation on a 4096-core cluster that demonstrates considerable performance advantages from these strategies. Pavan Balaji, William Gropp, Rajeev Thakur |
DASC | 3 |
| 2013 | MPI-Interoperable Generalized Active MessagesabstractData-intensive applications have become increasingly important in recent years, yet traditional data movement approaches for scientific computation are not well suited for such applications. The Active Message (AM) model is an alternative communication paradigm that is better suited for such applications by allowing computation to be dynamically moved closer to data. Given the wide usage of MPI in scientific computing, enabling an MPI-interoperable AM paradigm would allow traditional applications to incrementally start utilizing AMs in portions of their applications, thus eliminating the programming effort of rewriting entire applications. In our previous work, we extended the MPI ACCUMULATE and MPI GET ACCUMULATE operations in the MPI standard to support AMs. However, the semantics of accumulate-style AMs are fundamentally restricted by the semantics of MPI ACCUMULATE and MPI GET ACCUMULATE, which were not designed to support the AM model. In this paper, we present a new generalized framework for MPI-interoperable AMs that can alleviate those restrictions, thus providing a richer semantics to accommodate a wide variety of application computational patterns. Together with a new API, we present a detailed description of the correctness semantics of this functionality and a reference implementation that demonstrates how various API choices affect the flexibility provided to the MPI implementation and consequently its performance. Pavan Balaji, William Gropp, Rajeev Thakur |
ICPADS | 3 |
| 2013 | Performance Analysis of the Lattice Boltzmann Model Beyond Navier-StokesabstractThe lattice Boltzmann method is increasingly important in facilitating large-scale fluid dynamics simulations. To date, these simulations have been built on discretized velocity models of up to 27 neighbors. Recent work has shown that higher order approximations of the continuum Boltzmann equation enable not only recovery of the Navier-Stokes hydrodynamics, but also simulations for a wider range of Knudsen numbers, which is especially important in micro- and nanoscale flows. These higher-order models have significant impact on both the communication and computational complexity of the application. We present a performance study of the higher-order models as compared to the traditional ones, on both the IBM Blue Gene/P and Blue Gene/Q architectures. We study the tradeoffs of many optimizations methods such as the use of deep halo level ghost cells that, alongside hybrid programming models, reduce the impact of extended models and enable efficient modeling of extreme regimes of computational fluid dynamics. Amanda Randles, Vivek Kale, Jeff R. Hammond, William Gropp, Efthimios Kaxiras |
IPDPS | 4 |
| 2013 | Analysis of topology-dependent MPI performance on Gemini networksabstractCurrent HPC systems utilize a variety of interconnection networks, with varying features and communication characteristics. MPI normalizes these interconnects with a common interface used by most HPC applications. However, network properties can have a significant impact on application performance. We explore the impact of the interconnect on application performance on the Blue Waters supercomputer. Blue Waters uses a three-dimensional, Cray Gemini torus network, which provides twice the Y-dimension bandwidth in the X and Z dimensions. Through several benchmarks, including a halo-exchange example, we demonstrate that application-level mapping to the network topology yields significant performance improvements. Antonio J. Peña, Ralf G. Correa Carvalho, James Dinan, Pavan Balaji, Rajeev Thakur, William Gropp |
EuroMPI | 6 |
| 2012 | A Decoupled Execution Paradigm for Data-Intensive High-End ComputingabstractHigh-end computing (HEC) applications in critical areas of science and technology tend to be more and more data intensive. I/O has become a vital performance bottleneck of modern HEC practice. Conventional HEC execution paradigms, however, are computing-centric for computation intensive applications. They are designed to utilize memory and CPU performance and have inherent limitations in addressing the critical I/O bottleneck issues of HEC. In this study, we propose a decoupled execution paradigm (DEP) to address the challenging I/O bottleneck issues. DEP is the first paradigm enabling users to identify and handle data-intensive operations separately. It can significantly reduce costly data movement and is better than the existing execution paradigms for data-intensive applications. The initial experimental tests have confirmed its promising potential. Its data-centric architecture could have an impact in future HEC systems, programming models, and algorithms design and development. Yong Chen 0001, Xian-He Sun, William Gropp, Rajeev Thakur |
CLUSTER | 4 |
| 2012 | Modeling the Performance of an Algebraic Multigrid Cycle Using Hybrid MPI/OpenMPabstractThe rise of multicore cluster architectures has led to intense interest in using a combination of MPI and OpenMP to more effectively program these machines. We present a performance model for hybrid implementation of the solve cycle of algebraic multigrid (AMG), a popular iterative solver for large sparse linear systems and a key component of many scientific simulations. We validate the model on two leading parallel platforms, and discuss implications for applications programmed in a hybrid model on future machines. Hormozd Gahvari, William Gropp, Kirk E. Jordan, Martin Schulz 0001, Ulrike Meier Yang |
ICPP | 2 |
| 2012 | Hybrid Static/dynamic Scheduling for Already Optimized Dense Matrix FactorizationabstractWe present the use of a hybrid static/dynamic scheduling strategy of the task dependency graph for direct methods used in dense numerical linear algebra. This strategy provides a balance of data locality, load balance, and low dequeue overhead. We show that the usage of this scheduling in communication avoiding dense factorization leads to significant performance gains. On a 48 core AMD Opteron NUMA machine, our experiments show that we can achieve up to 64% improvement over a version of CALU that uses fully dynamic scheduling, and up to 30% improvement over the version of CALU that uses fully static scheduling. On a 16-core Intel Xeon machine, our hybrid static/dynamic scheduling approach is up to 8% faster than the version of CALU that uses a fully static scheduling or fully dynamic scheduling. Our algorithm leads to speedups over the corresponding routines for computing LU factorization in well known libraries. On the 48 core AMD NUMA machine, our best implementation is up to 110% faster than MKL, while on the 16 core Intel Xeon machine, it is up to 82% faster than MKL. Our approach also shows significant speedups compared with PLASMA on both of these systems. Simplice Donfack, Laura Grigori, William Gropp, Vivek Kale |
IPDPS | 3 |
| 2012 | Faster topology-aware collective algorithms through non-minimal communicationabstractKnown algorithms for two important collective communication operations, allgather and reduce-scatter, are minimal-communication algorithms; no process sends or receives more than the minimum amount of data. This, combined with the data-ordering semantics of the operations, limits the flexibility and performance of these algorithms. Our novel non-minimal, topology-aware algorithms deliver far better performance with the addition of a very small amount of redundant communication. We develop novel algorithms for Clos networks and single or multi-ported torus networks. Tests on a 32k-node BlueGene/P result in allgather speedups of up to 6x and reduce-scatter speedups of over 11x compared to the native IBM algorithm. Broadcast, reduce, and allreduce can be composed of allgather or reduce-scatter and other collective operations; our techniques also improve the performance of these algorithms. Paul Sack, William Gropp |
PPoPP | 2 |
| 2012 | Efficient Multithreaded Context ID Allocation in MPI
James Dinan, David Goodell, William Gropp, Rajeev Thakur, Pavan Balaji |
EuroMPI | 3 |
| 2012 | MPI 3 and Beyond: Why MPI Is Successful and What Challenges It Faces
William Gropp |
EuroMPI | 1 |
| 2012 | Advanced MPI Including New MPI-3 Features
William Gropp, Ewing L. Lusk, Rajeev Thakur |
EuroMPI | 1 |
| 2012 | Leveraging MPI's One-Sided Communication Interface for Shared-Memory Programming
Torsten Hoefler, James Dinan, Darius Buntinas, Pavan Balaji, Brian W. Barrett, Ron Brightwell, William Gropp, Vivek Kale, Rajeev Thakur |
EuroMPI | 7 |
| 2012 | Adaptive Strategy for One-Sided Communication in MPICH2
Gopalakrishnan Santhanaraman, William Gropp |
EuroMPI | 3 |
| 2011 | Weighted locality-sensitive scheduling for mitigating noise on multi-core clustersabstractRecent studies have shown that operating system (OS) interference, popularly called OS noise can be a significant problem as we scale to a large number of processors. One solution for mitigating noise is to turn off certain OS services on the machine. However, this is typically infeasible because full-scale OS services may be required for some applications. Furthermore, it is not a choice that an end user can make. Thus, we need an application-level solution. Building upon previous work that demonstrated the utility of within-node light-weight load balancing, we discuss the technique of weighted micro-scheduling and provide insights based on experimentation for two different machines with very different noise signatures. Through careful enumeration of the search space of scheduler parameters, we allow our weighted micro-scheduler to be dynamic, adaptive and tunable for a specific application running on a specific architecture. By doing this, we show how we can enable running scientific applications efficiently on a very large number of processors, even in the presence of noise. Vivek Kale, Abhinav Bhatele, William Gropp |
HiPC | 3 |
| 2011 | Modeling the performance of an algebraic multigrid cycle on HPC platformsabstractNow that the performance of individual cores has plateaued, future supercomputers will depend upon increasing parallelism for performance. Processor counts are now in the hundreds of thousands for the largest machines and will soon be in the millions. There is an urgent need to model application performance at these scales and to understand what changes need to be made to ensure continued scalability. This paper considers algebraic multigrid (AMG), a popular and highly efficient iterative solver for large sparse linear systems that is used in many applications. We discuss the challenges for AMG on current parallel computers and future exascale architectures, and we present a performance model for an AMG solve cycle as well as performance measurements on several massively-parallel platforms. Hormozd Gahvari, Allison H. Baker, Martin Schulz 0001, Ulrike Meier Yang, Kirk E. Jordan, William Gropp |
ICS | 6 |
| 2011 | Performance modeling as the key to extreme scale computingabstractParallel computing is primarily about achieving greater performance than is possible without using parallelism. Especially for the high-end, where systems cost tens to hundreds of millions of dollars, making the best use of these valuable and scarce systems is important. Yet few applications really understand how well they are performing with respect to the achievable performance on the system. The Blue Waters system, currently being installed at the University of Illinois, will offer sustained performance in excess of 1 PetaFLOPS for many applications. However, achieving this level of performance requires careful attention to many details, as this system has many features that must be used to get the best performance. To address this problem, the Blue Waters project is exploring the use of performance models that provide enough information to guide the development and tuning of applications, ranging from improving the performance of small loops to identifying the need for new algorithms. Using Blue Waters as an example of an extreme scale system, this talk will describe some of the challenges faced by applications at this scale, the role that performance modeling can play in preparing applications for extreme scale, and some ways in which performance modeling has guided performance enhancements for those applications. William Gropp |
ICS | 1 |
| 2011 | Architectural Constraints to Attain 1 Exaflop/s for Three Scientific Application ClassesabstractThe first Teraflop/s computer, the ASCI Red, became operational in 1997, and it took more than 11 years for a Petaflop/s performance machine, the IBM Roadrunner, to appear on the Top500 list. Efforts have begun to study the hardware and software challenges for building an exascale machine. It is important to understand and meet these challenges in order to attain Exaflop/s performance. This paper presents a feasibility study of three important application classes to formulate the constraints that these classes will impose on the machine architecture for achieving a sustained performance of 1 Exaflop/s. The application classes being considered in this paper are -- classical molecular dynamics, cosmological simulations and unstructured grid computations (finite element solvers). We analyze the problem sizes required for representative algorithms in each class to achieve 1 Exaflop/s and the hardware requirements in terms of the network and memory. Based on the analysis for achieving an Exaflop/s, we also discuss the performance of these algorithms for much smaller problem sizes. Abhinav Bhatele, Pritish Jetley, Hormozd Gahvari, Lukasz Wesolowski, William Gropp, Laxmikant V. Kalé |
IPDPS | 5 |
| 2011 | LACIO: A New Collective I/O Strategy for Parallel I/O SystemsabstractParallel applications benefit considerably from the rapid advance of processor architectures and the available massive computational capability, but their performance suffers from large latency of I/O accesses. The poor I/O performance has been attributed as a critical cause of the low sustained performance of parallel systems. Collective I/O is widely considered a critical solution that exploits the correlation among I/O accesses from multiple processes of a parallel application and optimizes the I/O performance. However, the conventional collective I/O strategy makes the optimization decision based on the logical file layout to avoid multiple file system calls and does not take the physical data layout into consideration. On the other hand, the physical data layout in fact decides the actual I/O access locality and concurrency. In this study, we propose a new collective I/O strategy that is aware of the underlying physical data layout. We confirm that the new Layout-Aware Collective I/O (LACIO) improves the performance of current parallel I/O systems effectively with the help of noncontiguous file system calls. It holds promise in improving the I/O performance for parallel systems. Yong Chen 0001, Xian-He Sun, Rajeev Thakur, Philip C. Roth, William Gropp |
IPDPS | 5 |
| 2011 | Scalable Memory Use in MPI: A Case Study with MPICH2
David Goodell, William Gropp, Rajeev Thakur |
EuroMPI | 2 |
| 2011 | Performance Expectations and Guidelines for MPI Derived Datatypes
William Gropp, Torsten Hoefler, Rajeev Thakur, Jesper Larsson Träff |
EuroMPI | 1 |
| 2011 | Multi-core and Network Aware MPI Topology Functions
Mohammad J. Rashti, Jonathan Green, Pavan Balaji, Ahmad Afsahi, William Gropp |
EuroMPI | 5 |
| 2011 | Avoiding hot-spots on two-level direct networksabstractA low-diameter, fast interconnection network is going to be a prerequisite for building exascale machines. A two-level direct network has been proposed by several groups as a scalable design for future machines. IBM's PERCS topology and the dragonfly network discussed in the DARPA exascale hardware study are examples of this design. The presence of multiple levels in this design leads to hot-spots on a few links when processes are grouped together at the lowest level to minimize total communication volume. This is especially true for communication graphs with a small number of neighbors per task. Routing and mapping choices can impact the communication performance of parallel applications running on a machine with a two-level direct topology. This paper explores intelligent topology aware mappings of different communication patterns to the physical topology to identify cases that minimize link utilization. We also analyze the trade-offs between using direct and indirect routing with different mappings. We use simulations to study communication and overall performance of applications since there are no installations of two-level direct networks yet. This study raises interesting issues regarding the choice of job scheduling, routing and mapping for future machines. Abhinav Bhatele, William Gropp, Laxmikant V. Kalé |
SC | 3 |
| 2010 | Enabling the Next Generation of Scalable ClustersabstractSummary form only given. Clusters revolutionized computing by making supercomputer capabilities widely available. But one of the main drivers of that revolution, the rapid doubling of processor clock rates, ran out of steam several years ago. To maintain (or even increase) the historic rate of improvement in computing power, processor designs are rapidly increasing parallelism at all levels, including more functional units, more cores, and ways to share resources among threads. Heterogeneous designs that use more specialized processors such as GPGPUs are becoming common. The scale of high-end systems is also getting larger, with 1000-core systems becoming commonplace and systems with over 300,000 cores planned for 2011. However, the software and algorithms for these systems are still basically the same as when the cluster revolution began. Drawing on experiences with the sustained PetaFLOPS system, called Blue Waters, to be installed at Illinois in 2011, and with exploratory work into Exascale system designs, this talk will discuss some of the challenges facing the cluster community as scalability becomes increasingly important and reviews some of the developments in algorithms, programming models, and software frameworks that must complement the evolution of cluster hardware. William Gropp |
CCGRID | 1 |
| 2010 | Minimizing MPI Resource Contention in Multithreaded Multicore EnvironmentsabstractWith the ever-increasing numbers of cores per node in high-performance computing systems, a growing number of applications are using threads to exploit shared memory within a node and MPI across nodes. This hybrid programming model needs efficient support for multithreaded MPI communication. In this paper, we describe the optimization of one aspect of a multithreaded MPI implementation: concurrent accesses from multiple threads to various MPI objects, such as communicators, datatypes, and requests. The semantics of the creation, usage, and destruction of these objects implies, but does not strictly require, the use of reference counting to prevent memory leaks and premature object destruction. We demonstrate how a naive multithreaded implementation of MPI object management via reference counting incurs a significant performance penalty. We then detail two solutions that we have implemented in MPICH2 to mitigate this problem almost entirely, including one based on a novel garbage collection scheme. In our performance experiments, this new scheme improved the multithreaded messaging rate by up to 31% over the naive reference counting method. David Goodell, Pavan Balaji, Darius Buntinas, Gábor Dózsa, William Gropp, Sameer Kumar 0001, Bronis R. de Supinski, Rajeev Thakur |
CLUSTER | 5 |
| 2010 | Extreme scale computing: Challenges and opportunitiesabstractAn extreme scale system is one that is one thousand times more capable than a current comparable system, with the same power and physical footprint. Intuitively, this means that the power consumption and physical footprint of a current departmental server should be enough to deliver petascale performance, and that a single, commodity chip should deliver terascale performance. In this panel, we will discuss the resulting challenges in energy/power efficiency, concurrency and locality, resiliency and programmability, and the research opportunities that may take us to extreme scale systems. Josep Torrellas, William Gropp, Vivek Sarkar, Jaime H. Moreno, Kunle Olukotun |
HPCA | 2 |
| 2010 | An introductory exascale feasibility study for FFTs and multigridabstractThe coming decade is going to see a push towards exascale computing. Assuming gigahertz cores, this means exascale systems will have between 100 million and 1 billion of them to achieve this level of performance. At this scale, some important questions need to be answered on the applications end. What applications are feasible at this scale? What needs to be done to make them scalable? How does the hardware have to adapt to meet application needs? In this paper, we introduce a new feasibility-based approach to answering these questions. Our approach involves finding upper and lower bounds on problem size and machine parameters to determine a feasibility region for the application in question. As the underlying architecture of a future exascale machine is currently unknown, we use LogP-based performance models and vary machine parameters to give architecture-indepenent hardware constraints. We consider both strong-scaling and weak-scaling scenarios, and present results for two applications, the Fast Fourier Transform and basic geometric multigrid. The results show substantial constraints that need to be satisfied to enable exascale performance. Hormozd Gahvari, William Gropp |
IPDPS | 2 |
| 2010 | An adaptive performance modeling tool for GPU architecturesabstractThis paper presents an analytical model to predict the performance of Sara S. Baghsorkhi, Matthieu Delahaye, Sanjay J. Patel, William Gropp, Wen-Mei W. Hwu |
PPoPP | 4 |
| 2010 | Extreme scale computing: challenges and opportunitiesabstractNo abstract available. Josep Torrellas, William Gropp, Jaime H. Moreno, Kunle Olukotun, Vivek Sarkar |
PPoPP | 2 |
| 2010 | PMI: A Scalable Parallel Process-Management Interface for Extreme-Scale Systems
Pavan Balaji, Darius Buntinas, David Goodell, William Gropp, Jayesh Krishna, Ewing L. Lusk, Rajeev Thakur |
EuroMPI | 4 |
| 2010 | Enabling Concurrent Multithreaded MPI Communication on Multicore Petascale Systems
Gábor Dózsa, Sameer Kumar 0001, Pavan Balaji, Darius Buntinas, David Goodell, William Gropp, Joe Ratterman, Rajeev Thakur |
EuroMPI | 6 |
| 2010 | Toward Performance Models of MPI Implementations for Understanding Application Scaling Issues
Torsten Hoefler, William Gropp, Rajeev Thakur, Jesper Larsson Träff |
EuroMPI | 2 |
| 2010 | Load Balancing for Regular Meshes on SMPs with MPI
Vivek Kale, William Gropp |
EuroMPI | 2 |
| 2010 | A Scalable MPI_Comm_split Algorithm for Exascale Computing
Paul Sack, William Gropp |
EuroMPI | 2 |
| 2010 | Formal methods applied to high-performance computing software design: a case study of MPI one-sided communication-based lockingabstractAbstract There is a growing need to address the complexity of verifying the numerous concurrent protocols employed in the high‐performance computing software. Today's approaches for verification consist of testing detailed implementations of these protocols. Unfortunately, this approach can seldom show the absence of bugs, and often results in serious bugs escaping into the deployed software. An approach calledModel Checkinghas been demonstrated to be eminently helpful in debugging these protocols early in the software life cycle by offering the ability to represent and exhaustively analyze simplified formal protocol models. The effectiveness of model checking has yet to be adequately demonstrated in high‐performance computing. This paper presents a case study of a concurrent protocol that was thought to be sufficiently well tested, but proved to contain two very non‐obvious deadlocks in them. These bugs were automatically detected through model checking. The protocol models in which these bugs were detected were also easy to create. Recent work in our group demonstrates that even this tedium of model creation can be eliminated by employing dynamic source‐code‐level analysis methods. Our case study comes from the important domain of Message Passing Interface (MPI)‐based programming, which is universally employed for simulating and predicting anything from the structural integrity of combustion chambers to the path of hurricanes. We argue that model checking must be taught as well as used widely within HPC, given this and similar success stories. Copyright © 2009 John Wiley & Sons, Ltd. Salman Pervez, Ganesh Gopalakrishnan, Robert M. Kirby, Rajeev Thakur, William Gropp |
Softw. Pract. Exp. | 5 |
| 2010 | Self-Consistent MPI Performance GuidelinesabstractMessage passing using the Message-Passing Interface (MPI) is at present the most widely adopted framework for programming parallel applications for distributed memory and clustered parallel systems. For reasons of (universal) implementability, the MPI standard does not state any specific performance guarantees, but users expect MPI implementations to deliver good and consistent performance in the sense of efficient utilization of the underlying parallel (communication) system. For performance portability reasons, users also naturally desire communication optimizations performed on one parallel platform with one MPI implementation to be preserved when switching to another MPI implementation on another platform. We address the problem of ensuring performance consistency and portability by formulating performance guidelines and conditions that are desirable for good MPI implementations to fulfill. Instead of prescribing a specific performance model (which may be realistic on some systems, under some MPI protocol and algorithm assumptions, etc.), we formulate these guidelines by relating the performance of various aspects of the semantically strongly interrelated MPI standard to each other. Common-sense expectations, for instance, suggest that no MPI function should perform worse than a combination of other MPI functions that implement the same functionality, no specialized function should perform worse than a more general function that can implement the same functionality, no function with weak semantic guarantees should perform worse than a similar function with stronger semantics, and so on. Such guidelines may enable implementers to provide higher quality MPI implementations, minimize performance surprises, and eliminate the need for users to make special, nonportable optimizations by hand. We introduce and semiformalize the concept of self-consistent performance guidelines for MPI, and provide a (nonexhaustive) set of such guidelines in a form that could be automatically verified by benchmarks and experiment management tools. We present experimental results that show cases where guidelines are not satisfied in common MPI implementations, thereby indicating room for improvement in today's MPI implementations. Jesper Larsson Träff, William Gropp, Rajeev Thakur |
IEEE Trans. Parallel Distributed Syst. | 2 |
| 2009 | Natively Supporting True One-Sided Communication inabstractAs high-end computing systems continue to grow in scale, the performance that applications can achieve on such large scale systems depends heavily on their ability to avoid explicitly synchronized communication with other processes in the system. Accordingly, several modern and legacy parallel programming models (such as MPI, UPC, global arrays) have provided many programming constructs that enable implicit communication using one-sided communication operations. While MPI is the most widely used communication model for scientific computing, the usage of one-sided communication is restricted; this is mainly owing to the inefficiencies in current MPI implementations that internally rely on synchronization between processes even during one-sided communication, thus losing the potential of such constructs. In our previous work, we had utilized native one-sided communication primitives offered by high-speed networks such as infiniband (IB) to allow for true one-sided communication in MPI. In this paper, we extend this work to natively take advantage of one-sided atomic operations on cache-coherent multi-core/multi-processor architectures while still utilizing the benefits of networks such as IB. Specifically, we present a sophisticated hybrid design that uses locks that migrate between IB hardware atomics and multi-core CPU atomics to take advantage of both. We demonstrate the capability of our proposed design with a wide range of experiments illustrating its benefits in performance as well as its potential to avoid explicit synchronization. Gopalakrishnan Santhanaraman, Pavan Balaji, Rajeev Thakur, William Gropp, Dhabaleswar K. Panda 0001 |
CCGRID | 5 |
| 2009 | Investigating High Performance RMA Interfaces for the MPI-3 StandardabstractThe MPI-2 Standard, released in 1997, defined an interface for one-sided communication, also known as remote memory access (RMA). It was designed with the goal that it should permit efficient implementations on multiple platforms and networking technologies, and also in heterogeneous environments and non-cache-coherent systems. Nonetheless, even 12 years after its existence, the MPI-2 RMA interface remains scarcely used for a number of reasons. This paper discusses the limitations of the MPI-2 RMA specification, outlines the goals and requirements for a new RMA API that would better meet the needs of both users and implementers, and presents a strawman proposal for such an API. We also study the tradeoffs facing the design of this new API and discuss how it may be implemented efficiently on both cache-coherent and non-cache-coherent systems. Vinod Tipparaju, William Gropp, Hubert Ritzdorf, Rajeev Thakur, Jesper Larsson Träff |
ICPP | 2 |
| 2009 | Test suite for evaluating performance of multithreaded MPI communication
Rajeev Thakur, William Gropp |
Parallel Comput. | 2 |
| 2008 | Communication Analysis of Parallel 3D FFT for Flat Cartesian Meshes on Large Blue Gene Systems
Pavan Balaji, William Gropp, Rajeev Thakur |
HiPC | 3 |
| 2008 | Improving the Performance of Tensor Matrix Vector Multiplication in Cumulative Reaction Probability Based Quantum Chemistry Codes
Dinesh K. Kaushik, William Gropp, Michael Minkoff, Barry Smith 0002 |
HiPC | 2 |
| 2008 | 2008 International Conference on Parallel Processing September 8-12, 2008 Portland, Oregon Exploring Parallel I/O Concurrency with Speculative PrefetchingabstractParallel applications can benefit greatly from massive computational capability, but their performance usually suffers due to large latency in I/O accesses. Conventional I/O prefetching techniques are conservative and are limited by low accuracy and coverage. As the processor performance has been increasing rapidly and the computing power is virtually free, we introduce a novel speculative approach for comprehensive and aggressive parallel I/O prefetching in this study. We present the design of our approach as well as challenges, solutions, and our prototype implementation. The experiments have shown promising results in reducing I/O access latency. Yong Chen 0001, Surendra Byna, Xian-He Sun, Rajeev Thakur, William Gropp |
ICPP | 5 |
| 2008 | Parallel I/O prefetching using MPI file caching and I/O signaturesabstractParallel I/O prefetching is considered to be effective in improving I/O performance. However, the effectiveness depends on determining patterns among future I/O accesses swiftly and fetching data in time, which is difficult to achieve in general. In this study, we propose an I/O signature-based prefetching strategy. The idea is to use a predetermined I/O signature of an application to guide prefetching. To put this idea to work, we first derived a classification of patterns and introduced a simple and effective signature notation to represent patterns. We then developed a toolkit to trace and generate I/O signatures automatically. Finally, we designed and implemented a thread-based client-side collective prefetching cache layer for MPI-IO library to support prefetching. A prefetching thread reads I/O signatures of an application and adjusts them by observing I/O accesses at runtime. Experimental results show that the proposed prefetching method improves I/O performance significantly for applications with complex patterns. Surendra Byna, Yong Chen 0001, Xian-He Sun, Rajeev Thakur, William Gropp |
SC | 5 |
| 2008 | Hiding I/O latency with pre-execution prefetching for parallel applicationsabstractParallel applications are usually able to achieve high computational performance but suffer from large latency in I/O accesses. I/O prefetching is an effective solution for masking the latency. Most of existing I/O prefetching techniques, however, are conservative and their effectiveness is limited by low accuracy and coverage. As the processor-I/O performance gap has been increasing rapidly, data-access delay has become a dominant performance bottleneck. We argue that it is time to revisit the ldquoI/O wallrdquo problem and trade the excessive computing power with data-access speed. We propose a novel pre-execution approach for masking I/O latency. We describe the pre-execution I/O prefetching framework, the pre-execution thread construction methodology, the underlying library support, and the prototype implementation in the ROMIO MPI-IO implementation in MPICH2. Preliminary experiments show that the pre-execution approach is promising in reducing I/O access latency and has real potential. Yong Chen 0001, Surendra Byna, Xian-He Sun, Rajeev Thakur, William Gropp |
SC | 5 |
| 2007 | Advanced Flow-control Mechanisms for the Sockets Direct Protocol over InfiniBandabstractThe Sockets Direct Protocol (SDP) is an industry standard to allow existing TCP/IP applications to be executed on high-speed networks such as InfiniBand (IB). Like many other high-speed networks, IB requires the receiver process to inform the network interface card (NIC), before the data arrives, about buffers in which incoming data has to be placed. To ensure that the receiver process is ready to receive data, the sender process typically performs flow-control on the data transmission. Existing designs of SDP flow-control are naive and do not take advantage of several interesting features provided by IB. Specifically, features such as RDMA are only used for performing zero-copy communication, although RDMA has more capabilities such as sender-side buffer management (where a sender process can manage SDP resources for the sender as well as the receiver). Similarly, IB also provides hardware flow-control capabilities that have not been studied in previous literature. In this paper, we utilize these capabilities to improve the SDP flow-control over IB using two designs: RDMA-based flow-control and NIC-assisted RDMA-based flow-control. We evaluate the designs using micro-benchmarks and real applications. Our evaluations reveal that these designs can improve the resource usage of SDP and consequently its performance by an order-of-magnitude in some cases. Moreover we can achieve 10-20% improvement in performance for various applications. Pavan Balaji, Sitha Bhagvat, Dhabaleswar K. Panda 0001, Rajeev Thakur, William Gropp |
ICPP | 5 |
| 2007 | Nonuniformly Communicating Noncontiguous Data: A Case Study with PETSc and MPIabstractDue to the complexity associated with developing parallel applications, scientists and engineers rely on high-level software libraries such as PETSc, ScaLAPACK and PESSL to ease this task. Such libraries assist developers by providing abstractions for mathematical operations, data representation and management of parallel layouts of the data, while internally using communication libraries such as MPI and PVM. With high-level libraries managing data layout and communication internally, it can be expected that they organize application data suitably for performing the library operations optimally. However, this places additional overhead on the underlying communication library by making the data layout noncontiguous in memory and communication volumes (data transferred by a process to each of the other processes) nonuniform. In this paper, we analyze the overheads associated with these two aspects (noncontiguous data layouts and nonuniform communication volumes) in the context of the PETSc software toolkit over the MPI communication library. We describe the issues with the current approaches used by MPICH2 (an implementation of MPI), propose different approaches to handle these issues and evaluate these approaches with micro-benchmarks as well as an application over the PETSc software library. Our experimental results demonstrate close to an order of magnitude improvement in the performance of a 3-D Laplacian multi-grid solver application when evaluated on a 128 processor cluster. Pavan Balaji, Darius Buntinas, Satish Balay, Barry Smith 0002, Rajeev Thakur, William Gropp |
IPDPS | 6 |
| 2007 | Analyzing the impact of supporting out-of-order communication on in-order performance with iWARPabstractDue to the growing need to tolerate network faults and congestion in high-end computing systems, supporting multiple network communication paths is becoming increasingly important. However, multi-path communication comes with the disadvantage of out-of-order arrival of packets (because packets may traverse different paths). While modern networking stacks such as the Internet Wide-Area RDMA Protocol (iWARP) over 10-Gigabit Ethernet (10GE) support multi-path communication, their current implementations do not handle out-of-order packets primarily owing to the overhead on in-order communication that it adds. Specifically, in iWARP, supporting out-of-order packets requires every packet to carry additional information causing significant overhead on packets that arrive in-order. Thus, in this paper, we analyze the trade-offs in designing a feature-complete iWARP stack, i.e., one that provides support for out-of-order arriving packets, and thus, multi-path systems, while focusing on the performance of in-order communication. We propose three feature-complete designs of iWARP and analyze the pros and cons of each of these designs using performance experiments based on several micro-benchmarks as well as an iso-surface visual rendering application. Our analysis reveals that the iWARP design providing the best overall performance depends on the particular characteristics of the upper layers and that different designs are optimal based on the metric of interest. Pavan Balaji, Wu-chun Feng, Sitha Bhagvat, Dhabaleswar K. Panda 0001, Rajeev Thakur, William Gropp |
SC | 6 |
| 2007 | Implementation and evaluation of shared-memory communication and synchronization operations in MPICH2 using the Nemesis communication subsystem
Darius Buntinas, Guillaume Mercier, William Gropp |
Parallel Comput. | 3 |
| 2007 | Thread-safety in an MPI implementation: Requirements and analysis
William Gropp, Rajeev Thakur |
Parallel Comput. | 1 |
| 2006 | Design and Evaluation of Nemesis, a Scalable, Low-Latency, Message-Passing Communication SubsystemabstractThis paper presents a new low-level communication subsystem called Nemesis. Nemesis has been designed and implemented to be scalable and efficient both in the intranode communication context using shared-memory and in the internode communication case using high-performance networks and is natively multimethod-enabled. Nemesis has been integrated in MPICH2 as a CH3 channel and delivers better performance than other dedicated communication channels in MPICH2. Furthermore, the resulting MPICH2 architecture outperforms other MPI implementations in point-to-point benchmarks. Darius Buntinas, Guillaume Mercier, William Gropp |
CCGRID | 3 |
| 2006 | High performance file I/O for the Blue Gene/L supercomputerabstractParallel I/O plays a crucial role for most data-intensive applications running on massively parallel systems like Blue Gene/L that provides the promise of delivering enormous computational capability. We designed and implemented a highly scalable parallel file I/O architecture for Blue Gene/L, which leverages the benefit of the hierarchical and functional partitioning design of the system software with separate computational and I/O cores. The architecture exploits the scalability aspect of GPFS (General Parallel File System) at the backend, while using MPI I/O as an interface between the application I/O and the file system. We demonstrate the impact of our high performance I/O solution for Blue Gene/L with a comprehensive evaluation that consists of a number of widely used parallel I/O benchmarks and I/O intensive applications. Our design and implementation is not only able to deliver at least one order of magnitude speed up in terms of I/O bandwidth for a real-scale application HOMME (achieving aggregate bandwidth of 1.8 GB/Sec and 2.3 GB/Sec for write and read accesses, respectively), but also supports high-level parallel I/O data interfaces such as parallel HDF5 and parallel NetCDF scaling up to a large number of processors. Hao Yu 0008, Ramendra K. Sahoo, C. Howson, Gheorghe Almási 0001, José G. Castaños, Manish Gupta 0002, José E. Moreira, Jeff Parker, Thomas Engelsiepen, Robert B. Ross, Rajeev Thakur, Robert Latham, William Gropp |
HPCA | 13 |
| 2006 | Data Transfers between Processes in an SMP System: Performance Study and Application to MPIabstractThis paper focuses on the transfer of large data in SMP systems. Achieving good performance for intranode communication is critical for developing an efficient communication system, especially in the context of SMP clusters. We evaluate the performance of five transfer mechanisms: shared-memory buffers, message queues, the Ptrace system call, kernel module-based copy, and a high-speed network. We evaluate each mechanism based on latency, bandwidth, its impact on application cache usage, and its suitability to support MPI two-sided and one-sided messages Darius Buntinas, Guillaume Mercier, William Gropp |
ICPP | 3 |
| 2006 | Collective communication on architectures that support simultaneous communication over multiple linksabstractTraditional collective communication algorithms are designed with the assumption that a node can communicate with only one other node at a time. On new parallel architectures such as the IBM Blue Gene/L, a node can communicate with multiple nodes simultaneously. We have redesigned and reimplemented many of the MPI collective communication algorithms to take advantage of this ability to send simultaneously, including broadcast, reduce(-to-one), scatter, gather, allgather, reduce-scatter, and allreduce. We show that these new algorithms have lower expected costs than the previously known lower bounds based on old models of parallel computation. Results are included comparing their performance to the default implementations in IBM's MPI. Ernie Chan, Robert A. van de Geijn, William Gropp, Rajeev Thakur |
PPoPP | 3 |
| 2006 | S01 - Advanced MPI: I/O and one-sided communicationabstractThis tutorial is about advanced use of MPI, in particular the parallel I/O and one-sided communication features added in MPI-2. Implementations are now available both from vendors and from open-source projects so that these MPI-2 capabilities can now really be used in practice. The tutorial will be heavily example-driven. For each example we introduce concepts, describe the problem being solved, then walk through the code and its execution.Examples were chosen to cover scenarios seen in real applications, such as 1D and 2D mesh decomposition, checkpointing of sparse data structures, and providing atomic access to shared memory data structures. Attendees will leave the tutorial with both an understanding of these advanced concepts and a collection of working example codes that they are familiar with and have seen in action. This will prepare them for applying these concepts in their own applications. William Gropp, Ewing L. Lusk, Rajeev Thakur, Robert B. Ross |
SC | 1 |
| 2006 | M01 - Application supercomputing and multiscale simulation techniquesabstractTeraflop performance is no longer something of the future as complex integrated and multiscale 3D simulations drive supercomputer development. This tutorial addresses computation at the highest end. An overview of architectures is given (BlueGene/L, Columbia, NEC SX-8, Cray and IBM lines, high-performing clusters) along with programming tools necessary for application development. Parallel programming concepts (MPI, OpenMP, HPF, UPC, CAF) are reviewed and compared. What are the major issues facing application code developers today? How do the challenges vary from cluster computing to the complex hybrid architectures with superscalar and vector processors? What are the barriers we must overcome to achieve true sustained petascale performance? We address these questions and give tips, tricks, and tools of the trade for large mulitscale application development. We discuss some advanced MPI including dynamic process management and optimization. We draw from a series of terascale and multiscale applications and discuss specific challenges and performance issues. Alice E. Koniges, William Gropp, Ewing L. Lusk, David C. Eder |
SC | 2 |
| 2006 | Awards & video - Awards sessionabstractIn this session the awards for: Best Paper, Best Student Paper, Best Poster, ACM Student Research Competition, Analytics Challenge, Bandwidth Challenge, and Storage Challenge will be presented. Daniel A. Reed, William Gropp, Allan Sussman, Jeffrey J. Evans, Paul Fussell, Debbie Montano, Raymond L. Paden |
SC | 2 |
| 2006 | Multi-core issues - Multi-Core for HPC: breakthrough or breakdown?abstractA dramatic trend in computing is the adoption of multi-core technology by the vendors from which our current and future HPC systems are being derived. Multi-core is offered as a path to continued reliance and benefits of Moore's Law while reining in the previously unfettered growth of power consumption and design complexity. Are we saved? or is it but a fools mission, trapping us in a technical cul de sac with no long term direction and no way to reinvent an alternative future. The panel will consider the following questions:* Can multi-core span the next decade of Moore's Law progression?* Are the pins and caches a strangle hold on the future effectiveness of multi-core?* Can innovative algorithmic techniques exploit the opportunities and address the challenges of multi-core?* How will programming models and supporting system software change to accommodate the unique properties and peculiarities of multi-core structures? Thomas L. Sterling, Peter M. Kogge, William J. Dally, Steve Scott, William Gropp, David E. Keyes, Pete Beckman |
SC | 5 |
| 2005 | Implementing MPI-IO atomic mode without file system supportabstractThe ROMIO implementation of the MPI-IO standard provides a portable infrastructure for use on top of any number of different underlying storage targets. These different targets vary widely in their capabilities, and in some cases, additional effort is needed within ROMIO to support the complete MPI-IO semantics. One aspect of the interface that can be problematic to implement is the MPI-IO atomic mode. This mode requires enforcing strict consistency semantics. For some file systems, native locks may be used to enforce these semantics, but not all file systems have lock support. In this work, we describe two algorithms for implementing efficient mutex locks using MPI-1 and MPI-2 capabilities. We then show how these algorithms may be used to implement a portable MPI-IO atomic mode for ROMIO. We evaluate the performance of these algorithms and show that they impose little additional overhead on the system. Because of the low-overhead nature of these algorithms, they are likely useful in a variety of situations where distributed locks are needed in the MPI-2 environment. Robert B. Ross, Robert Latham, William Gropp, Rajeev Thakur, Brian R. Toonen |
CCGRID | 3 |
| 2004 | High performance MPI-2 one-sided communication over InfiniBandabstractMany existing MPI-2 one-sided communication implementations are built on top of MPI send/receive operations. Although this approach can achieve good portability, it suffers front high communication overhead and dependency on remote process for communication progress. To address these problems, we propose a high performance MPI-2 one-sided communication design over the InfiniBand Architecture. In our design, MPI-2 one-sided communication operations such as MPI-Put, MPI-Get and MPI-Accumulate are directly mapped to InfiniBand Remote Direct Memory Access (RDMA) operations. Our design has been implemented based on MPICH2 over InfiniBand. We present detailed design issues for this approach and perform a set of microbenchmarks to characterize different aspects of its performance. Our performance evaluation shows that compared with the design based on MPI send/receive, our design can improve throughput up to 77%, and reduce latency and synchronization overhead up to 19% and 13%, respectively. Under certain process skew, the bad impact can be significantly reduced by new design, from 41% to nearly 0%. It also can achieve better overlap of communication and computation. Weihang Jiang, Jiuxing Liu, Hyun-Wook Jin, Dhabaleswar K. Panda 0001, William Gropp, Rajeev Thakur |
CCGRID | 5 |
| 2004 | Predicting memory-access cost based on data-access patternsabstractImproving memory performance at software level is more effective in reducing the rapidly expanding gap between processor and memory performance. Loop transformations (e.g. loop unrolling, loop tiling) and array restructuring optimizations improve the memory performance by increasing the locality of memory accesses. To find the best optimization parameters at runtime, we need a fast and simple analytical model to predict the memory access cost. Most of the existing models are complex and impractical to be integrated in the runtime tuning systems. In this paper, we propose a simple, fast and reasonably accurate model that is capable of predicting the memory access cost based on a wide range of data access patterns that appear in many scientific applications. Surendra Byna, Xian-He Sun, William Gropp, Rajeev Thakur |
CLUSTER | 3 |
| 2004 | Implementing MPI on the BlueGene/L Supercomputer
Gheorghe Almási 0001, Charles Archer, José G. Castaños, C. Christopher Erway, Philip Heidelberger, Xavier Martorell, José E. Moreira, Kurt W. Pinnow, Joe Ratterman, Nils Smeds, Burkhard D. Steinmacher-Burow, William Gropp, Brian R. Toonen |
Euro-Par | 12 |
| 2004 | Design and Implementation of MPICH2 over InfiniBand with RDMA SupportabstractSummary form only given. For several years, MPI has been the de facto standard for writing parallel applications. One of the most popular MPI implementations is MPICH. Its successor, MPICH2, features a completely new design that provides more performance and flexibility. To ensure portability, it has a hierarchical structure based on which porting can be done at different levels. In this paper, we present our experiences in designing and implementing MPICH2 over InfiniBand. Because of its high performance and open standard, InfiniBand is gaining popularity in the area of high-performance computing. Our study focuses on optimizing the performance of MPl-1 functions in MPICH2. One of our objectives is to exploit remote direct memory access (RDMA) in InfiniBand to achieve high performance. We have based our design on the RDMA channel interface provided by MP1CH2, which encapsulates architecture-dependent communication functionalities into a very small set of functions. Starting with a basic design, we apply different optimizations and also propose a zero-copy-based design. We characterize the impact of our optimizations and designs using microbenchmarks. We have also performed an application-level evaluation using the NAS parallel benchmarks. Our optimized MPICH2 implementation achieves 7.6/spl mu/s latency and 857 MB/s bandwidth, which are close to the raw performance of the underlying InfiniBand layer. Our study shows that the RDMA channel interface in MPICH2 provides a simple, yet powerful, abstraction that enables implementations with high performance by exploiting RDMA operations in InfiniBand. To the best of our knowledge, this is the first high-performance design and implementation ofMPICH2 on InfiniBand using RDMA support. Jiuxing Liu, Weihang Jiang, Pete Wyckoff, Dhabaleswar K. Panda 0001, David Ashton, Darius Buntinas, William Gropp, Brian R. Toonen |
IPDPS | 7 |
| 2003 | Noncontiguous I/O Accesses Through MPI-IOabstractI/O performance remains a weakness of parallel computing systems today. While this weakness is partly attributed to rapid advances in other system components, I/O interfaces available to programmers and the I/O methods supported by file systems have traditionally not matched efficiently with the types of I/O operations that scientific applications perform, particularly noncontiguous accesses. The MPI-IO interface allows for rich descriptions of the I/O patterns desired for scientific applications and implementations such as ROMIO have taken advantage of this ability while remaining limited by underlying file system methods. A method of noncontiguous data access, list I/O, was recently implemented in the Parallel Virtual File System (PVFS). We implement support for this interface in the ROMIO MPI-IO implementation. Through a suite of noncontiguous I/O tests we compared ROMIO list I/O to current methods of ROMIO noncontiguous access and found that the list I/O interface provides performance benefits in many noncontiguous cases. Avery Ching, Alok N. Choudhary, Kenin Coloma, Wei-keng Liao, Robert B. Ross, William Gropp |
CCGRID | 6 |
| 2003 | Improving the Performance of MPI Derived Datatypes by Optimizing Memory-Access CostabstractThe MPI Standard supports derived datatypes, which allow users to describe noncontiguous memory layout and communicate noncontiguous data with a single communication function. This feature enables an MPI implementation to optimize the transfer of noncontiguous data. In practice, however, few MPI implementations implement derived datatypes in a way that performs better than what the user can achieve by manually packing data into a contiguous buffer and then calling an MPI function. In this paper, we present a technique for improving the performance of derived datatypes by automatically using packing algorithms that are optimized for memory-access cost. The packing algorithms use memory-optimization techniques that the user cannot apply easily without advanced knowledge of the memory architecture. We present performance results for a matrix-transpose example that demonstrate that our implementation of derived datatypes significantly outperforms both manual packing by the user and the existing derived-datatype code in the MPI implementation (MPICH). Surendra Byna, William Gropp, Xian-He Sun, Rajeev Thakur |
CLUSTER | 2 |
| 2003 | Efficient Structured Data Access in Parallel File SystemsabstractParallel scientific applications store and retrieve very large, structured datasets. Directly supporting these structured accesses is an important step in providing high-performance I/O solutions for these applications. High-level interfaces such as HDF5 and Parallel netCDF provide convenient APIs for accessing structured datasets, and the MPI-IO interface also supports efficient access to structured data. However, parallel file systems do not traditionally support such access. In this work we present an implementation of structured data access support in the context of the parallel virtual file system (PVFS). We call this support "datatype I/O" because of its similarity to MPI datatypes. This support is built by using a reusable datatype-processing component from the MPICH2 MPI implementation. We describe how this component is leveraged to efficiently process structured data representations resulting from MPI-IO operations. We quantitatively assess the solution using three test applications. We also point to further optimizations in the processing path that could be leveraged for even more efficient operation. Avery Ching, Alok N. Choudhary, Wei-keng Liao, Robert B. Ross, William Gropp |
CLUSTER | 5 |
| 2003 | Using MPI-2: Advanced Features of the Message Passing Interface
William Gropp, Ewing L. Lusk, Robert B. Ross, Rajeev Thakur |
CLUSTER | 1 |
| 2003 | Toward Understanding Soft Faults in High Performance Cluster Networks
Jeffrey J. Evans, Seongbok Baik, Cynthia Hood 0001, William Gropp |
Integrated Network Management | 4 |
| 2003 | Exploring the Relationship Between Parallel Application Run-Time Variability and Network Performance in ClustersabstractHighly variable parallel application execution time is a persistent issue in cluster computing environments, and can be particularly acute in systems composed of networks of workstations (NOWs). We are looking at this issue in terms of consistency. In particular, we are focusing on network performance. Before we can use techniques from fault management to attain consistency, this paper presents our preliminary analysis of run-time variability from logs and experiments, exposing important issues related to systemic inconsistency in NOW clusters. The characterization of application sensitivity can be used to set network performance goals, thereby defining operational requirements. Network performance depends on the virtual topology imposed by the scheduler's allocation of nodes and the communication patterns of the set of running applications. Therefore it is important to look at both the network and the cluster's centralized node mapper (scheduler) as critical subsystems. Jeffrey J. Evans, Cynthia Hood 0001, William Gropp |
LCN | 3 |
| 2003 | Parallel netCDF: A High-Performance Scientific I/O InterfaceabstractDataset storage, exchange, and access play a critical role in scientific applications. For such purposes netCDF serves as a portable, efficient file format and programming interface, which is popular in numerous scientific application domains. However, the original interface does not provide an efficient mechanism for parallel data storage and access. In this work, we present a new parallel interface for writing and reading netCDF datasets. This interface is derived with minimal changes from the serial netCDF interface but defines semantics for parallel access and is tailored for high performance. The underlying parallel I/O is achieved through MPI-IO, allowing for substantial performance gains through the use of collective I/O optimizations. We compare the implementation strategies and performance with HDF5. Our tests indicate programming convenience and significant I/O performance improvement with this parallel netCDF (PnetCDF) interface. Wei-keng Liao, Alok N. Choudhary, Robert B. Ross, Rajeev Thakur, William Gropp, Robert Latham, Andrew R. Siegel, Brad Gallagher, Michael Zingale |
SC | 6 |
| 2002 | Noncontiguous I/O through PVFSabstractWith the tremendous advances in processor and memory technology, I/O has risen to become the bottleneck in high-performance computing for many applications. The development of parallel file systems has helped to ease the performance gap, but I/O still remains an area needing significant performance improvement. Research has found that noncontiguous I/O access patterns in scientific applications combined with current file system methods, to perform these accesses lead to unacceptable performance for large data sets. To enhance performance of noncontiguous I/O, we have created list I/O, a native version of noncontiguous I/O. We have used the Parallel Virtual File System (PVFS) to implement our ideas. Our research and experimentation shows that list I/O outperforms current noncontiguous I/O access methods in most I/O situations and can substantially enhance the performance of real-world scientific applications. Avery Ching, Alok N. Choudhary, Wei-keng Liao, Robert B. Ross, William Gropp |
CLUSTER | 5 |
| 2002 | Goals Guiding Design: PVM and MPabstractPVM and MPI, two systems for programming clusters, are often compared. The comparisons usually start with the unspoken assumption that PVM and MPI represent different solutions to the same problem. In this paper we show that, in fact, the two systems often are solving different problems. In cases where the problems do match but the solutions chosen by PVM and MPI are different, we explain the reasons for the differences. Usually such differences can be traced to explicit differences in the goals of the two systems, their origins, or the relationship between their specifications and their implementations. For example, we show that the requirement for portability and performance across many platforms caused MPI to choose approaches different from those made by PVM, which is able to exploit the similarities of network-connected systems. William Gropp, Ewing L. Lusk |
CLUSTER | 1 |
| 2002 | An Evaluation of Object-Based Data Transfers on High Performance NetworksabstractWe describe FOBS: a simple user-level communication protocol designed to take advantage of the available bandwidth in a high-bandwidth, high-delay network environment. We compare the performance of FOBS with that of TCP both with and without the so-called Large Window extensions designed to improve the performance of TCP in this type of network environment. It is shown that FOBS can obtain on the order of 90% of the available bandwidth across both short and long high-performance network connections. In the case of the long haul connection, this represents a bandwidth that is 1.8 times higher than that of the optimized TCP algorithm. Also, we demonstrate that the additional traffic placed on the network due to the greedy nature of the algorithm is quite reasonable, representing approximately 3% of the total data transferred. Phillip M. Dickens, William Gropp |
HPDC | 2 |
| 2002 | Prototype of AM3: Active Mapper and Monitoring Module for Myrinet EnvironmentsabstractSystem area networks (SAN) have been developed to address the needs of computing clusters. Myricom's Myrinet architecture is one of the predominant technologies in this area. One of the key issues for SANs is fault-tolerant routing. Myrinet provides Mapper software to discover and maintain network topology. Myrinet's Mapper is centralized, susceptible to probe packet deadlock, and does not incorporate host monitoring. We propose an alternative mapper called AM3. AM3 is hierarchical, reduces the number of probe packets, and incorporates host monitoring. We have implemented a prototype version of AM3 on Chiba City, Argonne National Lab's 512 CPU Linux cluster. Seongbok Baik, Cynthia Hood 0001, William Gropp |
LCN | 3 |
| 2002 | Optimizing noncontiguous accesses in MPI-IO
Rajeev Thakur, William Gropp, Ewing L. Lusk |
Parallel Comput. | 2 |
| 2001 | Advanced Cluster Programming with MP
William Gropp |
CLUSTER | 1 |
| 2001 | Learning from the Success of MPI
William Gropp |
HiPC | 1 |
| 2001 | Interfacing Parallel Jobs to Process ManagersabstractA variety of projects worldwide are developing what we call "heterogeneous MPI". These MPI implementations are designed to operate on multiple computers, perhaps of different types, ranging in complexity from a set of desktop workstations to several supercomputers connected via a wide area network. These considerations led us to investigate the feasibility of defining a common API that could be used within MPI implementations to access process startup, initialization, monitoring, and control functions provided by an underlying process management system. If various MPI implementations coded to that API, one could then develop multiple "process management" modules that could be reused within different MPI implementations, thus allowing partitioning of effort between different development groups. In pursuit of this goal, we have designed such an API, which we call BNR. The major goals of the BNR interface are outlined. Brian R. Toonen, David Ashton, Ewing L. Lusk, Ian T. Foster, William Gropp, Edgar Gabriel, Ralph M. Butler, Nicholas T. Karonis |
HPDC | 5 |
| 2001 | Components and interfaces of a process management system for parallel programs
Ralph M. Butler, William Gropp, Ewing L. Lusk |
Parallel Comput. | 2 |
| 2001 | High-performance parallel implicit CFD
William Gropp, Dinesh K. Kaushik, David E. Keyes, Barry Smith 0002 |
Parallel Comput. | 1 |
| 2000 | Analyzing the Parallel Scalability of an Implicit Unstructured Mesh CFD Code
William Gropp, Dinesh K. Kaushik, Barry Smith 0002, David E. Keyes |
HiPC | 1 |
| 2000 | Exploiting Hierarchy in Parallel Computer Networks to Optimize Collective Operation PerformanceabstractThe efficient implementation of collective communication operations has received much attention. Initial efforts modeled network communication and produced "optimal" trees based on those models. However, the models used by these initial efforts assumed equal point-to-point latencies between any two processes. This assumption is violated in heterogeneous systems such as clusters of SMPs and wide-area "computational grids", and as a result, collective operations that utilize the trees generated by these models perform suboptimally. In response, more recent work has focused on creating topology-aware trees for collective operations that minimize communication across slower channels (e.g., a wide-area network). While these efforts have significant communication benefits, they all limit their view of the network to only two layers. We present a strategy based upon a multilayer view of the network. By creating multilevel topology trees we take advantage of communication cost differences at every level in the network. We used this strategy to implement topology-aware versions of several MPI collective operations in MPICH-G, the Globus-enabled version of the popular MPICH implementation of the MPI standard. Using information about topology discovered by Globus, we construct these topology-aware trees automatically during execution, thus freeing the MPI application programmer from having to write special files or functions to describe the topology to the MPICH library. We present results demonstrating the advantages of our multilevel approach by comparing it to the default (topology-unaware) implementation provided by MPICH and a topology-aware two-layer implementation. Nicholas T. Karonis, Bronis R. de Supinski, Ian T. Foster, William Gropp, Ewing L. Lusk, John Bresnahan |
IPDPS | 4 |
| 2000 | Performance Modeling and Tuning of an Unstructured Mesh CFD ApplicationabstractThis paper describes performance tuning experiences with a three-dimensional unstructured grid Euler flow code from NASA, which we have reimplemented in the PETSc framework and ported to several large-scale machines, including the ASCI Red and Blue Pacific machines, the SGI Origin, the Cray T3E, and Beowulf clusters. The code achieves a respectable level of performance for sparse problems, typical of scientific and engineering codes based on partial differential equations, and scales well up to thousands of processors. Since the gap between CPU speed and memory access rate is widening, the code is analyzed from a memory-centric perspective (in contrast to traditional flop-orientation) to understand its sequential and parallel performance. Performance tuning is approached on three fronts: data layouts to enhance locality of reference, algorithmic parameters, and parallel programming model. This effort was guided partly by some simple performance models developed for the sparse matrix-vector product operation. William Gropp, Dinesh K. Kaushik, David E. Keyes, Barry Smith 0002 |
SC | 1 |
| 2000 | MPICH-GQ: Quality-of-Service for Message Passing ProgramsabstractParallel programmers typically assume that all resources required for a program’s execution are dedicated to that purpose. However, in local and wide area networks, contention for shared networks, CPUs, and I/O systems can result in significant variations in availability, with consequent adverse effects on overall performance. We describe a new message-passing architecture, MPICH-GQ, that uses quality of service (QoS) mechanisms to manage contention and hence improve performance of message passing interface (MPI) applications. MPICH-GQ combines new QoS specification, traffic shaping, QoS reservation, and QoS implementation techniques to deliver QoS capabilities to the high-bandwidth bursty flows, complex structures, and reliable protocols used in high-performance applications-characteristics very different from the low-bandwidth, constant bit-rate media flows and unreliable protocols for which QoS mechanisms were designed. Results obtained on a differentiated services testbed demonstrate our ability to maintain application performance in the face of heavy network contention. Alain J. Roy, Ian T. Foster, William Gropp, Nicholas T. Karonis, Volker Sander, Brian R. Toonen |
SC | 3 |
| 2000 | From Trace Generation to Visualization: A Performance Framework for Distributed Parallel SystemsabstractIn this paper we describe a trace analysis framework, from trace generation to visualization. It includes a unified tracing facility on IBMâ SPä systems, a self-defining interval file format, an API for framework extensions, utilities for merging and statistics generation, and a visualization tool with preview and multiple time-space diagrams. The trace environment is extremely scalable, and combines MPI events with system activities in the same set of trace files, one for each SMP node. Since the amount of trace data may be very large, utilities are developed to convert and merge individual trace files into a self-defining interval trace file with multiple frame directories. The interval format allows the development of multiple time-space diagrams, such as thread-activity view, processor-activity view, etc., from the same interval file. A visualization tool, Jumpshot, is modified to visualize these views. A statistics utility is developed using the API, along with its graphics viewer. Ching-Farn Eric Wu, Anthony Bolmarcich, Marc Snir, David Wootton, Farid Parpia, Ewing L. Lusk, William Gropp |
SC | 8 |
| 1999 | Achieving High Sustained Performance in an Unstructured Mesh CFD ApplicationabstractThis paper highlights a three-year project by an interdisciplinary team on a legacy F77 computational fluid dynamics code, with the aim of demonstrating that implicit unstructured grid simulations can execute at rates not far from those of explicit structured grid codes, provided attention is paid to data motion complexity and the reuse of data positioned at the levels of the memory hierarchy closest to the processor, in addition to traditional operation count complexity. The demonstration code is from NASA and the enabling parallel hardware and (freely available) software toolkit are from DOE, but the resulting methodology should be broadly applicable, and the hardware limitations exposed should allow programmers and vendors of parallel platforms to focus with greater encouragement on sparse codes with indirect addressing. This snapshot of ongoing work shows a performance of 15 microseconds per degree of freedom to steady-state convergence of Euler flow on a mesh with 2.8 million vertices using 3072 dual-processor nodes of ASCI Red, corresponding to a sustained floating-point rate of 0.227 Tflop/s. W. K. Anderson, William Gropp, Dinesh K. Kaushik, David E. Keyes, Barry Smith 0002 |
SC | 2 |
| 1999 | Parallel computation of three-dimensional nonlinear magnetostatic problemsabstractWe describe a general-purpose parallel code for computing accurate solutions to large computationally demanding, 3D, nonlinear magnetostatic problems. The code, CORAL, is based on a volume integral equation formulation. Using an IBM SP parallel computer and iterative solution methods, we successfully solved the dense linear systems inherent in such formulations. A key component of our work was the use of the PETSc library, which provides parallel portability and access to the latest linear algebra solution technology. Copyright © 1999 John Wiley & Sons, Ltd. William Gropp, Kimmo Forsman, Lauri Kettunen |
Concurr. Pract. Exp. | 2 |
| 1998 | A Case for Using MPI's Derived Datatypes to Improve I/O PerformanceabstractMPI-IO, the I/O part of the MPI-2 standard, is a promising new interface for parallel I/O. A key feature of MPI-IO is that it allows users to access several noncontiguous pieces of data from a file with a single I/O function call by defining file views with derived datatypes. We explain how critical this feature is for high performance, why users must create and use derived datatypes whenever possible, and how it enables implementations to perform optimizations. In particular, we describe two optimizations our MPI-IO implementation, ROMIO, performs: data sieving and collective I/O. We demonstrate the performance and portability of the approach with performance results on five different parallel machines: HP Exemplar, IBM SP, Intel Paragon, NEC SX-4, and SGI Origin2000. Rajeev Thakur, William Gropp, Ewing L. Lusk |
SC | 2 |
| 1998 | Wide-Area Implementation of the Message Passing Interface
Ian T. Foster, Jonathan Geisler, William Gropp, Nicholas T. Karonis, Ewing L. Lusk, George K. Thiruvathukal, Steven Tuecke |
Parallel Comput. | 3 |
| 1997 | A High-Performance MPI Implementation on a Shared-Memory Vector Supercomputer
William Gropp, Ewing L. Lusk |
Parallel Comput. | 1 |
| 1996 | A High-Performance, Portable Implementation of the MPI Message Passing Interface Standard
William Gropp, Ewing L. Lusk, Nathan E. Doss, Anthony Skjellum |
Parallel Comput. | 1 |
| 1993 | Panel - Software Tools for High-Performance Distributed Computing
Vaidy S. Sunderam, Geoffrey C. Fox, Al Geist, William Gropp, Bob Harrison, Adam Kolawa, Michael J. Quinn, Anthony Skjellum |
HPDC | 4 |
| 1993 | Applications-driven parallel I/OabstractNo abstract available. N. Galbreath, William Gropp |
SC | 2 |
| 1990 | CLAM and CLAMShell: An Interactive Front-End for Parallel Computing and Visualization
David E. Foulser, William Gropp |
ICPP (3) | 2 |
| 1990 | Krylov Methods Preconditioned with Incompletely Factored Matrices on the CM-2abstractThe performance is measured of the components of the key interative kernel of a preconditioned Krylov space interative linear system solver. In some sense, these numbers can be regarded as best case timings for these kernels. Sweeps were timed over meshes, sparse triangular solves, and inner products on a large 3-D model problem over a cube shaped domain discretized with a seven point template. The performance of the CM-2 is highly dependent on the use of very specialized programs. These programs mapped a regular problem domain onto the processor topology in a careful manner and used the optimized local NEWS communications network. The rather dramatic deterioration in performance was documented when these ideal conditions no longer apply. A synthetic workload generator was developed to produce and solve a parameterized family of increasingly irregular problems. Harry Berryman, Joel H. Saltz, William Gropp, Ravi Mirchandaney |
J. Parallel Distributed Comput. | 3 |
| 1987 | Solving PDEs on loosely-coupled parallel processors
William Gropp |
Parallel Comput. | 1 |