EDBT 2026 Demo / reviewers in the wild / expert
Xiaoyi Lu 0001
dblp:27/5660-1
· DBLP profile ↗
113ranked-venue papers
10as first author
32since 2021 · last 2026
0000-0001-7581-8905ORCID · verified
Domains — the database's venue-derived domains; a paper can count in several
Systems, architecture and hardware · 84 · 3 first-author · 27 since 2021Applied, interdisciplinary, general and emerging computing · 16 · 3 first-author · 1 since 2021Databases, data management, data science and information retrieval · 15 · 3 first-author · 2 since 2021Artificial intelligence and machine learning · 12 · 3 first-authorComputer networks · 1 · 1 since 2021Software engineering, systems software and programming languages · 1
| Year | Publication | Venue | Position |
|---|---|---|---|
| 2026 | eGPU: Production-Scale Elastic Sharing Over 10,000 GPUsabstractAs the cost of GPUs continues to rise, GPU-sharing solutions have become increasingly important for improving efficiency and maximizing resource utilization. At the same time, large-scale operational deployments of such solutions remain relatively less explored, especially in heterogeneous production environments where workload dynamics and orchestration complexity introduce new practical considerations. In this paper, we introduce eGPU, an elastic, efficient, and scalable GPU-sharing framework tailored for production-scale concurrent machine learning (ML) training and inference. eGPU enables fine-grained, runtime-adjustable sharing of GPUs across multiple jobs, while preserving high resource utilization and fault isolation. To address communication bottlenecks, eGPU supports native NVLink/NCCL-based communication between shared GPU instances, capabilities that are limited or unavailable in many existing designs. Built with production deployment in mind, eGPU integrates with Kubernetes (K8s) to support large-scale orchestration. It has been deployed and running stably in production clusters with over$\text{1 0, 0 0 0 ~ G P U s}$for five years. Our evaluation results show that eGPU achieves elastic and precise control over instance sizes, improves job efficiency by 21 % to 31% than SOTA sharing solutions, saves the number of GPUs required by up to$8 \times$, and improves cluster GPU utilization by more than$\mathrm{3} \times$. Xiaochuan Tang, Hao Qi 0008, Jianbo Dong, Yinghao Yu, Zhennan Xue, Daocheng Ying, Zheng Cao 0003, Xiaoyi Lu 0001 |
HPCA | 9 |
| 2026 | When RDMA Goes Long-Haul: Characterization, Modeling, and Verbs-Level Emulation with Implications for Federated LearningabstractLong-haul Remote Direct Memory Access (RDMA) is rapidly emerging as a viable mechanism for extending high-performance communication beyond traditional datacenter environments into wide-area networks (WANs). However, increased round-trip time (RTT) fundamentally alters RDMA behavior, shifting execution from open-loop injection toward a closed-loop, window-limited regime governed by resource constraints and concurrency effects. Consequently, performance characteristics diverge from established datacenter assumptions, rendering existing models insufficient for predicting performance at WAN scale. Understanding these dynamics requires real-world characterization, yet geographically distributed long-haul testbeds remain expensive, scarce, and difficult to reproduce. Yuke Li 0003, Zhonghao Chen, Xiaoyi Lu 0001 |
HPDC | 3 |
| 2026 | PACER: A Userspace Network Rate Controller in MPI with Adaptive Compression for Parallel Applications
Yuke Li 0003, Darren Ng, Arjun Kashyap, Sheng Di, Guanpeng Li, Xiaoyi Lu 0001 |
ICS | 6 |
| 2025 | PULSE: Power Usage Monitoring Leveraging Sensing of Electromagnetic Field
Rahul Sidramappa Hoskeri, Xiaoyi Lu 0001, Hua Huang 0003 |
EWSN | 2 |
| 2025 | DPU-KV: On the Benefits of DPU Offloading for In-Memory Key-Value Stores at the EdgeabstractIn-memory key-value stores (KVS) are widely used for edge data storage, where low latency and high throughput are essential. Data Processing Units (DPUs), with their low power use and offloading capabilities, suit resource-constrained edge computing. While DPUs offer a new design point for KVS, their offloading in edge environments remains underexplored and challenging. In this paper, we unveil the potential of offloading in-memory CPU-based KVS to SoC-based DPUs, specifically NVIDIA's BlueField-2 (BF-2) and BlueField-3 (BF-3), with the aim of enhancing KVS performance. We propose a principled exploration methodology of dividing a KVS (i.e., MICA) into its logical components and identifying the CPU-intensive KVS component (i.e., communication engine). Next, we perform fine-grained offloading analysis and explorations on DPUs. To maximize benefits in terms of latency and throughput from fine-grained KVS offloading on DPUs, we propose a series of significant performance optimizations, including a key-value-based queue-pair model, overlapped KV request/response processing, reduced DMA operations per KV batch, dual-communication engine, and a sharding-based design. Our key finding is that our proposed fine-grained KVS offloading designs on modern DPU architectures (i.e., BF-2 and BF-3) can provide much lower latency (up to 68%) and higher throughput (up to 36%) than MICA (CPU-only) and coarse-grained DPU offloading schemes at the edge. To our knowledge, this paper is the first to explore the performance benefits of fine-grained KVS offloading to DPUs at the edge. Arjun Kashyap, Yuke Li 0003, Xiaoyi Lu 0001 |
HPDC | 3 |
| 2025 | Understanding the Idiosyncrasies of Emerging BlueField DPUsabstractData Processing Units (DPUs) are becoming available in datacenter environments to offload/accelerate workloads from the host.However, a comprehensive analysis is required to help users determine how to effectively utilize DPUs for their workloads, considering the various configurations and generations available.To fill in this gap, we conduct a fair and rigorous characterization by performing 15 benchmarking tests to demonstrate the evolution of representative SoC-based DPUs, specifically NVIDIA's BlueField-1, BlueField-2, and BlueField-3.Our work surfaces several idiosyncrasies across three key characterization dimensions-network, DMA engine, and memory.For network, we exhaustively test two major DPU modes-on-path (and five submodes) and offpath modes.We develop DPUDMABench, a microbenchmark suite to systematically analyze different data exchange primitives supported by DPU's DMA engine.We also conduct two application case studies examining the DPU mode's performance impact on TCP/IP and RDMA-based key-value stores (MICA and HERD).Based on our multi-generational DPU characterization, we identify and summarize 14 major idiosyncrasies, along with providing guidelines for optimal system and future hardware design. Arjun Kashyap, Yuke Li 0003, Darren Ng, Xiaoyi Lu 0001 |
ICS | 4 |
| 2025 | FedDES: Discrete Event Based Performance Simulation for Federated Learning SystemsabstractFederated Learning (FL) is a scalable and privacy-preserving paradigm well-suited for edge computing. Real-world FL deployments face substantial systems challenges such as compute variability and communication delays, motivating researchers to leverage simulation before real deployment. Most existing FL simulators, however, struggle to scale efficiently and incur long runtimes even for small workloads. To address this, we present FedDES, a high-fidelity, framework-agnostic discrete-event simulation platform that accurately models the runtime behavior of FL systems, including client training, communication overhead, network dynamics, and aggregation strategies. FedDES supports flexible configurations and diverse aggregation approaches, achieving simulation error within 2% of real deployments and delivering over 1000× speedup compared to prior tools. Large-scale experiments with up to 131,072 clients further show that the aggregation strategy critically affects performance, especially under heterogeneous and variable network conditions typical of edge environments. Zhonghao Chen, Weicong Chen 0002, Kibaek Kim, Guanpeng Li, Sheng Di, Xiaoyi Lu 0001 |
SEC | 7 |
| 2025 | SPRT2: Scalable, Parallel, and Real-Time fMRI Data Analysis on Heterogeneous ArchitecturesabstractReal-time functional Magnetic Resonance Imaging (fMRI) data analysis using the Sequential Probability Ratio Test (SPRT) enables dynamic adjustments to experimental protocols and early session termination, improving data quality and reducing patient fatigue. However, implementing SPRT in real-time fMRI analysis presents significant challenges due to the need for large-scale, high-dimensional data processing within strict time constraints. Furthermore, the ongoing advancements in fMRI hardware are driving a data explosion in the field, necessitating solutions that scale effectively. Existing approaches fall short in meeting real-time requirements and fail to fully exploit High-Performance Computing (HPC) and Big Data technologies. In this paper, we introduce Scalable, Parallel, and Real-Time Sequential Probability Ratio Test ($\text{SPRT}^{2}$), a toolkit that integrates HPC and Big Data techniques to enable efficient real-time SPRT-based fMRI data analysis.$\text{SPRT}^{2}$combines novel performance optimizations, such as hint-assisted matrix chain multiplication and sparse matrix techniques on heterogeneous architectures (CPUs and GPUs), with an Apache Spark-based framework for scalability and fault tolerance. Evaluated across 23 human subject experiments,$\text{SPRT}^{2}$achieves real-time analysis within the 1 -second repetition time while minimizing computational resource utilization (just 180 CPU cores).$\text{SPRT}^{2}$reduces session lengths by up to 33% and improves data quality. Furthermore,$\text{SPRT}^{2}$demonstrates near-linear scalability, efficiently processing synthetic datasets (35.9 billion voxels) over HPC platforms with 1,000 CPU cores or 8 NVIDIA A100 GPUs. To the best of our knowledge,$\text{SPRT}^{2}$is the first solution to integrate HPC and Big Data technologies for real-time fMRI analysis, setting a new standard in computational neuroscience. This work highlights the convergence of HPC and Big Data technologies and opens new avenues for tackling complex computational challenges in scalable and real-time fMRI data analysis. Weicong Chen 0002, Sarah J. Carr, Curtis Tatsuoka, Xiaoyi Lu 0001 |
IPDPS | 5 |
| 2025 | SBMGT: Scaling Bayesian Multinomial Group TestingabstractGroup testing is a widely used binary classification method that efficiently distinguishes between samples with and without a binary-classifiable attribute by pooling and testing subsets of a group. Bayesian Group Testing (BGT) is the state-of-the-art approach, which integrates prior risk information into a Bayesian Boolean Lattice framework to minimize test counts and reduce false classifications. However, BGT, like other existing group testing techniques, struggles with multinomial group testing, where samples have multiple binary-classifiable attributes that can be individually distinguished simultaneously. We address this need by proposing Bayesian Multinomial Group Testing (BMGT), which includes a new Bayesian-based model and supporting theorems for an efficient and precise multinomial pooling strategy. We further design and develop SBMGT, a high-performance and scalable framework to tackle BMGT's computational challenges by proposing three key innovations: 1) a parallel binary-encoded product lattice model with up to 99.8% efficiency; 2) the Bayesian Balanced Partitioning Algorithm (BBPA), a multinomial pooling strategy optimized for parallel computation with up to 97.7% scaling efficiency on 4096 cores; and 3) a scalable multinomial group testing analytics framework, demonstrated in a real-world disease surveillance case study using AIDS and STDs datasets from Uganda, where SBMGT reduced tests by up to 54% and lowered false classification rates by 92% compared to BGT. Weicong Chen 0002, Hao Qi 0008, Curtis Tatsuoka, Xiaoyi Lu 0001 |
PPoPP | 4 |
| 2025 | DPAR: High-Performance, Secure, and Scalable Differential Privacy-based AllReduceabstractSecure, efficient, and scalable AllReduce-based data aggregation is essential for Artificial Intelligence (AI) and scientific applications on modern High-Performance Computing (HPC) and cloud infrastructures. As AllReduce is increasingly used across these distributed infrastructures, privacy has become a critical concern. State-of-the-art (SOTA) Homomorphic Encryption (HE)-based AllReduce solutions introduce high overhead, require secure key exchanges, and remain vulnerable to collusion. We propose DPAR, the first differentially private, collusion-resistant AllReduce framework optimized for large-scale HPC and AI workloads. DPAR introduces three key innovations: integrating Differential Privacy (DP) to eliminate collusion risks without key exchanges, scalable noise growth to preserve accuracy, and performance optimizations using a noise pooling mechanism. DPAR is a drop-in Message Passing Interface (MPI) AllReduce replacement, providing strong privacy with minimal performance cost. Evaluated on Delta and Frontier supercomputers with up to 8192 cores, DPAR outperforms the SOTA HE solution by up to 34.7% in modern AI workloads. Hao Qi 0008, Weicong Chen 0002, Chenghong Wang, Xiaoyi Lu 0001 |
SC | 4 |
| 2025 | HPC-R1: Characterizing R1-like Large Reasoning Models on HPCabstractLarge Reasoning Models (LRMs) are becoming increasingly popular as they offer advanced capabilities in logical inference, mathematical reasoning, and knowledge synthesis, even beyond those of standard language models. However, their complex training workflows present significant challenges in reproducibility, efficiency, and system-level optimization. This paper introduces HPC-R1, a comprehensive characterization of LRM training on the NERSC Perlmutter supercomputer, representing behavior on a Top500-ranked system. We analyze all major stages, including supervised fine-tuning (SFT), Group Relative Policy Optimization (GRPO)-based reinforcement learning (RL), autoregressive generation, and distillation using customized state-of-the-art frameworks. Our detailed performance analysis reveals key system inefficiencies and scaling behaviors. Through our in-depth analysis, we present 19 key observations across all stages, including 4 for SFT, 7 for GRPO-based RL, 6 for generation, and 2 for distillation. Based on these findings, we present several key recommendations to guide future HPC-AI system design. Adam Weingram, Zhonghao Chen, Hao Qi 0008, Xiaoyi Lu 0001 |
SC | 5 |
| 2025 | FedEFsz: Fair Cross-Silo Federated Learning System With Error-Bounded Lossy CompressionabstractCross-Silo federated learning systems have been identified as an efficient approach to scaling DNN training across geographically-distributed data silos to preserve the privacy of the training data. Communication efficiency and fairness are two major issues that need to be both satisfied when federated learning systems are deployed in practice. Simultaneously guaranteeing both of them, however, is exceptionally difficult because simply combining communication reduction and fairness optimization approaches often causes non-converged training or drastic accuracy degradation. To bridge this gap, we proposeFedEFsz. On the one hand, it integrates the state-of-the-art error-bounded lossy compressor SZ3 into cross-silo federated learning systems to significantly reduce communication traffic during the training. On the other hand, it achieves a high fairness (i.e., rather consistent model accuracy and performance across different clients) through a carefully designed heuristic algorithm that can tune the error-bound of SZ3 for different clients during the training. Extensive experimental results based on a GPU cluster with 65 GPU cards show thatFedEFszimproves the fairness across different benchmarks by up to$60.88\%$and meanwhile reduces the communication traffic by up to$315\times$. Sheng Di, Benben Liu, Zhuoran Ji, Guanpeng Li, Xiaoyi Lu 0001, Amelie Chi Zhou, Khalid Ayedh Alharthi, Jiannong Cao 0001 |
IEEE Trans. Parallel Distributed Syst. | 6 |
| 2024 | Kspeed: Beating I/O Bottlenecks of Data Provisioning for RDMA Training ClustersabstractThe rapidly-increasing computing power of GPUs has rendered the I/O subsystem a bottleneck for distributed deep learning (DL) training. Currently, substantial data preprocessing work (e.g., decoding) has to be conducted on CPUs for a wide range of training scenarios such as computer vision (CV) and audio. Unfortunately, the involvement of training nodes' host memory and/or CPUs on the critical path of loading data to GPUs incurs significant GPU stalls in modern RDMA training clusters, because CPUs are much slower than GPUs and the connection from PCIe switches to host memory tends to suffer from incast problems. Moreover, this also incurs high CPU usage and resource contention, which consequently causes data loading performance variation and stragglers. This paper presents KSpeed, a novel data provisioning framework for large-scale RDMA training clusters. As many data preprocessing tasks need to be done by CPUs, KSpeed organizes host memory and CPU resources in the cluster to build a disaggregated memory/CPU pool, where the nodes can read raw input data from backend storage to their host memory, preprocess the data by their CPUs if necessary, and write cached/preprocessed data (on demand) directly to the training workers' GPU memory to minimize GPU stalls. KSpeed leverages the multi-rail RDMA network to eliminate unnecessary memory copies, interference, and congestion. Evaluation on a 96-GPU cluster shows that KSpeed delivers$5.4 \times \sim 100 \times$higher data loading performance over the state-of-the-art designs (DPP and Alluxio). KSpeed achieves near-linear scalability as the GPU number increases from 8 to 512. Jianbo Dong, Hao Qi 0008, Tianjing Xu, Xiaoli Liu 0002, Rongyao Wang, Xiaoyi Lu 0001, Zheng Cao 0003, Binzhang Fu |
ICNP | 7 |
| 2024 | gZCCL: Compression-Accelerated Collective Communication Framework for GPU ClustersabstractGPU-aware collective communication has become a major bottleneck for modern computing platforms as GPU computing power rapidly rises. A traditional approach is to directly integrate lossy compression into GPU-aware collectives, which can lead to serious performance issues such as underutilized GPU devices and uncontrolled data distortion. In order to address these issues, in this paper, we propose gZCCL, a first-ever general framework that designs and optimizes GPU-aware, compression-enabled collectives with an accuracy-aware design to control error propagation. To validate our framework, we evaluate the performance on up to 512 NVIDIA A100 GPUs with real-world applications and datasets. Experimental results demonstrate that our gZCCL-accelerated collectives, including both collective computation (Allreduce) and collective data movement (Scatter), can outperform NCCL as well as Cray MPI by up to 4.5 × and 28.7 ×, respectively. Furthermore, our accuracy evaluation with an image-stacking application confirms the high reconstructed data quality of our accuracy-aware framework. Jiajun Huang 0001, Sheng Di, Xiaodong Yu 0001, Jinyang Liu 0003, Yafan Huang, Kenneth Raffenetti, Hui Zhou 0012, Kai Zhao 0008, Xiaoyi Lu 0001, Zizhong Chen, Franck Cappello, Yanfei Guo, Rajeev Thakur |
ICS | 10 |
| 2024 | Druto: Upper-Bounding Silent Data Corruption Vulnerability in GPU ApplicationsabstractDue to the increasing scale of high-performance computing (HPC) systems, transient hardware faults have become a major reliability concern. Consequently, Silent Data Corruptions (SDCs) due to these faults have been a common insidious consequence in GPU applications. Developers often measure the application resilience with a set of program test inputs available in the benchmark suite, assuming the resilience would not fluctuate much among different inputs. However, we observe that this assumption often results in an over-optimistic evaluation for GPU applications. As a result, the subsequent SDC protection following the evaluation can hardly meet the expected reliability bar in the production environment, where applications would run with potentially arbitrary input values. To this end, we propose Druto – a compiler-based automated technique that searches for inputs to incrementally approach the upper bound of a GPU application’s SDC probability. We develop Druto based on the property that the resilience profiles of a small group of representative threads in a GPU kernel can approximately rank various inputs in terms of the overall SDC probability. Therefore, Druto strategically steers the search towards new program inputs that efficiently portray the overall SDC probability. Evaluation shows that the SDC probability derived from Druto’s input generation is as much as 74× higher than that from existing techniques. Moreover, existing techniques cannot find our generated inputs even given 5× more search time. Md Hasanur Rahman 0001, Sheng Di, Shengjian Guo, Xiaoyi Lu 0001, Guanpeng Li, Franck Cappello |
IPDPS | 4 |
| 2024 | An Optimized Error-controlled MPI Collective Framework Integrated with Lossy CompressionabstractWith the ever-increasing computing power of supercomputers and the growing scale of scientific applications, the efficiency of MPI collective communications turns out to be a critical bottleneck in large-scale distributed and parallel processing. The large message size in MPI collectives is particularly concerning because it can significantly degrade the overall parallel performance. To address this issue, prior research simply applies the off-the-shelf fix-rate lossy compressors in the MPI collectives, leading to suboptimal performance, limited generalizability, and unbounded errors. In this paper, we propose a novel solution, called C-Coll, which leverages error-bounded lossy compression to significantly reduce the message size, resulting in a substantial reduction in communication cost. The key contributions are three-fold. (1) We develop two general, optimized lossy-compression-based frameworks for both types of MPI collectives (collective data movement as well as collective computation), based on their particular characteristics. Our framework not only reduces communication cost but also preserves data accuracy. (2) We customize SZx, an ultra-fast error-bounded lossy compressor, to meet the specific needs of collective communication. (3) We integrate C-Coll into multiple collectives, such as MPI Allreduce, MPI Scatter, and MPI Bcast, and perform a comprehensive evaluation based on real-world scientific datasets. Experiments show that our solution outperforms the original MPI collectives as well as multiple baselines and related efforts by 1.8–2.7×. Jiajun Huang 0001, Sheng Di, Xiaodong Yu 0001, Jinyang Liu 0003, Xiaoyi Lu 0001, Kenneth Raffenetti, Hui Zhou 0012, Kai Zhao 0008, Zizhong Chen, Franck Cappello, Yanfei Guo, Rajeev Thakur |
IPDPS | 7 |
| 2024 | Accelerating Lossy and Lossless Compression on Emerging BlueField DPU ArchitecturesabstractData compression has become a crucial technique in addressing performance bottlenecks caused by increasing data volumes in High-Performance Computing (HPC), Big Data, and Deep Learning (DL). Despite its potential to boost system performance, recent studies have identified significant challenges with existing compression methods, mainly due to their high computational demands amidst continuously growing data sizes. Concurrently, the advent of Data Processing Units (DPUs), equipped with programmable System-on-Chip (SoC) and specialized compression accelerators, offers a promising opportunity to alter the landscape of data compression. This paper explores the complexities and potential of leveraging NVIDIA BlueField DPUs to accelerate lossy and lossless compression. Towards this, we introduce PEDAL, an innovative library that leverages the hardware capabilities of DPUs to unify and optimize data compression designs. Moreover, we seamlessly co-design PEDAL with the popular MPICH MPI library, demonstrating up to 101x speedup in compression time and 88x decrease in communication latency. Drawing on these achievements, we share our experience with various research communities about accelerating data compression on DPUs in communication-oriented HPC scenarios. Yuke Li 0003, Arjun Kashyap, Weicong Chen 0002, Yanfei Guo, Xiaoyi Lu 0001 |
IPDPS | 5 |
| 2024 | NVMe-oPF: Designing Efficient Priority Schemes for NVMe-over-Fabrics with Multi-Tenancy SupportabstractResource disaggregation is prevalent in datacenters since it provides high resource utilization when compared to servers dedicated to either compute, memory, or storage. NVMe-over-Fabrics (NVMe-oF) is the standardized protocol for accessing disaggregated network storage. Currently, the NVMe-oF specification lacks semantics to prioritize I/O requests based on different application needs. Since applications have varying goals — latency-sensitive or throughput-critical I/O — we need to design efficient schemes to allow applications to specify the type of performance they wish to achieve. To this end, we propose a new NVMe-over-Priority-Fabrics (NVMe-oPF) protocol with multi-tenancy support that allows applications to specify whether to optimize for latency or throughput. NVMe-oPF proposes coalescing request completions, lock-free optimization, zero-copy queues, out-of-order request completion handling, and window size optimization for the specific I/O patterns, queue depths, and I/O sizes that yield the best performance. Our NVMe-oPF-10Gbps can achieve up to 2.94X improvement in throughput and reduces tail latency by up to 32.1% for highly concurrent multi-tenant read workloads when compared to the state-of-the-art userspace NVMe-oF runtime design in Intel Storage Performance Development Kit (SPDK). For write workloads with 100Gbps, NVMe-oPF achieves a 32.6% increase in throughput while maintaining low latency compared to SPDK. We also bring performance benefits to the application level with HDF5 by increasing write workload throughput by 25.2% in larger-scale experiments. Darren Ng, Andrew Lin, Arjun Kashyap, Guanpeng Li, Xiaoyi Lu 0001 |
IPDPS | 5 |
| 2024 | hZCCL: Accelerating Collective Communication with Co-Designed Homomorphic CompressionabstractAs network bandwidth struggles to keep up with rapidly growing computing capabilities, the efficiency of collective communication has become a critical challenge for exa-scale distributed and parallel applications. Traditional approaches directly utilize error-bounded lossy compression to accelerate collective computation operations, exposing unsatisfying performance due to the expensive decompression-operation-compression (DOC) workflow. To address this issue, we present a first-ever homomorphic compression-communication co-design, hZCCL, which enables operations to be performed directly on compressed data, saving the cost of time-consuming decompression and recompression. In addition to the co-design framework, we build a light-weight compressor, optimized specifically for multi-core CPU platforms. We also present a homomorphic compressor with a run-time heuristic to dynamically select efficient compression pipelines for reducing the cost of DOC handling. We evaluate $h \mathbf{Z C C L}$ with up to 512 nodes and across five application datasets. The experimental results demonstrate that our homomorphic compressor achieves a CPU throughput of up to $379.08 \mathrm{~GB} / \mathrm{s}$, surpassing the conventional DOC workflow by up to $36.53 \times$. Moreover, our $\boldsymbol{h Z C C L}$-accelerated collectives outperform two state-of-the-art baselines, delivering speedups of up to $2.12 \times$ and $6.77 \times$ compared to original MPI collectives in single-thread and multi-thread modes, respectively, while maintaining data accuracy. Jiajun Huang 0001, Sheng Di, Xiaodong Yu 0001, Jinyang Liu 0003, Zizhe Jian, Xin Liang 0001, Kai Zhao 0008, Xiaoyi Lu 0001, Zizhong Chen, Franck Cappello, Yanfei Guo, Rajeev Thakur |
SC | 9 |
| 2024 | Versatile Datapath Soft Error Detection on the Cheap for HPC ApplicationsabstractWith the ongoing reduction in technology sizes and voltage levels, modern microprocessors are increasingly susceptible to soft errors, corrupting datapath units during program execution. While these error types have received considerable attention recently, existing solutions either confine themselves to limited scopes or incur massive overheads in performance and power consumption, hindering practical usage. In this work, we propose CONDA, a novel error detection technique based on code transformation and static program analysis, achieving versatile datapath protection at low cost. At compile time, ConDa analyzes program characteristics and transforms the original program code without complicating its control-flow and memory access patterns. At runtime, ConDa detects datapath errors with low overhead and latency. The evaluation of 38 benchmarks and a parallel HPC simulation reveals that CONDA only incurs 57.79% runtime overhead, which is 41.84% faster than existing state-of-the-art, with the same level of error detection effectiveness and low detection latency. Yafan Huang, Sheng Di, Xiaoyi Lu 0001, Guanpeng Li |
SC | 4 |
| 2024 | On the Feasibility and Benefits of Extensive EvaluationabstractBenchmark and system parameters often have a significant impact on performance evaluation, which raises a long-lasting question about which settings we should use. This paper studies the feasibility and benefits of extensive evaluation. A full extensive evaluation, which tests all possible settings, is usually too expensive. This work investigates whether it is possible to sample a subset of the settings and, upon them, generate observations that match those from a full extensive evaluation. Towards this goal, we have explored the incremental sampling approach, which starts by measuring a small subset of random settings, builds a prediction model on these samples using the popular ANOVA approach, adds more samples if the model is not accurate enough, and terminates otherwise. To summarize our findings: 1) Enhancing a research prototype to support extensive evaluation mostly involves changing hard-coded configurations, which does not take much effort. 2) Some systems are highly predictable, which means that they can achieve accurate predictions with a low sampling rate, but some systems are less predictable. 3) We have not found a method that can consistently outperform random sampling + ANOVA. Based on these findings, we provide recommendations to improve artifact predictability and strategies for selecting parameter values during evaluation. Yujie Hui, Miao Yu 0023, Hao Qi 0008, Yifan Gan, Tianxi Li, Yuke Li 0003, Xueyuan Ren, Sixiang Ma, Xiaoyi Lu 0001, Yang Wang 0009 |
Proc. ACM Manag. Data | 9 |
| 2023 | Characterizing Lossy and Lossless Compression on Emerging BlueField DPU ArchitecturesabstractThe Data Processing Unit (DPU) (i.e., programmable SmartNICs with System-on-Chip or SoC cores) has emerged as a valuable supplementary resource to the host CPU. The DPU architecture has been attracting significant attention within High-Performance Computing (HPC) and data center clusters due to its advanced capabilities and accelerators, which include a hardware-based data compression engine. This positions the DPU as a prospective tool for accelerating and offloading compression workloads from the hosts, which can potentially speed up data-intensive applications. The convergence of Big Data, HPC, and Machine Learning (ML) systems has rendered large data volumes a major performance bottleneck in message communication and data storage. While compression can boost performance, recent studies reveal that compression techniques (e.g., lossy and lossless) are compute-intensive and time-consuming, particularly with larger data sizes. Consequently, this paper characterizes the performance of three lossy (SZ3) and lossless (DEFLATE and zlib) compression algorithms with seven real-world data sets on the popular NVIDIA’s BlueField DPUs to explore potential opportunities for offloading these workloads from the host. We find that compared to DPU’s SoC cores, DPU’s hardware compression engine can obtain up to 26.8x performance speedup. Furthermore, we discuss the challenges and opportunities associated with employing NVIDIA’s BlueField DPUs to accelerate lossy and lossless compression/decompression workloads. Our research discloses five important takeaways which shed light on future research directions for lossy and lossless compressions on DPUs. Yuke Li 0003, Arjun Kashyap, Yanfei Guo, Xiaoyi Lu 0001 |
HOTI | 4 |
| 2023 | Performance Characterization of Large Language Models on High-Speed InterconnectsabstractLarge Language Models (LLMs) have recently gained significant popularity due to their ability to generate human-like text and perform a wide range of natural language processing tasks. Training these models usually requires a large amount of computational resources and is often done in a distributed manner. The use of high-speed interconnects can significantly influence the efficiency of distributed training. Therefore, there poses a need for systematic studies to explore the distributed training characteristics of these models on high-speed interconnects. This paper presents a comprehensive performance characterization of representative large language models: GPT, BERT, and T5. We evaluate their training performance in terms of iteration time, interconnect utilization, and scalability, over different high-speed interconnects and communication protocols, including TCP/IP, IPoIB, and RDMA. We observe that interconnects play a vital role in LLM training. Specifically, RDMA-100 Gbps outperforms IPoIB-100 Gbps and TCP/IP-10 Gbps by an average of 2.51x and 4.79x regarding training iteration time, and scores the highest interconnect utilization (up to 60 Gbps) in both strong and weak scaling, compared to IPoIB with up to 20 Gbps and TCP/IP with up to 9 Gbps, leading to the shortest training time. We also observe that larger models tend to have higher requirements for communication bandwidth, especially for AllReduce during backward propagation, which can take up to 91.12% of training time. Through our evaluation, we envision opportunities to improve the communication time for better training performance of LLMs. We extensively explore and summarize the role communication plays in distributed LLM training. Hao Qi 0008, Liuyao Dai, Weicong Chen 0002, Zhen Jia 0001, Xiaoyi Lu 0001 |
HOTI | 5 |
| 2023 | SBGT: Scaling Bayesian-based Group Testing for Disease SurveillanceabstractThe COVID-19 pandemic underscored the necessity for disease surveillance using group testing. Novel Bayesian methods using lattice models were proposed, which offer substantial improvements in group testing efficiency by precisely quantifying uncertainty in diagnoses, acknowledging varying individual risk and dilution effects, and guiding optimally convergent sequential pooled test selections using a Bayesian Halving Algorithm. Computationally, however, Bayesian group testing poses considerable challenges as computational complexity grows exponentially with sample size. This can lead to shortcomings in reaching a desirable scale without practical limitations. We propose a new framework for scaling Bayesian group testing based on Spark: SBGT. We show that SBGT is lightning fast and highly scalable. In particular, SBGT is up to 376x, 1733x, and 1523x faster than the state-of-the-art framework in manipulating lattice models, performing test selections, and conducting statistical analyses, respectively, while achieving up to 97.9% scaling efficiency up to 4096 CPU cores. More importantly, SBGT fulfills our mission towards reaching applicable scale for guiding pooling decisions in wide-scale disease surveillance, and other large scale group testing applications. Weicong Chen 0002, Hao Qi 0008, Xiaoyi Lu 0001, Curtis Tatsuoka |
IPDPS | 3 |
| 2023 | xCCL: A Survey of Industry-Led Collective Communication Libraries for Deep Learning
Adam Weingram, Yuke Li 0003, Hao Qi 0008, Darren Ng, Liuyao Dai, Xiaoyi Lu 0001 |
J. Comput. Sci. Technol. | 6 |
| 2022 | HiBGT: High-Performance Bayesian Group Testing for COVID-19abstractThe COVID-19 pandemic has necessitated disease surveillance using group testing. Novel Bayesian methods using lattice models were proposed, which offer substantial improvements in group testing efficiency by precisely quantifying uncertainty in diagnoses, acknowledging varying individual risk and dilution effects, and guiding optimally convergent sequential pooled test selections. Computationally, however, Bayesian group testing poses considerable challenges as computational complexity grows exponentially with sample size. HPC and big data stacks are needed for assessing computational and statistical performance across fluctuating prevalence levels at large scales. Here, we study how to design and optimize critical computational components of Bayesian group testing, including lattice model representation, test selection algorithms, and statistical analysis schemes, under the context of parallel computing. To realize this, we propose a high-performance Bayesian group testing framework named HiBGT, based on Apache Spark, which systematically explores the design space of Bayesian group testing and provides comprehensive heuristics on how to achieve high-performance, highly scalable Bayesian group testing. We show that HiBGT can perform large-scale test selections (> 250state iterations) and accelerate statistical analyzes up to 15.9x (up to 363x with little trade-offs) through a varied selection of sophisticated parallel computing techniques while achieving near linear scalability using up to 924 CPU cores. Weicong Chen 0002, Curtis Tatsuoka, Xiaoyi Lu 0001 |
HIPC | 3 |
| 2022 | NVMe-oAF: Towards Adaptive NVMe-oF for IO-Intensive Workloads on HPC CloudabstractApplications running inside containers or virtual machines, traditionally use TCP/IP for communication in HPC clouds and data centers. The TCP/IP path usually becomes a major performance bottleneck for applications performing NVMe-over-Fabrics (NVMe-oF) based I/O operations in disaggregated storage settings. We propose an adaptive communication channel, called NVMe-over-Adaptive-Fabric (NVMe-oAF), that applications could leverage to eliminate the high-latency and low-bandwidth incurred by remote I/O requests over TCP/IP. NVMe-oAF accelerates I/O intensive applications using locality awareness along with optimized shared memory and TCP/IP paths. The adaptiveness of the fabric stems from the ability to adaptively select shared memory or TCP channel and further applying optimizations for the chosen channel. To evaluate NVMe-oAF, we co-design Intel's SPDK library with our designs and show up to 7.1x bandwidth improvement and up to 4.2x latency reduction for various workloads over commodity TCP/IP-based Ethernet networks (e.g., 10Gbps, 25Gbps, and 100Gbps). We achieve similar (or sometimes better) performance when compared to NVMe-over-RDMA by avoiding the cumbersome management of RDMA in HPC cloud environments. Finally, we also co-design NVMe-oAF with H5bench to showcase the benefit it brings to HDF5 applications. Our evaluation indicates up to a 7x bandwidth improvement when compared with the network file system (NFS). Arjun Kashyap, Xiaoyi Lu 0001 |
HPDC | 2 |
| 2022 | A Study of Database Performance Sensitivity to Experiment SettingsabstractTo allow performance comparison across different systems, our community has developed multiple benchmarks, such as TPC-C and YCSB, which are widely used. However, despite such effort, interpreting and comparing performance numbers is still a challenging task, because one can tune benchmark parameters, system features, and hardware settings, which can lead to very different system behaviors. Such tuning creates a long-standing question of whether the conclusion of a work can hold under different settings. This work tries to shed light on this question by reproducing 11 works evaluated under TPC-C and YCSB, measuring their performance under a wider range of settings, and investigating the reasons for the change of performance numbers. By doing so, this paper tries to motivate the discussion about whether and how we should address this problem. While this paper does not give a complete solution---this is beyond the scope of a single paper, it proposes concrete suggestions we can take to improve the state of the art. Yang Wang 0009, Miao Yu 0023, Yujie Hui, Xueyuan Ren, Tianxi Li, Xiaoyi Lu 0001 |
Proc. VLDB Endow. | 9 |
| 2021 | DStore: A Fast, Tailless, and Quiescent-Free Object Store for PMEMabstractThe advent of fast, byte-addressable persistent memory (PMEM) has fueled a renaissance in re-evaluating storage system design. Unfortunately, prior work has been unable to provide both consistent and fast performance because they rely on traditional cached or uncached approaches to system design, compromising at least one of the requirements. This paper presents DStore, a fast, tailless, and quiescent-free object store for non-volatile memory. To fulfill all three requirements, we propose a novel two-level approach, called DIPPER, which fully decouples the volatile frontend and persistent backend by leveraging the byte addressability and performance of PMEM. The novelty of our approach is in allowing the frontend and backend to operate independently and in parallel without affecting crash consistency. This not only avoids the need to quiesce the system but also allows for increased concurrency in the frontend through the use of observational equivalency. Using this approach, DStore achieves optimal scalability and low latency without compromising on crash consistency. Evaluation on Intel's Optane DC Persistent Memory Module (DCPMM) demonstrates that DStore can simultaneously provide fast performance, uninterrupted service, and low tail latency. Moreover, DStore can deliver up to 6x lower tail latency service level objectives (SLO) and up to 5x higher throughput SLO compared to state-of-the-art PMEM optimized systems. Shashank Gugnani, Xiaoyi Lu 0001 |
HPDC | 2 |
| 2021 | Characterizing and Accelerating End-to-End EdgeAI Inference Systems for Object Detection Applications
Yujie Hui, Jeffrey Lien, Xiaoyi Lu 0001 |
SEC | 3 |
| 2021 | NVMe-CR: A Scalable Ephemeral Storage Runtime for Checkpoint/Restart with NVMe-over-FabricsabstractEmerging SSDs with NVMe-over-Fabrics (NVMf) support provide new opportunities to significantly improve the performance of IO-intensive HPC applications. However, state-of-the-art parallel filesystems can not extract the best possible performance from fast NVMe SSDs and are not designed for latency-critical ephemeral IO tasks, such as checkpoint/restart. In this paper, we propose a powerful abstraction called microfs to peel away unnecessary software layers and eliminate namespace coordination. Building upon this abstraction, we present the design of NVMe-CR, a scalable ephemeral storage runtime for clusters with disaggregated compute and storage. NVMe-CR proposes techniques like metadata provenance, log record coalescing, and logically isolated shared device access, built around the microfs abstraction, to reduce the overhead of writing millions of concurrent checkpoint files. NVMe-CR utilizes high-density allflash arrays accessible via NVMf to absorb bursty checkpoint IO and increase the progress rates of applications obliviously. Using the ECP CoMD application as a use case, results show that our runtime can achieve near perfect (> 0.96) efficiency at 448 processes and reduce checkpoint overhead by as much as 2x compared to state-of-the-art storage systems. Shashank Gugnani, Tianxi Li, Xiaoyi Lu 0001 |
IPDPS | 3 |
| 2021 | HatRPC: hint-accelerated thrift RPC over RDMAabstractIn this paper, we propose a novel hint-accelerated Remote Procedure Call (RPC) framework based on Apache Thrift over Remote Direct Memory Access (RDMA) protocols, called HatRPC. HatRPC proposes a hierarchical hint scheme towards optimizing heterogeneous RPC services and functions. The proposed hint design is composed of service-granularity and function-granularity hints for achieving varied optimization goals and reducing design space for further optimizing the underneath RDMA communication engine. We co-design a key-value store called HatKV with HatRPC and LMDB. The effectiveness and efficiency of HatRPC are validated and evaluated with our proposed Apache Thrift Benchmarks (ATB), YCSB, and TPC-H workloads. Performance evaluations show that the proposed HatRPC approach can deliver up to 55% performance improvement for ATB benchmarks and up to 1.51X speedup for TPC-H queries compared with vanilla Thrift over IPoIB. In addition, the co-designed HatKV can achieve up to 85.5% improvement for YCSB workloads. Tianxi Li, Xiaoyi Lu 0001 |
SC | 3 |
| 2020 | RDMP-KV: designing remote direct memory persistence based key-value stores with PMEMabstractByte-addressable persistent memory (PMEM) can be directly manipulated by Remote Direct Memory Access (RDMA) capable networks. However, existing studies to combine RDMA and PMEM can not deliver the desired performance due to their PMEM-oblivious communication protocols. In this paper, we propose novel PMEM-aware RDMA-based communication protocols for persistent key-value stores, referred to as Remote Direct Memory Persistence based Key-Value stores (RDMPKV). RDMP-KV employs a hybrid `server-reply/server-bypass' approach to `durably' store individual key-value objects on PMEM-equipped servers. RDMP-KV's runtime can easily adapt to existing (server-assisted durability) and emerging (appliance durability) RDMA-capable interconnects, while ensuring server scalability through a lightweight consistency scheme. Performance evaluations show that RDMP-KV can improve the server-side performance with different persistent key-value storage architectures by up to 22x, as compared with PMEM-oblivious RDMA-`Server-Reply' protocols. Our evaluations also show that RDMP-KV outperforms a distributed PMEM-based filesystem by up to 65% and a recent RDMA-to-PMEM framework by up to 71%. Tianxi Li, Dipti Shankar, Shashank Gugnani, Xiaoyi Lu 0001 |
SC | 4 |
| 2020 | INEC: fast and coherent in-network erasure codingabstractErasure coding (EC) is a promising fault tolerance scheme that has been applied to many well-known distributed storage systems. The capability of Coherent EC Calculation and Networking on modern SmartNICs has demonstrated that EC will be an essential feature of in-network computing. In this paper, we propose a set of coherent in-network EC primitives, named INEC. Our analyses based on the proposed α-β performance model demonstrate that INEC primitives can enable different kinds of EC schemes to fully leverage the EC offload capability on modern SmartNICs. We implement INEC on commodity RDMA NICs and integrate it into five state-of-the-art EC schemes. Our experiments show that INEC primitives significantly reduce 50th, 95th, and 99thpercentile latencies, and accelerate the end-to-end throughput, write, and degraded read performance of the key-value store co-designed with INEC by up to 99.57%, 47.30%, and 49.55%, respectively. Xiaoyi Lu 0001 |
SC | 2 |
| 2020 | CirroData: Yet Another SQL-on-Hadoop Data Analytics Engine with High Performance
Zheng-Hao Jin, Ying-Xin Hu, Li Zha, Xiaoyi Lu 0001 |
J. Comput. Sci. Technol. | 5 |
| 2020 | Understanding the Idiosyncrasies of Real Persistent MemoryabstractHigh capacity persistent memory (PMEM) is finally commercially available in the form of Intel's Optane DC Persistent Memory Module (DCPMM). Researchers have raced to evaluate and understand the performance of DCPMM itself as well as systems and applications designed to leverage PMEM resulting from over a decade of research. Early evaluations of DCPMM show that its behavior is more nuanced and idiosyncratic than previously thought. Several assumptions made about its performance that guided the design of PMEM-enabled systems have been shown to be incorrect. Unfortunately, several peculiar performance characteristics of DCPMM are related to the memory technology (3D-XPoint) used and its internal architecture. It is expected that other technologies (such as STT-RAM, memristor, ReRAM, NVDIMM), with highly variable characteristics, will be commercially shipped as PMEM in the near future. Current evaluation studies fail to understand and categorize the idiosyncratic behavior of PMEM; i.e., how do the peculiarities of DCPMM related to other classes of PMEM. Clearly, there is a need for a study which can guide the design of systems and is agnostic to PMEM technology and internal architecture. In this paper, we first list and categorize the idiosyncratic behavior of PMEM by performing targeted experiments with our proposed PMIdioBench benchmark suite on a real DCPMM platform. Next, we conduct detailed studies to guide the design of storage systems, considering generic PMEM characteristics. The first study guides data placement on NUMA systems with PMEM while the second study guides the design of lock-free data structures, for both eADR- and ADR-enabled PMEM systems. Our results are often counter-intuitive and highlight the challenges of system design with PMEM. Shashank Gugnani, Arjun Kashyap, Xiaoyi Lu 0001 |
Proc. VLDB Endow. | 3 |
| 2019 | SCOR-KV: SIMD-Aware Client-Centric and Optimistic RDMA-Based Key-Value Store for Emerging CPU ArchitecturesabstractModern distributed key-value store-based applications rely on bulk-read operations like 'Multi-Get' (MGet) to accelerate their data serving phase. While state-of-the-art database systems employ SIMD-based techniques to optimize data-parallel operations on their in-memory structures, such as hash-tables, they have not been adapted into high-performance RDMA-accelerated key-value (KV) stores. In this paper, we present a holistic approach to designing high-performance SIMD-aware KV stores for emerging multi-core CPU architectures. Towards this, we first perform an in-depth study of the opportunities and challenges involved in leveraging AVX-512 vectorization-based parallel hash table designs with a state-of-the-art high-performance key-value store like RDMA-Memcached. Based on this, we propose a SIMD-Aware Client-Centric and Optimistic RDMA-based Key-Value Store, SCOR-KV, that optimally exploits 'RDMA+SIMD' to accelerate read-heavy MGet operations. SCOR-KV presents an SIMD-conscious KV store friendly hash table layout, that leverages the vertically vectorized N-way cuckoo hash table design with optimistic KV pair lookup schemes. To complement this, we propose RDMA-optimized SIMD-aware MGet communication protocols that offload the server-side pre-/post-processing overheads to the client, while enabling optimal end-to-end performance. Our performance evaluations over the latest Intel Skylake CPUs and IB EDR interconnects show that our proposed SCOR-KV can achieve up to 3.7-8.6x improvement in server-side Get throughput. Through our SIMD-aware RDMA schemes, SCOR-KV can also improve Multi-Get latencies for read-heavy YCSB workloads by about 2.2x, as compared to the RDMA-Memcached design running over the state-of-the-art CPU-optimized MemC3 hash table design. Dipti Shankar, Xiaoyi Lu 0001, Dhabaleswar K. Panda 0001 |
HiPC | 2 |
| 2019 | UMR-EC: A Unified and Multi-Rail Erasure Coding Library for High-Performance Distributed Storage SystemsabstractDistributed storage systems typically need data to be stored redundantly to guarantee data durability and reliability. While the conventional approach towards this objective is to store multiple replicas, today's unprecedented data growth rates encourage modern distributed storage systems to employ Erasure Coding (EC) techniques, which can achieve better storage efficiency. Various hardware-based EC schemes have been proposed in the community to leverage the advanced compute capabilities on modern data center and cloud environments. Currently, there is no unified and easy way for distributed storage systems to fully exploit multiple devices such as CPUs, GPUs, and network devices (i.e., multi-rail support) to perform EC operations in parallel; thus, leading to the under-utilization of the available compute power. In this paper, we first introduce an analytical model to analyze the design scope of efficient EC schemes in distributed storage systems. Guided by the performance model, we propose UMR-EC, a Unified and Multi-Rail Erasure Coding library that can fully exploit heterogeneous EC coders. Our proposed interface is complemented by asynchronous semantics with optimized metadata-free scheme and EC rate-aware task scheduling that can enable a highly-efficient I/O pipeline. To show the benefits and effectiveness of UMR-EC, we re-design HDFS 3.x write/read pipelines based on the guidelines observed in the proposed performance model. Our performance evaluations show that our proposed designs can outperform the write performance of replication schemes and the default HDFS EC coder by 3.7x - 6.1x and 2.4x - 3.3x, respectively, and can improve the performance of read with failure recoveries up to 5.1x compared with the default HDFS EC coder. Compared with the fastest available CPU coder (i.e., ISA-L), our proposed designs have an improvement of up to 66.0% and 19.4% for write and read with failure recoveries, respectively. Xiaoyi Lu 0001, Dipti Shankar, Dhabaleswar K. Panda 0001 |
HPDC | 2 |
| 2019 | C-GDR: High-Performance Container-Aware GPUDirect MPI Communication Schemes on RDMA NetworksabstractIn recent years, GPU-based platforms have received significant success for parallel applications. In addition to highly optimized computation kernels on GPUs, the cost of data movement on GPU clusters plays critical roles in delivering high performance for end applications. Many recent studies have been proposed to optimize the performance of GPUor CUDA-aware communication runtimes and these designs have been widely adopted in the emerging GPU-based applications. These studies mainly focus on improving the communication performance on native environments, i.e., physical machines, however GPU-based communication schemes on cloud environments are not well studied yet. This paper first investigates the performance characteristics of state-of-the-art GPU-based communication schemes on both native and container-based environments, which show a significant demand to design high-performance container-aware communication schemes in GPU-enabled runtimes to deliver near-native performance for end applications on clouds. Next, we propose the C-GDR approach to design high-performance Container-aware GPUDirect communication schemes on RDMA networks. C-GDR allows communication runtimes to successfully detect process locality, GPU residency, NUMA, architecture information, and communication pattern to enable intelligent and dynamic selection of the best communication and data movement schemes on GPU-enabled clouds. We have integrated C-GDR with the MVAPICH2 library. Our evaluations show that MVAPICH2 with C-GDR has clear performance benefits on container-based cloud environments, compared to default MVAPICH2-GDR and Open MPI. For instance, our proposed CGDR can outperform default MVAPICH2-GDR schemes by up to 66% on micro-benchmarks and up to 26% on HPC applications over a container-based environment. Jie Zhang 0045, Xiaoyi Lu 0001, Ching-Hsiang Chu, Dhabaleswar K. Panda 0001 |
IPDPS | 2 |
| 2019 | TriEC: tripartite graph based erasure coding NIC offloadabstractErasure Coding (EC) NIC offload is a promising technology for designing next-generation distributed storage systems. However, this paper has identified three major limitations of current-generation EC NIC offload schemes on modern SmartNICs. Thus, this paper proposes a new EC NIC offload paradigm based on the tripartite graph model, namely TriEC. TriEC supports both encode-and-send and receive-and-decode operations efficiently. Through theorem-based proofs, co-designs with memcached (i.e., TriEC-Cache), and extensive experiments, we show that TriEC is correct and can deliver better performance than the state-of-the-art EC NIC offload schemes (i.e., BiEC). Benchmark evaluations demonstrate that TriEC outperforms BiEC by up to 1.82x and 2.33x for encoding and recovering, respectively. With extended YCSB workloads, TriEC reduces the average write latency by up to 23.2% and the recovery time by up to 37.8%. TriEC outperforms BiEC by 1.32x for a full-node recovery with 8 million records. Xiaoyi Lu 0001 |
SC | 2 |
| 2019 | Performance analysis of deep learning workloads using roofline trajectories
M. Haseeb Javed, Khaled Z. Ibrahim, Xiaoyi Lu 0001 |
CCF Trans. High Perform. Comput. | 3 |
| 2019 | Exploiting Hardware Multicast and GPUDirect RDMA for Efficient BroadcastabstractBroadcast is a widely used operation in many streaming and deep learning applications to disseminate large amounts of data on emerging heterogeneous High-Performance Computing (HPC) systems. However, traditional broadcast schemes do not fully utilize hardware features for Graphics Processing Unit (GPU)-based applications. In this paper, a model-oriented analysis is presented to identify performance bottlenecks of existing broadcast schemes on GPU clusters. Next, streaming-based broadcast schemes are proposed to exploit InfiniBand hardware multicast (IB-MCAST) and NVIDIA GPUDirect technology for efficient message transmission. The proposed designs are evaluated in the context of using Message Passing Interface (MPI) based benchmarks and applications. The experimental results indicate improved scalability and up to 82 percent reduction of latency compared to the state-of-the-art solutions in the benchmark-level evaluation. Furthermore, compared to the state-of-the-art, the proposed design yields stable higher throughput for a synthetic streaming workload, and 1.3x faster training time for a deep learning framework. Ching-Hsiang Chu, Xiaoyi Lu 0001, Ammar Ahmad Awan, Hari Subramoni, Bracy Elton, Dhabaleswar K. Panda 0001 |
IEEE Trans. Parallel Distributed Syst. | 2 |
| 2018 | Spark-uDAPL: Cost-Saving Big Data Analytics on Microsoft Azure Cloud with RDMA Networks*abstractEfficient Big Data analytics on Cloud Computing systems is still full of challenges. One of the biggest hurdles is the unsatisfactory performance offered by underlying virtualized I/O devices such as networks. To address this issue, the modern cloud resource providers (e.g., Microsoft Azure) have deployed high-performance networks, such as Remote Direct Memory Access (RDMA) capable networks in their clouds. However, in this paper, we find that by far, the RDMA networks on Microsoft Azure cannot support either IPoIB or native standard Verbs-based RDMA protocols. Instead, applications need to use the uDAPL (i.e., user Direct Access Programming Library) interface to enable RDMA communication on Azure Cloud, which makes impossible for modern Big Data stacks to leverage these high-performance networks as none of them can support the uDAPL interface yet. To address this issue, we first design an efficient uDAPL-based communication library with the best combinations of uDAPL communication operations. Then, we adapt the designed uDAPL library into the Hadoop RPC ping-pong message passing engine and the Spark Shuffle engine for bulk data transferring. Through our designs, we can improve the performance of Big Data analytics workloads with Hadoop RPC and Spark on RDMA-enabled Azure VMs by up to 90% and 82%, respectively, and save users’ cloud resource renting cost by 4.24x. To the best of our knowledge, this is the first work to design a uDAPL-based RDMA communication engine for Big Data analytics stacks (e.g., Spark). Xiaoyi Lu 0001, Dipti Shankar, Dhabaleswar K. Panda 0001 |
IEEE BigData | 1 |
| 2018 | High-Performance Multi-Rail Erasure Coding Library over Modern Data Center Architectures: Early ExperiencesabstractVarious hardware-based Erasure Coding (EC) schemes have been proposed [5, 6, 8, 12-14] to leverage the advanced compute capabilities on modern data centers. Currently, there is no unified and easy way for distributed storage systems to fully exploit multiple devices such as CPUs, GPUs, and network devices (i.e., multi-rail support) to perform EC operations in parallel. In this paper, we validate that it is time to design an unified library to efficiently exploit heterogeneous EC coders. HDFS co-designed with our proposed library outperforms the write performance of replication scheme and the default HDFS EC coder by 2.7x - 6.1x and 2.4x - 3.3x, respectively, and improves the performance of read with failure recoveries by up to 2.6x and 5.1x compared to the replication scheme and the default HDFS EC coder, respectively. Xiaoyi Lu 0001, Dipti Shankar, Dhabaleswar K. Panda 0001 |
SoCC | 2 |
| 2018 | Cutting the Tail: Designing High Performance Message Brokers to Reduce Tail Latencies in Stream ProcessingabstractOver the last decade, organizations have become heavily reliant on providing near-instantaneous insights to the end user based on vast amounts of data collected from various sources in real-time. In order to accomplish this task, a stream processing pipeline is constructed, which in its most basic form, consists of a Stream Processing Engine (SPE) and a Message Broker (MB). The SPE is responsible for performing actual computations on the data and providing insights from it. MB, on the other hand, acts as an intermediate queue to which data is written by ephemeral sources and then fetched by the SPE to perform computations on. Due to the inherent real-time nature of such a pipeline, low latency is a highly desirable feature for them. Thus, several existing research works in the community focus on improving latency and throughput of the streaming pipeline. However, there is a dearth of studies optimizing the tail latencies of such pipelines. Moreover, the root cause of this high tail latency is still vague. In this paper, we propose a model-based approach to analyze in-depth the reasons behind high tail latency in streaming systems such as Apache Kafka. Having found the MB to be a major contributor of messages with high tail latencies in a streaming pipeline, we design and implement an RDMA-enhanced high-performance MB, called Frieda, with the higher goal of accelerating any arbitrary stream processing pipeline regardless of the SPE used. Our experiments show a reduction of up to 98% in 99.9th percentile latency for microbenchmarks and up to 31% for full-fledged stream processing pipeline constructed using Yahoo! Streaming Benchmark. M. Haseeb Javed, Xiaoyi Lu 0001, Dhabaleswar K. Panda 0001 |
CLUSTER | 2 |
| 2018 | OC-DNN: Exploiting Advanced Unified Memory Capabilities in CUDA 9 and Volta GPUs for Out-of-Core DNN TrainingabstractExisting frameworks cannot train large DNNs that do not fit the GPU memory without explicit memory management schemes. In this paper, we propose OC-DNN - a novel Out-of-Core DNN training framework that exploits new Unified Memory features along with new hardware mechanisms in Pascal and Volta GPUs. OC-DNN has two major design components — 1) OC-Caffe; an enhanced version of Caffe that exploits innovative UM features like asynchronous prefetching, managed page-migration, exploitation of GPU-based page faults, and the cudaMemAdvise interface to enable efficient out-of-core training for very large DNNs, and 2) an interception library to transpar-ently leverage these cutting-edge features for other frameworks. We provide a comprehensive performance characterization of our designs. OC-Caffe provides comparable performance (to Caffe) for regular DNNs. OC-Caffe-Opt is up to 1.9X faster than OC-Caffe-Naive and up to 5X faster than optimized CPU-based training for out-of-core workloads. OC-Caffe also allows scale-up (DGX-1) and scale-out on multi-GPU clusters. Ammar Ahmad Awan, Ching-Hsiang Chu, Hari Subramoni, Xiaoyi Lu 0001, Dhabaleswar K. Panda 0001 |
HiPC | 4 |
| 2018 | Accelerating TensorFlow with Adaptive RDMA-Based gRPCabstractGoogle's TensorFlow is one of the most popular Deep Learning frameworks nowadays. Distributed TensorFlow supports various channels to efficiently transfer tensors, such as gRPC over TCP/IP, gRPC+Verbs, and gRPC+MPI. At present, the community lacks a thorough characterization of distributed TensorFlow communication channels. This is critical because high-performance Deep Learning with TensorFlow needs an efficient communication runtime. Thus, we conduct a thorough analysis of the communication characteristics of distributed TensorFlow. Our studies show that none of the existing channels in TensorFlow can support adaptive and efficient communication for Deep Learning workloads with different message sizes. Moreover, the community needs to maintain these different channels while the users are also expected to tune these channels to get the desired performance. Therefore, this paper proposes a unified approach to have a single gRPC runtime (i.e., AR-gRPC) in TensorFlow with Adaptive and efficient RDMA protocols. In AR-gRPC, we propose designs such as hybrid communication protocols, message pipelining and coalescing, zero-copy transmission etc. to make our runtime be adaptive to different message sizes for Deep Learning workloads. Our performance evaluations show that AR-gRPC can significantly speedup gRPC performance by up to 4.1x and 2.3x compared to the default gRPC design on IPoIB and another RDMA-based gRPC design in the community. Comet supercomputer shows that AR-gRPC design can reduce the Point-to-Point latency by up to 75% compared to the default gRPC design. By integrating our AR-gRPC with TensorFlow, we can achieve up to 3x distributed training speedup over default gRPC-IPoIB based TensorFlow. Rajarshi Biswas, Xiaoyi Lu 0001, Dhabaleswar K. Panda 0001 |
HiPC | 2 |
| 2018 | Multi-Threading and Lock-Free MPI RMA Based Graph Processing on KNL and POWER ArchitecturesabstractIntel Knights Landing (KNL) and IBM POWER architectures are becoming widely deployed on modern supercomputing systems due to its powerful components. MPI Remote Memory Access (RMA) model that provides one-sided communication semantics has been seen as an attractive approach for developing High-Performance Data Analytics (HPDA) applications such as graph processing with irregular communication characteristics. To take advantage of a large number of hardware threads offered by KNL and POWER, HPDA applications and MPI RMA runtime need to be re-designed to get optimal performance. In this paper, we propose multi-threading and lock-free designs in the MPI runtime as well as Graph500 application on KNL and POWER architectures. At the micro-bench level, our proposed runtime-level designs are able to reduce the latency of uni-directional MPI_Put and MPI_Get by up to 3X compared to IntelMPI and Spectrum MPI. At the application level, with 1,024 processes on 32 KNL nodes, our proposed design could outperform IntelMPI library by 32%. With 512 processes on eight POWER nodes, our proposed design could outperform Spectrum MPI library by 19%. To the best of our knowledge, this is the first paper to design and evaluate MPI RMA-based graph processing applications on KNL and POWER architectures. Xiaoyi Lu 0001, Hari Subramoni, Dhabaleswar K. Panda 0001 |
EuroMPI | 2 |
| 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. | 3 |
| 2018 | Networking and communication challenges for post-exascale systemsabstractWith the significant advancement in emerging processor, memory, and networking technologies, exascale systems will become available in the next few years (2020–2022). As the exascale systems begin to be deployed and used, there will be a continuous demand to run next-generation applications with finer granularity, finer time-steps, and increased data sizes. Based on historical trends, next-generation applications will require postexascale systems during 2025–2035. In this study, we focus on the networking and communication challenges for post-exascale systems. Firstly, we present an envisioned architecture for post-exascale systems. Secondly, the challenges are summarized from different perspectives: heterogeneous networking technologies, high-performance communication and synchronization protocols, integrated support with accelerators and field-programmable gate arrays, fault-tolerance and quality-of-service support, energy-aware communication schemes and protocols, softwaredefined networking, and scalable communication protocols with heterogeneous memory and storage. Thirdly, we present the challenges in designing efficient programming model support for high-performance computing, big data, and deep learning on these systems. Finally, we emphasize the critical need for co-designing runtime with upper layers on these systems to achieve the maximum performance and scalability. Dhabaleswar K. Panda 0001, Xiaoyi Lu 0001, Hari Subramoni |
Frontiers Inf. Technol. Electron. Eng. | 2 |
| 2017 | Characterization of Big Data Stream Processing Pipeline: A Case Study using Flink and KafkaabstractIn recent years there has been a surge in applications focusing on streaming data to generate insights in real-time. Both academia, as well as industry, have tried to address this use case by developing a variety of Stream Processing Engines (SPEs) with a diverse feature set. On the other hand, Big Data applications have started to make use of High-Performance Computing (HPC) which possess superior memory, I/O, and networking resources compared to typical Big Data clusters. Recent studies evaluating the performance of SPEs have focused on commodity clusters. However, exhaustive studies need to be performed to profile individual stages of a stream processing pipeline and how best to optimize each of these stages to best leverage the resources provided by HPC clusters. To address this issue, we profile the performance of a big data streaming pipeline using Apache Flink as the SPE and Apache Kafka as the intermediate message queue. We break the streaming pipeline into two distinct phases and evaluate percentile latencies for two different networks, namely 40GbE and InfiniBand EDR (100Gbps), to determine if a typical streaming application is network intensive enough to benefit from a faster interconnect. Moreover, we explore whether the volume of input data stream has any effect on the latency characteristics of the streaming pipeline, and if so how does it compare for different stages in the streaming pipeline and different network interconnects. Our experiments show an increase of over 10x in 98 percentile latency when input stream volume is increased from 128MB/s to 256MB/s. Moreover, we find the intermediate stages of the stream pipeline to be a significant contributor to the overall latency of the system. M. Haseeb Javed, Xiaoyi Lu 0001, Dhabaleswar K. Panda 0001 |
BDCAT | 2 |
| 2017 | Characterizing and accelerating indexing techniques on distributed ordered tablesabstractIn recent years, most Web 2.0/3.0 applications have been built on top of distributed systems which allow data to be modeled as Distributed Ordered Tables (DOTs) such as Apache HBase. To analyze the stored data, SQL-like range queries over a DOT are fundamental requirements. However, range queries over existing DOT implementations are highly inefficient. Several secondary index techniques have been proposed to alleviate this issue, but they introduce additional overhead while creating and updating the index. Moreover, index techniques introduce several additional challenges for DOTs, particularly, network communication and thread models for concurrent request processing. In this paper, we first characterize the performance of index techniques on DOTs from a networking perspective. We then propose an RDMA-based high-performance communication framework which uses HBase as the underlying DOT implementation to accelerate these techniques. We propose several thread models for our RDMA-based design and compare their performance. We design a parallel insert operation to reduce index creation overhead. We also design several benchmarks to evaluate DOT-based systems. Experimental evaluations with state-of-the-art index techniques (CCIndex and Apache Phoenix) show that our design can reduce the insert overhead for secondary indices to just 23%. Evaluation with TPC-H queries demonstrates an increase in query throughput by up to 2x, while application evaluation with real-world workloads and data (100M records) provided by AdMaster Inc. show up to 35% reduction in execution time. Shashank Gugnani, Xiaoyi Lu 0001, Houliang Qi, Li Zha, Dhabaleswar K. Panda 0001 |
IEEE BigData | 2 |
| 2017 | Performance characterization and acceleration of big data workloads on OpenPOWER systemabstractIBM's POWER processor has been advocated as the high-performance architecture designed for processing Big Data workloads. With the collaborations through the OpenPOWER Foundation, more and more innovations for POWER architecture are emerging to solve Big Data challenges. For example, with the cooperation between IBM and Mellanox, the latest generation of Remote Direct Memory Access (RDMA) capable InfiniBand network can deliver tremendous performance on POWER processors. On the other hand, many RDMA-based designs and optimizations recently have been proposed in the community for accelerating big data processing systems (such as Apache Hadoop and Spark). However, these studies mostly focus on achieving higher performance over Intel Xeon or other x86 architectures. As OpenPOWER systems are getting momentum, we set out to answer the question how much can the RDMA-based communication runtime benefit Big Data processing middleware running over OpenPOWER systems as compared to the default TCP/IP-based designs. To answer this question, this paper first presents an extensive performance characterization on RDMA-based Hadoop RPC engine over OpenPOWER system. We further propose new designs to enable efficient CPU affinity policies and architecture-aware tuning in the RDMA-based communication engine for Hadoop and Spark. With these various accelerations, our performance evaluation shows that our proposed designs can achieve up to 2.73X performance improvement for Hadoop RPC benchmark as compared to default Hadoop running with IP-over-IB protocol on OpenPOWER systems. In addition, our proposed design can gain up to 29.37% performance improvement for Hadoop and Spark workloads as compared to the default RDMA designs running on an OpenPOWER cluster. Xiaoyi Lu 0001, Dipti Shankar, Dhabaleswar K. Panda 0001 |
IEEE BigData | 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 | 3 |
| 2017 | Swift-X: Accelerating OpenStack Swift with RDMA for Building an Efficient HPC CloudabstractRunning Big Data applications in the cloud has become extremely popular in recent times. To enable the storage of data for these applications, cloud-based distributed storage solutions are a must. OpenStack Swift is an object storage service which is widely used for such purposes. Swift is one of the main components of the OpenStack software package. Although Swift has become extremely popular in recent times, its proxy server based design limits the overall throughput and scalability of the cluster. Swift is based on the traditional TCP/IP sockets based communication which has known performance issues such as context-switch and buffer copies for each message transfer. Modern high-performance interconnects such as InfiniBand and RoCE offer advanced features such as RDMA and provide high bandwidth and low latency communication. In this paper, we propose two new designs to improve the performance and scalability of Swift. We propose changes to the Swift architecture and operation design. We propose high-performance implementations of network communication and I/O modules based on RDMA to provide the fastest possible object transfer. In addition, we use efficient hashing algorithms to accelerate object verification in Swift. Experimental evaluations with microbenchmarks, Swift stack benchmark (ssbench), and synthetic application workloads reveal up to 2x and 7.3x performance improvement with our two proposed designs for put and get operations. To the best of our knowledge, this is the first work towards accelerating OpenStack Swift with RDMA over high-performance interconnects in the literature. Shashank Gugnani, Xiaoyi Lu 0001, Dhabaleswar K. Panda 0001 |
CCGrid | 2 |
| 2017 | A Scalable Network-Based Performance Analysis Tool for MPI on Large-Scale HPC SystemsabstractStudying the interaction among applications, MPI runtimes, and the fabric they run on is critical to understanding application performance. There exists no high-performance and scalable tool that enables understanding this interplay on modern multi-petaflop systems. Designing such a tool is non-trivial and involves multiple components including 1) data profiling/collection from network/MPI library, 2) storing and, 3) rendering the data. Furthermore, achieving this with minimal overhead and scalability is a challenging task. We take up this challenge and propose a high-performance and scalable network-based performance analysis tool for MPI libraries operating on modern networks like InfiniBand and Omni-Path. Our designs facilitate caching and pre-rendering, allowing a cluster with 6,541 nodes, 764 switches and, 16,893 network links renders in just 30 seconds - a 44X speed up over non-prerendered solutions. The proposed lock-free and optimized memory-backed storage design enables the tool to handle over a quarter million inserts into the database every 45 seconds (data from 27,504 switch ports and 104,656 MPI processes). The tool has been successfully deployed and validated on HPC systems at OSC and on Comet at SDSC. Hari Subramoni, Xiaoyi Lu 0001, Dhabaleswar K. Panda 0001 |
CLUSTER | 2 |
| 2017 | MPI-LiFE: Designing High-Performance Linear Fascicle Evaluation of Brain Connectome with MPIabstractIn this paper, we combine high-performance computing science with computational neuroscience methods to show how to speed-up cutting-edge methods for mapping and evaluation of the large-scale network of brain connections. More specifically, we use a recent factorization method of the Linear Fascicle Evaluation model (i.e., LiFE [1], [2]) that allows for statistical evaluation of brain connectomes. The method called ENCODE [3], [4] uses a Sparse Tucker Decomposition approach to represent the LiFE model. We show that we can implement the optimization step of the ENCODE method using MPI and OpenMP programming paradigms. Our approach involves the parallelization of the multiplication step of the ENCODE method. We model our design theoretically and demonstrate empirically that the design can be used to identify optimal configurations for the LiFE model optimization via ENCODE method on different hardware platforms. In addition, we co-design the MPI runtime with the LiFE model to achieve profound speed-ups. Extensive evaluation of our designs on multiple clusters corroborates our theoretical model. We show that on a single node on TACC Stampede2, we can achieve speed-ups of up to 8.7x as compared to the original approach. Shashank Gugnani, Xiaoyi Lu 0001, Franco Pestilli, Cesar F. Caiafa, Dhabaleswar K. Panda 0001 |
HiPC | 2 |
| 2017 | Designing Registration Caching Free High-Performance MPI Library with Implicit On-Demand Paging (ODP) of InfiniBandabstractModern high-performance communication runtime systems have taken advantage of advanced features on highperformance networks (e.g. InfiniBand) to deliver optimal performance. High-performance communication over InfiniBand typically requires the communication buffers to be registered first. However, buffer registration and deregistration are costly operations, which leads to performance degradation if they happen frequently. To hide this overhead, many existing communication runtime choose to design a high-performance registration cache to reduce the number of buffer registrations, but such type of designs still need some amount of buffers to be registered and cached, which leads to multiple issues such as performance overhead, high memory consumption for bookkeeping, and code complexity for maintaining the registration cache. To solve these issues, a recently introduced feature for InfiniBand called Implicit OnDemand Paging (ODP) is getting momentum. This feature enables one process to register its complete memory address space for I/O accesses. To fully take advantage of Implicit-ODP, it is critical to fully understand the behavior and benefits of Implicit-ODP on InfiniBand and performance/memory trade-offs it presents. This paper first presents an analysis of the Implicit-ODP feature and studies its basic performance with InfiniBand verbs-level micro-benchmarks. Then, we describe the design tradeoffs with Implicit-ODP and the various optimizations at MPI runtime. We propose and design communication protocols that can leverage the Implicit-ODP feature at the MPI level. The experimental results at the micro-benchmark level and application level show that our proposed design can deliver comparable performance to the existing pin-down scheme, while it does not need registration cache in the MPI runtime. To the best of our knowledge, this is the first work to study and analyze the Implicit-ODP feature and design a registration caching free MPI library with it. Xiaoyi Lu 0001, Hari Subramoni, Dhabaleswar K. Panda 0001 |
HiPC | 2 |
| 2017 | High-Performance and Resilient Key-Value Store with Online Erasure Coding for Big Data WorkloadsabstractDistributed key-value store-based caching solutions are being increasingly used to accelerate Big Data applications on modern HPC clusters. This has necessitated incorporating fault-tolerance capabilities into high-performance key-value stores such as Memcached that are otherwise volatile in nature. In-memory replication is being used as the primary mechanism to ensure resilient data operations. However, this incurs increased network I/O with high remote memory requirements. On the other hand, Erasure Coding is being extensively explored for enabling data resilience, while achieving better storage efficiency. In this paper, we first perform an in-depth modeling-based analysis of the performance trade-offs of In-Memory Replication and Erasure Coding schemes for key-value stores, and explore the possibilities of employing Online Erasure Coding for enabling resilience in high-performance key-value stores for HPC clusters. We then design a non-blocking API-based engine to perform efficient Set/Get operations by overlapping the encoding/decoding involved in enabling Erasure Coding-based resilience with the request/response phases, by leveraging RDMA on high performance interconnects. Performance evaluations show that the proposed designs can outperform synchronous RDMA-based replication by about 2.8x, and can improve YCSB throughput and average read/write latencies by about 1.34x - 2.6x over asynchronous replication for larger key-value pair sizes (>16KB). We also demonstrate its benefits by incorporating it into a hybrid and resilient key-value store-based burst-buffer system over Lustre for accelerating Big Data I/O on HPC clusters. Dipti Shankar, Xiaoyi Lu 0001, Dhabaleswar K. Panda 0001 |
ICDCS | 2 |
| 2017 | Efficient and Scalable Multi-Source Streaming Broadcast on GPU Clusters for Deep LearningabstractBroadcast operations (e.g. MPI_Bcast) have been widely used in deep learning applications to exchange a large amount of data among multiple graphics processing units (GPUs). Recent studies have shown that leveraging the InfiniBand hardware-based multicast (IB-MCAST) protocol can enhance scalability of GPU-based broadcast operations. However, these initial designs with IB-MCAST are not optimized for multi-source broadcast operations with large messages, which is the common communication scenario for deep learning applications. In this paper, we first model existing broadcast schemes and analyze their performance bottlenecks on GPU clusters. Then, we propose a novel broadcast design based on message streaming to better exploit IB-MCAST and NVIDIA GPUDirect RDMA (GDR) technology for efficient large message transfer operation. The proposed design can provide high overlap among multi-source broadcast operations. Experimental results show up to 68% reduction of latency compared to state-of-the-art solutions in a benchmark-level evaluation. The proposed design also shows near-constant latency for a single broadcast operation as a system grows. Furthermore, it yields up to 24% performance improvement in the popular deep learning framework, Microsoft CNTK, which uses multi-source broadcast operations; notably, the performance gains are achieved without modifications to applications. Our model validation shows that the proposed analytical model and experimental results match within a 10% range. Our model also predicts that the proposed design outperforms existing schemes for multi-source broadcast scenarios with increasing numbers of broadcast sources in large-scale GPU clusters. Ching-Hsiang Chu, Xiaoyi Lu 0001, Ammar Ahmad Awan, Hari Subramoni, Jahanzeb Maqbool Hashmi, Bracy Elton, Dhabaleswar K. Panda 0001 |
ICPP | 2 |
| 2017 | High-Performance Virtual Machine Migration Framework for MPI Applications on SR-IOV Enabled InfiniBand ClustersabstractHigh-speed interconnects (e.g. InfiniBand) have been widely deployed on modern HPC clusters. With the emergence of HPC in the cloud, high-speed interconnects have paved their way into the cloud with recently introduced Single Root I/O Virtualization (SR-IOV) technology, which is able to provide efficient sharing of high-speed interconnect resources and achieve near-native I/O performance. However, recent studies have shown that SR-IOV-based virtual networks prevent virtual machine migration, which is an essential virtualization capability towards high availability and resource provisioning. Although several initial solutions have been pro- posed in the literature to solve this problem, our investigations show that there are still many restrictions on these proposed approaches, such as depending on specific network adapters and/or hypervisors, which will limit the usage scope of these solutions on HPC environments. In this paper, we propose a high-performance virtual machine migration framework for MPI applications on SR-IOV enabled InfiniBand clusters. Our proposed method does not need any modification to the hypervisor and InfiniBand drivers and it can efficiently handle virtual machine (VM) migration with SR-IOV IB device. Our evaluation results indicate that the proposed design is able to not only achieve fast VM migration speed but also guarantee the high performance for MPI applications during the migration in the HPC cloud. At the application level, for NPB LU benchmark running inside VM, our proposed design is able to completely hide the migration overhead through the computation and migration overlapping. Furthermore, our proposed design shows good scaling when migrating multiple VMs. Jie Zhang 0045, Xiaoyi Lu 0001, Dhabaleswar K. Panda 0001 |
IPDPS | 2 |
| 2017 | Scalable reduction collectives with data partitioning-based multi-leader designabstractExisting designs for MPI_Allreduce do not take advantage of the vast parallelism available in modern multi-/many-core processors like Intel Xeon/Xeon Phis or the increases in communication throughput and recent advances in high-end features seen with modern interconnects like InfiniBand and Omni-Path. In this paper, we propose a high-performance and scalable Data Partitioning-based Multi-Leader (DPML) solution for MPI_Allreduce that can take advantage of the parallelism offered by multi-/many-core architectures in conjunction with the high throughput and high-end features offered by InfiniBand and Omni-Path to significantly enhance the performance of MPI_Allreduce on modern HPC systems. We also model DPML-based designs to analyze the communication costs theoretically. Microbenchmark level evaluations show that the proposed DPML-based designs are able to deliver up to 3.5 times performance improvement for MPI_Allreduce for multiple HPC systems at scale. At the application-level, up to 35% and 60% improvement is seen in communication for HPCG and miniAMR respectively. Mohammadreza Bayatpour, Sourav Chakraborty 0003, Hari Subramoni, Xiaoyi Lu 0001, Dhabaleswar K. Panda 0001 |
SC | 4 |
| 2017 | Designing Locality and NUMA Aware MPI Runtime for Nested Virtualization based HPC Cloud with SR-IOV Enabled InfiniBandabstractHypervisor-based virtualization solutions reveal good security and isolation, while container-based solutions make applications and workloads more portable and distributed in an effective, standardized and repeatable way. Therefore, nested virtualization based computing environments (e.g., container over virtual machine), which inherit the capabilities from both solutions, are becoming more and more attractive in clouds (e.g., running Docker over Amazon EC2 VMs). Recent studies have shown that running applications in either VMs or containers still has significant overhead, especially for I/O intensive workloads. This motivates us to investigate whether the nested virtualization based solution can be adopted to build high-performance computing (HPC) clouds for running MPI applications efficiently and where the bottlenecks lie. To eliminate performance bottlenecks, we propose a high-performance two-layer locality and NUMA aware MPI library, which is able to dynamically detect co-resident containers inside one VM as well as detect co-resident VM inside one host at MPI runtime. Thus the MPI processes across different containers and VMs can communicate to each other by shared memory or Cross Memory Attach (CMA) channels instead of network channel if they are co-resident. We further propose an enhanced NUMA aware hybrid design to utilize InfiniBand loopback based channel to optimize large message transfer across containers when they are running on different sockets. Performance evaluations show that compared with the performance of the state-of-art (1Layer) design, our proposed enhance-hybrid design can bring up to 184%, 81% and 12% benefit on point-to-point, collective operations, and end applications. Compared with the default performance, our enhanced-hybrid design delivers up to 184%, 85% and 16% performance improvement. Jie Zhang 0045, Xiaoyi Lu 0001, Dhabaleswar K. Panda 0001 |
VEE | 2 |
| 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. | 3 |
| 2016 | Performance characterization of hadoop workloads on SR-IOV-enabled virtualized InfiniBand clustersabstractBig Data Systems are becoming increasingly complex and generally have very high operational costs. Cloud computing offers attractive solutions for managing large scale systems. However, one of the major bottlenecks in VM performance is virtualized I/O. Since Big Data applications and middleware rely heavily on high performance interconnects such as InfiniBand, the performance of virtualized InfiniBand interfaces is vital. Single Root I/O Virtualization (SR-IOV) is a hardware based approach which offers significant performance benefits as compared to software based I/O virtualization. With the increasing adoption of InfiniBand network for cloud computing, it is important to evaluate the performance benefits of SR-IOV for InfiniBand networks; especially to see the performance characteristics of Big Data applications and middleware under different scenarios. We characterize the main performance factors for different workloads through this study (such as map task scheduling, I/O, data replication, etc.). Our experimental evaluations show that the performance difference for a wide set of Big Data benchmarks and applications over SR-IOV with InfiniBand using RDMA-enabled Hadoop as compared to native InfiniBand network is just 5 -- 15%. In addition, with RDMA-enabled Hadoop, we see 20.9 -- 81.6% performance improvement for RDMA as compared to IPoIB. Shashank Gugnani, Xiaoyi Lu 0001, Dhabaleswar K. Panda 0001 |
BDCAT | 2 |
| 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 | 3 |
| 2016 | High-performance design of apache spark with RDMA and its benefits on various workloadsabstractThe in-memory data processing framework, Apache Spark, has been stealing the limelight for low-latency interactive applications, iterative and batch computations. Our early experience study [17] has shown that Apache Spark can be enhanced to leverage advanced features (e.g., RDMA) on high-performance networks (e.g., InfiniBand and RoCE) to improve the performance of shuffle phase. With the fast evolving of the Apache Spark ecosystem, the Spark architecture has been changing a lot. This motivates us to investigate whether the earlier RDMA design can be adapted and further enhanced for the new Apache Spark architecture. We also aim to improve the performance for various Spark workloads (e.g., Batch, Graph, SQL). In this paper, we present a detailed design for high-performance RDMA-based Apache Spark on high-performance networks. We conduct systematic performance evaluations on three modern clusters (Chameleon, SDSC Comet, and an in-house cluster) with cutting-edge InfiniBand technologies, such as latest IB EDR (100 Gbps) network, recently introduced Single Root I/O Virtualization (SR-IOV) technology for IB, etc. The evaluation results show that compared to the default Spark running with IP over InfiniBand (IPoIB), our proposed design can achieve up to 79% performance improvement for Spark RDD operation benchmarks (e.g., GroupBy, SortBy), up to 38% performance improvement for batch workloads (e.g., Sort and TeraSort in Intel HiBench), up to 46% performance improvement for graph processing workloads (e.g., PageRank), up to 32% performance improvement for SQL queries (e.g., Aggregation, Join) on varied scales (up to 1,536 cores) of bare-metal IB clusters. Performance evaluations on SR-IOV enabled IB clusters also show 37% improvement achieved by our RDMA-based design. Our RDMA-based Spark design is implemented as a pluggable module and it does not change any Spark APIs, which means that it can be combined with other existing enhanced designs for Apache Spark and Hadoop proposed in the community. To show this, we further evaluate the performance of a combined version of `RDMA-Spark+RDMA-HDFS' and the numbers show that the combination can achieve the best performance with up to 82% improvement for Intel HiBench Sort and TeraSort on SDSC Comet cluster. Xiaoyi Lu 0001, Dipti Shankar, Shashank Gugnani, Dhabaleswar K. Panda 0001 |
IEEE BigData | 1 |
| 2016 | Boldio: A hybrid and resilient burst-buffer over lustre for accelerating big data I/OabstractThe limitation of local storage space in the HPC environments has placed an unprecedented demand on the performance of the underlying shared parallel file systems. This has necessitated a scalable solution for running Big Data middleware (e.g., Hadoop) on HPC clusters. In this paper, we propose Boldio, a hybrid and resilient key-value store-based Burst-Buffer system Over Lustre for accelerating I/O-intensive Big Data workloads, that can leverage RDMA on high-performance interconnects and storage technologies such as PCIe-/NVMe-SSDs, etc. We demonstrate that Boldio can improve the performance of the I/O phase of Hadoop workloads running on HPC clusters; serving as a light-weight, high-performance, and resilient remote I/O staging layer between the application and Lustre. Performance evaluations show that Boldio can improve the TestDFSIO write performance over Lustre by up to 3x and TestDFSIO read performance by 7x, while reducing the execution time of Hadoop Sort benchmark by up to 30%. We demonstrate that we can significantly improve Hadoop I/O throughput over popular in-memory distributed storage systems such as Alluxio (formerly Tachyon), when high-speed local storage is limited. Dipti Shankar, Xiaoyi Lu 0001, Dhabaleswar K. Panda 0001 |
IEEE BigData | 2 |
| 2016 | Designing Virtualization-Aware and Automatic Topology Detection Schemes for Accelerating Hadoop on SR-IOV-Enabled CloudsabstractHadoop is gaining more and more popularity in virtualized environments because of the flexibility and elasticity offered by cloud-based systems. Hadoop supports topology-awareness through topology-aware designs in all of its major components. However, there exists no service that can automatically detect the underlying network topology in a scalable and efficient manner, and provide this information to the Hadoop framework. Moreover, the topology-aware designs in Hadoop are not optimized for virtualized platforms. In this paper, we propose a new library called Hadoop-Virt, based on RDMA-Hadoop, which provides with an automatic topology detection module and virtualization-aware designs in Hadoop to fully take the advantage of virtualized environments. Our experimental evaluations show that Hadoop-Virt delivers upto 34% better performance in the default execution mode and upto 52.6% better performance in the distributed mode as compared to default RDMA-Hadoop for SR-IOV-enabled virtualized clusters. Shashank Gugnani, Xiaoyi Lu 0001, Dhabaleswar K. Panda 0001 |
CloudCom | 2 |
| 2016 | Impact of HPC Cloud Networking Technologies on Accelerating Hadoop RPC and HBaseabstractThe performance of Hadoop components can be significantly improved by leveraging advanced features such as Remote Direct Memory Access (RDMA) on modern HPC clusters, where high-performance networks like InfiniBand (IB) and RoCE have been deployed widely. With the emergence of high-performance computing in the cloud (HPC Cloud), high-performance networks have paved their way into the cloud with recently introduced Single Root I/O Virtualization (SR-IOV) technology. With these advancements in HPC Cloud networking technologies, it is high time to investigate the design opportunities and impact of networking architectures (different generations of IB, 40GigE, 40G-RoCE) and protocols (TCP/IP, IPoIB, RC, UD, Hybrid) in accelerating Hadoop components over high-performance networks. In this paper, we propose a network architecture and multi-protocol aware Hadoop RPC design, that can take advantage of RC and UD protocols for IB and RoCE. A hybrid transport design with RC and UD is proposed which can deliver memory scalability and performance for Hadoop RPC. We present a comprehensive performance analysis on five bare-metal IB/RoCE clusters and one SR-IOV enabled cluster in the Chameleon Cloud. Our performance evaluations reveal that our proposed designs can achieve up to 12.5x performance improvement for Hadoop RPC over IPoIB. Further, we integrate our RPC engine into Apache HBase, and demonstrate that we can accelerate YCSB workloads by up to 3.6x. Other insightful observations on performance characteristics of different HPC Cloud networking technologies are also shared in this paper. Xiaoyi Lu 0001, Dipti Shankar, Shashank Gugnani, Hari Subramoni, Dhabaleswar K. Panda 0001 |
CloudCom | 1 |
| 2016 | Slurm-V: Extending Slurm for Building Efficient HPC Cloud with SR-IOV and IVShmem
Jie Zhang 0045, Xiaoyi Lu 0001, Sourav Chakraborty 0003, Dhabaleswar K. Panda 0001 |
Euro-Par | 2 |
| 2016 | Mizan-RMA: Accelerating Mizan Graph Processing Framework with MPI RMAabstractThe MPI programming model has been used by several graph processing systems. Out of these systems, only MPI two-sided programming model is used for datamovements. The MPI one-sided programming model that achieves better overlap ofcommunication and computation has been seen as advantageous for applicationswith irregular communication patterns. However, the benefits of MPI one-sidedprogramming model for graph processing systems is still not exploited. In thispaper, we choose the Mizan graph processing framework that uses MPI two-sidedprogramming model and analyze its performance bottlenecks. Based on ouranalysis, we propose Mizan-RMA which takes advantage of MPI one-sidedprogramming model to alleviate these bottlenecks. Experimental evaluations showthat our proposed designs could achieve up to 2.8X improvement compared with thedefault version. To the best of our knowledge, this is the first paper tore-design a graph processing system with MPI one-sided programming model. Xiaoyi Lu 0001, Khaled Hamidouche, Jie Zhang 0045, Dhabaleswar K. Panda 0001 |
HiPC | 2 |
| 2016 | High Performance MPI Library for Container-Based HPC Cloud on InfiniBand ClustersabstractVirtualization technology has grown rapidly over the past few decades. As a lightweight solution, container-based virtualization provides a promising approach to efficiently build HPC clouds. However, our study shows clear performance bottleneck when running MPI jobs on multi-container environments. This motivates us to first analyze the performance bottleneck for MPI jobs running in different container deployment scenarios. To eliminate performance bottleneck, we propose a high performance locality-aware MPI library, which is able to dynamically detect co-resident containers at runtime. Through this design, the MPI processes in co-resident containers can communicate to each other by shared memory and Cross Memory Attach (CMA) channels instead of the network channel. A comprehensive performance study indicates that compared with the default case, our proposed design can significantly improve the communication performance by up to 9X and 86% in terms of MPI point-to-point and collective operations, respectively. The results for applications demonstrate that the locality-aware design can reduce up to 16% of execution time. The evaluation results also show that by the help of locality-aware design, we can achieve near-native performance in container-based HPC cloud with minor overhead. The proposed locality-aware MPI design reveals significant potential to be utilized to efficiently build large scale container-based HPC clouds. Jie Zhang 0045, Xiaoyi Lu 0001, Dhabaleswar K. Panda 0001 |
ICPP | 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 | 3 |
| 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 | 2 |
| 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 | 3 |
| 2016 | Designing MPI library with on-demand paging (ODP) of infiniband: challenges and benefitsabstractExisting InfiniBand drivers require the communication buffers to be pinned in physical memory during communication. Most runtimes leave these buffers pinned until the end of the run. Such situation limits the swappable memory space for applications. To address these concerns, Mellanox has recently introduced the On-Demand Paging (ODP) feature for InfiniBand. With ODP, communication buffers are paged in when they are needed by the HCA and paged out when the OS needs to swap them. This paper presents a thorough analysis on ODP and studies its performance characteristics. With these studies, we propose novel designs of ODP-aware MPI communication protocols. To the best of our knowledge, this is the first work to study and analyze the ODP feature and design an ODP-aware MPI library. Performance evaluations with applications show that ODP-aware designs can reduce the size of pin-down buffers by 11X without performance degradation compared with the pin-down scheme. Khaled Hamidouche, Xiaoyi Lu 0001, Hari Subramoni, Jie Zhang 0045, Dhabaleswar K. Panda 0001 |
SC | 3 |
| 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. | 2 |
| 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 | 3 |
| 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 | 2 |
| 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 | 2 |
| 2015 | MVAPICH2 over OpenStack with SR-IOV: An Efficient Approach to Build HPC CloudsabstractCloud Computing with Virtualization offers attractive flexibility and elasticity to deliver resources by providing a platform for consolidating complex IT resources in a scalable manner. However, efficiently running HPC applications on Cloud Computing systems is still full of challenges. One of the biggest hurdles in building efficient HPC clouds is the unsatisfactory performance offered by underlying virtualized environments, more specifically, virtualized I/O devices. Recently, Single Root I/O Virtualization (SR-IOV) technology has been steadily gaining momentum for high-performance interconnects such as InfiniBand and 10GigE. Due to its near native performance for inter-node communication, many cloud systems such as Amazon EC2 have been using SR-IOV in their production environments. Nevertheless, recent studies have shown that the SR-IOV scheme lacks locality aware communication support, which leads to performance overheads for inter-VM communication within the same physical node. In this paper, we propose an efficient approach to build HPC clouds based on MVAPICH2 over Open Stack with SR-IOV. We first propose an extension for Open Stack Nova system to enable the IV Shmem channel in deployed virtual machines. We further present and discuss our high-performance design of virtual machine aware MVAPICH2 library over Open Stack-based HPC Clouds. Our design can fully take advantage of high-performance SR-IOV communication for inter-node communication as well as Inter-VM Shmem (IVShmem) for intra-node communication. A comprehensive performance evaluation with micro-benchmarks and HPC applications has been conducted on an experimental Open Stack-based HPC cloud and Amazon EC2. The evaluation results on the experimental HPC cloud show that our design and extension can deliver near bare-metal performance for implementing SR-IOV-based HPC clouds with virtualization. Further, compared with the performance on EC2, our experimental HPC cloud can exhibit up to 160X, 65X, 12X improvement potential in terms of point-to-point, collective and application for future HPC clouds. Jie Zhang 0045, Xiaoyi Lu 0001, Mark Daniel Arnold, Dhabaleswar K. Panda 0001 |
CCGRID | 2 |
| 2015 | High Performance MPI Datatype Support with User-Mode Memory Registration: Challenges, Designs, and BenefitsabstractNoncontiguous data communication has been heavily adopted in scientific applications, especially for those written with MPI. Common strategies to handle noncontiguous data, like packing/unpacking, incur significant performance overhead during communication, which could become as a barrier of using MPI derived datatypes. Recently, a novel feature of Mellanox InfiniBand, called User-mode Memory Registration (UMR), has been introduced for noncontiguous data communication. UMR has the potential to support MPI derived datatype communication efficiently without the overhead of packing/unpacking. In this paper, we analyze the UMR feature and study its basic performance with InfiniBand verbs-level micro-benchmarks. With this knowledge, we propose UMR-based schemes to support zero-copy datatype communication at MPI level. We show that a naive integration of UMR with an MPI stack could not bring performance benefits over existing schemes. Thus we propose two schemes -- UMR Pool and UMR Cache -- to enable high performance MPI datatype communication with UMR. To the best of our knowledge, this is the first paper to study, analyze, and design MPI noncontiguous data communication using the UMR feature. We propose and implement UMR-based designs on top of MVAPICH2 library. The experimental results at the microbenchmark level show that the proposed UMR-based design is able to deliver 4X performance improvement in latency for large message vector benchmarks over the packing/unpacking scheme. At the application level, for a 3D stencil communication kernel with MPI derived datatype on 512 processes, the optimized UMR-based design outperforms the packing/unpacking scheme by 27% in execution time. Hari Subramoni, Khaled Hamidouche, Xiaoyi Lu 0001, Dhabaleswar K. Panda 0001 |
CLUSTER | 4 |
| 2015 | High-Performance and Scalable Design of MPI-3 RMA on Xeon Phi Clusters
Khaled Hamidouche, Xiaoyi Lu 0001, Jian Lin 0006, Dhabaleswar K. Panda 0001 |
Euro-Par | 3 |
| 2015 | High Performance OpenSHMEM Strided Communication Support with InfiniBand UMRabstractExchanging data on noncontiguous user buffers has been a dominant communication pattern in many scientific applications. The OpenSHMEM specification introduces a new set of communication routines to support strided data communication. Most high performance implementations of the OpenSHMEM specification support strided data communication by either packing/unpacking or multiple reads/writes based scheme, which incurs significant performance overhead during communication. This performance overhead could prevent application developers from using OpenSHMEM strided data communication routines. Recently, Mellanox has introduced a novel feature, called User-mode Memory Registration (UMR), for noncontiguous data transfer. UMR has the potential to support efficient OpenSHMEM strided data communication. In this paper, we propose UMR-based schemes to support one-sided zero-copy strided data communication for OpenSHMEM. To the best of our knowledge, this is the first paper to design OpenSHMEM strided data communication using the UMR feature. We propose and implement UMR-based designs on top of MVAPICH2-X. Experimental results with shmem iget operation show 3X performance improvement over the multiple reads scheme in default MVAPICH2-X, and 20X performance improvement over the OpenSHMEM reference implementation configured with GASNet. At the application level, for a 3D stencil communication kernel with OpenSHMEM iget routines on 512 processes, the proposed UMR-based design outperforms the multiple reads scheme in default MVAPICH2-X by 20% in total execution time. Khaled Hamidouche, Xiaoyi Lu 0001, Jie Zhang 0045, Jian Lin 0006, Dhabaleswar K. Panda 0001 |
HiPC | 3 |
| 2015 | Accelerating Apache Hive with MPI for Data Warehouse SystemsabstractData warehouse systems, like Apache Hive, have been widely used in the distributed computing field. However, current generation data warehouse systems have not fully embraced High Performance Computing (HPC) technologies even though the trend of converging Big Data and HPC is emerging. For example, in traditional HPC field, Message Passing Interface (MPI) libraries have been optimized for HPC applications during last decades to deliver ultra-high data movement performance. Recent studies, like DataMPI, are extending MPI for Big Data applications to bridge these two fields. This trend motivates us to explore whether MPI can benefit data warehouse systems, such as Apache Hive. In this paper, we propose a novel design to accelerate Apache Hive by utilizing DataMPI. We further optimize the DataMPI engine by introducing enhanced non-blocking communication and parallelism mechanisms for typical Hive workloads based on their communication characteristics. Our design can fully and transparently support Hive workloads like Intel HiBench and TPC-H with high productivity. Performance evaluation with Intel HiBench shows that with the help of light-weight DataMPI library design, efficient job start up and data movement mechanisms, Hive on DataMPI performs 30% faster than Hive on Hadoop averagely. And the experiments on TPC-H with ORCFile show that the performance of Hive on DataMPI can improve 32% averagely and 53% at most more than that of Hive on Hadoop. To the best of our knowledge, Hive on DataMPI is the first attempt to propose a general design for fully supporting and accelerating data warehouse systems with MPI. Lu Chao, Chundian Li, Xiaoyi Lu 0001, Zhiwei Xu 0002 |
ICDCS | 4 |
| 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 | 3 |
| 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 | 2 |
| 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 | 2 |
| 2015 | Accelerating Iterative Big Data Computing Through MPI
Xiaoyi Lu 0001 |
J. Comput. Sci. Technol. | 2 |
| 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 | 2 |
| 2014 | High performance OpenSHMEM for Xeon Phi clusters: Extensions, runtime designs and application co-designabstractIntel Many Integrated Core (MIC) architectures are becoming an integral part of modern supercomputer architectures due to their high compute density and performance per watt. Partitioned Global Address Space (PGAS) programming models, such as OpenSHMEM, provide an attractive approach for developing scientific applications with irregular communication characteristics, by abstracting shared memory address space, along with one-sided communication semantics. However, the current OpenSHMEM standard does not efficiently support heterogeneous memory architectures such as Xeon Phi. Host and Xeon Phi cores have different memory capacities and compute characteristics. But, the global symmetric memory allocation in the current OpenSHMEM standard mandates that same amount of memory be allocated on every process. In this paper, we propose extensions to overcome this restriction and propose high performance runtime-level designs for efficient communication involving Xeon Phi processors. Further, we re-design applications to demonstrate the effectiveness of the proposed designs and extensions. Experimental evaluations indicate 4X to 7X reduction in OpenSHMEM data movement operation latencies, and 6X to 11X improvement in performance for collective operations. Application evaluations in symmetric mode indicate performance improvements of 28% at 1,024 processes. Further, application redesigns using the proposed extensions provide several magnitudes of performance improvement, as compared to the symmetric mode. To the best of our knowledge, this is the first research work that proposes high performance runtime designs for OpenSHMEM on Intel Xeon Phi clusters. Khaled Hamidouche, Xiaoyi Lu 0001, Sreeram Potluri, Jie Zhang 0045, Karen A. Tomko, Dhabaleswar K. Panda 0001 |
CLUSTER | 3 |
| 2014 | Scalable Graph500 design with MPI-3 RMAabstractThe MPI two-sided programming model has been widely used for scientific applications. However, the benefits of MPI one-sided communication are still not well exploited. Recently, MPI-3 Remote Memory Access (RMA) was introduced with several advanced features which provide better performance, programmability, and flexibility over MPI-2 RMA. However, few studies have shown the benefits of using MPI-3 RMA for scientific applications. In this paper, we take advantage of the new features from MPI-3 RMA to re-design a scalable Graph500 benchmark. Our design achieves much better overlap of communication and computation than the default two sided based implementation. The results show that the proposed design can achieve up to 2X improvement compared with the best MPI based implementation running with 4,096 cores. To the best of our knowledge, this is the first paper to re-design a high performance and scalable Graph500 with MPI-3 RMA. Xiaoyi Lu 0001, Sreeram Potluri, Khaled Hamidouche, Karen A. Tomko, Dhabaleswar K. Panda 0001 |
CLUSTER | 2 |
| 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 | 2 |
| 2014 | Can Inter-VM Shmem Benefit MPI Applications on SR-IOV Based Virtualized Infiniband Clusters?
Jie Zhang 0045, Xiaoyi Lu 0001, Rong Shi, Dhabaleswar K. Panda 0001 |
Euro-Par | 2 |
| 2014 | High performance MPI library over SR-IOV enabled infiniband clustersabstractVirtualization has become a central role in HPC Cloud due to easy management and low cost of computation and communication. Recently, Single Root I/O Virtualization (SR-IOV) technology has been introduced for high-performance interconnects such as InfiniBand and can attain near to native performance for inter-node communication. However, the SR-IOV scheme lacks locality aware communication support, which leads to performance overheads for inter-VM communication within a same physical node. To address this issue, this paper first proposes a high performance design of MPI library over SR-IOV enabled InfiniBand clusters by dynamically detecting VM locality and coordinating data movements between SR-IOV and Inter-VM shared memory (IVShmem) channels. Through our proposed design, MPI applications running in virtualized mode can achieve efficient locality-aware communication on SR-IOV enabled InfiniBand clusters. In addition, we optimize communications in IVShmem and SR-IOV channels by analyzing the performance impact of core mechanisms and parameters inside MPI library to deliver better performance in virtual machines. Finally, we conduct comprehensive performance studies by using point-to-point and collective benchmarks, and HPC applications. Experimental evaluations show that our proposed MPI library design can significantly improve the performance for point-to-point and collective operations, and MPI applications with different InfiniBand transport protocols (RC and UD) by up to 158%, 76%, 43%, respectively, compared with SR-IOV. To the best of our knowledge, this is the first study to offer a high performance MPI library that supports efficient locality aware MPI communication over SR-IOV enabled InfiniBand clusters. Jie Zhang 0045, Xiaoyi Lu 0001, Rong Shi, Dhabaleswar K. Panda 0001 |
HiPC | 2 |
| 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 | 2 |
| 2014 | HAND: A Hybrid Approach to Accelerate Non-contiguous Data Movement Using MPI Datatypes on GPU ClustersabstractAn increasing number of MPI applications are being ported to take advantage of the compute power offered by GPUs. Data movement continues to be the major bottleneck on GPU clusters, more so when data is non-contiguous, which is common in scientific applications. The existing techniques of optimizing MPI data type processing, to improve performance of non-contiguous data movement, handle only certain data patterns efficiently while incurring overheads for the others. In this paper, we first propose a set of optimized techniques to handle different MPI data types. Next, we propose a novel framework (HAND) that enables hybrid and adaptive selection among different techniques and tuning to achieve better performance with all data types. Our experimental results using the modified DDTBench suite demonstrate up to a 98% reduction in data type latency. We also apply this data type-aware design on an N-Body particle simulation application. Performance evaluation of this application on a 64 GPU cluster shows that our proposed approach can achieve up to 80% and 54% increase in performance by using struct and indexed data types compared to the existing best design. To the best of our knowledge, this is the first attempt to propose a hybrid and adaptive solution to integrate all existing schemes to optimize arbitrary non-contiguous data movement using MPI data types on GPU clusters. Rong Shi, Xiaoyi Lu 0001, Sreeram Potluri, Khaled Hamidouche, Jie Zhang 0045, Dhabaleswar K. Panda 0001 |
ICPP | 2 |
| 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 | 2 |
| 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 | 2 |
| 2014 | DataMPI: Extending MPI to Hadoop-Like Big Data ComputingabstractMPI has been widely used in High Performance Computing. In contrast, such efficient communication support is lacking in the field of Big Data Computing, where communication is realized by time consuming techniques such as HTTP/RPC. This paper takes a step in bridging these two fields by extending MPI to support Hadoop-like Big Data Computing jobs, where processing and communication of a large number of key-value pair instances are needed through distributed computation models such as MapReduce, Iteration, and Streaming. We abstract the characteristics of key-value communication patterns into a bipartite communication model, which reveals four distinctions from MPI: Dichotomic, Dynamic, Data-centric, and Diversified features. Utilizing this model, we propose the specification of a minimalistic extension to MPI. An open source communication library, DataMPI, is developed to implement this specification. Performance experiments show that DataMPI has significant advantages in performance and flexibility, while maintaining high productivity, scalability, and fault tolerance of Hadoop. Xiaoyi Lu 0001, Li Zha, Zhiwei Xu 0002 |
IPDPS | 1 |
| 2014 | Performance Characterization of Hadoop and Data MPI Based on Amdahl's Second LawabstractAmdahl's second law has been seen as a useful guideline for designing and evaluating balanced computer systems for decades. This law has been mainly used for hardware systems and peak capacities. This paper utilizes Amdahl's second law from a new angle, i.e., Evaluating the influence on systems performance and balance of the application framework software, a key component of big data systems. We compare two big data application framework software systems, Apache Hadoop and Data MPI, with three representative application benchmarks and various data sizes. System monitors and hardware performance counters are used to record the resource utilization, characteristics of instructions execution, memory accesses, and I/O rates. These numbers are used to reveal the three runtime metrics of Amdahl's second law: CPU speed (GIPS), memory capacity (GB), and I/O rate (Gbps). The experiment and evaluation results show that a Data MPI-based big data system has better performance and is more balanced than a Hadoop-based system. Xiaoyi Lu 0001, Zhiwei Xu 0002 |
NAS | 3 |
| 2014 | Initial study of multi-endpoint runtime for MPI+OpenMP hybrid programming model on multi-core systemsabstractState-of-the-art MPI libraries rely on locks to guarantee thread-safety. This discourages application developers from using multiple threads to perform MPI operations. In this paper, we propose a high performance, lock-free multi-endpoint MPI runtime, which can achieve up to 40\% improvement for point-to-point operation and one representative collective operation with minimum or no modifications to the existing applications. Miao Luo, Xiaoyi Lu 0001, Khaled Hamidouche, Krishna Chaitanya Kandalla, Dhabaleswar K. Panda 0001 |
PPoPP | 2 |
| 2013 | SR-IOV Support for Virtualization on InfiniBand Clusters: Early ExperienceabstractHigh Performance Computing (HPC) systems are becoming increasingly complex and are also associated with very high operational costs. The cloud computing paradigm, coupled with modern Virtual Machine (VM) technology offers attractive techniques to easily manage large scale systems, while significantly bringing down the cost of computation, memory and storage. However, running HPC applications on cloud systems still remains a major challenge. One of the biggest hurdles in realizing this objective is the performance offered by virtualized computing environments, more specifically, virtualized I/O devices. Since HPC applications and communication middlewares rely heavily on advanced features offered by modern high performance interconnects such as InfiniBand, the performance of virtualized InfiniBand interfaces is crucial. Emerging hardware-based solutions, such as the Single Root I/O Virtualization (SR-IOV), offer an attractive alternative when compared to existing software-based solutions. The benefits of SR-IOV have been widely studied for GigE and 10GigE networks. However, with InfiniBand networks being increasingly adopted in the cloud computing domain, it is critical to fully understand the performance benefits of SR-IOV in InfiniBand network, especially for exploring the performance characteristics and trade-offs of HPC communication middlewares (such as Message Passing Interface (MPI), Partitioned Global Address Space (PGAS)) and applications. To the best of our knowledge, this is the first paper that offers an in-depth analysis on SR-IOV with InfiniBand. Our experimental evaluations show that for the performance of MPI and PGAS point-to-point communication benchmarks over SR-IOV with InfiniBand is comparable to that of the native InfiniBand hardware, for most message lengths. However, we observe that the performance of MPI collective operations over SR-IOV with InfiniBand is inferior to native (non-virtualized) mode. We also evaluate the trade-offs of various VM to CPU mapping policies on modern multi-core architectures and present our experiences. Xiaoyi Lu 0001, Krishna Chaitanya Kandalla, Mark Daniel Arnold, Dhabaleswar K. Panda 0001 |
CCGRID | 3 |
| 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 | 2 |
| 2013 | A scalable and portable approach to accelerate hybrid HPL on heterogeneous CPU-GPU clustersabstractAccelerating High-Performance Linkpack (HPL) on heterogeneous clusters with multi-core CPUs and GPUs has attracted a lot of attention from the High Performance Computing community. It is becoming common for large scale clusters to have GPUs on only a subset of nodes in order to limit system costs. The major challenge for HPL in this case is to efficiently take advantage of all the CPU and GPU resources available on a cluster. In this paper, we present a novel two-level workload partitioning approach for HPL that distributes workload based on the compute power of CPU/GPU nodes across the cluster. Our approach also handles multi-GPU configurations. Unlike earlier approaches for heterogeneous clusters with CPU and GPU nodes, our design takes advantage of asynchronous kernel launches and CUDA copies to overlap computation and CPU-GPU data movement. It uses techniques such as process grid reordering to reduce MPI communication/contention while ensuring load balance across nodes. Our experimental results using 32 GPU and 128 CPU nodes of Oakley, a research cluster at Ohio Supercomputer Center, shows that our proposed approach can achieve more than 80% of combined actual peak performance of CPU and GPU nodes. This provides 47% and 63% increase in the HPL performance that can be reported using only CPU nodes and only GPU nodes, respectively. Rong Shi, Sreeram Potluri, Khaled Hamidouche, Xiaoyi Lu 0001, Karen A. Tomko, Dhabaleswar K. Panda 0001 |
CLUSTER | 4 |
| 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 | 1 |
| 2011 | Vega LingCloud: A Resource Single Leasing Point System to Support Heterogeneous Application Modes on Shared InfrastructureabstractIn large organizations or IDCs, different departments always occupy and maintain dedicated resources to satisfy their or their customers' heterogeneous application loads. This situation easily makes the infrastructure management a repeated and inefficient work. Even worse, it is difficult to share the resources owned by different departments even when they are idle, because the application modes on these resources are quite different. This paper introduces a live system, Vega Ling Cloud, which provides a Resource Single Leasing Point System for consolidated renting physical and virtual machines to support heterogeneous application modes on shared infrastructure. Furthermore, we present the asset-leasing model and the architecture of Vega Ling Cloud. According to the evaluation, Vega Ling Cloud is better than other systems like Open Nebula and Enomaly ECP in the aspects of uniformity, flexibility, security, usability, and efficiency. The experimental result of management overhead shows that the deployment speed of virtual machine in Vega Ling Cloud is 4.1 times of that in the Open Nebula and VIDA hybrid system for deploying 64 virtual machines concurrently. From a representative micro-cloud example in a research group, we show that the consolidated way of leasing physical and virtual machine in Vega Ling Cloud is approbatory. Up to now, Vega Ling Cloud has been deployed in real world environments, which include a private cloud of large organization in Beijing and a public cloud in the Dongguan IDC of China. The resource scale of the Dongguan cloud reaches about 504 cores, 625 GB memory, and 156 TB storage. The number of supported real applications in Beijing and Dongguan clouds has exceeded 35, and their modes involve high performance computing, large scale data processing, virtual machine leasing, data storage, and so on. Xiaoyi Lu 0001, Jian Lin 0006, Li Zha, Zhiwei Xu 0002 |
ISPA | 1 |
| 2010 | VegaWarden: A Uniform User Management System for Cloud ApplicationsabstractIn a virtual cluster based Cloud Computing environment, the sharing of infrastructure introduces two problems on user management: usability and security. Meanwhile, we observe that most conventional user management frameworks in the network environment are not fit for the scale expansion and interconnection of dynamic virtualization environment. In this paper, we propose VegaWarden, a uniform user management system to solve these problems. VegaWarden supplies a global user space for different virtual infrastructures and application services in one Cloud, and allows user system interconnection among homogeneous Cloud instances. A uniform authentication model enables the security isolation of different administrative domains, and a decentralized architecture ensures its scalability. We have implemented VegaWarden in an experimental Cloud-oriented infrastructure and a production Grid Computing environment. The functionality and performance of VegaWarden have been demonstrated. Jian Lin 0006, Xiaoyi Lu 0001, Yongqiang Zou, Li Zha |
NAS | 2 |
| 2010 | JAMILA: A Usable Batch Job Management System to Coordinate Heterogeneous Clusters and Diverse Applications over Grid or Cloud Infrastructure
Juan Peng, Xiaoyi Lu 0001, Boqun Cheng, Li Zha |
NPC | 2 |
| 2010 | Investigating, Modeling, and Ranking Interface Complexity of Web Services on the World Wide WebabstractAnalyzing factors of affecting Web Service invocation performance is a hot topic. Among the factors, service interface complexity is a key one investigated by much research work. However, these researches mainly analyze the impact on the performance of primitives, some simple data structures like mesh interface object, or array of them. For the complex data structures, these works lack of a systematic approach to characterize the impact. This paper firstly makes a detailed statistics of service interface complexity based on a large sample space (10K+ WSDL files) on the World Wide Web, and we find about 41.6% services contain complex data structures. The statistic results guide us to conduct many experiments for finding out the correlation of service interface complexity to invocation performance. We discover an interesting feature on commonly used Web Service platforms (Axis/gSOAP/.Net), called Data Structural Form Unaware. As each parameter or return value can be represented as a tree with two kinds of nodes, structure type node and primitive type node, then the feature means that the overhead caused by parameters or return values is independent of the organization of nodes in the trees, but it is only related with the number and the type of nodes. By this feature, we present a simple model to quantify the impact of service interface complexity. Using our model, a Service Interface Performance Vector can be calculated by parsing a WSDL file to estimate the overhead of service interface design. This vector, together with the invocation probability of each operation, generates the Service Interface Performance Score, a comprehensive complexity indicator which can be used to evaluate and rank the service interfaces. At last, a rank report of services on the WWW is shown. Our work can play a guiding role in service interface designing, service interface performance predicting, and ranking. Xiaoyi Lu 0001, Jian Lin 0006, Yongqiang Zou, Juan Peng, Xingwu Liu, Li Zha |
SERVICES | 1 |
| 2009 | ICOMC: Invocation Complexity Of Multi-Language Clients for Classified Web Services and its Impact on Large Scale SOA ApplicationsabstractTheoretically, multi-language clients invocating web services is no longer a problem due to XML-based interface descriptions by WSDL, but the reality is not so good. Some implementation level difficulties still exist when invoking web services from clients in different programming languages. These difficulties are caused by involving complex data structures in the service interface, carrying additional information such as WS-security headers in the SOAP messages, missing language features such as Reflection in C/C++ and so on, which make large scale multi-language SOA application development a time-consuming and buggy work. This paper proposes a new complexity ICOMC, short for Invocation Complexity Of Multi-language Clients, to quantify these difficulties, introduces implementation cost and runtime performance metrics for ICOMC, and indentifies three factors dominating the ICOMC: service interface, message context, and language feature. Consequently, the problem is formulated as finding out the correlation of the three factors to ICOMC. To simplify the problem, web services are classified into four categories: SISM, SICM, CISM and CICM according to service interface complexity and message context complexity. Furthermore, micro-benchmark experiments are done in C/C++/Java for all four categories. This paper also takes the GOS System Software of the China National Grid as a real large scale application to implement its C/C++ client APIs and compare them with the original Java APIs. Evaluations based on micro-benchmarks and real application show the correlations between the factors and ICOMC. Our results benefit web service interface designing, appropriate language adoption, and implementation cost / runtime performance estimation. Xiaoyi Lu 0001, Yongqiang Zou, Jian Lin 0006, Li Zha |
PDCAT | 1 |
| 2008 | An Experimental Analysis for Memory Usage of GOS CoreabstractAs grid software and grid applications grow in size and complexity, memory usage experiment and analysis have become more and more essential. In this paper, we investigate an Anti-TestCase-based approach to analyze the memory usage of GOS core software, and present a set of Anti-TestCases, which describe predefined memory problems of GOS core. This approach is applied with a set of automated test tools in controlled experiments including basic memory requirements, memory errors and memory usage of GOS core. With the experiment and analysis, we find several memory usage results of GOS core and they are valuable for other systems based on JVM, Tomcat and Axis. We also give some JVM memory configuration suggestions for better performance in the production server. Xiaoyi Lu 0001, Qiang Yue 0001, Yongqiang Zou |
PDCAT | 1 |