EDBT 2026 Demo / reviewers in the wild / expert
Weikuan Yu
dblp:50/1859
· DBLP profile ↗
79ranked-venue papers
18as first author
13since 2021 · last 2026
0000-0002-8754-0311ORCID · verified
Domains — the database's venue-derived domains; a paper can count in several
Systems, architecture and hardware · 73 · 18 first-author · 11 since 2021Software engineering, systems software and programming languages · 3Databases, data management, data science and information retrieval · 2 · 1 since 2021Applied, interdisciplinary, general and emerging computing · 2Artificial intelligence and machine learning · 1Computer networks · 1Security and privacy · 1 · 1 since 2021
| Year | Publication | Venue | Position |
|---|---|---|---|
| 2026 | VeX: Scaling HNSW-Based Vector Search with DPU Memory and Parallelism
Hyungsun Yoo, Woojung Kim, Donghyun Min, Myungcheol Lee, Jihoon Yang, Weikuan Yu, Youngjae Kim 0001 |
CCGrid | 7 |
| 2026 | Dual-Blade: Dual-Path NVMe-Direct KV-Cache Offloading for Edge LLM Inference
Bodon Jeong, Hongsu Byun, Youngjae Kim 0001, Weikuan Yu, Kyungkeun Lee, Jihoon Yang, Sungyong Park |
ICDCS | 4 |
| 2026 | 2FiA: Towards WiFi Sensing-Based Authentication with Unique Biometrics
Bofan Li, Zhankai Ye, Weikuan Yu, Yongning Tang, Liu Xiu |
SP | 3 |
| 2025 | MEMORYBRIDGE: Leveraging Cloud Resource Characteristics for Cost-Efficient Disk-Based GNN Training via Two-Level ArchitectureabstractGraph Neural Networks (GNNs) are machine learning models that process graph-structured data by learning relationships between vertices and edges, as well as graph-level characteristics. Recently, with the emergence of large graph datasets on a TB scale, dataset sizes have exceeded the memory capacity of single machines. As a result, traditional methods that load all graph data into memory have become unusable, leading to the emergence of disk-based GNN training that uses storage as a memory extension. Recent research has focused on reducing disk I/O bottlenecks in disk-based GNNs. However, disk-based GNNs face new challenges in cloud environments due to two main characteristics. First, compared to node-local environments, the significantly slower cloud storage I/O speed becomes the main bottleneck of the entire training process. Second, pre-defined virtual machines prevent users from freely utilizing desired memory sizes, bandwidth, and the latest GPU technologies. These limitations have made existing disk-based GNN research unusable in cloud environments. To overcome this, we propose MEMORYBRIDGE, a system that cost-effectively accelerates GNN training in cloud environments through a novel two-level architecture that utilizes affordable GPU resources as training nodes and remote memory resources without GPUs as memory nodes, instead of using a single expensive GPU resource. This architecture consists of two key components: (i) a mathematical solver that recommends the most cost-effective resource combination, and (ii) a cloud-specialized GNN framework that implements graphaware fixed caching and batch pipelining optimization. The experimental results show that MEMORYBRIDGE achieved a speed improvement of up to 32.7x compared to existing GNN training frameworks and a cost efficiency of 9.9x compared to alternative resource configuration strategies, effectively handling the unique problems that arise from the combination of cloud environments and GNN training. Yoochan Kim, Weikuan Yu, Hong-Yeon Kim, Youngjae Kim 0001 |
CCGrid | 2 |
| 2025 | Spatio-Temporal Resource Control for Cloud-Native GPU ProvisioningabstractModern cloud platforms, such as Kubernetes, provide a service-oriented resource abstraction for explicit Quality of Service (QoS) provisioning, which guarantees resource reservations and supports resource elasticity. In regards to CPU resources, tenants can explicitly specify a guaranteed base demand and an upper bound for resources. Thus, it is desirable that GPU resource provisioning should be analogous to mature cloud-native CPU resource abstraction and can support familiar QoS classes, such as Guaranteed, Burstable, and BestEffort. Although existing research efforts have primarily focused on maximizing GPU utilization by exploiting profiling, capabilities for guaranteed reservation, precise throttling, and elastic bursting are essential to support cloud-native GPU provisioning. It is challenging to provide such features due to the GPU's asynchronous and non-preemptive characteristics. To address this issue, we introduce a spatio-temporal provisioning framework that ensures both resource guarantee and elasticity for GPUs. We define a resource model to control both spatial and temporal dimensions of GPU resources and support familiar QoS classes. Our framework features a per-container Agent for transparent resource accounting and local quota enforcement, and the central Multi-Tenant Arbitrator for global, partition-aware fair scheduling. Performance evaluation demonstrates that our framework can provide accurate GPU resource guarantee and manages dynamic mixed-QoS workloads by honoring reserved GPU resources during contention while allowing tenants to elastically burst into idle capacity up to their limit. Hyeon-Jun Jang, Sang-Jae Kim, Weikuan Yu, Hyun-Wook Jin |
SoCC | 3 |
| 2025 | LENS: label sparsity-tolerant adversarial learning on spatial deceptive reviews
Sirish Prabakar, Haiquan Chen 0001, Zhe Jiang 0001, Carl Yang 0001, Weikuan Yu, Da Yan 0001 |
GeoInformatica | 5 |
| 2024 | DFTracer: An Analysis-Friendly Data Flow Tracer for AI-Driven WorkflowsabstractModern HPC workflows involve intricate coupling of simulation, data analytics, and artificial intelligence (AI) applications to improve time to scientific insight. These workflows require a cohesive set of performance analysis tools to provide a comprehensive understanding of data exchange patterns in HPC systems. However, current tools are not designed to work with an AI-based I/O software stack that requires tracing at multiple levels of the application. To this end, we developed a data flow tracer called DFTracer to capture data-centric events from workflows and the I/O stack to build a detailed understanding of the data exchange within AI-driven workflows. DFTracer has the following three novel features, including a unified interface to capture trace data from different layers in the software stack, a trace format that is analysis-friendly and optimized to support efficiently loading multi-million events in a few seconds, and the capability to tag events with workflow-specific context to perform domain-centric data flow analysis for workflows. Additionally, we demonstrate that DFTracer has a $1.44 x$ smaller runtime overhead and 1.3-7.1x smaller trace size than state-of-the-art tracing tools such as Score-P, Recorder, and Darshan. Moreover, with AI-driven workflows, Score-P, Recorder, and Darshan cannot find I/O accesses from dynamically spawned processes, and their load performance of 100 M events is three orders of magnitude slower than DFTracer. In conclusion, we demonstrate that DFTracer can capture multi-level performance data, including contextual event tagging with a low overhead of 1-5% from AI-driven workflows such as MuMMI and Microsoft’s Megatron Deepspeed running on large-scale HPC systems. Hariharan Devarajan, Loïc Pottier, Kaushik Velusamy, Huihuo Zheng, Izzet Yildirim, Olga Kogiou, Weikuan Yu, Antonios Kougkas, Xian-He Sun, Jae-Seung Yeom, Kathryn Mohror |
SC | 7 |
| 2022 | SVAGC: Garbage Collection with a Scalable Virtual Address Swapping TechniqueabstractManaged programming languages including Java and Scala are very popular for data analytics and mobile applications. However, they often face challenging issues due to the overhead caused by the automatic memory management to detect and reclaim free available memory. It has been observed that during their Garbage Collection (GC), excessively long pauses can account for up to 40 % of the total execution time. Therefore, mitigating the GC overhead has been an active research topic to satisfy today's application requirements. This paper proposes a new technique called SwapVA to improve data copying in the copying/moving phases of GCs and reduce the GC pause time, thereby mitigating the issue of GC overhead. Our contribution is twofold. First, a SwapVA system call is introduced as a zero-copy technique to accelerate the GC copying/moving phase. Second, for the demonstration of its effectiveness, we have integrated SwapVA into SVAGC as an implementation of scalable Full GC on multi-core systems. Based on our results, the proposed solutions can dramatically reduce the GC pause in applications with large objects by as much as 70.9% and 97%, respectively, in the Sparse.large/4 (one quarter of the default input size) and Sigverify benchmarks. Ismail Ataie, Weikuan Yu |
CLUSTER | 2 |
| 2022 | DFMan: A Graph-based Optimization of Dataflow Scheduling on High-Performance Computing SystemsabstractScientific research and development campaigns are materialized by workflows of applications executing on high-performance computing (HPC) systems. These applications con-sist of tasks that can have inter- or intra-application flows of data to achieve the research goals successfully. These dataflows create dependencies among the tasks and cause resource con-tention on shared storage systems, thus limiting the aggregated I/O bandwidth achieved by the workflow. However, these I/O performance issues are often solved by tedious and manual efforts that demand holistic knowledge about the data dependencies in the workflow and the information about the infrastructure being utilized. Taking this into consideration, we design DFMan, a graph-based dataflow management and optimization framework for maximizing I/O bandwidth by leveraging the powerful storage stack on HPC systems to manage data sharing optimally among the tasks in the workflows. In particular, we devise a graph-based optimization algorithm that can leverage an intuitive graph representation of dataflow- and system-related information, and automatically carry out co-scheduling of task and data placement. According to our experiments, DFMan optimizes a wide variety of scientific workflows such as Hurricane 3D on Cloud Model 1 (CM1), Montage Carina Nebula (NGC3372), and an emulated dataflow kernel of the Multiscale Machine-learned Modeling Infrastructure (MuMMI I/O) on the Lassen supercomputer, and improves their aggregated I/O bandwidth by up to 5.42 x, 2.12 x and 1.29 x, respectively, compared to the baseline bandwidth. Fahim Chowdhury, Francesco Di Natale, Adam Moody, Kathryn Mohror, Weikuan Yu |
IPDPS | 5 |
| 2022 | PhaST: Hierarchical Concurrent Log-Free Skip List for Persistent MemoryabstractSkip list (skiplist) is a competitive index structure that offers superior concurrency and excellent performance but with high memory overhead and low access locality. Emerging persistent memory (PM) technologies present an opportunity to mitigate the capacity constraint of DRAM. However, data consistency on PM typically results in excessive write overhead. In addition, fast concurrent access to an index is critical to the throughput on high-end contemporary computer systems. In this article, we propose a Partitioned HierArchical SkiplisT calledPhaST, which can simultaneously reduce the skiplist height and improve its access locality, through its hierarchy of component structures, while enabling fast parallel recovery in case of failure. To ensure high concurrency and fast data consistency, we also have developed writelock-free concurrent insert and log-free atomic split. Furthermore, we have developed a durable lock-free concurrent search that can discern transient structural inconsistencies and deliver highly concurrent read operations. We have conducted an extensive evaluation ofPhaSTcompared to state-of-the-art studies such as NV-Skiplist, wB+-Tree, FPTree, and FAST-FAIR. Our evaluation results showPhaSToutperforms other indexing structures by up to 4.05× and 2.87× in single-threaded inserts and searches, and 1.56× and 2.62× in concurrent inserts and searches. Zhenxin Li, Bing Jiao, Shuibing He, Weikuan Yu |
IEEE Trans. Parallel Distributed Syst. | 4 |
| 2021 | Compression of Time Evolutionary Image Data through Predictive Deep Neural NetworksabstractRecent advances in Deep Neural Networks (DNNs) have demonstrated a promising potential in predicting the temporal and spatial proximity of time evolutionary data. In this paper, we have developed an effective (de)compression framework called TEZIP that can support dynamic lossy and lossless compression of time evolutionary image frames with high compression ratio and speed. TEZIP first trains a Recurrent Neural Network called PredNet to predict future image frames based on base frames, and then derives the resulting differences between the predicted frames and the actual frames as more compressible delta frames. Next we equip TEZIP with techniques that can exploit spatial locality for the encoding of delta frames and apply lossless compressors on the resulting frames. Furthermore, we introduce window-based prediction algorithms and dynamically pinpoint the trade-off between the window size and the relative errors of predicted frames. Finally, we have conducted an extensive set of tests to evaluate TEZIP. Our experimental results show that, in terms of compression ratio, TEZIP outperforms existing lossless compressors such as x265 by up to 3.2x and lossy compressors such as SZ by up to 3.3x. Rupak Roy, Kento Sato, Subhadeep Bhattacharya, Xingang Fang, Yasumasa Joti, Takaki Hatsui, Toshiyuki Nishiyama Hiraki, Jian Guo 0004, Weikuan Yu |
CCGRID | 9 |
| 2021 | O(1) Communication for Distributed SGD through Two-Level Gradient AveragingabstractLarge neural network models present a hefty communication challenge to distributed Stochastic Gradient Descent (SGD), with a per-iteration communication complexity of $\mathcal{O}(n)$ per worker for a model of n parameters. Many sparsification and quantization techniques have been proposed to compress the gradients, some reducing the per-iteration communication complexity to $\mathcal{O}(k)$, where $k\ll n$. In this paper, we introduce a strategy called two-level gradient averaging (A2SGD) to consolidate all gradients down to merely two local averages per worker before the computation of two global averages for an updated model. A2SGD also retains local errors to maintain the variance for fast convergence. Our analysis shows that A2SGD converges similar to the default distributed SGD algorithm. Our evaluation validates the conclusion and demonstrates that A2SGD significantly reduces the communication traffic per worker, and improves the overall training time of LSTM-PTB by $3.2\times$ and $23.2\times$, compared to Top-K and QSGD, respectively. We evaluate the effectiveness of our approach using two kinds of optimizers, SGD and Adam. Also, our evaluation with various communication options demonstrates the strength of our approach both in terms of communication reduction and convergence. To the best of our knowledge, A2SGD is the first to achieve $\mathcal{O}$ (1) communication complexity per worker without incurring a significant accuracy degradation of DNN models while communicating only two scalars representing gradients per worker for distributed SGD. Subhadeep Bhattacharya, Weikuan Yu, Fahim Chowdhury, Kathryn Mohror |
CLUSTER | 2 |
| 2021 | ROBOTune: High-Dimensional Configuration Tuning for Cluster-Based Data AnalyticsabstractSpark is popular for its ability to enable high-performance data analytics applications on diverse systems. Its great versatility is achieved through numerous user- and system-level options, resulting in an exponential configuration space that, ironically, hinders data analytics’s optimal performance. The colossal complexity is caused by two main issues: the high dimensionality of configuration space and the expensive black-box configuration-performance relationship. In this paper, we design and develop a robust tuning framework called ROBOTune that can tackle both issues and tune Spark applications quickly for efficient data analytics. Specifically, it performs parameter selection through a Random Forests based model to reduce the dimensionality of analytics configuration space. In addition, ROBOTune employs Bayesian Optimization to overcome the complex nature of the configuration-performance relationship and balance exploration and exploitation to efficiently locate a globally optimal or near-optimal configuration. Furthermore, ROBOTune strengthens Latin Hypercube Sampling with caching and memoization to enhance the coverage and effectiveness in the generation of sample configurations. Our evaluation results demonstrate that ROBOTune finds similar or better performing configurations than contemporary tuning tools like BestConfig and Gunther while improving on search cost by 1.59 × and 1.53 × on average and up to 2.27 × and 1.71 × , respectively. Muhib Khan, Weikuan Yu |
ICPP | 2 |
| 2020 | Ad Hoc File Systems for High-Performance Computing
André Brinkmann, Kathryn Mohror, Weikuan Yu, Philip H. Carns, Toni Cortes, Scott Klasky, Alberto Miranda, Franz-Josef Pfreundt, Robert B. Ross, Marc-Andre Vef |
J. Comput. Sci. Technol. | 3 |
| 2019 | Efficient User-Level Storage Disaggregation for Deep LearningabstractOn large-scale high performance computing (HPC) systems, applications are provisioned with aggregated resources to meet their peak demands for brief periods. This results in resource underutilization because application requirements vary a lot during execution. This problem is particularly pronounced for deep learning applications that are running on leadership HPC systems with a large pool of burst buffers in the form of flash or non-volatile memory (NVM) devices. In this paper, we examine the I/O patterns of deep neural networks and reveal their critical need of loading many small samples randomly for successful training. We have designed a specialized Deep Learning File System (DLFS) that provides a thin set of APIs. Particularly, we design the metadata management of DLFS through an in-memory tree-based sample directory and its file services through the user-level SPDK protocol that can disaggregate the capabilities of NVM Express (NVMe) devices to parallel training tasks. Our experimental results show that DLFS can dramatically improve the throughput of training for deep neural networks on NVMe over Fabric, compared with the kernel-based Ext4 file system. Furthermore, DLFS achieves efficient user-level storage disaggregation with very little CPU utilization. Yue Zhu 0002, Weikuan Yu, Bing Jiao, Kathryn Mohror, Adam Moody, Fahim Chowdhury |
CLUSTER | 2 |
| 2019 | I/O Characterization and Performance Evaluation of BeeGFS for Deep LearningabstractParallel File Systems (PFSs) are frequently deployed on leadership High Performance Computing (HPC) systems to ensure efficient I/O, persistent storage and scalable performance. Emerging Deep Learning (DL) applications incur new I/O and storage requirements to HPC systems with batched input of small random files. This mandates PFSs to have commensurate features that can meet the needs of DL applications. BeeGFS is a recently emerging PFS that has grabbed the attention of the research and industry world because of its performance, scalability and ease of use. While emphasizing a systematic performance analysis of BeeGFS, in this paper, we present the architectural and system features of BeeGFS, and perform an experimental evaluation using cutting-edge I/O, Metadata and DL application benchmarks. Particularly, we have utilized AlexNet and ResNet-50 models for the classification of ImageNet dataset using the Livermore Big Artificial Neural Network Toolkit (LBANN), and ImageNet data reader pipeline atop TensorFlow and Horovod. Through extensive performance characterization of BeeGFS, our study provides a useful documentation on how to leverage BeeGFS for the emerging DL applications. Fahim Chowdhury, Yue Zhu 0002, Todd Heer, Saul Paredes, Adam Moody, Robin Goldstone, Kathryn Mohror, Weikuan Yu |
ICPP | 8 |
| 2019 | Exploration of memory hybridization for RDD caching in SparkabstractApache Spark is a popular cluster computing framework for iterative analytics workloads due to its use of Resilient Distributed Datasets (RDDs) to cache data for in-memory processing. We have revealed that the performance of Spark RDD cache can be severely limited if its capacity falls short to the needs of the workloads. In this paper, we have explored different memory hybridization strategies to leverage emergent Non-Volatile Memory (NVM) devices for Spark's RDD cache. We have found that a simple layered hybridization approach does not offer an effective solution. Therefore, we have designed a flat hybridization scheme to leverage NVM for caching RDD blocks, along with several architectural optimizations such as dynamic memory allocation for block unrolling, asynchronous migration with preemption, and opportunistic eviction to disk. We have performed an extensive set of experiments to evaluate the performance of our proposed flat hybridization strategy and found it to be robust in handling different system and NVM characteristics. Our proposed approach uses DRAM for a fraction of the hybrid memory system and yet manages to keep the increase in execution time to be within 10% on average. Moreover, our opportunistic eviction of blocks to disk improves performance by up to 7.5% when utilized alongside the current mechanism. Muhib Khan, Muhammad Ahad Ul Alam, Amit Kumar Nath, Weikuan Yu |
ISMM | 4 |
| 2019 | Multivariate modeling and two-level scheduling of analytic queries
Amit Kumar Nath, Xiaoning Ding, Huansong Fu, Muhib Khan, Weikuan Yu |
Parallel Comput. | 6 |
| 2018 | SHMEMGraph: Efficient and Balanced Graph Processing Using One-Sided CommunicationabstractState-of-the-art synchronous graph processing frameworks face both inefficiency and imbalance issues that cause their performance to be suboptimal. These issues include the inefficiency of communication and the imbalanced graph computation/communication costs in an iteration. We propose to replace their conventional two-sided communication model with the one-sided counterpart. Accordingly, we design SHMEMGraph, an efficient and balanced graph processing framework that is formulated across a global memory space and takes advantage of the flexibility and efficiency of one-sided communication for graph processing. Through an efficient one-sided communication channel, SHMEMGraph utilizes the high-performance operations with RDMA while minimizing the resource contention within a computer node. In addition, SHMEMGraph synthesizes a number of optimizations to address both computation imbalance and communication imbalance. By using a graph of 1 billion edges, our evaluation shows that compared to the state-of-the-art Gemini framework, SHMEMGraph achieves an average improvement of 35.5% in terms of job completion time for five representative graph algorithms. Huansong Fu, Manjunath Gorentla Venkata, Shaeke Salman, Neena Imam, Weikuan Yu |
CCGrid | 5 |
| 2018 | Entropy-Aware I/O Pipelining for Large-Scale Deep Learning on HPC SystemsabstractDeep neural networks have recently gained tremendous interest due to their capabilities in a wide variety of application areas such as computer vision and speech recognition. Thus it is important to exploit the unprecedented power of leadership High-Performance Computing (HPC) systems for greater potential of deep learning. While much attention has been paid to leverage the latest processors and accelerators, I/O support also needs to keep up with the growth of computing power for deep neural networks. In this research, we introduce an entropy-aware I/O framework called DeepIO for large-scale deep learning on HPC systems. Its overarching goal is to coordinate the use of memory, communication, and I/O resources for efficient training of datasets. DeepIO features an I/O pipeline that utilizes several novel optimizations: RDMA (Remote Direct Memory Access)-assisted in-situ shuffling, input pipelining, and entropy-aware opportunistic ordering. In addition, we design a portable storage interface to support efficient I/O on any underlying storage system. We have implemented DeepIO as a prototype for the popular TensorFlow framework and evaluated it on a variety of different storage systems. Our evaluation shows that DeepIO delivers significantly better performance than existing memory-based storage systems. Yue Zhu 0002, Fahim Chowdhury, Huansong Fu, Adam Moody, Kathryn Mohror, Kento Sato, Weikuan Yu |
MASCOTS | 7 |
| 2017 | High-Performance Key-Value Store On OpenSHMEMabstractRecently, there has been a growing interest in enabling fast data analytics by leveraging system capabilities from large-scale high-performance computing (HPC) systems. OpenSHMEM is a popular run-time system on HPC systems that has been used for large-scale compute-intensive scientific applications. In this paper, we propose to leverage OpenSHMEM to design a distributed in-memory key-value store for fast data analytics. Accordingly, we have developed SHMEMCache on top of OpenSHMEM to leverage its symmetric global memory, efficient one-sided communication operations and general portability. We have also evaluated SHMEMCache through extensive experimental studies. Our results show that SHMEMCache has accomplished significant performance improvements over hte original Memcached in terms of latency and throughput. Our evaluation on the Titan supercomputer has also demonstrated that SHMEMCache can scale to 1024 nodes. Huansong Fu, Manjunath Gorentla Venkata, Ahana Roy Choudhury, Neena Imam, Weikuan Yu |
CCGrid | 5 |
| 2017 | MetaKV: A Key-Value Store for Metadata Management of Distributed Burst BuffersabstractDistributed burst buffers are a promising storage architecture for handling I/O workloads for exascale computing. Their aggregate storage bandwidth grows linearly with system node count. However, although scientific applications can achieve scalable write bandwidth by having each process write to its node-local burst buffer, metadata challenges remain formidable, especially for files shared across many processes. This is due to the need to track and organize file segments across the distributed burst buffers in a global index. Because this global index can be accessed concurrently by thousands or more processes in a scientific application, the scalability of metadata management is a severe performance-limiting factor. In this paper, we propose MetaKV: a key-value store that provides fast and scalable metadata management for HPC metadata workloads on distributed burst buffers. MetaKV complements the functionality of an existing key-value store with specialized metadata services that efficiently handle bursty and concurrent metadata workloads: compressed storage management, supervised block clustering, and log-ring based collective message reduction. Our experiments demonstrate that MetaKV outperforms the state-of-the-art key-value stores by a significant margin. It improves put and get metadata operations by as much as 2.66× and 6.29×, respectively, and the benefits of MetaKV increase with increasing metadata workload demand. Teng Wang 0001, Adam Moody, Yue Zhu 0002, Kathryn Mohror, Kento Sato, Tanzima Z. Islam, Weikuan Yu |
IPDPS | 7 |
| 2017 | FARMS: Efficient mapreduce speculation for failure recovery in short jobs
Huansong Fu, Haiquan Chen 0001, Yue Zhu 0002, Weikuan Yu |
Parallel Comput. | 4 |
| 2017 | A case study of tuning MapReduce for efficient Bioinformatics in the cloud
Lizhen Shi, Zhong Wang 0003, Weikuan Yu, Xiandong Meng |
Parallel Comput. | 3 |
| 2016 | OAWS: Memory Occlusion Aware Warp SchedulingabstractWe have closely examined GPU resource utilization when executing memory-intensive benchmarks. Our detailed analysis of GPU global memory accesses reveals that divergent loads can lead to the occlusion of Load-Store units, resulting in quick consumption of MSHR entries. Such memory occlusion prevents other ready memory instructions from accessing L1 data cache, eventually stalling warp schedulers and degrading the overall performance. We have designed memory Occlusion Aware Warp Scheduling (OAWS) that can dynamically predict the demand of MSHR entries of divergent memory instructions, and maximize the number of concurrent warps such that their aggregate MSHR consumptions are within the MSHR capacity. Our dynamic OAWS policy can prevent memory occlusions and effectively leverage more MSHR entries for better IPC performance for GPU. Experimental results show that the static and dynamic versions of OAWS achieve 36.7% and 73.1% performance improvement, compared to the baseline GTO scheduling. Particularly, dynamic OAWS outperforms MASCAR, CCWS, and SWL-Best by 70.1%, 57.8%, and 11.4%, respectively. Bin Wang 0019, Yue Zhu 0002, Weikuan Yu |
PACT | 3 |
| 2016 | An ephemeral burst-buffer file system for scientific applicationsabstractBurst buffers are becoming an indispensable hardware resource on large-scale supercomputers to buffer the bursty I/O from scientific applications. However, there is a lack of software support for burst buffers to be efficiently shared by applications within a batch-submitted job and recycled across different batch jobs. In addition, burst buffers need to cope with a variety of challenging I/O patterns from data-intensive scientific applications. In this study, we have designed an ephemeral Burst Buffer File System (BurstFS) that supports scalable and efficient aggregation of I/O bandwidth from burst buffers while having the same life cycle as a batch-submitted job. BurstFS features several techniques including scalable metadata indexing, co-located I/O delegation, and server-side read clustering and pipelining. Through extensive tuning and analysis, we have validated that BurstFS has accomplished our design objectives, with linear scalability in terms of aggregated I/O bandwidth for parallel writes and reads. Teng Wang 0001, Kathryn Mohror, Adam Moody, Kento Sato, Weikuan Yu |
SC | 5 |
| 2016 | Exploiting Analytics Shipping with Virtualized MapReduce on HPC Backend Storage ServersabstractLarge-scale scientific applications on High-Performance Computing (HPC) systems are generating a colossal amount of data that need to be analyzed in a timely manner for new knowledge, but are too costly to transfer due to their sheer size. Many HPC systems have catered to in situ analytics solutions that can analyze temporary datasets as they are generated, i.e., without storing to long-term storage media. However, there is still an open question on how to conduct efficient analytics of permanent datasets that have been stored to the backend persistent storage because of their long-term value. To fill the void, we exploit the analytics shipping model for fast analysis of large-scale scientific datasets on HPC backend storage servers. Through an efficient integration of MapReduce and the popular Lustre storage system, we have developed a Virtualized Analytics Shipping (VAS) framework that can ship MapReduce programs to Lustre storage servers. The VAS framework includes three component techniques: (a) virtualized analytics shipping with fast network and disk I/O; (b) stripe-aligned data distribution and task scheduling and (c) pipelined intermediate data merging and reducing. The first technique provides necessary isolation between MapReduce analytics and Lustre I/O services. The second and third techniques optimize MapReduce on Lustre and avoid explicit shuffling. Our performance evaluation demonstrates that VAS offers an exemplary implementation of analytics shipping and delivers fast and virtualized MapReduce programs on backend Lustre storage servers. Cong Xu 0008, Robin Goldstone, Byron Neitzel, Weikuan Yu |
IEEE Trans. Parallel Distributed Syst. | 6 |
| 2015 | TRIO: Burst Buffer Based I/O OrchestrationabstractThe growing computing power on leadership HPC systems is often accompanied by ever-escalating failure rates. Checkpointing is a common defensive mechanism used by scientific applications for failure recovery. However, directly writing the large and bursty checkpointing dataset to parallel file systems can incur significant I/O contention on storage servers. Such contention in turn degrades bandwidth utilization of storage servers and prolongs the average job I/O time of concurrent applications. Recently burst buffers have been proposed as an intermediate layer to absorb the bursty I/O traffic from compute nodes to storage backend. But an I/O orchestration mechanism is still desirable to efficiently move checkpointing data from burst buffers to storage backend. In this paper, we propose a burst buffer based I/O orchestration framework, named TRIO, to intercept and reshape the bursty writes for better sequential write traffic to storage servers. Meanwhile, TRIO coordinates the flushing orders among concurrent burst buffers to alleviate the contention on storage server. Our experimental results demonstrated that TRIO could efficiently utilize storage bandwidth and reduce the average job I/O time by 37% on average for data-intensive applications in typical checkpointing scenarios. Teng Wang 0001, Sarp Oral, Michael Pritchard, Bin Wang 0019, Weikuan Yu |
CLUSTER | 5 |
| 2015 | Eliminating intra-warp conflict misses in GPU
Bin Wang 0019, Weikuan Yu |
DATE | 4 |
| 2015 | DaCache: Memory Divergence-Aware GPU Cache ManagementabstractThe lock-step execution model of GPU requires a warp to have the data blocks for all its threads before execution. However, there is a lack of salient cache mechanisms that can recognize the need of managing GPU cache blocks at the warp level for increasing the number of warps ready for execution. In addition, warp scheduling is very important for GPU-specific cache management to reduce both intra- and inter-warp conflicts and maximize data locality. In this paper, we propose a Divergence-Aware Cache (DaCache) management that can orchestrate L1D cache management and warp scheduling together for GPGPUs. In DaCache, the insertion position of an incoming data block depends on the fetching warp's scheduling priority. Blocks of warps with lower priorities are inserted closer to the LRU position of the LRU-chain so that they have shorter lifetime in cache. This fine-grained insertion policy is extended to prioritize coherent loads over divergent loads so that coherent loads are less vulnerable to both inter- and intra-warp thrashing. DaCache also adopts a constrained replacement policy with L1D bypassing to sustain a good supply of Fully Cached Warps (FCW), along with a dynamic mechanism to adjust FCW during runtime. Our experiments demonstrate that DaCache achieves 40.4% performance improvement over the baseline GPU and outperforms two state-of-the-art thrashing-resistant techniques RRIP and DIP by 40% and 24.9%, respectively. Bin Wang 0019, Weikuan Yu, Xian-He Sun |
ICS | 2 |
| 2015 | Cracking Down MapReduce Failure Amplification through Analytics Logging and MigrationabstractMapReduce is popular for big data analytics because it offers easy-to-use map and reduce user interfaces while hiding the complexity of system scalability and fault resiliency issues. While a large body of literature has focused on improving the performance and scalability of MapReduce, the issue of fault resiliency has thus far received little attention. In this paper, we take on an effort to investigate the fault resiliency of MapReduce using YARN (the next-generation Hadoop) as a case study. We reveal that the failures of a MapTask, a ReduceTask or a compute node can cause distinctly different impact to MapReduce programs. Particularly, YARN MapReduce is not able to gracefully handle failures that involve ReduceTasks, causing prolonged task execution, delayed job completion, and, more severely, failure amplifications due to the cascading effects to other tasks. These problems together cause the performance collapse of MapReduce jobs. In this paper, we introduce a new fault-tolerant framework that can crack down failure amplification and gracefully handle failure scenarios. It is designed with two key fault handling techniques: analytics logging and speculative fast migration. Analytics logging is a light-weight mechanism that logs the key progress information of MapReduce tasks, speculative fast migration handles node failures by proactively re-executing MapTasks, migrating ReduceTasks, and collective merging with a pipeline of shuffle/merge and reduce stages. Our performance evaluation demonstrates that these techniques can eliminate failure amplification and deliver fast job execution compared to the existing task re-execution mechanism in MapReduce. Yandong Wang 0001, Huansong Fu, Weikuan Yu |
IPDPS | 3 |
| 2015 | Preserving Row Buffer Locality for PCM Wear-Leveling under Massive ParallelismabstractPhase Change Memory (PCM) is a promising alternative technology for DRAM because of its advantages in terms of transistor density and energy consumption. It has been exploited to work in concert or alone inside various memory systems to meet the growing bandwidth needs of massive parallelism. PCM memory cells, however, have a common problem of limited write endurance. Various wear-leaving techniques have been employed for uniform distribution of memory writes, typically through address transformation schemes such as randomization to avoid hot writes. Unfortunately, such address transformation can have the undesirable consequence of disrupting the row buffer locality in sequential memory accesses, resulting in the loss of memory performance. Our analysis reveals that this situation is particularly severe under massive parallelism of manycore processors such as GPUs. In this paper, we introduce a combination of two techniques, matrix-based partial randomization and rowbuffer locality-aware rotation, to alleviate the locality disruption of address transformation and preserve the row buffer locality of PCM-based global memory in GPU. Our evaluation results show that, compared to existing techniques, our techniques can adequately preserve the row buffer locality and minimize the loss of memory performance, while achieving similar endurance and better energy efficiency for a variety of GPGPU applications. Bin Wang 0019, Weikuan Yu |
MASCOTS | 4 |
| 2015 | SFMapReduce: An optimized MapReduce framework for Small FilesabstractHadoop, an open-source implementation of MapReduce, is widely used because of its ease of programming, scalability, and availability. With the explosive development of cloud computing, business and scientific applications increasingly take advantage of Hadoop. The sizes of files stored and processed in Hadoop are not bound to very large files anymore. However, Hadoop cannot provide stable and efficient services for small files at both storage and processing levels. To solve these problems, we propose an optimized MapReduce framework for small files, SFMapReduce. In SFMapReduce, we present two techniques, Small File Layout (SFLayout) and customized MapReduce (CMR). SFLayout is used to solve the memory problem and improve I/O performance in HDFS. CMR provides an interface for MapReduce so that SFMapReduce can process MapReduce with SFLayout efficiently. Our experimental results show that SFMapReduce decreases the memory pressure on the Hadoop NameNode, and provides better loading and retrieving throughput. On average, SFMapReduce achieves an improvement on MapReduce processing by 14.5 times and 20.8 times, compared with the original Hadoop and HAR layout. Hai Pham, Jianhui Yue, Weikuan Yu |
NAS | 5 |
| 2015 | Virtual Shuffling for Efficient Data Movement in MapReduceabstractMapReduce is a popular parallel processing framework for large-scale data analytics. To keep up with the increasing volume of datasets, it requires efficient I/O capability from the underlying computer systems to process and analyze data in two phases (mapping and reducing). Between these phases, MapReduce requires a shuffling phase to globally exchange the intermediate data generated by the mapping phase. We reveal that data shuffling, by physically moving segments of intermediate data across disks, causes significant I/O contention and compounds the I/O problem. In this paper, we propose a novel virtual shuffling strategy to enable efficient data movement and reduce I/O for MapReduce shuffling, thereby reducing power consumption and conserving energy. Virtual shuffling is realized through a combination of three techniques including a three-level segment table, near-demand merging, and dynamic and balanced merging subtrees. Our experimental results show that virtual shuffling significantly speeds up data movement in MapReduce and achieves faster job execution. Particularly, its reduction in disk I/O accesses results in as much as 12% savings in power consumption for MapReduce programs. Weikuan Yu, Yandong Wang 0001, Xinyu Que, Cong Xu 0008 |
IEEE Trans. Computers | 1 |
| 2014 | BurstMem: A high-performance burst buffer system for scientific applicationsabstractThe growth of computing power on large-scale systems requires commensurate high-bandwidth I/O systems. Many parallel file systems are designed to provide fast sustainable I/O in response to applications' soaring requirements. To meet this need, a novel system is imperative to temporarily buffer the bursty I/O and gradually flush datasets to long-term parallel file systems. In this paper, we introduce the design of BurstMem, a high-performance burst buffer system. BurstMem provides a storage framework with efficient storage and communication management strategies. Our experiments demonstrate that BurstMem is able to speed up the I/O performance of scientific applications by up to 8.5× on leadership computer systems. Teng Wang 0001, Sarp Oral, Yandong Wang 0001, Bradley W. Settlemyer, Scott Atchley, Weikuan Yu |
IEEE BigData | 6 |
| 2014 | Characterization and Optimization of Memory-Resident MapReduce on HPC SystemsabstractMapReduce is a widely accepted framework for addressing big data challenges. Recently, it has also gained broad attention from scientists at the U.S. leadership computing facilities as a promising solution to process gigantic simulation results. However, conventional high-end computing systems are constructed based on the compute-centric paradigm while big data analytics applications prefer a data-centric paradigm such as MapReduce. This work characterizes the performance impact of key differences between compute- and data-centric paradigms and then provides optimizations to enable a dual-purpose HPC system that can efficiently support conventional HPC applications and new data analytics applications. Using a state-of-the-art MapReduce implementation Spark and the Hyperion system at Lawrence Livermore National Laboratory, we have examined the impact of storage architectures, data locality and task scheduling to the memory-resident MapReduce jobs. Based on our characterization and findings of the performance behaviors, we have introduced two optimization techniques, namely Enhanced Load Balancer and Congestion-Aware Task Dispatching, to improve the performance of Spark applications. Yandong Wang 0001, Robin Goldstone, Weikuan Yu, Teng Wang 0001 |
IPDPS | 3 |
| 2014 | Non-work-conserving effects in MapReduce: diffusion limit and criticalityabstractSequentially arriving jobs share a MapReduce cluster, each desiring a fair allocation of computing resources to serve its associated map and reduce tasks. The model of such a system consists of a processor sharing queue for the MapTasks and a multi-server queue for the ReduceTasks. These two queues are dependent through a constraint that the input data of each ReduceTask are fetched from the intermediate data generated by the MapTasks belonging to the same job. A more generalized form of MapReduce queueing model can capture the essence of other distributed data processing systems that contain interdependent processor sharing queues and multi-server queues. Jian Tan 0001, Yandong Wang 0001, Weikuan Yu, Li Zhang 0002 |
SIGMETRICS | 3 |
| 2014 | Hello ADIOS: the challenges and lessons of developing leadership class I/O frameworksabstractSUMMARY Applications running on leadership platforms are more and more bottlenecked by storage input/output (I/O). In an effort to combat the increasing disparity between I/O throughput and compute capability, we created Adaptable IO System (ADIOS) in 2005. Focusing on putting users first with a service oriented architecture, we combined cutting edge research into new I/O techniques with a design effort to create near optimal I/O methods. As a result, ADIOS provides the highest level of synchronous I/O performance for a number of mission critical applications at various Department of Energy Leadership Computing Facilities. Meanwhile ADIOS is leading the push for next generation techniques including staging and data processing pipelines. In this paper, we describe the startling observations we have made in the last half decade of I/O research and development, and elaborate the lessons we have learned along this journey. We also detail some of the challenges that remain as we look toward the coming Exascale era. Copyright © 2013 John Wiley & Sons, Ltd. Qing Liu 0002, Jeremy Logan, Yuan Tian 0004, Hasan Abbasi, Norbert Podhorszki, Jong Choi 0001, Scott Klasky, Roselyne Tchoua, Jay F. Lofstead, Ron A. Oldfield, Manish Parashar, Nagiza F. Samatova, Karsten Schwan, Arie Shoshani, Matthew Wolf, Kesheng Wu, Weikuan Yu |
Concurr. Comput. Pract. Exp. | 17 |
| 2014 | Design and Evaluation of Network-Levitated Merge for Hadoop AccelerationabstractHadoop is a popular open source implementation of the MapReduce programming model for cloud computing. However, it faces a number of issues to achieve the best performance from the underlying systems. These include a serialization barrier that delays the reduce phase, repetitive merges, and disk accesses, and the lack of portability to different interconnects. To keep up with the increasing volume of data sets, Hadoop also requires efficient I/O capability from the underlying computer systems to process and analyze data. We describe Hadoop-A, an acceleration framework that optimizes Hadoop with plug-in components for fast data movement, overcoming the existing limitations. A novel network-levitated merge algorithm is introduced to merge data without repetition and disk access. In addition, a full pipeline is designed to overlap the shuffle, merge, and reduce phases. Our experimental results show that Hadoop-A significantly speeds up data movement in MapReduce and doubles the throughput of Hadoop. In addition, Hadoop-A significantly reduces disk accesses caused by intermediate data. Weikuan Yu, Yandong Wang 0001, Xinyu Que |
IEEE Trans. Parallel Distributed Syst. | 1 |
| 2013 | Exploring hybrid memory for GPU energy efficiency through software-hardware co-designabstractHybrid memory designs, such as DRAM plus Phase Change Memory (PCM), have shown some promise for alleviating power and density issues faced by traditional memory systems. But previous studies have concentrated on CPU systems with a modest level of parallelism. This work studies the problem in a massively parallel setting. Specifically, it investigates the special implications to hybrid memory imposed by the massive parallelism in GPU. It empirically shows that, contrary to promising results demonstrated for CPU, previous designs of PCM-based hybrid memory result in significant degradation to the energy efficiency of GPU. It reveals that the fundamental reason comes from a multi-facet mismatch between those designs and the massive parallelism in GPU. It presents a solution that centers around a close cooperation between compiler-directed data placement and hardware-assisted runtime adaptation. The co-design approach helps tap into the full potential of hybrid memory for GPU without requiring dramatic hardware changes over previous designs, yielding 6% and 49% energy saving on average compared to pure DRAM and pure PCM respectively, and keeping performance loss less than 2%. Bin Wang 0019, Bo Wu 0002, Dong Li 0001, Xipeng Shen, Weikuan Yu, Yizheng Jiao, Jeffrey S. Vetter |
PACT | 5 |
| 2013 | SLOAVx: Scalable LOgarithmic AlltoallV Algorithm for Hierarchical Multicore SystemsabstractScientific applications use collective communication operations in Message Passing Interface (MPI) for global synchronization and data exchanges. Alltoall and AlltoallV are two important collective operations. They are used by MPI jobs to exchange messages among all of MPI processes. AlltoallV is a generalization of Alltoall, supporting messages of varying sizes. However, the existing MPI AlltoallV implementation has linear complexity, i.e., each process has to send messages to all other processes in the job. Such linear complexity can result in sub optimal scalability of MPI applications when they are deployed on millions of cores. To address above challenge, in this paper, we introduce a new Scalable LOgarithmic AlltoallV algorithm, named SLOAV, for MPI AlltoallV collective operation. SLOAV aims to achieve global exchange of small messages of different sizes in a logarithmic number of rounds. Furthermore, given the prevalence of multicore systems with shared memory, we design a hierarchical AlltoallV algorithm based on SLOAV by leveraging the advantages of shared memory, which is referred to as SLOAVx. Compared to SLOAV, SLOAVx significantly reduces the inter-node communication, thus improving the entire system performance and mitigating the impact of message latency. We have implemented and embedded both algorithms in Open MPI. Our evaluation on large-scale computer systems shows that for the 8-byte and 1024-process MPI Alltoallv operation, the SLOAV can reduce the latency by as much as 86.4%, when compared to the state-of-the-art, and SLOAVx can further optimize the SLOAV by up to 83.1% in terms of message latency on multicore systems. In addition, experiments with NAS Parallel Benchmark (NPB) demonstrate that our algorithms are very effective for real-world applications. Cong Xu 0008, Manjunath Gorentla Venkata, Richard L. Graham, Yandong Wang 0001, Weikuan Yu |
CCGRID | 6 |
| 2013 | A case of system-wide power management for scientific applicationsabstractThe advance of high-performance computing systems towards exascale will be constrained by the systems' energy consumption levels. Large numbers of processing components, memory, interconnects, and storage components must all be considered to achieve exascale performance within a targeted energy bound. While application-aware power allocation schemes for computing resources are well studied, a portable and scalable budget-constrained power management scheme for scientific applications on exascale systems is still required. Execution activities within scientific applications can be categorized as CPU-bound, I/O-bound and communication-bound. Such activities tend to be clustered into ‘phases’, offering opportunities to manage their power consumption separately. Our experiments have demonstrated that their performance and energy consumption are affected differently by CPU frequency, an opportunity to fine tune CPU frequency for a minimal impact on the total execution time but significant savings on the energy consumption. By exploiting this opportunity, we present a phase-aware hierarchical power management framework that can opportunistically deliver good tradeoffs between system power consumption and application performance under a power budget. Our hierarchical power management framework consists of two main techniques: Phase-Aware CPU Frequency Scaling (PAFS) and opportunistic provisioning for power-constrained performance optimization. We have performed a systematic evaluation using both simulations and representative scientific applications on real systems. Our results show that our techniques can achieve 4.3%–17% better energy efficiency for large-scale scientific applications. Jay F. Lofstead, Teng Wang 0001, Weikuan Yu |
CLUSTER | 4 |
| 2013 | Profiling and Improving I/O Performance of a Large-Scale Climate Scientific ApplicationabstractExascale computing systems are soon to emerge, which will pose great challenges on the huge gap between computing and I/O performance. Many large-scale scientific applications play an important role in our daily life. The huge amounts of data generated by such applications require highly parallel and efficient I/O management policies. In this paper, we adopt a mission-critical scientific application, GEOS-5, as a case to profile and analyze the communication and I/O issues that are preventing applications from fully utilizing the underlying parallel storage systems. Through in-detail architectural and experimental characterization, we observe that current legacy I/O schemes incur significant network communication overheads and are unable to fully parallelize the data access, thus degrading applications' I/O performance and scalability. To address these inefficiencies, we redesign its I/O framework along with a set of parallel I/O techniques to achieve high scalability and performance. Evaluation results on the NASA discover cluster show that our optimization of GEOS- 5 with ADIOS has led to significant performance improvements compared to the original GEOS-5 implementation. Bin Wang 0019, Teng Wang 0001, Yuan Tian 0004, Cong Xu 0008, Yandong Wang 0001, Weikuan Yu, Carlos A. Cruz, Shujia Zhou, Thomas L. Clune, Scott Klasky |
ICCCN | 7 |
| 2013 | JVM-Bypass for Efficient Hadoop ShufflingabstractHadoop employs Java-based network transport stack on top of the Java Virtual Machine (JVM) for its data shuffling and merging purposes. Our examination reveals that JVM introduces a significant amount of overhead to data processing capability of the native interface. Furthermore, JVM constrains the use of high-performance networking mechanisms such as RDMA (Remote Direct Memory Access) which has established itself as an effective data movement technology in many networking environments because of its low-latency, high bandwidth, low CPU utilization, and energy efficiency. In this paper, we introduce a plug-in library called JVM-Bypass Shuffling (JBS) for Hadoop data shuffling. JBS helps Hadoop data shuffling by avoiding Javabased transport protocols, removing the overhead and limitations of the JVM. In addition, we design JBS as a portable library that can leverage both TCP/IP and RDMA on different network systems such as InfiniBand and 1/10 Gigabit Ethernet. We have designed and implemented JBS as part of Hadoop acceleration. It has been transferred to Mellanox as the software product UDA (Unstructured Data Accelerator) and used to enable our studies on a variety of merging algorithms. Our performance evaluation demonstrates that JBS can effectively reduce the execution time of Hadoop jobs by up to 66.3% and lower the CPU utilization by 48.1%. Yandong Wang 0001, Cong Xu 0008, Weikuan Yu |
IPDPS | 4 |
| 2013 | A Versatile Performance and Energy Simulation Tool for Composite GPU Global MemoryabstractAs a cost-effective compute device, Graphic Processing Unit (GPU) has been widely embraced in the field of high performance computing. GPU is characterized by its massive thread-level parallelism and high memory bandwidth. Although GPU has exhibited tremendous potential, recent GPU architecture researches mainly focus on GPU compute units and full system exploration is rare due to the lack of accurate simulators that can reveal hardware organization of both GPU compute units and its memory system. In order to fill this void, we build a GPU simulator called VxGPUSim that can support the simulation with detailed performance, timing and power consumption statistics. Our experimental evaluation demonstrates that VxGPUSim can faithfully reveal the internal execution details of GPU global memory of various memory configurations. It can enable further research on the design of GPU global memory for performance and energy tradeoffs. Bin Wang 0019, Yizheng Jiao, Weikuan Yu, Xipeng Shen, Dong Li 0001, Jeffrey S. Vetter |
MASCOTS | 3 |
| 2013 | A lightweight I/O scheme to facilitate spatial and temporal queries of scientific data analyticsabstractIn the era of petascale computing, more scientific applications are being deployed on leadership scale computing platforms to enhance the scientific productivity. Many I/O techniques have been designed to address the growing I/O bottleneck on large-scale systems by handling massive scientific data in a holistic manner. While such techniques have been leveraged in a wide range of applications, they have not been shown as adequate for many mission critical applications, particularly in data postprocessing stage. One of the examples is that some scientific applications generate datasets composed of a vast amount of small data elements that are organized along many spatial and temporal dimensions but require sophisticated data analytics on one or more dimensions. Including such dimensional knowledge into data organization can be beneficial to the efficiency of data post-processing, which is often missing from exiting I/O techniques. In this study, we propose a novel I/O scheme named STAR (Spatial and Temporal AggRegation) to enable high performance data queries for scientific analytics. STAR is able to dive into the massive data, identify the spatial and temporal relationships among data variables, and accordingly organize them into an optimized multi-dimensional data structure before storing to the storage. This technique not only facilitates the common access patterns of data analytics, but also further reduces the application turnaround time. In particular, STAR is able to enable efficient data queries along the time dimension, a practice common in scientific analytics but not yet supported by existing I/O techniques. In our case study with a critical climate modeling application GEOS-5, the experimental results on Jaguar supercomputer demonstrate an improvement up to 73 times for the read performance compared to the original I/O method. Yuan Tian 0004, Scott Klasky, Bin Wang 0019, Hasan Abbasi, Shujia Zhou, Norbert Podhorszki, Thomas L. Clune, Jeremy Logan, Weikuan Yu |
MSST | 10 |
| 2013 | DynaM: Dynamic Multiresolution Data Representation for Large-Scale Scientific AnalysisabstractFast growing large-scale systems enable scientific applications to run at a much larger scale and accordingly produce gigantic volumes of simulation output. Such data imposes a grand challenge to post-processing tasks such as visualization and data analysis, because these tasks are often performed at a host machine that is remotely located and equipped with much less memory and storage resources. During the simulation runs, it is also desirable for scientists to be able to interactively monitor and steer the progress of simulation. This requires scientific data to be represented in an efficient form for initial exploration and computation steering. In this paper, we propose DynaM a software framework that can represent scientific data in a multiresolution form, and dynamically organize data blocks into an optimized layout for efficient scientific analysis. DynaM supports a convolution-based multiresolution data representation for abstracting scientific data for visualization at a wide spectrum of resolution. To support the efficient generation and retrieval of different data granularities from such representation, a dynamic data organization in DynaM is enabled to cater distinct peculiarities of different size data blocks for efficient and balanced I/O performance. Our experimental results demonstrate that DynaM can efficiently represent large scientific dataset and speed up the visualization of multidimensional scientific data. An up to 29 times speedup is achieved on Jaguar supercomputer at Oak Ridge National Laboratory. Yuan Tian 0004, Scott Klasky, Weikuan Yu, Bin Wang 0019, Hasan Abbasi, Norbert Podhorszki, Ray W. Grout |
NAS | 3 |
| 2013 | CooMR: cross-task coordination for efficient data management in MapReduce programsabstractHadoop is a widely adopted open source implementation of MapReduce programming model for big data processing. It represents system resources as available map and reduce slots and assigns them to various tasks. This execution model gives little regard to the need of cross-task coordination on the use of shared system resources on a compute node, which results in task interference. In addition, the existing Hadoop merge algorithm can cause excessive I/O. In this study, we undertake an effort to address both issues. Accordingly, we have designed a cross-task coordination framework called CooMR for efficient data management in MapReduce programs. CooMR consists of three component schemes including cross-task opportunistic memory sharing and log-structured I/O consolidation, which are designed to facilitate task coordination, and the key-based in-situ merge (KISM) algorithm which is designed to enable the sorting/merging of Hadoop intermediate data without actually moving the pairs. Our evaluation demonstrates that CooMR is able to increase task coordination, improve system resource utilization, and significantly speed up the execution time of MapReduce programs. Yandong Wang 0001, Yizheng Jiao, Cong Xu 0008, Weikuan Yu |
SC | 5 |
| 2012 | A system-aware optimized data organization for efficient scientific analyticsabstractLarge-scale scientific applications on High End Computing systems produce a large volume of highly complex datasets. Such data imposes a grand challenge to conventional storage systems for the need of efficient I/O solutions during both the simulation runtime and data post-processing phases. With the mounting needs of scientific discovery, the read performance of large-scale simulations has becomes a critical issue for the HPC community. In this study, we propose a system-aware optimized data organization strategy that can organize data blocks of multidimensional scientific data efficiently based on simulation output and the underlying storage systems, thereby enabling efficient scientific analytics. Our experimental results demonstrate a performance speedup up to 72 times for the combustion simulation S3D, compared to the logically contiguous data layout. Yuan Tian 0004, Scott Klasky, Weikuan Yu, Hasan Abbasi, Bin Wang 0019, Norbert Podhorszki, Ray W. Grout, Matthew Wolf |
HPDC | 3 |
| 2012 | Identifying Opportunities for Byte-Addressable Non-Volatile Memory in Extreme-Scale Scientific ApplicationsabstractFuture exascale systems face extreme power challenges. To improve power efficiency of future HPC systems, non-volatile memory (NVRAM) technologies are being investigated as potential alternatives to existing memories technologies. NVRAMs use extremely low power when in standby mode, and have other performance and scaling benefits. Although previous work has explored the integration of NVRAM into various architecture and system levels, an open question remains: do specific memory workload characteristics of scientific applications map well onto NVRAMs' features when used in a hybrid NVRAM-DRAM memory system? Furthermore, are there common classes of data structures used by scientific applications that should be frequently placed into NVRAM?In this paper, we analyze several mission-critical scientific applications in order to answer these questions. Specifically, we develop a binary instrumentation tool to statistically report memory access patterns in stack, heap, and global data. We carry out hardware simulation to study the impact of NVRAM for both memory power and system performance. Our study identifies many opportunities for using NVRAM for scientific applications. In two of our applications, 31% and 27% of the memory working sets are suitable for NVRAM. Our simulations suggest at least 27% possible power savings and reveal that the performance of some applications is insensitive to relatively long NVRAM write-access latencies. Dong Li 0001, Jeffrey S. Vetter, Gabriel Marin, Collin McCurdy, Cristian Cira, Weikuan Yu |
IPDPS | 7 |
| 2012 | PCM-Based Durable Write Cache for Fast Disk I/OabstractFlash based solid-state devices (FSSDs) have been adopted within the memory hierarchy to improve the performance of hard disk drive (HDD) based storage system. However, with the fast development of storage-class memories, new storage technologies with better performance and higher write endurance than FSSDs are emerging, e.g., phase-change memory (PCM). Understanding how to leverage these state-of the-art storage technologies for modern computing systems is important to solve challenging data intensive computing problems. In this paper, we propose to leverage PCM for a hybrid PCM-HDD storage architecture. We identify the limitations of traditional LRU caching algorithms for PCMbased caches, and develop a novel hash-based write caching scheme called HALO to improve random write performance of hard disks. To address the limited durability of PCM devices and solve the degraded spatial locality in traditional wear-leveling techniques, we further propose novel PCM management algorithms that provide effective wear-leveling while maximizing access parallelism. We have evaluated this PCM-based hybrid storage architecture using applications with a diverse set of I/O access patterns. Our experimental results demonstrate that the HALO caching scheme leads to an average reduction of 36.8% in execution time compared to the LRU caching scheme, and that the SFC wear leveling extends the lifetime of PCM by a factor of 21.6. Bin Wang 0019, Patrick Carpenter, Dong Li 0001, Jeffrey S. Vetter, Weikuan Yu |
MASCOTS | 6 |
| 2012 | SMART-IO: SysteM-AwaRe Two-Level Data Organization for Efficient Scientific AnalyticsabstractCurrent I/O techniques have pushed the write performance close to the system peak, but they usually overlook the read side of problem. With the mounting needs of scientific discovery, it is important to provide good read performance for many common access patterns. Such demand requires an organization scheme that can effectively utilize the underlying storage system. However, the mismatch between conventional data layout on disk and common scientific access patterns leads to significant performance degradation when a subset of data is accessed. To this end, we design a system-aware Optimized Chunking model, which aims to find an optimized organization that can strike for a good balance between data transfer efficiency and processing overhead. To enable such model for scientific applications, we propose SMART-IO, a two-level data organization framework that can organize the blocks of multidimensional data efficiently. This scheme can adapt data layouts based on data characteristics and underlying storage systems, and enable efficient scientific analytics. Our experimental results demonstrate that SMART-IO can significantly improve the read performance for challenging access patterns, and speed up data analytics. For a mission critical combustion simulation code S3D, Smart-IO achieves up to 72 times speedup for planar reads of a 3-D variable compared to the logically contiguous data layout. Yuan Tian 0004, Scott Klasky, Weikuan Yu, Hasan Abbasi, Bin Wang 0019, Norbert Podhorszki, Ray W. Grout, Matthew Wolf |
MASCOTS | 3 |
| 2012 | Classifying soft error vulnerabilities in extreme-scale scientific applications using a binary instrumentation toolabstractExtreme-scale scientific applications are at a significant risk of being hit by soft errors on supercomputers as the scale of these systems and the component density continues to increase. In order to better understand the specific soft error vulnerabilities in scientific applications, we have built an empirical fault injection and consequence analysis tool - BIFIT -that allows us to evaluate how soft errors impact applications. In particular, BIFIT is designed with capability to inject faults at very specific targets: an arbitrarily-chosen execution point and any specific data structure. We apply BIFIT to three mission-critical scientific applications and investigate the applications vulnerability to soft errors by performing thousands of statistical tests. We, then, classify each applications individual data structures based on their sensitivity to these vulnerabilities, and generalize these classifications across applications. Subsequently, these classifications can be used to apply appropriate resiliency solutions to each data structure within an application. Our study reveals that these scientific applications have a wide range of sensitivities to both the time and the location of a soft error; yet, we are able to identify intrinsic relationships between application vulnerabilities and specific types of data objects. In this regard, BIFIT enables new opportunities for future resiliency research. Dong Li 0001, Jeffrey S. Vetter, Weikuan Yu |
SC | 3 |
| 2012 | HiCOO: Hierarchical cooperation for scalable communication in Global Address Space programming models on Cray XT systems
Weikuan Yu, Xinyu Que, Vinod Tipparaju, Jeffrey S. Vetter |
J. Parallel Distributed Comput. | 1 |
| 2011 | Network-Friendly One-Sided Communication through Multinode Cooperation on Petascale Cray XT5 SystemsabstractOne-sided communication is important to enable asynchronous communication and data movement for Global Address Space (GAS) programming models. Such communication is typically realized through direct messages between initiator and target processes. For peta scale systems with 10,000s of nodes and 100,000s of cores, these direct messages require dedicated communication buffers and/or channels, which can lead to significant scalability challenges for GAS programming models. In this paper, we describe a network-friendly communication model, multinode cooperation, to enable indirect one-sided communication. Compute nodes work together to handle one-side requests through (1) request forwarding in which one node can intercept a request and forward it to a target node, and (2) request aggregation in which one node can aggregate many requests to a target node. We have implemented multinode cooperation for a popular GAS runtime library, Aggregate Remote Memory Copy Interface (ARMCI). Our experimental results on a large scale Cray XT5 system demonstrate that multinode cooperations able to greatly increase memory scalability by reducing communication buffers required on each node. In addition, multinode cooperation improves the resiliency of GAS runtime system to network contention. Furthermore, multinode cooperation can benefit the performance of scientific applications. In one case, it reduces the total execution time of an NWChem application by 52%. Xinyu Que, Weikuan Yu, Vinod Tipparaju, Jeffrey S. Vetter, Bin Wang 0019 |
CCGRID | 2 |
| 2011 | EDO: Improving Read Performance for Scientific Applications through Elastic Data OrganizationabstractLarge scale scientific applications are often bottlenecked due to the writing of checkpoint-restart data. Much work has been focused on improving their write performance. With the mounting needs of scientific discovery from these datasets, it is also important to provide good read performance for many common access patterns, which requires effective data organization. To address this issue, we introduce Elastic Data Organization (EDO), which can transparently enable different data organization strategies for scientific applications. Through its flexible data ordering algorithms, EDO harmonizes different access patterns with the underlying file system. Two levels of data ordering are introduced in EDO. One works at the level of data groups (a.k.a process groups). It uses Hilbert Space Filling Curves (SFC) to balance the distribution of data groups across storage targets. Another governs the ordering of data elements within a data group. It divides a data group into sub chunks and strikes a good balance between the size of sub chunks and the number of seek operations. Our experimental results demonstrate that EDO is able to achieve balanced data distribution across all dimensions and improve the read performance of multidimensional datasets in scientific applications. Yuan Tian 0004, Scott Klasky, Hasan Abbasi, Jay F. Lofstead, Ray W. Grout, Norbert Podhorszki, Qing Liu 0002, Yandong Wang 0001, Weikuan Yu |
CLUSTER | 9 |
| 2011 | BMF: Bitmapped Mass Fingerprinting for Fast Protein IdentificationabstractProtein identification is an important objective for proteomic and medical sciences as well as for pharmaceutical industry. With recent large-scale automation of genome sequencing and the explosion of protein databases, it is important to exploit latest data processing technologies and design highly scalable algorithms to speed up protein identification. In this study, we design, implement, and evaluate a new software tool, Bitmapped Mass Fingerprinting (BMF), that can efficiently construct a bitmap index for short peptides, and quickly identify candidate proteins from leading protein databases. BMF is developed by integrating the Fast Bit indexing technology and the popular Message Passing Interface (MPI) for parallelization. By exploiting Fast Bit for peptide mass fingerprinting across protein boundaries, we are able to accomplish parallel computation and I/O for a scalable implementation of protein identification. Our experimental results show that BMF brings dramatic performance improvement for protein identification from various protein databases. In particular, we demonstrate that BMF can effectively scale up to 8,192 cores on the Jaguar Supercomputer at Oak Ridge National Laboratory, achieving superb performance in identifying proteins from the National Center for Biotechnology Information (NCBI) non-redundant (NR) protein database. Weikuan Yu, Kesheng Wu, Wei-Shinn Ku, Cong Xu 0008, Juan Gao |
CLUSTER | 1 |
| 2011 | Virtual Topologies for Scalable Resource Management and Contention Attenuation in a Global Address Space Model on the Cray XT5abstractGlobal Address Space (GAS) programming models enable a convenient, shared-memory style addressing model, and support completely asynchronous data movement. Their underlying runtime systems face critical challenges in (1) scalably managing resources (such as memory for communication buffers), and (2) gracefully handling unpredictable communication patterns and any associated contention. In this research, we investigate these challenges for a popular GAS runtime library, Aggregate Remote Memory Copy Interface (ARMCI) on, large-scale Cray XT5 systems. We represent the management of communication resources as directed graphs, and propose two new scalable virtual topologies, Meshed Fully Connected Graphs (MFCG) and Cubic Fully Connected Graphs (CFCG), for scalable resource management and contention attenuation. To ensure deadlock-free communication in these multi-dimensional topologies, we design and develop Lowest Dimension First (LDF) forwarding to support fully- or partially-populated MFCG and CFCG on any number of nodes. We have extensively evaluated the benefits of these virtual topologies on the petascale Jaguar Cray XT5 system at Oak Ridge National Laboratory. Our experimental results demonstrate MFCG as the most suitable virtual topology because of its benefits in resource management, contention mitigation, and the resulting benefit to scientific applications. Weikuan Yu, Vinod Tipparaju, Xinyu Que, Jeffrey S. Vetter |
ICPP | 1 |
| 2011 | Hadoop acceleration through network levitated mergeabstractHadoop is a popular open-source implementation of the MapReduce programming model for cloud computing. However, it faces a number of issues to achieve the best performance from the underlying system. These include a serialization barrier that delays the reduce phase, repetitive merges and disk access, and lack of capability to leverage latest high speed interconnects. We describe Hadoop-A, an acceleration framework that optimizes Hadoop with plugin components implemented in C++ for fast data movement, overcoming its existing limitations. A novel network-levitated merge algorithm is introduced to merge data without repetition and disk access. In addition, a full pipeline is designed to overlap the shuffle, merge and reduce phases. Our experimental results show that Hadoop-A doubles the data processing throughput of Hadoop, and reduces CPU utilization by more than 36%. Yandong Wang 0001, Xinyu Que, Weikuan Yu, Dror Goldenberg, Dhiraj Sehgal |
SC | 3 |
| 2009 | Design, implementation, and evaluation of transparent pNFS on LustreabstractParallel NFS (pNFS) is an emergent open standard for parallelizing data transfer over a variety of I/O protocols. Prototypes of pNFS are actively being developed by industry and academia to examine its viability and possible enhancements. In this paper, we present the design, implementation, and evaluation of lpNFS, a Lustre-based parallel NFS. We achieve our primary objective in designing lpNFS as an enabling technology for transparent pNFS accesses to an opaque Lustre file system. We optimize the data flow paths in lpNFS by using two techniques: (a) fast memory coping for small messages, and (b) page sharing for zero-copy bulk data transfer. Our initial performance evaluation shows that the performance of lpNFS is comparable to that of original Lustre. Given these results, we assert that lpNFS is a promising approach to combining the benefits of pNFS and Lustre, and it exposes the underlying capabilities of Lustre file systems while transparently supporting pNFS clients. Weikuan Yu, Oleg Drokin, Jeffrey S. Vetter |
IPDPS | 1 |
| 2008 | Xen-Based HPC: A Parallel I/O PerspectiveabstractVirtualization using Xen-based virtual machine environment has yet to permeate the field of high performance computing (HPC). One major requirement for HPC is the availability of scalable and high performance I/O. Conventional wisdom suggests that virtualization of system services must lead to degraded performance. In this presentation, we take on a parallel I/O perspective to study the viability of Xen-based HPC for data-intensive programs. We have analyzed the overheads and migration costs for parallel I/O programs in a Xen-based virtual machine cluster. Our analysis covers PVFS-based parallel I/O over two different networking protocols: TCP-based Gigabit Ethernet and VMM-bypass InfiniBand. Our experimental results suggest that network processing in Xen-based virtualization can significantly impact the performance of Parallel I/O. By carefully tuning the networking layers, we have demonstrated the following for Xen-based HPC I/O: (I) TCP offloading can help achieve low overhead parallel I/O; (2) parallel reads and writes require different network tuning to achieve good I/O bandwidth; and (3) Xen-based HPC environment can support high performance parallel I/O with both negligible overhead and little migration cost. Weikuan Yu, Jeffrey S. Vetter |
CCGRID | 1 |
| 2008 | Empirical Analysis of a Large-Scale Hierarchical Storage System
Weikuan Yu, Sarp Oral, Shane Canon, Jeffrey S. Vetter, Ramanan Sankaran |
Euro-Par | 1 |
| 2008 | ParColl: Partitioned Collective I/O on the Cray XTabstractCollective I/O orchestrates I/O from parallel processes by aggregating fine-grained requests into large ones. However, its performance is typically a fraction of the potential I/O bandwidth on large scale platforms such as Cray XT. Based on our analysis, the time spent in global process synchronization dominates the actual time in file reads/writes, which imposes a 'collective wall' on the performance of collective I/O. In this paper, we introduce a novel technique called partitioned collective I/O (ParColl). ParColl augments the original two-phase collective I/O protocol with new mechanisms for file area partitioning, I/O aggregator distribution and intermediate file views. Through these mechanisms, a group of processes and their targeted file are consistently divided into a collection of small subgroups, each performing I/O aggregation in a disjoint manner. File consistency is maintained through intermediate file views when necessary. Together, these mechanisms greatly reduce the cost of global synchronization. Our experimental results demonstrate that ParColl significantly improves the performance and the scalability of collective I/O. In one case, we show a 416% improvement on 1024 processes for a visualization I/O benchmark. We also show that the I/O patterns in scientific applications can benefit significantly from this technique, e.g. BT-I/O and Flash I/O. Weikuan Yu, Jeffrey S. Vetter |
ICPP | 1 |
| 2008 | Performance characterization and optimization of parallel I/O on the Cray XTabstractThis paper presents an extensive characterization, tuning, and optimization of parallel I/O on the Cray XT supercomputer, named Jaguar, at Oak Ridge National Laboratory. We have characterized the performance and scalability for different levels of storage hierarchy including a single Lustre object storage target, a single S2A storage couplet, and the entire system. Our analysis covers both data- and metadata-intensive I/O patterns. In particular, for small, non-contiguous data- intensive I/O on Jaguar, we have evaluated several parallel I/O techniques, such as data sieving and two- phase collective I/O, and shed light on their effectiveness. Based on our characterization, we have demonstrated that it is possible, and often prudent, to improve the I/O performance of scientific benchmarks and applications by tuning and optimizing I/O. For example, we demonstrate that the I/O performance of the S3D combustion application can be improved at large scale by tuning the I/O system to avoid a bandwidth degradation of 49% with 8192 processes when compared to 4096 processes. We have also shown that the performance of Flash I/O can be improved by 34% by tuning the collective I/O parameters carefully. Weikuan Yu, Jeffrey S. Vetter, Sarp Oral |
IPDPS | 1 |
| 2008 | Early evaluation of IBM BlueGene/PabstractBlueGene/P (BG/P) is the second generation BlueGene architecture from IBM, succeeding BlueGene/L (BG/L). BG/P is a system-on-a-chip (SoC) design that uses four PowerPC 450 cores operating at 850 MHz with a double precision, dual pipe floating point unit per core. These chips are connected with multiple interconnection networks including a 3-D torus, a global collective network, and a global barrier network. The design is intended to provide a highly scalable, physically dense system with relatively low power requirements per flop. In this paper, we report on our examination of BG/P, presented in the context of a set of important scientific applications, and as compared to other major large scale supercomputers in use today. Our investigation confirms that BG/P has good scalability with an expected lower performance per processor when compared to the Cray XT4's Opteron. We also find that BG/P uses very low power per floating point operation for certain kernels, yet it has less of a power advantage when considering science driven metrics for mission applications. Sadaf R. Alam, Richard F. Barrett, M. Bast, Mark R. Fahey, Jeffery A. Kuehn, Collin McCurdy, James H. Rogers, Philip C. Roth, Ramanan Sankaran, Jeffrey S. Vetter, Patrick H. Worley, Weikuan Yu |
SC | 12 |
| 2008 | Wide-area performance profiling of 10GigE and InfiniBand technologiesabstractFor wide-area high-performance applications, light-paths provide 10Gbps connectivity, and multi-core hosts with PCI-Express can drive such data rates. However, sustaining such end-to-end application throughputs across connections of thousands of miles remains challenging, and the current performance studies of such solutions are very limited. We present an experimental study of two solutions to achieve such throughputs based on: (a) 10Gbps Ethernet with TCP/IP transport protocols, and (b) InfiniBand and its wide-area extensions. For both, we generate performance profiles over 10Gbps connections of lengths up to 8600 miles, and discuss the components, complexity, and limitations of sustaining such throughputs, using different connections and host configurations. Our results indicate that IB solution is better suited for applications with a single large flow, and 10GigE solution is better for those with multiple competing flows. Nageswara S. V. Rao, Weikuan Yu, William R. Wing, Stephen W. Poole, Jeffrey S. Vetter |
SC | 2 |
| 2007 | Exploiting Lustre File Joining for Effective Collective IOabstractLustre is a parallel file system that presents high aggregated IO bandwidth by striping file extents across many storage devices. However, our experiments indicate excessively wide striping can cause performance degradation. Lustre supports an innovative file joining feature that joins files in place. To mitigate striping overhead and benefit collective IO, we propose two techniques: split writing and hierarchical striping. In split writing, a file is created as separate subfiles, each of which is striped to only a few storage devices. They are joined as a single file at the file close time. Hierarchical striping builds on top of split writing and orchestrates the span of subfiles in a hierarchical manner to avoid overlapping and achieve the appropriate coverage of storage devices. Together, these techniques can avoid the overhead associated with large stripe width, while still being able to combine bandwidth available from many storage devices. We have prototyped these techniques in the ROMIO implementation of MPI-IO. Experimental results indicate that split writing and hierarchical striping can significantly improve the performance of Lustre collective IO in terms of both data transfer and management operations. On a Lustre file system configured with 46 object storage targets, our implementation improves collective write performance of a 16-process job by as much as 220%. Weikuan Yu, Jeffrey S. Vetter, Shane Canon, Song Jiang 0001 |
CCGRID | 1 |
| 2007 | FlexFetch: A History-Aware Scheme for I/O Energy Saving in Mobile ComputingabstractExtension of battery lifetime has always been a major issue for mobile computing. While more and more data are involved in mobile computing, energy consumption caused by I/O operations becomes increasingly large. In a pervasive computing environment, the requested data can be stored both on the local disk of a mobile computer by using the hoarding technique, and on the remote server, where data are accessible via wireless communication. Based on the current operational states of local disk (active or standby), the amount of data to be requested (small or large), and currently available wireless bandwidth (strong or weak reception), data access source can be adaptively selected to achieve maximum energy reduction. To this end, we propose a profile-based I/O management scheme, FlexFetch, that is aware of access history and adaptive to current access environment. Our simulation experiments driven by real-life traces demonstrate that the scheme can significantly reduce energy consumption in a mobile computer compared with existing representative schemes. Feng Chen 0005, Song Jiang 0001, Weisong Shi, Weikuan Yu |
ICPP | 4 |
| 2006 | Application-Transparent Checkpoint/Restart for MPI Programs over InfiniBandabstractUltra-scale computer clusters with high speed interconnects, such as InfiniBand, are being widely deployed for their excellent performance and cost effectiveness. However, the failure rate on these clusters also increases along with their augmented number of components. Thus, it becomes critical for such systems to be equipped with fault tolerance support. In this paper, we present our design and implementation of checkpoint/restart framework for MPI programs running over InfiniBand clusters. Our design enables low-overhead, application-transparent checkpointing. It uses coordinated protocol to save the current state of the whole MPI job to reliable storage, which allows users to perform rollback recovery if the system runs into faulty states later. Our solution has been incorporated into MVAPICH2, an open-source high performance MPI-2 implementation over InfiniBand. Performance evaluation of this implementation has been carried out using NAS benchmarks, HPL benchmark, and a real-world application called GROMACS. Experimental results indicate that in our design, the overhead to take checkpoints is low, and the performance impact for checkpointing applications periodically is insignificant. For example, time for checkpointing GROMACS is less than 0.3% of the execution time, and its performance only decreases by 4% with checkpoints taken every minute. To the best of our knowledge, this work is the first report of checkpoint/restart support for MPI over InfiniBand clusters in the literature Qi Gao 0004, Weikuan Yu, Wei Huang 0003, Dhabaleswar K. Panda 0001 |
ICPP | 2 |
| 2006 | High Performance Block I/O for Global File System (GFS) with InfiniBand RDMAabstractState-of-the-art network technology has evolved to 10Gbps. However, TCP's high processing overhead and redundant data copies remain a major bottleneck for applications to fully benefit from such high speed technology. Remote direct memory access (RDMA), as an emerging communication protocol, provides an opportunity for efficient storage system design by virtue of RDMA's semantics. Although RDMA based designs have been proposed to improve network file I/O protocols in several previous works, its benefit for cluster file system block I/O is not clear yet. We propose a new technique - "buffer management delegation", which offloads message buffer management to remote communication party. Using this technique, we design our zero copy RDMA based block transfer scheme for GNBD (global network block device), a block access protocol of Red Hat Global File System, to optimize cluster file system performance over 10Gbps InfiniBand network. We evaluate this new scheme in comparison with our copy based scheme and TCP over the same InfiniBand hardware. The evaluation quantifies the redundant copy impact for both bulk data transfer and file system meta-data operations. The results using open source file system benchmarks and widely used system utilities show that our implementation improves GFSperformance by up to 47% compared with copy based scheme, and by up to 136% compared with TCP Weikuan Yu, Dhabaleswar K. Panda 0001 |
ICPP | 2 |
| 2006 | Adaptive connection management for scalable MPI over InfiniBandabstractSupporting scalable and efficient parallel programs is a major challenge in parallel computing with the widespread adoption of large-scale computer clusters and supercomputers. One of the pronounced scalability challenges is the management of connections between parallel processes, especially over connection-oriented interconnects such as VIA and InfiniBand. In this paper, we take on the challenge of designing efficient connection management for parallel programs over InfiniBand clusters. We propose adaptive connection management (ACM) to dynamically control the establishment of InfiniBand reliable connections (RC) based on the communication frequency between MPI processes. We have investigated two different ACM algorithms: an on-demand algorithm that starts with no InfiniBand RC connections; and a partial static algorithm with only 2 * logN number of InfiniBand RC connections initially. We have designed and implemented both ACM algorithms in MVAPICH to study their benefits. Two mechanisms have been exploited for the establishment of new RC connections: one using InfiniBand unreliable datagram and the other using InfiniBand connection management. For both mechanisms, MPI communication issues, such as progress rules, reliability and race conditions are handled to ensure efficient and lightweight connection management. Our experimental results indicate that ACM algorithms can benefit parallel programs in terms of the process initiation time, the number of active connections, and the resource usage. For parallel programs on a 16-node cluster, they can reduce the process initiation time by 15% and the initial memory usage by 18%. Weikuan Yu, Qi Gao 0004, Dhabaleswar K. Panda 0001 |
IPDPS | 1 |
| 2006 | Benefits of high speed interconnects to cluster file systems: a case study with LustreabstractCluster file systems and storage area networks (SAN) make use of network IO to achieve higher IO bandwidth. Effective integration of networking mechanisms is important to their performance. In this paper, we perform an evaluation of a popular cluster file system, Lustre, over two of the leading high speed cluster interconnects: InfiniBand and Quadrics. Our evaluation is performed with both sequential IO and parallel IO benchmarks in order to explore the capacity of Lustre under different communication characteristics. Experimental results show that direct implementations of Lustre over both interconnects can improve its performance, compared to an IP emulation over InfiniBand (IPoIB). The performance of Lustre over Quadrics is comparable to that of Lustre over InfiniBand with the platforms we have. Latest InfiniBand products can embrace latest technologies, such as PCI-Express and DDR, and provide higher capacity. Our results show that over a Lustre file system with two object storage fervers (OSSs), InfiniBand with PCI-Express technology can improve Lustre write performance by 24%. Furthermore, our experimental results indicate that Lustre meta-data operations do not scale with an increasing number of OSSs, in spite of using high performance interconnects. Weikuan Yu, Ranjit Noronha, Dhabaleswar K. Panda 0001 |
IPDPS | 1 |
| 2005 | Head-to-TOE Evaluation of High-Performance Sockets over Protocol Offload EnginesabstractDespite the performance drawbacks of Ethernet, it still possesses a sizable footprint in cluster computing because of its low cost and backward compatibility to existing Ethernet infrastructure. In this paper, we demonstrate that these performance drawbacks can be reduced (and in some cases, arguably eliminated) by coupling TCP offload engines (TOEs) with 10-Gigabit Ethernet (10GigE). Although there exists significant research on individual network technologies such as 10GigE, InfiniBand (IBA), and Myrinet; to the best of our knowledge, there has been no work that compares the capabilities and limitations of these technologies with the recently introduced 10GigE TOEs in a homogeneous experimental testbed. Therefore, we present performance evaluations across 10GigE, IBA, and Myrinet (with identical cluster-compute nodes) in order to enable a coherent comparison with respect to the sockets interface. Specifically, we evaluate the network technologies at two levels: (i) a detailed micro-benchmark evaluation and (ii) an application-level evaluation with sample applications from different domains, including a bio-medical image visualization tool known as the Virtual Microscope, an iso-surface oil reservoir simulator, a cluster file-system known as the parallel virtual file-system (PVFS), and a popular cluster management tool known as Ganglia. In addition to 10GigE's advantage with respect to compatibility to wide-area network infrastructures, e.g., in support of grids, our results show that 10GigE also delivers performance that is comparable to traditional high-speed network technologies such as IBA and Myrinet in a system-area network environment to support clusters and that 10GigE is particularly well-suited for sockets-based applications Pavan Balaji, Wu-chun Feng, Qi Gao 0004, Ranjit Noronha, Weikuan Yu, Dhabaleswar K. Panda 0001 |
CLUSTER | 5 |
| 2005 | High performance support of parallel virtual file system (PVFS2) over QuadricsabstractParallel I/O needs to keep pace with the demand of high performance computing applications on systems with ever-increasing speed. Exploiting high-end interconnect technologies to reduce the network access cost and scale the aggregated bandwidth is one of the ways to increase the performance of storage systems. In this paper, we explore the challenges of supporting parallel file system with modern features of Quadrics, including user-level communication and RDMA operations. We design and implement a Quadrics-capable version of a parallel file system (PVFS2). Our design overcomes the challenges imposed by Quadrics static communication model to dynamic client/server architectures. Quadrics QDMA and RDMA mechanisms are integrated and optimized for high performance data communication. Zero-copy PVFS2 list IO is achieved with a Single Event Associated MUltiple RDMA (SEAMUR) mechanism. Experimental results indicate that the performance of PVFS2, with Quadrics user-level protocols and RDMA operations, is significantly improved in terms of both data transfer and management operations. With four IO server nodes, our implementation improves PVFS2 aggregated read bandwidth by up to 140% compared to PVFS2 over TCP on top of Quadrics IP implementation. Moreover, it delivers significant performance improvement to application benchmarks such as mpi-tile-io [24] and BTIO [26]. To the best of our knowledge, this is the first work in the literature to report the design of a high performance parallel file system over Quadrics user-level communication protocols. Weikuan Yu, Dhabaleswar K. Panda 0001 |
ICS | 1 |
| 2004 | Scalable, high-performance NIC-based all-to-all broadcast over Myrinet/GMabstractAll-to-all broadcast is one of the common collective operations that involve dense communication between all processes in a parallel program. Previously, programmable network interface cards (NICs) have been leveraged to efficiently support collective operations, including barrier, broadcast, and reduce. This work explores the characteristics of all-to-all broadcast and proposes new algorithms to exploit the potential advantages of NIC programmability. Along with these algorithms, salient strategies have been used to provide scalable topology management, global buffer management, efficient communication processing, and message reliability. The algorithms have been incorporated into a NIC-based collective protocol over Myrinet/GM. The NIC-based all-to-all broadcast operations improve all-to-all broadcast bandwidth over 16 nodes by a factor of 3, compared to host-based all-to-all broadcast operation. Furthermore, the NIC-based operations have been demonstrated to achieve better scalability to large systems and very low host CPU utilization. Weikuan Yu, Dhabaleswar K. Panda 0001, Darius Buntinas |
CLUSTER | 1 |
| 2004 | Fast and Scalable Startup of MPI Programs in InfiniBand Clusters
Weikuan Yu, Jiesheng Wu, Dhabaleswar K. Panda 0001 |
HiPC | 1 |
| 2004 | Efficient and Scalable Barrier over Quadrics and Myrinet with a New NIC-Based Collective Message Passing ProtocolabstractSummary form only given. Modern interconnects often have programmable processors in the network interface that can be utilized to offload communication processing from host CPU. We explore different schemes to support collective operations at the network interface and propose a new collective protocol. With barrier as an initial case study, we have demontrated that much of the communication processing can be greatly simplified with this collective protocol. Accordingly, we have designed and implemented efficient and scalable NIC-based barrier operations over two high performance interconnects, Quadrics and Myrinet. Our evaluation shows that, over a Quadrics cluster of 8 nodes with ELan3 network, the NIC-based barrier operation achieves a barrier latency of only 5.60/spl mu/s. This result is a 2.48 factor of improvement over the Elanlib tree-based barrier operation. Over a Myrinet cluster of 8 nodes with LANai-XP NIC cards, a barrier latency of 14.20/spl mu/s over 8 nodes is achieved. This is a 2.64 factor of improvement over the host-based barrier algorithm. Furthermore, an analytical model developed for the proposed scheme indicates that a NIC-based barrier operation on a 1024-node cluster can be performed with only 22.13/spl mu/s latency over Quadrics and with 38.94/spl mu/s latency over Myrinet. These results indicate the potential for developing high performance communication subsystems for next generation clusters. Weikuan Yu, Darius Buntinas, Richard L. Graham, Dhabaleswar K. Panda 0001 |
IPDPS | 1 |
| 2003 | High Performance and Reliable NIC-Based Multicast over Myrinet/GM-2abstractMulticast is an important collective operation for parallel programs. Some network interface cards (NICs), such as Myrinet, have programmable processors that can be programmed to support multicast. We propose a high performance and reliable NIC-based multicast scheme, in which a NIC-based multisend mechanism is used to send multiple replicas of a message to different destinations, and a NIC-based forwarding mechanism to forward the received packets without intermediate host involvement. We have explored different design alternatives and implemented the proposed scheme with the set of best alternatives over Myrinet/GM-2. MPICH-GM has also been modified to take advantage of this scheme. At the GM-level, the NIC-based multicast improves the multicast latency by a factor up to 1.48 for messages les 512 bytes, and a factor up to 1.86 for 16KB messages over 16 nodes compared to the traditional host-based multicast. Similar improvements are also achieved at the MPI level. In addition, it is demonstrated that NIC-based multicast is tolerant to process skew and has significant benefits for large systems Weikuan Yu, Darius Buntinas, Dhabaleswar K. Panda 0001 |
ICPP | 1 |
| 2003 | Performance Comparison of MPI Implementations over InfiniBand, Myrinet and QuadricsabstractIn this paper, we present a comprehensive performance comparison of MPI implementations over Infini-Band, Myrinet and Quadrics. Our performance evaluation consists of two major parts. The first part consists of a set of MPI level micro-benchmarks that characterize different aspects of MPI implementations. The second part of the performance evaluation consists of application level benchmarks. We have used the NAS Parallel Benchmarks and the sweep3D benchmark. We not only present the overall performance results, but also relate application communication characteristics to the information we acquired from the micro-benchmarks. Our results show that the three MPI implementations all have their advantages and disadvantages. For our 8-node cluster, InfiniBand can offer significant performance improvements for a number of applications compared with Myrinet and Quadrics when using the PCI-X bus. Even with just the PCI bus, InfiniBand can still perform better if the applications are bandwidth-bound. Jiuxing Liu, B. Chandrasekaran 0001, Jiesheng Wu, Weihang Jiang, Sushmitha P. Kini, Weikuan Yu, Darius Buntinas, Pete Wyckoff, Dhabaleswar K. Panda 0001 |
SC | 6 |