José E. Moreira

dblp:12/1593 · DBLP profile ↗
← Back
69ranked-venue papers
10as first author
8since 2021 · last 2023
0000-0001-7029-6327ORCID · corroborated

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

Systems, architecture and hardware · 58 · 9 first-author · 7 since 2021Software engineering, systems software and programming languages · 5 · 1 first-author · 2 since 2021Security and privacy · 2Artificial intelligence and machine learning · 1Databases, data management, data science and information retrieval · 1
YearPublicationVenuePosition
2023 Fast matrix multiplication via compiler-only layered data reorganization and intrinsic lowering
abstract
Abstract The resurgence of machine learning has increased the demand for high‐performance basic linear algebra subroutines (BLAS), which have long depended on libraries to achieve peak performance on commodity hardware. High‐performance BLAS implementations rely on a layered approach that consists of tiling and packing layers—for data (re)organization—and micro kernels that perform the actual computations. The algorithm for the tiling and packing layers is target independent but is parameterized to the memory hierarchy and register‐file size. The creation of high‐performance micro kernels requires significant development effort to write tailored assembly code for each architecture. This hand optimization task is complicated by the recent introduction of matrix engines by 's (Matrix Multiply Assist—MMA), (Advanced Matrix eXtensions—AMX), and (Matrix Extensions—ME) to deliver high‐performance matrix operations. This article presents a compiler‐only alternative to the use of high‐performance libraries by incorporating, to the best of our knowledge and for the first time, the automatic generation of the layered approach into LLVM, a production compiler. Modular design of the algorithm, such as the use of LLVM's matrix‐multiply intrinsic for a clear interface between the tiling and packing layers and the micro kernel, makes it easy to retarget the code generation to multiple accelerators. The parameterization of the tiling and packing layers is demonstrated in the generation of code for the MMA unit on IBM's POWER10. This article also describes an algorithm that lowers the matrix‐multiply intrinsic to the MMA unit. The use of intrinsics enables a comprehensive performance study. In processors without hardware matrix engines, the tiling and packing delivers performance up to (Intel)—for small matrices—and more than (POWER9)—for large matrices—faster than PLuTo, a widely used polyhedral optimizer. The performance also approaches high‐performance libraries and is only slower than OpenBLAS and on‐par with Eigen for large matrices. With MMA in POWER10 this solution is, for large matrices, over faster the vector‐extension solution, matches Eigen performance, and achieves up to ofBLASpeak performance.
Braedy Kuzma, Ivan Korostelev, João P. L. de Carvalho, José E. Moreira, Christopher Barton, Guido Araujo, José Nelson Amaral
Softw. Pract. Exp.4
2023 Advancing Direct Convolution Using Convolution Slicing Optimization and ISA Extensions
abstract
Convolution is one of the most computationally intensive operations that must be performed for machine learning model inference. A traditional approach to computing convolutions is known as the Im2Col + BLAS method. This article proposes SConv: a direct-convolution algorithm based on an MLIR/LLVM code-generation toolchain that can be integrated into machine-learning compilers. This algorithm introduces: (a) Convolution Slicing Analysis (CSA)—a convolution-specific 3D cache-blocking analysis pass that focuses on tile reuse over the cache hierarchy; (b) Convolution Slicing Optimization—a code-generation pass that uses CSA to generate a tiled direct-convolution macro-kernel; and (c) Vector-based Packing—an architecture-specific optimized input-tensor packing solution based on vector-register shift instructions for convolutions with unitary stride. Experiments conducted on 393 convolutions from full ONNX-MLIR machine learning models indicate that the elimination of the Im2Col transformation and the use of fast packing routines result in a total packing time reduction, on full model inference, of 2.3×–4.0× on Intel x86 and 3.3×–5.9× on IBM POWER10. The speed-up over an Im2Col + BLAS method based on current BLAS implementations for end-to-end machine-learning model inference is in the range of 11%–27% for Intel x86 and 11%–34% for IBM POWER10 architectures. The total convolution speedup for model inference is 13%–28% on Intel x86 and 23%–39% on IBM POWER10. SConv also outperforms BLAS GEMM, when computing pointwise convolutions in more than 82% of the 219 tested instances.
Victor Ferrari, Rafael C. F. Sousa, Márcio Machado Pereira, João P. L. de Carvalho, José Nelson Amaral, José E. Moreira, Guido Araujo
ACM Trans. Archit. Code Optim.6
2023 YaConv: Convolution with Low Cache Footprint
abstract
This article introduces YaConv , a new algorithm to compute convolution using GEMM microkernels from a Basic Linear Algebra Subprograms library that is efficient for multiple CPU architectures. Previous approaches either create a copy of each image element for each filter element or reload these elements into cache for each GEMM call, leading to redundant instances of the image elements in cache. Instead, YaConv loads each image element once into the cache and maximizes the reuse of these elements. The output image is computed by scattering results of the GEMM microkernel calls to the correct locations in the output image. The main advantage of this new algorithm—which leads to better performance in comparison to the existing im2col approach on several architectures—is a more efficient use of the memory hierarchy. The experimental evaluation on convolutional layers from PyTorch, along with a parameterized study, indicates an average 24% speedup over im2col convolution. Increased performance comes as a result of 3× reduction in L3 cache accesses and 2× fewer branch instructions.
Ivan Korostelev, João P. L. de Carvalho, José E. Moreira, José Nelson Amaral
ACM Trans. Archit. Code Optim.3
2022 Return-oriented programming protection in the IBM POWER10
abstract
Return-oriented programming (ROP) is a technique for hijacking the control-flow of a program and forcing it to perform computations that were never originally intended. ROP is achieved by modifying the values of return addresses saved to memory, causing a failure of control-flow integrity. In the POWER10 processor, we have adopted a cryptographic mechanism for ROP protection. Return addresses are cryptographically hashed when control-flow enters a function and the hash is saved in memory. The hash is recomputed and compared to the saved value just before a return from the function. Any mismatch is flagged as a violation and generates an exception to the supervisor. POWER10 was augmented with instructions to generate and verify those hashes. We minimize performance impact on running programs by implementing the cryptographic hash with dedicated functional units.
José E. Moreira, Debapriya Chatterjee, Kattamuri Ekanadham, Arnold Flores
CF1
2022 Dense dynamic blocks: optimizing SpMM for processors with vector and matrix units using machine learning techniques
abstract
Recent processors have been augmented with matrix-multiply units that operate on small matrices, creating a functional unit-rich environment. These units have been successfully employed on dense matrix operations such as those found in the Basic Linear Algebra Subprograms (BLAS). In this work, we exploit these new matrix-multiply facilities to speed up Sparse Matrix Dense Matrix Multiplications (SpMM) for highly sparse matrices.
Serif Yesil, José E. Moreira, Josep Torrellas
ICS2
2022 Modeling Matrix Engines for Portability and Performance
abstract
Matrix engines, also known as matrix-multiplication accel-erators, capable of computing on 2D matrices of various data types are traditionally found only on GPUs. However, they are increasingly being introduced into CPU architectures to support AI/ML computations. Unlike traditional SIMD functional units, these accelerators require both the input and output data to be packed into a specific 2D-data layout that is often dependent on the input and output data types. Due to the large variety of supported data types and architectures, a common abstraction is required to unify these seemingly disparate accelerators and more efficiently produce high-performance code. In this paper, we show that the hardware characteristics of a vast array of different matrix engines can be unified using a single analytical model that casts matrix engines as an accumulation of multiple outer-products (also known as rank-k updates). This allows us to easily and quickly develop high-performance kernels using matrix engines for different architectures. We demonstrate our matrix engine model and its portability by applying it to two distinct architectures. Using our model, we show that high-performance computational kernels and packing routines required for high-performance dense linear algebra libraries can be easily designed. Furthermore, we show that the performance attained by our implementations is around 90–99 % (80–95 % on large problems) of the theoretical peak throughput of the matrix engines.
Nicholai Tukanov, Rajalakshmi Srinivasaraghavan, José E. Moreira, Tze Meng Low
IPDPS3
2021 Energy Efficiency Boost in the AI-Infused POWER10 Processor
abstract
We present the novel micro-architectural features, supported by an innovative and novel pre-silicon methodology in the design of POWER10. The resulting projected energy efficiency boost over POWER9 is 2.6x at core level (for SPECint) and up to 3x at socket level. In addition, a new feature supporting inline AI acceleration was added to the POWER ISA and incorporated into the POWER10 processor core design. The resulting boost in SIMD/AI socket performance is projected to be up to 10x for FP32 and 21x for INT8 models of ResNet-50 and BERT-Large. In this paper, we describe the novel methodology deployed and used not only to obtain these efficiency boosts for traditional workloads, but also to infuse AI/ML/HPC capability directly into the POWER10 core.
Brian W. Thompto, Dung Q. Nguyen, José E. Moreira, Ramon Bertran Monfort, Hans M. Jacobson, Richard J. Eickemeyer, Rahul M. Rao, Michael Goulet, Marcy Byers, Christopher J. Gonzalez, Karthik Swaminathan, Nagu R. Dhanwada, Silvia M. Müller, Satish Kumar Sadasivam, Robert K. Montoye, William J. Starke, Christian G. Zoellin, Michael S. Floyd, Jeffrey Stuecheli, Nandhini Chandramoorthy, John-David Wellman, Alper Buyuktosunoglu, Matthias Pflanz, Balaram Sinharoy, Pradip Bose
ISCA3
2021 KernelFaRer: Replacing Native-Code Idioms with High-Performance Library Calls
abstract
Well-crafted libraries deliver much higher performance than code generated by sophisticated application programmers using advanced optimizing compilers. When a code pattern for which a well-tuned library implementation exists is found in the source code of an application, the highest performing solution is to replace the pattern with a call to the library. Idiom-recognition solutions in the past either required pattern matching machinery that was outside of the compilation framework or provided a very brittle solution that would fail even for minor variants in the pattern source code. This article introduces Kernel Find & Replacer ( KernelFaRer ), an idiom recognizer implemented entirely in the existing LLVM compiler framework. The versatility of KernelFaRer is demonstrated by matching and replacing two linear algebra idioms, general matrix-matrix multiplication (GEMM), and symmetric rank-2k update (SYR2K). Both GEMM and SYR2K are used extensively in scientific computation, and GEMM is also a central building block for deep learning and computer graphics algorithms. The idiom recognition in KernelFaRer is much more robust than alternative solutions, has a much lower compilation overhead, and is fully integrated in the broadly used LLVM compilation tools. KernelFaRer replaces existing GEMM and SYR2K idioms with computations performed by BLAS, Eigen, MKL (Intel’s x86), ESSL (IBM’s PowerPC), and BLIS (AMD). Gains in performance that reach 2000× over hand-crafted source code compiled at the highest optimization level demonstrate that replacing application code with library call is a performant solution.
João P. L. de Carvalho, Braedy Kuzma, Ivan Korostelev, José Nelson Amaral, Christopher Barton, José E. Moreira, Guido Araujo
ACM Trans. Archit. Code Optim.6
2018 GraphBLAS: handling performance concerns in large graph analytics
abstract
Emerging applications in health-care, social media analytics, cyber-security, homeland security, and marketing require large graph analytics. Attaining good performance on these applications on modern day hardware is challenging because of the complex pipelines and deep memory hierarchy of these machines. In this paper, we review the linear algebra formulation of graph-analytics and show that it effectively handles the separation of performance concerns, best handled by system developers, from application logic concerns.
Manoj Kumar 0006, José E. Moreira, Pratap Pattnaik
CF2
2016 Speeding Up Stencil Computations with Kernel Convolution
abstract
A technique to speed up stencil computation is introduced. Computation and data reuse schemes are developed for its application to 1- and 3-dimensional stencils. The approach traverses the data domain fewer times than a state-of-the-art, straightforward iterative stencil implementation would. Performance results are shown for a variety of platforms, exemplifying how it can be straightforwardly applied with existing techniques and frameworks. The technique, named Aggregate Stencil-Loop Iteration (ASLI), works by applying a stencil obtained by the original stencil operator convolved with itself one or more times. This more complex operator creates new opportunities for in-register data reuse and increases the FLOPs-to-load ratio. The total number of FLOPs decreases for 1D but increases for 2D and 3D star-shaped stencils. In both scenarios, speed-up relative to the state-of-the-art is achieved. ASLI is relatively easy to implement and works synergistically with existing methods to optimize stencil computations.
Guilherme Carvalho Januario, Bryan S. Rosenburg, Yoonho Park, Michael Perrone, José E. Moreira, Tereza Cristina M. B. Carvalho
SBAC-PAD5
2016 Workshop on high-performance computational finance
abstract
The purpose of this special issue is to collate a selection of representative research articles that were primarily presented at the Seventh Workshop on High-Performance Computational Finance, held in conjunction with SC'14. This annual workshop brings together practitioners, researchers, vendors, and scholars from the complementary fields of computational finance and high-performance computing in order to promote exchange of ideas, discuss future collaborations and develop new research directions. Financial companies increasingly rely on high performance computers to analyze high volumes of financial data, automatically execute trades, and manage risk.
Matthew Dixon, José E. Moreira, David Daly
Concurr. Comput. Pract. Exp.2
2013 Design and Implementation of a Scalable Membership Service for Supercomputer Resiliency-Aware Runtime
Yoav Tock, Benjamin Mandler, José E. Moreira, Terry R. Jones
Euro-Par3
2013 Latest trends in computer architectures and parallel and distributed technologies
abstract
ABSTRACT This special issue focuses on new developments in high performance applications, as well as the latest trends in computer architecture and parallel and distributed technologies, and is based on extended, thoroughly revised papers from the 22nd International Symposium on Computer Architecture and High Performance Computing. The authors were invited to provide extended versions of their original papers, taking into account comments and suggestions raised during the peer review process and comments from the audience during the conference. Copyright © 2012 John Wiley & Sons, Ltd.
Bruno Schulze, Vinod E. F. Rebello, José E. Moreira
Concurr. Comput. Pract. Exp.3
2012 Accelerating business analytics applications
abstract
Business text analytics applications have seen rapid growth, driven by the mining of data for various decision making processes. Regular expression processing is an important component of these applications, consuming as much as 50% of their total execution time. While prior work on accelerating regular expression processing has focused on Network Intrusion Detection Systems, business analytics applications impose different requirements on regular expression processing efficiency. We present an analytical model of accelerators for regular expression processing, which includes memory bus-, I/O bus-, and network-attached accelerators with a focus on business analytics applications. Based on this model, we advocate the use of vector-style processing for regular expressions in business analytics applications, leveraging the SIMD hardware available in many modern processors. In addition, we show how SIMD hardware can be enhanced to improve regular expression processing even further. We demonstrate a realized speedup better than 1.8 for the entire range of data sizes of interest. In comparison, the alternative strategies deliver only marginal improvement for large data sizes, while performing worse than the SIMD solution for small data sizes.
Valentina Salapura, Tejas Karkhanis, Priya Nagpurkar, José E. Moreira
HPCA4
2012 Special Issue for the Workshop on High Performance Computational Finance
M. F. Dixon, David Daly, Maria Eleftheriou, José E. Moreira, Kyung Dong Ryu
Concurr. Comput. Pract. Exp.4
2009 Fifth International Workshop on System Management Techniques, Processes, and Services (SMTPS)
abstract
It is our pleasure to welcome all speakers, authors and participants to this fifth edition of the International Workshop on System Management Techniques, Processes, and Services (SMTPS), organized as a 2009 IPDPS (IEEE International Parallel & Distributed Processing Symposium) Workshop.
José E. Moreira
IPDPS1
2008 Scalable server provisioning with HOP-SCOTCH
abstract
The problem of provisioning servers in a cluster infrastructure includes the issues of coordinating access and sharing of physical resources, loading servers with the appropriate software images, supporting storage access to users and applications, and providing basic monitoring and control services for those servers. We have developed a system called HOP-SCOTCH that automates the provisioning process for large clusters. Our solution relies on directory services to implement access control. It uses network boot and managed root disks to control the image of each server. We leverage IBM's global storage architecture to provide storage to users and applications. Finally, interfaces are provided to access the services both programmatically and interactively. We demonstrate the scalable behavior of HOP-SCOTCH by experimenting with a cluster of 40 blade servers. We can provision all 40 servers with a brand new image in under 15 minutes.
David Daly, Marcio A. Silva, José E. Moreira
IPDPS3
2008 Multitoroidal Interconnects For Tightly Coupled Supercomputers
abstract
The processing elements of many modern tightly coupled multicomputers are connected via mesh or toroidal networks. Such interconnects are simple and highly scalable, but suffer from high fragmentation, low utilization, and insufficient fault tolerance when the resources allocated to each job are dedicated. High-dimensional interconnects may be more efficient in certain cases, but are based on complex and expensive components and scale poorly. We present a novel hardware/software architectural approach that detaches the processing elements of the system from the interconnect and augments the traditional toroidal topology to provide additional connectivity options and additional link redundancy. We explore the properties of the new "multitoroidal" topology and the improvements it offers in resource utilization and failure tolerance. We present the results of extensive simulation studies to show that for practically important types of workloads, the resource utilization may be increased by 50 percent and, in certain cases, as much as 100 percent compared to toroidal machines and is, in fact, close to the theoretically optimal case of a full crossbar interconnect. The combined hardware/software architectural innovation is a major significant improvement in resource utilization on top of the state of the art in scheduling algorithm research. Also, multitoroidal multicomputers are able to work under link failure rates of 0.002 failures per week that would shut down toroidal machines. A variant of the multitoroidal architecture is implemented in the Blue Gene/L supercomputer.
Yariv Aridor, Tamar Domany, Oleg Goldshmidt, Yevgeny Kliteynik, Edi Shmueli, José E. Moreira
IEEE Trans. Parallel Distributed Syst.6
2007 Experiences Understanding Performance in a Commercial Scale-Out Environment
Robert W. Wisniewski, Mathieu Desnoyers, Maged M. Michael, José E. Moreira, Doron Shiloach, Livio B. Soares
Euro-Par5
2007 Scalability of the Nutch search engine
abstract
Nutch is an open source search engine that is gaining increasing popularity in the commercial world. The Nutch architecture leads itself to a wide range of parallelization techniques. Multiple backend servers can be used to both partition the corpus of search data, thus increasing the rate of queries serviced, and to increase the size of the search data while preserving the service rate. Alternatively, multiple search engines can operate in parallel, further increasing the query rate. In this paper, we analyze the performance and scalability of various configurations of Nutch. The configurations were implemented as part of the Commercial Scale Out project at IBM Research, and were used to investigate the applicability of scale-out architectures in commercial environments. We conclude that Nutch is highly scalable, with the different configurations behaving differently from a performance perspective.
José E. Moreira, Maged M. Michael, Dilma Da Silva, Doron Shiloach, Parijat Dube, Li Zhang 0002
ICS1
2007 Base Operating System Provisioning and Bringup for a Commercial Supercomputer
abstract
Commercial scale-out is a new research project at IBM research. Its main goal is to investigate and develop technologies for the use of large scale parallelism in commercial applications, eventually leading to a commercial supercomputer. The project leverages and explores the features of IBM's BladeCenter family of products. A significant challenge in using a large cluster of servers is the installation and provisioning of the base operating system in those servers. Compounding this problem is the issue of maintenance of the software image in each server after its provisioning. This paper describes the system we developed to manage the installation, provisioning, and maintenance process for a cluster of blades, providing a base level of functionality to be used by higher level management tools. The system leverages the management facilitation features of BladeCenter, and exploits the network and storage architecture of the commercial scale-out prototype cluster. It uses a single shared root filesystem image to reduce management complexity, and completely automates the process of bringing a new blade into the cluster upon its insertion into a BladeCenter chassiss.
David Daly, Jong Hyuk Choi, José E. Moreira, Amos Waterland
IPDPS3
2007 Scale-up x Scale-out: A Case Study using Nutch/Lucene
abstract
Scale-up solutions in the form of large SMPs have represented the mainstream of commercial computing for the past several years. The major server vendors continue to provide increasingly larger and more powerful machines. More recently, scale-out solutions, in the form of clusters of smaller machines, have gained increased acceptance for commercial computing. Scale-out solutions are particularly effective in high-throughput Web-centric applications. In this paper, we investigate the behavior of two competing approaches to parallelism, scale-up and scale-out, in an emerging search application. Our conclusions show that a scale-out strategy can be the key to good performance even on a scale-up machine. Furthermore, scale-out solutions offer better price/performance, although at an increase in management complexity.
Maged M. Michael, José E. Moreira, Doron Shiloach, Robert W. Wisniewski
IPDPS2
2007 Performance Studies of a WebSphere Application, Trade, in Scale-out and Scale-up Environments
abstract
Scale-out approach, in contrast to scale-up approach (exploring increasing performance by utilizing more powerful shared-memory servers), refers to deployment of applications on a large number of small, inexpensive, but tightly packaged and tightly interconnected servers. Recently, there has been an increasing interest in scale-out approach. The purpose of this study is to discover advantages or disadvantages of scale-out systems with a typical enterprise workload, IBM Trade Performance Benchmark Sample for Websphere application server (a.k.a. Trade6). In this work, through cross system performance comparison, we show that for such workload, scale-out approach has better performance/cost effect. In term of scalability, we show that Websphere application server packages for distributed environment scale well while the possible bottleneck of the application deployment is the database tier. We present preliminary results to show that both database partitioning feature (DPF) and federated database server approaches are not exactly suitable for providing scale-out solution for the database tier of workloads similar to Trade (small tables and short transactions). In addition, we discuss our on-going effort on further performance study: (1) studies of performance/scalability for larger deployments by adopting the IBM AMBIENCE queuing network modeling tool, (2) performance breakdowns utilizing IBM ACTC hardware counter library.
Hao Yu 0008, José E. Moreira, Parijat Dube, I-Hsin Chung, Li Zhang 0002
IPDPS2
2007 Performance Evaluation of a Commercial Application, Trade, in Scale-out Environments
abstract
Scale-out approach, in contrast to scale-up approach (exploring increasing performance by utilizing more powerful shared-memory servers), refers to deployment of applications on a large number of small, inexpensive, but tightly packaged and tightly interconnected servers. The purpose of this study is to understand the performance of scale-out architectures with a typical enterprise workload, IBM Trade Performance Benchmark Sample for WebSphere Application Server (a.k.a. Trade). We describe a performance evaluation methodology that gives accurate predictions of application performance and system utilization by utilizing experimental data driven model development. Through experiments and extrapolation from the derived model, we show that for such workload, WebSphere Application Server packages for distributed environments scale well while the possible bottleneck of the application deployment is the database tier.
Parijat Dube, Hao Yu 0008, Li Zhang 0002, José E. Moreira
MASCOTS4
2006 High performance file I/O for the Blue Gene/L supercomputer
abstract
Parallel 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
HPCA7
2006 A database-centric approach to system management in the Blue Gene/L supercomputer
abstract
In designing the management system for Blue Gene/L, we adopted a database-centric approach. All configuration and operational data for a particular Blue Gene/L system are stored in a relational database that is kept in the system's service node. The database also serves as the communication bus for the various processes implementing the management system. This design offers many advantages, including the ability to use SQL commands to retrieve reliability, availability, and serviceability (RAS) information about the system. Information about machine partitioning and user jobs can be obtained the same way. Leveraging the database, we have developed a Web interface for system management. This management system has been successfully implemented and deployed in all 19 Blue Gene/L installations at the time of this writing
Ralph Bellofatto, Paul G. Crumley, David Darrington, Brant Knudson, Mark Megerian, José E. Moreira, Alda S. Ohmacht, John Orbeck, Don Reed, Greg Stewart
IPDPS6
2006 Blue Gene system software - Design and implementation of a one-sided communication interface for the IBM eServer Blue Gene® supercomputer
abstract
This paper discusses the design and implementation of a one-sided communication interface for the IBM Blue Gene/L supercomputer. This interface facilitates ARMCI and the Global Arrays toolkit and can be used by other one-sided communication libraries. New protocols, interrupt driven communication, and compute node kernel enhancements were required to enable these libraries. Three possible methods for enabling ARMCI on the Blue Gene/L software stack are discussed. A detailed look into the development process shows how the implementation of the one-sided communication interface was completed. This was accomplished on a compressed time scale with the collaboration of various organizations within IBM and open source communities. In addition to enabling the one-sided libraries, bandwidth enhancements were made for communication along a diagonal on the Blue Gene/L torus network. The maximum bandwidth improved by a factor of three. This work will enable a variety of one-sided applications to run on Blue Gene/L.
Michael Blocksome, Charles Archer, Todd Inglett, Patrick McCarthy, Michael B. Mundy, Joe Ratterman, A. Sidelnik, Brian E. Smith, Gheorghe Almási 0001, José G. Castaños, Derek Lieber, José E. Moreira, Sriram Krishnamoorthy, Vinod Tipparaju, Jarek Nieplocha
SC12
2006 Blue Gene system software - Designing a highly-scalable operating system: the Blue Gene/L story
abstract
Blue Gene/L is currently the world's fastest and most scalable supercomputer. It has demonstrated essentially linear scaling all the way to 131,072 processors in several benchmarks and real applications. The operating systems for the compute and I/O nodes of Blue Gene/L, are among the components responsible for that scalability. Compute nodes are dedicated to running application processes, whereas I/O nodes are dedicated to performing system functions. The operating systems adopted for each of these nodes reflect this separation of function. Compute nodes run a lightweight operating system called the compute node kernel. I/O nodes run a port of the Linux operating system. This paper discusses the architecture and design of this solution for Blue Gene/L in the context of the hardware characteristics that led to the design decisions. It also explains and demonstrates how those decisions are instrumental in achieving the performance and scalability for which Blue Gene/L is famous.
José E. Moreira, Michael Brutman, José G. Castaños, Thomas Engelsiepen, Mark Giampapa, Thomas Gooding, Roger L. Haskin, Todd Inglett, Derek Lieber, Patrick McCarthy, Michael B. Mundy, Jeff Parker, Brian P. Wallenfelt
SC1
2006 Blue Gene system software - Topology mapping for Blue Gene/L supercomputer
abstract
Mapping virtual processes onto physical processos is one of the most important issues in parallel computing. The problem of mapping of processes/tasks onto processors is equivalent to the graph embedding problem which has been studied extensively. Although many techniques have been proposed for embeddings of two-dimensional grids, hypercubes, etc., there are few efforts on embeddings of three-dimensional grids and tori. Motivated for better support of task mapping for Blue Gene/L supercomputer, in this paper, we present embedding and integration techniques for the embeddings of three-dimensional grids and tori. The topology mapping library that based on such techniques generates high-quality embeddings of two/three-dimensional grids/tori. In addition, the library is used in BG/L MPI library for scalable support of MPI topology functions. With extensive empirical studies on large scale systems against popular benchmarks and real applications, we demonstrate that the library can significantly improve the communication performance and the scalability of applications.
Hao Yu 0008, I-Hsin Chung, José E. Moreira
SC3
2005 Filtering Failure Logs for a BlueGene/L Prototype
abstract
The growing computational and storage needs of several scientific applications mandate the deployment of extreme-scale parallel machines, such as IBM's BlueGene/L, which can accommodate as many as 128K processors. In this paper, we present our experiences in collecting and filtering error event logs from a 8192 processor BlueGene/L prototype at IBM Rochester, which is currently ranked #8 in the Top-500 list. We analyze the logs collected from this machine over a period of 84 days starting from August 26, 2004. We perform a three-step filtering algorithm on these logs: extracting and categorizing failure events; temporal filtering to remove duplicate reports from the same location; and finally coalescing failure reports of the same error across different locations. Using this approach, we can substantially compress these logs, removing over 99.96% of the 828,387 original entries, and more accurately portray the failure occurrences on this system.
Yinglung Liang, Yanyong Zhang, Anand Sivasubramaniam, Ramendra K. Sahoo, José E. Moreira, Manish Gupta 0002
DSN5
2005 Probabilistic QoS Guarantees for Supercomputing Systems
abstract
Supercomputing systems must be able to reliably and efficiently complete their assigned workloads, even in the presence of failures. This paper proposes a system that allows the system and users to negotiate a mutually desirable risk strategy; in order to accomplish this, the system makes probabilistic guarantees on quality of service (QoS), of the form, "Job j can be completed by deadline d with probability p". In order to make such guarantees, the system uses event prediction (forecasting) in conjunction with fault-aware job scheduling and cooperative checkpointing strategies. Using job logs and failure traces from actual high performance computing systems, we employ trace-based simulations to assess the effects of the prediction accuracy (a) and user risk strategy (U) on a variety of performance metrics. Compared to a system that does not use event prediction, a high forecasting accuracy resulted in QoS and utilization improvements of as much as 6%, along with an 89% reduction in the amount of lost work. Therefore, our results show that a system that makes probabilistic QoS guarantees using a market-based scheduling approach can increase both system performance and reliability.
Adam J. Oliner, Larry Rudolph, Ramendra K. Sahoo, José E. Moreira, Manish Gupta 0002
DSN4
2005 Early Experience with Scientific Applications on the Blue Gene/L Supercomputer
Gheorghe Almási 0001, Gyan Bhanot, Dong Chen 0005, Maria Eleftheriou, Blake G. Fitch, Alan Gara, Robert S. Germain, John A. Gunnels, Manish Gupta 0002, Philip Heidelberger, Michael Pitman, Aleksandr Rayshubskiy, James C. Sexton, Frank Suits, Pavlos Vranas, Robert Walkup, T. J. Christopher Ward, Yuriy Zhestkov, Alessandro Curioni, Wanda Andreoni, Charles Archer, José E. Moreira, Richard Loft, Henry M. Tufo, Theron Voran, Katherine Riley
Euro-Par22
2005 The Evolution of the Blue Gene/L Supercomputer
José E. Moreira
Euro-Par1
2005 Scaling physics and material science applications on a massively parallel Blue Gene/L system
abstract
Blue Gene/L represents a new way to build supercomputers, using a large number of low power processors, together with multiple integrated interconnection networks. Whether real applications can scale to tens of thousands of processors (on a machine like Blue Gene/L) has been an open question. In this paper, we describe early experience with several physics and material science applications on a 32,768 node Blue Gene/L system, which was installed recently at the Lawrence Livermore National Laboratory. Our study shows some problems in the applications and in the current software implementation, but overall, excellent scaling of these applications to 32K nodes on the current Blue Gene/L system. While there is clearly room for improvement, these results represent the first proof point that MPI applications can effectively scale to over ten thousand processors. They also validate the scalability of the hardware and software architecture of Blue Gene/L.
Gheorghe Almási 0001, Gyan Bhanot, Alan Gara, Manish Gupta 0002, James C. Sexton, Robert Walkup, Vasily V. Bulatov, Andrew W. Cook, Bronis R. de Supinski, James N. Glosli, Jeffrey A. Greenough, François Gygi, Alison Kubota, Steve Louis, Thomas E. Spelce, Frederick H. Streitz, Peter L. Williams, Robert K. Yates, Charles Archer, José E. Moreira, Charles A. Rendleman
ICS20
2005 Optimization of MPI collective communication on BlueGene/L systems
abstract
BlueGene/L is currently the world's fastest supercomputer. It consists of a large number of low power dual-processor compute nodes interconnected by high speed torus and collective networks, Because compute nodes do not have shared memory, MPI is the the natural programming model for this machine. The BlueGene/L MPI library is a port of MPICH2.In this paper we discuss the implementation of MPI collectives on BlueGene/L. The MPICH2 implementation of MPI collectives is based on point-to-point communication primitives. This turns out to be suboptimal for a number of reasons. Machine-optimized MPI collectives are necessary to harness the performance of BlueGene/L. We discuss these optimized MPI collectives, describing the algorithms and presenting performance results measured with targeted micro-benchmarks on real BlueGene/L hardware with up to 4096 compute nodes.
Gheorghe Almási 0001, Philip Heidelberger, Charles Archer, Xavier Martorell, C. Christopher Erway, José E. Moreira, Burkhard D. Steinmacher-Burow, Yili Zheng
ICS6
2005 Open Job Management Architecture for the Blue Gene/L Supercomputer
Yariv Aridor, Tamar Domany, Oleg Goldshmidt, Yevgeny Kliteynik, José E. Moreira, Edi Shmueli
JSSPP5
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-Par7
2004 Adaptive incremental checkpointing for massively parallel systems
abstract
Given the scale of massively parallel systems, occurrence of faults is no longer an exception but a regular event. Periodic checkpointing is becoming increasingly important in these systems. However, huge memory footprints of parallel applications place severe limitations on scalability of normal checkpointing techniques. Incremental checkpointing is a well researched technique that addresses scalability concerns, but most of the implementations require paging support from hardware and the underlying operating system, which may not be always available. In this paper, we propose a software based adaptive incremental checkpoint technique which uses a secure hash function to uniquely identify changed blocks in memory. Our algorithm is the first self-optimizing algorithm that dynamically computes the optimal block boundaries, based on the history of changed blocks. This provides better opportunities for minimizing checkpoint file size. Since the hash is computed in software, we do not need any system support for this. We have implemented and tested this mechanism on the BlueGene/L system. Our results on several well-known benchmarks are encouraging, both in terms of reduction in average checkpoint file size and adaptivity towards application’s memory access patterns.
Rahul Garg 0001, Meeta Sharma Gupta, José E. Moreira
ICS4
2004 Fault-Aware Job Scheduling for BlueGene/L Systems
abstract
Summary form only given. Large-scale systems like BlueGene/L are susceptible to a number of software and hardware failures that can affect system performance. We evaluate the effectiveness of a previously developed job scheduling algorithm for BlueGene/L in the presence of faults. We have developed two new job-scheduling algorithms considering failures while scheduling the jobs. We have also evaluated the impact of these algorithms on average bounded slowdown, average response time and system utilization, considering different levels of proactive failure prediction and prevention techniques reported in the literature. Our simulation studies show that the use of these new algorithms with even trivial fault prediction confidence or accuracy levels (as low as 10%) can significantly improve the performance of the BlueGene/L system.
Adam J. Oliner, Ramendra K. Sahoo, José E. Moreira, Manish Gupta 0002, Anand Sivasubramaniam
IPDPS3
2004 The BlueGene/L pseudo cycle-accurate simulator
abstract
The design and development of a new computer system is a lengthy process, with a considerable amount of time elapsed between the beginning of development and first hardware availability. Hence, fast and reasonably accurate simulation of processor architecture has become critical as an enabling mechanism for software engineers to develop and tune system software and applications. In this paper, we present the time-stamped timing model extensions to the BlueGene/L functional simulator. These extensions were implemented to create a pseudo cycle-accurate simulator capable of providing tracing capabilities for detection of bottlenecks and for performance tuning of applications, before the actual hardware became available. Our validation tests, using the DAXPY kernel and the serial version of the NAS benchmarks, show that our pseudo cycle-accurate simulator provides timing information within 15% of the times measured using the actual BlueGene/L hardware. In addition, we present a couple of case studies, which describes how this simulator can be used for identification of performance bottlenecks and for application tuning.
Leonardo R. Bachega, José R. Brunheroto, Luiz De Rose, Pedro Mindlin, José E. Moreira
ISPASS5
2004 Multi-toroidal Interconnects: Using Additional Communication Links to Improve Utilization of Parallel Computers
Yariv Aridor, Tamar Domany, Oleg Goldshmidt, Edi Shmueli, José E. Moreira, Larry Stockmeier
JSSPP5
2004 Unlocking the Performance of the BlueGene/L Supercomputer
abstract
The BlueGene/L supercomputer is expected to deliver new levels of application performance by providing a combination of good single-node computational performance and high scalability. To achieve good single-node performance, the BlueGene/L design includes a special dual floating-point unit on each processor and the ability to use two processors per node. BlueGene/L also includes both a torus and a tree network to achieve high scalability. We demonstrate how benchmarks and applications can take advantage of these architectural features to get the most out of BlueGene/L.
Gheorghe Almási 0001, Siddhartha Chatterjee, Alan Gara, John A. Gunnels, Manish Gupta 0002, Amy Henning, José E. Moreira, Robert Walkup
SC7
2003 An Overview of the Blue Gene/L System Software Organization
Gheorghe Almási 0001, Ralph Bellofatto, José R. Brunheroto, Calin Cascaval, José G. Castaños, Luis Ceze, Paul Crumley, C. Christopher Erway, Joseph Gagliano, Derek Lieber, Xavier Martorell, José E. Moreira, Alda Sanomiya, Karin Strauss
Euro-Par12
2003 Obtaining Hardware Performance Metrics for the BlueGene/L Supercomputer
Pedro Mindlin, José R. Brunheroto, Luiz De Rose, José E. Moreira
Euro-Par4
2003 A Volumetric FFT for BlueGene/L
Maria Eleftheriou, José E. Moreira, Blake G. Fitch, Robert S. Germain
HiPC2
2003 Gang Scheduling Extensions for I/O Intensive Workloads
Yanyong Zhang, Antony Yang, Anand Sivasubramaniam, José E. Moreira
JSSPP4
2003 Critical event prediction for proactive management in large-scale computer clusters
abstract
As the complexity of distributed computing systems increases, systems management tasks require significantly higher levels of automation; examples include diagnosis and prediction based on real-time streams of computer events, setting alarms, and performing continuous monitoring. The core of autonomic computing, a recently proposed initiative towards next-generation IT-systems capable of 'self-healing', is the ability to analyze data in real-time and to predict potential problems. The goal is to avoid catastrophic failures through prompt execution of remedial actions.This paper describes an attempt to build a proactive prediction and control system for large clusters. We collected event logs containing various system reliability, availability and serviceability (RAS) events, and system activity reports (SARs) from a 350-node cluster system for a period of one year. The 'raw' system health measurements contain a great deal of redundant event data, which is either repetitive in nature or misaligned with respect to time. We applied a filtering technique and modeled the data into a set of primary and derived variables. These variables used probabilistic networks for establishing event correlations through prediction algorithms. We also evaluated the role of time-series methods, rule-based classification algorithms and Bayesian network models in event prediction.Based on historical data, our results suggest that it is feasible to predict system performance parameters (SARs) with a high degree of accuracy using time-series models. Rule-based classification techniques can be used to extract machine-event signatures to predict critical events with up to 70% accuracy.
Ramendra K. Sahoo, Adam J. Oliner, Irina Rish, Manish Gupta 0002, José E. Moreira, Sheng Ma, Ricardo Vilalta, Anand Sivasubramaniam
KDD5
2003 Enabling Dual-Core Mode in BlueGene/L: Challenges and Solutions
abstract
BlueGene/L is a massively parallel computer system with 65536 dual-processor compute nodes. The peak performance of BlueGene/L is in excess of 360 TFLOP/s if both processor cores in a node are used for computation. The main challenge of deploying this dual-core mode of operation is that the L1 caches in each core are not hardware coherent. This forces a software-based approach to cache coherence and guides our design of a programming model for dual-core mode. We describe the design, implementation, and performance evaluation of system software for enabling the use of dual-core mode on BlueGene/L. Our preliminary performance results show that our approach to dual-core mode is effective for key numerical kernels.
George S. Almási, Leonardo R. Bachega, Siddhartha Chatterjee, Manish Gupta 0002, Derek Lieber, Xavier Martorell, José E. Moreira
SBAC-PAD7
2003 Supporting multidimensional arrays in Java
abstract
Abstract The lack of direct support for multidimensional arrays in JavaTM has been recognized as a major deficiency in the language's applicability to numerical computing. It has been shown that, when augmented with multidimensional arrays, Java can achieve very high‐performance for numerical computing through the use of compiler techniques and efficient implementations of aggregate array operations. Three approaches have been discussed in the literature for extending Java with support for multidimensional arrays: class libraries that implement these structures; extending the Java language with new syntactic constructs for multidimensional arrays that are directly translated to bytecode; and relying on the Java Virtual Machine to recognize those arrays of arrays that are being used to simulate multidimensional arrays. This paper presents a balanced presentation of the pros and cons of each technique in the areas of functionality, language and virtual machine impact, implementation effort, and effect on performance. We show that the best choice depends on the relative importance attached to the different metrics, and thereby provide a common ground for a rational discussion and comparison of the techniques. Copyright © 2003 John Wiley & Sons, Ltd.
José E. Moreira, Samuel P. Midkiff, Manish Gupta 0002
Concurr. Comput. Pract. Exp.1
2003 An Integrated Approach to Parallel Scheduling Using Gang-Scheduling, Backfilling, and Migration
abstract
Effective scheduling strategies to improve response times, throughput, and utilization are an important consideration in large supercomputing environments. Parallel machines in these environments have traditionally used space-sharing strategies to accommodate multiple jobs at the same time by dedicating the nodes to a single job until it completes. This approach, however, can result in low system utilization and large job wait times. This paper discusses three techniques that can be used beyond simple space-sharing to improve the performance of large parallel systems. The first technique we analyze is backfilling, the second is gang-scheduling, and the third is migration. The main contribution of this paper is an analysis of the effects of combining the above techniques. Using extensive simulations based on detailed models of realistic workloads, the benefits of combining the various techniques are shown over a spectrum of performance criteria.
Yanyong Zhang, Hubertus Franke, José E. Moreira, Anand Sivasubramaniam
IEEE Trans. Parallel Distributed Syst.3
2002 Blue Gene/L, a System-On-A-Chip
abstract
Summary form only given. Large powerful networks coupled to state-of-the-art processors have traditionally dominated supercomputing. As technology advances, this approach is likely to be challenged by a more cost-effective System-On-A-Chip approach, with higher levels of system integration. The scalability of applications to architectures with tens to hundreds of thousands of processors is critical to the success of this approach. Significant progress has been made in mapping numerous compute-intensive applications, many of them grand challenges, to parallel architectures. Applications hoping to efficiently execute on future supercomputers of any architecture must be coded in a manner consistent with an enormous degree of parallelism. The BG/L program is developing a peak nominal 180 TFLOPS (360 TFLOPS for some applications) supercomputer to serve a broad range of science applications. BG/L generalizes QCDOC, the first System-On-A-Chip supercomputer that is expected in 2003. BG/L consists of 65,536 nodes, and contains five integrated networks: a 3D torus, a combining tree, a Gb Ethernet network, barrier/global interrupt network and JTAG.
George S. Almási, Daniel K. Beece, Ralph Bellofatto, Gyan Bhanot, Randy Bickford, Matthias A. Blumrich, Arthur A. Bright, José R. Brunheroto, Calin Cascaval, José G. Castaños, Luis Ceze, Paul Coteus, Siddhartha Chatterjee, Dong Chen 0005, George L.-T. Chiu, Thomas M. Cipolla, Paul Crumley, Alina Deutsch, Marc Boris Dombrowa, Wilm E. Donath, Maria Eleftheriou, Blake G. Fitch, Joseph Gagliano, Alan Gara, Robert S. Germain, Mark Giampapa, Manish Gupta 0002, Fred G. Gustavson, Shawn Hall, Ruud A. Haring, David F. Heidel, Philip Heidelberger, Lorraine M. Herger, Dirk Hoenicke, T. Jamal-Eddine, Gerard V. Kopcsay, Alphonso P. Lanzetta, Derek Lieber, M. Lu, Mark P. Mendell, Lawrence S. Mok, José E. Moreira, Ben J. Nathanson, Matthew Newton, Martin Ohmacht, Rick A. Rand, Richard D. Regan, Ramendra K. Sahoo, Alda Sanomiya, Eugen Schenfeld, Sarabjeet Singh, Peilin Song, Burkhard D. Steinmacher-Burow, Karin Strauss, Richard A. Swetz, Todd Takken, R. Brett Tremaine, Mickey Tsao, Pavlos Vranas, T. J. Christopher Ward, Michael E. Wazlowski, J. Brown, Thomas A. Liebsch, A. Schram, G. Ulsh
CLUSTER42
2002 Job Scheduling for the BlueGene/L System (Research Note)
Elie Krevat, José G. Castaños, José E. Moreira
Euro-Par3
2002 Evaluation of a Multithreaded Architecture for Cellular Computing
abstract
Cyclops is a new architecture for high-performance parallel computers that is being developed at the IBM T. J. Watson Research Center. The basic cell of this architecture is a single-chip SMP (symmetric multiprocessor) system with multiple threads of execution, embedded memory and integrated communications hardware. Massive intra-chip parallelism is used to tolerate memory and functional unit latencies. Large systems with thousands of chips can be built by replicating this basic cell in a regular pattern. In this paper, we describe the Cyclops architecture and evaluate two of its new hardware features: a memory hierarchy with a flexible cache organization and fast barrier hardware. Our experiments with the STREAM benchmark show that a particular design can achieve a sustainable memory bandwidth of 40 GB/s, equal to the peak hardware bandwidth and similar to the performance of a 128-processor SGI Origin 3800. For small vectors, we have observed in-cache bandwidth above 80 GB/s. We also show that the fast barrier hardware can improve the performance of the Splash-2 FFT kernel by up to 10%. Our results demonstrate that the Cyclops approach of integrating a large number of simple processing elements and multiple memory banks in the same chip is an effective alternative for designing high-performance systems.
Calin Cascaval, José G. Castaños, Luis Ceze, Monty Denneau, Manish Gupta 0002, Derek Lieber, José E. Moreira, Karin Strauss, Henry S. Warren Jr.
HPCA7
2002 Job Scheduling for the BlueGene/L System
Elie Krevat, José G. Castaños, José E. Moreira
JSSPP3
2002 An overview of the BlueGene/L Supercomputer
abstract
This paper gives an overview of the BlueGene/L Supercomputer. This is a jointly funded research partnership between IBM and the Lawrence Livermore National Laboratory as part of the United States Department of Energy ASCI Advanced Architecture Research Program. Application performance and scaling studies have recently been initiated with partners at a number of academic and government institutions,including the San Diego Supercomputer Center and the California Institute of Technology. This massively parallel system of 65,536 nodes is based on a new architecture that exploits system-on-a-chip technology to deliver target peak processing power of 360 teraFLOPS (trillion floating-point operations per second). The machine is scheduled to be operational in the 2004-2005 time frame, at price/performance and power consumption/performance targets unobtainable with conventional architectures.
Narasimha R. Adiga, Gheorghe Almási 0001, George S. Almási, Yariv Aridor, Rajkishore Barik, Daniel K. Beece, Ralph Bellofatto, Gyan Bhanot, Randy Bickford, Matthias A. Blumrich, Arthur A. Bright, José R. Brunheroto, Calin Cascaval, José G. Castaños, Waiman Chan, Luis Ceze, Paul Coteus, Siddhartha Chatterjee, Dong Chen 0005, George L.-T. Chiu, Thomas M. Cipolla, Paul Crumley, K. M. Desai, Alina Deutsch, Tamar Domany, Marc Boris Dombrowa, Wilm E. Donath, Maria Eleftheriou, C. Christopher Erway, J. Esch, Blake G. Fitch, Joseph Gagliano, Alan Gara, Rahul Garg 0001, Robert S. Germain, Mark Giampapa, Balaji Gopalsamy, John A. Gunnels, Manish Gupta 0002, Fred G. Gustavson, Shawn Hall, Ruud A. Haring, David F. Heidel, Philip Heidelberger, Lorraine M. Herger, Dirk Hoenicke, R. D. Jackson, T. Jamal-Eddine, Gerard V. Kopcsay, Elie Krevat, Manish P. Kurhekar, Alphonso P. Lanzetta, Derek Lieber, L. K. Liu, M. Lu, Mark P. Mendell, A. Misra, Yosef Moatti, Lawrence S. Mok, José E. Moreira, Ben J. Nathanson, Matthew Newton, Martin Ohmacht, Adam J. Oliner, Vinayaka Pandit, R. B. Pudota, Rick A. Rand, Richard D. Regan, Bradley Rubin, Albert E. Ruehli, Silvius Vasile Rus, Ramendra K. Sahoo, Alda Sanomiya, Eugen Schenfeld, M. Sharma, Edi Shmueli, Sarabjeet Singh, Peilin Song, Vijay Srinivasan, Burkhard D. Steinmacher-Burow, Karin Strauss, Christopher W. Surovic, Richard A. Swetz, Todd Takken, R. Brett Tremaine, Mickey Tsao, Arun R. Umamaheshwaran, P. Verma, Pavlos Vranas, T. J. Christopher Ward, Michael E. Wazlowski, W. Barrett, C. Engel, B. Drehmel, B. Hilgart, D. Hill, F. Kasemkhani, David J. Krolak, Chun-Tao Li 0001, Thomas A. Liebsch, James A. Marcella, A. Muff, A. Okomo, M. Rouse, A. Schram, M. Tubbs, G. Ulsh, Charles D. Wait, J. Wittrup, Myung Bae, Kenneth A. Dockser, Lynn Kissel, Mark K. Seager, Jeffrey S. Vetter, K. Yates
SC60
2002 Modeling and analysis of dynamic coscheduling in parallel and distributed environments
abstract
Scheduling in large-scale parallel systems has been and continues to be an important and challenging research problem. Several key factors, including the increasing use of off-the-shelf clusters of workstations to build such parallel systems, have resulted in the emergence of a new class of scheduling strategies, broadly referred to as dynamic coscheduling. Unfortunately, the size of both the design and performance spaces of these emerging scheduling strategies is quite large, due in part to the numerous dynamic interactions among the different components of the parallel computing environment as well as the wide range of applications and systems that can comprise the parallel environment. This in turn makes it difficult to fully explore the benefits and limitations of the various proposed dynamic coscheduling approaches for large-scale systems solely with the use of simulation and/or experimentation.To gain a better understanding of the fundamental properties of different dynamic coscheduling methods, we formulate a general mathematical model of this class of scheduling strategies within a unified framework that allows us to investigate a wide range of parallel environments. We derive a matrix-analytic analysis based on a stochastic decomposition and a fixed-point iteration. A large number of numerical experiments are performed in part to examine the accuracy of our approach. These numerical results are in excellent agreement with detailed simulation results. Our mathematical model and analysis is then used to explore several fundamental design and performance tradeoffs associated with the class of dynamic coscheduling policies across a broad spectrum of parallel computing environments.
Mark S. Squillante, Yanyong Zhang, Anand Sivasubramaniam, Natarajan Gautam, Hubertus Franke, José E. Moreira
SIGMETRICS6
2001 Demonstrating the scalability of a molecular dynamics application on a Petaflop computer
abstract
The IBM Blue Gene project has endeavored into the development of a cellular architecture computer with millions of concurrent threads of execution. One of the major challenges of this project is demonstrating that applications can successfully exploit this massive amount of parallelism. Starting from the sequential version of a well known molecular dynamics code, we developed a new application that exploits the multiple levels of parallelism in the Blue Gene cellular architecture. We perform both analytical and simulation studies of the behavior of this application when executed on a very large number of threads. As a result, we demonstrate that this class of applications can execute efficiently on a large cellular machine.
George S. Almási, Calin Cascaval, José G. Castaños, Monty Denneau, Wilm E. Donath, Maria Eleftheriou, Mark Giampapa, C. T. Howard Ho, Derek Lieber, José E. Moreira, Dennis M. Newns, Marc Snir, Henry S. Warren Jr.
ICS10
2001 An Integrated Approach to Parallel Scheduling Using Gang-Scheduling, Backfilling, and Migration
Yanyong Zhang, Hubertus Franke, José E. Moreira, Anand Sivasubramaniam
JSSPP3
2001 Impact of Workload and System Parameters on Next Generation Cluster Scheduling Mechanisms
abstract
Scheduling of processes onto processors of a parallel machine has always been an important and challenging area of research. The issue becomes even more crucial and difficult as we gradually progress to the use of off-the-shelf workstations, operating systems, and high bandwidth networks to build cost-effective clusters for demanding applications. Clusters are gaining acceptance not just in scientific applications that need supercomputing power, but also in domains such as databases, web service, and multimedia which place diverse Quality-of-Service (QoS) demands on the underlying system. Further, these applications have diverse characteristics in terms of their computation, communication, and I/O requirements, making conventional parallel scheduling solutions, such as space sharing or gang scheduling, unattractive. At the same time, leaving it to the native operating system of each node to make decisions independently can lead to ineffective use of system resources whenever there is communication. Instead, an emerging class of dynamic coscheduling mechanisms that attempt to take remedial actions to guide the system toward coscheduled execution without requiring explicit synchronization offers a lot of promise for cluster scheduling. Using a detailed simulator, this paper evaluates the pros and cons of different dynamic coscheduling alternatives while comparing their advantages over traditional gang scheduling (and not performing any coordinated scheduling at all). The impact of dynamic job arrivals, job characteristics, and different system parameters on these alternatives is evaluated in terms of several performance criteria. In addition, heuristics to enhance one of the alternatives even further are identified, classified, and evaluated. It is shown that these heuristics can significantly outperform the other alternatives over a spectrum of workload and system parameters and is thus a much better option for clusters than conventional gang scheduling.
Yanyong Zhang, Anand Sivasubramaniam, José E. Moreira, Hubertus Franke
IEEE Trans. Parallel Distributed Syst.3
2000 The Impact of Migration on Parallel Job Scheduling for Distributed Systems
Yanyong Zhang, Hubertus Franke, José E. Moreira, Anand Sivasubramaniam
Euro-Par3
2000 Automatic loop transformations and parallelization for Java
abstract
From a software engineering perspective, the Java programming language provides an attractive platform for writing numerically intensive applications. A major drawback hampering its widespread adoption in this domain has been its poor performance on numerical codes. This paper describes a prototype Java compiler which demonstrates that it is possible to achieve performance levels approaching those of current state-of-the-art C, C++ and Fortran compilers on numerical codes. We describe a new transformation called alias versioning that takes advantage of the simplicity of pointers in Java. This transformation, combined with other techniques that we have developed, enables the compiler to perform high order loop transformations (for better data locality) and parallelization completely automatically. We believe that our compiler is the first to have such capabilities of optimizing numerical Java codes. We achieve, with Java, between 80 and 100% of the performance of highly optimized Fortran code in a variety of benchmarks. Furthermore, the automatic parallelization achieves speedups of up to 3.8 on four processors. Combining this compiler technology with packages containing the features expected by programmers of numerical applications would enable Java to become a serious contender for implementing new numerical applications.
Pedro V. Artigas, Manish Gupta 0002, Samuel P. Midkiff, José E. Moreira
ICS4
2000 A simulation-based study of scheduling mechanisms for a dynamic cluster environment
abstract
Scheduling of processes onto processors of a parallel machine has always been an important and challenging area of research. The issue becomes even more crucial and difficult as we gradually progress to the use of off-the-shelf workstations, operating systems, and high bandwidth networks to build cost-effective clusters for demanding applications. Clusters are gaining acceptance not just in scientific applications that need supercomputing power, but also in domains such as databases, web service and multimedia, which place diverse Quality-of-Service (QoS) demands on the underlying system. Further, these applications have diverse characteristics in terms of their computation, communication and I/O requirements, making conventional parallel scheduling solutions, such as space sharing or coscheduling, an unattractive option. At the same time, leaving it to the native operating system of each node to make decisions independently can lead to ineffective use of system resources whenever there is communication. Instead, an emerging class of dynamic coscheduling mechanisms, that attempt to take remedial actions to guide the system towards coscheduled execution without requiring explicit synchronization, offer a lot of promise for cluster scheduling. Using a detailed simulator, this paper evaluates the pros and cons of different dynamic coscheduling alternatives, while comparing their advantages over traditional coscheduling (and not performing any coordinated scheduling at all). The impact of dynamic job arrivals, job characteristics and different system parameters on these alternatives are evaluated in terms of several performance criteria.
Yanyong Zhang, Anand Sivasubramaniam, José E. Moreira, Hubertus Franke
ICS3
2000 Improving Parallel Job Scheduling by Combining Gang Scheduling and Backfilling Techniques
abstract
Two different approaches have been commonly used to address problems associated with space sharing scheduling strategies: (a) augmenting space sharing with backfilling, which performs out of order job scheduling; and (b) augmenting space sharing with time sharing, using a technique called coscheduling or gang scheduling. With three important experimental results-impact of priority queue order on backfilling, impact of overestimation of job execution times, and comparison of scheduling techniques-this paper presents an integrated strategy that combines backfilling with gang scheduling. Using extensive simulations based on detailed models of realistic workloads, the benefits of combining backfilling and gang scheduling are clearly demonstrated over a spectrum of performance criteria.
Yanyong Zhang, Anand Sivasubramaniam, Hubertus Franke, José E. Moreira
IPDPS4
2000 From flop to megaflops: Java for technical computing
abstract
Although there has been some experimentation with Java as a language for numerically intensive computing, there is a perception by many that the language is unsuited for such work because of performance deficiencies. In this article we show how optimizing array bounds checks and null pointer checks creates loop nests on which aggressive optimizations can be used. Applying these optimizations by hand to a simple matrix-multiply test case leads to Java-compliant programs whose performance is in excess of 500 Mflops on a four-processor 332MHz RS/6000 model F50 computer. We also report in this article the effect that various optimizations have on the performance of six floating-point-intensive benchmarks. Through these optimizations we have been able to achieve with Java at least 80% of the peak Fortran performance on the same benchmarks. Since all of these optimizations can be automated, we conclude that Java will soon be a serious contender for numerically intensive computing.
José E. Moreira, Samuel P. Midkiff, Manish Gupta 0002
ACM Trans. Program. Lang. Syst.1
1999 An Evaluation of Parallel Job Scheduling for ASCI Blue-Pacific
abstract
In this paper we analyze the behavior of a gang-scheduling strategy that we are developing for the ASCI Blue-Pacific machines. Using actual job logs for one of the ASCI machines we generate a statistical model of the current workload with hyper Erlang distributions. We then vary the parameters of those distributions to generate various workloads, representative of different operating points of the machine. Through simulation we obtain performance parameters for three different scheduling strategies: (i) first-come first-serve, (ii) gang-scheduling, and (iii) backfilling. Our results show that backfilling, can be very effective for the common operating points in the 60-70% utilization range. However, for higher utilization rates, time-sharing techniques such as gang-scheduling offer much better performance.
Hubertus Franke, Joefon Jann, José E. Moreira, Pratap Pattnaik, Morris A. Jette
SC3
1999 High Performance Computing with the Array Package for Java: A Case Study using Data Mining
abstract
This paper discusses several techniques used in developing a parallel, production quality data mining application in Java.We started by developing three sequential versions of a product recommendation data mining application: (i) a Fortran 90 version used as a performance reference, (ii) a plain Java implementation that only uses the primitive array structures from the language, and (iii) a baseline Java implementation that uses our Array package for Java.This Array package provides parallelism at the level of individual Array and BLAS operations.Using this Array package, we also developed two parallel Java versions of the data mining application: one that relies entirely on the implicit parallelism provided by the Array package, and another that is explicitly parallel at the application level.We discuss the design of the Array package, as well as the design of the data mining application.We compare the trade-offs between performance and the abstraction level the different Java versions present to the application programmer.Our studies show that, although a plain Java implementation performs poorly, the Java implementation with the Array package is quite competitive in performance with Fortran.We achieve a single processor performance of 109 Mflops, or 91% of Fortran performance, on a 332 MHz PowerPC 604e processor.Both the implicitly and explicitly parallel forms of our Java implementations also parallelize well.On an SMP with four of those PowerPC processors, the implicitly parallel form achieves 290 Mflops with no effort from the application programmer, while the explicitly parallel form achieves 340 Mflops.1
José E. Moreira, Samuel P. Midkiff, Manish Gupta 0002, Rick Lawrence
SC1
1998 An Infrastructure for Efficient Parallel Job Execution in Terascale Computing Environments
abstract
Recent Terascale computing environments, such as those in the Department of Energy Accelerated Strategic Computing Initiative, present a new challenge to job scheduling and execution systems. The traditional way to concurrently execute multiple jobs in such large machines is through space-sharing: each job is given dedicated use of a pool of processors. Previous work in this area has demonstrated the benefits of sharing the parallel machine's resources not only spatially but also temporally. Time-sharing creates virtual processors for the execution of jobs. The scheduling is typically performed cyclically and each time-slice of the cycle can be considered an independent virtual machine. When all tasks of a parallel job are scheduled to run on the same time-slice (same virtual machine), gang-scheduling is accomplished. Research has shown that gang-scheduling can greatly improve system utilization and job response time in large parallel systems. We are developing GangLL, a research prototype system for performing gang-scheduling on the ASCI Blue-Pacific machine, an IBM RS/6000 SP to be installed at Lawrence Livermore National Laboratory. This machine consists of several hundred nodes, interconnected by a high-speed communication switch. GangLL is organized as a centralized scheduler that performs global decision-making, and a local daemon in each node that controls job execution according to those decisions. The centralized scheduler builds an Ousterhout matrix that precisely defines the temporal and spatial allocation of tasks in the system. Once the matrix is built, it is distributed to each of the local daemons using a scalable hierarchical distributions scheme. A two-phase commit is used in the distribution scheme to guarantee that all local daemons have consistent information. The local daemons enforce the schedule dedicated by the Ousterhout matrix in their corresponding nodes. This requires suspending and resuming execution of tasks and multiplexing access to the communication switch. Large supercomputing centers tend to have their own job scheduling systems, to handle site specific conditions. Therefore, we are designing GangLL so that it can interact with an external site scheduler. The goal is to let the site scheduler control spatial allocation of jobs, if so desired, and to decide when jobs run. GangLL then performs the detailed temporal allocation and controls the actual execution of jobs. The site scheduler can control the fraction of a shared processor that a job receives through an execution factor parameter. To quantify the benefits of our gang-scheduling system to job execution in a large parallel system, we simulate the system with a realistic workload. We measure performance parameters under various degrees of time-sharing, characterized by the multiprogramming level. Our results show that higher multiprogramming levels lead to higher system utilization and lower job response times. We also report some results from the initial deployment of GangLL on a small multiprocessor system.
José E. Moreira, Waiman Chan, Liana L. Fong, Hubertus Franke, Morris A. Jette
SC1
1998 Dynamic Data Distribution and Processor Repartitioning for Irregularly Structured Computations
José E. Moreira, Vijay K. Naik, Samuel P. Midkiff
J. Parallel Distributed Comput.1
1997 A Checkpointing Strategy for Scalable Recovery on Distributed Parallel Systems
abstract
In this paper, we describe a new scheme for checkpointing parallel applications on message-passing scalable distributed memory systems. The novelty of our scheme is that a checkpointed application can be restored, from its checkpointed state, in a reconfigured form. Thus, a parallel application may be checkpointed while executing with t1 tasks on p1 processors, and then restarted from the checkpointed state with t2 tasks on p2 processors. As a result, applications can recover from partial failures in the underlying system. Also, the reconfigurable checkpointed states can be migrated from one parallel system to another even if they do not have the same number of processors. We describe a new programming model for implementing a reconfigurable checkpointing scheme for parallel programs. This new model is derived from the DRMS programming model, developed in the context of run-time reconfiguration of parallel applications. A key component of our implementation is the distribution-independent representation of application array data structures in persistent storage. For further optimizing the performance of checkpoint/restart operations, we provide parallel array section streaming operations for such distributed arrays. We present performance data for the reconfigurable checkpointing and restarting of parallel applications and compare that with the performance of conventional forms of checkpointing. Our results demonstrate the advantages of the new scheme we describe.
Vijay K. Naik, Samuel P. Midkiff, José E. Moreira
SC3