Wu-chun Feng

dblp:82/6853 · also Wu-Chun Feng · DBLP profile ↗
← Back
173ranked-venue papers
20as first author
15since 2021 · last 2026
0000-0002-6015-0727ORCID · corroborated

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

Systems, architecture and hardware · 129 · 13 first-author · 12 since 2021Computer networks · 16 · 3 first-author · 1 since 2021Applied, interdisciplinary, general and emerging computing · 11 · 2 first-author · 2 since 2021Software engineering, systems software and programming languages · 9 · 3 first-authorSecurity and privacy · 3Human-computer interaction and ubiquitous computing · 3 · 1 first-authorArtificial intelligence and machine learning · 2 · 1 first-authorDatabases, data management, data science and information retrieval · 1Graphics, computer vision, multimedia, augmented reality and games · 1
YearPublicationVenuePosition
2026 On the Efficacy of PyTorch for High-Performance Computing: A Case Study in Computational Physics
abstract
Python has become a popular language for scientific computing due to its productivity and ease of use, with frameworks like PyTorch delivering portable performance without exposing users to the complexity of high-performance languages like C, C++, and Fortran. While PyTorch is well established in machine learning (ML), its suitability for non-ML scientific applications remains unclear.
Beau Johnston, Niteya Shah, Wu-chun Feng
CF3
2026 Looking for (Genomic) Needles in a Haystack: Sparsity-Driven Search for Identifying Correlated Genetic Mutations in Cancer
Ritvik Prabhu, Emil Vatai, Bernard Moussad, Emmanuel Jeannot, Ramu Anandakrishnan, Wu-chun Feng, Mohamed Wahib
IPDPS6
2026 Optimizing Management of Persistent Data Structures in High-Performance Analytics
abstract
Large-scale data analytics workflows ingest massive input data into various data structures, including graphs and key-value datastores. These data structures undergo multiple transformations and computations and are typically reused in incremental and iterative analytics workflows. Persisting in-memory views of these data structures enables reusing them beyond the scope of a single program run while avoiding repetitive raw data ingestion overheads. Memory-mapped I/O enables persisting in-memory data structures without data serialization and deserialization overheads. However, memory-mapped I/O lacks the key feature of persisting consistent snapshots of these data structures for incremental ingestion and processing. The obstacles to efficient virtual memory snapshots using memory-mapped I/O include background writebacks outside the application's control, and the significantly high storage footprint of such snapshots. To address these limitations, we presentPrivateer, a memory and storage management tool that enables storage-efficient virtual memory snapshotting while also optimizing snapshot I/O performance. We integratedPrivateerintoMetall, a state-of-the-art persistent memory allocator for C++, and the Lightning Memory-Mapped Database (LMDB), a widely-used key-value datastore in data analytics and machine learning.Privateeroptimized application performance by 1.22× when storing data structure snapshots to node-local storage, and up to 16.7× when storing snapshots to a parallel file system.Privateeralso optimizes storage efficiency of incremental data structure snapshots by up to 11× using data deduplication and compression.
Karim Youssef, Keita Iwabuchi, Maya B. Gokhale, Wu-chun Feng, Roger A. Pearce
IEEE Trans. Parallel Distributed Syst.4
2025 A Graph-Based Approach for Early and Explainable Health Risk Assessment
abstract
Early and explainable health risk assessment provides clinicians with an early indication of risk and explainability of its causes to perform proactive and informed medical intervention. However, clinical data is typically recorded as multivariate time-series (MTS) that are sparse, irregular, and misaligned, making risk assessment challenging. Most existing methods preprocess MTS to make it dense, regular, and aligned, but this (1) delays assessment due to preprocessing overhead and (2) compromises explainability by introducing synthetic values that erode clinical trust. Methods that operate on raw MTS often overlook the joint need for early detection and explainability. To address the above challenges, we present GrEEnHeRA (Graph-based approach for Early and Explainable Health Risk Assessment) that directly operates on raw MTS. As clinical observations arrive, GrEEnHeRA incrementally transforms a patient's MTS into a temporal sequence of graphs to enable early health risk assessment by serving as input to a spatio-temporal learning model: a graph isomorphism network (GIN) that captures intra-graph spatial structure and a gated recurrent unit (GRU) that models temporal dynamics across the graph sequence. Moreover, its design enables clinically meaningful explainability by identifying influential multi-variable correlations, temporal trends, and anomalous clinical observations that contributed to the risk. Using two PhysioNet Challenge datasets, we demonstrate that GrEEnHeRA is more accurate than existing methods and identifies risk hours before the onset of an adverse condition.
Sonal Jha, Wu-chun Feng
BIBE2
2025 Top-Down SBP: Turning Graph Clustering Upside Down
abstract
Stochastic block partitioning (SBP) is a statistical inference-based algorithm for clustering vertices within a graph. It has been shown to be statistically robust and highly accurate even on graphs with a complex structure, but its poor scalability limits its usability to smaller-sized graphs. In this manuscript we argue that one reason for its poor scalability is the agglomerative, or bottom-up, nature of SBP's algorithmic design; the agglomerative computations cause high memory usage and create a large search space that slows down statistical inference, particularly in the algorithm's initial iterations. To address this bottleneck, we propose Top-Down SBP, a novel algorithm that replaces the agglomerative (bottom-up) block merges in SBP with a block-splitting operation. This enables the algorithm to start with all vertices in one cluster and subdivide them over time into smaller clusters. We show that Top-Down SBP is up to 7.7× faster than Bottom-Up SBP without sacrificing accuracy and can process larger graphs than Bottom-Up SBP on the same hardware due to an up to 4.1× decrease in memory usage. Additionally, we adapt existing methods for accelerating Bottom-Up SBP to the Top-Down approach, leading to up to 13.2× speedup over accelerated Bottom-Up SBP and up to 403× speedup over sequential Bottom-Up SBP on 64 compute nodes. Thus, Top-Down SBP represents substantial improvements to the scalability of SBP, enabling the analysis of larger datasets on the same hardware.
Frank Wanye, Vitaliy Gleyzer, Edward K. Kao, Wu-chun Feng
HPDC4
2025 STAGS: A Graph-Sampling Approach for GNN-based Network Anomaly Detection
Saikat Dey, Mark K. Gardner, Jeffry Lang, Wu-chun Feng
Networking4
2025 Scalable and Maintainable Distributed Sequence Alignment Using Spark
abstract
The exponential growth of genomic data presents a challenge to bioinformatics research. NCBI BLAST, a popular pairwise sequence alignment tool, does not scale with the hundreds of gigabytes (GB) of sequenced data. Therefore, mpiBLAST was widely adopted and scaled up to 65,536 processors. However, mpiBLAST is tightly coupled with an obsolete NCBI BLAST version, creating a challenge to upgrading mpiBLAST with the ever-changing NCBI BLAST code. Recent parallel BLAST implementations, like SparkBLAST, use parallelism wrappers separate from NCBI BLAST to overcome this issue. However, query partitioning, a parallel method that duplicates the genome database on each compute node, makes SparkBLAST scale poorly with databases larger than a single node's memory. Thus, no parallel BLAST utility simultaneously addresses performance, scalability, and software maintainability. To fill this gap, we introduce SparkLeBLAST, a parallel BLAST tool that uses the Spark framework and efficient data partitioning to combine mpiBLAST's performance and scalability with SparkBLAST's simplicity and maintainability. SparkLeBLAST democratizes scalable genomic analysis for domain scientists without extensive distributed computing experience. SparkLeBLAST runs up to 6.68× faster than SparkBLAST. SparkLeBLAST also accelerates taxonomic assignment of COVID-19 genomic diversity analysis by 20.9× as it speeds up the BLAST search component by 88.6× using 128 compute nodes.
Karim Youssef, Yusuf Elnady, Eli Tilevich, Wu-chun Feng
IEEE Trans. Comput. Biol. Bioinform.4
2024 G2A2: Graph Generator with Attributes and Anomalies
abstract
Many data-mining applications use dynamic attributed graphs to represent relational information; but due to security and privacy concerns, there is a dearth of publicly available datasets that can be represented as dynamic attributed graphs. Even when such datasets are available, they do not have ground truth that can be useful for classification problems, e.g., anomaly detection. Thus, researchers commonly generate synthetic graphs using either statistical or deep generative (DG) methods. However, neither approach produces ground truth. Statistical methods struggle to replicate intricate patterns found in real-world dynamic attributed graphs, while DG methods require a significant number of graphs for training.
Saikat Dey, Sonal Jha, Wu-chun Feng
CF3
2024 BLP: Block-Level Pipelining for GPUs
abstract
Programming models like OpenMP offer expressive interfaces to program graphics processing units (GPUs) via directive-based offload. By default, these models copy data to or from the device without overlapping computation, thus impacting performance. Rather than leave the onerous task of manually pipelining and tuning data communication and computation to the end user, we propose an OpenMP extension that supports block-level pipelining and, in turn, present our block-level pipelining (BLP) approach that overlaps data communication and computation in a single kernel. BLP uses persistent thread blocks with cooperative thread groups to process sub-tasks on different streaming multiprocessors and uses GPU flag arrays to enforce task dependencies without CPU involvement.
Wu-chun Feng, Xuewen Cui, Thomas Scogland, Bronis R. de Supinski
CF1
2024 Welcome Message from the IEEE Cluster 2024 Program Chairs
abstract
We are thrilled to share the program for this year's IEEE Cluster conference, showcasing a diverse range of topics in cluster computing and emphasizing the field's ongoing significance. The program strikes a balance between established subjects like architectures, software environments, and scientific applications, and emerging areas such as data analytics and deep learning.
Yutong Lu, Wu-chun Feng, Mohamed Wahib
CLUSTER2
2024 Optimizing and Scaling the 3D Reconstruction of Single-Particle Imaging
abstract
An X-ray free electron laser (XFEL) facility can produce on the order of 1,000,000 extremely bright X-ray light pulses per second. Using an XFEL to image the atomic structure of a molecule requires fast analysis of an enormous amount of data, estimated to exceed one terabyte per second and requiring petabytes of storage. The SpiniFEL application provides such analysis by determining the 3D structure of proteins from single-particle imaging (SPI) experiments performed using XFELs, but it needs significantly better performance and efficiency to scale and keep up with the terabyte-per-second data production. Thus, this paper addresses the high-performance computing optimizations and scaling needed to improve this 3D reconstruction of SPI data. First, we optimize data movement, memory efficiency, and algorithms to improve the per-node computational efficiency and deliver a 5.28× speedup over the baseline GPU implementation.In addition, we achieved a 485× speedup for the post-analysis reconstruction resolution, which previously took as long as the 3D reconstruction of SPI data. Second, we present a novel distributed shared-memory computational algorithm to hide data latency and load-balance network traffic, thus enabling the processing of 128× more orientations than previously possible. Third, we conduct an exploratory study over the hyperparameter space for the SpiniFEL application to identify the optimal parameters for our underlying target hardware, which ultimately led to an up to 1.25× speedup for the number of streams. Overall, we achieve a 6.6× speedup (i.e., 5.28×1.25) over the previous fastest GPUMPI-based SpiniFEL realization.
Niteya Shah, Christine Sweeney, Vinay Ramakrishnaiah, Jeffrey Donatelli, Wu-chun Feng
IPDPS5
2023 Exact Distributed Stochastic Block Partitioning
abstract
Stochastic block partitioning (SBP) is a community detection algorithm that is highly accurate even on graphs with a complex community structure, but its inherently serial nature hinders its widespread adoption by the wider scientific community. To make it practical to analyze large real-world graphs with SBP, there is a growing need to parallelize and distribute the algorithm. The current state-of-the-art distributed SBP algorithm is a divide-and-conquer approach that limits communication between compute nodes until the end of inference. This leads to the breaking of computational dependencies, which causes convergence issues as the number of compute nodes increases and when the graph is sufficiently sparse. To address this shortcoming, we introduce EDiSt — an exact distributed stochastic block partitioning algorithm. Under EDiSt, compute nodes periodically share community assignments during inference. Due to this additional communication, EDiSt improves upon the divide-and-conquer algorithm by allowing it to scale out to a larger number of compute nodes without suffering from convergence issues, even on sparse graphs. We show that EDiSt provides speedups of up to 26.9× over the divide-and-conquer approach and speedups up to 44.0× over shared memory parallel SBP when scaled out to 64 compute nodes.
Frank Wanye, Vitaliy Gleyzer, Edward K. Kao, Wu-chun Feng
CLUSTER4
2022 On the Parallelization of MCMC for Community Detection
abstract
The rapid growth in size of real-world graph datasets necessitates the design of parallel and scalable graph analytics algorithms for large graphs. Community detection is a graph analysis technique with use cases in many domains from bioinformatics to network security. Markov chain Monte Carlo (MCMC)-based methods for performing community detection, such as the stochastic block partitioning (SBP) algorithm, are robust to graphs with a complex structure, but have traditionally been difficult to parallelize due to the serial nature of the underlying MCMC algorithm. This paper presents hybrid SBP (H-SBP), a novel hybrid method to parallelize the inherently sequential computation within each MCMC chain, for SBP. H-SBP processes a fraction of the most influential graph vertices serially and the remaining majority of the vertices in parallel using asynchronous Gibbs. We empirically show that H-SBP speeds up the MCMC computations by up to 5.6 × on real-world graphs while maintaining accuracy.
Frank Wanye, Vitaliy Gleyzer, Edward K. Kao, Wu-chun Feng
ICPP4
2021 ComputeCOVID19+: Accelerating COVID-19 Diagnosis and Monitoring via High-Performance Deep Learning on CT Images
abstract
The COVID-19 pandemic has highlighted the importance of diagnosis and monitoring as early and accurately as possible. However, the reverse-transcription polymerase chain reaction (RT-PCR) test results in two issues: (1) protracted turnaround time from sample collection to testing result and (2) compromised test accuracy, as low as 67%, due to when and how the samples are collected, packaged, and delivered to the lab to conduct the RT-PCR test. Thus, we present ComputeCOVID19+, our computed tomography-based framework to improve the testing speed and accuracy of COVID-19 (plus its variants) via a deep learning-based network for CT image enhancement called DDnet, short for DenseNet and Deconvolution network. To demonstrate its speed and accuracy, we evaluate ComputeCOVID19+ across several sources of computed tomography (CT) images and on many heterogeneous platforms, including multi-core CPU, many-core GPU, and even FPGA. Our results show that ComputeCOVID19+ can significantly shorten the turnaround time from days to minutes and improve the testing accuracy to 91%.
Garvit Goel, Atharva Gondhalekar, Jingyuan Qi, Zhicheng Zhang 0005, Wu-chun Feng
ICPP6
2021 Scaling Out a Combinatorial Algorithm for Discovering Carcinogenic Gene Combinations to Thousands of GPUs
abstract
Cancer is a leading cause of death in the US, second only to heart disease. It is primarily a result of a combination of an estimated two-nine genetic mutations (multi-hit combinations). Although a body of research has identified hundreds of cancer-causing genetic mutations, we don't know the specific combination of mutations responsible for specific instances of cancer for most cancer types. An approximate algorithm for solving the weighted set cover problem was previously adapted to identify combinations of genes with mutations that may be responsible for individual instances of cancer. However, the algorithm's computational requirement scales exponentially with the number of genes, making it impractical for identifying more than three-hit combinations, even after the algorithm was parallelized and scaled up to a V100 GPU. Since most cancers have been estimated to require more than three hits, we scaled out the algorithm to identify combinations of four or more hits using 1000 nodes (6000 V100 GPUs with ≈ 48×106processing cores) on the Summit supercomputer at Oak Ridge National Laboratory. Efficiently scaling out the algorithm required a series of algorithmic innovations and optimizations for balancing an exponentially divergent workload across processors and for minimizing memory latency and inter-node communication. We achieved an average strong scaling efficiency of 90.14% (80.96%-97.96% for 200 to 1000 nodes), compared to a 100 node run, with 84.18% scaling efficiency for 1000 nodes. With experimental validation, the multi-hit combinations identified here could provide further insight into the etiology of different cancer subtypes and provide a rational basis for targeted combination therapy.
Sajal Dash, Qais Al-Hajri, Wu-chun Feng, Harold R. Garner, Ramu Anandakrishnan
IPDPS3
2020 Approximate Pattern Matching for On-Chip Interconnect Traffic Prediction
abstract
Emerging multi-chip module GPUs (MCM-GPUs) expend over 17% of the total power budget on chip interconnects and this fraction is expected to increase as chip size increases. Towards proactively managing the power consumption of these interconnects, we propose approximate pattern matching to predict future interconnect traffic from past observations. Compared to past prediction techniques such as Markov model (MM) and history table (HT), our proposed technique reduces average prediction error to 2.66% from 7.11% and 3.83% for MM and HT, respectively.
Vignesh Adhinarayanan, Wu-chun Feng
PACT2
2020 Alleviating Load Imbalance in Data Processing for Large-Scale Deep Learning
abstract
Scalable deep learning remains an onerous challenge, as it is constrained by many factors, including those related to load imbalance. For many deep-learning software systems, multiple data-processing components-including neural network training, graph scheduling, input pipeline, and gradient synchronization-execute simultaneously and asynchronously. Such execution can cause the various data-processing components to contend with one another for the hardware resources, leading to severe load imbalance and, in turn, degraded scalability. In this paper, we present an in-depth analysis of state-of-the-art deep-learning software, TensorFlow and Horovod, to understand their scalability limitations. Based on this analysis, we propose four novel solutions that minimize resource contention and improve deep-learning performance by up to 35% for training various neural networks on 24,576 GPUs of the Summit supercomputer at Oak Ridge National Laboratory.
Sarunya Pumma, Daniele Buono, Fabio Checconi, Xinyu Que, Wu-chun Feng
CCGRID5
2020 SparkLeBLAST: Scalable Parallelization of BLAST Sequence Alignment Using Spark
abstract
The exponential growth of genomic data presents challenges in analyzing and computing on such biological data at scale. While NCBI’s BLAST is a widely used pairwise sequence alignment tool, it does not scale to large datasets that are hundreds of gigabytes (GB) in size. To address this scalability problem, mpiBLAST emerged and became widely used, enabling scaling to 65,536 processes. However, mpiBLAST suffers from being tightly coupled with a specific implementation of BLAST, rendering it difficult to upgrade with the ever-evolving NCBI BLAST code. To address this shortcoming, recent parallel BLAST tools, such as SparkBLAST, consist of wrappers that are decoupled from the BLAST code but suffer from poor scalability with large sequence databases. Thus, there does not exist any parallel BLAST tool that can simultaneously address the issues of performance, scalability, programmability, and upgradability. To address this void, we propose SparkLeBLAST, a parallel BLAST tool that leverages our performance modeling and the Spark framework to deliver the performance and scalability of mpiBLAST and the ease of programming and upgradability of SparkBLAST, respectively. Ultimately, SparkLeBLAST delivers a 10x speedup relative to the state-of-the-art SparkBLAST and nearly a 2x speedup relative to the latest version of mpiBLAST.
Karim Youssef, Wu-chun Feng
CCGRID2
2020 Exploring FPGA Optimizations in OpenCL for Breadth-First Search on Sparse Graph Datasets
abstract
Breath-first search (BFS) is a fundamental building block in many graph-based applications. It is challenging to optimize due to its irregular memory-access pattern. Prior work, based on hardware description languages (HDLs) and high-level synthesis (HLS), address the memory-access bottleneck by using techniques such as edge-centric traversal, data alignment, and compute-unit (CU) replication. While these optimizations work well for dense graph datasets, optimizing BFS on sparse graphs remains a significant challenge due to the kernel launch overhead and poor workload distribution across processing elements. As a complement to the prior work, we present and evaluate optimizations in OpenCL for BFS on sparse graphs. Specifically, we explore application-specific and architecture-aware optimizations aimed at mitigating the irregular global-memory access bottleneck in sparse graphs. In our kernel design, we consider factors such as choice of data structure between queue and array, number of memory banks, and kernel launch configuration. We evaluate the impact of proposed optimizations on a diverse set of sparse graphs. In comparison with the state-of-the-art OpenCL implementation for FPGA, we achieve 5.7X-22.3X speedup on Stratix 10 SX 2800 FPGA for the graphs that are most sensitive to our optimization scheme.
Atharva Gondhalekar, Wu-chun Feng
FPL2
2020 ETH: An Architecture for Exploring the Design Space of In-situ Scientific Visualization
abstract
As high-performance computing (HPC) moves towards the exascale era, large-scale scientific simulations are generating enormous datasets. Many techniques (e.g., in-situ methods, data sampling, and compression) have been proposed to help visualize these large datasets under various constraints such as storage, power, and energy. However, evaluating these techniques and understanding the trade-offs (e.g., performance, efficiency, and quality) remains a challenging task.To enable exploration of the design space across such trade-offs, we propose the Exploration Test Harness (ETH), an architecture for the early-stage exploration of visualization and rendering approaches, job layout, and visualization pipelines. ETH covers a broader parameter space than current large-scale visualization applications such as ParaView and VisIt. It also promotes the study of simulation-visualization coupling strategies through a data-centric approach, rather than requiring coupling with a specific scientific simulation code. Furthermore, with experimentation on an extensively instrumented supercomputer, we study more metrics of interest than was previously possible. Importantly, ETH will help to answer important what-if scenarios and trade-off questions in the early stages of pipeline development, helping scientists to make informed choices about how to best couple a simulation code with visualization at extreme scale.
Greg Abram, Vignesh Adhinarayanan, Wu-chun Feng, David H. Rogers 0001, James P. Ahrens
IPDPS3
2020 Towards insight-driven sampling for big data visualisation
abstract
Creating an interactive, accurate, and low-latency big data visualisation is challenging due to the volume, variety, and velocity of the data. Visualisation options range from visualising the entire big dataset, which could take a long time and be taxing to the system, to visualising a small subset of the dataset, which could be fast and less taxing to the system but could also lead to a less-beneficial visualisation as a result of information loss. The main research questions investigated by this work are what effect sampling has on visualisation insight and how to provide guidance to users in navigating this trade-off. To investigate these issues, we study an initial case of simple estimation tasks on histogram visualisations of sampled big data, in hopes that these results may generalise. Leveraging sampling, we generate subsets of large datasets and create visualisations for a crowd-sourced study involving a simple cognitive visualisation task. Using the results of this study, we quantify insight, sampling, visualisation, and perception error in comparison to the full dataset. We use these results to model the relationship between sample size and insight error, and we propose the use of our model to guide big data visualisation sampling.
Moeti Masiane, Anne Driscoll, Wu-chun Feng, John E. Wenskovitch, Chris North 0001
Behav. Inf. Technol.3
2019 Adaptive Task Aggregation for High-Performance Sparse Solvers on GPUs
abstract
Sparse solvers are heavily used in computational fluid dynamics (CFD), computer-aided design (CAD), and other important application domains. These solvers remain challenging to execute on massively parallel architectures, due to the sequential dependencies between the fine-grained application tasks. In particular, parallel sparse solvers typically suffer from substantial scheduling and dependency-management overheads relative to the compute operations. We propose adaptive task aggregation (ATA) to efficiently execute such irregular computations on GPU architectures via hierarchical dependency management and low-latency task scheduling. On a gamut of representative problems with different data-dependency structures, ATA significantly outperforms existing GPU task-execution approaches, achieving a geometric mean speedup of 2.2X to 3.7X across different sparse kernels (with speedups of up to two orders of magnitude).
Ahmed E. Helal, Ashwin M. Aji, Michael L. Chu, Bradford M. Beckmann, Wu-chun Feng
PACT5
2019 Iterative machine learning (IterML) for effective parameter pruning and tuning in accelerators
abstract
With the rise of accelerators (e.g., GPUs, FPGAs, and APUs) in computing systems, the parallel computing community needs better tools and mechanisms with which to productively extract performance. While modern compilers provide flags to activate different optimizations to improve performance, the effectiveness of such automated optimization depends on the algorithm and its mapping to the underlying accelerator architecture. Currently, however, extracting the best performance from an algorithm on an accelerator requires significant expertise and manual effort to exploit both spatial and temporal sharing of computing resources in order to improve overall performance. In particular, maximizing the performance on an algorithm on an accelerator requires extensive hyperparameter (e.g., thread-block size) selection and tuning. Given the myriad of hyperparameter dimensions to optimize across, the search space of optimizations is generally extremely large, making it infeasible to exhaustively evaluate each optimization configuration.
Xuewen Cui, Wu-chun Feng
CF2
2019 On the Portability of CPU-Accelerated Applications via Automated Source-to-Source Translation
abstract
Over the past decade, accelerator-based supercomputers have grown from 0% to 42% performance share on the TOP500. Ideally, GPU-accelerated code on such systems should be "write once, run anywhere," regardless of the GPU device (or for that matter, any parallel device, e.g., CPU or FPGA). In practice, however, portability can be significantly more limited due to the sheer volume of code implemented in non-portable languages. For example, the tremendous success of CUDA, as evidenced by the vast cornucopia of CUDA-accelerated applications, makes it infeasible to manually rewrite all these applications to achieve portability. Consequently, we achieve portability by using our automated CUDA-to-OpenCL source-to-source translator called CU2CL. To demonstrate the state of the practice, we use CU2CL to automatically translate three medium-to-large, CUDA-optimized codes to OpenCL, thus enabling the codes to run on other GPU-accelerated systems (as well as CPU- or FPGA-based systems). These automatically translated codes deliver performance portability, including as much as three-fold performance improvement, on a GPU device not supported by CUDA.
Paul Sathre, Mark K. Gardner, Wu-chun Feng
HPC Asia3
2018 GPU power prediction via ensemble machine learning for DVFS space exploration
abstract
A software-based approach to achieve high performance within a power budget often involves dynamic voltage and frequency scaling (DVFS). Thus, accurately predicting the power consumption of an application at different DVFS levels (or more generally, different processor configurations) is paramount for the energy-efficient functioning of a high-performance computing (HPC) system. The increasing prevalence of graphics processing units (GPUs) in HPC systems presents new challenges in power management, and machine learning presents an unique way to improve the software-based power management of these systems. As such, we explore the problem of GPU power prediction at different DVFS states via machine learning. Specifically, we propose a new ensemble technique that incorporates three machine-learning techniques --- sequential minimal optimization regression, simple linear regression, and decision tree --- to reduce the mean absolute error (MAE) to 3.5%.
Bishwajit Dutta, Vignesh Adhinarayanan, Wu-chun Feng
CF3
2018 Taming irregular applications via advanced dynamic parallelism on GPUs
abstract
On recent GPU architectures, dynamic parallelism, which enables the launching of kernels from the GPU without CPU involvement, provides a way to improve the performance of irregular applications by generating child kernels dynamically to reduce workload imbalance and improve GPU utilization. However, in practice, dynamic parallelism does not improve performance due to high kernel launch overhead and low child kernel occupancy. Consequently, most existing studies focus on mitigating the kernel launch overhead. As the kernel launch overhead has decreased due to algorithmic redesigns and hardware architectural innovations, the organization of subtasks to child kernels becomes a new performance bottleneck.
Jing Zhang 0039, Ashwin M. Aji, Michael L. Chu, Hao Wang 0002, Wu-chun Feng
CF5
2018 CommAnalyzer: automated estimation of communication cost and scalability on HPC clusters from sequential code
abstract
To deliver scalable performance to large-scale scientific and data analytic applications, HPC cluster architectures adopt the distributed-memory model. The performance and scalability of parallel applications on such systems are limited by the communication cost across compute nodes. Therefore, projecting the minimum communication cost and maximum scalability of the user applications plays a critical role in assessing the benefits of porting these applications to HPC clusters as well as developing efficient distributed-memory implementations. Unfortunately, this task is extremely challenging for end users, as it requires comprehensive knowledge of the target application and hardware architecture and demands significant effort and time for manual system analysis.
Ahmed E. Helal, Changhee Jung, Wu-chun Feng, Yasser Y. Hanafy
HPDC3
2018 A Framework for Auto-Parallelization and Code Generation: An Integrative Case Study with Legacy FORTRAN Codes
abstract
GLAF, short for Grid-based Language and Auto-parallelization Framework, is a programming framework that seeks to democratize parallel programming by facilitating better productivity in parallel computing via an intuitive graphical programming interface (GPI) that automatically parallelizes and generates code in many languages. Originally, GLAF addressed program development from scratch via the GPI; but this unduly restricted GLAF's utility to creating new codes only. Thus, this paper extends GLAF by enabling program development from pre-existing kernels of interest, which can then be easily and transparently integrated into existing legacy codes. Specifically, we address the theoretical and practical limitations of integration and interoperability of auto-generated parallel code within existing FORTRAN codes; enhance GLAF to overcome these limitations; and present an integrative case study and evaluation of the enhanced GLAF via the implementation of important kernels in two NASA codes: (1) the Synoptic Surface & Atmospheric Radiation Budget (SARB), part of the Clouds and the Earth's Radiant Energy System (CERES), and (2) the Fully Unstructured Navier-Stokes (FUN3D) suite for computational fluid dynamics.
Konstantinos Krommydas, Paul Sathre, Ruchira Sasanka, Wu-chun Feng
ICPP4
2018 Highly Efficient Compensation-Based Parallelism for Wavefront Loops on GPUs
abstract
Wavefront loops are widely used in many scientific applications, e.g., partial differential equation (PDE) solvers and sequence alignment tools. However, due to the data dependencies in wavefront loops, it is challenging to fully utilize the abundant compute units of GPUs and to reuse data through their memory hierarchy. Existing solutions can only optimize for these factors to a limited extent. For example, tiling-based methods optimize memory access but may result in load imbalance; while compensation-based methods, which change the original order of computation to expose more parallelism and then compensate for it, suffer from both global synchronization overhead and limited generality. In this paper, we first prove under which circumstances that breaking data dependencies and properly changing the sequence of computation operators in our compensation-based method does not affect the correctness of results. Based on this analysis, we design a highly efficient compensation-based parallelism on GPUs. Our method provides weighted scan-based GPU kernels to optimize the computation and combines with the tiling method to optimize memory access and synchronization. The performance results on the NVIDIA K80 and P100 GPU platforms demonstrate that our method can achieve significant improvements for four types of real-world application kernels over the state-of-the-art research.
Kaixi Hou, Hao Wang 0002, Wu-chun Feng, Jeffrey S. Vetter, Seyong Lee
IPDPS3
2018 A Framework for the Automatic Vectorization of Parallel Sort on x86-Based Processors
abstract
The continued growth in the width of vector registers and the evolving library of intrinsics on the modern x86 processors make manual optimizations for data-level parallelism tedious and error-prone. In this paper, we focus on parallel sorting, a building block for many higher-level applications, and propose a framework for the Automatic SIMDization of Parallel Sorting (ASPaS) on x86-based multiand many-core processors. That is, ASPaS takes any sorting network and a given instruction set architecture (ISA) as inputs and automatically generates vector code for that sorting network. After formalizing the sort function as a sequence of comparators and the transpose and merge functions as sequences of vector-matrix multiplications, ASPaS can map these functions to operations from a selected “pattern pool” that is based on the characteristics of parallel sorting, and then generate the vector code with the real ISA intrinsics. The performance evaluation on the Intel Ivy Bridge and Haswell CPUs, and Knights Corner MIC illustrates that automatically generated sorting codes from ASPaS can outperform the widely used sorting tools, achieving up to 5.2x speedup over the single-threaded implementations from STL and Boost and up to 6.7x speedup over the multi-threaded parallel sort from Intel TBB.
Kaixi Hou, Hao Wang 0002, Wu-chun Feng
IEEE Trans. Parallel Distributed Syst.3
2017 Robotomata: A framework for approximate pattern matching of big data on an automata processor
abstract
Approximate pattern matching (APM) has been widely used in big data applications, e.g., genome data analysis, speech recognition, fraud detection, computer vision, etc. Although an automata-based approach is an efficient way to realize APM, the inherent sequentiality of automata deters its implementation on general-purpose parallel platforms, e.g., multicore CPUs and many-core GPUs. Recently, however, Micron has proposed its Automata Processor (AP), a processing-in-memory (PIM) architecture dedicated for non-deterministic automata (NFA) simulation. It has nominally achieved thousands-fold speedup over a multicore CPU for many big data applications. Alas, the AP ecosystem suffers from two major problems. First, the current APIs of AP require manual manipulations of all computational elements. Second, multiple rounds of time-consuming compilation are needed for large datasets. Both problems hinder programmer productivity and end-to-end performance. Therefore, we propose a paradigm-based approach to hierarchically generate automata on AP and use this approach to create Robotomata, a framework for APM on AP. By taking in the following inputs — the types of APM paradigms, desired pattern length, and allowed number of errors as input — our framework can generate the optimized APM-automata codes on AP, so as to improve programmer productivity. The generated codes can also maximize the reuse of pre-compiled macros and significantly reduce the time for reconfiguration. We evaluate Robotomata by comparing it to two state-of-the-art APM implementations on AP with real-world datasets. Our experimental results show that our generated codes can achieve up to 30.5x and 12.8x speedup with respect to configuration while maintaining the computational performance. Compared to the counterparts on CPU, our codes achieve up to 393x overall speedup, even when including the reconfiguration costs. We highlight the importance of counting the configuration time towards the overall performance on AP, which would provide better insight in identifying essential hardware features, specifically for large-scale problem sizes.
Xiaodong Yu 0001, Kaixi Hou, Hao Wang 0002, Wu-chun Feng
IEEE BigData4
2017 Developing Dynamic Profiling and Debugging Support in OpenCL for FPGAs
abstract
With FPGAs emerging as a promising accelerator for general-purpose computing, there is a strong demand to make them accessible to software developers. Recent advances in OpenCL compilers for FPGAs pave the way for synthesizing FPGA hardware from OpenCL kernel code. To enable broader adoption of this paradigm, significant challenges remain. This paper presents our efforts in developing dynamic profiling and debugging support in OpenCL for FPGAs. We first propose primitive code patterns, including a timestamp and an event-ordering function, and then develop a framework, which can be plugged easily into OpenCL kernels, to dynamically collect and process run-time information.
Anshuman Verma, Huiyang Zhou, Skip Booth, Robbie King, James Coole, Andy Keep, John Marshall, Wu-chun Feng
DAC8
2017 Parallel I/O Optimizations for Scalable Deep Learning
abstract
As deep learning systems continue to grow in importance, researchers have been analyzing approaches to make such systems efficient and scalable on high-performance computing platforms. As computational parallelism increases, however, data I/O becomes the major bottleneck limiting the overall system scalability. In this paper, we continue our efforts to improve LMDB, the I/O subsystem of the Caffe deep learning framework. In a previous paper we presented LMDBIO---an optimized I/O plugin for Caffe that takes into account the data access pattern of Caffe in order to vastly improve I/O performance. Nevertheless, LMDBIO's optimizations, which we henceforth call LMM (localized mmap), are limited to intranode performance, and these optimizations do little to minimize the I/O inefficiencies in distributed-memory environments. In this paper, we propose LMDBIO-DM, an enhanced version of LMDBIO-LMM that optimizes the I/O access of Caffe in distributed-memory environments. We present several sophisticated data I/O techniques that allow for significant improvement in such environments. Our experimental results show that LMDBIO-DM can improve the overall execution time of Caffe by more than 30-fold compared with LMDB and by 2-fold compared with LMDBIO-LMM.
Sarunya Pumma, Min Si, Wu-chun Feng, Pavan Balaji
ICPADS3
2017 Fast segmented sort on GPUs
abstract
Segmented sort, as a generalization of classical sort, orders a batch of independent segments in a whole array. Along with the wider adoption of manycore processors for HPC and big data applications, segmented sort plays an increasingly important role than sort. In this paper, we present an adaptive segmented sort mechanism on GPUs. Our mechanisms include two core techniques: (1) a differentiated method for different segment lengths to eliminate the irregularity caused by various workloads and thread divergence; and (2) a register-based sort method to support N-to-M data-thread binding and in-register data communication. We also implement a shared memory-based merge method to support non-uniform length chunk merge via multiple warps. Our segmented sort mechanism shows great improvements over the methods from CUB, CUSP and ModernGPU on NVIDIA K80-Kepler and TitanX-Pascal GPUs. Furthermore, we apply our mechanism on two applications, i.e., suffix array construction and sparse matrix-matrix multiplication, and obtain obvious gains over state-of-the-art implementations.
Kaixi Hou, Weifeng Liu 0002, Hao Wang 0002, Wu-chun Feng
ICS4
2017 Demystifying automata processing: GPUs, FPGAs or Micron's AP?
abstract
Many established and emerging applications perform at their core some form of pattern matching, a computation that maps naturally onto finite automata abstractions. As a consequence, in recent years there has been a substantial amount of work on high-speed automata processing, which has led to a number of implementations targeting a variety of parallel platforms: CPUs, GPUs, FPGAs, ASICs, and Network Processors. More recently, Micron has announced its Automata Processor (AP), a DRAM-based accelerator of non-deterministic finite automata (NFA). Despite the abundance of work in this domain, the advantages and disadvantages of different automata processing accelerators and the innovation space in this area are still unclear.
Marziyeh Nourian, Xiaodong Yu 0001, Wu-chun Feng, Michela Becchi
ICS4
2017 Characterizing and Modeling Power and Energy for Extreme-Scale In-Situ Visualization
abstract
Plans for exascale computing have identified power and energy as looming problems for simulations running at that scale. In particular, writing to disk all the data generated by these simulations is becoming prohibitively expensive due to the energy consumption of the supercomputer while it idles waiting for data to be written to permanent storage. In addition, the power cost of data movement is also steadily increasing. A solution to this problem is to write only a small fraction of the data generated while still maintaining the cognitive fidelity of the visualization. With domain scientists increasingly amenable towards adopting an in-situ framework that can identify and extract valuable data from extremely large simulation results and write them to permanent storage as compact images, a large-scale simulation will commit to disk a reduced dataset of data extracts that will be much smaller than the raw results, resulting in a savings in both power and energy. The goal of this paper is two-fold: (i) to understand the role of in-situ techniques in combating power and energy issues of extreme-scale visualization and (ii) to create a model for performance, power, energy, and storage to facilitate what-if analysis. Our experiments on a specially instrumented, dedicated 150-node cluster show that while it is difficult to achieve power savings in practice using in-situ techniques, applications can achieve significant energy savings due to shorter write times for in-situ visualization. We present a characterization of power and energy for in-situ visualization; an application-aware, architecture-specific methodology for modeling and analysis of such in-situ workflows; and results that uncover indirect power savings in visualization workflows for high-performance computing (HPC).
Vignesh Adhinarayanan, Wu-chun Feng, David H. Rogers 0001, James P. Ahrens, Scott Pakin
IPDPS2
2017 Directive-Based Partitioning and Pipelining for Graphics Processing Units
abstract
The community needs simpler mechanisms to access the performance available in accelerators, such as GPUs, FPGAs, and APUs, due to their increasing use in state-of-the-art supercomputers. Programming models like CUDA, OpenMP, OpenACC and OpenCL can efficiently offload compute-intensive workloads to these devices. By default these models naively offload computation without overlapping it with communication (copying data to or from the device). Achieving performance can require extensive refactoring and hand-tuning to apply optimizations such as pipelining. Further, users must manually partition the dataset whenever its size is larger than device memory, which can be especially difficult when the device memory size is not exposed to the user. We propose a directive-based partitioning and pipelining extension for accelerators appropriate for either OpenMP or OpenACC. Its interface supports overlap of data transfers and kernel computation without explicit user splitting of data. It can map data to a pre-allocated device buffer and automate memory-constrained array indexing and sub-task scheduling. We evaluate a prototype implementation with four different applications. The experimental results show that our approach can reduce memory usage by 52% to 97% while delivering a 1.41× to 1.65× speedup over the naive offload model.
Xuewen Cui, Thomas Scogland, Bronis R. de Supinski, Wu-chun Feng
IPDPS4
2017 PaPar: A Parallel Data Partitioning Framework for Big Data Applications
abstract
Today, big data applications can generate large-scale data sets at an unprecedented rate; and scientists have turned to parallel and distributed systems for data analysis. Although many big data processing systems provide advanced mechanisms to partition data and tackle the computational skew, it is difficult to efficiently implement skew-resistant mechanisms, because the runtime of different partitions not only depends on input data size but also algorithms that will be applied on data. As a result, many research efforts have been undertaken to explore user-defined partitioning methods for different types of applications and algorithms. However, manually writing application-specific partitioning methods requires significant coding effort, and finding the optimal data partitioning strategy is particularly challenging even for developers that have mastered sufficient application knowledge. In this paper, we propose PaPar, a Parallel data Partitioning framework for big data applications, to simplify the implementations of data partitioning algorithms. PaPar provides a set of computational operators and distribution strategies for programmers to describe desired data partitioning methods. Taking an input data configuration file and a workflow configuration file as the input, PaPar can automatically generate the parallel partitioning codes by formalizing the user-defined workflow as a sequence of key-value operations and matrix-vector multiplications, and efficiently mapping to the parallel implementations with MPI and MapReduce. We apply our approach on two applications: muBLAST, a MPI implementation of BLAST algorithms for biological sequence search; and PowerLyra, a computation and partitioning method for skewed graphs. The experimental results show that compared to the partitioning methods of applications, the codes generated by PaPar can produce the same data partitions with comparable or less partitioning time.
Hao Wang 0002, Jing Zhang 0039, Da Zhang 0004, Sarunya Pumma, Wu-chun Feng
IPDPS5
2017 Eliminating Irregularities of Protein Sequence Search on Multicore Architectures
abstract
Finding regions of local similarity between biological sequences is a fundamental task in computational biology. BLAST is the most widely-used tool for this purpose, but it suffers from irregularities due to its heuristic nature. To achieve fast search, recent approaches construct the index from the database instead of the input query. However, database indexing introduces more challenges in the design of index structure and algorithm, especially for data access through the memory hierarchy on modern multicore processors. In this paper, based on existing heuristic algorithms, we design and develop a database indexed BLAST with the identical sensitivity as query indexed BLAST (i.e., NCBI-BLAST). Then, we identify that existing heuristic algorithms of BLAST can result in serious irregularities in database indexed search. To eliminate irregularities in BLAST algorithm, we propose muBLASTP, that uses multiple optimizations to improve data locality and parallel efficiency for multicore architectures and multi-node systems. Experiments on a single node demonstrate up to a 5.1-fold speedup over the multi-threaded NCBI BLAST. For the inter-node parallelism, we achieve nearly linear scaling on up to 128 nodes and gain up to 8.9-fold speedup over mpiBLAST.
Jing Zhang 0039, Sanchit Misra, Hao Wang 0002, Wu-chun Feng
IPDPS4
2017 A runtime estimation framework for ALICE
Sarunya Pumma, Wu-chun Feng, Phond Phunchongharn, Sylvain Chapeland, Tiranee Achalakul
Future Gener. Comput. Syst.2
2017 Parallel programming with pictures is a Snap!
Annette C. Feng, Mark K. Gardner, Wu-chun Feng
J. Parallel Distributed Comput.3
2017 cuBLASTP: Fine-Grained Parallelization of Protein Sequence Search on CPU+GPU
abstract
BLAST, short for Basic Local Alignment Search Tool, is a ubiquitous tool used in the life sciences for pairwise sequence search. However, with the advent of next-generation sequencing (NGS), whether at the outset or downstream from NGS, the exponential growth of sequence databases is outstripping our ability to analyze the data. While recent studies have utilized the graphics processing unit (GPU) to speedup the BLAST algorithm for searching protein sequences (i.e., BLASTP), these studies use coarse-grained parallelism, where one sequence alignment is mapped to only one thread. Such an approach does not efficiently utilize the capabilities of a GPU, particularly due to the irregularity of BLASTP in both execution paths and memory-access patterns. To address the above shortcomings, we present a fine-grained approach to parallelize BLASTP, where each individual phase of sequence search is mapped to many threads on a GPU. This approach, which we refer to as cuBLASTP, reorders data-access patterns and reduces divergent branches of the most time-consuming phases (i.e., hit detection and ungapped extension). In addition, cuBLASTP optimizes the remaining phases (i.e., gapped extension and alignment with trace back) on a multicore CPU and overlaps their execution with the phases running on the GPU.
Jing Zhang 0039, Hao Wang 0002, Wu-chun Feng
IEEE ACM Trans. Comput. Biol. Bioinform.3
2016 O3FA: A Scalable Finite Automata-based Pattern-Matching Engine for Out-of-Order Deep Packet Inspection
abstract
To match the signatures of malicious traffic across packet boundaries, network-intrusion detection (and prevention) systems (NIDS) typically perform pattern matching after flow reassembly or packet reordering. However, this may lead to the need for large packet buffers, making detection vulnerable to denial-of-service (DoS) attacks, whereby attackers exhaust the buffer capacity by sending long sequences of out-of-order packets. While researchers have proposed solutions for exact-match patterns, regular-expression matching on out-of-order packets is still an open problem. Specifically, a key challenge is the matching of complex sub-patterns (such as repetitions of wildcards matched at the boundary between packets). Our proposed approach leverages the insight that various segments matching the same repetitive sub-pattern are logically equivalent to the regular-expression matching engine, and thus, inter-changing them would not affect the final result. In this paper, we present O3FA, a new finite automata-based, deep packet-inspection engine to perform regular-expression matching on out-of-order packets without requiring flow reassembly. O3FA consists of a deterministic finite automaton (FA) coupled with a set of prefix-/suffix-FA, which allows processing out-of-order packets on the fly. We present our design, optimization, and evaluation for the O3FA engine. Our experiments show that our design requires 20x-4000x less buffer space than conventional buffering-and-reassembling schemes on various datasets and that it can process packets in real-time, i.e., without reassembly.
Xiaodong Yu 0001, Wu-chun Feng, Danfeng Yao, Michela Becchi
ANCS2
2016 Bridging the FPGA programmability-portability Gap via automatic OpenCL code generation and tuning
abstract
Programming FPGAs has been an arduous task that requires extensive knowledge of hardware design languages (HDLs), such as Verilog or VHDL, and low-level hardware details. With OpenCL support for FPGAs, the design, prototyping and implementation of an FPGA is increasingly moving towards a much higher level of abstraction, when compared to the intrinsically low-level nature of HDLs. On the other hand, in the context of traditional (i.e., CPU) software development, OpenCL is still considered to be low-level and complex because the programmer needs to manually expose parallelism in the code. In this work, we present our approach to enhancing FPGA programmability via GLAF, a visual programming framework, to automatically generate synthesizable OpenCL code with an array of FPGA-specific optimizations. We find that our tool facilitates the development process and produces functionally correct and well-performing code on the FPGA for our molecular modeling, gene sequence search, and filtering algorithms.
Konstantinos Krommydas, Ruchira Sasanka, Wu-chun Feng
ASAP3
2016 Online Power Estimation of Graphics Processing Units
abstract
Accurate power estimation at runtime is essential for the efficient functioning of a power management system. While years of research have yielded accurate power models for the online prediction of instantaneous power for CPUs, such power models for graphics processing units (GPUs) are lacking. GPUs rely on low-resolution power meters that only nominally support basic power management. To address this, we propose an instantaneous power model, and in turn, a power estimator, that uses performance counters in a novel way so as to deliver accurate power estimation at runtime. Our power estimator runs on two real NVIDIA GPUs to show that accurate runtime estimation is possible without the need for the high-fidelity details that are assumed on simulation-based power models. To construct our power model, we first use correlation analysis to identify a concise set of performance counters that work well despite GPU device limitations. Next, we explore several statistical regression techniques and identify the best one. Then, to improve the prediction accuracy, we propose a novel application-dependent modeling technique, where the model is constructed online at runtime, based on the readings from a low-resolution, built-in GPU power meter. Our quantitative results show that a multi-linear model, which produces a mean absolute error of 6%, works the best in practice. An application-specific quadratic model reduces the error to nearly 1%. We show that this model can be constructed with low overhead and high accuracy at runtime. To the best of our knowledge, this is the first work attempting to model the instantaneous power of a real GPU system, earlier related work focused on average power.
Vignesh Adhinarayanan, Balaji Subramaniam, Wu-chun Feng
CCGrid3
2016 cuART: Fine-Grained Algebraic Reconstruction Technique for Computed Tomography Images on GPUs
abstract
Algebraic reconstruction technique (ART) is an iterative algorithm for computed tomography (CT) image reconstruction. Due to the high computational cost, researchers turn to modern HPC systems with GPUs to accelerate the ART algorithm. However, the existing proposals suffer from inefficient designs of compressed data structure and computational kernel on GPUs. In this paper, we identify the computational patterns in the ART as the product of a sparse matrix (and its transpose) with multiple vectors (SpMV and SpMV_T). Because the implementations with well-tuned libraries, including cuSPARSE, BRC, and CSR5, underperform the expectations, we propose cuART, a complete compression and parallelization solution for the ART-based CT on GPUs. Based on the physical characteristics, i.e., the symmetries in the system matrix, we propose the symmetry-based CSR format (SCSR), which can further compress data storage by removing symmetric but redundant non-zero elements. Leveraging the sparsity patterns of X-ray projection, wetransform the CSR format to multiple dense sub-matrices in SCSR. We then design a transposition-free kernel to optimize the data access for both SpMV and SpMV_T. The experimental results illustrate that our mechanism can reduce memory usage significantly and make practical datasets fit into a single GPU. Our results also illustrate the superior performance of cuART compared to the existing methods on CPU and GPU.
Xiaodong Yu 0001, Hao Wang 0002, Wu-chun Feng, Hao Gong 0001
CCGrid3
2016 Directive-Based Pipelining Extension for OpenMP
abstract
Programming models like CUDA, OpenMP, OpenACC and OpenCL are designed to offload compute-intensive workloads to accelerators efficiently. However, the naive offload model, which synchronously copies and executes in sequence, requires extensive hand-tuning of techniques, such as pipelining to overlap computation and communication. Therefore, we propose an easy-to-use, directive-based pipelining extension for OpenMP to overlap data transfers and kernel computation. This extension can map data to a pre-allocated device buffer and can automate memory-constrained array indexing and sub-task scheduling. We evaluate a prototype implementation of our approach with three different applications. The experimental results show that our approach can reduce memory usage by 52% to 97% while delivering a 1:41X to 1:65X speedup over the naive offload model.
Xuewen Cui, Thomas Scogland, Bronis R. de Supinski, Wu-chun Feng
CLUSTER4
2016 Bridging the Performance-Programmability Gap for FPGAs via OpenCL: A Case Study with OpenDwarfs
abstract
For decades, the streaming architecture of FPGAs has delivered accelerated performance across many application domains, such as option pricing solvers in finance, computational fluid dynamics in oil and gas, and packet processing in network routers and firewalls. However, this performance has come at the significant expense of programmability, i.e., the performance-programmability gap. In particular, FPGA developers use a hardware design language (HDL) to implement the application data path and to design hardware modules for computation pipelines, memory management, synchronization, and communication. This process requires extensive low-level knowledge of the target FPGA architecture and consumes significant development time and effort. To address this lack of programmability of FPGAs, OpenCL provides an easy-to-use and portable programming model for CPUs, GPUs, APUs, and now, FPGAs. However, this significantly improved programmability can come at the expense of performance, that is, there still remains a performance-programmability gap. To improve the performance of OpenCL kernels on FPGAs, and thus, bridge the performance-programmability gap, we apply and evaluate the effect of various optimization techniques on GEM, an N-body method from the OpenDwarfs benchmark suite.
Konstantinos Krommydas, Ahmed E. Helal, Anshuman Verma, Wu-chun Feng
FCCM4
2016 Telescoping Architectures: Evaluating Next-Generation Heterogeneous Computing
abstract
Architectural innovation has telescoped the HPC community from the commodity (Beowulf) cluster in a machine room, i.e., a multi-node system with Ethernet interconnect, to a commodity cluster on a chip, i.e., multicore CPU with an on-die interconnect. We project that this notion of "telescoping architecture" will apply more broadly to heterogeneous computing, namely from heterogeneous clusters like Tianhe-2 in a machine room to on a chip. To that end, we present an experimental studythat extends the notion of telescoping architectures to identify the ideal mixture of compute engines (CEs) and the number of such CEs on a chip to create a heterogeneous "cluster on a chip" (CoC). Specifically, we experiment with heterogeneous architectures that contain single or multiple instances of CPUs, GPUs, Intel MICs, and FPGAs to demonstrate their performance efficacy given continuing advances in hardware technology, software, tools, and run-time support.
Konstantinos Krommydas, Wu-chun Feng
HiPC2
2016 Parallel Transposition of Sparse Data Structures
abstract
Many applications in computational sciences and social sciences exploit sparsity and connectivity of acquired data. Even though many parallel sparse primitives such as sparse matrix-vector (SpMV) multiplication have been extensively studied, some other important building blocks, e.g., parallel transposition for sparse matrices and graphs, have not received the attention they deserve.
Hao Wang 0002, Weifeng Liu 0002, Kaixi Hou, Wu-chun Feng
ICS4
2016 AAlign: A SIMD Framework for Pairwise Sequence Alignment on x86-Based Multi-and Many-Core Processors
abstract
Pairwise sequence alignment algorithms, e.g., Smith-Waterman and Needleman-Wunsch, with adjustable gap penalty systems are widely used in bioinformatics. The strong data dependencies in these algorithms, however, prevents compilers from effectively auto-vectorizing them. When programmers manually vectorize them on multi-and many-core processors, two vectorizing strategies are usually considered, both of which initially ignore data dependencies and then appropriately correct in a subsequent stage: (1) iterate, which vectorizes and then compensates the scoring results with multiple rounds of corrections and (2) scan, which vectorizes and then corrects the scoring results primarily via one round of parallel scan. However, manually writing such vectorizing code efficiently is non-trivial, even for experts, and the code may not be portable across ISAs. In addition, even highly vectorized and optimized codes may not achieve optimal performance because selecting the best vectorizing strategy depends on the algorithms, configurations (gap systems), and input sequences. Therefore, we propose a framework called AAlign to automatically vectorize pairwise sequence alignment algorithms across ISAs. AAlign ingests a sequential code (which follows our generalized paradigm for pairwise sequence alignment) and automatically generates efficient vector code for iterate and scan. To reap the benefits of both vectorization strategies, we propose a hybrid mechanism where AAlign automatically selects the best vectorizing strategy at runtime no matter which algorithms, configurations, and input sequences are specified. On Intel Haswell and MIC, the generated codes for Smith-Waterman and Needleman-Wunsch achieve up to a 26-fold speedup over their sequential counterparts. Compared to the highly optimized and multi-threaded sequence alignment tools, e.g., SWPS3 and SWAPHI, our codes can deliver up to 2.5-fold and 1.6-fold speedups, respectively.
Kaixi Hou, Hao Wang 0002, Wu-chun Feng
IPDPS3
2016 An automated framework for characterizing and subsetting GPGPU workloads
abstract
Graphics processing units (GPUs) are becoming increasingly common in today's computing systems due to their superior performance and energy efficiency relative to their cost. To further improve these desired characteristics, researchers have proposed several software and hardware techniques. Evaluation of these proposed techniques could be tricky due to the ad-hoc nature in which applications are selected for evaluation. Sometimes researchers spend unnecessary time evaluating redundant workloads, which is particularly problematic for time-consuming studies involving simulation. Other times, they fail to expose the shortcomings of their proposed techniques when too few workloads are chosen for evaluation. To overcome these problems, we propose an automated framework that characterizes and subsets GPGPU workloads, depending on a user-chosen set of performance metrics/counters. This framework internally uses principal component analysis (PCA) to reduce the dimensionality of the chosen metrics and then uses hierarchical clustering to identify similarity among the workloads. In this study, we use our framework to identify redundancy in the recently released SPEC ACCEL OpenCL benchmark suite using a few architecture-dependent metrics. Our analysis shows that a subset of eight applications provides most of the diversity in the 19-application benchmark suite. We also subset the Parboil, Rodinia, and SHOC benchmark suites and then compare them against each other to identify “gaps” in these suites. As an example, we show that SHOC has many applications that are similar to each other and could benefit from adding four applications from Parboil to improve its diversity.
Vignesh Adhinarayanan, Wu-chun Feng
ISPASS2
2016 Characterizing Performance and Power towards Efficient Synchronization of GPU Kernels
abstract
There is a lack of support for explicit synchronization in GPUs between the streaming multiprocessors (SMs) adversely impacts the performance of the GPUs to efficiently perform inter-block communication. In this paper, we present several approaches to inter-block synchronization using explicit/implicit CPU-based and dynamic parallelism (DP) mechanisms. Although this topic has been addressed in previous research studies, there has been neither a solid quantification of such overhead, nor guidance on when to use each of the different approaches. Therefore, we quantify the synchronization overhead relative to the number of kernel launches and the input data sizes. The quantification, in turn, provides insight as to when to use each of the aforementioned synchronization mechanisms in a target application. Our results show that implicit CPU synchronization has a significant overhead that hurts the application performance when using medium to large data sizes with relatively large number of kernel launches (i.e. ~1100-5000). Hence, it is recommended to use explicit CPU synchronization with these configurations. In addition, among the three different approaches, we conclude that dynamic parallelism (DP) is the most efficient with small data sizes (i.e., ~128k bytes), regardless of the number of kernel launches. Also, Dynamic Parallelism (DP), implicitly, performs inter-block (i.e. global) synchronization with no CPU intervention. Therefore, DP significantly reduces the power consumed by the CPU and PCIe for global synchronization. Our findings show that DP reduces the power consumption by ~8-10%. However, DP-based synchronization is a trade-off, in which it is accompanied by ~2-5% performance loss.
Islam Harb, Wu-chun Feng
MASCOTS2
2016 MetaMorph: a library framework for interoperable kernels on multi- and many-core clusters
abstract
To attain scalable performance efficiently, the HPC community expects future exascale systems to consist of multiple nodes, each with different types of hardware accelerators. In addition to GPUs and Intel MICs, additional candidate accelerators include embedded multiprocessors and FPGAs. End users need appropriate tools to efficiently use the available compute resources in such systems, both within a compute node and across compute nodes. As such, we present MetaMorph, a library framework designed to (automatically) extract as much computational capability as possible from HPC systems. Its design centers around three core principles: abstraction, interoperability, and adaptivity. To demonstrate its efficacy, we present a case study that uses the structured grids design pattern, which is heavily used in computational fluid dynamics. We show how MetaMorph significantly reduces the development time, while delivering performance and interoperability across an array of heterogeneous devices, including multicore CPUs, Intel MICs, AMD GPUs, and NVIDIA GPUs.
Ahmed E. Helal, Paul Sathre, Wu-chun Feng
SC3
2016 muBLASTP: database-indexed protein sequence search on multicore CPUs
abstract
BACKGROUND: The Basic Local Alignment Search Tool (BLAST) is a fundamental program in the life sciences that searches databases for sequences that are most similar to a query sequence. Currently, the BLAST algorithm utilizes a query-indexed approach. Although many approaches suggest that sequence search with a database index can achieve much higher throughput (e.g., BLAT, SSAHA, and CAFE), they cannot deliver the same level of sensitivity as the query-indexed BLAST, i.e., NCBI BLAST, or they can only support nucleotide sequence search, e.g., MegaBLAST. Due to different challenges and characteristics between query indexing and database indexing, the existing techniques for query-indexed search cannot be used into database indexed search. RESULTS: muBLASTP, a novel database-indexed BLAST for protein sequence search, delivers identical hits returned to NCBI BLAST. On Intel Haswell multicore CPUs, for a single query, the single-threaded muBLASTP achieves up to a 4.41-fold speedup for alignment stages, and up to a 1.75-fold end-to-end speedup over single-threaded NCBI BLAST. For a batch of queries, the multithreaded muBLASTP achieves up to a 5.7-fold speedups for alignment stages, and up to a 4.56-fold end-to-end speedup over multithreaded NCBI BLAST. CONCLUSIONS: With a newly designed index structure for protein database and associated optimizations in BLASTP algorithm, we re-factored BLASTP algorithm for modern multicore processors that achieves much higher throughput with acceptable memory footprint for the database index.
Jing Zhang 0039, Sanchit Misra, Hao Wang 0002, Wu-chun Feng
BMC Bioinform.4
2016 MultiCL: Enabling automatic scheduling for task-parallel workloads in OpenCL
Ashwin M. Aji, Antonio J. Peña, Pavan Balaji, Wu-chun Feng
Parallel Comput.4
2016 Fast Detection of Transformed Data Leaks
abstract
The leak of sensitive data on computer systems poses a serious threat to organizational security. Statistics show that the lack of proper encryption on files and communications due to human errors is one of the leading causes of data loss. Organizations need tools to identify the exposure of sensitive data by screening the content in storage and transmission, i.e., to detect sensitive information being stored or transmitted in the clear. However, detecting the exposure of sensitive information is challenging due to data transformation in the content. Transformations (such as insertion and deletion) result in highly unpredictable leak patterns. In this paper, we utilize sequence alignment techniques for detecting complex data-leak patterns. Our algorithm is designed for detecting long and inexact sensitive data patterns. This detection is paired with a comparable sampling algorithm, which allows one to compare the similarity of two separately sampled sequences. Our system achieves good detection accuracy in recognizing transformed leaks. We implement a parallelized version of our algorithms in graphics processing unit that achieves high analysis throughput. We demonstrate the high multithreading scalability of our data leak detection method required by a sizable organization.
Xiaokui Shu, Jing Zhang 0039, Danfeng Yao, Wu-chun Feng
IEEE Trans. Inf. Forensics Secur.4
2016 MPI-ACC: Accelerator-Aware MPI for Scientific Applications
abstract
Data movement in high-performance computing systems accelerated by graphics processing units (GPUs) remains a challenging problem. Data communication in popular parallel programming models, such as the Message Passing Interface (MPI), is currently limited to the data stored in the CPU memory space. Auxiliary memory systems, such as GPU memory, are not integrated into such data movement standards, thus providing applications with no direct mechanism to perform end-to-end data movement. We introduce MPI-ACC, an integrated and extensible framework that allows end-to-end data movement in accelerator-based systems. MPI-ACC provides productivity and performance benefits by integrating support for auxiliary memory spaces into MPI. MPI-ACC supports data transfer among CUDA, OpenCL and CPU memory spaces and is extensible to other offload models as well. MPI-ACC's runtime system enables several key optimizations, including pipelining of data transfers, scalable memory management techniques, and balancing of communication based on accelerator and node architecture. MPI-ACC is designed to work concurrently with other GPU workloads with minimum contention. We describe how MPI-ACC can be used to design new communication-computation patterns in scientific applications from domains such as epidemiology simulation and seismology modeling, and we discuss the lessons learned. We present experimental results on a state-of-the-art cluster with hundreds of GPUs; and we compare the performance and productivity of MPI-ACC with MVAPICH, a popular CUDA-aware MPI solution. MPI-ACC encourages programmers to explore novel application-specific optimizations for improved overall cluster utilization.
Ashwin M. Aji, Lokendra S. Panwar, Karthik Murthy, Milind Chabbi, Pavan Balaji, Keith R. Bisset, James Dinan, Wu-chun Feng, John M. Mellor-Crummey, Xiaosong Ma, Rajeev Thakur
IEEE Trans. Parallel Distributed Syst.9
2015 Automatic Command Queue Scheduling for Task-Parallel Workloads in OpenCL
abstract
OpenCL is a portable interface that can be used to program cluster nodes with heterogeneous compute devices. The OpenCL specification tightly binds its workflow abstraction, or "command queue," to a specific device for the entire program. For best performance, the user has to find the ideal queue -- device mapping at command queue creation time, an effort that requires a thorough understanding of the match between the characteristics of all the underlying device architectures and the kernels in the program. In this paper, we propose to add scheduling attributes to the OpenCL context and command queue objects that can be leveraged by an intelligent runtime scheduler to automatically perform ideal queue - device mapping. Our proposed extensions enable the average OpenCL programmer to focus on the algorithm design rather than scheduling and automatically gain performance without sacrificing programmability. As an example, we design and implement an OpenCL runtime for task-parallel workloads, called MultiCL, which efficiently schedules command queues across devices. Within MultiCL, we implement several key optimizations to reduce runtime overhead. Our case studies include the SNU-NPB OpenCL benchmark suite and a real-world seismology simulation. We show that, on average, users have to apply our proposed scheduler extensions to only four source lines of code in existing OpenCL applications in order to automatically benefit from our runtime optimizations. We also show that MultiCL always maps command queues to the optimal device set with negligible runtime overhead.
Ashwin M. Aji, Antonio J. Peña, Pavan Balaji, Wu-chun Feng
CLUSTER4
2015 Rapid Screening of Transformed Data Leaks with Efficient Algorithms and Parallel Computing
abstract
The leak of sensitive data on computer systems poses a serious threat to organizational security. Organizations need to identify the exposure of sensitive data by screening the content in storage and transmission, i.e., to detect sensitive information being stored or transmitted in the clear. However, detecting the exposure of sensitive information is challenging due to data transformation in the content. Transformations (such as insertion, deletion) result in highly unpredictable leak patterns. Existing automata-based string matching algorithms are impractical for detecting transformed data leaks, because of its formidable complexity when modeling the required regular expressions. We design two new algorithms for detecting long and transformed data leaks. Our system achieves high detection accuracy in recognizing transformed leaks compared to the state-of-the-art inspection methods. We parallelize our prototype on graphics processing unit and demonstrate the strong scalability of our detection solution required by a sizable organization.
Xiaokui Shu, Jing Zhang 0039, Danfeng Yao, Wu-chun Feng
CODASPY4
2015 GLAF: A Visual Programming and Auto-tuning Framework for Parallel Computing
abstract
The past decade's computing revolution has delivered parallel hardware to the masses. However, the ability to exploit its capabilities and ignite scientific breakthrough at a proportionate level remains a challenge due to the lack of parallel programming expertise. Although different solutions have been proposed to facilitate harvesting the seeds of parallel computing, most target seasoned programmers and ignore the special nature of a target audience like domain experts. This paper addresses the challenge of realizing a programming abstraction and implementing an integrated development framework for this audience. We present GLAF -- a grid-based language and auto-parallelizing, auto-tuning framework. Its key elements are its intuitive visual programming interface, which attempts to render expressing and validating an algorithm easier for domain experts, and its ability to automatically generate efficient serial and parallel Fortran and C code, including potentially beneficial code modifications (e.g., With respect to data layout). We find that the above features assist novice programmers to avoid common programming pitfalls and provide fast implementations.
Konstantinos Krommydas, Ruchira Sasanka, Wu-chun Feng
ICPP3
2015 ASPaS: A Framework for Automatic SIMDization of Parallel Sorting on x86-based Many-core Processors
abstract
Due to the difficulty that modern compilers have in vectorizing applications on vector-extension architectures, programmers resort to manually programming vector registers with intrinsics in order to achieve better performance. However, the continued growth in the width of registers and the evolving library of intrinsics make such manual optimizations tedious and error-prone. Hence, we propose a framework for the Automatic SIMDization of Parallel Sorting (ASPaS) on x86-based multicore and manycore processors. That is, ASPaS takes any sorting network and a given instruction set architecture (ISA) as inputs and automatically generates vectorized code for that sorting network.
Kaixi Hou, Hao Wang 0002, Wu-chun Feng
ICS3
2015 Design and Evaluation of Scalable Concurrent Queues for Many-Core Architectures
abstract
As core counts increase and as heterogeneity becomes more common in parallel computing, we face the prospect of programming hundreds or even thousands of concurrent threads in a single shared-memory system. At these scales, even highly-efficient concurrent algorithms and data structures can become bottlenecks, unless they are designed from the ground up with throughput as their primary goal.
Thomas Scogland, Wu-chun Feng
ICPE2
2015 Accelerating Bioinformatics Applications via Emerging Parallel Computing Systems
abstract
The papers in this issue focus on advanced parallel computing systems for bioinformatics applications. This papers provide a forum to publish recent advances in the improvement of handling bioinformatics problems on emerging parallel computing systems. These systems can be characterized by exploiting different types of parallelism, including fine-grained versus coarse-grained and thread-level parallelism versus datalevel parallelism versus request-level parallelism. Hence, parallel computing systems based on multi- and many-core CPUs, many-core GPUs, vector processors, or FPGAs offer the promise to massively accelerate many bioinformatics algorithms and applications, ranging from computeintensive to data-intensive. Such computing systems are increasingly ubiquitous, ranging from “big iron” datacenter supercomputers and datacenter cloud computing down to GPU-accelerated smartphones and laptops.
Juan Antonio Gómez Pulido, Bertil Schmidt, Wu-chun Feng
IEEE ACM Trans. Comput. Biol. Bioinform.3
2015 CoreTSAR: Core Task-Size Adapting Runtime
abstract
Heterogeneity continues to increase at all levels of computing, with the rise of accelerators such as GPUs, FPGAs, and other co-processors into everything from desktops to supercomputers. As a consequence, efficiently managing such disparate resources has become increasingly complex. CoreTSAR seeks to reduce this complexity by adaptively worksharing parallel-loop regions across compute resources without requiring any transformation of the code within the loop. Our results show performance improvements of up to three-fold over a current state-of-the-art heterogeneous task scheduler as well as linear performance scaling from a single GPU to four GPUs for many codes. In addition, CoreTSAR demonstrates a robust ability to adapt to both a variety of workloads and underlying system configurations.
Thomas Scogland, Wu-chun Feng, Barry Rountree, Bronis R. de Supinski
IEEE Trans. Parallel Distributed Syst.2
2014 Locality-aware memory association for multi-target worksharing in OpenMP
abstract
No abstract available.
Thomas Scogland, Wu-chun Feng
PACT2
2014 On the characterization of OpenCL dwarfs on fixed and reconfigurable platforms
abstract
The proliferation of heterogeneous computing platforms presents the parallel computing community with new challenges. One such challenge entails evaluating the efficacy of such parallel architectures and identifying the architectural innovations that ultimately benefit applications. To address this challenge, we need benchmarks that capture the execution patterns (i.e., dwarfs or motifs) of applications, both present and future, in order to guide future hardware design. Furthermore, we desire a common programming model for the benchmarks that facilitates code portability across a wide variety of different processors (e.g., CPU, APU, GPU, FPGA, DSP) and computing environments (e.g., embedded, mobile, desktop, server). As such, we present the latest release of OpenDwarfs, a benchmark suite that currently realizes the Berkeley dwarfs in OpenCL, a vendor-agnostic and open-standard computing language for parallel computing. Using OpenDwarfs, we characterize a diverse set of fixed and reconfigurable parallel platforms: multicore CPUs, discrete and integrated GPUs, Intel Xeon Phi coprocessor, as well as a FPGA. We describe the computation and communication patterns exposed by a representative set of dwarfs, obtain relevant profiling data and execution information, and draw conclusions that highlight the complex interplay between dwarfs' patterns and the underlying hardware architecture of modern parallel platforms.
Konstantinos Krommydas, Wu-chun Feng, Muhsen Owaida, Christos D. Antonopoulos, Nikolaos Bellas
ASAP2
2014 Runtime Adaptation for Autonomic Heterogeneous Computing
abstract
Heterogeneity is increasing at all levels of computing, certainly with the rise in general purpose computing with GPUs in everything from phones to supercomputers. More quietly it is increasing with the rise of NUMA systems, hierarchical caching, OS noise, and a myriad of other factors. As heterogeneity becomes a fact of life at every level of computing, efficiently managing heterogeneous compute resources is becoming a critical task. The focus of my dissertation is developing methods and systems to allow software to adapt to the heterogeneous hardware it finds at runtime. The goal is to make the complex functions of heterogeneous computing autonomic, handling load balancing, memory coherence and other performance critical factors in the runtime. The investigation began by studying heterogeneity caused by system topology and resource contention in MPI applications. Since then the focus has shifted to work-sharing across CPU and GPU resources for accelerated OpenMP, and automatically managing the hardware capability imbalances between these resources. Moving forward, I propose to produce a system extending upon both previous approaches to offer work-sharing, topology aware affinity management, as well as novel automated memory transformations to reduce communication and increase memory access efficiency.
Thomas Scogland, Wu-chun Feng
CCGRID2
2014 Enabling Efficient Power Provisioning for Enterprise Applications
abstract
The increasing demand for computation and the commensurate rise in the power density of data centers have led to increased costs associated with constructing and operating a data center. Exacerbating such costs, data centers are often over-provisioned to avoid costly outages associated with the potential overloading of electrical circuitry. However, such over-provisioning is often unnecessary since a data center rarely operates at its maximum capacity. It is imperative that we maximize the use of the available power budget in order to enhance the efficiency of data centers. On the other hand, introducing power constraints to improve the efficiency of a data center can cause unacceptable violation of performance agreements (i.e., throughput and response time constraints). As such, we present a thorough empirical study of performance under power constraints as well as a runtime system to set appropriate power constraints for meeting strict performance targets. In this paper, we design a runtime system based on a load prediction model and an optimization framework to set the appropriate power constraints to meet specific performance targets. We then present the effects of our runtime system on energy proportionality, average power, performance, and instantaneous power consumption of enterprise applications. Our results shed light on mechanisms to tune the power provisioned for a server under strict performance targets and opportunities to improve energy proportionality and instantaneous power consumption via power limiting.
Balaji Subramaniam, Wu-chun Feng
CCGRID2
2014 SLAM: scalable locality-aware middleware for I/O in scientific analysis and visualization
abstract
Whereas traditional scientific applications are computationally intensive, recent applications require more data-intensive analysis and visualization. As the computational power and size of compute clusters continue to increase, the I/O read rates and associated network cost for these data-intensive applications create a serious performance bottleneck when faced with the massive data sets of today's "big data" era.
Jiangling Yin, Jun Wang 0001, Wu-chun Feng, Xuhong Zhang 0002, Junyao Zhang 0007
HPDC3
2014 Petascale Application of a Coupled CPU-GPU Algorithm for Simulation and Analysis of Multiphase Flow Solutions in Porous Medium Systems
abstract
Large-scale simulation can provide a wide range of information needed to develop and validate theoretical models for multiphase flow in porous medium systems. In this paper, we consider a coupled solution in which a multiphase flow simulator is coupled to an analysis approach used to extract the interfacial geometries as the flow evolves. This has been implemented using MPI to target heterogeneous nodes equipped with GPUs. The GPUs evolve the multiphase flow solution using the lattice Boltzmann method while the CPUs compute up scaled measures of the morphology and topology of the phase distributions and their rate of evolution. Our approach is demonstrated to scale to 4,096 GPUs and 65,536 CPU cores to achieve a maximum performance of 244,754 million-lattice-node updates per second (MLUPS) in double precision execution on Titan. In turn, this approach increases the size of systems that can be considered by an order of magnitude compared with previous work and enables detailed in situ tracking of averaged flow quantities at temporal resolutions that were previously impossible. Furthermore, it virtually eliminates the need for post-processing and intensive I/O and mitigates the potential loss of data associated with node failures.
James E. McClure, Hao Wang 0002, Jan F. Prins, Cass T. Miller, Wu-chun Feng
IPDPS5
2014 cuBLASTP: Fine-Grained Parallelization of Protein Sequence Search on a GPU
abstract
BLAST, short for Basic Local Alignment Search Tool, is a fundamental algorithm in the life sciences that compares biological sequences. However, with the advent of next-generation sequencing (NGS) and increase in sequence read-lengths, whether at the outset or downstream from NGS, the exponential growth of sequence databases is arguably outstripping our ability to analyze the data. Though several recent studies have utilized the graphics processing unit (GPU) to speedup the BLAST algorithm for searching protein sequences (i.e., BLASTP), these studies used coarse-grained parallel approaches, where one sequence alignment is mapped to only one thread. Moreover, due to the irregular memory access patterns in BLASTP, there remain significant challenges to map the most time-consuming phases (i.e., hit detection and ungapped extension) to the GPU using a fine-grained multithreaded approach. To address the above issues, we propose cuBLASTP, an efficient fine-grained BLASTP implementation for the GPU using CUDA. Our cuBLASTP realization encompasses many research contributions, including (1) memory-access reordering to reorder hits from column-major order to diagonal-major order, (2) position-based indexing to map a hit with a packed data structure to a bin, (3) aggressive hit filtering to eliminate hits beyond the threshold distance along the diagonal, (4) diagonal-based parallelism and hit-based parallelism for ungapped extension to extend sequences with different lengths in databases, and (5) hierarchical buffering to reduce memory-access overhead for the core data structures. The experimental results show that on a NVIDIA Kepler GPU, cuBLASTP delivers up to a 5.0-fold speedup over sequential FSA-BLAST and a 3.7-fold speedup over multithreaded NCBI-BLAST for the overall program execution. In addition, compared with GPU-BLASTP (the fastest GPU implementation of BLASTP to date), cuBLASTP achieves up to a 2.8-fold speedup for the kernel execution on the GPU and a 1.8-fold speedup for the overall program execution.
Jing Zhang 0039, Hao Wang 0002, Heshan Lin, Wu-chun Feng
IPDPS4
2014 Aeromancer: A Workflow Manager for Large-Scale MapReduce-Based Scientific Workflows
abstract
The Hadoop framework has gained significant attention from the scientific community due to its applicability to large-scale data analysis in many areas. This analysis often involves multiple stages of processing, which in turn, constitutes a workflow. While some stages of a workflow are mandatory, others are subject to the type of analysis to be done. In addition, a workflow may possess data dependencies between stages that must be enforced, and it may exhibit varying levels of sensitivity. The resources needed for such data analysis can range from a laptop to in-house clusters (or private cloud) to a public cloud. Managing such workflows, while using such a gamut of computing resources, is an unnecessarily arduous task for domain scientists. To address the above challenges, we present Aeromancer, a feature-rich workflow manager for running Map Reduce-based workflows that utilizes both client and cloud resources. Aeromancer offers an ensemble of features, including the simultaneous use of client resources (e.g., On-premises clusters) and public cloud resources, automatic data-dependency and data-transfer handling, intra-flow, on-demand cluster provisioning, and support for directed-acyclic graphs (DAGs). To demonstrate its functionality, we apply Aeromancer to several bioinformatics pipelines, as part of a "big data" case study in the life sciences, which seeks to increase the adoption of hybrid computing environments, including the emerging "client cloud" computing model, for running data-intensive workflows.
Mohamed Nabeel, Nabanita Maji, Jing Zhang 0039, Nataliya Timoshevskaya, Wu-chun Feng
TrustCom5
2014 A power-measurement methodology for large-scale, high-performance computing
abstract
Improvement in the energy efficiency of supercomputers can be accelerated by improving the quality and comparability of efficiency measurements. The ability to generate accurate measurements at extreme scale are just now emerging. The realization of system-level measurement capabilities can be accelerated with a commonly adopted and high quality measurement methodology for use while running a workload, typically a benchmark. This paper describes a methodology that has been developed collaboratively through the Energy Efficient HPC Working Group to support architectural analysis and comparative measurements for rankings, such as the Top500 and Green500. To support measurements with varying amounts of effort and equipment required we present three distinct levels of measurement, which provide increasing levels of accuracy. Level 1 is similar to the Green500 run rules today, a single average power measurement extrapolated from a subset of a machine. Level 2 is more comprehensive, but still widely achievable. Level 3 is the most rigorous of the three methodologies but is only possible at a few sites. However, the Level 3 methodology generates a high quality result that exposes details that the other methodologies may miss. In addition, we present case studies from the Leibniz Supercomputing Centre (LRZ), Argonne National Laboratory (ANL) and Calcul Québec Université Laval that explore the benefits and difficulties of gathering high quality, system-level measurements on large-scale machines.
Thomas Scogland, Craig P. Steffen, Torsten Wilde, Florent Parent, Susan Coghlan, Natalie J. Bates, Wu-chun Feng, Erich Strohmaier
ICPE7
2014 SDAFT: A novel scalable data access framework for parallel BLAST
Jiangling Yin, Junyao Zhang 0007, Jun Wang 0001, Wu-chun Feng
Parallel Comput.4
2013 Optimizing Burrows-Wheeler Transform-Based Sequence Alignment on Multicore Architectures
abstract
Computational biology sequence alignment tools using the Burrows-Wheeler Transform (BWT) are widely used in next-generation sequencing (NGS) analysis. However, despite extensive optimization efforts, the performance of these tools still cannot keep up with the explosive growth of sequencing data. Through an in-depth performance analysis of BWA, a popular BWT-based aligner on multicore architectures, we demonstrate that such tools are limited by memory bandwidth due to their irregular memory access patterns. We then propose a locality-aware implementation of BWA that aims at optimizing its performance by better exploiting the caching mechanisms of modern multicore processors. Experimental results show that our improved BWA implementation can reduce last-level cache (LLC) misses by 30% and translation look aside buffer (TLB) misses by 20%, resulting in up to 2.6-fold speedup over the original BWA implementation.
Jing Zhang 0039, Heshan Lin, Pavan Balaji, Wu-chun Feng
CCGRID4
2013 Cascaded TCP: Applying pipelining to TCP for efficient communication over wide-area networks
abstract
The bandwidth utilization in traditional TCP protocols (e.g., TCP New Reno) suffers over high-latency and high-bandwidth links due to the inherent characteristics of TCP congestion control. Conventional methods of improving throughput cannot be applied per se for streaming applications. The challenge is exacerbated by “big data” applications such as with the Long Wavelength Array data that is generated at a rate of up to 4 terabytes per hour. To improve bandwidth utilization, we introduce layer-4 relay(s) that enable the pipelining of TCP connections. That is, a traditional end-to-end connection is split into independent streams, each with shorter latencies, that are then concatenated (or cascaded) together to form an equivalent end-to-end TCP connection. This addresses the root cause by decreasing the latency over which the congestion-control protocol operates. To understand when relays are beneficial, we present an analytical model, empirical data and its analyses, to validate our argument and to characterize the impact of latency and available bandwidth on throughput. We also provide insight into how relays may be setup to achieve better bandwidth utilization.
Umar Kalim, Mark K. Gardner, Eric J. Brown, Wu-chun Feng
GLOBECOM4
2013 On the efficacy of GPU-integrated MPI for scientific applications
Ashwin M. Aji, Lokendra S. Panwar, Milind Chabbi, Karthik Murthy, Pavan Balaji, Keith R. Bisset, James Dinan, Wu-chun Feng, John M. Mellor-Crummey, Xiaosong Ma, Rajeev Thakur
HPDC9
2013 Accelerating fast Fourier Transform for wideband channelization
abstract
Wideband channelization is a compute-intensive task with performance requirements that are arguably greater than what current multi-core CPUs can provide. To date, researchers have used dedicated hardware such as field programmable gate arrays (FPGAs) to address the performance-critical aspects of the channelizer. In this work, we assess the viability of the graphics processing unit (GPU) to achieve the necessary performance. In particular, we focus on the fast Fourier Transform (FFT) stage of a wideband channelizer. While there exists previous work for FFT on a NVIDIA GPU, the substantially higher peak floating-point performance of an AMD GPU has been less explored. Thus, we consider three generations of AMD GPUs and provide insight into the optimization of FFT on these platforms. Our architecture-aware approach across three different generations of AMD GPUs outperforms a multithreaded Intel Sandy Bridge CPU with vector extensions by factors of 4.3, 4.9, and 6.6 on the Radeon HD 5870, 6970, and 7970, respectively.
Carlo C. del Mundo, Vignesh Adhinarayanan, Wu-chun Feng
ICC3
2013 Seamless Migration of Virtual Machines across Networks
abstract
Current technologies that support live migration require that the virtual machine (VM) retain its IP network address. As a consequence, VM migration is oftentimes restricted to movement within an IP subnet or entails interrupted network connectivity to allow the VM to migrate. Thus, migrating VMs beyond subnets becomes a significant challenge for the purposes of load balancing, moving computation close to data sources, or connectivity recovery during natural disasters. Conventional approaches use tunneling, routing, and layer-2 expansion methods to extend the network to geographically disparate locations, thereby transforming the problem of migration between subnets to migration within a subnet. These approaches, however, increase complexity and involve considerable human involvement. The contribution of our paper is to address the aforementioned shortcomings by enabling VM migration across subnets and doing so with uninterrupted network connectivity. We make the case that decoupling IP addresses from the notion of transport endpoints is the key to solving a host of problems, including seamless VM migration and mobility. We demonstrate that VMs can be migrated seamlessly between different subnets - without losing network state - by presenting a backward-compatible prototype implementation and a case study.
Umar Kalim, Mark K. Gardner, Eric J. Brown, Wu-chun Feng
ICCCN4
2013 pVOCL: Power-Aware Dynamic Placement and Migration in Virtualized GPU Environments
abstract
Power-hungry Graphics processing unit (GPU) accelerators are ubiquitous in high performance computing data centers today. GPU virtualization frameworks introduce new opportunities for effective management of GPU resources by decoupling them from application execution. However, power management of GPU-enabled server clusters faces significant challenges. The underlying system infrastructure shows complex power consumption characteristics depending on the placement of GPU workloads across various compute nodes, power-phases and cabinets in a datacenter. GPU resources need to be scheduled dynamically in the face of time-varying resource demand and peak power constraints. We propose and develop a power-aware virtual OpenCL (pVOCL) framework that controls the peak power consumption and improves the energy efficiency of the underlying server system through dynamic consolidation and power-phase topology aware placement of GPU workloads. Experimental results show that pVOCL achieves significant energy savings compared to existing power management techniques for GPU-enabled server clusters, while incurring negligible impact on performance. It drives the system towards energy-efficient configurations by taking an optimal sequence of adaptation actions in a virtualized GPU environment and meanwhile keeps the power consumption below the peak power budget.
Palden Lama, Yan Li 0005, Ashwin M. Aji, Pavan Balaji, James Dinan, Shucai Xiao, Yunquan Zhang, Wu-chun Feng, Rajeev Thakur, Xiaobo Zhou 0002
ICDCS8
2013 Wideband Channelization for Software-Defined Radio via Mobile Graphics Processors
abstract
Wideband channelization is a computationally intensive task within software-defined radio (SDR). To support this task, the underlying hardware should provide high performance and allow flexible implementations. Traditional solutions use field-programmable gate arrays (FPGAs) to satisfy these requirements. While FPGAs allow for flexible implementations, realizing a FPGA implementation is a difficult and time-consuming process. On the other hand, multicore processors while more programmable, fail to satisfy performance requirements. Graphics processing units (GPUs) overcome the above limitations. However, traditional GPUs are power-hungry and can consume as much as 350 watts, making them ill-suited for many SDR environments, particularly those that are battery-powered. Here we explore the viability of low-power mobile graphics processors to simultaneously overcome the limitations of performance, flexibility, and power. Via execution profiling and performance analysis, we identify major bottlenecks in mapping the wideband channelization algorithm onto these devices and adopt several optimization techniques to achieve multiplicative speed-up over a multithreaded implementation. Overall, our approach delivers a speedup of up to 43-fold on the discrete AMD Radeon HD 6470M GPU and 27-fold on the integrated AMD Radeon HD 6480G GPU, when compared to a vectorized and multithreaded version running on the AMD A4-3300M CPU.
Vignesh Adhinarayanan, Wu-chun Feng
ICPADS2
2013 On the Portability of the OpenCL Dwarfs on Fixed and Reconfigurable Parallel Platforms
abstract
The proliferation of heterogeneous computing systems presents the parallel computing community with the challenge of porting legacy and emerging applications to multiple processors with diverse programming abstractions. OpenCL is a vendor-agnostic and industry-supported programming model that offers code portability on heterogeneous platforms, allowing applications to be developed once and deployed "anywhere." In this paper, we use the OpenCL implementation of the Open Dwarfs, a benchmark suite that captures patterns of computation and communication common to classes of important applications, as delineated by Berkeley's Dwarfs. We evaluate portability across multicore CPU, GPU, APU (CPUs+GPUs on a die), the Intel Xeon Phi co-processor, and the FPGA. To realize FPGA portability, we exploit SOpenCL (Silicon OpenCL), a CAD tool that automatically converts OpenCL kernels to customizable hardware accelerators. We show that a single, unmodified OpenCL code base, i.e., Open Dwarfs, can be effectively used to target multiple, architecturally diverse platforms.
Konstantinos Krommydas, Muhsen Owaida, Christos D. Antonopoulos, Nikolaos Bellas, Wu-chun Feng
ICPADS5
2013 On the Programmability and Performance of Heterogeneous Platforms
abstract
General-purpose computing on an ever-broadening array of parallel devices has led to an increasingly complex and multi-dimensional landscape with respect to programmability and performance optimization. The growing diversity of parallel architectures presents many challenges to the domain scientist, including device selection, programming model, and level of investment in optimization. All of these choices influence the balance between programmability and performance. In this paper, we characterize the performance achievable across a range of optimizations, along with their programmability, for multi- and many-core platforms - specifically, an Intel Sandy Bridge CPU, Intel Xeon Phi co-processor, and NVIDIA Kepler K20 GPU - in the context of an n-body, molecular-modeling application called GEM. Our systematic approach to optimization delivers implementations with speed-ups of 194.98×, 885.18×, and 1020.88× on the CPU, Xeon Phi, and GPU, respectively, over the naive serial version. Beyond the speed-ups, we characterize the incremental optimization of the code from naive serial to fully hand-tuned on each platform through four distinct phases of increasing complexity to expose the strengths and weaknesses of the programming models offered by each platform.
Konstantinos Krommydas, Thomas Scogland, Wu-chun Feng
ICPADS3
2013 Online Performance Projection for Clusters with Heterogeneous GPUs
abstract
We present a fully automated approach to project the relative performance of an OpenCL program over different GPUs. Performance projections can be made within a small amount of time, and the projection overhead stays relatively constant with the input data size. As a result, the technique can help runtime tools make dynamic decisions about which GPU would run faster for a given kernel. Usage cases of this technique include scheduling or migrating GPU workloads over a heterogeneous cluster with different types of GPUs.
Lokendra S. Panwar, Ashwin M. Aji, Jiayuan Meng, Pavan Balaji, Wu-chun Feng
ICPADS5
2013 Towards energy-proportional computing for enterprise-class server workloads
abstract
Massive data centers housing thousands of computing nodes have become commonplace in enterprise computing, and the power consumption of such data centers is growing at an unprecedented rate. Adding to the problem is the inability of the servers to exhibit energy proportionality, i.e., provide energy-efficient execution under all levels of utilization, which diminishes the overall energy efficiency of the data center. It is imperative that we realize effective strategies to control the power consumption of the server and improve the energy efficiency of data centers. With the advent of Intel Sandy Bridge processors, we have the ability to specify a limit on power consumption during runtime, which creates opportunities to design new power-management techniques for enterprise workloads and make the systems that they run on more energy proportional.
Balaji Subramaniam, Wu-chun Feng
ICPE2
2013 Characterizing the challenges and evaluating the efficacy of a CUDA-to-OpenCL translator
Mark K. Gardner, Paul Sathre, Wu-chun Feng, Gabriel Martinez
Parallel Comput.3
2012 Transparent Accelerator Migration in a Virtualized GPU Environment
abstract
This paper presents a framework to support transparent, live migration of virtual GPU accelerators in a virtualized execution environment. Migration is a critical capability in such environments because it provides support for fault tolerance, on-demand system maintenance, resource management, and load balancing in the mapping of virtual to physical GPUs. Techniques to increase responsiveness and reduce migration overhead are explored. The system is evaluated by using four application kernels and is demonstrated to provide low migration overheads. Through transparent load balancing, our system provides a speedup of 1.7 to 1.9 for three of the four application kernels.
Shucai Xiao, Pavan Balaji, James Dinan, Rajeev Thakur, Susan Coghlan, Heshan Lin, Gaojin Wen, Jue Hong, Wu-chun Feng
CCGRID10
2012 Heterogeneous Task Scheduling for Accelerated OpenMP
abstract
Heterogeneous systems with CPUs and computational accelerators such as GPUs, FPGAs or the upcoming Intel MIC are becoming mainstream. In these systems, peak performance includes the performance of not just the CPUs but also all available accelerators. In spite of this fact, the majority of programming models for heterogeneous computing focus on only one of these. With the development of Accelerated Open MP for GPUs, both from PGI and Cray, we have a clear path to extend traditional Open MP applications incrementally to use GPUs. The extensions are geared toward switching from CPU parallelism to GPU parallelism. However they do not preserve the former while adding the latter. Thus computational potential is wasted since either the CPU cores or the GPU cores are left idle. Our goal is to create a runtime system that can intelligently divide an accelerated Open MP region across all available resources automatically. This paper presents our proof-of-concept runtime system for dynamic task scheduling across CPUs and GPUs. Further, we motivate the addition of this system into the proposed Open MP for Accelerators standard. Finally, we show that this option can produce as much as a two-fold performance improvement over using either the CPU or GPU alone.
Thomas Scogland, Barry Rountree, Wu-chun Feng, Bronis R. de Supinski
IPDPS3
2012 Automatic NUMA characterization using Cbench
abstract
Clusters of seemingly homogeneous compute nodes are increasingly heterogeneous within each node due to replication and distribution of node-level subsystems. This intra-node heterogeneity can adversely affect program execution performance by inflicting additional data-access costs when accessing non-local data. In this work-in-progress paper, we present extensions to the Cbench Scalable Testing Framework for analyzing main memory and PCIe data-access performance in modern NUMA architectures. The information provided by this tool will be of use for task scheduling, performance modeling, and evaluation of NUMA systems.
Ryan K. Braithwaite, Wu-chun Feng, Patrick S. McCormick
ICPE2
2012 OpenCL and the 13 dwarfs: a work in progress
abstract
In the past, evaluating the architectural innovation of parallel computing devices relied on a benchmark suite based on existing programs, e.g., EEMBC or SPEC. However, with the growing ubiquity of parallel computing devices, we argue that it is unclear how best to express parallel computation, and hence, a need exists to identify a higher level of abstraction for reasoning about parallel application requirements. Therefore, the goal of this combination "Work-in-Progress and Vision" paper is to delineate application requirements in a manner that is not overly specific to individual applications or the optimizations used for certain hardware platforms, so that we can draw broader conclusions about hardware requirements. Our initial effort, dubbed "OpenCL and the 13 Dwarfs" or OCD for short, realizes Berkeley's 13 computational dwarfs of scientific computing in OpenCL, where each dwarf captures a pattern of computation and communication that is common to a class of important applications.
Wu-chun Feng, Heshan Lin, Thomas Scogland, Jing Zhang 0039
ICPE1
2012 Multi-dimensional characterization of electrostatic surface potential computation on graphics processors
abstract
BACKGROUND: Calculating the electrostatic surface potential (ESP) of a biomolecule is critical towards understanding biomolecular function. Because of its quadratic computational complexity (as a function of the number of atoms in a molecule), there have been continual efforts to reduce its complexity either by improving the algorithm or the underlying hardware on which the calculations are performed. RESULTS: We present the combined effect of (i) a multi-scale approximation algorithm, known as hierarchical charge partitioning (HCP), when applied to the calculation of ESP and (ii) its mapping onto a graphics processing unit (GPU). To date, most molecular modeling algorithms perform an artificial partitioning of biomolecules into a grid/lattice on the GPU. In contrast, HCP takes advantage of the natural partitioning in biomolecules, which in turn, better facilitates its mapping onto the GPU. Specifically, we characterize the effect of known GPU optimization techniques like use of shared memory. In addition, we demonstrate how the cost of divergent branching on a GPU can be amortized across algorithms like HCP in order to deliver a massive performance boon. CONCLUSIONS: We accelerated the calculation of ESP by 25-fold solely by parallelization on the GPU. Combining GPU and HCP, resulted in a speedup of at most 1,860-fold for our largest molecular structure. The baseline for these speedups is an implementation that has been hand-tuned SSE-optimized and parallelized across 16 cores on the CPU. The use of GPU does not deteriorate the accuracy of our results.
Mayank Daga, Wu-chun Feng
BMC Bioinform.2
2011 Performance Characterization and Optimization of Atomic Operations on AMD GPUs
abstract
Atomic operations are important building blocks in supporting general-purpose computing on graphics processing units (GPUs). For instance, they can be used to coordinate execution between concurrent threads, and in turn, assist in constructing complex data structures such as hash tables or implementing GPU-wide barrier synchronization. While the performance of atomic operations has improved substantially on the latest NVIDIA Fermi-based GPUs, system-provided atomic operations still incur significant performance penalties on AMD GPUs. A memory-bound kernel on an AMD GPU, for example, can suffer severe performance degradation when including an atomic operation, even if the atomic operation is never executed. In this paper, we first quantify the performance impact of atomic instructions to application kernels on AMD GPUs. We then propose a novel software-based implementation of atomic operations that can significantly improve the overall kernel performance. We evaluate its performance against the system-provided atomic using two micro-benchmarks and four real applications. The results show that using our software based atomic operations on an AMD GPU can speedup an application kernel by 67-fold over the same application kernel but with the (default) system-provided atomic operations.
Marwa K. Elteir, Heshan Lin, Wu-chun Feng
CLUSTER3
2011 Energy-efficient E-puting everywhere
abstract
Throughout the 1990s and much of the 2000s, the halls of high-performance computing (HPC) echoed with sentiments like the following: "In HPC, no one cares about energy efficiency or power consumption, and no one ever will." While such extreme talk has subsided, computational performance (or speed) via parallelism still rule the roost. Conversely, one could argue that the consumer electronics space has taken a complementary approach, where energy efficiency and power consumption have been first-order design constraints, with speed only needing to be "good enough" for ordinary daily tasks. However, the increasing computational demands that end users will place on (consumer) electronics, such as computations for personalized medicine, point to the need for "supercomputing in small spaces" (http://sss.cs.vt.edu/). This trend, in turn, will elevate performance to be a first-order design constraint in consumer electronics, on par with energy efficiency and power consumption. This talk will discuss how a "trickle-up" approach will deliver supercomputing in small spaces via an increasingly converged world of energy-efficient (consumer) electronics and computing, or e-puting.
Wu-chun Feng
HPDC1
2011 Restoring End-to-End Resilience in the Presence of Middleboxes
abstract
The philosophy upon which the Internet was built places the intelligence close to the edge. As the Internet has matured, intermediate devices or middleboxes, such as firewalls or application gateways, have been introduced, thereby weakening the end-to-end nature of the network. As a result, applications must often modify their behavior to accommodate the middleboxes. This is is especially true in the case of transient failure of stateful devices. The failure of a middlebox causes it to lose the state it maintained, causing the failure of the associated TCP connections. Rather than assign the responsibility for recovery to applications, we incorporate a mechanism called an isolation boundary into TCP itself. The isolation boundary maintains a small amount of state across TCP connections, thus enabling reconnection. Furthermore, it does so without breaking backward compatibility with existing TCP. We present an implementation of the isolation boundary in the FreeBSD kernel and demonstrate its backward compatibility with TCP. We quantify the performance impact of the proposed mechanism on the establishment of new and resumed connections for both legacy and extended TCP stacks.
Eric J. Brown, Mark K. Gardner, Umar Kalim, Wu-chun Feng
ICCCN4
2011 AVS video decoder on multicore systems: Optimizations and tradeoffs
abstract
Newer video compression standards provide high video quality and greater compression efficiency, compared to their predecessors. Their increased complexity can be outbalanced by leveraging all the levels of available parallelism, task- and data-level, using available off-the-shelf hardware, such as current generation's chip multiprocessors. As we move to more cores though, scalability issues arise and need to be tackled in order to take advantage of the abundant computational power. In this paper we evaluate a previously implemented parallel version of the AVS video decoder on the experimental 32-core Intel Manycore Testing Lab. We examine this previous version's performance bottlenecks and scalability issues and introduce a distributed queue implementation as the proposed solution. Finally, we provide insight on separate optimizations regarding inter macroblocks and investigate performance variations and tradeoffs, when combined with a distributed queue scheme.
Konstantinos Krommydas, Christos D. Antonopoulos, Nikolaos Bellas, Wu-chun Feng
ICME4
2011 Architecture-Aware Mapping and Optimization on a 1600-Core GPU
abstract
The graphics processing unit (GPU) continues to make in-roads as a computational accelerator for high-performance computing (HPC). However, despite its increasing popularity, mapping and optimizing GPU code remains a difficult task, it is a multi-dimensional problem that requires deep technical knowledge of GPU architecture. Although substantial literature exists on how to map and optimize GPU performance on the more mature NVIDIA CUDA architecture, the converse is true for OpenCL on an AMD GPU, such as the 1600-core AMD Radeon HD 5870 GPU. Consequently, we present and evaluate architecture-aware mapping and optimizations for the AMD GPU. The most prominent of which include (i) explicit use of registers, (ii) use of vector types, (iii) removal of branches, and (iv) use of image memory for global data. We demonstrate the efficacy of our AMD GPU mapping and optimizations by applying each in isolation as well as in concert to a large-scale, molecular modeling application called GEM. Via these AMD-specific GPU optimizations, our optimized OpenCL implementation on an AMD Radeon HD 5870 delivers more than a four-fold improvement in performance over the basic OpenCL implementation. In addition, it outperforms our optimized CUDA version on an NVIDIA GTX280 by 12%. Overall, we achieve a speedup of 371-fold over a serial but hand-tuned SSE version of our molecular modeling application, and in turn, a 46-fold speedup over an ideal scaling on an 8-core CPU.
Mayank Daga, Thomas Scogland, Wu-chun Feng
ICPADS3
2011 StreamMR: An Optimized MapReduce Framework for AMD GPUs
abstract
MapReduce is a programming model from Google that facilitates parallel processing on a cluster of thousands of commodity computers. The success of MapReduce in cluster environments has motivated several studies of implementing MapReduce on a graphics processing unit (GPU), but generally focusing on the NVIDIA GPU. Our investigation reveals that the design and mapping of the MapReduce framework needs to be revisited for AMD GPUs due to their notable architectural differences from NVIDIA GPUs. For instance, current state-of-the-art MapReduce implementations employ atomic operations to coordinate the execution of different threads. However, atomic operations can implicitly cause inefficient memory access, and in turn, severely impact performance. In this paper, we propose Streamer, an OpenCL MapReduce framework optimized for AMD GPUs. With efficient atomic-free algorithms for output handling and intermediate result shuffling, Stream MR is superior to atomic-based MapReduce designs and can outperform existing atomic-free MapReduce implementations by nearly five-fold on an AMD Radeon HD 5870.
Marwa K. Elteir, Heshan Lin, Wu-chun Feng, Thomas Scogland
ICPADS3
2011 CU2CL: A CUDA-to-OpenCL Translator for Multi- and Many-Core Architectures
abstract
The use of graphics processing units (GPUs) in high-performance parallel computing continues to become more prevalent, often as part of a heterogeneous system. For years, CUDA has been the de facto programming environment for nearly all general-purpose GPU (GPGPU) applications. In spite of this, the framework is available only on NVIDIA GPUs, traditionally requiring reimplementation in other frameworks in order to utilize additional multi- or many-core devices. On the other hand, OpenCL provides an open and vendor-neutral programming environment and runtime system. With implementations available for CPUs, GPUs, and other types of accelerators, OpenCL therefore holds the promise of a "write once, run anywhere" ecosystem for heterogeneous computing. Given the many similarities between CUDA and OpenCL, manually porting a CUDA application to OpenCL is typically straightforward, albeit tedious and error-prone. In response to this issue, we created CU2CL, an automated CUDA-to-OpenCL source-to-source translator that possesses a novel design and clever reuse of the Clang compiler framework. Currently, the CU2CL translator covers the primary constructs found in CUDA runtime API, and we have successfully translated many applications from the CUDA SDK and Rodinia benchmark suite. The performance of the automatically translated applications via CU2CL is on par with their manually ported counterparts.
Gabriel Martinez, Mark K. Gardner, Wu-chun Feng
ICPADS3
2011 Optimizing Dynamic Programming on Graphics Processing Units via Adaptive Thread-Level Parallelism
abstract
Dynamic programming (DP) is an important computational method for solving a wide variety of discrete optimization problems such as scheduling, string editing, packaging, and inventory management. In general, DP is classified into four categories based on the characteristics of the optimization equation. Because applications that are classified in the same category of DP have similar program behavior, the research community has sought to propose general solutions for parallelizing each category of DP. However, most existing studies focus on running DP on CPU-based parallel systems rather than on accelerating DP algorithms on the graphics processing unit (GPU). This paper presents the GPU acceleration of an important category of DP problems called nonserial polyadic dynamic programming (NPDP). In NPDP applications, the degree of parallelism varies significantly in different stages of computation, making it difficult to fully utilize the compute power of hundreds of processing cores in a GPU. To address this challenge, we propose a methodology that can adaptively adjust the thread-level parallelism in mapping a NPDP problem onto the GPU, thus providing sufficient and steady degrees of parallelism across different compute stages. We realize our approach in a real-world NPDP application -- the optimal matrix parenthesization problem. Experimental results demonstrate our method can achieve a speedup of 13.40 over the previously published GPU algorithm.
Chao-Chin Wu, Jenn-Yang Ke, Heshan Lin, Wu-chun Feng
ICPADS4
2011 Accelerating Protein Sequence Search in a Heterogeneous Computing System
abstract
The "Basic Local Alignment Search Tool'' (BLAST) is arguably the most widely used computational tool in bioinformatics. However, the computational power required for routine BLAST analysis has been outstripping Moore's Law due to the exponential growth in the size of the genomic sequence databases that BLAST searches on. To address the above issue, we propose the design and optimization of the BLAST algorithm for searching protein sequences (i.e., BLASTP) in a heterogeneous computing system. The end result is a BLASTP implementation that delivers a seven-fold speedup over the sequential BLASTP for the most computationally intensive phase (i.e., hit detection and ungapped extension) on a NVIDIA Fermi C2050 GPU. In addition, when pipelining the processing on a dual-core CPU and the NVIDIA Fermi GPU, our implementation can achieve a six-fold speedup for the overall program execution.
Shucai Xiao, Heshan Lin, Wu-chun Feng
IPDPS3
2011 Coordinating Computation and I/O in Massively Parallel Sequence Search
abstract
With the explosive growth of genomic information, the searching of sequence databases has emerged as one of the most computation and data-intensive scientific applications. Our previous studies suggested that parallel genomic sequence-search possesses highly irregular computation and I/O patterns. Effectively addressing these runtime irregularities is thus the key to designing scalable sequence-search tools on massively parallel computers. While the computation scheduling for irregular scientific applications and the optimization of noncontiguous file accesses have been well-studied independently, little attention has been paid to the interplay between the two. In this paper, we systematically investigate the computation and I/O scheduling for data-intensive, irregular scientific applications within the context of genomic sequence search. Our study reveals that the lack of coordination between computation scheduling and I/O optimization could result in severe performance issues. We then propose an integrated scheduling approach that effectively improves sequence-search throughput by gracefully coordinating the dynamic load balancing of computation and high-performance noncontiguous I/O.
Heshan Lin, Xiaosong Ma, Wu-chun Feng, Nagiza F. Samatova
IEEE Trans. Parallel Distributed Syst.3
2010 MOON: MapReduce On Opportunistic eNvironments
abstract
MapReduce offers an ease-of-use programming paradigm for processing large data sets, making it an attractive model for distributed volunteer computing systems. However, unlike on dedicated resources, where MapReduce has mostly been deployed, such volunteer computing systems have significantly higher rates of node unavailability. Furthermore, nodes are not fully controlled by the MapReduce framework. Consequently, we found the data and task replication scheme adopted by existing MapReduce implementations woefully inadequate for resources with high unavailability.
Heshan Lin, Xiaosong Ma, Jeremy S. Archuleta, Wu-chun Feng, Mark K. Gardner, Zhe Zhang 0005
HPDC4
2010 On the Goodput of TCP NewReno in Mobile Networks
abstract
Next-generation wireless networks such as LTE and WiMax can achieve throughputs of several Mbps with TCP. These higher throughputs, however, can easily be destroyed by frequent handoffs, which occur in urban environments due to shadowing. A primary reason for the throughput drop during handoffs is the out of order arrival of packets at the receiver. As a result, in this paper, we model the precise effect of packet-reordering on the goodput of TCP NewReno. Specifically, we develop a TCP NewReno model that captures the goodput of TCP as a function of round-trip time, average time duration between packet-reorder events, average number of packets reordered during every reorder event, and the congestion window threshold of TCP NewReno. We also developed an emulator that runs on a router to implement packet reordering events from time to time. We validate our NewReno model by comparing the goodput results obtained by transferring data between two hosts connected via the emulator to the goodput results that our model predicts.
Sushant Sharma, Donald W. Gillies, Wu-chun Feng
ICCCN3
2010 Enhancing MapReduce via Asynchronous Data Processing
abstract
The Map Reduce programming model simplifies large-scale data processing on commodity clusters by having users specify a map function that processes input key/value pairs to generate intermediate key/value pairs, and a reduce function that merges and converts intermediate key/value pairs into final results. Typical Map Reduce implementations such as Hadoop enforce barrier synchronization between the map and reduce phases, i.e., the reduce phase does not start until all map tasks are finished. In turn, this synchronization requirement can cause inefficient utilization of computing resources and can adversely impact performance. Thus, we present and evaluate two different approaches to cope with the synchronization drawback of existing Map Reduce implementations. The first approach, hierarchical reduction, starts a reduce task as soon as a predefined number of map tasks completes, it then aggregates the results of different reduce tasks following a tree structure. The second approach, incremental reduction, starts a predefined number of reduce tasks from the beginning and has each reduce task incrementally reduce records collected from map tasks. Together with our performance modeling, we evaluate different reducing approaches with two real applications on a 32-node cluster. The experimental results have shown that incremental reduction outperforms hierarchical reduction in general. Also, incremental reduction can speed-up the original Hadoop implementation by up to 35.33% for the word count application and 57.98% for the grep application. In addition, incremental reduction outperforms the original Hadoop in an emulated cloud environment with heterogeneous compute nodes.
Marwa K. Elteir, Heshan Lin, Wu-chun Feng
ICPADS3
2010 Inter-block GPU communication via fast barrier synchronization
abstract
The graphics processing unit (GPU) has evolved from a fixed-function processor with programmable stages to a programmable processor with many fixed-function components that deliver massive parallelism. Consequently, GPUs increasingly take advantage of the programmable processing power for general-purpose, non-graphics tasks, i.e., general-purpose computation on graphics processing units (GPGPU). However, while the GPU can massively accelerate data parallel (or task parallel) applications, the lack of explicit support for inter-block communication on the GPU hampers its broader adoption as a general-purpose computing device. Inter-block communication on the GPU occurs via global memory and then requires a barrier synchronization across the blocks, i.e., inter-block GPU communication via barrier synchronization. Currently, such synchronization is only available via the CPU, which in turn, incurs significant overhead. Thus, we seek to propose more efficient methods for inter-block communication. To systematically address this problem, we first present a performance model for the execution of kernels on GPUs. This performance model partitions the kernel’s execution time into three phases: (1) kernel launch to the GPU, (2) computation on the GPU, and (3) inter-block GPU communication via barrier synchronization. Using three well-known algorithms — FFT, dynamic programming, and bitonic sort — we show that the latter phase, i.e., inter-block GPU communication, can consume more than 50% of the overall execution time. Therefore, we propose three new approaches to inter-block GPU communication via barrier synchronization, all of which run only on the GPU: GPU simple synchronization, GPU tree-based synchronization, and GPU lock-free synchronization. We then evaluate the efficacy of each of these approaches in isolation via a micro-benchmark as well as integrated with the three aforementioned algorithms. For the micro-benchmark, the experimental results show that our GPU lock-free synchronization performs 7.8 times faster than CPU explicit synchronization and 3.7 times faster than CPU implicit synchronization. When integrated with the FFT, dynamic programming, and bitonic sort algorithms, our GPU lock-free synchronization improves the performance by 8%, 24%, and 39%, respectively, when compared to the more efficient CPU implicit synchronization.
Shucai Xiao, Wu-chun Feng
IPDPS2
2010 To GPU synchronize or not GPU synchronize?
abstract
The graphics processing unit (GPU) has evolved from being a fixed-function processor with programmable stages into a programmable processor with many fixed-function components that deliver massive parallelism. By modifying the GPU's stream processor to support “general-purpose computation” on the GPU (GPGPU), applications that perform massive vector operations can realize many orders-of-magnitude improvement in performance over a traditional processor, i.e., CPU. However, the breadth of general-purpose computation that can be efficiently supported on a GPU has largely been limited to highly dataparallel or task-parallel applications due to the lack of explicit support for communication between streaming multiprocessors (SMs) on the GPU. Such communication can occur via the global memory of a GPU, but it then requires a barrier synchronization across the SMs of the GPU in order to complete the communication between SMs. Although our previous work demonstrated that implementing barrier synchronization on the GPU itself can significantly improve performance and deliver correct results in critical bioinformatics applications, guaranteeing the correctness of inter-SM communication is only possible if a memory consistency model is assumed. To address this problem, NVIDIA recently introduced the _threadfence() function in CUDA 2.2, a function that can guarantee the correctness of GPU-based inter-SM communication. However, this function currently introduces so much overhead that when using it in (direct) GPU synchronization, GPU synchronization actually performs worse than indirect synchronization via the CPU, thus raising the question of whether “to GPU synchronize or not GPU synchronize?”
Wu-chun Feng, Shucai Xiao
ISCAS1
2010 Broadening accessibility to computer science for K-12 education
abstract
Enrollments in computer science and computer engineering have decreased dramatically since the dot-com bubble burst in 2000 even though it is projected that nearly three quarters of all science and engineering jobs in the future will be in these fields. Meeting this demand will require a substantial effort to inspire and motivate students as early as in the elementary school years. The challenge is to provide motivational access to computer science training, particularly for females and minorities in disadvantaged areas.
Mark K. Gardner, Wu-chun Feng
ITiCSE2
2010 Missing genes in the annotation of prokaryotic genomes
abstract
BACKGROUND: Protein-coding gene detection in prokaryotic genomes is considered a much simpler problem than in intron-containing eukaryotic genomes. However there have been reports that prokaryotic gene finder programs have problems with small genes (either over-predicting or under-predicting). Therefore the question arises as to whether current genome annotations have systematically missing, small genes. RESULTS: We have developed a high-performance computing methodology to investigate this problem. In this methodology we compare all ORFs larger than or equal to 33 aa from all fully-sequenced prokaryotic replicons. Based on that comparison, and using conservative criteria requiring a minimum taxonomic diversity between conserved ORFs in different genomes, we have discovered 1,153 candidate genes that are missing from current genome annotations. These missing genes are similar only to each other and do not have any strong similarity to gene sequences in public databases, with the implication that these ORFs belong to missing gene families. We also uncovered 38,895 intergenic ORFs, readily identified as putative genes by similarity to currently annotated genes (we call these absent annotations). The vast majority of the missing genes found are small (less than 100 aa). A comparison of select examples with GeneMark, EasyGene and Glimmer predictions yields evidence that some of these genes are escaping detection by these programs. CONCLUSIONS: Prokaryotic gene finders and prokaryotic genome annotations require improvement for accurate prediction of small genes. The number of missing gene families found is likely a lower bound on the actual number, due to the conservative criteria used to determine whether an ORF corresponds to a real gene.
Andrew S. Warren, Jeremy S. Archuleta, Wu-chun Feng, João Carlos Setubal
BMC Bioinform.3
2010 Global-scale distributed I/O with ParaMEDIC
abstract
Abstract Achieving high performance for distributed I/O on a wide‐area network continues to be an elusive holy grail. Despite enhancements in network hardware as well as software stacks, achieving high‐performance remains a challenge. In this paper, our worldwide team took a completely new and non‐traditional approach to distributed I/O, calledParaMEDIC: Parallel Metadata Environment for Distributed I/O and Computing, by utilizing application‐specifictransformationof data to orders of magnitude smaller metadata before performing the actual I/O. Specifically, this paper details our experiences in deploying a large‐scale system to facilitate the discovery of missing genes and constructing a genome similarity tree by encapsulating the mpiBLAST sequence‐search algorithm into ParaMEDIC. The overall project involved nine computational sites spread across the U.S. and generated more than a petabyte of data that was ‘teleported’ to a large‐scale facility in Tokyo for storage. Copyright © 2010 John Wiley & Sons, Ltd.
Pavan Balaji, Wu-chun Feng, Heshan Lin, Jeremy S. Archuleta, Satoshi Matsuoka, Andrew S. Warren, João Carlos Setubal, Ewing L. Lusk, Rajeev Thakur, Ian T. Foster, Daniel S. Katz, Shantenu Jha, K. Shinpaugh, Susan Coghlan, Daniel A. Reed
Concurr. Comput. Pract. Exp.2
2009 On the Robust Mapping of Dynamic Programming onto a Graphics Processing Unit
abstract
Graphics processing units (GPUs) have been widely used to accelerate algorithms that exhibit massive data parallelism or task parallelism. When such parallelism is not inherent in an algorithm, computational scientists resort to simply replicating the algorithm on every multiprocessor of a NVIDIA GPU, for example, to create such parallelism, resulting in embarrassingly parallel ensemble runs that deliver significant aggregate speed-up. However, the fundamental issue with such ensemble runs is that the problem size to achieve this speed-up is limited to the available shared memory and cache of a GPU multiprocessor. An example of the above is dynamic programming (DP), one of the Berkeley 13 dwarfs. All known DP implementations to date use the coarse-grained approach of embarrassingly parallel ensemble runs because a fine-grained parallelization on the GPU would require extensive communication between the multiprocessors of a GPU, which could easily cripple performance as communication between multiprocessors is not natively supported in a GPU. Consequently, we address the above by proposing a fine-grained parallelization of a single instance of the DP algorithm that is mapped to the GPU. Our parallelization incorporates a set of techniques aimed to substantially improve GPU performance: matrix re-alignment, coalesced memory access, tiling, and GPU (rather than CPU) synchronization. The specific DP algorithm that we parallelize is called Smith-Waterman (SWat), which is an optimal local-sequence alignment algorithm. We then use this SWat algorithm as a baseline to compare our GPU implementation, i.e., CUDA-SWat, to our implementation on the cell broadband engine, i.e., Cell-SWat.
Shucai Xiao, Ashwin M. Aji, Wu-chun Feng
ICPADS3
2009 GePSeA: A General-Purpose Software Acceleration Framework for Lightweight Task Offloading
abstract
Hardware-acceleration techniques continue to be used to speed-up the execution of scientific codes. To do so, software developers identify portions of these codes that are amenable for offloading and map them to hardware accelerators. However, offloading such tasks to specialized hardware accelerators is non-trivial. Furthermore, these accelerators can add significant cost to a computing system. Consequently, we propose a framework called GePSeA (General Purpose Software Acceleration Framework), which uses a small fraction of the computational power on multi-core architectures to ``onload'' complex application-specific tasks. Specifically, GePSeA provides a lightweight process that acts as a helper agent to the application by executing application-specific tasks asynchronously and efficiently. We then apply the GePSeA framework to a real application, namely, an open-source computational biology application, and demonstrate significant application-level benefits.
Pavan Balaji, Wu-chun Feng
ICPP3
2009 Multi-dimensional characterization of temporal data mining on graphics processors
abstract
Through the algorithmic design patterns of data parallelism and task parallelism, the graphics processing unit (GPU) offers the potential to vastly accelerate discovery and innovation across a multitude of disciplines. For example, the exponential growth in data volume now presents an obstacle for high-throughput data mining in fields such as neuroscience and bioinformatics. As such, we present a characterization of a MapReduced-based data-mining application on a general-purpose GPU (GPGPU). Using neuroscience as the application vehicle, the results of our multi-dimensional performance evaluation show that a ldquoone-size-fits-allrdquo approach maps poorly across different GPGPU cards. Rather, a high-performance implementation on the GPGPU should factor in the 1) problem size, 2) type of GPU, 3) type of algorithm, and 4) data-access method when determining the type and level of parallelism. To guide the GPGPU programmer towards optimal performance within such a broad design space, we provide eight general performance characterizations of our data-mining application.
Jeremy S. Archuleta, Yong Cao 0003, Thomas Scogland, Wu-chun Feng
IPDPS4
2009 The Green500 List: Year one
abstract
The latest release of the Green500 List in November 2008 marked its one-year anniversary. As such, this paper aims to provide an analysis and retrospective examination of the Green500 List in order to understand how the list has evolved and what trends have emerged. In addition, we present community feedback on the Green500 List, particularly from two Green500 birds-of-a-feather (BoF) sessions at the International Supercomputing Conference in June 2008 and SC|08 in November 2008, respectively.
Wu-chun Feng, Thomas Scogland
IPDPS1
2009 On the energy efficiency of graphics processing units for scientific computing
abstract
The graphics processing unit (GPU) has emerged as a computational accelerator that dramatically reduces the time to discovery in high-end computing (HEC). However, while today's state-of-the-art GPU can easily reduce the execution time of a parallel code by many orders of magnitude, it arguably comes at the expense of significant power and energy consumption. For example, the NVIDIA GTX 280 video card is rated at 236 watts, which is as much as the rest of a compute node, thus requiring a 500-W power supply. As a consequence, the GPU has been viewed as a ldquonon-greenrdquo computing solution. This paper seeks to characterize, and perhaps debunk, the notion of a ldquopower-hungry GPUrdquo via an empirical study of the performance, power, and energy characteristics of GPUs for scientific computing. Specifically, we take an important biological code that runs in a traditional CPU environment and transform and map it to a hybrid CPU+GPU environment. The end result is that our hybrid CPU+GPU environment, hereafter referred to simply as GPU environment, delivers an energy-delay product that is multiple orders of magnitude better than a traditional CPU environment, whether unicore or multicore.
Shucai Xiao, Wu-chun Feng
IPDPS3
2008 Optimizing performance, cost, and sensitivity in pairwise sequence search on a cluster of PlayStations
abstract
The Smith-Waterman algorithm is a dynamic programming method for determining optimal local alignments between nucleotide or protein sequences. However, it suffers from quadratic time and space complexity. As a result, many algorithmic and architectural enhancements have been proposed to solve this problem, but at the cost of reduced sensitivity in the algorithms or significant expense in hardware, respectively.
Ashwin M. Aji, Wu-chun Feng
BIBE2
2008 Achieving Edge-Based Fairness in a Multi-Hop Environment
abstract
We propose efficient buffer-accounting algorithms that achieve edge-based max-min and proportional fairness in a multi-hop (MH), multi-bottleneck network environment by extending and generalizing an existing proactive queue-management scheme called GREEN. We call our scheme GREEN-MH. We envision deploying GREEN-MH at an institutional gateway in the context of a larger multi-hop and multi-bottleneck network environment. GREEN-MH uses a dynamic buffer-accounting algorithm on a per-flow basis such that certain edge-based fairness policies (e.g., max-min and proportional) are enforced among the competing TCP flows.
Mustafa Arisoylu, Wu-chun Feng
CCNC2
2008 Making a Case for Proactive Flow Control in Optical Circuit-Switched Networks
Mithilesh Kumar 0002, Vineeta Chaube, Pavan Balaji, Wu-chun Feng, Hyun-Wook Jin
HiPC4
2008 Semantic-based distributed i/o with the paramedic framework
abstract
Many large-scale applications simultaneously rely on multiple resources for efficient execution. For example, such applications may require both large compute and storage resources; however, very few supercomputing centers can provide large quantities of both. Thus, data generated at the compute site oftentimes has to be moved to a remote storage site for either storage or visualization and analysis. Clearly, this is not an efficient model, especially when the two sites are distributed over a wide-area network.
Pavan Balaji, Wu-chun Feng, Heshan Lin
HPDC2
2008 Impact of Network Sharing in Multi-Core Architectures
abstract
As commodity components continue to dominate the realm of high-end computing, two hardware trends have emerged as major contributors-high-speed networking technologies and multi-core architectures. Communication middleware such as the Message Passing Interface (MPI) uses the network technology for communicating between processes that reside on different physical nodes, while using shared memory for communicating between processes on different cores within the same node. Thus, two conflicting possibilities arise: (i) with the advent of multi-core architectures, the number of processes that reside on the same physical node and hence share the same physical network can potentially increase significantly, resulting in increased network usage, and (ii) given the increase in intra-node shared-memory communication for processes residing on the same node, the network usage can potentially decrease significantly. In this paper, we address these two conflicting possibilities and study the behavior of network usage in multi-core environments with sample scientific applications. Specifically, we analyze trends that result in increase or decrease of network usage, and we derive insights into application performance based on these. We also study the sharing of different resources in the system in multi-core environments and identify the contribution of the network in this mix. In addition, we study different process allocation strategies and analyze their impact on such network sharing.
Ganesh Narayanaswamy, Pavan Balaji, Wu-chun Feng
ICCCN3
2008 Semantics-based distributed I/O for mpiBLAST
abstract
BLAST is a widely used software toolkit for genomic sequence search. mpiBLAST is a freely available, open-source parallelization of BLAST that uses database segmentation to allow different worker processes to search (in parallel) unique segments of the database. After searching, the workers write their output to a filesystem. While mpiBLAST has been shown to achieve high performance in clusters with fast local filesystems, its I/O processing remains a concern for scalability, especially in systems having limited I/O capabilities such as distributed filesystems spread across a wide-area network. Thus, we present ParaMEDIC---a novel environment that uses application-specific semantic information to compress I/O data and improve performance in distributed environments. Specifically, for mpiBLAST, ParaMEDIC partitions worker processes into compute and I/O workers. Compute workers, instead of directly writing the output to the filesystem, the workers process the output using semantic knowledge about the application to generate metadata and write the metadata to the filesystem. I/O workers, which physically reside closer to the actual storage, then process this metadata to re-create the actual output and write it to the filesystem. This approach allows ParaMEDIC to reduce I/O time, thus accelerating mpiBLAST by as much as 25-fold.
Pavan Balaji, Wu-chun Feng, Jeremy S. Archuleta, Heshan Lin, Rajkumar Kettimuthu, Rajeev Thakur, Xiaosong Ma
PPoPP2
2008 Massively parallel genomic sequence search on the Blue Gene/P architecture
abstract
This paper presents our first experiences in mapping and optimizing genomic sequence search onto the massively parallel IBM Blue Gene/P (BG/P) platform. Specifically, we performed our work on mpiBLAST, a parallel sequence-search code that has been optimized on numerous supercomputing environments. In doing so, we identify several critical performance issues. Consequently, we propose and study different approaches for mapping sequence-search and parallel I/O tasks on such massively parallel architectures.We demonstrate that our optimizations can deliver nearly linear scaling (93% efficiency) on up to 32,768 cores of BG/P. In addition, we show that such scalability enables us to complete a large-scale bioinformatics problem - sequence searching a microbial genome database against itself to support the discovery of missing genes in genomes - in only a few hours on BG/P. Previously, this problem was viewed as computationally intractable in practice.
Heshan Lin, Pavan Balaji, Ruth Poole, Carlos P. Sosa, Xiaosong Ma, Wu-chun Feng
SC6
2008 Asymmetric interactions in symmetric multi-core systems: analysis, enhancements and evaluation
abstract
Multi-core architectures have spurred the rapid growth in high-end computing systems. While the vast majority of such multi-core processors contain symmetric hardware components, their interaction with systems software, in particular the communication stack, results in a remarkable amount of asymmetry in the effective capability of the different cores. In this paper, we analyze such interactions and propose a novel management library called SyMMer (Systems Mapping Manager) that monitors these interactions and dynamically manages the mapping of processes on processor cores to transparently improve application performance. Together with a detailed description of the SyMMer library, we also present performance evaluation comparing SyMMer to a vanilla communication library using various micro-benchmarks as well as popular applications and scientific libraries. Experimental results demonstrate more than a two-fold improvement in communication time and 10-15% improvement in overall application performance.
Thomas Scogland, Pavan Balaji, Wu-chun Feng, Ganesh Narayanaswamy
SC3
2008 Algorithms for Integrated Routing and Scheduling for Aggregating Data from Distributed Resources on a Lambda Grid
abstract
In many e-science applications, there exists an important need to aggregate information from data repositories distributed around the world. In an effort to better link these resources in a unified manner, many lambda-grid networks, which provide end-to-end dedicated optical-circuit-switched connections, have been investigated. In this context, we consider the problem of aggregating files from distributed databases at a (grid) computing node over a lambda grid. The challenge is (1) to identify routes (that is, circuits) in the lambda-grid network, along which files should be transmitted, and (2) to schedule the transfers of these files over their respective circuits. To address this challenge, we propose a hybrid approach that combines offline and online scheduling. We define the Time-Path Scheduling Problem (TPSP) for offline scheduling. We prove that TPSP is NP-complete, develop a Mixed Integer Linear Program (MILP) formulation for TPSP, and then propose a greedy approach to solve TPSP because the MILP does not scale well. We compare the performance of the greedy approach on a few representative lambda-grid network topologies. One key input to the offline schedule is the file transfer time. Due to dynamics at the receiving end host, which is hard to model precisely, the actual file transfer time may vary. We first propose a model for estimating the file transfer time. Then, we propose online reconfiguration algorithms so that as files are transferred, the offline schedule may be modified online, depending on the amount of time that it actually took to transfer the file. This helps in reducing the total time to transfer all the files, which is an important metric. To demonstrate the effectiveness of our approach, we present results on an emulated lambda-grid network testbed.
Amitabha Banerjee, Wu-chun Feng, Dipak Ghosal, Biswanath Mukherjee
IEEE Trans. Parallel Distributed Syst.2
2007 CPU MISER: A Performance-Directed, Run-Time System for Power-Aware Clusters
abstract
Performance and power are critical design constraints in today's high-end computing systems. Reducing power consumption without impacting system performance is a challenge for the HPC community. We present a runtime system (CPU MISER) and an integrated performance model for performance-directed, power-aware cluster computing. CPU MISER supports system-wide, application-independent, fine-grain, dynamic voltage and frequency scaling (DVFS) based power management for a generic power-aware cluster. Experimental results show that CPU MISER can achieve as much as 20% energy savings for the NAS parallel benchmarks. In addition to energy savings, CPU MISER is able to constrain performance loss for most applications within user-specified limits. These constraints are achieved through accurate performance modeling and prediction, coupled with advanced control techniques.
Rong Ge 0002, Xizhou Feng, Wu-chun Feng, Kirk W. Cameron
ICPP3
2007 A Maintainable Software Architecture for Fast and Modular Bioinformatics Sequence Search
abstract
Bioinformaticists use the Basic Local Alignment Search Tool (BLAST) to characterize an unknown sequence by comparing it against a database of known sequences, thus detecting evolutionary relationships and biological properties. mpiBLAST is a widely-used, high-performance, open-source parallelization of BLAST that runs on a computer cluster delivering super-linear speedups. However, the Achilles heel of mpiBLAST is its lack of modularity, thus adversely affecting maintainability and extensibility. Alleviating this shortcoming requires an architectural refactoring to improve maintenance and extensibility while preserving high performance. Toward that end, this paper evaluates five different software architectures and details how each satisfies our design objectives. In addition, we introduce a novel approach to using mixin layers to enable mixing-and-matching of modules in constructing sequence-search applications for a variety of high-performance computing systems. Our design, which we call "mixin layers with refined roles", utilizes mixin layers to separate functionality into complementary modules and the refined roles in each layer improve the inherently modular design by precipitating flexible and structured parallel development, a necessity for an open-source application. We believe that this new software architecture for mpiBLAST-2.0 will benefit both the users and developers of the package and that our evaluation of different software architectures will be of value to other software engineers faced with the challenges of creating maintainable and extensible, high-performance, bioinformatics software.
Jeremy S. Archuleta, Eli Tilevich, Wu-chun Feng
ICSM3
2007 Green Supercomputing in a Desktop Box
abstract
The advent of the Beowulf cluster in 1994 provided dedicated compute cycles, i.e., supercomputing for the masses, as a cost-effective alternative to large supercomputers, i.e., supercomputing for the few. However as the cluster movement matured, these clusters became like their large-scale supercomputing brethren - a shared (and power-hungry) datacenter resource that must reside in a actively-cooled machine room in order to operate properly. The above observation, coupled with the increasing performance gap between the PC and supercomputer, provides the motivation for a "green supercomputer" in a desktop box. Thus, this paper presents and evaluates such an architectural solution: a 12-node personal desktop supercomputer that offers an interactive environment for developing parallel codes and achieves 14 Gflops on Linpack but sips only 185 watts of power at load - all this in the approximate form factor of a Sun SPARCstation 1 pizza box.
Wu-chun Feng, Avery Ching, Chung-Hsing Hsu
IPDPS1
2007 Analyzing the impact of supporting out-of-order communication on in-order performance with iWARP
abstract
Due to the growing need to tolerate network faults and congestion in high-end computing systems, supporting multiple network communication paths is becoming increasingly important. However, multi-path communication comes with the disadvantage of out-of-order arrival of packets (because packets may traverse different paths). While modern networking stacks such as the Internet Wide-Area RDMA Protocol (iWARP) over 10-Gigabit Ethernet (10GE) support multi-path communication, their current implementations do not handle out-of-order packets primarily owing to the overhead on in-order communication that it adds. Specifically, in iWARP, supporting out-of-order packets requires every packet to carry additional information causing significant overhead on packets that arrive in-order. Thus, in this paper, we analyze the trade-offs in designing a feature-complete iWARP stack, i.e., one that provides support for out-of-order arriving packets, and thus, multi-path systems, while focusing on the performance of in-order communication. We propose three feature-complete designs of iWARP and analyze the pros and cons of each of these designs using performance experiments based on several micro-benchmarks as well as an iso-surface visual rendering application. Our analysis reveals that the iWARP design providing the best overall performance depends on the particular characteristics of the upper layers and that different designs are optimal based on the metric of interest.
Pavan Balaji, Wu-chun Feng, Sitha Bhagvat, Dhabaleswar K. Panda 0001, Rajeev Thakur, William Gropp
SC2
2007 High-performance computing using accelerators
Wu-chun Feng, Dinesh Manocha
Parallel Comput.1
2006 A Feedback Mechanism for Network Scheduling in LambdaGrids
abstract
Next-generation e-Science applications will require the ability to transfer information at high data rates between distributed computing centers and data repositories. A Lambda-Grid offers dedicated, optical, circuit-switched, point-to-point connections, which may be reserved exclusively for an application. Though such dedicated high-speed connections eliminate congestion in the network, they effectively push the network congestion out to the end systems, as processing speeds have not kept up with networking speeds. Therefore, developing an efficient transport protocol over such highspeed dedicated circuits is of critical importance. In this work, we propose the idea of a lightweight end-system protocol, based on performance monitoring, to significantly improve the performance of data transport over a LambdaGrid. In particular, we focus on dynamically monitoring the OS task scheduling at the receiving end-system so that potential end-system congestion may be detected early and appropriate feedback can be transmitted back to the sending end-system to avoid packet losses. One example of such an evasive action is to suspend transmission for certain duration of time during which the OS on the receiving end-system must handle other computational processes. With this in mind, we propose to extend the Reliable-Blast UDP (RBUDP) protocol to take such evasive action by using a simple feedback mechanism that is activated via performance monitoring. The new protocol, named RBUDP dramatically improves the performance of data transfer over LambdaGrids. We demonstrate the effectiveness of our proposed protocol and illustrate the performance gains achieved via network emulation.
Pallab Datta, Sushant Sharma, Wu-chun Feng
CCGRID3
2006 Exploring I/O Strategies for Parallel Sequence-Search Tools with S3aSim
abstract
Parallel sequence-search tools are rising in popularity among computational biologists. With the rapid growth of sequence databases, database segmentation is the trend of the future for such search tools. While I/O currently is not a significant bottleneck for parallel sequence-search tools, future technologies including faster processors, customized computational hardware such as FPGAs, improved search algorithms, and exponentially growing databases emphasize an increasing need for efficient parallel I/O in future parallel sequence-search tools. Our paper focuses on examining different I/O strategies for these future tools in a modern parallel file system (PVFS2). Because implementing and comparing various I/O algorithms in every search tool is labor-intensive and time-consuming, we introduce S3aSim, a general simulation framework for sequence-search which allows us to quickly implement, test, and profile various I/O strategies. We examine a variety of I/O strategies (e.g., master-writing and various worker-writing strategies using individual and collective I/O methods) for storing result data in sequence-search tools such as mpiBLAST, pioBLAST, and parallel HMMer. Our experiments fully detail the interaction of computing and I/O within a full application simulation as opposed to typical I/O-only benchmarks
Avery Ching, Wu-chun Feng, Heshan Lin, Xiaosong Ma, Alok N. Choudhary
HPDC2
2006 When Optical Networking Meets Grid Computing?
abstract
High-speed optical networking infrastructures such as National LambdaRail promise to revolutionize the way that we approach grid computing. With optical-fiber speeds doubling every 9 months and computing speeds doubling only every 18-24 months over the past few decades, we have finally reached a crossroads where network speeds have outstripped the ability of processors to keep up. Thus, we may be entering a new world where the central architectural element in grid computing is the optical network, not the end-host computer(s). The goal of this panel is to discuss whether such a new world is emerging, and if so, what the research challenges will be.
Wu-chun Feng, Mark K. Gardner, Gigi Karmous-Edwards, Jerry Sobieski, Malathi Veeraraghavan
ICCCN1
2006 RAPID: an end-system aware protocol for intelligent data transfer over lambda grids
abstract
Next-generation e-science applications will require the ability to transfer information at high data rates between distributed computing centers and data repositories. To support such applications, lambda grid networks have been built to provide large, on-demand bandwidth between end-points that are interconnected via optical circuit-switched lambdas. It is extremely important to develop an efficient transport protocol over such high-capacity, dedicated circuits. Because lambdas provide dedicated bandwidth between endpoints, they obviate the need for network congestion control. Consequently, past research has demonstrated that rate-based transport protocols, such as RBUDP, are more effective than TCP in transferring data over lambdas. However, while lambdas eliminate congestion in the network, they ultimately push the congestion to the endpoints - congestion that current rate-based transport protocols are ill-suited to handle. In this paper we introduce a "rate-adaptive protocol for intelligent delivery (RAPID)" of data that is lightweight and end-system performance-aware, so as to maximize end-to-end throughput while minimizing packet loss. Based on self monitoring of the dynamic task-priority at the receiving end-system, our protocol enables the receiver to proactively deliver feedback to the sender, so that the sender may adapt its sending rate to avoid congestion at the receiving end-system. This avoids large bursts of packet losses typically observed in current rate-based transport protocols. Over a 10-Gigabit link emulation of an optical circuit, RAPID reduces file-transfer time, and hence improves end-to-end throughput by as much as 25%.
Amitabha Banerjee, Wu-chun Feng, Biswanath Mukherjee, Dipak Ghosal
IPDPS2
2006 Making a case for a Green500 list
abstract
For decades now, the notion of "performance" has been synonymous with "speed" (as measured in FLOPS, short for floating-point operations per second). Unfortunately, this particular focus has led to the emergence of supercomputers that consume egregious amounts of electrical power and produce so much heat that extravagant cooling facilities must be constructed to ensure proper operation. In addition, the emphasis on speed as the performance metric has caused other performance metrics to be largely ignored, e.g., reliability, availability, and usability. As a consequence, all of the above has led to an extraordinary increase in the total cost of ownership (TCO) of a supercomputer. Despite the importance of the TOP500 List, we argue that the list makes it much more difficult for the high-performance computing (UPC) community to focus on performance metrics other than speed. Therefore, to raise awareness to other performance metrics of interest, e.g., energy efficiency for improved reliability, we propose a Green500 List and discuss the potential metrics that would be used to rank supercomputing systems on such a list.
Sushant Sharma, Chung-Hsing Hsu, Wu-chun Feng
IPDPS3
2006 Grid networks and portals - End-system aware, rate-adaptive protocol for network transport in LambdaGrid environments
abstract
Next-generation e-Science applications will require the ability to transfer information at high data rates between distributed computing centers and data repositories. A LambdaGrid offers dedicated, optical, circuit-switched, point-to-point connections that can be reserved exclusively for such applications. These dedicated high-speed connections eliminate network congestion as seen in traditional Internet, but they effectively push the network congestion to the end systems, as processing speeds cannot keep up with networking speeds. Thus, developing an efficient transport protocol over such high-speed dedicated circuits is of critical importance.We propose the idea of a end-system aware, rate-adaptive protocol for network transport, based on end-system performance monitoring. Our proposed protocol significantly improves the performance of data transfer over LambdaGrids by intelligently adapting the sending rate based on end-system constraints. We demonstrate the effectiveness of our proposed protocol and illustrate the performance gains achieved via wide-area network emulation.
Pallab Datta, Wu-chun Feng, Sushant Sharma
SC2
2006 Grid applications - Parallel genomic sequence-searching on an ad-hoc grid: experiences, lessons learned, and implications
abstract
The Basic Local Alignment Search Tool (BLAST) allows bioinformaticists to characterize an unknown sequence by comparing it against a database of known sequences. The similarity between sequences enables biologists to detect evolutionary relationships and infer biological properties of the unknown sequence.mpiBLAST, our parallel BLAST, decreases the search time of a 300 KB query on the current NT database from over two full days to under 10 minutes on a 128-processor cluster and allows larger query files to be compared. Consequently, we propose to compare the largest query available, the entire NT database, against the largest database available, the entire NT database. The result of this comparison will provide critical information to the biology community, including insightful evolutionary, structural, and functional relationships between every sequence and family in the NT database.Preliminary projections indicated that to complete the above task in a reasonable length of time required more processors than were available to us at a single site. Hence, we assembled GreenGene, an ad-hoc grid that was constructed "on the fly" from donated computational, network, and storage resources during last year's SC|05. GreenGene consisted of 3048 processors from machines that were distributed across the United States. This paper presents a case study of mpiBLAST on GreenGene --- specifically, a pre-run characterization of the computation, the hardware and software architectural design, experimental results, and future directions.
Mark K. Gardner, Wu-chun Feng, Jeremy S. Archuleta, Heshan Lin, Xiaosong Ma
SC2
2005 Head-to-TOE Evaluation of High-Performance Sockets over Protocol Offload Engines
abstract
Despite 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
CLUSTER2
2005 A Feasibility Analysis of Power Awareness in Commodity-Based High-Performance Clusters
abstract
We present a feasibility study of a power-reduction scheme that reduces the thermal power of processors by lowering frequency and voltage in the context of high-performance computing. The study revolves around a 16-processor Opteron-based Beowulf cluster, configured as four nodes of quad-processors, and shows that one can easily reduce a significant amount of CPU and system power dissipation and its associated energy costs while still maintaining high performance. Specifically, our study shows that a 5% performance slowdown can be traded off for an average of 19% system energy savings and 24% system power reduction. These preliminary empirical results, via real measurements, are encouraging because hardware failures often occur when the cluster is running hot, i.e, when the workload is heavy, and the new power-reduction scheme can effectively reduce a cluster's power demands during these busy periods
Chung-Hsing Hsu, Wu-chun Feng
CLUSTER2
2005 Q-Composer and CpR: a probabilistic synthesizer and regulator of traffic (a probabilistic control of buffer occupancy)
abstract
We present and show the correctness of two algorithms called Q-Composer and CpR. Q-Composer is a probabilistic traffic-synthesizer and CpR is a probabilistic traffic-regulator. Given a cumulative distribution function F, Q-Composer synthesizes a flow that when fed into a single-input single-output network element, the distribution of the queue-size probability at the element closely follows F. CpR regulates an arbitrary traffic so that when the regulated traffic (i.e. the output of CpR) is fed into a single-input single-output network element, the distribution of the queue-size probability at the element closely follows a pre-specified cdf F. CpR can be viewed as a probabilistic generalization of deterministic Leaky-bucket regulators. Q-Composer and CpR are straightforward algorithms to implement and have applications in providing end-to-end probabilistic quality-of-service guarantees, multimedia encoding/decoding, and resource allocation and in simulation studies, beside other areas.
Sami Ayyorgun, Sarut Vanichpun, Wu-chun Feng
INFOCOM3
2005 A Power-Aware Run-Time System for High-Performance Computing
abstract
For decades, the high-performance computing (HPC) community has focused on performance, where performance is defined as speed. To achieve better performance per compute node, microprocessor vendors have not only doubled the number of transistors (and speed) every 18-24 months, but they have also doubled the power densities. Consequently, keeping a large-scale HPC system functioning properly requires continual cooling in a largemachine room, thus resulting in substantial operational costs. Furthermore, the increase in power densities has led (in part) to a decrease in system reliability, thus leading to lost productivity. To address these problems, we propose a power-aware algorithm that automatically and transparently adapts its voltage and frequency settings to achieve significant power reduction and energy savings with minimal impact on performance. Specifically, we leverage a commodity technology called "dynamic voltage and frequency scaling" to implement our power-aware algorithm in the run-time system of commodity HPC systems.
Chung-Hsing Hsu, Wu-chun Feng
SC2
2005 Analyzing MPI performance over 10-Gigabit ethernet
Justin Gus Hurwitz, Wu-chun Feng
J. Parallel Distributed Comput.2
2005 Anatomy of UDP and M-VIA for cluster communication
Laxmi N. Bhuyan, Wu-chun Feng
J. Parallel Distributed Comput.3
2004 A Multimodal Interface for the Immediate Transcription of Radiology Dictation
abstract
We present the design and implementation of an integrated multimodal interface that delivers instant turnaround on transcribing a radiology dictation. This instant turnaround time virtually eliminates a hospital's liability with respect to improper transcriptions of oral dictations and all but eliminates the need for transcribers. The multimodal interface seamlessly integrates three modes of input - speech, handwriting, and written gestures to provide an easy-to-use system for the radiologist.
Wu-chun Feng
CBMS1
2004 A Systematic Approach for Providing End-to-End Probabilistic QoS Guarantees
abstract
We propose a probabilistic characterization of network traffic. This characterization can handle traffic with heavy-tailed distributions in performance analysis. We show that queue size, output traffic, virtual delay, aggregate traffic, etc. at various points in a network can easily be characterized within the framework. This characterization is measurable and allows for a simple probabilistic method for regulating network traffic. All of these properties of the proposed characterization enable a systematic approach for providing end-to-end probabilistic QoS guarantees.
Sami Ayyorgun, Wu-chun Feng
ICCCN2
2004 Re-Architecting Flow Control Adaptation for Grid Environments
abstract
Summary form only given. The performance of TCP in wide-area networks (WANs) is becoming increasingly important with the deployment of computational and data grids. In WAN environments, TCP does not provide good performance for data-intensive applications without the tuning of flow-control buffer sizes. Manual adjustment of buffer sizes is tedious even for network experts. For application scientists, tuning is often an impediment to getting work done. Thus, buffer tuning should be automated. Existing techniques for automatic buffer tuning only measure the bandwidth-delay product (BDP) during connection establishment. This ignores the large fluctuation of the BDP over the lifetime of the connection. In contrast, the dynamic right-sizing algorithm dynamically changes buffer sizes in response to changing network conditions. We describe a new user-space implementation of dynamic right-sizing in FTP (drsFTP) that supports third-party data transfers, a mainstay of scientific computing. In addition to comparing the performance of the new implementation with the old in a WAN-emulated environment, we give performance results over a live WAN. In this particular WAN environment, the new implementation produces transfer rates of up to five times higher than untuned FTP.
Adam Engelhart, Mark K. Gardner, Wu-chun Feng
IPDPS3
2004 User-space auto-tuning for TCP flow control in computational grids
Mark K. Gardner, Sunil Thulasidasan, Wu-chun Feng
Comput. Commun.3
2003 MAGNET: A Tool for Debugging, Analyzing and Adapting Computing Systems
abstract
As computing systems grow in complexity, the cluster and grid communities require more sophisticated tools to diagnose, debug and analyze such systems. We have developed a toolkit called MAGNET (Monitoring Apparatus for General kerNel-Event Tracing) that provides a detailed look at operating-system kernel events with very low overhead. Using the fine-grained information that MAGNET exports from kernel space, challenging problems become amenable to identification and correction. In this paper, we first present the design, implementation and evaluation of MAGNET. Then, we show its use as a diagnostic tool, an online-monitoring tool and a tool for building adaptive applications in clusters and grids.
Mark K. Gardner, Wu-chun Feng, Michael Broxton, Adam Engelhart, Justin Gus Hurwitz
CCGRID2
2003 Optimizing GridFTP through Dynamic Right-Sizing
abstract
In this paper, we describe the integration of dynamic right-sizing - an automatic and scalable buffer management technique for enhancing TCP (transport control protocol) performance - into GridFTP, a subsystem of the Globus Toolkit for managing bulk data transfers across computational Grids. Such Grids are often characterized by networks with large bandwidth-delay products. Unfortunately, many of today's Grid applications use only a small fraction of available bandwidth because the default buffer sizes in TCP are tuned for yesterday's WAN (wide access network) speeds. Buffer sizes can be manually tuned to allow TCP flow control to adapt to high-speed WAN environments, but this is a tedious process. Although recent work has shown how to automatically tune system buffers during connection set-up, these values may not be appropriate for the connection's lifetime due to varying network delay and throughput. We show how using the technique of dynamic right-sizing (DRS) in GridFTP helps us optimize memory usage while maintaining high throughput over the lifetime of the connection. We also show how DRS enhances important GridFTP features such as striped and third-party data transfers in a scalable way. The technique is implemented entirely in user space so that end users do not have to modify the kernel.
Sunil Thulasidasan, Wu-chun Feng, Mark K. Gardner
HPDC2
2003 Optimizing 10-Gigabit Ethernet for Networks of Workstations, Clusters, and Grids: A Case Study
abstract
This paper presents a case study of the 10-Gigabit Ethernet (10GbE) adapter from Intel R . Specifically, with appropriate optimizations to the configurations of the 10GbE adapter and TCP, we demonstrate that the 10GbE adapter can perform well in local-area, storage-area, system-area, and wide-area networks. For local-area, storage-area, and system-area networks in support of networks of workstations, network-attached storage, and clusters, respectively, we can achieve over 7-Gb/s end-to-end throughput and 12-µs end-to-end latency between applications running on Linux-based PCs. For the wide-area network in support of grids, we broke the recently-set Internet2 Land Speed Record by 2.5 times by sustaining an end-to-end TCP/IP throughput of 2.38 Gb/s between Sunnyvale, California and Geneva, Switzerland (i.e., 10,037 kilometers) to move over a terabyte of data in less than an hour. Thus, the above results indicate that 10GbE may be a cost-effective solution across a multitude of computing environments.
Wu-chun Feng, Justin Gus Hurwitz, Harvey B. Newman, Sylvain Ravot, Roger Les Cottrell, Olivier Martin 0003, Fabrizio Coccetti, Cheng Jin 0009, David X. Wei, Steven H. Low
SC1
2003 Automatic Flow-Control Adaptation for Enhancing Network Performance in Computational Grids
Wu-chun Feng, Mark K. Gardner, Mike Fisk, Eric Weigle
J. Grid Comput.1
2003 Scheduling and Transport for File Transfers on High-Speed Optical Circuits
Malathi Veeraraghavan, Wu-chun Feng, Edwin K. P. Chong
J. Grid Comput.3
2002 The Bladed Beowulf: A Cost-Effective Alternative to Traditional Beowulf
abstract
We present a new twist to the Beowulf cluster - the Bladed Beowulf. In contrast to traditional Beowulfs which typically use Intel or AMD processors, our Bladed Beowulf uses Trans-meta processors in order to keep thermal power dissipation low and reliability and density high while still achieving comparable performance to Intel- and AMD-based clusters. Given the ever increasing complexity of traditional supercomputers and Beowulf clusters; the issues of size, reliability power consumption, and ease of administration and use will be "the" issues of this decade for high-performance computing. Bigger and faster machines are simply not good enough anymore. To illustrate, we present the results of performance benchmarks on our Bladed Beowulf and introduce two performance metrics that contribute to the total cost of ownership (TCO) of a computing system - performance/power and performance/space.
Wu-chun Feng, Michael S. Warren, Eric Weigle
CLUSTER1
2002 GREEN: proactive queue management over a best-effort network
abstract
We present a proactive queue-management (PQM) algorithm called GREEN (generalized random early evasion network) that applies knowledge of the steady-state behavior of TCP connections to drop packets intelligently and proactively, thus preventing congestion from ever occurring and ensuring a higher degree of fairness between flows. This congestion-prevention approach is in contrast to the congestion-avoidance approach of traditional active queue-management (AQM) schemes where congestion is actively detected early and then reacted to. In addition to enhancing fairness, GREEN keeps packet-queue lengths relatively low and reduces bandwidth and latency jitter. These characteristics are particularly beneficial to real-time multimedia applications. Further, GREEN achieves the above while maintaining high link utilization and low packet loss.
Wu-chun Feng, Apu Kapadia, Sunil Thulasidasan
GLOBECOM1
2002 Dynamic Right-Sizing in FTP (drsFTP): Enhancing Grid Performance in User-Space
abstract
With the advent of computational grids, networking performance over the wide-area network (WAN) has become a critical component in the grid infrastructure. Unfortunately, many high-performance grid applications only use a small fraction of the available bandwidth because operating systems and their associated protocol stacks are still tuned for yesterdays WAN speeds. As a result, network gurus undertake the tedious process of manually, tuning system buffers to allow TCP flow control to scale to today's WAN grid environments. Although recent research has shown how to set the size of these system buffers automatically at connection set-up, the buffer sizes are only appropriate at the beginning of the connection's lifetime. To address these problems, we describe an automated and scalable technique called dynamic right-sizing. We implement this technique in user space (in particular for bulk-data transfer) so that end users do not have to modify the kernel to achieve a significant increase in throughput.
Mark K. Gardner, Wu-chun Feng, Mike Fisk
HPDC2
2002 A Comparison of TCP Automatic Tuning Techniques for Distributed Computing
abstract
Rather than painful, manual, static, per-connection optimization of TCP buffer sizes simply to achieve acceptable performance for distributed applications, many researchers have proposed techniques to perform this tuning automatically. This paper first discusses the relative merits of the various approaches in theory, and then provides substantial experimental data concerning two competing implementations-the buffer autotuning already present in Linux 2.4.x and "dynamic right-sizing." The paper reveals heretofore unknown aspects of the problem and current solutions, provides insight into the proper approach for different circumstances, and points toward ways to further improve performance.
Eric Weigle, Wu-chun Feng
HPDC2
2002 On the transient behavior of TCP Vegas
abstract
Research has shown that TCP Vegas performs better than TCP Reno with respect to overall network utilization, stability, fairness, throughput, packet loss, and burstiness. We analyze and improve the transient behavior of TCP Vegas, an important issue in today's large "bandwidth-delay product" networks. To quantify of our analysis, we introduce a new metric that captures the transient performance of TCP, namely, the (normalized) convergence time. We then consider the slow-start mechanism in TCP Vegas and show that, with a properly configured /spl gamma/ parameter, the transient behavior of TCP Vegas improves with respect to convergence time.
Sarut Vanichpun, Wu-chun Feng
ICCCN2
2002 Honey, I Shrunk the Beowulf!
abstract
In this paper, we present a novel twist on the Beowulf cluster - the Bladed Beowulf. Designed by RLX Technologies and integrated and configured at Los Alamos National Laboratory, our Bladed Beowulf consists of compute nodes made from commodity off-the-shelf parts mounted on motherboard blades measuring 14.7" /spl times/ 4.7" /spl times/ 0.58". Each motherboard blade (node) contains a 633 MHz Trans-meta TM5600/spl trade/ CPU, 256 MB memory, 10 GB hard disk, and three 100-Mb/s Fast Ethernet network interfaces. Using a chassis provided by RLX, twenty-four such nodes mount side-by-side in a vertical orientation to fit in a rack-mountable 3U space, i.e., 19" in width and 5.25" in height. A Bladed Beowulf can reduce the total cost of ownership (TCO) of a traditional Beowulf by a factor of three while providing Beowulf-like performance. Accordingly, rather than use the traditional definition of price-performance ratio where price is the cost of acquisition, we introduce a new metric called ToPPeR: total price-performance ratio, where total price encompasses TCO. We also propose two related (but more concrete) metrics: performance-space ratio and performance-power ratio.
Wu-chun Feng, Michael S. Warren, Eric Weigle
ICPP1
2002 High-density computing: a 240-processor Beowulf in one cubic meter
abstract
We present results from computations on Green Destiny, a 240-processor Beowulf cluster which is contained entirely within a single 19-inch wide 42U rack. The cluster consists of 240 Transmeta TM5600 667-MHz CPUs mounted on RLX Technologies motherboard blades. The blades are mounted side-by-side in an RLX 3U rack-mount chassis, which holds 24 blades. The overall cluster contains 10 chassis and associated Fast and Gigabit Ethernet switches. The system has a footprint of 0.5 meter 2 (6 square feet), a volume of 0.85 meter 3 (30 cubic feet) and a measured power dissipation under load of 5200 watts (including network switches). We have measured the performance of the cluster using a gravitational treecode N-body simulation of galaxy formation using 200 million particles, which sustained an average of 38.9 Gflops on 212 nodes of the system. We also present results from a three-dimensional hydrodynamic simulation of a core-collapse supernova.
Michael S. Warren, Eric Weigle, Wu-chun Feng
SC3
2002 The MAGNeT Toolkit: Design, Implementation and Evaluation
Wu-chun Feng, Mark K. Gardner, Jeffrey R. Hay
J. Supercomput.1
2002 Packet Spacing: An Enabling Mechanism for Delivering Multimedia Content in Computational Grids
Annette C. Feng, Apu Kapadia, Wu-chun Feng, Geneva G. Belford
J. Supercomput.3
2001 A Case for TCP Vegas in High-Performance Computational Grids
abstract
Computational grids such as the Information Power Grid (Johnston et al., 1999), Particle Physics Data Grid , and Earth System Grid depend on TCP to provide reliable communication between nodes across a wide-area network (WAN). Of the available TCP implementations, TCP Reno and its variants are the most widely deployed; however, Reno's performance in computational grids is mediocre at best. Due to conflicting results in the evaluation of TCP implementations, we present a detailed simulation study that unifies the conflicting results and demonstrates the limitations of earlier work. We focus on the two most debated versions of TCP-Reno and Vegas. Using real traffic distributions, we show that Vegas performs well over modern high-performance links and better than Reno with the proper selection of the Vegas parameters /spl alpha/ and /spl beta/. Our results exhibit ways to significantly enhance the performance of distributed computational grids that rely on TCP.
Eric Weigle, Wu-chun Feng
HPDC2
2001 MAGNeT: monitor for application-generated network traffic
abstract
Over the last decade, network practitioners have focused on monitoring, measuring, and characterizing traffic in the network to gain insight into building critical network components (from the protocol stack to routers and switches to network interface cards). Previous research shows that additional insight can be obtained by monitoring traffic at the application level (i.e., before application-sent traffic is modulated by the protocol stack) rather than in the network (i.e., after it is modulated by the protocol stack). Consequently, this paper describes a monitor for application-generated network traffic (MAGNeT) that captures traffic generated by the application rather than traffic in the network. MAGNeT consists of application programs as well as modifications to the standard Linux kernel. Together, these tools provide the capability of monitoring an application's network behavior and protocol state information in production systems. The use of MAGNeT will enable the research community to construct a library of real traces of application-generated traffic from which researchers can more realistically test network protocol designs and theory. MAGNeT can also be used to verify the correct operation of protocol enhancements and to troubleshoot and tune protocol implementations.
Wu-chun Feng, Jeffrey R. Hay, Mark K. Gardner
ICCCN1
2001 Dynamic right-sizing: a simulation study
abstract
Virtually all network applications requiring reliable end-to-end communication depend on TCP. Unfortunately, the performance of any stock TCP is abysmal over wide-area networks (WANs) and even over local-area networks (LANs) with very high-bandwidth links. Currently, network researchers manually optimize TCP buffer sizes to achieve acceptable performance over a given connection. Unfortunately, this manual optimization requires changes to the kernel on both end hosts involved in the network connection (changes that are only effective for connections between these two hosts). Furthermore, because two administrative domains must be coordinated to perform this optimization, this process can be tedious and time consuming. To address these problems, this paper illustrates the benefits of a new technique called dynamic right-sizing. This technique dynamically and automatically determines the best buffer size, and hence flow-control window size in TCP. Our simulation study shows that dynamic right-sizing can improve the performance of flows by two orders of magnitude over stock TCP implementations that have static flow-control windows.
Eric Weigle, Wu-chun Feng
ICCCN2
2001 The Effects of Inter-packet Spacing on the Delivery of Multimedia Content
abstract
Streaming multimedia content with UDP has become increasingly popular over distributed systems such as the Internet. However, because UDP does not possess any congestion control mechanism and most best-effort traffic is served by the congestion-controlled TCP, UDP flows steal bandwidth from TCP to the point that TCP flows can starve for network resources. Furthermore, such applications may cause the Internet infrastructure to eventually suffer from congestion collapse because UDP traffic does not self-regulate itself. To address this problem, next-generation Internet routers will implement active queue management schemes to punish malicious traffic, e.g. non-adaptive UDP flows, and to the improve the performance of congestion-controlled traffic, e.g. TCP flows. The arrival of such routers will cripple the performance of today's UDP-based multimedia applications. So, in this paper, we introduce the notion of inter-packet spacing with control feedback to enable these UDP-based applications to perform well in the next-generation Internet while being adaptive and self-regulating. When compared with traditional UDP-based multimedia streaming, we illustrate that our counter-intuitive inter-packet spacing scheme with control feedback can reduce packet loss by 90% without adversely affecting the delivered throughput.
Apu Kapadia, Annette C. Feng, Wu-chun Feng
ICDCS3
2001 Performance Evaluation of the Quadrics Interconnection Network
abstract
In this paper we present an in-depth description of the Quadrics interconnection network (QsNET) and an experimental performance evaluation on a 64-node AlphaServer cluster. We explore several performance dimensions and scaling properties of the network by using a collection of benchmarks, based on different traffic patterns. Experiments with permutation patterns and uniform traffic are conducted to illustrate the basic characteristics of the interconnect under conditions commonly created by parallel scientific applications. Moreover, the behavior of the QsNET under I/O traffic, and the influence of the placement of the I/O servers are analyzed. The effects of using dedicated I/O nodes or shared I/O nodes are also exposed. In addition, we evaluate how background I/O traffic interferes with other parallel applications running concurrently. The experimental results indicate that the QsNET provides excellent performance in most cases, with excellent contention resolution mechanisms. Some important guidelines for applications and I/O servers mapping on large-scale clusters are also given.
Fabrizio Petrini, Adolfy Hoisie, Wu-chun Feng, Richard Graham
IPDPS3
2000 Scheduling with Global Information in Distributed Systems
abstract
Buffered coscheduling is a distributed scheduling methodology for time-sharing communicating processes in a distributed system, e.g., PC cluster. The principle mechanisms involved in this methodology are communication buffering and strobing. With communication buffering, communication generated by each processor is buffered and performed at the end of regular intervals (or time slices) to amortize communication and scheduling overhead. This regular communication structure is then leveraged by introducing a strobing mechanism which performs a total exchange of information at the end of each time slice. Thus, a distributed system can rely on this global information to more efficiently schedule communicating processes rather than rely on isolated or implicit information gathered from local events between processors. We describe how buffered coscheduling is implemented in the context of our SMART simulator. We then present performance measurements for two synthetic workloads and demonstrate the effectiveness of buffered coscheduling under different computational granularities, context-switch times and time-slice granularities.
Fabrizio Petrini, Wu-chun Feng
ICDCS2
2000 On the Burstiness of the TCP Congestion-Control Mechanism in a Distributed Computing System
abstract
Several studies in network traffic characterization have concluded that network traffic is self-similar and therefore not readily amenable to statistical multiplexing in a distributed computing system. This paper examines the effects of the TCP protocol stack on network traffic via an experimental study on the different implementations of TCP. We show that even when aggregate application traffic smooths out as more applications' traffic are multiplexed, TCP introduces burstiness into the aggregate traffic load, reducing network performance when statistical multiplexing is used within the network gateways.
Peerapol Tinnakornsrisuphap, Wu-chun Feng, Ian R. Philp
ICDCS2
2000 The Adverse Impact of the TCP Congestion-Control Mechanism in Heterogeneous Computing Systems
abstract
Via experimental study, we illustrate how TCP modulates application traffic in such a way as to adversely affect network performance in a heterogeneous computing system. Even when aggregate application traffic smooths out as more applications' traffic are multiplexed, TCP induces burstiness into the aggregate traffic load, and thus hurts network performance. This burstiness is particularly bad in TCP Reno, and even worse when RED gateways are employed. Based on the results of this experimental study, we then develop a stochastic model for TCP Reno to demonstrate how the burstiness in TCP Reno can be modeled.
Wu-chun Feng, Peerapol Tinnakornsrisuphap
ICPP1
2000 Buffered Coscheduling: A New Methodology for Multitasking Parallel Jobs on Distributed Systems
abstract
Buffered coscheduling is a scheduling methodology for time-sharing communicating processes in parallel and distributed systems. The methodology has two primary features: communication buffering and strobing. With communication buffering, communication generated by each processor is buffered and performed at the end of regular intervals to amortize communication and scheduling overhead. This infrastructure is then leveraged by a strobing mechanism to perform a total exchange of information at the end of each interval, thus providing global information to more efficiently schedule communicating processes. This paper describes how buffered coscheduling can optimize resource utilization by analyzing workloads with varying computational granularities, load imbalances, and communication patterns. The experimental results, performed using a detailed simulation model, show that buffered coscheduling is very effective on fast SANs such as Myrinet as well as slower switch-based LANs.
Fabrizio Petrini, Wu-chun Feng
IPDPS2
2000 Time-Sharing Parallel Jobs in the Presence of Multiple Resource Requirements
Fabrizio Petrini, Wu-chun Feng
JSSPP2
2000 The Failure of TCP in High-Performance Computational Grids
abstract
Distributed computational grids depend on TCP to ensure reliable end-to-end communication between nodes across the wide-area network (WAN). Unfortunately, TCP performance can be abysmal even when buffers on the end hosts are manually optimized. Recent studies blame the self-similar nature of aggregate network traffic for TCP’s poor performance because such traffic is not readily amenable to statistical multiplexing in the Internet, and hence computational grids. In this paper, we identify a source of self-similarity previously ignored, a source that is readily controllable - TCP. Via an experimental study, we examine the effects of the TCP stack on network traffic using different implementations of TCP. We show that even when aggregate application traffic ought to smooth out as more applications’ traffic are multiplexed, TCP induces burstiness into the aggregate traffic load, thus adversely impacting network performance. Furthermore, our results indicate that TCP performance will worsen as WAN speeds continue to increase.
Wu-chun Feng, Peerapol Tinnakornsrisuphap
SC1
1999 Dynamic Client-Side Scheduling in a Real-Time CORBA System
abstract
CORBA allows objects to communicate, independent of the specific techniques, languages, and platforms used to implement the objects. However, due to the multilevel software layering needed to provide this independence, CORBA cannot support real time applications since it lacks essential quality of service (QoS) features. Recent work on real time CORBA includes an offline scheduled, hard, real time system based on rate-monotonic scheduling and an online scheduled, best-effort, real time system based on the earliest-deadline-first algorithm. The former provides QoS guarantees at the expense of run time scheduling flexibility while the latter provides the complement. We propose an approach which provides the advantages of both, that is, QoS guarantees and run time scheduling flexibility.
Wu-chun Feng
COMPSAC1
1997 Algorithms for Scheduling Real-Time Tasks with Input Error and End-to-End Deadlines
abstract
This paper describes algorithms for scheduling preemptive, imprecise, composite tasks in real-time. Each composite task consists of a chain of component tasks, and each component task is made up of a mandatory part and an optional part. Whenever a component task uses imprecise input, the processing times of its mandatory and optional parts may become larger. The composite tasks are scheduled by a two-level scheduler. At the high level, the composite tasks are scheduled preemptively on one processor, according to an existing algorithm for scheduling simple imprecise tasks. The low-level scheduler then distributes the time budgeted for each composite task across its component tasks so as to minimize the output error of the composite task.
Wu-chun Feng, Jane W.-S. Liu
IEEE Trans. Software Eng.1