VLDB 2026 Research / reviewers in the wild / expert
Pavan Balaji
dblp:70/2491
· DBLP profile ↗
162ranked-venue papers
26as first author
12since 2021 · last 2026
—ORCID · conflict
Domains — the database's venue-derived domains; a paper can count in several
Systems, architecture and hardware · 141 · 22 first-author · 10 since 2021Software engineering, systems software and programming languages · 4 · 2 first-author · 1 since 2021Computer networks · 3 · 1 first-author · 1 since 2021Artificial intelligence and machine learning · 1 · 1 since 2021Security and privacy · 1Databases, data management, data science and information retrieval · 1Applied, interdisciplinary, general and emerging computing · 1
| Year | Publication | Venue | Position |
|---|---|---|---|
| 2026 | Connecting 100K+ GPUs: Building the Communication Stack for Large-Scale LLM TrainingabstractThe arrival of 100K+ GPU clusters marks a new frontier in AI infrastructure. Standard communication stack meets new challenges as physical topologies span multiple datacenter buildings, introducing high bandwidth-delay product links where latency increases by up to 30× compared to intra-rack traffic. Furthermore, the transition toward Mixture-of-Experts architectures generating bursty all-to-all patterns that create transient congestion hotspots. These constraints, combined with an operational environment where hardware failures shift from anomalies to frequent occurrences, renders traditionally lightweight operations like initialization and resource management challenging. Hongyi Zeng, Min Si, Pavan Balaji, Yongzhou Chen, Ching-Hsiang Chu, Adithya Gangidi, Prashanth Kannan, Bingzhe Liu, Saif Hasan, Deep Shah, Ashmitha Jeevaraj Shetty, Gregory R. Steinbrecher, Srikanth Sundaresan, Yulun Wang, Yexin Wu, Mingran Yang, Kenny Yu, Minlan Yu, Cen Zhao, Shengbao Zheng, Wesley Bland, Denis Boyda, Suman Gumudavelli, Subodh Iyengar, Cristian Lumezanu, Rui Miao 0001, Venkat Ramesh, Jingliang Ren, Maxim Samoylov, Jan Seidel, Qiye Tan, Xinfeng Xie, Yimeng Zhao, Shuqiang Zhang, Art Zhu |
SIGCOMM | 3 |
| 2025 | HALoS: Hierarchical Asynchronous Local SGD over Slow Networks for Geo-Distributed Large Language Model TrainingabstractTraining large language models (LLMs) increasingly relies on geographically distributed accelerators, causing prohibitive communication costs across regions and uneven utilization of heterogeneous hardware. We propose HALoS, a hierarchical asynchronous optimization framework that tackles these issues by introducing local parameter servers (LPSs) within each region and a global parameter server (GPS) that merges updates across regions. This hierarchical design minimizes expensive inter-region communication, reduces straggler effects, and leverages fast intra-region links. We provide a rigorous convergence analysis for HALoS under non-convex objectives, including theoretical guarantees on the role of hierarchical momentum in asynchronous training. Empirically, HALoS attains up to 7.5× faster convergence than synchronous baselines in geo-distributed LLM training and improves upon existing asynchronous methods by up to 2.1×. Crucially, HALoS preserves the model quality of fully synchronous SGD—matching or exceeding accuracy on standard language modeling and downstream benchmarks—while substantially lowering total training time. These results demonstrate that hierarchical, server-side update accumulation and global model merging are powerful tools for scalable, efficient training of new-era LLMs in heterogeneous, geo-distributed environments. Geon-Woo Kim, Shashidhar Gandham, Omar Baldonado, Adithya Gangidi, Pavan Balaji, Zhangyang Wang, Aditya Akella |
ICML | 6 |
| 2025 | Scaling Llama 3 Training with Efficient Parallelism StrategiesabstractLlama is a widely used open-source large language model.This paper presents the design and implementation of the parallelism techniques used in Llama 3 pre-training.To achieve efficient training on tens of thousands of GPUs, Llama 3 employs a combination of four-dimensional parallelism: fully sharded data parallelism, tensor parallelism, pipeline parallelism, and context parallelism.Beyond achieving efficiency through parallelism and model co-design, we Weiwei Chu, Xinfeng Xie, Jiecao Yu, Jie Wang 0022, Amar Phanishayee, Chunqiang Tang, Yuchen Hao, Muhammet Mustafa Ozdal, Vedanuj Goswami, Naman Goyal 0001, Abhishek Kadian, Andrew Gu, Chris Cai, Xiaodong Wang 0020, Min Si, Pavan Balaji, Ching-Hsiang Chu, Jongsoo Park |
ISCA | 19 |
| 2024 | UpDLRM: Accelerating Personalized Recommendation using Real-World PIM ArchitectureabstractDeep Learning Recommendation Models (DLRMs) have gained popularity in recommendation systems due to their effectiveness in handling large-scale recommendation tasks. The embedding layers of DLRMs have become the performance bottleneck due to their intensive needs on memory capacity and memory bandwidth. In this paper, we propose UpDLRM, which utilizes real-world processing-in-memory (PIM) hardware, UPMEM DPU, to boost the memory bandwidth and reduce recommendation latency. The parallel nature of the DPU memory can provide high aggregated bandwidth for the large number of irregular memory accesses in embedding lookups, thus offering great potential to reduce the inference latency. To fully utilize the DPU memory bandwidth, we further studied the embedding table partitioning problem to achieve good workload-balance and efficient data caching. Evaluations using real-world datasets show that, UpDLRM achieves much lower inference time for DLRM compared to both CPU-only and CPU-GPU hybrid counterparts. Sitian Chen, Haobin Tan, Amelie Chi Zhou, Yusen Li, Pavan Balaji |
DAC | 5 |
| 2024 | Accelerating Communication in Deep Learning Recommendation Model Training with Dual-Level Adaptive Lossy CompressionabstractDLRM is a state-of-the-art recommendation system model that has gained widespread adoption across various industry applications. The large size of DLRM models, however, necessitates the use of multiple devices/GPUs for efficient training. A significant bottleneck in this process is the time-consuming all-to-all communication required to collect embedding data from all devices. To mitigate this, we introduce a method that employs error-bounded lossy compression to reduce the communication data size and accelerate DLRM training. We develop a novel error-bounded lossy compression algorithm, informed by an in-depth analysis of embedding data features, to achieve high compression ratios. Moreover, we introduce a dual-level adaptive strategy for error-bound adjustment, spanning both table-wise and iteration-wise aspects, to balance the compression benefits with the potential impacts on accuracy. We further optimize our compressor for PyTorch tensors on GPUs, minimizing compression overhead. Evaluation shows that our method achieves a 1.38 × training speedup with a minimal accuracy impact. Boyuan Zhang 0002, Fanjiang Ye, Min Si, Ching-Hsiang Chu, Jiannan Tian, Chunxing Yin, Summer Deng, Yuchen Hao, Pavan Balaji, Tong Geng, Dingwen Tao |
SC | 10 |
| 2023 | Near-Lossless MPI Tracing and Proxy Application AutogenerationabstractTraces of MPI communications are used by many performance analysis and visualization tools. Storing exhaustive traces of large-scale MPI applications is infeasible, however, because of their large volume. Aggregated or lossy MPI traces are smaller but provide much less information. In this paper we present Pilgrim, a near-lossless MPI tracing tool that, by using sophisticated compression techniques, generates small trace files at large scales and incurs only moderate overheads. We perform comprehensive studies of various compression techniques used for storing timestamps associated with each call. This timing information is essential for analysis purposes such as skews study. To demonstrate the usefulness of the detailed information stored by Pilgrim, we present a proxy application generator that can generate proxy apps that preserve original communication patterns from the Pilgrim traces. Chen Wang 0004, Yanfei Guo, Pavan Balaji, Marc Snir |
IEEE Trans. Parallel Distributed Syst. | 3 |
| 2021 | RMACXX: An Efficient High-Level C++ Interface over MPI-3 RMAabstractParallel scientific applications can benefit from decoupling communication and synchronization. One-sided programming abstractions, which separate communication from synchronization, have in fact served as a motivation for partitioned global address space (PGAS) models. However, the use of PGAS models in application codes in a manner that fully exploits the benefit of these programming models requires significant development effort. Meanwhile, a vast majority of scientific codes already use the Message Passing Interface (MPI) and need convenient features to support application-specific one-sided communication scenarios. MPI Remote Memory Access (RMA) can be employed for this purpose. MPI is a low-level API, however, and developing applications with MPI RMA requires programmers to be well versed in its nuances. We present RMACXX, a compact set of C++ bindings to MPI-3 RMA, to ease the use of MPI RMA. Unlike other PGAS models, which may have interoperability issues with MPI, RMACXX is written on top of MPI and uses the same runtime as MPI. The basic functionality of RMACXX adds only a relatively small number of extra instructions (about 20) to the critical communication path. Moreover, RMACXX provides an intuitive API for building a wide variety of scientific applications while enjoying performance matching handwritten MPI-3 RMA codes. Yanfei Guo, Pavan Balaji, Assefaw Hadish Gebremedhin |
CCGRID | 3 |
| 2021 | Daps: A Dynamic Asynchronous Progress Stealing Model for MPI CommunicationabstractMPI provides nonblocking point-to-point and one-sided communication models to help applications achieve communication and computation overlap. These models provide the opportunity for MPI to offload data transfer to low level network hardware while the user process is computing. In practice, however, MPI implementations have to often handle complex data transfer in software due to limited capability of network hardware. Therefore, additional asynchronous progress is necessary to ensure prompt progress of these software-handled communication. Traditional mechanisms either spawn an additional background thread on each MPI process or launch a fixed number of helper processes on each node. Both mechanisms may degrade performance in user computation due to statically occupied CPU resources. The user has to fine-tune the progress resource deployment to gain overall performance. For complex multiphase applications, unfortunately, severe performance degradation may occur due to dynamically changing communication characteristics and thus changed progress requirement. This paper proposes a novel Dynamic Asynchronous Progress Stealing model, called Daps, to completely address the asynchronous progress complication. Daps is implemented inside the MPI runtime. It dynamically leverages idle MPI processes to steal communication progress tasks from other busy computing processes located on the same node. The basic concept of Daps is straightforward; however, various implementation challenges have to be resolved due to the unique requirements of interprocess data and code sharing. We present our design that ensures high performance while maintaining strict program correctness. We compare Daps with state-of-the-art asynchronous progress approaches by utilizing both microbenchmarks and HPC proxy applications. Kaiming Ouyang, Min Si, Atsushi Hori, Zizhong Chen, Pavan Balaji |
CLUSTER | 5 |
| 2021 | Lightweight preemptive user-level threadsabstractMany-to-many mapping models for user- to kernel-level threads (or "M:N threads") have been extensively studied for decades as a lightweight substitute for current Pthreads implementations that provide a simple one-to-one mapping ("1:1 threads"). M:N threads derive performance from their ability to allow users to context switch between threads and control their scheduling entirely in user space with no kernel involvement. This same ability, however, causes M:N threads to lose the kernel-provided ability of implicit OS preemption---threads have to explicitly yield control for other threads to be scheduled. Hence, programs over nonpreemptive M:N threads can cause core starvation, loss of prioritization, and, sometimes, deadlock unless programs are written to explicitly yield in proper places. This paper explores two techniques for M:N threads to efficiently achieve implicit preemption similar to 1:1 threads: signal-yield and KLT-switching. Overheads of these techniques, with our optimizations, can be less than 1% compared with nonpreemptive M:N threads. Our evaluation with three applications demonstrates that our preemption techniques for M:N threads improve core utilization and enhance the performance by utilizing lightweight context switching and flexible scheduling of M:N threads. Shumpei Shiina, Shintaro Iwasaki, Kenjiro Taura, Pavan Balaji |
PPoPP | 4 |
| 2021 | Pilgrim: scalable and (near) lossless MPI tracingabstractTraces of MPI communications are used by many performance analysis and visualization tools. Storing exhaustive traces of large scale MPI applications is infeasible, due to their large volume. Aggregated or lossy MPI traces are smaller, but provide much less information. In this paper, we present Pilgrim, a near lossless MPI tracing tool that incurs moderate overheads and generates small trace files at large scales, by using sophisticated compression techniques. Furthermore, for codes with regular communication patterns, Pilgrim can store their traces in constant space regardless of the problem size, the number of processors, and the number of iterations. In comparison with existing tools, Pilgrim preserves more information with less space in all the programs we tested. Chen Wang 0004, Pavan Balaji, Marc Snir |
SC | 2 |
| 2021 | Guest Editorial
Pavan Balaji, Jidong Zhai, Min Si |
IEEE Trans. Parallel Distributed Syst. | 1 |
| 2021 | Logically Parallel Communication for Fast MPI+Threads ApplicationsabstractSupercomputing applications are increasingly adopting the MPI+threads programming model over the traditional “MPI everywhere” approach to better handle the disproportionate increase in the number of cores compared with other on-node resources. In practice, however, most applications observe a slower performance with MPI+threads primarily because of poor communication performance. Recent research efforts on MPI libraries address this bottleneck by mapping logically parallel communication, that is, operations that are not subject to MPI's ordering constraints to the underlying network parallelism. Domain scientists, however, typically do not expose such communication independence information because the existing MPI-3.1 standard's semantics can be limiting. Researchers had initially proposed user-visible endpoints to combat this issue, but such a solution requires intrusive changes to the standard (new APIs). The upcoming MPI-4.0 standard, on the other hand, allows applications to relax unneeded semantics and provides them with many opportunities to express logical communication parallelism. In this article, we show how MPI+threads applications can achieve high performance with logically parallel communication. Through application case studies, we compare the capabilities of the new MPI-4.0 standard with those of the existing one and user-visible endpoints (upper bound). Logical communication parallelism can boost the overall performance of an application by over 2×. Rohit Zambre, Damodar Sahasrabudhe, Hui Zhou 0012, Martin Berzins, Aparna Chandramowlishwaran, Pavan Balaji |
IEEE Trans. Parallel Distributed Syst. | 6 |
| 2020 | How I learned to stop worrying about user-visible endpoints and love MPIabstractMPI+threads is gaining prominence as an alternative to the traditional "MPI everywhere" model in order to better handle the disproportionate increase in the number of cores compared with other on-node resources. However, the communication performance of MPI+threads can be 100x slower than that of MPI everywhere. Both MPI users and developers are to blame for this slowdown. MPI users traditionally have not exposed logical communication parallelism. Consequently, MPI libraries have used conservative approaches, such as a global critical section, to maintain MPI's ordering constraints for MPI+threads, thus serializing access to the underlying parallel network resources and limiting performance. Rohit Zambre, Aparna Chandramowlishwaran, Pavan Balaji |
ICS | 3 |
| 2020 | CAB-MPI: exploring interprocess work-stealing towards balanced MPI communicationabstractLoad balance is essential for high-performance applications. Unbalanced communication can cause severe performance degradation, even in computation-balanced BSP applications. Designing communication-balanced applications is challenging, however, because of the diverse communication implementations at the underlying runtime system. In this paper, we address this challenge through an interprocess workstealing scheme based on process-memory-sharing techniques. We present CAB-MPI, an MPI implementation that can identify idle processes inside MPI and use these idle resources to dynamically balance communication workload on the node. We design throughput-optimized strategies to ensure efficient stealing of the data movement tasks. We demonstrate the benefit of work stealing through several internal processes in MPI, including intranode data transfer, pack/unpack for noncontiguous communication, and computation in one-sided accumulates. The implementation is evaluated through a set of microbenchmarks and proxy applications on Intel Xeon and Xeon Phi platforms. Kaiming Ouyang, Min Si, Atsushi Hori, Zizhong Chen, Pavan Balaji |
SC | 5 |
| 2020 | Analysis of Threading Libraries for High Performance ComputingabstractWith the appearance of multi-/many core machines, applications and runtime systems have evolved in order to exploit the new on-node concurrency brought by new software paradigms. POSIX threads (Pthreads) was widely-adopted for that purpose and it remains as the most used threading solution in current hardware. Lightweight thread (LWT) libraries emerged as an alternative offering lighter mechanisms to tackle the massive concurrency of current hardware. In this article, we analyze in detail the most representative threading libraries including Pthread- and LWT-based solutions. In addition, to examine the suitability of LWTs for different use cases, we develop a set of microbenchmarks consisting of OpenMP patterns commonly found in current parallel codes, and we compare the results using threading libraries and OpenMP implementations. Moreover, we study the semantics offered by threading libraries in order to expose the similarities among different LWT application programming interfaces and their advantages over Pthreads. This article exposes that LWT libraries outperform solutions based on operating system threads when tasks and nested parallelism are required. Adrián Castelló 0001, Rafael Mayo 0002, Pavan Balaji, Enrique S. Quintana-Ortí, Antonio J. Peña |
IEEE Trans. Computers | 4 |
| 2020 | Memory-Efficient and Skew-Tolerant MapReduce Over MPI for Supercomputing SystemsabstractData analytics has become an integral part of large-scale scientific computing. Among various data analytics frameworks, MapReduce has gained the most traction. Although some efforts have been made to enable efficient MapReduce for supercomputing systems, they are often limited to fairly homogeneous workloads where equal partitioning of input data across tasks results in essentially equal output or temporary data generated on each task. For workloads that are more skewed, however, current implementations can result in imbalance in memory usage and, consequently, can cause a slowdown in execution time and a loss in data scalability. To tackle this problem, we enhance a previously published memory-conscious MapReduce over MPI framework called Mimir. Our enhancements to Mimir include combiner and dynamic repartition optimizations to minimize and balance memory usage and to achieve close to optimal balance of the memory usage across processes and to reduce the execution time by up to 12 times. Experimental results show that Mimir can scale to at least 3072 processes on the Tianhe-2 supercomputer on skewed datasets. Yanfei Guo, Boyu Zhang 0002, Pietro Cicotti, Yutong Lu, Pavan Balaji, Michela Taufer |
IEEE Trans. Parallel Distributed Syst. | 6 |
| 2020 | Analyzing the Performance Trade-Off in Implementing User-Level ThreadsabstractUser-level threads have been widely adopted as a means of achieving lightweight concurrent execution without the costs of OS-level threads. Nevertheless, the costs of managing user-level threads represent a performance barrier that dictates how fine grained the concurrency exposed by an application can be without incurring significant overheads; this in turn may translate into insufficient parallelism to exploit highly parallel systems. This article is a deep dive into the fundamental costs in implementing user-level threads. We first identify that one of the highest sources of fork-join overheads stems from deviations, events that incur context switching during the execution of a thread and disrupt a run-to-completion execution. We then conduct an in-depth investigation of a wide spectrum of methods with respect to how they handle deviations while covering both parent- and child-first scheduling policies. Our methodology involves a comprehensive instruction- and cache-level analysis of all methods on several modern CPU architectures. The primary finding of our evaluation is that dynamic promotion methods that assume the absence of deviation and dynamically provide context-switching support offer the best trade-off between performance and capability when the likelihood of deviation is low. Shintaro Iwasaki, Abdelhalim Amer, Kenjiro Taura, Pavan Balaji |
IEEE Trans. Parallel Distributed Syst. | 4 |
| 2019 | BOLT: Optimizing OpenMP Parallel Regions with User-Level ThreadsabstractOpenMP is widely used by a number of applications, computational libraries, and runtime systems. As a result, multiple levels of the software stack use OpenMP independently of one another, often leading to nested parallel regions. Although exploiting such nested parallelism is a potential opportunity for performance improvement, it often causes destructive performance with leading OpenMP runtimes because of their reliance on heavyweight OS-level threads. User-level threads (ULTs) are more lightweight alternatives but existing ULT-based runtimes suffer from several shortcomings: 1) thread management costs remain significant and outweigh the benefits from additional parallelism; 2) the shift to ULTs often hurts the more common flat parallelism case; and 3) absence of user control over thread-to-CPU binding, a critical feature on modern systems. This paper presents BOLT, a practical ULT-based OpenMP runtime system that efficiently supports both flat and nested parallelism. This is accomplished on three fronts: 1) advanced data reuse and thread synchronization strategies; 2) thread coordination that adapts to the level of oversubscription; and 3) an implementation of the modern OpenMP thread-to-CPU binding interface tailored to ULT-based runtimes. The result is a highly optimized runtime that transparently achieves similar performance compared with leading state-of-the-art widely used OpenMP runtimes under flat parallelism, while outperforming all existing runtimes under nested parallelism. Shintaro Iwasaki, Abdelhalim Amer, Kenjiro Taura, Pavan Balaji |
PACT | 5 |
| 2019 | Optimized Execution of Parallel Loops via User-Defined Scheduling PoliciesabstractOn-node parallelism continues to increase in importance for high-performance computing and most newly deployed supercomputers have tens of processor cores per node. These higher levels of on-node parallelism exacerbate the impact of load imbalance and locality in parallel computations, and current programming systems notably lack features to enable efficient use of these large numbers of cores or require users to modify codes significantly. Our work is motivated by the need to address application-specific load balance and locality requirements with minimal changes to application codes. Seonmyeong Bak, Yanfei Guo, Pavan Balaji, Vivek Sarkar |
ICPP | 3 |
| 2019 | Software combining to mitigate multithreaded MPI contentionabstractEfforts to mitigate lock contention from concurrent threaded accesses to MPI have reduced contention through fine-grained locking, avoided locking altogether by offloading communication to dedicated threads, or alleviated negative side effects from contention by using better lock management protocols. The blocking nature of lock-based methods, however, wastes the asynchrony benefits of nonblocking MPI operations, and the offloading model sacrifices CPU resources and incurs unnecessary software offloading overheads under low contention. Abdelhalim Amer, Charles Archer, Michael Blocksome, Chongxiao Cao, Michael Chuvelev, Hajime Fujita 0002, María Jesús Garzarán, Yanfei Guo, Jeff R. Hammond, Shintaro Iwasaki, Kenneth Raffenetti, Mikhail Shiryaev, Min Si, Kenjiro Taura, Sagar Thapaliya, Pavan Balaji |
ICS | 16 |
| 2019 | Special issue on the message passing interface
Pavan Balaji, Marc Casas |
Parallel Comput. | 1 |
| 2019 | Foreword to the special issue for the Workshop on Parallel Programming Models and Systems Software for High-End Computing (P2S2 2017)
Pavan Balaji, Abhinav Vishnu, Yong Chen 0001 |
Parallel Comput. | 1 |
| 2019 | International workshop on programming models and applications for multicores and manycores (PMAM 2018)
Min Si, Zhiyi Huang 0001, Pavan Balaji |
Parallel Comput. | 3 |
| 2019 | Guest Editor's Introduction: P2S2: SI 2016
Abhinav Vishnu, Pavan Balaji, Yong Chen 0001 |
Parallel Comput. | 2 |
| 2018 | Process-in-process: techniques for practical address-space sharingabstractThe two most common parallel execution models for many-core CPUs today are multiprocess (e.g., MPI) and multithread (e.g., OpenMP). The multiprocess model allows each process to own a private address space, although processes can explicitly allocate shared-memory regions. The multithreaded model shares all address space by default, although threads can explicitly move data to thread-private storage. In this paper, we present a third model called process-in-process (PiP), where multiple processes are mapped into a single virtual address space. Thus, each process still owns its process-private storage (like the multiprocess model) but can directly access the private storage of other processes in the same virtual address space (like the multithread model). Atsushi Hori, Min Si, Balazs Gerofi, Masamichi Takagi, Jai Dayal, Pavan Balaji, Yutaka Ishikawa |
HPDC | 6 |
| 2018 | On the Power of Combiner Optimizations in MapReduce Over MPI WorkflowsabstractAnalyzing large volumes of data is becoming more and more important in various scientific computing domains. MapReduce over MPI frameworks are an appealing solution to enable scalable big data analytics on supercomputing systems. These systems can further leverage features of MapReduce applications by merging (key/value) pairs before the reduce function in combiner optimizations. In this paper, we propose a pipeline combiner workflow and integrate it into Mimir, a cutting-edge implementation of Map Reduce over MPI. Our results with real datasets on the Tianhe-2 supercomputer prove that our pipeline combiner workflow can reduce memory usage up to 51% and improve the overall performance up to 61%. Yanfei Guo, Boyu Zhang 0002, Pietro Cicotti, Yutong Lu, Pavan Balaji, Michela Taufer |
ICPADS | 6 |
| 2018 | Scalable Communication Endpoints for MPI+Threads ApplicationsabstractHybrid MPI+threads programming is gaining prominence as an alternative to the traditional “MPI everywhere” model to better handle the disproportionate increase in the number of cores compared with other on-node resources. Current implementations of these two models represent the two extreme cases of communication resource sharing in modern MPI implementations. In the MPI-everywhere model, each MPI process has a dedicated set of communication resources (also known as endpoints), which is ideal for performance but is resource wasteful. With MPI+threads, current MPI implementations share a single communication endpoint for all threads, which is ideal for resource usage but is hurtful for performance. In this paper, we explore the tradeoff space between performance and communication resource usage in MPI+threads environments. We first demonstrate the two extreme cases-one where all threads share a single communication endpoint and another where each thread gets its own dedicated communication endpoint (similar to the MPI-everywhere model) and showcase the inefficiencies in both these cases. Next, we perform a thorough analysis of the different levels of resource sharing in the context of Mellanox InfiniBand. Using the lessons learned from this analysis, we design an improved resource-sharing model to produce scalable communication endpoints that can achieve the same performance as with dedicated communication resources per thread but using just a third of the resources. Rohit Zambre, Aparna Chandramowlishwaran, Pavan Balaji |
ICPADS | 3 |
| 2018 | Characterization of MPI usage on a production supercomputer
Sudheer Chunduri, Scott Parker, Pavan Balaji, Kevin Harms, Kalyan Kumaran |
SC | 3 |
| 2018 | Lessons learned from analyzing dynamic promotion for user-level threading
Shintaro Iwasaki, Abdelhalim Amer, Kenjiro Taura, Pavan Balaji |
SC | 4 |
| 2018 | On the adequacy of lightweight thread approaches for high-level parallel programming models
Adrián Castelló 0001, Rafael Mayo 0002, Kevin Sala, Vicenç Beltran 0001, Pavan Balaji, Antonio J. Peña |
Future Gener. Comput. Syst. | 5 |
| 2018 | Exploring the interoperability of remote GPGPU virtualization using rCUDA and directive-based programming models
Adrián Castelló 0001, Antonio J. Peña, Rafael Mayo 0002, Judit Planas, Enrique S. Quintana-Ortí, Pavan Balaji |
J. Supercomput. | 6 |
| 2018 | Argobots: A Lightweight Low-Level Threading and Tasking FrameworkabstractIn the past few decades, a number of user-level threading and tasking models have been proposed in the literature to address the shortcomings of OS-level threads, primarily with respect to cost and flexibility. Current state-of-the-art user-level threading and tasking models, however, either are too specific to applications or architectures or are not as powerful or flexible. In this paper, we present Argobots, a lightweight, low-level threading and tasking framework that is designed as a portable and performant substrate for high-level programming models or runtime systems. Argobots offers a carefully designed execution model that balances generality of functionality with providing a rich set of controls to allow specialization by end users or high-level programming models. We describe the design, implementation, and performance characterization of Argobots and present integrations with three high-level models: OpenMP, MPI, and colocated I/O services. Evaluations show that (1) Argobots, while providing richer capabilities, is competitive with existing simpler generic threading runtimes; (2) our OpenMP runtime offers more efficient interoperability capabilities than production OpenMP runtimes do; (3) when MPI interoperates with Argobots instead of Pthreads, it enjoys reduced synchronization costs and better latency-hiding capabilities; and (4) I/O services with Argobots reduce interference with colocated applications while achieving performance competitive with that of a Pthreads approach. Abdelhalim Amer, Pavan Balaji, Cyril Bordage, George Bosilca, Alex Brooks, Philip H. Carns, Adrián Castelló 0001, Damien Genet, Thomas Hérault, Shintaro Iwasaki, Prateek Jindal, Laxmikant V. Kalé, Sriram Krishnamoorthy, Jonathan Lifflander, Huiwei Lu, Esteban Meneses, Marc Snir, Yanhua Sun, Kenjiro Taura, Pete Beckman |
IEEE Trans. Parallel Distributed Syst. | 3 |
| 2018 | Dynamic Adaptable Asynchronous Progress Model for MPI RMA Multiphase ApplicationsabstractCasper is a process-based asynchronous progress model for MPI one-sided communication on multi- and many-core architectures. The one-sided communication is not truly one-sided in most MPI implementations: the target process still relies on software progress to complete incoming operations. Casper allows the user to specify an arbitrary number of cores dedicated to background ghost processes and transparently redirects the RMA operations to ghost processes by utilizing the PMPI redirection and MPI-3 shared-memory technologies. Although Casper benefits applications that suffer from lack of asynchronous progress, the operation redirection design might not support complex multiphase applications effectively, which often involve dynamically changing communication density and computing workloads. In this paper, we present an adaptive mechanism in Casper to address the limitation of static asynchronous progress in multiphase applications. We exploit two adaptive strategies, a user-guided strategy and a fully transparent and automatic strategy based on self-profiling and prediction, to dynamically reconfigure the asynchronous progress in Casper according to real-time performance characteristics during multiphase execution. We evaluate the adaptive approaches in both microbenchmarks and a real quantum chemistry application suite, NWChem, on the Cray XC30 supercomputer and an Intel Omni-Path cluster. Min Si, Antonio J. Peña, Jeff R. Hammond, Pavan Balaji, Masamichi Takagi, Yutaka Ishikawa |
IEEE Trans. Parallel Distributed Syst. | 4 |
| 2017 | Advanced Thread Synchronization for Multithreaded MPI ImplementationsabstractConcurrent multithreaded access to the Message Passing Interface (MPI) is gaining importance to support emerging hybrid MPI applications. The interoperability between threads and MPI, however, is complex and renders efficient implementations nontrivial. Prior studies showed that threads waiting for communication progress (waiting threads) often interfere with others (active threads) and degrade their progress. This situation occurs when both classes of threads compete for the same MPI resource and ownership passing to waiting threads does not guarantee communication to advance. The best-known practical solution prioritizes active threads and adapts first-in-first-out arbitration within each class. This approach, however, suffers from residual wasted resource acquisitions (waste) and ignores data locality, thus resulting in poor scalability. In this work, we propose thread synchronization improvements to eliminate waste while preserving data locality in a production MPI implementation. First, we leverage MPI knowledge and a fast synchronization method to eliminate waste and accelerate progress. Second, we rely on a cooperative progress model that dynamically elects and restricts a single waiting thread to drive a communication context for improved data locality. Third, we prioritize active threads and synchronize them with a locality-preserving lock that is hierarchical and exploits unbounded bias for high throughput. Results show significant improvement in synthetic microbenchmarks and two MPI+OpenMP applications. Hoang-Vu Dang, Abdelhalim Amer, Pavan Balaji |
CCGrid | 4 |
| 2017 | Scalable Assembly for Massive Genomic GraphsabstractScientists increasingly want to assemble large genomes, metagenomes, and large numbers of individual genomes. In order to meet the demand for processing these huge datasets, parallel genome assembly is a vital step. Among all the parallel genome assemblers, de Bruijn graph based ones are most popular. However, the size of de Bruijn graph is determined by the number of distinct kmers used in the algorithm, thus redundant kmers in the genome datasets donot contribute to the graph size. The scalability of genome assemblers is influenced directly by the distinct kmers in the dataset or de Bruijn graph size, rather than the input dataset size. In order to assembly large genomes, we have artificially created 16 datasets of 4 Terabytes in total from the human reference genome. The human reference genome is firstly mutated with a 5% mutation rate, and then subjected to a genome sequencing data simulator ART. The simulated datasets have linearly increasing number of distinct kmers as the size/number of the combined datasets increases. We then evaluate all five time-consuming steps of the SWAP-Assembler 2.0 (SWAP2) using these 16 simulated datasets. Compared with our previous experiment on 1000 human dataset with fixed de Bruijn graph size, the weak-scaling test shows that SWAP2 can scale well from 1024 cores using one dataset to 16,384 cores. The percentage of time usage for all five steps of SWAP2 is fixed, and total time usage is also constant. The result showed that the time usage of graph simplification occupied almost 75% of the total time usage, which will be subject to further optimization for future work. Jintao Meng 0001, Jianqiu Ge, Yanjie Wei, Pavan Balaji, Bingqiang Wang |
CCGrid | 5 |
| 2017 | A Performance Study of UCX over InfiniBandabstractUCX is an open-source communication framework with a two-level API design targeted at addressing the needs of large supercomputing systems. The lower-level interface, UCT, adds minimal overhead to data transfer but requires considerable effort from the user. The higher-level interface, UCP, is easier to use, but adds some overhead to the communication. This work focuses on charting the performance of UCX over InfiniBand, motivated by the usage of UCX as middleware for high-level communication libraries. We analyze performance shortcomings that stem from the two-level design and the sources of these performance losses. In particular, we target basic functions of UCP, evaluate their performance over InfiniBand, and analyze sources of overheads compared with UCT and Verbs. We propose and evaluate some fixes to minimize these overheads, in order to enhance UCP performance and scalability. Nikela Papadopoulou, Lena Oden, Pavan Balaji |
CCGrid | 3 |
| 2017 | S-Aligner: Ultrascalable Read Mapping on Sunway Taihu LightabstractThe availability and amount of sequenced genomes have been rapidly growing in recent years because of the adoption of next-generation sequencing (NGS) technologies that enable high-throughput short-read generation at highly competitive cost. Since this trend is expected to continue in the foreseeable future, the design and implementation of efficient and scalable NGS bioinformatics algorithms are important to research and industrial applications. In this paper, we introduce S-Aligner–a highly scalable read mapper designed for the Sunway Taihu Light supercomputer and its fourth-generationShenWei many-core architecture (SW26010). S-Aligner employs a combination of optimization techniques to overcome both the memory-bound and the compute-bound bottlenecks in the read mapping algorithm. In order to make full use of the compute power of Sunway Taihu Light, our design employs three levels of parallelism: (1) internode parallelism using MPI based on a task-grid pattern, (2) intranode parallelism using multithreading and asynchronous data transfer to fully utilize all 260 cores of the SW26010 many-core processor, and (3) vectorization to exploit the available 256-bit SIMD vector registers. Moreover, we have employed asynchronous access patterns and data-sharing strategies during file I/O to overcome bandwidth limitations of the network file system. Our performance evaluation demonstrates that S-Aligner scales almost linearly with approximately 95% efficiency for up to 13,312 nodes (concurrently harnessing more than 3 millioncompute cores). Furthermore, our implementation on a single node outperforms the established RazerS3 mapper running on a platform with eight Intel Xeon E7-8860v3 CPUs while achieving highly competitive alignment accuracy. Xiaohui Duan, Yuandong Chan, Christian Hundt 0002, Bertil Schmidt, Pavan Balaji |
CLUSTER | 6 |
| 2017 | GLT: A Unified API for Lightweight Thread Libraries
Adrián Castelló 0001, Rafael Mayo 0002, Pavan Balaji, Enrique S. Quintana-Ortí, Antonio J. Peña |
Euro-Par | 4 |
| 2017 | Exploiting Common Neighborhoods to Optimize MPI Neighborhood CollectivesabstractNeighborhood collectives were added to the Message Passing Interface (MPI) to better support sparse communication patterns found in many applications. These new collectives encourage more scalable programming styles, and greatly extend the scope of MPI collectives by allowing users to define their own collective communication patterns. In this paper, we describe a new, distributed algorithm for computing improved communication schedules for neighborhood collectives. We show how to discover common process neighborhoods in fully general MPI distributed graph topologies, and how to exploit this information to build message-combining communication schedules for the MPI neighborhood collectives. Our experimental results show considerable performance improvements for application communication topologies of various shapes and sizes. On average, the performance gain is around 50%, but it can also be as much as 71% for topologies with larger numbers of neighbors. Seyed Hessam Mirsadeghi, Jesper Larsson Träff, Pavan Balaji, Ahmad Afsahi |
HiPC | 3 |
| 2017 | Bloomfish: A Highly Scalable Distributed K-mer Counting FrameworkabstractK-mer counting is a fundamental operation in DNA research and genome analytics; its application includes estimating genome assembly, understanding similarities in genomic samples, and merging a newly processed genome with a reference genome. As the genome dataset becomes larger and larger, designing a highly optimized distributed-memory implementation becomes more and more important. Current distributed-memory solutions have two limitations: they have a high memory footprint, and they do not provide advanced optimizations for loading enormous genome datasets into memory. Based on these observations, we present Bloomfish, a distributed, memory-efficient, scalable solution to the limits of current work. To keep a low memory footprint, Bloomfish leverages the compact hash array design of the single-node Jellyfish system and the optimized workflow of the high-performance MapReduce framework Mimir. We have also codesigned Mimir's I/O to efficiently load enormous datasets. We ran Bloomfish on the Tianhe-2 supercomputer with large sequence datasets (up to 24 TB). Our results show that Bloomfish achieves unprecedented scalability in genome analytics. Yanfei Guo, Yanjie Wei, Bingqiang Wang, Yutong Lu, Pietro Cicotti, Pavan Balaji, Michela Taufer |
ICPADS | 7 |
| 2017 | Portable Topology-Aware MPI-I/OabstractRecent advances in storage devices are opening new opportunities in high-performance computing (HPC). Technologies such as solid-state drives (SSD) and non-volatile memories (NVM) are becoming increasingly popular because of the important gains they can represent for HPC. Indeed, novel architectures with deeper storage hierarchies populated with SSDs and/or NVM offer new ways to improve applications' performance. For instance, fast multilevel checkpointing or in-situ data analysis are some of the techniques that can be greatly improved thanks to these new technologies. However, optimizations made for one system can impose performance costs in another machine due to topology differences. To take advantage of increasingly complex systems, we propose extensions to MPI enabling codes to determine which nodes of a system share common features. Our approach provides a portable mechanism for resource discovery. It also lays the foundation for additional optimizations in checkpointing and in ROMIO. In this paper we present the design and implementation of such a feature and test it with multiple benchmarks. Our results demonstrate the benefits of this portable resource discovery functionality. Robert Latham, Leonardo Arturo Bautista-Gomez, Pavan Balaji |
ICPADS | 3 |
| 2017 | Hexe: A Toolkit for Heterogeneous Memory ManagementabstractHeterogeneity in memory is becoming increasingly common in high-end computing. Several modern supercomputers, such as those based on the Intel Knights Landing or NVIDIA P100 GPU architectures, already showcase multiple memory domains that are directly accessible by user applications, including on-chip high-bandwidth memory and off-chip traditional DDR memory. The next generation of supercomputers is expected to take this architectural trend one step further by including NVRAM as an additional byte-addressable memory option. Despite these trends, allocating and managing such memory are still tedious tasks. In this paper, we present hexe, a highly flexible and portable memory allocation toolkit. Unlike other memory allocation tools such as malloc, memkind, and cudaMallocManaged, hexe presents a rich and portable memory allocation framework that allows applications to carefully and precisely manage their memory across the various memory subsystems available on the system. Together with a detailed description of the design and capabilities of hexe, we present several case studies where the flexible memory allocation in hexe allows applications to achieve superior performance compared with that of other memory allocation tools. Lena Oden, Pavan Balaji |
ICPADS | 2 |
| 2017 | Parallel I/O Optimizations for Scalable Deep LearningabstractAs deep learning systems continue to grow in importance, researchers have been analyzing approaches to make such systems efficient and scalable on high-performance computing platforms. As computational parallelism increases, however, data I/O becomes the major bottleneck limiting the overall system scalability. In this paper, we continue our efforts to improve LMDB, the I/O subsystem of the Caffe deep learning framework. In a previous paper we presented LMDBIO---an optimized I/O plugin for Caffe that takes into account the data access pattern of Caffe in order to vastly improve I/O performance. Nevertheless, LMDBIO's optimizations, which we henceforth call LMM (localized mmap), are limited to intranode performance, and these optimizations do little to minimize the I/O inefficiencies in distributed-memory environments. In this paper, we propose LMDBIO-DM, an enhanced version of LMDBIO-LMM that optimizes the I/O access of Caffe in distributed-memory environments. We present several sophisticated data I/O techniques that allow for significant improvement in such environments. Our experimental results show that LMDBIO-DM can improve the overall execution time of Caffe by more than 30-fold compared with LMDB and by 2-fold compared with LMDBIO-LMM. Sarunya Pumma, Min Si, Wu-chun Feng, Pavan Balaji |
ICPADS | 4 |
| 2017 | GLTO: On the Adequacy of Lightweight Thread Approaches for OpenMP ImplementationsabstractOpenMP is the de facto standard application programming interface (API) for on-node parallelism. The most popular OpenMP runtimes rely on POSIX threads (pthreads) implementations that offer an excellent performance for coarse-grained parallelism and match perfectly with the current hardware. However, a recent trend in runtimes/applications points in the direction of leveraging massive on-node parallelism in conjunction with fine-grained and dynamic scheduling paradigms. It has been demonstrated that lightweight thread (LWT) solutions are more appropriate for these new parallel paradigms. We have developed GLTO, an OpenMP implementation over the recently-emerged Generic Lightweight Threads (GLT) API. GLT exports a common API for LWT libraries that offers the possibility of running the same application over different native LWT solutions. In this paper we use GLTO to analyze different scenarios where OpenMP implementations may benefit from the use of either LWT or pthreads. Our study reveals that none of the threading approaches obtains the best performance in all the scenarios, but that there are important gaps among them. Adrián Castelló 0001, Rafael Mayo 0002, Pavan Balaji, Enrique S. Quintana-Ortí, Antonio J. Peña |
ICPP | 4 |
| 2017 | Mimir: Memory-Efficient and Scalable MapReduce for Large Supercomputing SystemsabstractIn this paper we present Mimir, a new implementation of MapReduce over MPI. Mimir inherits the core principles of existing MapReduce frameworks, such as MR-MPI, while redesigning the execution model to incorporate a number of sophisticated optimization techniques that achieve similar or better performance with significant reduction in the amount of memory used. Consequently, Mimir allows significantly larger problems to be executed in memory, achieving large performance gains. We evaluate Mimir with three benchmarks on two highend platforms to demonstrate its superiority compared with that of other frameworks. Yanfei Guo, Boyu Zhang 0002, Pietro Cicotti, Yutong Lu, Pavan Balaji, Michela Taufer |
IPDPS | 6 |
| 2017 | Memory Compression Techniques for Network Address Management in MPIabstractMPI allows applications to treat processes as a logical collection of integer ranks for each MPI communicator, while internally translating these logical ranks into actual network addresses. In current MPI implementations the management and lookup of such network addresses use memory sizes that are proportional to the number of processes in each communicator. In this paper, we propose a new mechanism, called AV-Rankmap, for managing such translation. AV-Rankmap takes advantage of logical patterns in rank-address mapping that most applications naturally tend to have, and it exploits the fact that some parts of network address structures are naturally more performance critical than others. It uses this information to compress the memory used for network address management. We demonstrate that AV-Rankmap can achieve performance similar to or better than that of other MPI implementations while using significantly less memory. Yanfei Guo, Charles Archer, Michael Blocksome, Scott Parker, Wesley Bland, Kenneth Raffenetti, Pavan Balaji |
IPDPS | 7 |
| 2017 | Why is MPI so slow?: analyzing the fundamental limits in implementing MPI-3.1abstractThis paper provides an in-depth analysis of the software overheads in the MPI performance-critical path and exposes mandatory performance overheads that are unavoidable based on the MPI-3.1 specification. We first present a highly optimized implementation of the MPI-3.1 standard in which the communication stack---all the way from the application to the low-level network communication API---takes only a few tens of instructions. We carefully study these instructions and analyze the root cause of the overheads based on specific requirements from the MPI standard that are unavoidable under the current MPI standard. We recommend potential changes to the MPI standard that can minimize these overheads. Our experimental results on a variety of network architectures and applications demonstrate significant benefits from our proposed changes. Kenneth Raffenetti, Abdelhalim Amer, Lena Oden, Charles Archer, Wesley Bland, Hajime Fujita 0002, Yanfei Guo, Tomislav Janjusic, Dmitry Durnov, Michael Blocksome, Min Si, Akhil Langer, Gengbin Zheng, Masamichi Takagi, Paul K. Coffman, Sayantan Sur, Alexander Sannikov, Sergey Oblomov, Michael Chuvelev, Masayuki Hatanaka, Paul F. Fischer, Thilina Ratnayaka, Matthew Otten, Misun Min, Pavan Balaji |
SC | 28 |
| 2017 | Foreword to the Special Issue of the workshop on the seventh international workshop on programming models and applications for multicores and manycores (PMAM 2016)abstractForeword to the Special Issue of the workshop on the seventh international workshop on programming models and applications for Pavan Balaji, Kai-Cheung Leung |
Concurr. Comput. Pract. Exp. | 1 |
| 2017 | Enabling scalable and accurate clustering of distributed ligand geometries on supercomputers
Boyu Zhang 0002, Trilce Estrada, Pietro Cicotti, Pavan Balaji, Michela Taufer |
Parallel Comput. | 4 |
| 2016 | A Review of Lightweight Thread Approaches for High Performance ComputingabstractHigh-level, directive-based solutions are becoming the programming models (PMs) of the multi/many-core architectures. Several solutions relying on operating system (OS) threads perfectly work with a moderate number of cores. However, exascale systems will spawn hundreds of thousands of threads in order to exploit their massive parallel architectures and thus conventional OS threads are too heavy for that purpose. Several lightweight thread (LWT) libraries have recently appeared offering lighter mechanisms to tackle massive concurrency. In order to examine the suitability of LWTs in high-level runtimes, we develop a set of microbenchmarks consisting of commonly-found patterns in current parallel codes. Moreover, we study the semantics offered by some LWT libraries in order to expose the similarities between different LWT application programming interfaces. This study reveals that a reduced set of LWT functions can be sufficient to cover the common parallel code patterns andthat those LWT libraries perform better than OS threads-based solutions in cases where task and nested parallelism are becoming more popular with new architectures. Adrián Castelló 0001, Antonio J. Peña, Rafael Mayo 0002, Pavan Balaji, Enrique S. Quintana-Ortí |
CLUSTER | 5 |
| 2016 | Compiler-Assisted Overlapping of Communication and Computation in MPI ApplicationsabstractThe performance of distributed-memory applications, many of which are written in MPI, critically depends on how well the applications can ameliorate the long latency of data movement by overlapping them with ongoing computations, thereby minimizing wait time. This paper presents a study of the various optimization techniques to enable such overlapping in large MPI applications and presents a framework that uses an analytical performance model and an optimizing compiler to systematically enable a majority of such optimizations. In particular, we first generate an analytical performance model of the application execution flow to automatically identify potential communication hot spots that may induce long wait time. Next, for each communication hot spot, we search the execution flow graph to find surrounding loops that include sufficient local computation to overlap with the communication. Then, blocking MPI communications are decoupled into non-blocking operations when necessary, and their surrounding loop is transformed to hide the communication latencies behind local computations. We evaluated our framework using 7 MPI applications from the NAS benchmark suite. Our optimizations can attain 3% to 72% speedup over the original implementations. Jichi Guo, Qing Yi, Jiayuan Meng, Junchao Zhang 0002, Pavan Balaji |
CLUSTER | 5 |
| 2016 | One-Sided Interface for Matrix Operations Using MPI-3 RMA: A Case Study with ElementalabstractA one-sided programming model separates communication from synchronization, and is the driving principle behind partitioned global address space (PGAS) libraries such as Global Arrays (GA) and SHMEM. PGAS models expose a rich set of functionality that a developer needs in order to implement mathematical algorithms that require frequent multidimensional array accesses. However, use of existing PGAS libraries in application codes often requires significant development effort in order to fully exploit these programming models. On the other hand, a vast majority of scientific codes use MPI either directly or indirectly via third-party scientific computation libraries, and need features to support application-specific communication requirements (e.g., asynchronous update of distributed sparse matrices, commonly arising in machine learning workloads). For such codes it is often impractical to completely shift programming models in favor of special one-sided communication middleware. Instead, an elegant and productive solution is to exploit the one-sided functionality already offered by MPI-3 RMA (Remote Memory Access). We designed a general one-sided interface using the MPI-3 passive RMA model for remote matrix operations in the linear algebra library Elemental, we call the interface we designed RMAInterface. Elemental is an open source library for distributed-memory dense and sparse linear algebra and optimization. We employ RMAInterface to construct a Global Arrays-like API and demonstrate its performance scalability and competitivity with that of the existing GA (with ARMCI-MPI) for a quantum chemistry application. Jeff R. Hammond, Antonio J. Peña, Pavan Balaji, Assefaw Hadish Gebremedhin, Barbara M. Chapman |
ICPP | 4 |
| 2016 | SWAP-Assembler 2: Optimization of De Novo Genome Assembler at Extreme ScaleabstractIn this paper, we analyze and optimize the most time-consuming steps of the SWAP-Assembler, a parallel genome assembler, so that it can scale to a large number of cores for huge genomes with sequencing data ranging from terabyes to petabytes. Performance analysis results show that the most time-consuming steps are input parallelization, k-mer graph construction, and graph simplification (edge merging). For the input parallelization, the input data is divided into virtual fragments with nearly equal size, and the start position and end position of each fragment are automatically separated at the beginning of the reads. In k-mer graph construction, in order to improve the communication efficiency, the message size is kept constant between any two processes by proportionally increasing the number of nucleotides to the number of processes in the input parallelization step for each round. The memory usage is also decreased because only a small part of the input data is processed in each round. With graph simplification, the communication protocol reduces the number of communication loops from four to two loops and decreases the idle communication time. The optimized assembler is denoted SWAP-Assembler 2 (SWAP2). In our experiments using a 1000 Genomes project dataset of 4 terabytes (the largest dataset ever used for assembling) on the supercomputer Mira, the results show that SWAP2 scales to 131,072 cores with an efficiency of 40%. We also compared our work with both the HipMer assembler and the SWAP-Assembler. On the Yanhuang dataset of 300 gigabytes, SWAP2 shows a 3X speedup and 4X better scalability compared with the HipMer assembler and is 45 times faster than the SWAP-Assembler. The SWAP2 software is available at https://sourceforge.net/projects/swapassembler. Jintao Meng 0001, Pavan Balaji, Yanjie Wei, Bingqiang Wang, Shengzhong Feng |
ICPP | 3 |
| 2016 | Scalability Challenges in Current MPI One-Sided ImplementationsabstractMPI one-sided or remote memory access (RMA) communication provides a different execution model from traditional two-sided or group communication and is better suited for some classes of applications. However, current implementations of MPI RMA are notorious for their inability to scale to large systems or problem sizes. In this paper, we present a study of the RMA infrastructure in popular open-source MPI implementations. Our objective is to identify critical scalability limitations with respect to memory usage in these implementations. We then perform a thorough evaluation on two cluster computers to demonstrate those scalability limitations, and we provide suggestions on how they can be alleviated. Pavan Balaji, William Gropp |
ISPDC | 2 |
| 2016 | Work stealing for GPU-accelerated parallel programs in a global address space frameworkabstractSummary Task parallelism is an attractive approach to automatically load balance the computation in a parallel system and adapt to dynamism exhibited by parallel systems. Exploiting task parallelism through work stealing has been extensively studied in shared and distributed‐memory contexts. In this paper, we study the design of a system that uses work stealing for dynamic load balancing of task‐parallel programs executed on hybrid distributed‐memory CPU‐graphics processing unit (GPU) systems in a global‐address space framework. We take into account the unique nature of the accelerator model employed by GPUs, the significant performance difference between GPU and CPU execution as a function of problem size, and the distinct CPU and GPU memory domains. We consider various alternatives in designing a distributed work stealing algorithm for CPU‐GPU systems, while taking into account the impact of task distribution and data movement overheads. These strategies are evaluated using microbenchmarks that capture various execution configurations as well as the state‐of‐the‐art CCSD(T) application module from the computational chemistry domain. Copyright © 2016 John Wiley & Sons, Ltd. Humayun Arafat, James Dinan, Sriram Krishnamoorthy, Pavan Balaji, P. Sadayappan |
Concurr. Comput. Pract. Exp. | 4 |
| 2016 | Programming models and applications for multicores and manycoresabstractRapid advancements in multicore and manycore chips have been a revolution within chip manufacturing, almost eradicating single-core processors. From high-end servers to mobile phones, multicore and manycore chips are steadily entering every single aspect of information technology. However, programming multicore and manycore architectures remains challenging today. To fully utilize these chips, parallel programming models that allow sequential programs and programs utilizing limited parallelism to transition to architectures with massive parallelism, while maintaining good performance and productive development, are urgently needed. This special issue contains six articles selected from the 2014 International Workshop on Programming Models and Applications for Multicores and Manycores (PMAM 2014). These articles cover issues in parallel programming languages and models, work-stealing in parallel runtime environments, heterogeneous computing, and their implications on computational patterns. In ‘Adaptive Demand-aware Work-Stealing in Multi-programmed Multi-core Architectures’ 1, Chen et al. discuss how work-stealing models can be extended to support environments where multiple parallel programs might be concurrently executing. They propose a new work-stealing algorithm called demand-aware work-stealing, which alleviates this issue by allowing programs to donate and steal idle cores as needed. In ‘Palirria: Accurate On-line Parallelism Estimation for Adaptive Work-Stealing’ 2, Varisteas and Brorsson present a self-adapting work-stealing scheduling method for nested fork/join parallelism called ‘Palirria’. The proposed approach can be used to estimate the number of utilizable workers and self-adapt accordingly. In ‘Efficient CPU-GPU cooperative computing for solving the subset-sum problem’ 3, Wan et al. study the impact of heterogeneous computing with CPUs and Graphics Processing Unit (GPUs) for solving the subset-sum problem. The authors observe that while heterogeneous computing is prominent, not much study has been performed in the simultaneous usage of both CPU and GPU resources during computation. This paper proposes an efficient CPU–GPU cooperative computing scheme for solving the subset-sum problem, which enables the full utilization of all the computing power of both CPUs and GPUs. In ‘Dynamic Partitioning-based JPEG Decompression on Heterogeneous Multicore Architectures’ 4, Sodsong et al. introduce a novel JPEG decoding scheme for heterogeneous architectures consisting of a CPU and a general-purpose GPU. They employ an offline profiling step to determine the performance of a system's CPU and GPU with respect to JPEG decoding. Then the runtime partitioning and scheduling scheme exploits task, data, and pipeline parallelism by scheduling the non-parallelizable entropy decoding task on the CPU, whereas inverse discrete cosine transformations, color conversions, and upsampling are conducted on both the CPU and the GPU. In ‘Compiler Transformation of Nested Loops for GPGPUs’ 5, Tian et al. present their experiences in creating an open-source OpenACC compiler in an industrial framework (OpenUH as a branch of Open64). They also discuss in detail the techniques that they developed for loop-scheduling reduction operations on General Purpose Graphics Processing Unit (GPGPUs). In ‘Vectorizing Unstructured Mesh Computations for Many-core Architectures’ 6, Reguly et al. present results on achieving high performance through vectorization on CPUs and the Xeon-Phi on a key class of irregular applications: unstructured mesh computations. Using Single Instruction Multiple Threads (SIMT) and Single Instruction Multiple Data (SIMD) programming models, they show how unstructured mesh computations map to OpenCL or vector intrinsics through the use of code generation techniques in the OP2 Domain Specific Library and explore how irregular memory accesses and race conditions can be organized on different hardware. We hope that the articles in this special issue will provide readers with relevant insights into the emerging parallel programming models for multicore and manycore systems. Pavan Balaji holds appointments as a Computer Scientist and Group Lead at the Argonne National Laboratory, as an Institute Fellow of the Northwestern-Argonne Institute of Science and Engineering at Northwestern University, and as a Research Fellow of the Computation Institute at the University of Chicago. He leads the Programming Models and Runtime Systems group at Argonne. His research interests include parallel programming models and runtime systems for communication and I/O on extreme-scale supercomputing systems, modern system architecture, cloud computing systems, data-intensive computing, and big-data sciences. He has nearly 150 publications in these areas and has delivered nearly 150 talks and tutorials at various conferences and research institutes. Dr Balaji is a recipient of several awards including the U.S. Department of Energy Early Career award in 2012, TEDxMidwest Emerging Leader award in 2013, Crain's Chicago 40 under 40 award in 2012, Los Alamos National Laboratory Director's Technical Achievement award in 2005, Ohio State University Outstanding Researcher award in 2005, six best paper awards, one best paper finalist, and one best poster finalist. He has served as a chair or editor for nearly 50 journals, conferences, and workshops and as a technical program committee member in numerous conferences and workshops. He is a senior member of the IEEE and a professional member of the ACM. More details about Dr Balaji are available at http://www.mcs.anl.gov/balaji. Contact him at [email protected] Zhiyi Huang is an Associate Professor at the University of Otago, New Zealand. His research interests include parallel and distributed computing, multicore architectures, parallel programming models and environments, task scheduling, operating systems, green computing, and computer networks. More details about Prof. Zhiyi Huang are available at http://www.cs.otago.ac.nz/staffpriv/hzy. Contact him at [email protected] Pavan Balaji, Zhiyi Huang 0001 |
Concurr. Comput. Pract. Exp. | 1 |
| 2016 | An implementation and evaluation of the MPI 3.0 one-sided communication interfaceabstractSummary The Message Passing Interface (MPI) 3.0 standard includes a significant revision to MPI's remote memory access (RMA) interface, which provides support for one‐sided communication. MPI‐3 RMA is expected to greatly enhance the usability and performance of MPI RMA. We present the first complete implementation of MPI‐3 RMA and document implementation techniques and performance optimization opportunities enabled by the new interface. Our implementation targets messaging‐based networks and is publicly available in the latest release of the MPICH MPI implementation. Using this implementation, we explore the performance impact of new MPI‐3 functionality and semantics. Results indicate that the MPI‐3 RMA interface provides significant advantages over the MPI‐2 interface by enabling increased communication concurrency through relaxed semantics in the interface and additional routines that provide new window types, synchronization modes, and atomic operations. Copyright © 2016 John Wiley & Sons, Ltd. James Dinan, Pavan Balaji, Darius Buntinas, David Goodell, William Gropp, Rajeev Thakur |
Concurr. Comput. Pract. Exp. | 2 |
| 2016 | Performance analysis of data intensive cloud systems based on data management and replication: a survey
Saif Ur Rehman Malik, Samee Ullah Khan, Sam J. Ewen, Nikos Tziritas, Joanna Kolodziej, Albert Y. Zomaya, Sajjad Ahmad Madani, Nasro Min-Allah, Lizhe Wang 0001, Cheng-Zhong Xu 0001, Qutaibah M. Malluhi, Johnatan E. Pecero, Pavan Balaji, Abhinav Vishnu, Rajiv Ranjan 0001, Sherali Zeadally, Hongxiang Li 0001 |
Distributed Parallel Databases | 13 |
| 2016 | MultiCL: Enabling automatic scheduling for task-parallel workloads in OpenCL
Ashwin M. Aji, Antonio J. Peña, Pavan Balaji, Wu-chun Feng |
Parallel Comput. | 3 |
| 2016 | Special Issue on Parallel Programming Models and Systems Software for High-End Computing
Pavan Balaji, Abhinav Vishnu, Yong Chen 0001 |
Parallel Comput. | 1 |
| 2016 | A data-oriented profiler to assist in data partitioning and distribution for heterogeneous memory in HPC
Antonio J. Peña, Pavan Balaji |
Parallel Comput. | 2 |
| 2016 | Special Issue on Cluster Computing
Michela Taufer, Pavan Balaji, Satoshi Matsuoka |
Parallel Comput. | 2 |
| 2016 | MPI-ACC: Accelerator-Aware MPI for Scientific ApplicationsabstractData movement in high-performance computing systems accelerated by graphics processing units (GPUs) remains a challenging problem. Data communication in popular parallel programming models, such as the Message Passing Interface (MPI), is currently limited to the data stored in the CPU memory space. Auxiliary memory systems, such as GPU memory, are not integrated into such data movement standards, thus providing applications with no direct mechanism to perform end-to-end data movement. We introduce MPI-ACC, an integrated and extensible framework that allows end-to-end data movement in accelerator-based systems. MPI-ACC provides productivity and performance benefits by integrating support for auxiliary memory spaces into MPI. MPI-ACC supports data transfer among CUDA, OpenCL and CPU memory spaces and is extensible to other offload models as well. MPI-ACC's runtime system enables several key optimizations, including pipelining of data transfers, scalable memory management techniques, and balancing of communication based on accelerator and node architecture. MPI-ACC is designed to work concurrently with other GPU workloads with minimum contention. We describe how MPI-ACC can be used to design new communication-computation patterns in scientific applications from domains such as epidemiology simulation and seismology modeling, and we discuss the lessons learned. We present experimental results on a state-of-the-art cluster with hundreds of GPUs; and we compare the performance and productivity of MPI-ACC with MVAPICH, a popular CUDA-aware MPI solution. MPI-ACC encourages programmers to explore novel application-specific optimizations for improved overall cluster utilization. Ashwin M. Aji, Lokendra S. Panwar, Karthik Murthy, Milind Chabbi, Pavan Balaji, Keith R. Bisset, James Dinan, Wu-chun Feng, John M. Mellor-Crummey, Xiaosong Ma, Rajeev Thakur |
IEEE Trans. Parallel Distributed Syst. | 6 |
| 2015 | Characterizing MPI and Hybrid MPI+Threads Applications at Scale: Case Study with BFSabstractWith the increasing prominence of many-core architectures and decreasing per-core resources on large supercomputers, a number of applications developers are investigating the use of hybrid MPI+threads programming to utilize computational units while sharing memory. An MPI-only model that uses one MPI process per system core is capable of effectively utilizing the processing units, but it fails to fully utilize the memory hierarchy and relies on fine-grained internodes communication. Hybrid MPI+threads models, on the other hand, can handle internodes parallelism more effectively and alleviate some of the overheads associated with internodes communication by allowing more coarse-grained data movement between address spaces. The hybrid model, however, can suffer from locking and memory consistency overheads associated with data sharing. In this paper, we use a distributed implementation of the breadth-first search algorithm in order to understand the performance characteristics of MPI-only and MPI+threads models at scale. We start with a baseline MPI-only implementation and propose MPI+threads extensions where threads independently communicate with remote processes while cooperating for local computation. We demonstrate how the coarse-grained communication of MPI+threads considerably reduces time and space overheads that grow with the number of processes. At large scale, however, these overheads constitute performance barriers for both models and require fixing the root causes, such as the excessive polling for communication progress and inefficient global synchronizations. To this end, we demonstrate various techniques to reduce such overheads and show performance improvements on up to 512K cores of a Blue Gene/Q system. Abdelhalim Amer, Huiwei Lu, Pavan Balaji, Satoshi Matsuoka |
CCGRID | 3 |
| 2015 | Lessons Learned Implementing User-Level Failure Mitigation in MPICHabstractUser-level failure mitigation (ULFM) is becoming the front-running solution for process fault tolerance in MPI. While not yet adopted into the MPI standard, it is being used by applications and libraries and is being considered by the MPI Forum for future inclusion into MPI itself. In this paper, we introduce an implementation of ULFM in MPICH, a high-performance and widely portable implementation of the MPI standard. We demonstrate that while still a reference implementation, the runtime cost of the new API calls introduced is relatively low. Wesley Bland, Huiwei Lu, Pavan Balaji |
CCGRID | 4 |
| 2015 | SWAP-Assembler 2: Scalable Genome Assembler towards Millions of Cores - Practice and ExperienceabstractThere is widening gap between the throughput of massive parallel sequencing machines and the ability to analyze these huge sequencing data, which can be Tara bytes or even Peta bytes. Previously our assembly tool, SWAP-Assembler, can scale to 2048 cores on TianHe 1A for human Yanhuang genome. This work is to further scale SWAP-Assembler to millions of cores on Mira. SWAP-Assembler can be divided into 5 steps, and the most time consuming steps are input parallelization, kmer graph construction, graph simplification (edge merging). We optimize these three steps to keep the percentage of time usage in each step constant when the number of cores increases. For the input parallelization step, the input data is divided into virtual fragments with almost equal size, the begin position and end position for each fragment is automatically separated at the beginning symbol of reads. This data blocking strategy plays a central role in adjusting the data size to keep the communication and memory efficiency for the subsequent steps. In kmer graph construction, to prevent the communication efficiency degradation, the message size is kept constant (about 8k bytes) between any two processes by proportionally increasing the number of nucleotides to the number of processes in the input parallelization step in each round. The memory usage can be also benefited, as only a small part of the input data is processed in each round. Within graph simplification, the major improvement is to combine messages sending & receiving between its two neighbors into one loop in the communication protocol. After integrated with the above optimizations, the new assembly tool is denoted as SWAP-Assembler 2 or SWAP2 for short. In our experiment for 1k human genome dataset, the modified SWAP-Assembler 2 can scale to 16k cores with parallel efficiency of 70%. Jintao Meng 0001, Yanjie Wei, Pavan Balaji |
CCGRID | 4 |
| 2015 | Understanding Data Access Patterns Using Object-Differentiated Memory ProfilingabstractThe information provided by commonly used code-oriented profilers can be complemented by means of data-oriented profiling techniques. Based on a data-oriented approach, in this study we leverage techniques developed in previous papers to analyze the data access characteristics of a range of U.S. Department of Energy applications representative of different application domains. By analyzing object-differentiated memory access profiles, we identify markedly different access patterns across application stages. We find read-only and read-write periods, relatively large periods without accessing particular objects, and a variety of data access rates. This information is useful for devising software optimizations, for software and hardware code sign, and for data distribution and partitioning in heterogeneous memory systems. Antonio J. Peña, Pavan Balaji |
CCGRID | 2 |
| 2015 | Toward Implementing Robust Support for Portals 4 Networks in MPICHabstractThe Portals 4 network specification is a low-levelAPI for high-performance networks developed by Sandia National Laboratories, Intel Corporation, and the University of NewMexico. Portals 4 is specifically designed to support both the MPIand PGAS programming models efficiently by providing building blocks upon which to implement their particular features. In this paper we discuss our ongoing efforts to add efficient and robust support for Portals 4 networks inside MPICH, and we describe how the API semantics influenced our design. In particular, we found the lack of reliability guarantees from the Portals4 layer challenging to address. To tackle this situation, we implemented an intermediate layer - Rportals (reliable Portals), which modularizes the reliability functionality within our Portals network module for MPICH. In this paper we present theRportals design and its performance impact. Kenneth Raffenetti, Antonio J. Peña, Pavan Balaji |
CCGRID | 3 |
| 2015 | Implementation and Evaluation of MPI Nonblocking Collective I/OabstractThe well-known gap between relative CPU speeds and storage bandwidth results in the need for new strategies for managing I/O demands. In large-scale MPI applications, collective I/O has long been an effective way to achieve higher I/O rates, but it poses two constraints. First, although overlapping collective I/O and computation represents the next logical step toward a faster time to solution, MPI's existing collective I/O API provides only limited support for doing so. Second, collective routines (both for I/O and communication) impose a synchronization cost in addition to a communication cost. The upcoming MPI 3.1 standard will provide a new set of nonblocking collective I/O operations to satisfy the need of applications. We present here initial work on the implementation of MPI nonblocking collective I/O operations in the MPICH MPI library. Our implementation begins with the extended two-phase algorithm used in ROMIO's collective I/O implementation. We then utilize a state machine and the extended generalized request interface to maintain the progress of nonblocking collective I/O operations. The evaluation results indicate that our implementation performs as well as blocking collective I/O in terms of I/O bandwidth and is capable of overlapping I/O and other operations. We believe that our implementation can help users try nonblocking collective I/O operations in their applications. Robert Latham, Junchao Zhang 0002, Pavan Balaji |
CCGRID | 4 |
| 2015 | Techniques for Enabling Highly Efficient Message Passing on Many-Core ArchitecturesabstractMany-core architecture provides a massively parallel environment with dozens of cores and hundreds of hardware threads. Scientific application programmers are increasingly looking at ways to utilize such large numbers of lightweight cores for various programming models. Efficiently executing these models on massively parallel many-core environments is not easy, however and performance may be degraded in various ways. The first author's doctoral research focuses on exploiting the capabilities of many-core architectures on widely used MPI implementations. While application programmers have studied several approaches to achieve better parallelism and resource sharing, many of those approaches still face communication problems that degrade performance. In the thesis, we investigate the characteristics of MPI on such massively threaded architectures and propose two efficient strategies -- a multi-threaded MPI approach and a process-based asynchronous model -- to optimize MPI communication for modern scientific applications. Min Si, Pavan Balaji, Yutaka Ishikawa |
CCGRID | 2 |
| 2015 | Scaling NWChem with Efficient and Portable Asynchronous Communication in MPI RMAabstractNWChem is one of the most widely used computational chemistry application suites for chemical and biological systems. Despite its vast success, the computational efficiency of NWChem is still low. This is especially true in higher accuracy methods such as the CCSD(T) coupled cluster method, where it currently achieves a mere 50% computational efficiency when run at large scales. In this paper, we demonstrate the most computationally efficient scaling of NWChem CCSD(T) to date, and use it to solve large water clusters. We use our recently proposed process-based asynchronous progress framework for MPI RMA, called Casper, to scale the computation on water clusters at near-100% computational efficiency on up to 12288 cores. Min Si, Antonio J. Peña, Jeff R. Hammond, Pavan Balaji, Yutaka Ishikawa |
CCGRID | 4 |
| 2015 | Accurate Scoring of Drug Conformations at the Extreme ScaleabstractWe present a scalable method to extensively search for and accurately select pharmaceutical drug candidates in large spaces of drug conformations computationally generated and stored across the nodes of a large distributed system. For each legend conformation in the dataset, our method first extracts relevant geometrical properties and transforms the properties into a single metadata point in the three-dimensional space. Then, it performs an ochre-based clustering on the metadata to search for predominant clusters. Our method avoids the need to move legend conformations among nodes because it extracts relevant data properties locally and concurrently. By doing so, we can perform accurate and scalable distributed clustering analysis on large distributed datasets. We scale the analysis of our pharmaceutical datasets a factor of 400X higher in performance and 500X larger in size than ever before. We also show that our clustering achieves higher accuracy compared with that of traditional clustering methods and conformational scoring based on minimum energy. Boyu Zhang 0002, Trilce Estrada, Pietro Cicotti, Pavan Balaji, Michela Taufer |
CCGRID | 4 |
| 2015 | Runtime Support for Irregular Computation in MPI-Based ApplicationsabstractIn recent years more and more applications have been using irregular computation models in various domains such as bioinformatics and social network analysis. Traditional data movement approaches are not well suited for such applications because of the irregular communication patterns, sparse data structures, fast growth rate of data movement as system size or problem size rises, and so forth. Active Messages (AM) is an alternative programming paradigm that is more suitable for irregular computations. It allows small pieces of data to be dynamically moved to the remote process and certain computation to be triggered, and the remote process does not need to explicitly receive the data. In this paper, an outline of the first author's Ph.D. thesis, focusing on runtime support for irregular computation, is presented. In the first part, we combine the capability of AM with traditional MPI data movement patterns, and we propose a generalized MPI-interoperable AM framework (MPI-AM). In the second part, we extend the MPI-AM framework to provide a model of dynamic task parallelism for data-driven computation. In each part we describe critical issues, demonstrate the current status of the work and performance gain, and discuss remaining challenges to be solved. Pavan Balaji, William Gropp |
CCGRID | 2 |
| 2015 | Analyzing MPI-3.0 Process-Level Shared Memory: A Case Study with Stencil ComputationsabstractThe recently released MPI-3.0 standard introduced a process-level shared-memory interface which enables processes within the same node to have direct load/store access to each others' memory. Such an interface allows applications to declare data structures that are shared by multiple MPI processes on the node. In this paper, we study the capabilities and performance implications of using MPI-3.0 shared memory, in the context of a five-point stencil computation. Our analysis reveals that the use of MPI-3.0 shared memory has several unforeseen performance implications including disrupting certain compiler optimizations and incorrectly using suboptimal page sizes inside the OS. Based on this analysis, we propose several methodologies for working around these issues and improving communication performance by 40-85% compared to the current MPI-1.0 based approach. Junchao Zhang 0002, Kazutomo Yoshii, Shigang Li 0002, Yunquan Zhang, Pavan Balaji |
CCGRID | 6 |
| 2015 | Automatic Command Queue Scheduling for Task-Parallel Workloads in OpenCLabstractOpenCL is a portable interface that can be used to program cluster nodes with heterogeneous compute devices. The OpenCL specification tightly binds its workflow abstraction, or "command queue," to a specific device for the entire program. For best performance, the user has to find the ideal queue -- device mapping at command queue creation time, an effort that requires a thorough understanding of the match between the characteristics of all the underlying device architectures and the kernels in the program. In this paper, we propose to add scheduling attributes to the OpenCL context and command queue objects that can be leveraged by an intelligent runtime scheduler to automatically perform ideal queue - device mapping. Our proposed extensions enable the average OpenCL programmer to focus on the algorithm design rather than scheduling and automatically gain performance without sacrificing programmability. As an example, we design and implement an OpenCL runtime for task-parallel workloads, called MultiCL, which efficiently schedules command queues across devices. Within MultiCL, we implement several key optimizations to reduce runtime overhead. Our case studies include the SNU-NPB OpenCL benchmark suite and a real-world seismology simulation. We show that, on average, users have to apply our proposed scheduler extensions to only four source lines of code in existing OpenCL applications in order to automatically benefit from our runtime optimizations. We also show that MultiCL always maps command queues to the optimal device set with negligible runtime overhead. Ashwin M. Aji, Antonio J. Peña, Pavan Balaji, Wu-chun Feng |
CLUSTER | 3 |
| 2015 | Flexible Error Recovery Using Versions in Global View ResilienceabstractWe present the Global View Resilience (GVR) system, a library that enables applications to add resilience in a portable, application-controlled fashion using versioned distributed arrays. We briefly describe GVR's interfaces for distributed arrays, versioning, and cross-layer error recovery. We illustrate how GVR can be used for rollback recovery and a wide range additional error recovery techniques including forward recovery for latent errors or silent data corruptions. Application results demonstrate that GVR's interfaces and implementation are portable, flexible (support a variety of recovery models), efficient and create a gentle-slope path to tolerate growing error rates in future systems. Nan Dun, Hajime Fujita 0002, Aiman Fang, Andrew A. Chien, Pavan Balaji, Kamil Iskra, Wesley Bland, Andrew R. Siegel |
CLUSTER | 6 |
| 2015 | Empirical Comparison of Three Versioning ArchitecturesabstractFuture supercomputer systems will face serious reliability challenges. Among failure scenarios, latent errors are some of the most serious and concerning. Preserving multiple versions of critical data is a promising approach to deal with such errors. We are developing the Global View Resilience (GVR) library, with multi-version global arrays as one of the key features. This paper presents three array versioning architectures: flat array, flat array with change tracking, and log-structured array. We use a synthetic workload comparing the three array architectures in terms of runtime performance and memory requirements. The experiments show that the flat array with change tracking is the best architecture in terms of runtime performance, for versioning frequencies of 10-5ops-1or higher matching the second best architecture or beating it by over 8 times, whereas the log-structured array is preferable for low memory usage, since it saves up to 88% of memory compared with a flat array. Hajime Fujita 0002, Kamil Iskra, Pavan Balaji, Andrew A. Chien |
CLUSTER | 3 |
| 2015 | Exploring the Suitability of Remote GPGPU Virtualization for the OpenACC Programming Model Using rCUDAabstractOpenACC is an application programming interface (API) that aims to unleash the power of heterogeneous systems composed of CPUs and accelerators such as graphic processing units (GPUs) or Intel Xeon Phi coprocessors. This directive-based programming model is intended to enable developers to accelerate their application's execution with much less effort. Coprocessors offer significant computing power but in many cases these devices remain largely underused because not all parts of applications match the accelerator architecture. Remote accelerator virtualization frameworks introduce a means to address this problem. In particular, the remote CUDA virtualization middleware rCUDA provides transparent remote access to any GPU installed in a cluster. Combining these two technologies, OpenACC and rCUDA, in a single scenario is naturally appealing. In this work we explore how the different OpenACC directives behave on top of a remote GPGPU virtualization technology in two different hardware configurations. Our experimental evaluation reveals favorable performance results when the two technologies are combined, showing low overhead and similar scaling factors when executing OpenACC-enabled directives. Adrián Castelló 0001, Antonio J. Peña, Rafael Mayo 0002, Pavan Balaji, Enrique S. Quintana-Ortí |
CLUSTER | 4 |
| 2015 | Versioning Architectures for Local and Global MemoryabstractFuture supercomputer systems will face serious reliability challenges. Among failure scenarios, latent errors are some of the most serious and concerning. Preserving multiple versions of critical data is a promising approach to deal with such errors. We are developing the Global View Resilience (GVR) library, with multi-version global arrays as one of the key features. This paper presents three array versioning architectures: flat array, flat array with change tracking, and log-structured array. We use a synthetic workload that mimics the memory access patterns of radix sort, N-body simulation, and matrix multiplication, comparing the three array architectures in terms of runtime performance, memory requirements, and version restoration costs. The experiments show that the flat array with change tracking is the best architecture in terms of runtime performance, for versioning frequencies of 10-5ops-1or higher matching the second best architecture or beating it by up to 23 times, whereas the log-structured array is preferable for low memory usage, since it saves up to 98% of memory compared with a flat array. Hajime Fujita 0002, Kamil Iskra, Pavan Balaji, Andrew A. Chien |
ICPADS | 3 |
| 2015 | Casper: An Asynchronous Progress Model for MPI RMA on Many-Core ArchitecturesabstractIn this paper we present "Casper," a process-based asynchronous progress solution for MPI one-sided communication on multi- and many-core architectures. Casper uses transparent MPI call redirection through PMPI and MPI-3 shared-memory windows to map memory from multiple user processes into the address space of one or more ghost processes, thus allowing for asynchronous progress where needed while allowing native hardware-based communication where available. Unlike traditional thread- and interrupt-based asynchronous progress models, Casper provides the capability to dedicate an arbitrary number of ghost processes for asynchronous progress, thus balancing application requirements with the capabilities of the underlying MPI implementation. We present a detailed design of the proposed architecture including several techniques for maintaining correctness per the MPI-3 standard as well as performance optimizations where possible. We also compare Casper with traditional thread- and interrupt-based asynchronous progress models and demonstrate its performance improvements with a variety of micro benchmarks and a production chemistry application. Min Si, Antonio J. Peña, Jeff R. Hammond, Pavan Balaji, Masamichi Takagi, Yutaka Ishikawa |
IPDPS | 4 |
| 2015 | MPI+Threads: runtime contention and remediesabstractHybrid MPI+Threads programming has emerged as an alternative model to the “MPI everywhere” model to better handle the increasing core density in cluster nodes. While the MPI standard allows multithreaded concurrent communication, such flexibility comes with the cost of maintaining thread safety within the MPI implementation, typically implemented using critical sections. In contrast to previous works that studied the importance of critical-section granularity in MPI implementations, in this paper we investigate the implication of critical-section arbitration on communication performance. We first analyze the MPI runtime when multithreaded concurrent communication takes place on hierarchical memory systems. Our results indicate that the mutex-based approach that most MPI implementations use today can incur performance penalties due to unfair arbitration. We then present methods to mitigate these penalties with a first-come, first-served arbitration and a priority locking scheme that favors threads doing useful work. Through evaluations using several benchmarks and applications, we demonstrate up to 5-fold improvement in performance. Abdelhalim Amer, Huiwei Lu, Yanjie Wei, Pavan Balaji, Satoshi Matsuoka |
PPoPP | 4 |
| 2015 | Fault tolerant MapReduce-MPI for HPC clustersabstractBuilding MapReduce applications using the Message-Passing Interface (MPI) enables us to exploit the performance of large HPC clusters for big data analytics. However, due to the lacking of native fault tolerance support in MPI and the incompatibility between the MapReduce fault tolerance model and HPC schedulers, it is very hard to provide a fault tolerant MapReduce runtime for HPC clusters. We propose and develop FT-MRMPI, the first fault tolerant MapReduce framework on MPI for HPC clusters. We discover a unique way to perform failure detection and recovery by exploiting the current MPI semantics and the new proposal of user-level failure mitigation. We design and develop the checkpoint/restart model for fault tolerant MapReduce in MPI. We further tailor the detect/resume model to conserve work for more efficient fault tolerance. The experimental results on a 256-node HPC cluster show that FT-MRMPI effectively masks failures and reduces the job completion time by 39%. Yanfei Guo, Wesley Bland, Pavan Balaji, Xiaobo Zhou 0002 |
SC | 3 |
| 2015 | VOCL-FT: introducing techniques for efficient soft error coprocessor recoveryabstractPopular accelerator programming models rely on offloading computation operations and their corresponding data transfers to the coprocessors, leveraging synchronization points where needed. In this paper we identify and explore how such a programming model enables optimization opportunities not utilized in traditional checkpoint/restart systems, and we analyze them as the building blocks for an efficient fault-tolerant system for accelerators. Although we leverage our techniques to protect from detected but uncorrected ECC errors in the device memory in OpenCL-accelerated applications, coprocessor reliability solutions based on different error detectors and similar API semantics can directly adopt the techniques we propose. Adding error detection and protection involves a tradeoff between runtime overhead and recovery time. Although optimal configurations depend on the particular application, the length of the run, the error rate, and the temporary storage speed, our test cases reveal a good balance with significantly reduced runtime overheads. Antonio J. Peña, Wesley Bland, Pavan Balaji |
SC | 3 |
| 2015 | Improving concurrency and asynchrony in multithreaded MPI applications using software offloadingabstractWe present a new approach for multithreaded communication and asynchronous progress in MPI applications, wherein we offload communication processing to a dedicated thread. The central premise is that given the rapidly increasing core counts on modern systems, the improvements in MPI performance arising from dedicating a thread to drive communication outweigh the small loss of resources for application computation, particularly when overlap of communication and computation can be exploited. Our approach allows application threads to make MPI calls concurrently, enqueuing these as communication tasks to be processed by a dedicated communication thread. This not only guarantees progress for such communication operations, but also reduces load imbalance. Our implementation additionally significantly reduces the overhead of mutual exclusion seen in existing implementations for applications using MPI_THREAD_MULTIPLE. Our technique requires no modification to the application, and we demonstrate significant performance improvement (up to 2X) for QCD, 1-D FFT and deep learning CNN applications. Karthikeyan Vaidyanathan, Dhiraj D. Kalamkar, Kiran Pamnany, Jeff R. Hammond, Pavan Balaji, Dipankar Das 0002, Jongsoo Park, Bálint Joó |
SC | 5 |
| 2015 | Introduction Special Section of ICCCN 2014 Conference
Pavan Balaji, Lisong Xu, Changjun Jiang 0002, Xiaobo Zhou 0002 |
Comput. Commun. | 1 |
| 2015 | Scalable connectionless RDMA over unreliable datagrams
Ryan E. Grant, Mohammad J. Rashti, Pavan Balaji, Ahmad Afsahi |
Parallel Comput. | 3 |
| 2014 | Toward the efficient use of multiple explicitly managed memory subsystemsabstractThe increasing number of memory technologies offering different features such as optimized access patterns or capacity/speed ratios lead us to advocate for future HPC compute nodes equipped with heterogeneous memory subsystems. The aim is to alleviate further the ever-increasing gap between computation and memory access speeds, by taking advantage of the benefits these memory technologies provide. Compute nodes equipped with memory technologies such as scratchpad memory, on-chip 3D-stacked memory, or NVRAM-based memory are already a reality. Careful use of the different memory subsystems is mandatory in order to exploit the potential of such super-computers. While most multiple-memory models concentrate on extending the depth of the memory hierarchy by incorporating more levels of hardware-managed memories, we advocate for compute nodes equipped with heterogeneous software-managed memory subsystems. Although the exact approach to efficiently exploit them is still uncertain, a software ecosystem clearly is required in order to assist in an efficient data distribution. We address this problem at the memory object granularity. In this paper we use an object-differentiated profiling tool we have developed on top of the Valgrind instrumentation framework, in order to assess the most suitable memory subsystem for the different memory objects of two miniapplications from the Mantevo codesign project. Our results considering two different memory configurations as use cases reveal the potential benefits of carefully placing the different memory objects of an application among the different memory subsystems. Antonio J. Peña, Pavan Balaji |
CLUSTER | 2 |
| 2014 | WorkQ: A many-core producer/consumer execution model applied to PGAS computationsabstractPartitioned global address space (PGAS) applications, such as the Tensor Contraction Engine (TCE) in NWChem, often apply a one-process-per-core mapping in which each process iterates through the following work-processing cycle: (1) determine a work-item dynamically, (2) get data via one-sided operations on remote blocks, (3) perform computation on the data locally, (4) put (or accumulate) resultant data into an appropriate remote location, and (5) repeat the cycle. However, this simple flow of execution does not effectively hide communication latency costs despite the opportunities for making asynchronous progress. Utilizing nonblocking communication calls is not sufficient unless care is taken to efficiently manage a responsive queue of outstanding communication requests. This paper presents a new runtime model and its library implementation for managing tunable “work queues” in PGAS applications. Our runtime execution model, called WorkQ, assigns some number of on-node “producer” processes to primarily do communication (steps 1, 2, 4, and 5) and the other “consumer” processes to do computation (step 3); but processes can switch roles dynamically for the sake of performance. Load balance, synchronization, and overlap of communication and computation are facilitated by a tunable nodewise FIFO message queue protocol. Our WorkQ library implementation enables an MPI+X hybrid programming model where the X comprises SysV message queues and the user's choice of SysV, POSIX, and MPI shared memory. We develop a simplified software mini-application that mimics the performance behavior of the TCE at arbitrary scale, and we show that the WorkQ engine outperforms the original model by about a factor of 2. We also show performance improvement in the TCE coupled cluster module of NWChem. David Ozog, Allen D. Malony, Jeff R. Hammond, Pavan Balaji |
ICPADS | 4 |
| 2014 | MT-MPI: multithreaded MPI for many-core environmentsabstractMany-core architectures, such as the Intel Xeon Phi, provide dozens of cores and hundreds of hardware threads. To utilize such architectures, application programmers are increasingly looking at hybrid programming models, where multiple threads interact with the MPI library (frequently called "MPI+X" models). A common mode of operation for such applications uses multiple threads to parallelize the computation, while one of the threads also issues MPI operations (i.e., MPI FUNNELED or SERIALIZED thread-safety mode). In MPI+OpenMP applications, this is achieved, for example, by placing MPI calls in OpenMP critical sections or outside the OpenMP parallel regions. However, such a model often means that the OpenMP threads are active only during the parallel computation phase and idle during the MPI calls, resulting in wasted computational resources. In this paper, we present MT-MPI, an internally multithreaded MPI implementation that transparently coordinates with the threading runtime system to share idle threads with the application. It is designed in the context of OpenMP and requires modifications to both the MPI implementation and the OpenMP runtime in order to share appropriate information between them. We demonstrate the benefit of such internal parallelism for various aspects of MPI processing, including derived datatype communication, shared-memory communication, and network I/O operations. Min Si, Antonio J. Peña, Pavan Balaji, Masamichi Takagi, Yutaka Ishikawa |
ICS | 3 |
| 2014 | Portable, MPI-interoperable coarray fortranabstractThe past decade has seen the advent of a number of parallel programming models such as Coarray Fortran (CAF), Unified Parallel C, X10, and Chapel. Despite the productivity gains promised by these models, most parallel scientific applications still rely on MPI as their data movement model. One reason for this trend is that it is hard for users to incrementally adopt these new programming models in existing MPI applications. Because each model use its own runtime system, they duplicate resources and are potentially error-prone. Such independent runtime systems were deemed necessary because MPI was considered insufficient in the past to play this role for these languages. Wesley Bland, John M. Mellor-Crummey, Pavan Balaji |
PPoPP | 4 |
| 2014 | MC-Checker: Detecting Memory Consistency Errors in MPI One-Sided ApplicationsabstractOne-sided communication decouples data movement and synchronization by providing support for asynchronous reads and updates of distributed shared data. While such interfaces can be extremely efficient, they also impose challenges in properly performing asynchronous accesses to shared data. This paper presents MC-Checker, a new tool that detects memory consistency errors in MPI one-sided applications. MCChecker first performs online instrumentation and captures relevant dynamic events, such as one-sided communications and load/store operations. MC-Checker then performs analysis to detect memory consistency errors. When found, errors are reported along with useful diagnostic information. Experiments indicate that MC-Checker is effective at detecting and diagnosing memory consistency bugs in MPI one-sided applications, with low overhead, ranging from 24.6% to 71.1%, with an average of 45.2%. Zhezhe Chen, James Dinan, Pavan Balaji, Hua Zhong 0001, Jun Wei 0001, Tao Huang 0001 |
SC | 4 |
| 2014 | Nonblocking Epochs in MPI One-Sided CommunicationabstractThe synchronization model of the MPI one-sided communication paradigm can lead to serialization and latency propagation. For instance, a process can propagate non-RMA communication-related latencies to remote peers waiting in their respective epoch-closing routines in matching epochs. In this work, we discuss six latency issues that were documented for MPI-2.0 and show how they evolved in MPI-3.0. Then, we propose entirely nonblocking RMA synchronizations that allow processes to avoid waiting even in epoch-closing routines. The proposal provides contention avoidance in communication patterns that require back to back RMA epochs. It also fixes the latency propagation issues. Moreover, it allows the MPI progress engine to orchestrate aggressive schedulings to cut down the overall completion time of sets of epochs without introducing memory consistency hazards. Our test results show noticeable performance improvements for a lower-upper matrix decomposition as well as an application pattern that performs massive atomic updates. Judicael A. Zounmevo, Pavan Balaji, William Gropp, Ahmad Afsahi |
SC | 3 |
| 2014 | SWAP-Assembler: scalable and efficient genome assembly towards thousands of coresabstractBACKGROUND: There is a widening gap between the throughput of massive parallel sequencing machines and the ability to analyze these sequencing data. Traditional assembly methods requiring long execution time and large amount of memory on a single workstation limit their use on these massive data. RESULTS: This paper presents a highly scalable assembler named as SWAP-Assembler for processing massive sequencing data using thousands of cores, where SWAP is an acronym for Small World Asynchronous Parallel model. In the paper, a mathematical description of multi-step bi-directed graph (MSG) is provided to resolve the computational interdependence on merging edges, and a highly scalable computational framework for SWAP is developed to automatically preform the parallel computation of all operations. Graph cleaning and contig extension are also included for generating contigs with high quality. Experimental results show that SWAP-Assembler scales up to 2048 cores on Yanhuang dataset using only 26 minutes, which is better than several other parallel assemblers, such as ABySS, Ray, and PASHA. Results also show that SWAP-Assembler can generate high quality contigs with good N50 size and low error rate, especially it generated the longest N50 contig sizes for Fish and Yanhuang datasets. CONCLUSIONS: In this paper, we presented a highly scalable and efficient genome assembly software, SWAP-Assembler. Compared with several other assemblers, it showed very good performance in terms of scalability and contig quality. This software is available at: https://sourceforge.net/projects/swapassembler. Jintao Meng 0001, Bingqiang Wang, Yanjie Wei, Shengzhong Feng, Pavan Balaji |
BMC Bioinform. | 5 |
| 2014 | Special issue on programming models and applications for multicores and manycores - Guest Editors' Introduction
Pavan Balaji, Zhiyi Huang 0001 |
Parallel Comput. | 1 |
| 2014 | Processing MPI Derived Datatypes on Noncontiguous GPU-Resident DataabstractDriven by the goals of efficient and generic communication of noncontiguous data layouts in GPU memory, for which solutions do not currently exist, we present a parallel, noncontiguous data-processing methodology through the MPI datatypes specification. Our processing algorithm utilizes a kernel on the GPU to pack arbitrary noncontiguous GPU data by enriching the datatypes encoding to expose a fine-grained, data-point level of parallelism. Additionally, the typically tree-based datatype encoding is preprocessed to enable efficient, cached access across GPU threads. Using CUDA, we show that the computational method outperforms DMA-based alternatives for several common data layouts as well as more complex data layouts for which reasonable DMA-based processing does not exist. Our method incurs low overhead for data layouts that closely match best-case DMA usage or that can be processed by layout-specific implementations. We additionally investigate usage scenarios for data packing that incur resource contention, identifying potential pitfalls for various packing strategies. We also demonstrate the efficacy of kernel-based packing in various communication scenarios, showing multifold improvement in point-to-point communication and evaluating packing within the context of the SHOC stencil benchmark and HACC mesh analysis. John Jenkins, James Dinan, Pavan Balaji, Tom Peterka, Nagiza F. Samatova, Rajeev Thakur |
IEEE Trans. Parallel Distributed Syst. | 3 |
| 2013 | Optimizing Burrows-Wheeler Transform-Based Sequence Alignment on Multicore ArchitecturesabstractComputational biology sequence alignment tools using the Burrows-Wheeler Transform (BWT) are widely used in next-generation sequencing (NGS) analysis. However, despite extensive optimization efforts, the performance of these tools still cannot keep up with the explosive growth of sequencing data. Through an in-depth performance analysis of BWA, a popular BWT-based aligner on multicore architectures, we demonstrate that such tools are limited by memory bandwidth due to their irregular memory access patterns. We then propose a locality-aware implementation of BWA that aims at optimizing its performance by better exploiting the caching mechanisms of modern multicore processors. Experimental results show that our improved BWA implementation can reduce last-level cache (LLC) misses by 30% and translation look aside buffer (TLB) misses by 20%, resulting in up to 2.6-fold speedup over the original BWA implementation. Jing Zhang 0039, Heshan Lin, Pavan Balaji, Wu-chun Feng |
CCGRID | 3 |
| 2013 | Toward Asynchronous and MPI-Interoperable Active MessagesabstractMany new large-scale applications have emerged recently and become important in areas such as bioinformatics and social networks. These applications are often data-intensive and involve irregular communication patterns and complex operations on remote processes. Active messages have proven effective for parallelizing such nontraditional applications. However, most current active messages frameworks are low-level and system specific, do not efficiently support asynchronous progress, and are not interoperable with two-sided and collective communications. In this paper, we present the design and implementation of an active messages framework inside MPI to provide portability and programmability, and we explore challenges when asynchronously handling active messages and other messages from the network as well as from shared memory. We test our implementation with a set of comprehensive benchmarks. Evaluation results show that our framework has the advantages of overlapping and interoperability, while introducing only a modest overhead. Darius Buntinas, Judicael A. Zounmevo, James Dinan, David Goodell, Pavan Balaji, Rajeev Thakur, Ahmad Afsahi, William Gropp |
CCGRID | 6 |
| 2013 | Optimization Strategies for MPI-Interoperable Active MessagesabstractData-intensive applications, such as those in bioinformatics and social network analysis, differ from traditional scientific applications in that they often involve data-driven and irregular computation/communication patterns, making them ill-suited for traditional data movement approaches. Active Messages (AM) is an alternative programming model that allows dynamically moving computation closer to data, rather than moving the data to the local process. In our previous work, we proposed an MPI-interoperable AM framework that allows existing MPI applications to incrementally take advantage of AM capabilities. While that work presented a baseline implementation of how AMs semantically interact with the rest of the MPI infrastructure, it had several performance shortcomings. In this paper, we analyze these performance shortcomings and propose three optimization strategies: one implicitly derived by the MPI implementation and two explicitly hinted to by the application user. In addition to the detailed description of these optimization strategies, the paper presents a thorough performance evaluation on a 4096-core cluster that demonstrates considerable performance advantages from these strategies. Pavan Balaji, William Gropp, Rajeev Thakur |
DASC | 2 |
| 2013 | Topic 15: GPU and Accelerator Computing - (Introduction)
Naoya Maruyama, Leif Kobbelt, Pavan Balaji, Nikola Puzovic, Samuel Thibault |
Euro-Par | 3 |
| 2013 | On the efficacy of GPU-integrated MPI for scientific applications
Ashwin M. Aji, Lokendra S. Panwar, Milind Chabbi, Karthik Murthy, Pavan Balaji, Keith R. Bisset, James Dinan, Wu-chun Feng, John M. Mellor-Crummey, Xiaosong Ma, Rajeev Thakur |
HPDC | 6 |
| 2013 | pVOCL: Power-Aware Dynamic Placement and Migration in Virtualized GPU EnvironmentsabstractPower-hungry Graphics processing unit (GPU) accelerators are ubiquitous in high performance computing data centers today. GPU virtualization frameworks introduce new opportunities for effective management of GPU resources by decoupling them from application execution. However, power management of GPU-enabled server clusters faces significant challenges. The underlying system infrastructure shows complex power consumption characteristics depending on the placement of GPU workloads across various compute nodes, power-phases and cabinets in a datacenter. GPU resources need to be scheduled dynamically in the face of time-varying resource demand and peak power constraints. We propose and develop a power-aware virtual OpenCL (pVOCL) framework that controls the peak power consumption and improves the energy efficiency of the underlying server system through dynamic consolidation and power-phase topology aware placement of GPU workloads. Experimental results show that pVOCL achieves significant energy savings compared to existing power management techniques for GPU-enabled server clusters, while incurring negligible impact on performance. It drives the system towards energy-efficient configurations by taking an optimal sequence of adaptation actions in a virtualized GPU environment and meanwhile keeps the power consumption below the peak power budget. Palden Lama, Yan Li 0005, Ashwin M. Aji, Pavan Balaji, James Dinan, Shucai Xiao, Yunquan Zhang, Wu-chun Feng, Rajeev Thakur, Xiaobo Zhou 0002 |
ICDCS | 4 |
| 2013 | Online Performance Projection for Clusters with Heterogeneous GPUsabstractWe present a fully automated approach to project the relative performance of an OpenCL program over different GPUs. Performance projections can be made within a small amount of time, and the projection overhead stays relatively constant with the input data size. As a result, the technique can help runtime tools make dynamic decisions about which GPU would run faster for a given kernel. Usage cases of this technique include scheduling or migrating GPU workloads over a heterogeneous cluster with different types of GPUs. Lokendra S. Panwar, Ashwin M. Aji, Jiayuan Meng, Pavan Balaji, Wu-chun Feng |
ICPADS | 4 |
| 2013 | MPI-Interoperable Generalized Active MessagesabstractData-intensive applications have become increasingly important in recent years, yet traditional data movement approaches for scientific computation are not well suited for such applications. The Active Message (AM) model is an alternative communication paradigm that is better suited for such applications by allowing computation to be dynamically moved closer to data. Given the wide usage of MPI in scientific computing, enabling an MPI-interoperable AM paradigm would allow traditional applications to incrementally start utilizing AMs in portions of their applications, thus eliminating the programming effort of rewriting entire applications. In our previous work, we extended the MPI ACCUMULATE and MPI GET ACCUMULATE operations in the MPI standard to support AMs. However, the semantics of accumulate-style AMs are fundamentally restricted by the semantics of MPI ACCUMULATE and MPI GET ACCUMULATE, which were not designed to support the AM model. In this paper, we present a new generalized framework for MPI-interoperable AMs that can alleviate those restrictions, thus providing a richer semantics to accommodate a wide variety of application computational patterns. Together with a new API, we present a detailed description of the correctness semantics of this functionality and a reference implementation that demonstrates how various API choices affect the flexibility provided to the MPI implementation and consequently its performance. Pavan Balaji, William Gropp, Rajeev Thakur |
ICPADS | 2 |
| 2013 | Enhancing Performance Portability of MPI Applications through Annotation-Based TransformationsabstractMPI is the de facto standard for portable parallel programming on high-end systems. However, while the MPI standard provides functional portability, it does not provide sufficient performance portability across platforms. We present a framework that enables users to provide hints about communication patterns used within MPI applications. These annotations are then used by an automated program transformation system to leverage different MPI operations that better match each system's capabilities. Our framework currently supports three automated transformations: coalescing of operations in MPI one-sided communications, transformation of blocking communications to nonblocking, which enables communication-computation overlap, and selection of the appropriate communication operators based on the cache-coherence support of the underlying platform. We use our annotation-based approach to optimize several benchmark kernels, and we demonstrate that the framework is effective at automatically improving performance portability for MPI applications. Md. Ziaul Haque, Qing Yi, James Dinan, Pavan Balaji |
ICPP | 4 |
| 2013 | Inspector-Executor Load Balancing Algorithms for Block-Sparse Tensor ContractionsabstractDeveloping effective yet scalable load-balancing methods for irregular computations is critical to the successful application of simulations in a variety of disciplines at petascale and beyond. This paper explores a set of static and dynamic scheduling algorithms for block-sparse tensor contractions within the NWChem computational chemistry code for different degrees of sparsity (and therefore load imbalance). In this particular application, a relatively large amount of task information can be obtained at minimal cost, which enables the use of static partitioning techniques that take the entire task list as input. However, fully static partitioning is incapable of dealing with dynamic variation of task costs, such as from transient network contention or operating system noise, so we also consider hybrid schemes that utilize dynamic scheduling within subgroups. These two schemes, which have not been previously implemented in NWChem or its proxies (i.e. quantum chemistry mini-apps) are compared to the original centralized dynamic load-balancing algorithm as well as improved centralized scheme. In all cases, we separate the scheduling of tasks from the execution of tasks into an inspector phase and an executor phase. The impact of these methods upon the application is substantial on a large InfiniBand cluster: execution time is reduced by as much as 50% at scale. The technique is applicable to any scientific application requiring load balance where performance models or estimations of kernel execution times are available. David Ozog, Jeff R. Hammond, James Dinan, Pavan Balaji, Sameer Shende, Allen D. Malony |
ICPP | 4 |
| 2013 | Inspector/executor load balancing algorithms for block-sparse tensor contractionsabstractDeveloping effective yet scalable load-balancing methods for irregular computations is critical to the successful application of simulations in a variety of disciplines at petascale and beyond. This paper explores a set of static and dynamic scheduling algorithms for block-sparse tensor contractions within the NWChem computational chemistry code for different degrees of sparsity (and therefore load imbalance). In this particular application, a relatively large amount of task information can be obtained at minimal cost, which enables the use of static partitioning techniques that take the entire task list as input. However, fully static partitioning is incapable of dealing with dynamic variation of task costs, such as from transient network contention or operating system noise, so we also consider hybrid schemes that utilize dynamic scheduling within subgroups. These two schemes, which have not been previously implemented in NWChem or its proxies (i.e. quantum chemistry mini-apps) are compared to the original centralized dynamic load-balancing algorithm as well as improved centralized scheme. In all cases, we separate the scheduling of tasks from the execution of tasks into an inspector phase and an executor phase. The impact of these methods upon the application is substantial on a large InfiniBand cluster: execution time is reduced by as much as 50% at scale. The technique is applicable to any scientific application requiring load balance where performance models or estimations of kernel execution times are available. David Ozog, Sameer Shende, Allen D. Malony, Jeff R. Hammond, James Dinan, Pavan Balaji |
ICS | 6 |
| 2013 | Enabling MPI interoperability through flexible communication endpointsabstractThe current MPI model defines a one-to-one relationship between MPI processes and MPI ranks. This model captures many use cases effectively, such as one MPI process per core and one MPI process per node. However, this semantic has limited interoperability between MPI and other programming models that use threads within a node. In this paper, we describe an extension to MPI that introduces communication endpoints as a means to relax the one-to-one relationship between processes and threads. Endpoints enable a greater degree interoperability between MPI and other programming models, and we illustrate their potential for additional performance and computation management benefits through the decoupling of ranks from processes. James Dinan, Pavan Balaji, David Goodell, Douglas Miller, Marc Snir, Rajeev Thakur |
EuroMPI | 2 |
| 2013 | Analysis of topology-dependent MPI performance on Gemini networksabstractCurrent HPC systems utilize a variety of interconnection networks, with varying features and communication characteristics. MPI normalizes these interconnects with a common interface used by most HPC applications. However, network properties can have a significant impact on application performance. We explore the impact of the interconnect on application performance on the Blue Waters supercomputer. Blue Waters uses a three-dimensional, Cray Gemini torus network, which provides twice the Y-dimension bandwidth in the X and Z dimensions. Through several benchmarks, including a halo-exchange example, we demonstrate that application-level mapping to the network topology yields significant performance improvements. Antonio J. Peña, Ralf G. Correa Carvalho, James Dinan, Pavan Balaji, Rajeev Thakur, William Gropp |
EuroMPI | 4 |
| 2013 | Guest editors' introduction: Special issue on Cluster, Grid, and Cloud Computing
Pavan Balaji, Rajkumar Buyya |
Future Gener. Comput. Syst. | 1 |
| 2013 | Special issue on programming models, systems software, and tools for High-End Computing
Yong Chen 0001, Pavan Balaji, Abhinav Vishnu |
Parallel Comput. | 2 |
| 2013 | A survey on resource allocation in high performance distributed computing systems
Hameed Hussain, Saif Ur Rehman Malik, Abdul Hameed, Samee Ullah Khan, Gage Bickler, Nasro Min-Allah, Muhammad Bilal Qureshi, Yongji Wang 0002, Nasir Ghani, Joanna Kolodziej, Albert Y. Zomaya, Cheng-Zhong Xu 0001, Pavan Balaji, Abhinav Vishnu, Frédéric Pinel, Johnatan E. Pecero, Dzmitry Kliazovich, Pascal Bouvry, Hongxiang Li 0001, Lizhe Wang 0001, Dan Chen 0001, Ammar Rayes |
Parallel Comput. | 14 |
| 2013 | Guest Editors' introduction
Abhinav Vishnu, Pavan Balaji, Yong Chen 0001 |
J. Supercomput. | 2 |
| 2013 | Designing energy efficient communication runtime systems: a view from PGAS models
Abhinav Vishnu, Shuaiwen Song, Andrés Márquez 0001, Kevin J. Barker, Darren J. Kerbyson, Kirk W. Cameron, Pavan Balaji |
J. Supercomput. | 7 |
| 2012 | Transparent Accelerator Migration in a Virtualized GPU EnvironmentabstractThis paper presents a framework to support transparent, live migration of virtual GPU accelerators in a virtualized execution environment. Migration is a critical capability in such environments because it provides support for fault tolerance, on-demand system maintenance, resource management, and load balancing in the mapping of virtual to physical GPUs. Techniques to increase responsiveness and reduce migration overhead are explored. The system is evaluated by using four application kernels and is demonstrated to provide low migration overheads. Through transparent load balancing, our system provides a speedup of 1.7 to 1.9 for three of the four application kernels. Shucai Xiao, Pavan Balaji, James Dinan, Rajeev Thakur, Susan Coghlan, Heshan Lin, Gaojin Wen, Jue Hong, Wu-chun Feng |
CCGRID | 2 |
| 2012 | Enabling Fast, Noncontiguous GPU Data Movement in Hybrid MPI+GPU EnvironmentsabstractLack of efficient and transparent interaction with GPU data in hybrid MPI+GPU environments challenges GPU acceleration of large-scale scientific computations. A particular challenge is the transfer of noncontiguous data to and from GPU memory. MPI implementations currently do not provide an efficient means of utilizing data types for noncontiguous communication of data in GPU memory. To address this gap, we present an MPI data type-processing system capable of efficiently processing arbitrary data types directly on the GPU. We present a means for converting conventional data type representations into a GPU-amenable format. Fine-grained, element-level parallelism is then utilized by a GPU kernel to perform in-device packing and unpacking of noncontiguous elements. We demonstrate a several-fold performance improvement for noncontiguous column vectors, 3D array slices, and 4D array sub volumes over CUDA-based alternatives. Compared with optimized, layout-specific implementations, our approach incurs low overhead, while enabling the packing of data types that do not have a direct CUDA equivalent. These improvements are demonstrated to translate to significant improvements in end-to-end, GPU-to-GPU communication time. In addition, we identify and evaluate communication patterns that may cause resource contention with packing operations, providing a baseline for adaptively selecting data-processing strategies. John Jenkins, James Dinan, Pavan Balaji, Nagiza F. Samatova, Rajeev Thakur |
CLUSTER | 3 |
| 2012 | Supporting the Global Arrays PGAS Model Using MPI One-Sided CommunicationabstractThe industry-standard Message Passing Interface (MPI) provides one-sided communication functionality and is available on virtually every parallel computing system. However, it is believed that MPI's one-sided model is not rich enough to support higher-level global address space parallel programming models. We present the first successful application of MPI one-sided communication as a runtime system for a PGAS model, Global Arrays (GA). This work has an immediate impact on users of GA applications, such as NW Chem, who often must wait several months to a year or more before GA becomes available on a new architecture. We explore challenges present in the application of MPI-2 to PGAS models and motivate new features in the upcoming MPI-3 standard. The performance of our system is evaluated on several popular high-performance computing architectures through communication benchmarking and application benchmarking using the NW Chem computational chemistry suite. James Dinan, Pavan Balaji, Jeff R. Hammond, Sriram Krishnamoorthy, Vinod Tipparaju |
IPDPS | 2 |
| 2012 | Efficient Multithreaded Context ID Allocation in MPI
James Dinan, David Goodell, William Gropp, Rajeev Thakur, Pavan Balaji |
EuroMPI | 5 |
| 2012 | Leveraging MPI's One-Sided Communication Interface for Shared-Memory Programming
Torsten Hoefler, James Dinan, Darius Buntinas, Pavan Balaji, Brian W. Barrett, Ron Brightwell, William Gropp, Vivek Kale, Rajeev Thakur |
EuroMPI | 4 |
| 2011 | Building algorithmically nonstop fault tolerant MPI programsabstractWith the growing scale of high-performance computing (HPC) systems, today and more so tomorrow, faults are a norm rather than an exception. HPC applications typically tolerate fail-stop failures under the stop-and-wait scheme, where even if only one processor fails, the whole system has to stop and wait for the recovery of the corrupted data. It is now a more-or-less accepted fact that the stop-and-wait scheme will not scale to the next generation of HPC systems. Inspired by the previous stop-and-wait algorithm-based fault tolerance (ABFT) recovery technique, we propose in this paper a nonstop fault tolerance scheme at the application level and describe its implementation. When failure occurs during the execution of applications, we do not stop to wait for the recovery of the corrupted node; instead, we replace it with the corresponding redundant node and continue the execution. At the end of execution, the correct solution can be recovered algorithmically at a very low cost. In order to implement the scheme, some new fault-tolerant features of the Message Passing Interface (MPI) have been investigated and utilized in the MPICH implementation of MPI. We also describe a case study using High Performance Linpack (HPL) with these new features and evaluate the performance of both our new scheme and ABFT recovery. Experimental results show the advantage of our new scheme over ABFT recovery even in a small scale. Erlin Yao, Mingyu Chen 0001, Guangming Tan, Pavan Balaji, Darius Buntinas |
HiPC | 5 |
| 2011 | RDMA Capable iWARP over DatagramsabstractiWARP is a state of the art high-speed connection-based RDMA networking technology for Ethernet networks to provide InfiniBand-like zero-copy and one-sided communication capabilities over Ethernet. Despite the benefits offered by iWARP, many data center and web-based applications, such as stock-market trading and media-streaming applications, that rely on data gram-based semantics (mostly through UDP/IP) cannot take advantage of it because the iWARP standard is only defined over reliable, connection-oriented transports. This paper presents an RDMA model that functions over reliable and unreliable data grams. The ability to use data grams significantly expands the application space serviced by iWARP and can bring the scalability advantages of a connectionless transport to iWARP. In our previous work, we had developed an iWARP data gram solution using send/receive semantics showing excellent memory scalability and performance benefits over the current TCP-based iWARP. In this paper, we demonstrate an improved iWARP design that provides true RDMA semantics over data grams. Specifically, because traditional RDMA semantics do not map well to unreliable communication, we propose RDMA Write-Record, the first and the only method capable of supporting RDMA Write over both unreliable and reliable data grams. We demonstrate through a proof-of-concept software implementation that data gram-iWARP is feasible for real-world applications. Our proposed RDMA Write-Record method has been designed with data loss in mind and can provide superior performance under conditions of packet loss. It is shown through micro-benchmarks that by using RDMA capable data gram-iWARP a maximum of 256% increase in large message bandwidth and a maximum of 24.4\% improvement in small message latency can be achieved over traditional iWARP. For application results we focus on streaming applications, showing a 24% improvement in memory usage and up to a 74% improvement in performance, although the proposed approach is also applicable to the HPC domain. Ryan E. Grant, Mohammad J. Rashti, Ahmad Afsahi, Pavan Balaji |
IPDPS | 4 |
| 2011 | Noncollective Communicator Creation in MPI
James Dinan, Sriram Krishnamoorthy, Pavan Balaji, Jeff R. Hammond, Manojkumar Krishnan, Vinod Tipparaju, Abhinav Vishnu |
EuroMPI | 3 |
| 2011 | Multi-core and Network Aware MPI Topology Functions
Mohammad J. Rashti, Jonathan Green, Pavan Balaji, Ahmad Afsahi, William Gropp |
EuroMPI | 3 |
| 2010 | Minimizing MPI Resource Contention in Multithreaded Multicore EnvironmentsabstractWith the ever-increasing numbers of cores per node in high-performance computing systems, a growing number of applications are using threads to exploit shared memory within a node and MPI across nodes. This hybrid programming model needs efficient support for multithreaded MPI communication. In this paper, we describe the optimization of one aspect of a multithreaded MPI implementation: concurrent accesses from multiple threads to various MPI objects, such as communicators, datatypes, and requests. The semantics of the creation, usage, and destruction of these objects implies, but does not strictly require, the use of reference counting to prevent memory leaks and premature object destruction. We demonstrate how a naive multithreaded implementation of MPI object management via reference counting incurs a significant performance penalty. We then detail two solutions that we have implemented in MPICH2 to mitigate this problem almost entirely, including one based on a novel garbage collection scheme. In our performance experiments, this new scheme improved the multithreaded messaging rate by up to 31% over the naive reference counting method. David Goodell, Pavan Balaji, Darius Buntinas, Gábor Dózsa, William Gropp, Sameer Kumar 0001, Bronis R. de Supinski, Rajeev Thakur |
CLUSTER | 2 |
| 2010 | iWARP redefined: Scalable connectionless communication over high-speed EthernetabstractiWARP represents the leading edge of high performance Ethernet technologies. By utilizing an asynchronous communication model, iWARP brings the advantages of OS bypass and RDMA technology to Ethernet. The current specification of iWARP is only defined over connection-oriented transports such as TCP. The memory requirements of many connections along with TCP's flow and reliability controls lead to scalability and performance issues for large-scale HPC and datacenter applications. In this research, we propose guidelines to extend iWARP over datagrams to provide better scalability and performance. While the proposed extension is designed for use in both HPC and datacenters, the emphasis of this paper is on HPC applications. We present our software implementation of datagram-iWARP over UDP and MPI over datagram-iWARP. Our microbenchmark and MPI application results show performance and memory usage benefits for MPI applications, promoting the use of datagram-iWARP for large-scale HPC applications. Mohammad J. Rashti, Ryan E. Grant, Ahmad Afsahi, Pavan Balaji |
HiPC | 4 |
| 2010 | Fault-tolerant communication runtime support for data-centric programming modelsabstractThe largest supercomputers in the world today consist of hundreds of thousands of processing cores and many more other hardware components. At such scales, hardware faults are a commonplace, necessitating fault-resilient software systems. While different fault-resilient models are available, most focus on allowing the computational processes to survive faults. On the other hand, we have recently started investigating fault resilience techniques for data-centric programming models such as the partitioned global address space (PGAS) models. The primary difference in data-centric models is the decoupling of computation and data locality. That is, data placement is decoupled from the executing processes, allowing us to view process failure (a physical node hosting a process is dead) separately from data failure (a physical node hosting data is dead). In this paper, we take a first step toward data-centric fault resilience by designing and implementing a fault-resilient, one-sided communication runtime framework using Global Arrays and its communication system, ARMCI. The framework consists of a fault-resilient process manager; low-overhead and network-assisted remote-node fault detection module; non-data-moving collective communication primitives; and failure semantics and err or codes for one-sided communication runtime systems. Our performance evaluation indicates that the framework incurs little overhead compared to state-of-the-art designs and provides a fundamental framework of fault resiliency for PGAS models. Abhinav Vishnu, Huub J. J. Van Dam, Wibe de Jong, Pavan Balaji, Shuaiwen Song |
HiPC | 4 |
| 2010 | A study of hardware assisted IP over InfiniBand and its impact on enterprise data center performanceabstractHigh-performance sockets implementations such as the Sockets Direct Protocol (SDP) have traditionally showed major performance advantages compared to the TCP/IP stack over InfiniBand (IPoIB). These stacks bypass the kernel-based TCP/IP and take advantage of network hardware features, providing enhanced performance. SDP has excellent performance but limited utility as only applications relying on the TCP/IP sockets API can use it and other IP stack uses (IPSec, UDP, SCTP) or TCP layer modifications (iSCSI) cannot benefit from it. Recently, newer generations of InfiniBand adapters, such as ConnectX from Mellanox, have provided hardware support for the IP stack itself, such as Large Send Offload and Large Receive Offload. As such high performance socket networks are likely to be deployed or converged with existing Ethernet networking solutions, the performance of such technologies is important to assess. In this paper we take a first look at the performance advantages provided by these offload techniques and compare them to SDP. Our micro-benchmarks and enterprise data-center experiments show that hardware assisted IPoIB can provide competitive performance with SDP and even outperform it in some cases. Ryan E. Grant, Pavan Balaji, Ahmad Afsahi |
ISPASS | 2 |
| 2010 | PMI: A Scalable Parallel Process-Management Interface for Extreme-Scale Systems
Pavan Balaji, Darius Buntinas, David Goodell, William Gropp, Jayesh Krishna, Ewing L. Lusk, Rajeev Thakur |
EuroMPI | 1 |
| 2010 | Enabling Concurrent Multithreaded MPI Communication on Multicore Petascale Systems
Gábor Dózsa, Sameer Kumar 0001, Pavan Balaji, Darius Buntinas, David Goodell, William Gropp, Joe Ratterman, Rajeev Thakur |
EuroMPI | 3 |
| 2010 | Implementing MPI on Windows: Comparison with Common Approaches on Unix
Jayesh Krishna, Pavan Balaji, Ewing L. Lusk, Rajeev Thakur, Fabian Tiller |
EuroMPI | 2 |
| 2010 | Global-scale distributed I/O with ParaMEDICabstractAbstract Achieving high performance for distributed I/O on a wide‐area network continues to be an elusive holy grail. Despite enhancements in network hardware as well as software stacks, achieving high‐performance remains a challenge. In this paper, our worldwide team took a completely new and non‐traditional approach to distributed I/O, calledParaMEDIC: Parallel Metadata Environment for Distributed I/O and Computing, by utilizing application‐specifictransformationof data to orders of magnitude smaller metadata before performing the actual I/O. Specifically, this paper details our experiences in deploying a large‐scale system to facilitate the discovery of missing genes and constructing a genome similarity tree by encapsulating the mpiBLAST sequence‐search algorithm into ParaMEDIC. The overall project involved nine computational sites spread across the U.S. and generated more than a petabyte of data that was ‘teleported’ to a large‐scale facility in Tokyo for storage. Copyright © 2010 John Wiley & Sons, Ltd. Pavan Balaji, Wu-chun Feng, Heshan Lin, Jeremy S. Archuleta, Satoshi Matsuoka, Andrew S. Warren, João Carlos Setubal, Ewing L. Lusk, Rajeev Thakur, Ian T. Foster, Daniel S. Katz, Shantenu Jha, K. Shinpaugh, Susan Coghlan, Daniel A. Reed |
Concurr. Comput. Pract. Exp. | 1 |
| 2009 | Natively Supporting True One-Sided Communication inabstractAs high-end computing systems continue to grow in scale, the performance that applications can achieve on such large scale systems depends heavily on their ability to avoid explicitly synchronized communication with other processes in the system. Accordingly, several modern and legacy parallel programming models (such as MPI, UPC, global arrays) have provided many programming constructs that enable implicit communication using one-sided communication operations. While MPI is the most widely used communication model for scientific computing, the usage of one-sided communication is restricted; this is mainly owing to the inefficiencies in current MPI implementations that internally rely on synchronization between processes even during one-sided communication, thus losing the potential of such constructs. In our previous work, we had utilized native one-sided communication primitives offered by high-speed networks such as infiniband (IB) to allow for true one-sided communication in MPI. In this paper, we extend this work to natively take advantage of one-sided atomic operations on cache-coherent multi-core/multi-processor architectures while still utilizing the benefits of networks such as IB. Specifically, we present a sophisticated hybrid design that uses locks that migrate between IB hardware atomics and multi-core CPU atomics to take advantage of both. We demonstrate the capability of our proposed design with a wide range of experiments illustrating its benefits in performance as well as its potential to avoid explicit synchronization. Gopalakrishnan Santhanaraman, Pavan Balaji, Rajeev Thakur, William Gropp, Dhabaleswar K. Panda 0001 |
CCGRID | 2 |
| 2009 | Understanding Network Saturation Behavior on Large-Scale Blue Gene/P SystemsabstractAs researchers continue to architect massive-scale systems, it is becoming clear that these systems will utilize a significant amount of shared hardware between processing units. Systems such as the IBM Blue Gene (BG) and Cray XT have started utilizing flat (i.e., scalable) networks, which differ from switched fabrics in that they use a 3D torus or similar topology. This allows the network to grow only linearly with system scale, instead of the super linear growth needed for full fat-tree switched topologies, but at the cost of increased network sharing between processing nodes. While in many cases a full fat-tree is an over estimate of the needed bisectional bandwidth, it is not clear whether the other extreme of a flat topology is sufficient to move data around the network efficiently. In this paper, we study the network behavior of the IBM BG/P using several application communication kernels, and we monitor network congestion behavior based on detailed hardware counters. Our studies scale from small systems to 8 racks (32,768 cores) of BG/P and provide insights into the network communication characteristics of the system. Pavan Balaji, Harish Naik, Narayan Desai |
ICPADS | 1 |
| 2009 | Evaluation of ConnectX Virtual Protocol Interconnect for Data CentersabstractWith the emergence of new technologies such as Virtual Protocol Interconnect (VPI) for the modern data center, the separation between commodity networking technology and high-performance interconnects is shrinking. With VPI, a single network adapter on a data center server can easily be configured to use one port to interface with Ethernet traffic and another port to interface with high-bandwidth, low-latency InfiniBand technology. In this paper, we evaluate ConnectX VPI using microbenchmarks as well as real traces from a three-tier data center architecture. We find that with VPI each network segment in the data center can use the most optimal configuration (whether InfiniBand or Ethernet) without having to fall back to the lowest common denominator, as is currently the case. Our results show a maximum 26.7% increase in bandwidth, a 54.5% reduction in latency, and a 5% increase in real data center throughput. Ryan E. Grant, Ahmad Afsahi, Pavan Balaji |
ICPADS | 3 |
| 2009 | Improving Resource Availability by Relaxing Network Allocation Constraints on Blue Gene/PabstractHigh-end computing (HEC) systems have passed the petaflop barrier and continue to move toward the next frontier of {exascale} computing. As companies and research institutes continue to work toward architecting these enormous systems, it is becoming increasingly clear that these systems will utilize a significant amount of shared hardware between processing units, including shared caches, memory management engines, and network infrastructure. While these systems are optimized to use all of the hardware available in a dedicated manner to achieve the best performance, in practice, the shared nature of this hardware makes scheduling applications on it difficult and wasteful. For example, while the IBM Blue Gene/P system has been designed to use a torus network for efficient communication, some of the torus links (especially those connecting different racks) are shared between multiple racks. Thus, a job running on one rack, might preclude another job from running on a second rack in spite of having its compute resources completely idle. In this paper, we assess the relative performance degradation noticed by real applications when such shared network hardware is completely unutilized for some cases. Our measurements on Intrepid, one of the largest Blue Gene/P installations in the world, demonstrate less than 5% degradation for several leadership applications commonly run on the Intrepid system. Further, we demonstrate that the additional scheduling flexibility offered by not sharing such hardware can improve the overall job turnaround time by nearly 40% in some cases. Narayan Desai, Darius Buntinas, Daniel Buettner, Pavan Balaji |
ICPP | 4 |
| 2009 | GePSeA: A General-Purpose Software Acceleration Framework for Lightweight Task OffloadingabstractHardware-acceleration techniques continue to be used to speed-up the execution of scientific codes. To do so, software developers identify portions of these codes that are amenable for offloading and map them to hardware accelerators. However, offloading such tasks to specialized hardware accelerators is non-trivial. Furthermore, these accelerators can add significant cost to a computing system. Consequently, we propose a framework called GePSeA (General Purpose Software Acceleration Framework), which uses a small fraction of the computational power on multi-core architectures to ``onload'' complex application-specific tasks. Specifically, GePSeA provides a lightweight process that acts as a helper agent to the application by executing application-specific tasks asynchronously and efficiently. We then apply the GePSeA framework to a real application, namely, an open-source computational biology application, and demonstrate significant application-level benefits. Pavan Balaji, Wu-chun Feng |
ICPP | 2 |
| 2008 | Are nonblocking networks really needed for high-end-computing workloads?abstractHigh-speed interconnects are frequently used to provide scalable communication on increasingly large high-end computing systems. Often, these networks are nonblocking, where there exist independent paths between all pairs of nodes in the system allowing for simultaneous communication with zero network contention. This performance, however, comes at a heavy cost as the number of components needed (and hence cost) increases superlinearly with the number of nodes in the system. In this paper, we study the behavior of real and synthetic supercomputer workloads to understand the impact of the networkpsilas nonblocking capability on overall performance. Starting from a fully nonblocking network, we begin by assessing the worse-case performance degradation caused by removing interstage communication links, resulting in over provisioning and hence potentially blocking in the communication network.We also study the impact of several factors on this behavior, including system workloads, multicore processors, and switch crossbar sizes. Our observations show that a significant reduction in the number of interstage links can be tolerated on all of the workloads analyzed, causing less than 5% overall loss of performance. Narayan Desai, Pavan Balaji, P. Sadayappan |
CLUSTER | 2 |
| 2008 | Sockets Direct Protocol for Hybrid Network Stacks: A Case Study with iWARP over 10G Ethernet
Pavan Balaji, Sitha Bhagvat, Rajeev Thakur, Dhabaleswar K. Panda 0001 |
HiPC | 1 |
| 2008 | Communication Analysis of Parallel 3D FFT for Flat Cartesian Meshes on Large Blue Gene Systems
Pavan Balaji, William Gropp, Rajeev Thakur |
HiPC | 2 |
| 2008 | Making a Case for Proactive Flow Control in Optical Circuit-Switched Networks
Mithilesh Kumar 0002, Vineeta Chaube, Pavan Balaji, Wu-chun Feng, Hyun-Wook Jin |
HiPC | 3 |
| 2008 | Semantic-based distributed i/o with the paramedic frameworkabstractMany large-scale applications simultaneously rely on multiple resources for efficient execution. For example, such applications may require both large compute and storage resources; however, very few supercomputing centers can provide large quantities of both. Thus, data generated at the compute site oftentimes has to be moved to a remote storage site for either storage or visualization and analysis. Clearly, this is not an efficient model, especially when the two sites are distributed over a wide-area network. Pavan Balaji, Wu-chun Feng, Heshan Lin |
HPDC | 1 |
| 2008 | Impact of Network Sharing in Multi-Core ArchitecturesabstractAs commodity components continue to dominate the realm of high-end computing, two hardware trends have emerged as major contributors-high-speed networking technologies and multi-core architectures. Communication middleware such as the Message Passing Interface (MPI) uses the network technology for communicating between processes that reside on different physical nodes, while using shared memory for communicating between processes on different cores within the same node. Thus, two conflicting possibilities arise: (i) with the advent of multi-core architectures, the number of processes that reside on the same physical node and hence share the same physical network can potentially increase significantly, resulting in increased network usage, and (ii) given the increase in intra-node shared-memory communication for processes residing on the same node, the network usage can potentially decrease significantly. In this paper, we address these two conflicting possibilities and study the behavior of network usage in multi-core environments with sample scientific applications. Specifically, we analyze trends that result in increase or decrease of network usage, and we derive insights into application performance based on these. We also study the sharing of different resources in the system in multi-core environments and identify the contribution of the network in this mix. In addition, we study different process allocation strategies and analyze their impact on such network sharing. Ganesh Narayanaswamy, Pavan Balaji, Wu-chun Feng |
ICCCN | 2 |
| 2008 | Semantics-based distributed I/O for mpiBLASTabstractBLAST is a widely used software toolkit for genomic sequence search. mpiBLAST is a freely available, open-source parallelization of BLAST that uses database segmentation to allow different worker processes to search (in parallel) unique segments of the database. After searching, the workers write their output to a filesystem. While mpiBLAST has been shown to achieve high performance in clusters with fast local filesystems, its I/O processing remains a concern for scalability, especially in systems having limited I/O capabilities such as distributed filesystems spread across a wide-area network. Thus, we present ParaMEDIC---a novel environment that uses application-specific semantic information to compress I/O data and improve performance in distributed environments. Specifically, for mpiBLAST, ParaMEDIC partitions worker processes into compute and I/O workers. Compute workers, instead of directly writing the output to the filesystem, the workers process the output using semantic knowledge about the application to generate metadata and write the metadata to the filesystem. I/O workers, which physically reside closer to the actual storage, then process this metadata to re-create the actual output and write it to the filesystem. This approach allows ParaMEDIC to reduce I/O time, thus accelerating mpiBLAST by as much as 25-fold. Pavan Balaji, Wu-chun Feng, Jeremy S. Archuleta, Heshan Lin, Rajkumar Kettimuthu, Rajeev Thakur, Xiaosong Ma |
PPoPP | 1 |
| 2008 | Massively parallel genomic sequence search on the Blue Gene/P architectureabstractThis paper presents our first experiences in mapping and optimizing genomic sequence search onto the massively parallel IBM Blue Gene/P (BG/P) platform. Specifically, we performed our work on mpiBLAST, a parallel sequence-search code that has been optimized on numerous supercomputing environments. In doing so, we identify several critical performance issues. Consequently, we propose and study different approaches for mapping sequence-search and parallel I/O tasks on such massively parallel architectures.We demonstrate that our optimizations can deliver nearly linear scaling (93% efficiency) on up to 32,768 cores of BG/P. In addition, we show that such scalability enables us to complete a large-scale bioinformatics problem - sequence searching a microbial genome database against itself to support the discovery of missing genes in genomes - in only a few hours on BG/P. Previously, this problem was viewed as computationally intractable in practice. Heshan Lin, Pavan Balaji, Ruth Poole, Carlos P. Sosa, Xiaosong Ma, Wu-chun Feng |
SC | 2 |
| 2008 | Asymmetric interactions in symmetric multi-core systems: analysis, enhancements and evaluationabstractMulti-core architectures have spurred the rapid growth in high-end computing systems. While the vast majority of such multi-core processors contain symmetric hardware components, their interaction with systems software, in particular the communication stack, results in a remarkable amount of asymmetry in the effective capability of the different cores. In this paper, we analyze such interactions and propose a novel management library called SyMMer (Systems Mapping Manager) that monitors these interactions and dynamically manages the mapping of processes on processor cores to transparently improve application performance. Together with a detailed description of the SyMMer library, we also present performance evaluation comparing SyMMer to a vanilla communication library using various micro-benchmarks as well as popular applications and scientific libraries. Experimental results demonstrate more than a two-fold improvement in communication time and 10-15% improvement in overall application performance. Thomas Scogland, Pavan Balaji, Wu-chun Feng, Ganesh Narayanaswamy |
SC | 2 |
| 2007 | Designing high-end computing systems with InfiniBand and10-Gigabit Ethernet iWARPabstractInfiniBand (IB) architecture and 10-Gigabit Ethernet iWARP (10GE) are generating a lot of excitement towards building next generation High-End Computing (HEC) systems. This tutorial will provide an indepth look at this emerging trend and examine the suitability of these standards for prime-time HEC. It will start with a brief overview of IB and 10GE iWARP, and their architectural features. An overview of the emerging OpenFabrics stack which encapsulates both IB and 10GE iWARP in a unified manner will be presented. IB and 10GE iWARP hardware/software solutions and the market trends will be highlighted. Challenges in designing different kinds of systems using these standards on multi-core platforms for performance, scalability, portability and reliability will be covered. Specifically, case studies and experiences in designing HPC clusters (with MPI-1 and MPI-2 programming models), Parallel FileSystems, Networked File Systems (NFS), Storage Protocols, Multi-tier Datacenters, and Virtualization schemes will be presented together with the associated performance numbers and comparisons. Dhabaleswar K. Panda 0001, Pavan Balaji |
CLUSTER | 2 |
| 2007 | Advanced Flow-control Mechanisms for the Sockets Direct Protocol over InfiniBandabstractThe Sockets Direct Protocol (SDP) is an industry standard to allow existing TCP/IP applications to be executed on high-speed networks such as InfiniBand (IB). Like many other high-speed networks, IB requires the receiver process to inform the network interface card (NIC), before the data arrives, about buffers in which incoming data has to be placed. To ensure that the receiver process is ready to receive data, the sender process typically performs flow-control on the data transmission. Existing designs of SDP flow-control are naive and do not take advantage of several interesting features provided by IB. Specifically, features such as RDMA are only used for performing zero-copy communication, although RDMA has more capabilities such as sender-side buffer management (where a sender process can manage SDP resources for the sender as well as the receiver). Similarly, IB also provides hardware flow-control capabilities that have not been studied in previous literature. In this paper, we utilize these capabilities to improve the SDP flow-control over IB using two designs: RDMA-based flow-control and NIC-assisted RDMA-based flow-control. We evaluate the designs using micro-benchmarks and real applications. Our evaluations reveal that these designs can improve the resource usage of SDP and consequently its performance by an order-of-magnitude in some cases. Moreover we can achieve 10-20% improvement in performance for various applications. Pavan Balaji, Sitha Bhagvat, Dhabaleswar K. Panda 0001, Rajeev Thakur, William Gropp |
ICPP | 1 |
| 2007 | Analyzing and Minimizing the Impact of Opportunity Cost in QoS-aware Job SchedulingabstractQuality of service (QoS) mechanisms allowing users to request for turn-around time guarantees for their jobs have recently generated much interest. In our previous work we had designed a framework, QoPS, to allow for such QoS. This framework provides an admission control mechanism that only accepts jobs whose requested deadlines can be met and, once accepted, guarantees these deadlines. However, the framework is completely blind to the revenue these jobs can fetch for the supercomputer center. By accepting a job, the supercomputer center might relinquish its capability to accept some future arriving (and potentially more expensive) jobs. In other words, while each job pays an explicit price to the system for running it, the system may also be viewed as paying an implicit opportunity cost by accepting the job. Thus, accepting a job is profitable only when the job's price is higher than its opportunity cost. In this paper we analyze the impact such opportunity cost can have on the overall revenue of the supercomputer center and attempt to minimize it through predictive techniques. Specifically, we propose two extensions to QoPS, value-aware QoPS (VQoPS) and dynamic value-aware QoPS (DVQoPS), to provide such capabilities. We present detailed analysis of these schemes and demonstrate using simulation that they not only achieve several factors improvement in system revenue, but also good service differentiation as a much desired side-effect. Pavan Balaji, Gerald Sabin, P. Sadayappan |
ICPP | 2 |
| 2007 | Nonuniformly Communicating Noncontiguous Data: A Case Study with PETSc and MPIabstractDue to the complexity associated with developing parallel applications, scientists and engineers rely on high-level software libraries such as PETSc, ScaLAPACK and PESSL to ease this task. Such libraries assist developers by providing abstractions for mathematical operations, data representation and management of parallel layouts of the data, while internally using communication libraries such as MPI and PVM. With high-level libraries managing data layout and communication internally, it can be expected that they organize application data suitably for performing the library operations optimally. However, this places additional overhead on the underlying communication library by making the data layout noncontiguous in memory and communication volumes (data transferred by a process to each of the other processes) nonuniform. In this paper, we analyze the overheads associated with these two aspects (noncontiguous data layouts and nonuniform communication volumes) in the context of the PETSc software toolkit over the MPI communication library. We describe the issues with the current approaches used by MPICH2 (an implementation of MPI), propose different approaches to handle these issues and evaluate these approaches with micro-benchmarks as well as an application over the PETSc software library. Our experimental results demonstrate close to an order of magnitude improvement in the performance of a 3-D Laplacian multi-grid solver application when evaluated on a 128 processor cluster. Pavan Balaji, Darius Buntinas, Satish Balay, Barry Smith 0002, Rajeev Thakur, William Gropp |
IPDPS | 1 |
| 2007 | Designing Efficient Systems Services and Primitives for Next-Generation Data-CentersabstractCurrent data-centers lack in efficient support for intelligent services, such as requirements for caching documents and cooperation of caching servers, efficiently monitoring and managing the limited physical resources, load-balancing, controlling overload scenarios, that are becoming a common requirement today. On the other hand, the system area network (SAN) technology is making rapid advances during the recent years. Besides high performance, these modern interconnects are providing a range of novel features and their support in hardware (e.g., RDMA, atomic operations). In this paper, we extend our previously proposed framework comprising of three layers (communication protocol support, data-center service primitives and advanced data-center services) that work together to tackle the issues associated with existing data-centers. We present the performance results using data-center services such as cooperative caching and active resource monitoring and data-center primitives such as distributed data sharing substrate and distributed lock manager, which demonstrate significant performance benefits achievable by our framework as compared to existing data-centers in several cases. Karthikeyan Vaidyanathan, Sundeep Narravula, Pavan Balaji, Dhabaleswar K. Panda 0001 |
IPDPS | 3 |
| 2007 | Analyzing the impact of supporting out-of-order communication on in-order performance with iWARPabstractDue to the growing need to tolerate network faults and congestion in high-end computing systems, supporting multiple network communication paths is becoming increasingly important. However, multi-path communication comes with the disadvantage of out-of-order arrival of packets (because packets may traverse different paths). While modern networking stacks such as the Internet Wide-Area RDMA Protocol (iWARP) over 10-Gigabit Ethernet (10GE) support multi-path communication, their current implementations do not handle out-of-order packets primarily owing to the overhead on in-order communication that it adds. Specifically, in iWARP, supporting out-of-order packets requires every packet to carry additional information causing significant overhead on packets that arrive in-order. Thus, in this paper, we analyze the trade-offs in designing a feature-complete iWARP stack, i.e., one that provides support for out-of-order arriving packets, and thus, multi-path systems, while focusing on the performance of in-order communication. We propose three feature-complete designs of iWARP and analyze the pros and cons of each of these designs using performance experiments based on several micro-benchmarks as well as an iso-surface visual rendering application. Our analysis reveals that the iWARP design providing the best overall performance depends on the particular characteristics of the upper layers and that different designs are optimal based on the metric of interest. Pavan Balaji, Wu-chun Feng, Sitha Bhagvat, Dhabaleswar K. Panda 0001, Rajeev Thakur, William Gropp |
SC | 1 |
| 2006 | Asynchronous zero-copy communication for synchronous sockets in the sockets direct protocol (SDP) over InfiniBandabstractSockets direct protocol (SDP) is an industry standard pseudo sockets-like implementation to allow existing sockets applications to directly and transparently take advantage of the advanced features of current generation networks such as InfiniBand. The SDP standard supports two kinds of sockets semantics, viz., synchronous sockets (e.g., used by Linux, BSD, Windows) and asynchronous sockets (e.g., used by Windows, upcoming support in Linux). Due to the inherent benefits of asynchronous sockets, the SDP standard allows several intelligent approaches such as source-avail and sink-avail based zero-copy for these sockets. Unfortunately, most of these approaches are not beneficial for the synchronous sockets interface. Further, due to its portability, ease of use and support on a wider set of platforms, the synchronous sockets interface is the one used by most sockets applications today. Thus, a mechanism by which the approaches proposed for asynchronous sockets can be used for synchronous sockets is highly desirable. In this paper, we propose one such mechanism, termed as AZ-SDP (asynchronous zero-copy SDP), where we memory-protect application buffers and carry out communication asynchronously while maintaining the synchronous sockets semantics. We present our detailed design in this paper and evaluate the stack with an extensive set of benchmarks. The experimental results demonstrate that our approach can provide an improvement of close to 35% for medium-message unidirectional throughput and up to a factor of 2 benefit for computation-communication overlap tests and multi-connection benchmarks Pavan Balaji, Sitha Bhagvat, Hyun-Wook Jin, Dhabaleswar K. Panda 0001 |
IPDPS | 1 |
| 2006 | Designing next generation data-centers with advanced communication protocols and systems servicesabstractCurrent data-centers rely on TCP/IP over fast- and gigabit-Ethernet for data communication even within the cluster environment for cost-effective designs, thus limiting their maximum capacity. Together with raw performance, such data-centers also lack in efficient support for intelligent services, such as requirements for caching documents, managing limited physical resources, load-balancing, controlling overload scenarios, and prioritization and QoS mechanisms, that are becoming a common requirement today. On the other hand, the system area network (SAN) technology is making rapid advances during the recent years. Besides high performance, these modern interconnects are providing a range of novel features and their support in hardware (e.g., RDMA, atomic operations, QoS support). In this paper, we address the capabilities of these current generation SAN technologies in addressing the limitations of existing data-centers. Specifically, we present a novel framework comprising of three layers (communication protocol support, data-center service primitives and advanced data-center services) that work together to tackle the issues associated with existing data-centers. We also present preliminary results in the various aspects of the framework, which demonstrate close to an order of magnitude performance benefits achievable by our framework as compared to existing data-centers in several cases. Pavan Balaji, Karthikeyan Vaidyanathan, Sundeep Narravula, Hyun-Wook Jin, Dhabaleswar K. Panda 0001 |
IPDPS | 1 |
| 2005 | Architecture for caching responses with multiple dynamic dependencies in multi-tier data-centers over InfiniBandabstractIt has been well acknowledged in the research community that in order to design a data-center environment which is efficient and offers high performance, one of the critical issues that needs to be addressed is the effective reuse of cache content stored away from the origin server. However, for caching dynamically changing content (e.g., content involved in online banking, Internet auctions, etc.). consistency and coherency issues need to be addressed. In addition, most current real world requests have multiple dynamic dependencies, i.e., these requests might depend on multiple data objects. Further, these requests are not entirely independent; several requests might have common dependencies. While there have been previous research solutions on maintaining coherent caches for dynamic content, these solutions have several shortcomings including inability to adapt to server load or handle multiple dynamic dependencies. In this paper, we propose a load resilient architecture using one sided operations supported by several high performance interconnects such as InfiniBand, while maintaining multiple dynamic dependencies per response. Our experimental results show that our schemes to tackle the multi-dependency issue efficiently and significantly outperform the existing approaches. Further, our results demonstrate that the proposed load resilient architecture can possibly improve the performance of loaded data-centers by over an order of magnitude. Sundeep Narravula, Pavan Balaji, Karthikeyan Vaidyanathan, Hyun-Wook Jin, Dhabaleswar K. Panda 0001 |
CCGRID | 2 |
| 2005 | Head-to-TOE Evaluation of High-Performance Sockets over Protocol Offload EnginesabstractDespite the performance drawbacks of Ethernet, it still possesses a sizable footprint in cluster computing because of its low cost and backward compatibility to existing Ethernet infrastructure. In this paper, we demonstrate that these performance drawbacks can be reduced (and in some cases, arguably eliminated) by coupling TCP offload engines (TOEs) with 10-Gigabit Ethernet (10GigE). Although there exists significant research on individual network technologies such as 10GigE, InfiniBand (IBA), and Myrinet; to the best of our knowledge, there has been no work that compares the capabilities and limitations of these technologies with the recently introduced 10GigE TOEs in a homogeneous experimental testbed. Therefore, we present performance evaluations across 10GigE, IBA, and Myrinet (with identical cluster-compute nodes) in order to enable a coherent comparison with respect to the sockets interface. Specifically, we evaluate the network technologies at two levels: (i) a detailed micro-benchmark evaluation and (ii) an application-level evaluation with sample applications from different domains, including a bio-medical image visualization tool known as the Virtual Microscope, an iso-surface oil reservoir simulator, a cluster file-system known as the parallel virtual file-system (PVFS), and a popular cluster management tool known as Ganglia. In addition to 10GigE's advantage with respect to compatibility to wide-area network infrastructures, e.g., in support of grids, our results show that 10GigE also delivers performance that is comparable to traditional high-speed network technologies such as IBA and Myrinet in a system-area network environment to support clusters and that 10GigE is particularly well-suited for sockets-based applications Pavan Balaji, Wu-chun Feng, Qi Gao 0004, Ranjit Noronha, Weikuan Yu, Dhabaleswar K. Panda 0001 |
CLUSTER | 1 |
| 2005 | Supporting iWARP Compatibility and Features for Regular Network AdaptersabstractWith several recent initiatives in the protocol offloading technology present on network adapters, the user market is now distributed amongst various technology levels including regular Ethernet network adapters, TCP Offload Engines (TOEs) and the recently introduced iWARP-capable networks. While iWARP-capable networks provide all the features provided by their predecessors (TOEs and regular Ethernet network adapters) and a new richer programming interface, they lack with respect to backward compatibility. In this aspect, two important issues need to be considered. First, not all network adapters support iWARP; thus, software compatibility for regular network adapters (which have no offloaded protocol stack) with iWARP capable network adapters needs to be achieved. Second, several applications on top of regular Ethernet as well as TOE based adapters have been written with the sockets interface; rewriting such applications using the new iWARP interface is cumbersome and impractical. Thus, it is desirable to have an interface which provides a two-fold benefit: (i) it allows existing applications to run directly without any modifications and (ii) it exposes the richer feature set of iWARP to the applications to be utilized with minimal modifications. In this paper, we design and implement a software stack to handle these issues. Specifically, (i) the software stack emulates the functionality of the iWARP stack in software to provide compatibility for regular Ethernet adapters with iWARP capable networks and (ii) it provides applications with an extended sockets interface that provides the traditional sockets functionality as well as functionality extended with the rich iWARP features Pavan Balaji, Hyun-Wook Jin, Karthikeyan Vaidyanathan, Dhabaleswar K. Panda 0001 |
CLUSTER | 1 |
| 2005 | On the provision of prioritization and soft qos in dynamically reconfigurable shared data-centers over infinibandabstractIn the past few years several researchers have proposed and configured data-centers providing multiple independent services, known as shared data-centers. For example, several ISPs and other Web service providers host multiple unrelated Web-sites on their data-centers allowing potential differentiation in the service provided to each of them. Such differentiation becomes essential in several scenarios in a shared data-center environment. In this paper, we extend our previously proposed scheme on dynamic re-configurability to allow service differentiation in the shared data-center environment. In particular, we point out the issues associated with the basic dynamic configurability scheme and propose two extensions to it, namely (i) dynamic reconfiguration with prioritization and (ii) dynamic reconfiguration with prioritization and QoS. Our experimental results show that our extensions can allow the dynamic reconfigurability scheme to attain a performance improvement of up to five times for high priority Web sites irrespective of any background low priority requests. Also, these extensions are able to significantly improve the performance of low priority requests when there are minimal or no high priority requests in the system. Further, they can achieve a similar performance as a static scheme with up to 43% lesser nodes in some cases Pavan Balaji, Sundeep Narravula, Karthikeyan Vaidyanathan, Hyun-Wook Jin, Dhabaleswar K. Panda 0001 |
ISPASS | 1 |
| 2005 | Exploiting NIC architectural support for enhancing IP-based protocols on high-performance networks
Hyun-Wook Jin, Pavan Balaji, Chuck Yoo, Dhabaleswar K. Panda 0001 |
J. Parallel Distributed Comput. | 2 |
| 2004 | Towards provision of quality of service guarantees in job schedulingabstractConsiderable research has focused on the problem of scheduling dynamically arriving independent parallel jobs on a given set of resources. There has also been some recent work in the direction of providing differentiated service to different classes of jobs using statically or dynamically calculated priorities assigned to the jobs. However, the potential and usability of a quality of service based scheme has not been much studied. In This work, we extend a previously proposed scheme (QoPS) to provide quality of service to submitted jobs; we propose extensions to the algorithm in multiple aspects: (i) studying the effect of user tolerance towards missed deadlines on the overall profit attainable by the supercomputer center, (it) providing artificial slack to some jobs to maximize the overall profit and (hi) utilizing a kill-and-restart mechanism to further improve the profit attainable. Pavan Balaji, P. Sadayappan, Dhabaleswar K. Panda 0001 |
CLUSTER | 2 |
| 2004 | Sockets Direct Protocol over InfiniBand in clusters: is it beneficial?abstractThe Sockets Direct Protocol (SDP) had been proposed recently in order to enable sockets based applications to take advantage of the enhanced features provided by InfiniBand architecture. In this paper, we study the benefits and limitations of an implementation of SDP. We first analyze the performance of SDP based on a detailed suite of micro-benchmarks. Next, we evaluate it on two different real application domains: (1) A multitier data-center environment and (2) A Parallel Virtual File System (PVFS). Our micro-benchmark results show that SDP is able to provide up to 2.7 times better bandwidth as compared to the native sockets implementation over InfiniBand (IPoIB) and significantly better latency for large message sizes. Our experimental results also show that SDP is able to achieve a considerably higher performance (improvement of up to 2.4 times) as compared to IPoIB in the PVFS environment. In the data-center environment, SDP outperforms IPoIB for large file transfers inspite of currently being limited by a high connection setup time. However, this limitation is entirely implementation specific and as the InfiniBand software and hardware products are rapidly maturing, we expect this limitation to be overcome soon. Based on this, we have shown that the projected performance for SDP, without the connection setup time, can outperform IPoIB for small message transfers as well. Pavan Balaji, Sundeep Narravula, Karthikeyan Vaidyanathan, Savitha Krishnamoorthy, Jiesheng Wu, Dhabaleswar K. Panda 0001 |
ISPASS | 1 |
| 2003 | Impact of High Performance Sockets on Data Intensive ApplicationsabstractThe challenging issues in supporting data intensive applications on clusters include efficient movement of large volumes of data between processor memories and efficient coordination of data movement and processing by a runtime support to achieve high performance. Such applications have several requirements such as guarantees in performance, scalability with these guarantees and adaptability to heterogeneous environments. With the advent of user-level protocols like the Virtual Interface Architecture (VIA) and the modern InfiniBand Architecture, the latency and bandwidth experienced by applications has approached to that of the physical network on clusters. In order to enable applications written on top of TCP/IP to take advantage of the high performance of these user-level protocols, researchers have come up with a number of techniques including User Level Sockets Layers over high performance protocols. In this paper, we study the performance and limitations of such substrate, referred to here as SocketVIA, using a component framework designed to provide runtime support for data intensive applications. The experimental results show that by reorganizing certain components of an application (in our case, the partitioning of a dataset into smaller data chunks), we can make significant improvements in application performance. This leads to a higher scalability of applications with performance guarantees. It also allows fine grained load balancing, hence making applications more adaptable to heterogeneity in resource availability. The experimental results also show that the different performance characteristics of SocketVIA allow a more efficient partitioning of data at the source nodes, thus improving the performance of the application up to an order of magnitude in some cases. Pavan Balaji, Jiesheng Wu, Tahsin M. Kurç, Ümit V. Çatalyürek, Dhabaleswar K. Panda 0001, Joel H. Saltz |
HPDC | 1 |
| 2003 | QoPS: A QoS Based Scheme for Parallel Job Scheduling
Pavan Balaji, P. Sadayappan, Dhabaleswar K. Panda 0001 |
JSSPP | 2 |
| 2002 | High Performance User Level Sockets over Gigabit EthernetabstractWhile a number of user-level protocols have been developed to reduce the gap between the performance capabilities of the physical network and the performance actually available, applications that have already been developed on kernel based protocols such as TCP have largely been ignored. There is a need to make these existing TCP applications take advantage of the modern user-level protocols such as EMP or VIA which feature both low-latency and high bandwidth. We have designed, implemented and evaluated a scheme to support such applications written using the sockets API to run over EMP without any changes to the application itself. Using this scheme, we are able to achieve a latency of 28.5 /spl mu/s for the Datagram sockets and 37 /spl mu/s for Data Streaming sockets compared to a latency of 120 /spl mu/s obtained by TCP for 4-byte messages. This scheme attains a peak bandwidth of around 840 Mbps. Both the latency and the throughput numbers are close to those achievable by EMP. The ftp application shows twice as much benefit on our sockets interface while the Web server application shows up to six times performance enhancement as compared to TCP. To the best of our knowledge, this is the first such design and implementation for Gigabit Ethernet. Pavan Balaji, Piyush Shivam, Pete Wyckoff, Dhabaleswar K. Panda 0001 |
CLUSTER | 1 |