EDBT 2026 Demo / reviewers in the wild / expert
Md. Wasi-ur-Rahman
dblp:15/10332
· DBLP profile ↗
27ranked-venue papers
10as first author
0since 2021 · last 2018
—ORCID · none
Domains — the database's venue-derived domains; a paper can count in several
Systems, architecture and hardware · 19 · 7 first-authorArtificial intelligence and machine learning · 5 · 1 first-authorDatabases, data management, data science and information retrieval · 5 · 1 first-authorApplied, interdisciplinary, general and emerging computing · 5 · 1 first-authorSoftware engineering, systems software and programming languages · 2 · 1 first-author
Expertise — from the expertise taxonomy: the topics of the expert's papers under the CCF categories. A weight counts papers with recency: 1 for a paper about the topic, 0.3 when the topic is its context, halved every five years.
| Computer architecture, parallel and distributed computing, and storage systems
4 papers |
Storage systems · 50% Distributed systems · 30% High-performance computing · 13% | |
| Computer networks
1 paper |
Datacenter networks · 50% Internet architecture and protocols · 50% |
Topics — the 14 heaviest of 15, each with the papers that count most for it
| Topic | Weight | Papers | Last | Evidence papers |
|---|---|---|---|---|
Storage systems › file systems
distributed file system |
0.3 | 2 | 2014 | SOR-HDFS: a SEDA-based approach to maximize overlapping in RDMA-enhanced HDFS · HPDC 2014 High performance RDMA-based design of HDFS over InfiniBand · SC 2012 |
Storage systems › file systems › distributed file system › parallel file system
lustre file system |
0.3 | 1 | 2017 | A Comprehensive Study of MapReduce Over Lustre for Intermediate Data Placement and Shuffle Strategies on HPC Clusters · IEEE Trans. Parallel Distributed Syst. 2017 |
High-performance computing › data-intensive computing
MapReduce for HPC |
0.3 | 1 | 2017 | A Comprehensive Study of MapReduce Over Lustre for Intermediate Data Placement and Shuffle Strategies on HPC Clusters · IEEE Trans. Parallel Distributed Syst. 2017 |
Storage systems › file systems › distributed file system
parallel file system |
0.3 | 1 | 2017 | A Comprehensive Study of MapReduce Over Lustre for Intermediate Data Placement and Shuffle Strategies on HPC Clusters · IEEE Trans. Parallel Distributed Syst. 2017 |
Distributed systems › fault tolerance
checkpointing |
0.2 | 1 | 2014 | MIC-Check: a distributed check pointing framework for the intel many integrated cores architecture · HPDC 2014 |
Distributed systems › fault tolerance › checkpointing
distributed checkpointing |
0.2 | 1 | 2014 | MIC-Check: a distributed check pointing framework for the intel many integrated cores architecture · HPDC 2014 |
Distributed systems
fault tolerance |
0.2 | 1 | 2014 | MIC-Check: a distributed check pointing framework for the intel many integrated cores architecture · HPDC 2014 |
Datacenter networks
RDMA |
0.1 | 1 | 2012 | High performance RDMA-based design of HDFS over InfiniBand · SC 2012 |
Storage systems › file systems › distributed file system
HDFS |
0.1 | 1 | 2012 | High performance RDMA-based design of HDFS over InfiniBand · SC 2012 |
Distributed systems › distributed data processing
shuffle optimization |
0.1 | 1 | 2017 | A Comprehensive Study of MapReduce Over Lustre for Intermediate Data Placement and Shuffle Strategies on HPC Clusters · IEEE Trans. Parallel Distributed Syst. 2017 |
Interconnection networks and networks-on-chip › cluster interconnect
infiniband |
0.1 | 1 | 2014 | SOR-HDFS: a SEDA-based approach to maximize overlapping in RDMA-enhanced HDFS · HPDC 2014 |
Storage systems
i/o optimization |
0.1 | 1 | 2014 | SOR-HDFS: a SEDA-based approach to maximize overlapping in RDMA-enhanced HDFS · HPDC 2014 |
Processor architecture and microarchitecture
many-core architecture |
0.1 | 1 | 2014 | MIC-Check: a distributed check pointing framework for the intel many integrated cores architecture · HPDC 2014 |
Interconnection networks and networks-on-chip
remote direct memory access |
0.1 | 1 | 2014 | SOR-HDFS: a SEDA-based approach to maximize overlapping in RDMA-enhanced HDFS · HPDC 2014 |
Methods — techniques the papers use, named apart from their topics
RDMA · 0.8JNI · 0.3priority directory selection · 0.3online profiling · 0.3staged event-driven architecture · 0.2SEDA · 0.2
| Year | Publication | Venue | Position |
|---|---|---|---|
| 2018 | MR-Advisor: A comprehensive tuning, profiling, and prediction tool for MapReduce execution frameworks on HPC clusters
Md. Wasi-ur-Rahman, Nusrat S. Islam, Xiaoyi Lu 0001, Dipti Shankar, Dhabaleswar K. Panda 0001 |
J. Parallel Distributed Comput. | 1 |
| 2017 | NVMD: Non-volatile memory assisted design for accelerating MapReduce and DAG execution frameworks on HPC systemsabstractIn this paper, we propose an accelerated execution framework (NVMD) for MapReduce and Directed Acyclic Graph (DAG) based processing engines to leverage the benefits of Non-Volatile Memory (NVM). Through NVMD, novel features for MapReduce, such as a hybrid push and pull shuffle mechanism, non-blocking send and receive operations, and dynamic adaptation to the network congestion have been presented. The design has been adopted in two different data intensive computing middleware: Hadoop and Tez. Performance results illustrate that NVMD can out-perform the current best execution frameworks by a significant margin. Md. Wasi-ur-Rahman, Nusrat S. Islam, Xiaoyi Lu 0001, Dhabaleswar K. Panda 0001 |
IEEE BigData | 1 |
| 2017 | A Comprehensive Study of MapReduce Over Lustre for Intermediate Data Placement and Shuffle Strategies on HPC ClustersabstractWith high performance interconnects and parallel file systems, running MapReduce over modern High Performance Computing (HPC) clusters has attracted much attention due to its uniqueness of solving data analytics problems with a combination of Big Data and HPC technologies. Since the MapReduce architecture relies heavily on the availability of local storage media, the Lustre-based global storage in HPC clusters poses many new opportunities and challenges. In this paper, we perform a comprehensive study on different MapReduce over Lustre deployments and propose a novel high-performance design of YARN MapReduce on HPC clusters by utilizing Lustre as the additional storage provider for intermediate data. With a deployment architecture where both local disks and Lustre are utilized for intermediate data storage, we propose a novel priority directory selection scheme through which RDMA-enhanced MapReduce can choose the best intermediate storage during runtime by on-line profiling. Our results indicate that, we can achieve 44 percent performance benefit for shuffle-intensive workloads in leadership-class HPC systems. Our priority directory selection scheme can improve the job execution time by 63 percent over default MapReduce while executing multiple concurrent jobs. To the best of our knowledge, this is the first such comprehensive study for YARN MapReduce with Lustre and RDMA. Md. Wasi-ur-Rahman, Nusrat S. Islam, Xiaoyi Lu 0001, Dhabaleswar K. Panda 0001 |
IEEE Trans. Parallel Distributed Syst. | 1 |
| 2016 | Efficient data access strategies for Hadoop and Spark on HPC cluster with heterogeneous storageabstractThe most popular Big Data processing frameworks of these days are Hadoop MapReduce and Spark. Hadoop Distributed File System (HDFS) is the primary storage for these frameworks. Big Data frameworks like Hadoop MapReduce and Spark launch tasks based on data locality. In the presence of heterogeneous storage devices, when different nodes have different storage characteristics, only locality-aware data access cannot always guarantee optimal performance. Rather, storage type becomes important, specially when high performance SSD and in-memory storage devices along with high performance interconnects are available. Therefore, in this paper, we propose efficient data access strategies (e.g. Greedy (prioritizes storage type over locality), Hybrid (balances the load for locality and high performance storage), etc.) for Hadoop and Spark considering both data locality and storage types. We re-design HDFS to accommodate the enhanced access strategies. Our evaluations show that, the proposed data access strategies can improve the read performance of HDFS by up to 33% compared to the default locality-aware data access. The execution times of Hadoop and Spark Sort are also reduced by up to 32% and 17%. The performances of Hadoop and Spark TeraSort are also improved by up to 11% through our design. Nusrat S. Islam, Md. Wasi-ur-Rahman, Xiaoyi Lu 0001, Dhabaleswar K. Panda 0001 |
IEEE BigData | 2 |
| 2016 | High Performance Design for HDFS with Byte-Addressability of NVM and RDMAabstractNon-Volatile Memory (NVM) offers byte-addressability with DRAM like performance along with persistence. Thus, NVMs provide the opportunity to build high-throughput storage systems for data-intensive applications. HDFS (Hadoop Distributed File System) is the primary storage engine for MapReduce, Spark, and HBase. Even though HDFS was initially designed for commodity hardware, it is increasingly being used on HPC (High Performance Computing) clusters. The outstanding performance requirements of HPC systems make the I/O bottlenecks of HDFS a critical issue to rethink its storage architecture over NVMs. In this paper, we present a novel design for HDFS to leverage the byte-addressability of NVM for RDMA (Remote Direct Memory Access)-based communication. We analyze the performance potential of using NVM for HDFS and re-design HDFS I/O with memory semantics to exploit the byte-addressability fully. We call this design NVFS (NVM- and RDMA-aware HDFS). We also present cost-effective acceleration techniques for HBase and Spark to utilize the NVM-based design of HDFS by storing only the HBase Write Ahead Logs and Spark job outputs to NVM, respectively. We also propose enhancements to use the NVFS design as a burst buffer for running Spark jobs on top of parallel file systems like Lustre. Performance evaluations show that our design can improve the write and read throughputs of HDFS by up to 4x and 2x, respectively. The execution times of data generation benchmarks are reduced by up to 45%. The proposed design also reduces the overall execution time of the SWIM workload by up to 18% over HDFS with a maximum benefit of 37% for job-38. For Spark TeraSort, our proposed scheme yields a performance gain of up to 11%. The performances of HBase insert, update, and read operations are improved by 21%, 16%, and 26%, respectively. Our NVM-based burst buffer can improve the I/O performance of Spark PageRank by up to 24% over Lustre. To the best of our knowledge, this paper is the first attempt to incorporate NVM with RDMA for HDFS. Nusrat S. Islam, Md. Wasi-ur-Rahman, Xiaoyi Lu 0001, Dhabaleswar K. Panda 0001 |
ICS | 2 |
| 2016 | High-Performance Hybrid Key-Value Store on Modern Clusters with RDMA Interconnects and SSDs: Non-blocking Extensions, Designs, and BenefitsabstractHigh-performance, distributed key-value store-based caching solutions, such as Memcached, have played a crucial role in enhancing the performance of many Online and Offline Big Data applications. The advent of high-performance storage (e.g. NVMe SSD) and interconnects (e.g. InfiniBand) on modern clusters has directed several efforts towards employing 'RAM+SSD' hybrid storagearchitectures for key-value stores running over RDMA, in order to achieve high data retention, while maintaining low latency and high throughput. In this paper, we first perform a detailed analysis of the behavior of hybrid Memcached designs, and identify two major bottlenecks: the client-side wait for request completion and the server-side SSD I/O overhead. Based on this analysis, we propose new non-blocking API extensions for Memcached Set and Get operations, to support high data retention while trying to achieve near in-memory speeds. We enhance the existing runtime designs on both the client and the server, and propose an adaptive slab manager with different I/O schemes for higher throughput. We demonstrate that Libmemcached-based applications can achieve high performance by exploiting the communication/computation overlap that is made possible by the proposed non-blocking API extensions, with either In-memory or SSD-assisted designs of RDMA-based Memcached. Performance evaluations show that the proposed extensions and designs can achieve up to 16x improvement for Memcached Set/Get latency over current hybrid design for RDMA-Memcached when all data does not fit in memory, and up to 3.6x improvement over pure in-memory design of default Memcached over 'IP-over-IB' when all data can fit in memory. Dipti Shankar, Xiaoyi Lu 0001, Nusrat S. Islam, Md. Wasi-ur-Rahman, Dhabaleswar K. Panda 0001 |
IPDPS | 4 |
| 2016 | MR-Advisor: A Comprehensive Tuning Tool for Advising HPC Users to Accelerate MapReduce Applications on SupercomputersabstractMapReduce is the most popular parallel computing framework for big data processing which allows massive scalability across distributed computing environment. Advanced RDMA-based design of Hadoop MapReduce has been proposed that alleviates the performance bottlenecks in default Hadoop MapReduce by leveraging the benefits from RDMA. On the other hand, data processing engine, Spark, provides fast execution of MapReduce applications through in-memory processing. Performance optimization for these contemporary big data processing frameworks on modern High-Performance Computing (HPC) systems is a formidable task because of the numerous configuration possibilities in each of them. In this paper, we propose MR-Advisor, a comprehensive tuning tool for MapReduce. MR-Advisor is generalized to provide performance optimizations for Hadoop, Spark, and RDMA-enhanced Hadoop MapReduce designs over different file systems such as HDFS, Lustre, and Tachyon. Performance evaluations reveal that, with MR-Advisor's suggested values, the job execution performance can be enhanced by a maximum of 58% over the current best-practice values for user-level configuration parameters. To the best of our knowledge, this is the first tool that supports tuning for both Apache Hadoop and Spark, as well as the RDMA and Lustre-based advanced designs. Md. Wasi-ur-Rahman, Nusrat S. Islam, Xiaoyi Lu 0001, Dipti Shankar, Dhabaleswar K. Panda 0001 |
SBAC-PAD | 1 |
| 2016 | Characterizing and benchmarking stand-alone Hadoop MapReduce on modern HPC clusters
Dipti Shankar, Xiaoyi Lu 0001, Md. Wasi-ur-Rahman, Nusrat S. Islam, Dhabaleswar K. Panda 0001 |
J. Supercomput. | 3 |
| 2015 | Performance characterization and acceleration of in-memory file systems for Hadoop and Spark applications on HPC clustersabstractFor data-intensive computing, the low throughput of the existing disk-bound storage systems is a major bottleneck. Recent emergence of the in-memory file systems with heterogeneous storage support mitigates this problem to a great extent. Parallel programming frameworks, e.g. Hadoop MapReduce and Spark are increasingly being run on such high-performance file systems. However, no comprehensive study has been done to analyze the impacts of the in-memory file systems on various Big Data applications. This paper characterizes two file systems in literature, Tachyon [17] and Triple-H [13] that support in-memory and heterogeneous storage, and discusses the impacts of these two architectures on the performance and fault tolerance of Hadoop MapReduce and Spark applications. We present a complete methodology for evaluating MapReduce and Spark workloads on top of in-memory file systems and provide insights about the interactions of different system components while running these workloads. We also propose advanced acceleration techniques to adapt Triple-H for iterative applications and study the impact of different parameters on the performance of MapReduce and Spark jobs on HPC systems. Our evaluations show that, although Tachyon is 5x faster than HDFS for primitive operations, Triple-H performs 47% and 2.4x better than Tachyon for MapReduce and Spark workloads, respectively. Triple-H also accelerates K-Means by 15% over HDFS and 9% over Tachyon. Nusrat S. Islam, Md. Wasi-ur-Rahman, Xiaoyi Lu 0001, Dipti Shankar, Dhabaleswar K. Panda 0001 |
IEEE BigData | 2 |
| 2015 | Benchmarking key-value stores on high-performance storage and interconnects for web-scale workloadsabstractLeveraging a distributed key-value based caching layer has proven to be invaluable for scalable data-intensive web applications. With the emergence of high-performance storage (e.g. SSD) and interconnects (e.g. InfiniBand) on modern clusters, several efforts are being made to design high-performance key-value stores that can operate well with `RAM+SSD' hybrid storage architecture. This has made it essential for us to design micro-benchmarks that are tailored to evaluate these upcoming, hybrid designs. In this paper, we study popular web-scale and cloud serving workloads, to identify different application-specific aspects, including commonly occurring data request distributions, update patterns, and environmental factors, that affect the performance of hybrid key-value stores. Based on these characterization studies, we propose a micro-benchmark suite that can be used to study high-performance, hybrid key-value stores on modern clusters, from the perspectives of both the application and the key-value store. We demonstrate its ease-of-use using database-integrated and stand-alone execution modes. Performance evaluations with different Memcached distributions, such as SSD-Assisted RDMA-Memcached, fatcache, and twemcache, over different networks/protocols, show that `SSD+RDMA' can significantly enhance the performance of Memcached for various read-only and read-heavy workloads, that are representative of several common web-scale workloads. Dipti Shankar, Xiaoyi Lu 0001, Md. Wasi-ur-Rahman, Nusrat S. Islam, Dhabaleswar K. Panda 0001 |
IEEE BigData | 3 |
| 2015 | Triple-H: A Hybrid Approach to Accelerate HDFS on HPC Clusters with Heterogeneous Storage ArchitectureabstractHDFS (Hadoop Distributed File System) is the primary storage of Hadoop. Even though data locality offered by HDFS is important for Big Data applications, HDFS suffers from huge I/O bottlenecks due to the tri-replicated data blocks and cannot efficiently utilize the available storage devices in an HPC (High Performance Computing) cluster. Moreover, due to the limitation of local storage space, it is challenging to deploy HDFS in HPC environments. In this paper, we present a hybrid design (Triple-H) that can minimize the I/O bottlenecks in HDFS and ensure efficient utilization of the heterogeneous storage devices (e.g. RAM, SSD, and HDD) available on HPC clusters. We also propose effective data placement policies to speed up Triple-H. Our design integrated with parallel file system (e.g. Lustre) can lead to significant storage space savings and guarantee fault-tolerance. Performance evaluations show that Triple-H can improve the write and read throughputs of HDFS by up to 7x and 2x, respectively. The execution times of data generation benchmarks are reduced by up to 3x. Our design also improves the execution time of the Sort benchmark by up to 40% over default HDFS and 54% over Lustre. The alignment phase of the Cloudburst application is accelerated by 19%. Triple-H also benefits the performance of SequenceCount and Grep in PUMA over both default HDFS and Lustre. Nusrat S. Islam, Xiaoyi Lu 0001, Md. Wasi-ur-Rahman, Dipti Shankar, Dhabaleswar K. Panda 0001 |
CCGRID | 3 |
| 2015 | Accelerating I/O Performance of Big Data Analytics on HPC Clusters through RDMA-Based Key-Value StoreabstractHadoop Distributed File System (HDFS) is the underlying storage engine of many Big Data processing frameworks such as Hadoop MapReduce, HBase, Hive, and Spark. Even though HDFS is well-known for its scalability and reliability, the requirement of large amount of local storage space makes HDFS deployment challenging on HPC clusters. Moreover, HPC clusters usually have large installation of parallel file system like Lustre. In this study, we propose a novel design to integrate HDFS with Lustre through a high performance key-value store. We design a burst buffer system using RDMA-based Mem cached and present three schemes to integrate HDFS with Lustre through this buffer layer, considering different aspects of I/O, data-locality, and fault-tolerance. Our proposed schemes can ensure performance improvement for Big Data applications on HPC clusters. At the same time, they lead to reduced local storage requirement. Performance evaluations show that, our design can improve the write performance of Test DFSIO by up to 2.6x over HDFS and 1.5x over Lustre. The gain in read throughput is up to 8x. Sort execution time is reduced by up to 28% over Lustre and 19% over HDFS. Our design can also significantly benefit I/O-intensive workloads compared to both HDFS and Lustre. Nusrat S. Islam, Dipti Shankar, Xiaoyi Lu 0001, Md. Wasi-ur-Rahman, Dhabaleswar K. Panda 0001 |
ICPP | 4 |
| 2015 | High-Performance Design of YARN MapReduce on Modern HPC Clusters with Lustre and RDMAabstractThe viability and benefits of running MapReduce over modern High Performance Computing (HPC) clusters, with high performance interconnects and parallel file systems, have attracted much attention in recent times due to its uniqueness of solving data analytics problems with a combination of Big Data and HPC technologies. Most HPC clusters follow the traditional Beowulf architecture with a separate parallel storage system (e.g. Lustre) and either no, or very limited, local storage. Since the MapReduce architecture relies heavily on the availability of local storage media, the Lustre-based global storage system in HPC clusters poses many new opportunities and challenges. In this paper, we propose a novel high-performance design for running YARN MapReduce on such HPC clusters by utilizing Lustre as the storage provider for intermediate data. We identify two different shuffle strategies, RDMA and Lustre Read, for this architecture and provide modules to dynamically detect the best strategy for a given scenario. Our results indicate that due to the performance characteristics of the underlying Lustre setup, one shuffle strategy may outperform another in different HPC environments, and our dynamic detection mechanism can deliver best performance based on the performance characteristics obtained during runtime of job execution. Through this design, we can achieve 44% performance benefit for shuffle-intensive workloads in leadership-class HPC systems. To the best of our knowledge, this is the first attempt to exploit performance characteristics of alternate shuffle strategies for YARN MapReduce with Lustre and RDMA. Md. Wasi-ur-Rahman, Xiaoyi Lu 0001, Nusrat S. Islam, Raghunath Rajachandrasekar, Dhabaleswar K. Panda 0001 |
IPDPS | 1 |
| 2015 | Can RDMA benefit online data processing workloads on memcached and MySQL?abstractAt the onset of the widespread usage of social networking services in the Web 2.0/3.0 era, leveraging a distributed and scalable caching layer like Memcached is often invaluable to application server performance. Since a majority of the existing clusters today are equipped with modern high speed interconnects such as InfiniBand, that offer high bandwidth and low latency communication, there is potential to improve the response time and throughput of the application servers, by taking advantage of advanced features like RDMA. We explore the potential of employing RDMA to improve the performance of Online Data Processing (OLDP) workloads on MySQL using Memcached for real-world web applications. Dipti Shankar, Xiaoyi Lu 0001, Md. Wasi-ur-Rahman, Nusrat S. Islam, Dhabaleswar K. Panda 0001 |
ISPASS | 4 |
| 2014 | In-memory I/O and replication for HDFS with Memcached: Early experiencesabstractHadoop is the de-facto standard platform for large-scale data analytic applications. In spite of high availability and reliability guarantees, Hadoop Distributed File System (HDFS) suffers from huge I/O bottlenecks for storing the tri-replicated data blocks. The I/O overheads intrinsic to the HDFS architecture degrade the application performance. In this paper, we present a novel design (MEM-HDFS) to perform intelligent caching and replication of HDFS data blocks in Memcached that can significantly improve the I/O performance. In this design, we consider different deployment strategies for the Memcached servers (local and remote) and guarantee persistence of the Memcached data to HDFS on cache replacements. Performance evaluations show that MEM-HDFS can increase the read and write throughput of HDFS by up to 3.9x and 3.3x, respectively. Our design can also significantly speed up the data loading (to HDFS) phase. It reduces the execution times of data generation benchmarks like, TeraGen, RandomTextWriter, and RandomWriter by up to 50%, 39%, and 48%, respectively. The performances of other benchmarks like TeraSort and Grep are also improved by the proposed design. Nusrat S. Islam, Xiaoyi Lu 0001, Md. Wasi-ur-Rahman, Raghunath Rajachandrasekar, Dhabaleswar K. Panda 0001 |
IEEE BigData | 3 |
| 2014 | MapReduce over Lustre: Can RDMA-Based Approach Benefit?
Md. Wasi-ur-Rahman, Xiaoyi Lu 0001, Nusrat S. Islam, Raghunath Rajachandrasekar, Dhabaleswar K. Panda 0001 |
Euro-Par | 1 |
| 2014 | SOR-HDFS: a SEDA-based approach to maximize overlapping in RDMA-enhanced HDFSabstractIn this paper, we propose SOR-HDFS, a SEDA (Staged Event-Driven Architecture)-based approach to improve the performance of HDFS Write operation. This design not only incorporates RDMA-based communication over InfiniBand but also maximizes overlapping among different stages of data transfer and I/O. Performance evaluations show that, the new design improves the aggregated write throughput of Enhanced DFSIO benchmark in Intel HiBench by up to 64% and reduces the job execution time by 37% compared to IPoIB (IP over InfiniBand). Compared to the previous best RDMA-enhanced design [4], the improvements in throughput and execution time are 30% and 20%, respectively. Our design can also improve the performance of HBase Put operation by up to 53% over IPoIB and 29% compared to the previous best RDMA-enhanced HDFS. To the best of our knowledge, this is the first design of SEDA-based HDFS in the literature. Nusrat S. Islam, Xiaoyi Lu 0001, Md. Wasi-ur-Rahman, Dhabaleswar K. Panda 0001 |
HPDC | 3 |
| 2014 | MIC-Check: a distributed check pointing framework for the intel many integrated cores architectureabstractThe advent of many-core architectures like Intel MIC is enabling the design of increasingly capable supercomputers within reasonable power budgets. Fault-tolerance is becoming more important with the increased number of components and the complexity in these heterogeneous clusters. Checkpoint-restart mechanisms have been traditionally used to enhance the dependability of applications, and to enable dynamic task rescheduling in the face of system failures. Naive checkpointing protocols, which are predominantly I/O-intensive, face severe performance bottlenecks on the Xeon Phi architecture due to several inherent and acquired limitations. Consequently, existing checkpointing frameworks are not capable of serving distributed MPI applications that leverage heterogeneous hardware architectures. This paper discusses the I/O limitations on the Xeon Phi system, and describes the architecture and design of a novel distributed checkpointing framework, namely MIC-Check, for HPC applications running on it. Raghunath Rajachandrasekar, Sreeram Potluri, Akshay Venkatesh, Khaled Hamidouche, Md. Wasi-ur-Rahman, Dhabaleswar K. Panda 0001 |
HPDC | 5 |
| 2014 | Performance Modeling for RDMA-Enhanced Hadoop MapReduceabstractHadoop MapReduce is a popular parallel programming paradigm that allows scalable and fault-tolerant solutions to data-intensive applications on modern clusters. However, the performance behavior of this framework shows its inability to take advantage of high-performance interconnects. Recent studies show that by leveraging the benefits of high-performance interconnects, the overall performance of MapReduce jobs can be greatly enhanced by using additional features like in-memory merge, pipelined merge and reduce, and pre-fetching and caching of map outputs. Existing performance models are not sufficient to predict the performance behavior for RDMA-enhanced MapReduce with these features. In this paper, we propose a detailed mathematical model of RDMA-enhanced MapReduce based on a number of cluster-wide and job-level configuration parameters. We also propose a simplified version of this model for prediction of large-scale MapReduce job executions and validate it in various system and workload configurations. Results derived from the proposed model match the experimental results within a 2-11% range. To the best of our knowledge, this is the first model that correctly predicts the behavior for RDMA-enhanced Hadoop MapReduce. Md. Wasi-ur-Rahman, Xiaoyi Lu 0001, Nusrat S. Islam, Dhabaleswar K. Panda 0001 |
ICPP | 1 |
| 2014 | HOMR: a hybrid approach to exploit maximum overlapping in MapReduce over high performance interconnectsabstractHadoop MapReduce is the most popular open-source parallel programming model extensively used in Big Data analytics. Although fault tolerance and platform independence make Hadoop MapReduce the most popular choice for many users, it still has huge performance improvement potentials. Recently, RDMA-based design of Hadoop MapReduce has alleviated major performance bottlenecks with the implementation of many novel design features such as in-memory merge, prefetching and caching of map outputs, and overlapping of merge and reduce phases. Although these features reduce the overall execution time for MapReduce jobs compared to the default framework, further improvement is possible if shuffle and merge phases can also be overlapped with the map phase during job execution. In this paper, we propose HOMR (a Hybrid approach to exploit maximum Overlapping in MapReduce), that incorporates not only the features implemented in RDMA-based design, but also exploits maximum possible overlapping among all different phases compared to current best approaches. Our solution introduces two key concepts: Greedy Shuffle Algorithm and On-demand Shuffle Adjustment, both of which are essential to achieve significant performance benefits over the default MapReduce framework. Architecture of HOMR is generalized enough to provide performance efficiency both over different Sockets interface as well as previous RDMA-based designs over InfiniBand. Performance evaluations show that HOMR with RDMA over InfiniBand can achieve performance benefits of 54% and 56% compared to default Hadoop over IPoIB (IP over InfiniBand) and 10GigE, respectively. Compared to the previous best RDMA-based designs, this benefit is 29%. HOMR over Sockets also achieves a maximum of 38-40% benefit compared to default Hadoop over Sockets interface. We also evaluate our design with real-world workloads like SWIM and PUMA, and observe benefits of up to 16% and 18%, respectively, over the previous best-case RDMA-based design. To the best of our knowledge, this is the first approach to achieve maximum possible overlapping for MapReduce framework. Md. Wasi-ur-Rahman, Xiaoyi Lu 0001, Nusrat S. Islam, Dhabaleswar K. Panda 0001 |
ICS | 1 |
| 2013 | Does RDMA-based enhanced Hadoop MapReduce need a new performance model?abstractRecent studies [17, 12] show that leveraging benefits of high performance interconnects like InfiniBand, MapReduce performance in terms of job execution time can be greatly enhanced by using additional features like in-memory merge, pipelined merge and reduce, and prefetching and caching of map outputs. In this paper, we validate that it is time to have a new performance model for the RDMA-based design of MapReduce over high performance interconnects. Our initial results derived from the proposed analytical model matches the experimental results within a 3--5% range. Md. Wasi-ur-Rahman, Xiaoyi Lu 0001, Nusrat S. Islam, Dhabaleswar K. Panda 0001 |
SoCC | 1 |
| 2013 | High-Performance Design of Hadoop RPC with RDMA over InfiniBandabstractHadoop RPC is the basic communication mechanism in the Hadoop ecosystem. It is used with other Hadoop components like MapReduce, HDFS, and HBase in real world data-centers, e.g. Facebook and Yahoo!. However, the current Hadoop RPC design is built on Java sockets interface, which limits its potential performance. The High Performance Computing community has exploited high throughput and low latency networks such as InfiniBand for many years. In this paper, we first analyze the performance of current Hadoop RPC design by unearthing buffer management and communication bottlenecks, that are not apparent on the slower speed networks. Then we propose a novel design (RPCoIB) of Hadoop RPC with RDMA over InfiniBand networks. RPCoIB provides a JVM-bypassed buffer management scheme and utilizes message size locality to avoid multiple memory allocations and copies in data serialization and deserialization. Our performance evaluations reveal that the basic ping-pong latencies for varied data sizes are reduced by 42%-49% and 46%-50% compared with 10GigE and IPoIB QDR (32Gbps), respectively, while the RPCoIB design also improves the peak throughput by 82% and 64% compared with 10GigE and IPoIB. As compared to default Hadoop over IPoIB QDR, our RPCoIB design improves the performance of the Sort benchmark on 64 compute nodes by 15%, while it improves the performance of CloudBurst application by 10%. We also present thorough, integrated evaluations of our RPCoIB design with other research directions, which optimize HDFS and HBase using RDMA over InfiniBand. Compared with their best performance, we observe 10% improvement for HDFS-IB, and 24% improvement for HBase-IB. To the best of our knowledge, this is the first such design of the Hadoop RPC system over high performance networks such as InfiniBand. Xiaoyi Lu 0001, Nusrat S. Islam, Md. Wasi-ur-Rahman, Hari Subramoni, Hao Wang 0002, Dhabaleswar K. Panda 0001 |
ICPP | 3 |
| 2012 | Scalable Memcached Design for InfiniBand Clusters Using Hybrid TransportsabstractMem cached is a general-purpose key-value based distributed memory object caching system. It is widely used in data-center domain for caching results of database calls, API calls or page rendering. An efficient Mem cached design is critical to achieve high transaction throughput and scalability. Previous research in the field has shown that the use of high performance interconnects like InfiniBand can dramatically improve the performance of Mem cached. The Reliable Connection (RC) is the most commonly used transport model for InfiniBand implementations. However, it has been shown that RC transport imposes scalability issues due to high memory consumption per connection. Such a characteristic is not favorable for middle wares like Mem cached, where the server is required to serve thousands of clients. The Unreliable Datagram (UD) transport offers higher scalability, but has several other limitations, which need to be efficiently handled. In this context, we introduce a hybrid transport model which takes advantage of the best features of RC and UD to deliver scalability and performance higher than that of a single-transport. To the best of our knowledge, this is the first effort aimed at studying the impact of using a hybrid of multiple transport protocols on Mem cached performance. We present comprehensive performance analysis using micro benchmarks, application benchmarks and realistic industry workloads. Our performance evaluations reveal that our Hybrid transport delivers performance comparable to that of RC, while maintaining a steady memory footprint. Mem cached Get latency for 4byte data size, is 4.28μs and 4.86μs for RC and hybrid transports, respectively. This represents a factor of twelve improvement over the performance of SDP. In evaluations using Apache Olio benchmark with 1,024 clients, Mem cached execution time using RC, UD and hybrid transports are 1.61, 1.96 and 1.70 seconds, respectively. Further, our scalability analysis with 4,096 client connections reveal that our proposed hybrid transport achieves good memory scalability. Hari Subramoni, Krishna Chaitanya Kandalla, Md. Wasi-ur-Rahman, Hao Wang 0002, Sundeep Narravula, Dhabaleswar K. Panda 0001 |
CCGRID | 4 |
| 2012 | High-Performance Design of HBase with RDMA over InfiniBandabstractHBase is an open source distributed Key/Value store based on the idea of Big Table. It is being used in many data-center Papplications (e.g. Face book, Twitter, etc.) because of its portability and massive scalability. For this kind of system, low latency and high throughput is expected when supporting services for large scale concurrent accesses. However, the existing HBase implementation is built upon Java Sockets Interface that provides sub-optimal performance due to the overhead to provide cross-platform portability. The byte-stream oriented Java sockets semantics confine the possibility to leverage new generations of network technologies. This makes it hard to provide high performance services for data-intensive applications. High Performance Computing (HPC) domain has exploited high performance and low latency networks such as Infini Band for many years. These interconnects provide advanced network features, such as Remote Direct Memory Access (RDMA), to achieve high throughput and low latency along with low CPU utilization. RDMA follows memory-block semantics, which can be adopted efficiently to satisfy the object transmission primitives used in HBase. In this paper, we present a novel design of HBase for RDMA capable networks via Java Native Interface (JNI). Our design extends the existing open-source HBase software and makes it RDMA capable. Our performance evaluation reveals that latency of HBase Get operations of 1KB message size can be reduced to 43.7μs with the new design on QDR platform (32 Gbps). This is about a factor of 3.5 improvement over 10 Gigabit Ethernet (10 GigE) network with TCP Offload. Throughput evaluations using four HBase region servers and 64 clients indicate that the new design boosts up throughput by 3 X times over 1 GigE and 10 GigE networks. To the best of our knowledge, this is first HBase design utilizing high performance RDMA capable interconnects. Jian Huang 0006, Xiangyong Ouyang, Md. Wasi-ur-Rahman, Hao Wang 0002, Miao Luo, Hari Subramoni, Chet Murthy, Dhabaleswar K. Panda 0001 |
IPDPS | 4 |
| 2012 | Understanding the communication characteristics in HBase: What are the fundamental bottlenecks?abstractHBase is an open source, distributed, column-oriented Key/Value database. In this paper, we focus on analyzing the performance aspects of HBase. Existing literature on HBase provides high level descriptions of the operations and present overall performance results. We conducted comprehensive experiments and identified different factors contributing to the overall latency of Get and Put operations. Our experimental results reveal that communication time is about 67% and 45% for a 1 KB Get request over 1 Gigabit Ethernet (1 GigE) and 10 Gigabit Ethernet (10 GigE) networks, respectively, for in-memory workloads. Our results show that HBase communication stack and associated operations need to be re-designed for high-performance networks like InfiniBand and its features. Md. Wasi-ur-Rahman, Jian Huang 0006, Xiangyong Ouyang, Hao Wang 0002, Nusrat S. Islam, Hari Subramoni, Chet Murthy, Dhabaleswar K. Panda 0001 |
ISPASS | 1 |
| 2012 | High performance RDMA-based design of HDFS over InfiniBandabstractHadoop Distributed File System (HDFS) acts as the primary storage of Hadoop and has been adopted by reputed organizations (Facebook, Yahoo! etc.) due to its portability and fault-tolerance. The existing implementation of HDFS uses Javasocket interface for communication which delivers suboptimal performance in terms of latency and throughput. For dataintensive applications, network performance becomes key component as the amount of data being stored and replicated to HDFS increases. In this paper, we present a novel design of HDFS using Remote Direct Memory Access (RDMA) over InfiniBand via JNI interfaces. Experimental results show that, for 5GB HDFS file writes, the new design reduces the communication time by 87% and 30% over 1Gigabit Ethernet (1GigE) and IP-over-InfiniBand (IPoIB), respectively, on QDR platform (32Gbps). For HBase, the Put operation performance is improved by 26% with our design. To the best of our knowledge, this is the first design of HDFS over InfiniBand networks. Nusrat S. Islam, Md. Wasi-ur-Rahman, Raghunath Rajachandrasekar, Hao Wang 0002, Hari Subramoni, Chet Murthy, Dhabaleswar K. Panda 0001 |
SC | 2 |
| 2011 | Memcached Design on High Performance RDMA Capable InterconnectsabstractMemcached is a key-value distributed memory object caching system. It is used widely in the data-center environment for caching results of database calls, API calls or any other data. Using Memcached, spare memory in data-center servers can be aggregated to speed up lookups of frequently accessed information. The performance of Memcached is directly related to the underlying networking technology, as workloads are often latency sensitive. The existing Memcached implementation is built upon BSD Sockets interface. Sockets offers byte-stream oriented semantics. Therefore, using Sockets, there is a conversion between Memcached's memory-object semantics and Socket's byte-stream semantics, imposing an overhead. This is in addition to any extra memory copies in the Sockets implementation within the OS. Over the past decade, high performance interconnects have employed Remote Direct Memory Access (RDMA) technology to provide excellent performance for the scientific computation domain. In addition to its high raw performance, the memory-based semantics of RDMA fits very well with Memcached's memory-object model. While the Sockets interface can be ported to use RDMA, it is not very efficient when compared with low-level RDMA APIs. In this paper, we describe a novel design of Memcached for RDMA capable networks. Our design extends the existing open-source Memcached software and makes it RDMA capable. We provide a detailed performance comparison of our Memcached design compared to unmodified Memcached using Sockets over RDMA and 10 Gigabit Ethernet network with hardware-accelerated TCP/IP. Our performance evaluation reveals that latency of Memcached Get of 4 KB size can be brought down to 12 μs using ConnectX InfiniBand QDR adapters. Latency of the same operation using older generation DDR adapters is about 20 μs. These numbers are about a factor of four better than the performance obtained by using 10 GigE with TCP Offload. In addition, these latencies of Get requests over a range of message sizes are better by a factor of five to ten compared to IP over InfiniBand and Sockets Direct Protocol over InfiniBand. Further, throughput of small Get operations can be improved by a factor of six when compared to Sockets over 10 Gigabit Ethernet network. Similar factor of six improvement in throughput is observed over Sockets Direct Protocol using ConnectX QDR adapters. To the best of our knowledge, this is the first such memcached design on high performance RDMA capable interconnects. Hari Subramoni, Miao Luo, Minjia Zhang, Jian Huang 0006, Md. Wasi-ur-Rahman, Nusrat S. Islam, Xiangyong Ouyang, Hao Wang 0002, Sayantan Sur, Dhabaleswar K. Panda 0001 |
ICPP | 6 |