Sameer Kumar 0001

dblp:73/6026-1 · DBLP profile ↗
← Back
29ranked-venue papers
11as first author
0since 2021 · last 2018
0000-0001-8697-7370ORCID · corroborated

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

Systems, architecture and hardware · 25 · 8 first-author

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

Computer architecture, parallel and distributed computing, and storage systems
6 papers
Interconnection networks and networks-on-chip · 37% High-performance computing · 31% Parallel and multicore computing · 24%
Interdisciplinary, comprehensive, and emerging computing
1 paper
Computational science and engineering · 50% Bioinformatics and computational biology · 50%

Topics — the 17 heaviest of 20, each with the papers that count most for it

TopicWeightPapersLastEvidence papers
Interconnection networks and networks-on-chip › network topology
torus network
0.322012
Looking under the hood of the IBM blue gene/Q network · SC 2012
Collective algorithms for sub-communicators · PPoPP 2012
High-performance computing
collective communication
0.112012
Collective algorithms for sub-communicators · PPoPP 2012
Interconnection networks and networks-on-chip
network topology
0.112012
A divide and conquer strategy for scaling weather simulations with multiple regions of interest · SC 2012
Parallel and multicore computing
parallel programming models and scheduling
0.112012
A divide and conquer strategy for scaling weather simulations with multiple regions of interest · SC 2012
Parallel and multicore computing
processor allocation
0.112012
A divide and conquer strategy for scaling weather simulations with multiple regions of interest · SC 2012
High-performance computing
scientific computing systems
0.112012
A divide and conquer strategy for scaling weather simulations with multiple regions of interest · SC 2012
Interconnection networks and networks-on-chip › network topology
topology-aware mapping
0.112012
A divide and conquer strategy for scaling weather simulations with multiple regions of interest · SC 2012
Interconnection networks and networks-on-chip
network interface
0.112011
The IBM Blue Gene/Q interconnection network and message unit · SC 2011
High-performance computing
supercomputer architecture
0.112011
The IBM Blue Gene/Q interconnection network and message unit · SC 2011
Parallel and multicore computing
parallel programming runtimes
0.122006
Performance evaluation of adaptive MPI · PPoPP 2006
NAMD: biomolecular simulation on thousands of processors · SC 2002
Cloud and datacenter computing › virtualization › resource virtualization
processor virtualization
0.112006
Performance evaluation of adaptive MPI · PPoPP 2006
Performance modeling and evaluation
benchmarking
0.012012
Looking under the hood of the IBM blue gene/Q network · SC 2012
Parallel and multicore computing › parallel programming models
message passing
0.012012
Collective algorithms for sub-communicators · PPoPP 2012
High-performance computing
performance optimization at scale
0.012002
NAMD: biomolecular simulation on thousands of processors · SC 2002
Parallel and multicore computing › parallel programming models › message passing
MPI implementation
0.012006
Performance evaluation of adaptive MPI · PPoPP 2006
Bioinformatics and computational biology › molecular informatics › molecular modeling
biomolecular simulation
0.012002
NAMD: biomolecular simulation on thousands of processors · SC 2002
Computational science and engineering › computational chemistry › molecular simulation
molecular dynamics
0.012002
NAMD: biomolecular simulation on thousands of processors · SC 2002

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

torus interconnect mapping · 0.1performance prediction · 0.1offloading · 0.1dynamic routing · 0.1communication threads · 0.1collective algorithm design · 0.1routing algorithm · 0.1packet injection parallelization · 0.1dynamic load balancing · 0.1checkpoint/restart · 0.1spatial decomposition · 0.0predictive load balancing · 0.0particle mesh ewald · 0.0
YearPublicationVenuePosition
2018 Efficient Training of Convolutional Neural Nets on Large Distributed Systems
abstract
Deep Neural Networks (DNNs) have achieved impressive accuracy in many application domains including im-age classification. Training of DNNs is an extremely compute-intensive process and is solved using variants of the stochastic gradient descent (SGD) algorithm. A lot of recent research has focused on improving the performance of DNN training. In this paper, we present optimization techniques to improve the performance of the data parallel synchronous SGD algorithm using the Torch framework: (i) we maintain data in-memory to avoid file I/O overheads, (ii) we propose optimizations to the Torch data parallel table framework that handles multi-threading, and (iii) we present MPI optimization to minimize communication overheads. We evaluate the performance of our optimizations on a Power 8 Minsky cluster with 64 nodes and 256 NVidia Pascal P100 GPUs. With our optimizations, we are able to train 90 epochs of the ResNet-50 model on the Imagenet-1k dataset using 256 GPUs in just 48 minutes. This significantly improves on the previously best known performance of training 90 epochs of the ResNet-50 model on the same dataset using the same number of GPUs in 65 minutes. To the best of our knowledge, this is the best known training performance demonstrated for the Imagenet-1k dataset using 256 GPUs.
Dheeraj Sreedhar, Vaibhav Saxena, Yogish Sabharwal, Ashish Verma 0001, Sameer Kumar 0001
CLUSTER5
2016 Optimization of Message Passing Services on POWER8 InfiniBand Clusters
abstract
We present scalability and performance enhancements to MPI libraries on POWER8 InfiniBand clusters. We explore optimizations in the Parallel Active Messaging Interface (PAMI) libraries. We bypass IB VERBS via low level inline calls resulting in low latencies and high message rates. MPI is enabled on POWER8 by extension of both MPICH and Open MPI to call PAMI libraries. The IBM POWER8 nodes have GPU accelerators to optimize floating throughput of the node. We explore optimized algorithms for GPU-to-GPU communication with minimal processor involvement. We achieve a peak MPI message rate of 186 million messages per second. We also present scalable performance in the QBOX and AMG applications.
Sameer Kumar 0001, Robert Blackmore, Sameh Sharkawi, K. A. Nysal Jan, Amith R. Mamidala, T. J. Christopher Ward
EuroMPI1
2016 Space Performance Tradeoffs in Compressing MPI Group Data Structures
abstract
MPI is a popular programming paradigm on parallel machines today. MPI libraries sometimes use O(N) data structures to implement MPI functionality. The IBM Blue Gene/Q machine has 16 GB memory per node. If each node runs 32 MPI processes, only 512 MB is available per process, requiring the MPI library to be space efficient. This scenario will become severe in a future Exascale machine with tens of millions of cores and MPI endpoints. We explore techniques to compress the dense O(N) mapping data structures that map the logical process ID to the global rank. Our techniques minimize topological communicator mapping state by replacing table lookups with a mapping function. We also explore caching schemes with performance results to optimize overheads of the mapping functions for recent translations in multiple MPI micro-benchmarks, and the 3D FFT and Algebraic Multi Grid application benchmarks.
Sameer Kumar 0001, Philip Heidelberger, Craig B. Stunkel
EuroMPI1
2013 Optimization of MPI_Allreduce on the blue Gene/Q supercomputer
abstract
The IBM Blue Gene/Q supercomputer has a 5D torus network where each node is connected to ten bi-directional links. In this paper we present techniques to optimize the MPI_Allreduce collective operation by building ten different edge disjoint spanning trees on the ten torus links. We accelerate summing of network packets with local buffers by the use of Quad Processing SIMD unit in the BG/Q cores and executing the sums on multiple communication threads created by the PAMI libraries. The net gain we achieve is a peak throughput of 6.3 GB/sec for double precision floating point sum allreduce, that is a speedup of 3.75x over the collective network based algorithm in the product MPI stack on BG/Q.
Sameer Kumar 0001, Daniel Faraj
EuroMPI1
2012 Performance Evaluation and Optimization of Nested High Resolution Weather Simulations
Preeti Malakar, Vaibhav Saxena, Thomas George, Rashmi Mittal, Sameer Kumar 0001, Abdul Ghani Naim, Saiful Azmi bin Hj Husain
Euro-Par5
2012 Collective algorithms for sub-communicators
abstract
Collective communication over a group of processors is an integral and time consuming component in many high performance computing applications. Many modern day super- computers are based on torus interconnects and near optimal algorithms have been developed for collective communication over regular communicators on these systems. However, for an irregular communicator comprising of a subset of processors, the algorithms developed so far are not contention free in general and hence non-optimal. In this paper, we present a novel contention-free algorithm to perform collective operations over a subset of processors in a torus network. We also extend previous work on regular communicators to handle special cases of irregular communicators that occur frequently in parallel scientific applications. For the generic case where multiple node disjoint sub-communicators communicate simultaneously in a loosely synchronous fashion, we propose a novel cooperative approach to route the data for individual sub- communicators without contention. Empirical results demon- strate that our algorithms outperform the optimized MPI collective implementation on IBM's Blue Gene/P supercomputer for large data sizes and random node distributions.
Anshul Mittal, Thomas George, Yogish Sabharwal, Sameer Kumar 0001
ICS5
2012 PAMI: A Parallel Active Message Interface for the Blue Gene/Q Supercomputer
abstract
The Blue Gene/Q machine is the next generation in the line of IBM massively parallel supercomputers, designed to scale to 262144 nodes and sixteen million threads. With each BG/Q node having 68 hardware threads, hybrid programming paradigms, which use message passing among nodes and multi-threading within nodes, are ideal and will enable applications to achieve high throughput on BG/Q. With such unprecedented massive parallelism and scale, this paper is a groundbreaking effort to explore the design challenges for designing a communication library that can match and exploit such massive parallelism In particular, we present the Parallel Active Messaging Interface (PAMI) library as our BG/Q library solution to the many challenges that come with a machine at such scale. PAMI provides (1) novel techniques to partition the application communication overhead into many contexts that can be accelerated by communication threads, (2) client and context objects to support multiple and different programming paradigms, (3) lockless algorithms to speed up MPI message rate, and (4) novel techniques leveraging the new BG/Q architectural features such as the scalable atomic primitives implemented in the L2 cache, the highly parallel hardware messaging unit that supports both point-to-point and collective operations, and the collective hardware acceleration for operations such as broadcast, reduce, and all reduce. We experimented with PAMI on 2048 BG/Q nodes and the results show high messaging rates as well as low latencies and high throughputs for collective communication operations.
Sameer Kumar 0001, Amith R. Mamidala, Daniel Faraj, Brian E. Smith, Michael Blocksome, Bob Cernohous, Douglas Miller, Jeff Parker, Joe Ratterman, Philip Heidelberger, Dong Chen 0005, Burkhard D. Steinmacher-Burow
IPDPS1
2012 Collective algorithms for sub-communicators
abstract
Collective communication over a group of processors is an integral and time consuming component in many HPC applications. Many modern day supercomputers are based on torus interconnects. On such systems, for an irregular communicator comprising of a subset of processors, the algorithms developed so far are not contention free in general and hence non-optimal.
Anshul Mittal, Thomas George, Yogish Sabharwal, Sameer Kumar 0001
PPoPP5
2012 Looking under the hood of the IBM blue gene/Q network
abstract
This paper explores the performance and optimization of the IBM Blue Gene/Q (BG/Q) five dimensional torus network on up to 16K nodes. The BG/Q hardware supports multiple dynamic routing algorithms and different traffic patterns may require different algorithms to achieve best performance. Between 85% to 95% of peak network performance is achieved for all-to-all traffic, while over 85% of peak is obtained for challenging bisection pairings. A new software-controlled algorithm is developed for bisection traffic that selects which hardware algorithm to employ and achieves better performance than any individual hardware algorithm. The benefit of dynamic routing is shown for a highly non-uniform "transpose" traffic pattern. To evaluate memory and network performance, the HPCC Random Access benchmark was tuned for BG/Q and achieved 858 Giga Updates per Second (GUPS) on 16K nodes. To further accelerate message processing, the message libraries on BG/Q enable the offloading of messaging overhead onto dedicated communication threads. Several applications, including Algebraic Multigrid (AMG), exhibit from 3 to 20% gain using communication threads.
Dong Chen 0005, Noel Eisley, Philip Heidelberger, Sameer Kumar 0001, Amith R. Mamidala, Fabrizio Petrini, Robert M. Senger, Yutaka Sugawara, Robert Walkup, Burkhard D. Steinmacher-Burow, Anamitra R. Choudhury, Yogish Sabharwal, Swati Singhal, Jeff Parker
SC4
2012 A divide and conquer strategy for scaling weather simulations with multiple regions of interest
abstract
Accurate and timely prediction of weather phenomena, such as hurricanes and flash floods, require high-fidelity compute intensive simulations of multiple finer regions of interest within a coarse simulation domain. Current weather applications execute these nested simulations sequentially using all the available processors, which is sub-optimal due to their sub-linear scalability. In this work, we present a strategy for parallel execution of multiple nested domain simulations based on partitioning the 2-D processor grid into disjoint rectangular regions associated with each domain. We propose a novel combination of performance prediction, processor allocation methods and topology-aware mapping of the regions on torus interconnects. Experiments on IBM Blue Gene systems using WRF show that the proposed strategies result in performance improvement of up to 33% with topology-oblivious mapping and up to additional 7% with topology-aware mapping over the default sequential strategy.
Preeti Malakar, Thomas George, Sameer Kumar 0001, Rashmi Mittal, Vijay Natarajan, Yogish Sabharwal, Vaibhav Saxena, Sathish S. Vadhiyar
SC3
2011 The IBM Blue Gene/Q interconnection network and message unit
abstract
This is the first paper describing the IBM Blue Gene/Q interconnection network and message unit. The Blue Gene/Q system is the third generation in the IBM Blue Gene line of massively parallel supercomputers. The Blue Gene/Q architecture can be scaled to 20 PF/s and beyond. The network and the highly parallel message unit, which provides the functionality of a network interface, are integrated onto the same chip as the processors and cache memory, and consume 8% of the chip's area. For better application scalability and performance, we describe new routing algorithms and new techniques to parallelize the injection and reception of packets in the network interface. Measured hardware performance results are also presented.
Dong Chen 0005, Noel Eisley, Philip Heidelberger, Robert M. Senger, Yutaka Sugawara, Sameer Kumar 0001, Valentina Salapura, David L. Satterfield, Burkhard D. Steinmacher-Burow, Jeff Parker
SC6
2010 Minimizing MPI Resource Contention in Multithreaded Multicore Environments
abstract
With 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
CLUSTER6
2010 Optimization of applications with non-blocking neighborhood collectives via multisends on the Blue Gene/P supercomputer
abstract
We explore the multisend interface as a data mover interface to optimize applications with neighborhood collective communication operations. One of the limitations of the current MPI 2.1 standard is that the vector collective calls require counts and displacements (zero and nonzero bytes) to be specified for all the processors in the communicator. Further, all the collective calls in MPI 2.1 are blocking and do not permit overlap of communication with computation. We present the record replay persistent optimization to the multisend interface that minimizes the processor overhead of initiating the collective. We present four different case studies with the multisend API on Blue Gene/P (i) 3D-FFT, (ii) 4D nearest neighbor exchange as used in Quantum Chromodynamics, (iii) NAMD and (iv) neural network simulator NEURON. Performance results show 1.9× speedup with 32(3) 3D-FFTs, 1.9× speedup for 4D nearest neighbor exchange with the 2(4) problem, 1.6× speedup in NAMD and almost 3× speedup in NEURON with 256K cells and 1k connections/cell.
Sameer Kumar 0001, Philip Heidelberger, Dong Chen 0005, Michael L. Hines
IPDPS1
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
EuroMPI2
2009 Dynamic topology aware load balancing algorithms for molecular dynamics applications
abstract
Molecular Dynamics applications enhance our understanding of biological phenomena through bio-molecular simulations. Large-scale parallelization of MD simulations is challenging because of the small number of atoms and small time scales involved. Load balancing in parallel MD programs is crucial for good performance on large parallel machines. This paper discusses load balancing algorithms deployed in a MD code called NAMD. It focuses on new schemes deployed in the load balancers and provides an analysis of the performance benefits achieved. Specifically, the paper presents the technique of topology-aware mapping on 3D mesh and torus architectures, used to improve scalability and performance. These techniques have a wide applicability for latency intolerant applications.
Abhinav Bhatele, Laxmikant V. Kalé, Sameer Kumar 0001
ICS3
2009 MPI collective communications on the blue gene/p supercomputer: algorithms and optimizations
abstract
The IBM Blue Gene/P (BG/P) system is a massively parallel supercomputer succeeding BG/L, and it is based on orders of magnitude in system size and significant power consumption efficiency. BG/P comes with many enhancements to the machine design and new architectural features at the hardware and software levels. In this work, we demonstrate techniques to leverage the architectural features of BG/P to deliver high performance MPI collective communication primitives.
Ahmad Faraj, Sameer Kumar 0001, Brian E. Smith, Amith R. Mamidala, John A. Gunnels, Philip Heidelberger
ICS2
2008 Optimization of All-to-All Communication on the Blue Gene/L Supercomputer
abstract
All-to-all communication is a well known performance bottleneck for many applications, such as the ones that use the Fast-Fourier-transform (FFT) algorithm. We analyze the performance of all-to-all communication on the BlueGene/L torus interconnect that has link contention even for all-to-all operations with short messages. We observed that the performance of all-to-all depends on the shape of the processor partition. We present a performance analysis of all-to-all on partitions of various shapes. We then present optimization schemes that substantially improve the performance of all-to-all with short and large messages.In particular, throughput improved from 64% to over 99% of peak on the 65,536 (64 times 32 times 32) node Blue Gene/L machine at the Lawrence Livermore National Lab. We show the impact of the all-to-all performance optimizations in 1-D and 3-D FFT benchmarks. We achieved a performance of over 2.8 TF for the HPC Challenge 1D FFT benchmark with our optimized all-to-all.
Sameer Kumar 0001, Yogish Sabharwal, Rahul Garg 0001, Philip Heidelberger
ICPP1
2008 The deep computing messaging framework: generalized scalable message passing on the blue gene/P supercomputer
abstract
We present the architecture of the Deep Computing Messaging Framework (DCMF), a message passing runtime designed for the Blue Gene/P machine and other HPC architectures. DCMF has been designed to easily support several programming paradigms such as the Message Passing Interface (MPI), Aggregate Remote Memory Copy Interface (ARMCI), Charm++, and others. This support is made possible as DCMF provides an application programming interface (API) with active messages and non-blocking collectives. DCMF is being open sourced and has a layered component based architecture with multiple levels of abstraction, allowing the members of the community to contribute new components to its design at the various layers. The DCMF runtime can be extended to other architectures through the development of architecture specific implementations of interface classes. The production DCMF runtime on Blue Gene/P takes advantage of the direct memory access (DMA) hardware to offload message passing work and achieve good overlap of computation and communication. We take advantage of the fact that the Blue Gene/P node is a symmetric multi-processor with four cache-coherent cores and use multi-threading to optimize the performance on the collective network. We also present a performance evaluation of the DCMF runtime on Blue Gene/P and show that it delivers performance close to hardware limits.
Sameer Kumar 0001, Gábor Dózsa, Gheorghe Almási 0001, Philip Heidelberger, Dong Chen 0005, Mark Giampapa, Michael Blocksome, Ahmad Faraj, Jeff Parker, Joe Ratterman, Brian E. Smith, Charles Archer
ICS1
2008 Evaluating the effect of replacing CNK with linux on the compute-nodes of blue gene/l
abstract
The Blue Gene machines in production today run a small single-user, single-process kernel (CNK) having a limited functionality. Motivated by the desire to provide applications with a much richer operating environment, we evaluate the effect of replacing CNK with a standard Linux kernel on the compute nodes of Blue Gene/L. We show that with a relatively small amount of effort we were able to improve benchmark performance under Linux up to a level that is comparable to CNK.
Edi Shmueli, Gheorghe Almási 0001, José R. Brunheroto, José G. Castaños, Gábor Dózsa, Sameer Kumar 0001, Derek Lieber
ICS6
2008 Overcoming scaling challenges in biomolecular simulations across multiple platforms
abstract
NAMD is a portable parallel application for biomolecular simulations. NAMD pioneered the use of hybrid spatial and force decomposition, a technique now used by most scalable programs for biomolecular simulations, including Blue Matter and Desmond developed by IBM and D. E. Shaw respectively. NAMD has been developed using Charm++ and benefits from its adaptive communication-computation overlap and dynamic load balancing. This paper focuses on new scalability challenges in biomolecular simulations: using much larger machines and simulating molecular systems with millions of atoms. We describe new techniques developed to overcome these challenges. Since our approach involves automatic adaptive runtime optimizations, one interesting issue involves dealing with harmful interaction between multiple adaptive strategies. NAMD runs on a wide variety of platforms, ranging from commodity clusters to supercomputers. It also scales to large machines: we present results for up to 65,536 processors on IBM’s Blue Gene/L and 8,192 processors on Cray XT3/XT4. In addition, we present performance results on NCSA’s Abe, SDSC’s DataStar and TACC’s LoneStar cluster, to demonstrate efficient portability. We also compare NAMD with Desmond and Blue Matter.
Abhinav Bhatele, Sameer Kumar 0001, James C. Phillips, Gengbin Zheng, Laxmikant V. Kalé
IPDPS2
2006 Achieving strong scaling with NAMD on Blue Gene/L
abstract
NAMD is a scalable molecular dynamics application, which has demonstrated its performance on several parallel computer architectures. Strong scaling is necessary for molecular dynamics as problem size is fixed, and a large number of iterations need to be executed to understand interesting biological phenomenon. The Blue Gene/L machine is a massive source of compute power. It consists of tens of thousands of embedded Power PC 440 processors. In this paper, we present several techniques to scale NAMD to 8192 processors of Blue Gene/L. These include topology specific optimizations, new messaging protocols, load-balancing, and overlap of computation and communication. We were able to achieve 1.2 TF of peak performance for cutoff simulations and 0.99 TF with PME.
Sameer Kumar 0001, Chao Huang 0029, Gheorghe Almási 0001, Laxmikant V. Kalé
IPDPS1
2006 Performance evaluation of adaptive MPI
abstract
Processor virtualization via migratable objects is a powerful technique that enables the runtime system to carry out intelligent adaptive optimizations like dynamic resource management. CHARM++ is an early language/system that supports migratable objects. This paper describes Adaptive MPI (or AMPI), an MPI implementation and extension, that supports processor virtualization. AMPI implements virtual MPI processes (VPs), several of which may be mapped to a single physical processor. AMPI includes a powerful runtime support system that takes advantage of the degree of freedom afforded by allowing it to assign VPs onto processors. With this runtime system, AMPI supports such features as automatic adaptive overlapping of communication and computation, automatic load balancing, flexibility of running on arbitrary number of processors, and checkpoint/restart support. It also inherits communication optimization from CHARM++ framework. This paper describes AMPI, illustrates its performance benefits through a series of benchmarks, and shows that AMPI is a portable and mature MPI implementation that offers various performance benefits to dynamic applications.
Chao Huang 0029, Gengbin Zheng, Laxmikant V. Kalé, Sameer Kumar 0001
PPoPP4
2006 Scaling applications to massively parallel machines using Projections performance analysis tool
Laxmikant V. Kalé, Gengbin Zheng, Chee Wai Lee, Sameer Kumar 0001
Future Gener. Comput. Syst.4
2005 Improved Point-to-Point and Collective Communication Performance with Output-Queued High-Radix Routers
Sameer Kumar 0001, Craig B. Stunkel, Laxmikant V. Kalé
HiPC1
2004 Scaling All-to-All Multicast on Fat-tree Networks
Sameer Kumar 0001, Laxmikant V. Kalé
ICPADS1
2004 Faucets: Efficient Resource Allocation on the Computational Grid
abstract
The idea of a "computational grid" suggests that high end computational power can be thought of as a utility, similar to electricity or water. Making this metaphor work requires a sophisticated "power distribution" infrastructure. We present the Faucets framework that aims at providing (a) user-friendly compute power distribution across the grid, (b) market-driven selection of compute servers for each job, resulting in effective utilization of resources across the grid, and (c) improved utilization within individual compute servers. Utilization of individual compute servers is improved by the notions of adaptive jobs and smarter job schedulers. Server selection is facilitated by quality-of-service (QoS) contracts for parallel jobs. Market efficiencies are then attained by a bidding and evaluation system that makes the compute servers compete for every job by submitting bids, thus transforming the computational grid into a free market. Job submission and monitoring is simplified by several tools and databases within the Faucets system. We describe the overall architecture of the system. All the essential components of the system have been implemented, which are described In the work. We also discuss ongoing work and future research issues.
Laxmikant V. Kalé, Sameer Kumar 0001, Mani Potnuru, Jayant DeSouza, Sindhura Bandhakavi
ICPP2
2004 Opportunities and Challenges of Modern Communication Architectures: Case Study with QsNet
abstract
Summary form only given. We describe our efforts to scale message driven applications to a large number of processors on an Alpha cluster interconnected by QsNet. The clustering technology QsNet has a network interface with a communication coprocessor. The presence of the coprocessor minimizes main processor participation in message passing. We show the advantages of the communication coprocessor for message driven applications. To scale fine-grained message driven applications to a large number of processors we had to overcome several hardware, software and operating system hindrances. We describe them in detail and present solutions for them. We use NAMD, a molecular dynamics program, as a case study for many of our performance optimizations.
Sameer Kumar 0001, Laxmikant V. Kalé
IPDPS1
2002 A Malleable-Job System for Timeshared Parallel Machines
abstract
Malleable jobs are parallel programs that can change the number of processors on which they are executing at run time in response to an external command. One of the advantages of such jobs is that a job scheduler for malleable jobs can provide improved system utilization and average response time over a scheduler for traditional jobs. In this paper, we present a programming system for creating malleable jobs that is more general than other current malleable systems. In particular, it is not limited to the master-worker paradigm or the Fortran SPMD programming model, but can also support general purpose parallel programs including those written in MPI and Charm++, and has built-in migration and load-balancing, among other features.
Laxmikant V. Kalé, Sameer Kumar 0001, Jayant DeSouza
CCGRID2
2002 NAMD: biomolecular simulation on thousands of processors
abstract
NAMD is a fully featured, production molecular dynamics program for high performance simulation of large biomolecular systems. We have previously, at SC2000, presented scaling results for simulations with cutoff electrostatics on up to 2048 processors of the ASCI Red machine, achieved with an object-based hybrid force and spatial decomposition scheme and an aggressive measurement-based predictive load balancing framework. We extend this work by demonstrating similar scaling on the much faster processors of the PSC Lemieux Alpha cluster, and for simulations employing efficient (order N log N) particle mesh Ewald full electrostatics. This unprecedented scalability in a biomolecular simulation code has been attained through latency tolerance, adaptation to multiprocessor nodes, and the direct use of the Quadrics Elan library in place of MPI by the Charm++/Converse parallel runtime system.
James C. Phillips, Gengbin Zheng, Sameer Kumar 0001, Laxmikant V. Kalé
SC3