Judy Qiu

dblp:68/8357 · DBLP profile ↗
← Back
40ranked-venue papers
6as first author
2since 2021 · last 2023
0000-0002-3900-715XORCID · corroborated

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

Systems, architecture and hardware · 22 · 4 first-author · 1 since 2021Applied, interdisciplinary, general and emerging computing · 12 · 1 first-authorSoftware engineering, systems software and programming languages · 5Databases, data management, data science and information retrieval · 5 · 1 first-author · 1 since 2021Artificial intelligence and machine learning · 3

Expertise — from the expertise taxonomy: the topics of the expert's papers under the CCF categories. A weight counts papers with recency: 1 for a paper about the topic, 0.3 when the topic is its context, halved every five years.

Computer architecture, parallel and distributed computing, and storage systems
4 papers
Parallel and multicore computing · 53% Cloud and datacenter computing · 32% High-performance computing · 15%
Artificial intelligence
1 paper
Generative modeling · 100%
Human-computer interaction and pervasive computing
1 paper
Learning and educational technologies · 100%
Interdisciplinary, comprehensive, and emerging computing
4 papers
Bioinformatics and computational biology · 52% Computing education · 48%
Databases, data mining, and information retrieval
2 papers
Data mining · 87% Information retrieval · 13%
Computer graphics and multimedia
2 papers
Visualization and visual analytics · 100%

Topics — the 20 heaviest of 24, each with the papers that count most for it

TopicWeightPapersLastEvidence papers
Machine learning › Generative modeling
generative adversarial network
0.712023
ExamGAN and Twin-ExamGAN for Exam Script Generation · IEEE Trans. Knowl. Data Eng. 2023
Parallel and multicore computing › data-parallel programming
mapreduce
0.222011
Cloud Technologies for Bioinformatics Applications · IEEE Trans. Parallel Distributed Syst. 2011
Twister: a runtime for iterative MapReduce · HPDC 2010
Visualization and visual analytics
high-dimensional data visualization
0.122010
Browsing large scale cheminformatics data with dimension reduction · HPDC 2010
Dimension reduction and visualization of large high-dimensional data via interpolation · HPDC 2010
Cloud and datacenter computing
job scheduling
0.112011
Cloud Technologies for Bioinformatics Applications · IEEE Trans. Parallel Distributed Syst. 2011
High-performance computing › high-throughput computing
many-task computing
0.112011
Cloud Technologies for Bioinformatics Applications · IEEE Trans. Parallel Distributed Syst. 2011
Data mining
dimensionality reduction
0.112010
Dimension reduction and visualization of large high-dimensional data via interpolation · HPDC 2010
Data mining › dimensionality reduction
multidimensional scaling
0.112010
Dimension reduction and visualization of large high-dimensional data via interpolation · HPDC 2010
Visualization and visual analytics
dimensionality reduction
0.112010
Browsing large scale cheminformatics data with dimension reduction · HPDC 2010
Parallel and multicore computing › data parallelism
data-parallel applications
0.112010
Twister: a runtime for iterative MapReduce · HPDC 2010
Cloud and datacenter computing › cluster computing framework
mapreduce framework
0.112010
Cloud computing paradigms for pleasingly parallel biomedical applications · HPDC 2010
Parallel and multicore computing › parallel programming runtimes
mapreduce runtime
0.112010
Twister: a runtime for iterative MapReduce · HPDC 2010
Parallel and multicore computing
parallel programming runtimes
0.112010
Twister: a runtime for iterative MapReduce · HPDC 2010
Cloud and datacenter computing
utility computing
0.112010
Cloud computing paradigms for pleasingly parallel biomedical applications · HPDC 2010
Bioinformatics and computational biology › sequence analysis › sequence assembly
EST assembly
0.012011
Cloud Technologies for Bioinformatics Applications · IEEE Trans. Parallel Distributed Syst. 2011
Bioinformatics and computational biology
sequence alignment
0.012011
Cloud Technologies for Bioinformatics Applications · IEEE Trans. Parallel Distributed Syst. 2011
Bioinformatics and computational biology › sequence analysis
sequence assembly
0.012011
Cloud Technologies for Bioinformatics Applications · IEEE Trans. Parallel Distributed Syst. 2011
Bioinformatics and computational biology › molecular informatics
cheminformatics
0.012010
Browsing large scale cheminformatics data with dimension reduction · HPDC 2010
Bioinformatics and computational biology › sequence analysis › sequence assembly
genome assembly
0.012010
Cloud computing paradigms for pleasingly parallel biomedical applications · HPDC 2010
Information retrieval › interactive information retrieval
browsing
0.012010
Browsing large scale cheminformatics data with dimension reduction · HPDC 2010
High-performance computing › data-intensive computing
parallel data analysis
0.012010
Dimension reduction and visualization of large high-dimensional data via interpolation · HPDC 2010

Methods — techniques the papers use, named apart from their topics

generative adversarial network · 2.0deep learning · 2.0generative topographic mapping · 0.7parallelization · 0.3multidimensional scaling · 0.3interpolation · 0.3virtualization · 0.2hadoop · 0.2MPI · 0.2DryadLINQ · 0.2iterative mapreduce · 0.1
YearPublicationVenuePosition
2023 ExamGAN and Twin-ExamGAN for Exam Script Generation
abstract
Nowadays, the learning management system (LMS) has been widely used in different educational stages from primary to tertiary education for student administration, documentation, tracking, reporting, and delivery of educational courses, training programs, or learning and development programs. Towards effective learning outcome assessment, the exam script generation problem has attracted many attentions recently. But the research in this field is still in its early stage. Two essential issues have been ignored largely by existing solutions. First, given a course, it is unknown yet how to generate an quality exam script which concurrently has (i) the proper difficulty level, (ii) the coverage of essential knowledge points, (iii) the capability to distinguish academic performances between students, and (iv) the student scores in normal distribution. Second, while frequently encountered in practice, it is unknown so far how to generate a pair of high quality exam scripts which are equivalent in assessment (i.e., the student scores are comparable by taking either of them) but have significantly different sets of questions. To fill the gap, this paper proposes ExamGAN (Exam Script Generative Adversarial Network) to generate high quality exam scripts, and then extends ExamGAN to T-ExamGAN (Twin-ExamGAN) to generate a pair of high quality exam scripts. Based on extensive experiments on three benchmark datasets, it has verified the superiority of proposed solutions in various aspects against the state-of-the-art. Moreover, we have conducted a case study which demonstrated the effectiveness of proposed solution in the real teaching scenarios.
Zhengyang Wu 0001, Judy Qiu, Yong Tang 0001
IEEE Trans. Knowl. Data Eng.3
2021 Rank Position Forecasting in Car Racing
abstract
Rank position forecasting in car racing is a challenging problem when using a Deep Learning-based model over time-series data. It is featured with highly complex global dependency among the racing cars, with uncertainty resulted from existing and external factors; and it is also a problem with data scarcity. Existing methods, including statistical models, machine learning regression models, and several state-of-the-art deep forecasting models all perform not well on this problem. By an elaborate analysis of pit stop events, we find it critical to decompose the cause-and-effect relationship and model the rank position and pit stop events separately. In choosing a sub-model from different neural network models, we find the model with weak assumptions on the global dependency structure performs the best. Based on these observations, we propose RankNet, a combination of the encoder-decoder network and a separate Multilayer Perception network that is capable of delivering probabilistic forecasting to model the pit stop events and rank position in car racing. Further with the help of feature optimizations, RankNet demonstrates a significant performance improvement, where MAE improves 19% in two laps forecasting task and 7% in the stint forecasting task over the best baseline and is also more stable when adapting to unseen new data. Details of the model optimizations and performance profiling are presented. It is promising to provide useful interactions of neural networks in forecasting racing cars and shine a light on solutions to similar challenging issues in general forecasting problems.
Bo Peng 0011, Selahattin Akkas, Takuya Araki, Ohno Yoshiyuki, Judy Qiu
IPDPS6
2019 Anomaly Detection over Streaming Data: Indy500 Case Study
abstract
Sports racing is attracting billions of audiences each year. It is powered and transformed by the latest data analysis technologies, from race car design, driving skill improvements to audience engagement on social media. However, most of the data processing are off-line and retrospective analysis. The emerging real-time data analysis from the Internet of Things (IoT) result in fast data streams generated from distributed sensors. Applying advanced Machine Learning/Artificial Intelligence over such data streams to discover new information, predict future insights and make control decision is a crucial process. In this paper, we start by articulating racing car big data characteristics and present time-critical anomaly detection of the racing cars with the real-time sensors of cars and the tracks from actual racing events. We build a scalable system infrastructure based on neuro-morphic Hierarchical Temporal Memory Algorithm (HTM) algorithm and Storm stream processing engine. By courtesy of historical Indy500 racing logs, evaluation experiments on this prototype system demonstrate good performance in terms of anomaly detection accuracy and service level objective (SLO) of latency for a real-world streaming application.
Chathura Widanage, Sahil Tyagi, Ravi Teja, Bo Peng 0011, Supun Kamburugamuve, Dan Baum, Dayle Smith, Judy Qiu, Jon Koskey
CLOUD9
2019 A Fast Video Image Detection using TensorFlow Mobile Networks for Racing Cars
abstract
With the growth of the Internet of Things, we see an increase in the importance of analysis of data from the edge, often with the results needed in real-time. Indy Car series is one of the well-known racing series in North America. All cars are equipped with multiple cameras. The video streams captured by these cameras can be used for detection and predictive tasks to increase race safety and develop better strategies to win the race. Moreover, the data can be used together with the telemetry data to provide better analysis and predictions for the drivers and the teams. In a lot of video analytics tasks, the tasks begin with object detection as its foundation. The existing pretrained object detection models are inadequate to detect IndyCar race cars. Therefore, we have created a new dataset and have compared three different Single Shot Multibox Detector models from TensorFlow Detection Model Zoo. We run experiments on CPU and GPU. Since transferring the data from edge devices to a server, running inference, and sending the result back is time and resource consuming, we also test mobile detection models on an Edge TPU, which is a Google Coral Dev Board. Our initial results show that the Edge TPU gives the best inference time, and it is more suitable for a real-time machine learning task.
Selahattin Akkas, Sahaj Singh Maini, Judy Qiu
IEEE BigData3
2019 SubGraph2Vec: Highly-Vectorized Tree-like Subgraph Counting
abstract
Subgraph counting aims to count occurrences of a template T in a given network G (V, E). It is a powerful graph analysis tool and has found real-world applications in diverse domains. Scaling subgraph counting problems is known to be memory bounded and computationally challenging with exponential complexity. Although scalable parallel algorithms are known for several graph problems such as Triangle Counting and PageRank, this is not common for counting complex subgraphs. Here we address this challenge and study connected acyclic graphs or trees. We propose a novel vectorized subgraph counting algorithm, named SUBGRAPH2VEC, as well as both shared memory and distributed implementations: 1) reducing algorithmic complexity by minimizing neighbor traversal; 2) achieving a highly-vectorized implementation upon linear algebra kernels to significantly improve performance and hardware utilization. 3) SUBGRAPH2VEC improves the overall performance over the state-of-the-art work by orders of magnitude and up to 660x on a single node. 4) SUBGRAPH2VEC in distributed mode can scale up the template size to 20 and maintain good strong scalability. 5) enabling portability to both CPU and GPU.
Langshi Chen, Süleyman Cenk Sahinalp, Madhav V. Marathe, Anil Vullikanti, Andrey Nikolaev, Egor Smirnov, Ruslan Israfilov, Judy Qiu
IEEE BigData9
2019 HarpGBDT: Optimizing Gradient Boosting Decision Tree for Parallel Efficiency
abstract
Gradient Boosting Decision Tree (GBDT) is a widely used machine learning algorithm, whose training involves both irregular computation and random memory access and is challenging for system optimizations. In this paper, we conduct a comprehensive performance analysis of two state-of-the-art systems, XGBoost and LightGBM. They represent two typical parallel implementations for GBDT; one is data parallel and the other one is parallel over features. Substantial thread synchronization overhead, as well as the inefficiency of random memory access, is identified. We propose HarpGBDT, a new GBDT system designed from the perspective of parallel efficiency optimization. Firstly, we adopt a new tree growth method that selects the top K candidates of tree nodes to enable the use of more levels of parallelism without sacrificing the algorithm's accuracy. Secondly, we organize the training data and model data in blocks and propose a block-wise approach as a general model that enables the exploration of various parallelism options. Thirdly, we propose a mixed mode to utilize the advantages of a different mode of parallelism in different phases of training. By changing the configuration of the block size and parallel mode, HarpGBDT is able to attain better parallel efficiency. By extensive experiments on four datasets with different statistical characteristics on the Intel(R) Xeon(R) E5-2699 server, HarpGBDT on average performs 8x faster than XGBoost and 2.6x faster than LightGBM.
Bo Peng 0011, Judy Qiu, Langshi Chen, Selahattin Akkas, Egor Smirnov, Ruslan Israfilov, Sergey Khekhnev, Andrey Nikolaev
CLUSTER2
2018 Performance Characterization of Multi-threaded Graph Processing Applications on Many-Integrated-Core Architecture
abstract
In the age of Big Data, parallel graph processing has been a critical technique to analyze and understand connected data. Meanwhile, Moore's Law continues by integrating more cores into a single chip in the deep-nano regime. Many-Integrated-Core (MIC) processors emerge as a promising solution to process large graphs. In this paper, we empirically evaluate various computing platforms including an Intel Xeon E5 CPU, an Nvidia Tesla P40 GPU and a Xeon Phi 7210 MIC processor codenamed Knights Landing (KNL) in the domain of parallel graph processing. We show that the KNL gains encouraging performance and power efficiency when processing graphs, so that it can become an auspicious alternative to traditional CPUs and GPUs. We further characterize the impact of KNL architectural enhancements on the performance of a state-of-the-art graph framework. We have four key observations: 1 Different graph applications require distinctive numbers of threads to reach the peak performance. For the same application, various datasets need even different numbers of threads to achieve the best performance. 2 Not all graph applications actually benefit from high bandwidth MCDRAMs, while some of them favor low latency DDR4 DRAMs. 3 Vector processing units executing AVX512 SIMD instructions on KNLs are underutilized when running the state-of-the-art graph framework. 4 The sub-NUMA cache clustering mode offering the lowest local memory access latency hurts the performance of graph benchmarks that are lack of NUMA awareness. At last, we suggest future works including system auto-tuning tools and graph framework optimizations to fully exploit the potential of KNL for parallel graph processing.
Lei Jiang 0001, Langshi Chen, Judy Qiu
ISPASS3
2017 Benchmarking Harp-DAAL: High Performance Hadoop on KNL Clusters
abstract
Data analytics is undergoing a revolution in many scientific domains, and demands cost-effective parallel data analysis techniques. Traditional Java-based Big Data processing tools like Hadoop MapReduce are designed for commodity CPUs. In contrast, emerging manycore processors like the Xeon Phi have an order of magnitude greater computation power and memory bandwidth. To harness their computing capabilities, we propose the Harp-DAAL framework. We show that enhanced versions of MapReduce can be replaced by Harp, a Hadoop plug-in, that offers useful data abstractions for both high-performance iterative computation and MPI-quality communication, as well as drive Intel's native DAAL library. We select a subset of three machine learning algorithms and implement them within Harp-DAAL. Our scalability benchmarks ran on Knights Landing (KNL) clusters and achieved up to 2.5 times speedup of performance over the HPC solution in NOMAD and 15 to 40 times speedup over Java-based solutions in Spark. We further quantify the workloads on single node KNL with a performance breakdown at the micro-architecture level.
Langshi Chen, Bo Peng 0011, Bingjing Zhang, Xu T. Liu, Yiming Zou, Lei Jiang 0001, Robert Henschel, Craig A. Stewart, Emily McCallum, Tom Zahniser, Jon Omer, Judy Qiu
CLOUD13
2017 HarpLDA+: Optimizing latent dirichlet allocation for parallel efficiency
abstract
Latent Dirichlet Allocation (LDA) is a widely used machine learning technique in topic modeling and data analysis. Training large LDA models on big datasets involves dynamic and irregular computation patterns and is a major challenge to both algorithm optimization and system design. In this paper, we present a comprehensive benchmarking of our novel synchronized LDA training system HarpLDA+ based on Hadoop and Java. It demonstrates impressive performance when compared to three other MPI/C++ based state-of-the-art systems, which are LightLDA, F+NomadLDA, and WarpLDA. HarpLDA+ uses optimized collective communication with a timer control for load balance, leading to stable scalability in both shared-memory and distributed systems. We demonstrate in the experiments that HarpLDA+ is effective in reducing synchronization and communication overhead and outperforms the other three LDA training systems.
Bo Peng 0011, Bingjing Zhang, Langshi Chen, Mihai Avram 0002, Robert Henschel, Craig A. Stewart, Shaojuan Zhu, Emily McCallum, Lisa Smith, Tom Zahniser, Jon Omer, Judy Qiu
IEEE BigData12
2015 HPC-ABDS High Performance Computing Enhanced Apache Big Data Stack
abstract
We review the High Performance Computing Enhanced Apache Big Data Stack HPC-ABDS and summarize the capabilities in 21 identified architecture layers. These cover Message and Data Protocols, Distributed Coordination, Security & Privacy, Monitoring, Infrastructure Management, DevOps, Interoperability, File Systems, Cluster & Resource management, Data Transport, File management, NoSQL, SQL (NewSQL), Extraction Tools, Object-relational mapping, In-memory caching and databases, Inter-process Communication, Batch Programming model and Runtime, Stream Processing, High-level Programming, Application Hosting and PaaS, Libraries and Applications, Workflow and Orchestration. We summarize status of these layers focusing on issues of importance for data analytics. We highlight areas where HPC and ABDS have good opportunities for integration.
Geoffrey C. Fox, Judy Qiu, Supun Kamburugamuve, Shantenu Jha, André Luckow
CCGRID2
2015 Parallel Clustering of High-Dimensional Social Media Data Streams
abstract
We introduce Cloud DIKW (Data, Information, Knowledge, Wisdom) as an analysis environment supporting scientific discovery through integrated parallel batch and streaming processing, and apply it to one representative domain application: social media data stream clustering. In this context, recent work demonstrated that high-quality clusters can be generated by representing the data points using high-dimensional vectors that reflect textual content and social network information. However, due to the high cost of similarity computation, sequential implementations of even single-pass algorithms cannot keep up with the speed of real-world streams. This paper presents our efforts in meeting the constraints of realtimesocial media stream clustering through parallelization in Cloud DIKW. Specifically, we focus on two system-level issues. Firstly, most stream processing engines such as Apache Storm organize distributed workers in the form of a directed acyclic graph (DAG), which makes it difficult to dynamically synchronize the state of parallel clustering workers. We tackle this challenge by creating a separate synchronization channel using a pub-sub messaging system (ActiveMQ in our case). Secondly, due to the sparsity of the high-dimensional vectors, the size of cancroids grows quickly as new data points are assigned tithe clusters. As a result, traditional synchronization that directly broadcasts cluster cancroids becomes too expensive and limits the scalability of the parallel algorithm. We address this problem by communicating only dynamic changes of the clusters rather than the whole centred vectors. Our algorithm under Cloud DIKWcan process the Twitter 10% data stream ("gardenhose") in realtimewith 96-way parallelism. By natural improvements to CloudDIKW, including advanced collective communication techniques developed in our Harp project, we will be able to process the full Twitter data stream in real-time with 1000-way parallelism. Our use of powerful general software subsystems will enable many other applications that need integration of streaming and batch data analytics.
Xiaoming Gao, Emilio Ferrara, Judy Qiu
CCGRID3
2015 Harp: Collective Communication on Hadoop
abstract
Big data processing tools have evolved rapidly in recent years. MapReduce has proven very successful but is not optimized for many important analytics, especially those involving iteration. In this regard, Iterative MapReduce frameworks improve performance of MapReduce job chains through caching. Further, Pregel, Giraph and Graph Lab abstract data as a graph and process it in iterations. But all these tools are designed with a fixed data abstraction and have limited collective communication support to synchronize application data and algorithm control states among parallel processes. In this paper, we introduce a collective communication abstraction layer which provides efficient collective communication operations on several common data abstractions such as arrays, key-values and graphs, and define a Map Collective programming model which serves the diverse collective communication demands in different parallel algorithms. We implement a library called Harp to provide the features above and plug it into Hadoop so that applications abstracted in Map Collective model can be easily developed on top of MapReduce framework and conveniently integrated with other tools in Apache Big Data Stack. With improved expressiveness in the abstraction and excellent performance on the implementation, we can simultaneously support various applications from HPC to Cloud systems together with high performance.
Bingjing Zhang, Yang Ruan 0001, Judy Qiu
IC2E3
2014 Supporting Queries and Analyses of Large-Scale Social Media Data with Customizable and Scalable Indexing Techniques over NoSQL Databases
abstract
Social media data analysis demonstrates two special characteristics in Big Data processing. First, most analyses focus on data subsets related to specific social events or activities instead of the whole dataset. Second, analysis workflows consist of multiple stages, and algorithms applied in each stage may use different computation and communication patterns depending on processing frameworks. This paper presents our efforts in supporting the data storage and processing requirements for such characteristics. To achieve efficient queries about target data subsets, we propose a general customizable and scalable indexing framework that can be built over distributed NoSQL databases. This framework allows users to define suitable customized index structures for their query patterns against social media data, and supports scalable indexing of both historical and streaming data. We implement this framework on HBase, and name it IndexedHBase. Starting from IndexedHBase, we build a distributed analysis stack based on YARN to support analysis algorithms using different processing frameworks, such as Hadoop MapReduce, Harp, and Giraph. This analysis stack is used to host the Truthy social media data observatory, and we have applied the customized index structures in supporting both query evaluation and sophisticated analysis algorithms. Performance tests show that our solutions outperform implementations using both direct raw data scans and current indexing mechanisms in existing NoSQL databases.
Xiaoming Gao, Judy Qiu
CCGRID2
2014 Towards a Collective Layer in the Big Data Stack
abstract
We generalize MapReduce, Iterative MapReduce and data intensive MPI runtime as a layered Map-Collective architecture with Map-All Gather, Map-All Reduce, MapReduce Merge Broadcast and Map-Reduce Scatter patterns as the initial focus. Map-collectives improve the performance and efficiency of the computations while at the same time facilitating ease of use for the users. These collective primitives can be applied to multiple runtimes and we propose building high performance robust implementations that cross cluster and cloud systems. Here we present results for two collectives shared between Hadoop (where we term our extension H-Collectives) on clusters and the Twister4Azure Iterative MapReduce for the Azure Cloud. Our prototype implementations of Map-All Gather and Map-All Reduce primitives achieved up to 33% performance improvement for K-means Clustering and up to 50% improvement for Multi-Dimensional Scaling, while also improving the user friendliness. In some cases, use of Map-collectives virtually eliminated almost all the overheads of the computations.
Thilina Gunarathne, Judy Qiu, Dennis Gannon
CCGRID2
2014 CINET 2.0: A CyberInfrastructure for Network Science
abstract
Analysis of structural properties and dynamics of networks is currently a central topic in many disciplines including Social Sciences, Biology and Business. CINET, a cyber infrastructure for such studies, introduced the concept of supporting network analysis as a service. The basic idea is to allow experts in various disciplines to focus on obtaining domain-specific insights from the results of network analyses instead of worrying about programming details and allocation of computational resources needed to carry out the analyses. A basic version of CINET was released in May 2012. This paper discusses CINET 2.0, a significantly enhanced version that supports complex network analyses through a web portal. CINET 2.0 has already been used for teaching courses related to Network Science at several US universities. In this paper, we discuss how CINET 2.0 significantly extends CINET 1.0 through enhancements to some components and the addition of new components.
Sherif Hanie El Meligy Abdelhamid, Md. Maksudul Alam, Richard A. Aló, S. M. Arifuzzaman, Pete Beckman, Tirtha Bhattacharjee, Md Hasanuzzaman Bhuiyan, Keith R. Bisset, Stephen G. Eubank, Albert C. Esterline, Edward A. Fox, Geoffrey C. Fox, S. M. Shamimul Hasan, Harshal Hayatnagarkar, Maleq Khan, Chris J. Kuhlman, Madhav V. Marathe, Natarajan Meghanathan, Henning S. Mortveit, Judy Qiu, S. S. Ravi, Zalia Shams, Ongard Sirisaengtaksin, Samarth Swarup, Anil Vullikanti, Tak-Lon Wu
eScience20
2014 Hierarchical MapReduce: towards simplified cross-domain data processing
abstract
SUMMARY The MapReduce programming model has proven useful for data‐driven high throughput applications. However, the conventional MapReduce model limits itself to scheduling jobs within a single cluster. As job sizes become larger, single‐cluster solutions grow increasingly inadequate. We present a hierarchical MapReduce framework that utilizes computation resources from multiple clusters simultaneously to run MapReduce job across them. The applications implemented in this framework adopt theMap–Reduce–GlobalReducemodel where computations are expressed as three functions: Map, Reduce, and GlobalReduce. Two scheduling algorithms are proposed, one that targets compute‐intensive jobs and another data‐intensive jobs, evaluated using a life science application, AutoDock, and a simple Grep. Data management is explored through analysis of the Gfarm file system.Copyright © 2012 John Wiley & Sons, Ltd.
Beth Plale, Zhenhua Guo 0004, Wilfred W. Li, Judy Qiu, Yiming Sun 0001
Concurr. Comput. Pract. Exp.5
2014 Emerging Computational Methods for the Life Sciences Workshop 2012
abstract
Computing systems are rapidly changing with multicore, graphics processing units (GPUs), clusters, volunteer systems, clouds, and grids offering a confusing dazzling array of opportunities. New programming paradigms such as Google MapReduce and many-task computing have joined the traditional repertoire of workflow and parallel computing for the highest performance systems. Meanwhile, the life sciences are continuing to expand in data generated with continuing improvement in the instruments for high-throughput analysis. This ‘fourth paradigm’ (data driven science) is joined by complex systems or biocomplexity that can build phenomenological models of biological systems and processes. This special issue for Emerging Computational Methods for the Life Sciences Workshop ECMLS2012 1, juxtaposes these trends seeking those computational methods that will enhance scientific discovery. Within this overall scope, this special issue encouraged researchers to submit and present original work related to the latest trends in parallel and distributed high-performance systems applied to life science problems. Weber et al. 2 note that GPUs and multicore processors are now pervasive in computational sciences and high-performance computing. Their high-arithmetic throughput and memory bandwidth combined with their ever increasing programmability make them suitable for a widening variety of applications. They give a high-level overview of Specmaster, a MyriMatch port that can use every Open Computing Language (OpenCL) device available in a machine to identify peptides in tandem mass spectrometry data. Then they highlight device-specific optimizations for multicore CPUs and GPUs as well as describing framework for implementing these optimizations while still using a single code-base. They also provide performance results of Specmaster running on four different architectures and compare these numbers to MyriMatch. Finally, they improve on their existing work by showing Specmaster dynamically load balancing on 3 AMD Radeon 7970s and 32 AMD Opteron Interlagos 6272 cores as well as comparing the quality of their search results to MyriMatch. Yang et al. 3 study a difficulty in building a mechanistic model of biological systems coming from determination of correct parameter values. Their paper proposes a novel parameter estimation method to infer unknown parameters, such as kinetic rates, from noisy experimental observations. Derived from the ABC sequential Monte Carlo algorithm, their method predicts the distribution of each parameter rather than a single value via several intermediate distributions. Motivated by the computational intensity of the method, they improve the ABC sequential Monte Carlo method in two aspects. First, to increase efficiency, a windowing method is developed to reduce the parameter searching space, and an adaptive sampling weight mechanism is introduced to make the intermediate distributions converge to the target distributions in a much quicker manner. Second, to speedup the estimation process, they implement their method in a parallel computing environment to speedup the sampling process. Elllingson et al. 4 present the current state of high-throughput virtual screening. They describe a case study of using a task-parallel MPI version of Autodock4 to run a virtual high-throughput screen of one million compounds on the Jaguar Cray XK6 Supercomputer at Oak Ridge National Laboratory. The paper includes a description of scripts developed to increase the efficiency of the predocking file preparation and postdocking analysis. A detailed tutorial, scripts, and source code for this MPI version of Autodock4 are available online at http://www.bio.utk.edu/baudrylab/autodockmpi.htm Hamacher et al. 5 present a massively parallel implementation of the computation of coevolutionary signals from biomolecular sequence alignments based on mutual information (MI) and a normalization procedure to neutral evolution. The MI is computed for two-point and three-point correlations within any multiple sequence alignment. The high-computational demand in the normalization procedure is met efficiently with an implementation on GPUs using NVIDIA's CUDA framework. In particular, the normalization of the MI for three-point ‘cliques’ of amino acids or nucleotides requires large sampling numbers in the normalization that is achieved by using GPUs. GPU computation serves as an enabling technology here insofar as MI normalization is also possible using traditional computational methods or cluster computation, but only GPU computation makes MI normalization for sequence analysis feasible in a statistically sufficient sample and in acceptable time given affordable commodity hardware. They illustrate a) the computational efficiency and b) the biological usefulness of two-point and three-point MI by applications to the well-known protein calmodulin and the variable surface glycoprotein of Trypanosoma brucei which are subject to involved evolutionary pressure. Here, they find striking coevolutionary patterns and distinct information on the molecular evolution of these molecules that question previous work that relied on inefficient MI computations. Cushing et al. 6 note that task farming is often used to enable parameter sweep for exploration of large sets of initial conditions for large scale complex simulations. Such applications occur very often in life sciences. Available solutions enable us to perform parameter sweep by creating multiple job submissions with different parameters. This paper presents an approach to farm workflows employing service oriented paradigms using the WS-VLAM workflow manager from University of Amsterdam, which provides ways to create, control, and monitor workflows applications and their components. They present two service oriented approaches for workflow farming: task-level, whereby task harness acts as services by being invoked on which task to load, and data-level where the actual task is invoked as a service with different chunks of data to process. An experimental evaluation of the presented solution is performed with a biomedical application for which 3000 simulations were required to perform a Monte Carlo study. Finally, Stanberry et al. 7 note that modern biology is experiencing a rapid increase in data volumes that challenges analytical skills and existing cyberinfrastructure. Exponential expansion of the protein sequence universe (PSU), the protein sequence space, together with the costs and complexities of manual curation creates a major bottleneck in life sciences research. Existing resources lack scalable visualization tools that are instrumental for functional annotation. They describe a new visualization tool using multidimensional scaling to create a 3D embedding of the protein space. The advantages of the proposed PSU method include the ability to scale to large numbers of sequences, integrate different similarity measures with other functional and experimental data, and facilitate protein annotation. They apply the method to visualize the prokaryotic PSU by using sequence alignment scores. As an annotation example, they use an interpolation approach to map the set of annotated archaeal proteins into the prokaryotic PSU. Transdisciplinary approaches akin to the one described in this paper are urgently needed to quickly and efficiently translate the influx of new data into tangible innovations and groundbreaking discoveries We would like to thank the authors for contributing papers on their research on latest trends in data intensive technologies and applications for this special issue, and thank all the reviewers for providing constructive reviews and in helping to shape this special issue. Finally, we would like to thank the editors of Concurrency and Computation: Practice and Experience for providing us an opportunity to bring this special issue to the research community.
Judy Qiu, Ian T. Foster, Carole A. Goble
Concurr. Comput. Pract. Exp.1
2014 Special Issue for Emerging Computational Methods for the Life Sciences Workshop
abstract
Computing systems are rapidly changing with multicores, graphics processing units, clusters, volunteer systems, clouds, and grids, offering a confusing dazzling array of opportunities. New programming paradigms such as MapReduce and many-task computing have joined the traditional repertoire of workflow and parallel computing for the highest-performance systems. Meanwhile, the life sciences are continuing to expand in data generated, with continuing improvement in the instruments for high-throughput analysis. This ‘fourth paradigm’ (observationally driven science) is joined by complex systems or biocomplexity that can build phenomenological models of biological systems and processes. This special issue juxtaposes these trends, seeking those computational methods that will enhance scientific discovery. Within this overall scope, this special issue encouraged researchers to submit and present original work related to the latest trends in parallel and distributed high-performance systems applied to life science problems. Mitchel et al. 1 presents parallel implementations of two popular microarray data analysis techniques: exploratory clustering analyses using the random forest classifier and feature selection through identification of differentially expressed genes using the rank product method. The authors have parallelized these two applications using the SPRINT, which is a library for R that aims to reduce the complexity of using HPC systems by providing biostatisticians with drop-in parallelized replacements of existing R functions. The paper demonstrates how one can parallelize R routines with minimum changes to the existing codes with the help of SPRINT, speeding up serialized and time-consuming analysis procedures written in R. Authors also implemented a tree-reduction algorithm for parallel combining of the results, which showed a surprisingly large effect on the overall performance. The paper also provides experimental results achieving 40 times speedup over serialized codes by using 128 processes. Lanc et al. 2 describes the adaptation and parallelization process of Paired-End Mapper structural variation pipeline and the Burrows–Wheeler alignment tool for execution on clusters, grids, and clouds using the weaver/starch/makeflow workflow stack. Authors describe the application of previous obtained lessons to a new workflow with and without shared file storage to tract the intractable sequential running times of these applications on large datasets. Authors present lessons and results for refactoring bioinformatics tools for elastic scaling on personal clouds and describe the various challenges faced when constructing such a workflow, from dealing failure detection to managing dependencies and handling the quirks of the underlying operating systems. Authors scale the workflows on hundreds of processors, reducing the run times of the two workflows to hours from days with high speedup. The lessons and the experiences presented in this paper can lower the barrier to scalable execution of workflows, allowing users to better harness the power of heterogeneously distributed systems for their own tools. Luo et al. 3 describes an enhanced MapReduce-based programming model ‘Map-Reduce-GlobalReduce’, where the computations are expressed as three functions: Map, Reduce, and GlobalReduce. The authors name this model as ‘Hierarchical MapReduce’. The hierarchical MapReduce framework divides the MapReduce computations and utilizes computation resources from multiple clusters simultaneously to execute a MapReduce job across them. The design is a powerful extension to MapReduce, especially to provide additional processing power for very large computations. Two static prior-knowledge-based scheduling algorithms are proposed, one that targets compute-intensive jobs and another that targets data-intensive jobs, evaluated using a life science application, AutoDock, and a simple Grep. The authors demonstrate the utility of their design and the performance metrics by greatly accelerating the application AutoDock across three large clusters. Jha et al. 4 presents a runtime environment, Distributed Application Runtime Environment (DARE), that supports the scalable, flexible, and extensible composition of capabilities exploring the interoperability among heterogeneously distributed computing environments for pleasingly parallel applications. DARE is a SAGA-BigJob-based framework motivated by the next-generation sequencing (NGS) analysis and other similar data-intensive applications. The proposed framework would enable NGS-like applications to run automatically on different infrastructures. DARE can utilize HPC, grid, and cloud infrastructures through a unified framework to achieve task-level concurrency. In this work, authors use BFAST as a representative stand-alone tool used for NGS data analysis and a ChIP-Seq pipeline as a representative pipeline-based approach. This paper represents the initial steps in the design and development of a general-purpose, scalable, and extensible infrastructure to support NGS (gene) data analytics. Ellingson et al. 5 describes their experience porting the AutoDock molecular docking program to run within the open-source Hadoop MapReduce framework. Virtual molecular docking is a task parallel computational method used in computer-aided drug discovery that calculates the binding affinity of a small-molecule drug candidate to a target protein. Authors evaluate the performance of the AutoDock Hadoop implementation on the 1088-core Kandinsky cluster located at the Oak Ridge National Laboratory. In this environment, the authors were able to achieve an impressive 450-fold speedup over a serial execution, reducing >1 year of work to ~1 day. Finally, Schatz 6 introduces the current research topics on computational methods of genomics, the complexity of biological applications and computational assays, and the increasing demands of improving algorithms and parallel systems. The challenges brought by the ever-increasing amount of data produced by advanced instruments are elaborated systematically in detail. The author discusses how parallel computing and cloud computing have been used to run large-scale biological applications and list the challenges of cloud computing for digital genomics. Issues such as big data, data security and privacy, and cost of cloud utility are discussed. The author also discusses the advantage of using hardware accelerators to empower the genomics analysis and speculates several future trends of digital demands of genomics that can potentially help researchers to reshape their thinking. We would like to thank the authors for contributing papers on their research on latest trends in data-intensive technologies and applications for this special issue and all the reviewers for providing constructive reviews and in helping to shape this special issue. Finally, we would like to thank the editors of Concurrency and Computation: Practice and Experience for providing us an opportunity to bring this special issue to the research community.
Judy Qiu, Ian T. Foster, Ronald C. Taylor
Concurr. Comput. Pract. Exp.1
2014 Visualizing the Protein Sequence Universe
abstract
SUMMARY Modern biology is experiencing a rapid increase in data volumes that challenges our analytical skills and existing cyberinfrastructure. Exponential expansion of the protein sequence universe (PSU), the protein sequence space, together with the costs and complexities of manual curation creates a major bottleneck in life sciences research. Existing resources lack scalable visualization tools that are instrumental for functional annotation. Here, we describe a new visualization tool using multidimensional scaling to create a 3D embedding of the protein space. The advantages of the proposed PSU method include the ability to scale to large numbers of sequences, integrate different similarity measures with other functional and experimental data, and facilitate protein annotation. We applied the method to visualize the prokaryotic PSU using sequence alignment scores. As an annotation example, we used the interpolation approach to map the set of annotated archaeal proteins into the prokaryotic PSU. Transdisciplinary approaches akin to the one described in this paper are urgently needed to quickly and efficiently translate the influx of new data into tangible innovations and groundbreaking discoveries. Copyright © 2013 John Wiley & Sons, Ltd.
Larissa Stanberry, Roger Higdon, Winston Haynes, Natali Kolker, William Broomall, Saliya Ekanayake, Adam Hughes, Yang Ruan 0001, Judy Qiu, Eugene Kolker, Geoffrey C. Fox
Concurr. Comput. Pract. Exp.9
2013 High performance clustering of social images in a map-collective programming model
abstract
Large-scale iterative computations are common in many important data mining and machine learning algorithms. Most of these applications can be specified as iterations of MapReduce computations, leading to the Iterative MapReduce programming model [1] for efficient execution of data-intensive iterative computations interoperably between HPC and cloud environments. We observe that a systematic approach to collective communication is essential but notably missing in the current model. Thus we generalize the iterative MapReduce concept to Map-Collective on the premise that large collectives are a distinctive feature of data intensive and data mining applications. To show the necessity of Map-Collective model, this paper studies the implications of large-scale social image clustering problems, where 10--100 million images represented as points in a high dimensional (up to 2048) vector space are required to be divided into 1--10 million clusters.
Bingjing Zhang, Judy Qiu
SoCC2
2013 Scalable parallel computing on clouds using Twister4Azure iterative MapReduce
Thilina Gunarathne, Bingjing Zhang, Tak-Lon Wu, Judy Qiu
Future Gener. Comput. Syst.4
2012 CINET: A cyberinfrastructure for network science
abstract
Networks are an effective abstraction for representing real systems. Consequently, network science is increasingly used in academia and industry to solve problems in many fields. Computations that determine structure properties and dynamical behaviors of networks are useful because they give insights into the characteristics of real systems. We introduce a newly built and deployed cyberinfrastructure for network science (CINET) that performs such computations, with the following features: (i) it offers realistic networks from the literature and various random and deterministic network generators; (ii) it provides many algorithmic modules and measures to study and characterize networks; (iii) it is designed for efficient execution of complex algorithms on distributed high performance computers so that they scale to large networks; and (iv) it is hosted with web interfaces so that those without direct access to high performance computing resources and those who are not computing experts can still reap the system benefits. It is a combination of application design and cyberinfrastructure that makes these features possible. To our knowledge, these capabilities collectively make CINET novel. We describe the system and illustrative use cases, with a focus on the CINET user.
Sherif Elmeligy Abdelhamid, Richard A. Aló, S. M. Arifuzzaman, Pete Beckman, Md Hasanuzzaman Bhuiyan, Keith R. Bisset, Edward A. Fox, Geoffrey C. Fox, Kevin Hall, S. M. Shamimul Hasan, Anurodh Joshi, Maleq Khan, Chris J. Kuhlman, Spencer J. Lee, Jonathan Leidig, Hemanth Makkapati, Madhav V. Marathe, Henning S. Mortveit, Judy Qiu, S. S. Ravi, Zalia Shams, Ongard Sirisaengtaksin, Rajesh Subbiah, Samarth Swarup, Nick Trebon, Anil Vullikanti
eScience19
2012 Mining hidden mixture context with ADIOS-P to improve predictive pre-fetcher accuracy
abstract
Predictive pre-fetcher, which predicts future data access events and loads the data before users requests, has been widely studied, especially in file systems or web contents servers, to reduce data load latency. Especially in scientific data visualization, pre-fetching can reduce the IO waiting time. In order to increase the accuracy, we apply a data mining technique to extract hidden information. More specifically, we apply a data mining technique for discovering the hidden contexts in data access patterns and make prediction based on the inferred context to boost the accuracy. In particular, we performed Probabilistic Latent Semantic Analysis (PLSA), a mixture model based algorithm popular in the text mining area, to mine hidden contexts from the collected user access patterns and, then, we run a predictor within the discovered context. We further improve PLSA by applying the Deterministic Annealing (DA) method to overcome the local optimum problem. In this paper we demonstrate how we can apply PLSA and DA optimization to mine hidden contexts from users data access patterns and improve predictive pre-fetcher performance.
Jong Choi 0001, Hasan Abbasi, David Pugmire, Norbert Podhorszki, Scott Klasky, Cristian Capdevila, Manish Parashar, Matthew Wolf, Judy Qiu, Geoffrey C. Fox
eScience9
2012 Interpolative multidimensional scaling techniques for the identification of clusters in very large sequence sets
abstract
BACKGROUND: Modern pyrosequencing techniques make it possible to study complex bacterial populations, such as 16S rRNA, directly from environmental or clinical samples without the need for laboratory purification. Alignment of sequences across the resultant large data sets (100,000+ sequences) is of particular interest for the purpose of identifying potential gene clusters and families, but such analysis represents a daunting computational task. The aim of this work is the development of an efficient pipeline for the clustering of large sequence read sets. METHODS: Pairwise alignment techniques are used here to calculate genetic distances between sequence pairs. These methods are pleasingly parallel and have been shown to more accurately reflect accurate genetic distances in highly variable regions of rRNA genes than do traditional multiple sequence alignment (MSA) approaches. By utilizing Needleman-Wunsch (NW) pairwise alignment in conjunction with novel implementations of interpolative multidimensional scaling (MDS), we have developed an effective method for visualizing massive biosequence data sets and quickly identifying potential gene clusters. RESULTS: This study demonstrates the use of interpolative MDS to obtain clustering results that are qualitatively similar to those obtained through full MDS, but with substantial cost savings. In particular, the wall clock time required to cluster a set of 100,000 sequences has been reduced from seven hours to less than one hour through the use of interpolative MDS. CONCLUSIONS: Although work remains to be done in selecting the optimal training set size for interpolative MDS, substantial computational cost savings will allow us to cluster much larger sequence sets in the future.
Adam Hughes, Yang Ruan 0001, Saliya Ekanayake, Seung-Hee Bae, Qunfeng Dong, Mina Rho, Judy Qiu, Geoffrey C. Fox
BMC Bioinform.7
2012 Performance of windows multicore systems on threading and MPI
abstract
SUMMARY We present performance results on a Windows cluster with up to 768 cores using Message Passing Interface (MPI) and two variants of threading—Concurrency and Coordination Runtime (CCR) and Task Parallel Library (TPL). CCR presents a message‐based interface, while TPL allows for loops to be automatically parallelized. MPI is used between the cluster nodes (up to 32) and either threading or MPI for parallelism on the 24 cores of each node. We look at the performance of two significant bioinformatics applications; gene clustering and dimension reduction. We find that the two threading runtimes offer similar performance with MPI outperforming both at low levels of parallelism but threading much better when the grain size (problem size per process/thread) is small. We develop simple models for the performance of the clustering code. Copyright © 2011 John Wiley & Sons, Ltd.
Judy Qiu, Seung-Hee Bae
Concurr. Comput. Pract. Exp.1
2012 Guest Editor's Introduction: Special Section on Challenges and Solutions in Multicore and Many-Core Computing
abstract
It is our honor to serve as guest editors of this special section of the journal of Concurrency and Computation: Practice and Experience on Frontiers of GPU, Multi- and Many-Core Systems (FGMMS). We are pleased to present nine high-quality contributions in this special issue, where they were first presented at the Frontiers of GPU, Multi- and Many-Core Systems Workshop in conjunction with the 10 th IEEE/ACM International Symposium on Cluster, Cloud and Grid Computing (CCGrid 2010), held from May 17 to 20, 2010, in Melbourne, Victoria, Australia. The invited papers in this special issue represent augmented works drafted at the beginning of 2010 that addressed the issues below. Multicore and many-core microprocessors are being deployed in a broad spectrum of applications including clusters, clouds, and grids. Both conventional multicore and many-core processors, such as Intel Nehalem and IBM Power7 processors, and unconventional many-core processors, such as NVIDIA Tesla and AMD FireStream graphics processing units (GPUs), hold the promise of increasing performance through parallelism. However, GPU approaches in parallelism are distinctly different from those of conventional multicore and many-core processors, which raises new challenges. For example, how do we optimize applications for conventional multicore and many-core processors? How do we re-engineer applications to take advantage of GPUs’ tremendous computing power in a reasonable cost–benefit ratio? What are effective ways of using GPUs as accelerators? Enormous and rapid progress has been made in accelerator computing over the last two years, but nevertheless we believe the themes developed for the FGMMS workshop and that are represented here in this special issue are still very important ones. In the last two years we have seen the continued rise of GPU computing and indeed its uptake as multi-GPU systems now dominates the top 10 entries within the Top 500 supercomputing systems list 1. Although there has been something of a shakeout of the accelerator technologies that were prevalent three years ago, we have also seen the continued and steady rise of uptake of multicore conventional CPU devices, and an exciting future with combined multicore CPU and GPU devices seems likely. As the articles in this special issue suggest, there are still challenges ahead for application developers to make best use of these future and hybrid highly concurrent systems. There are however still many really important applications that will continue to drive interest, investment, and development of these technologies. The goals of this special issue are to discuss these and other issues and bring together developers of application algorithms and experts in utilizing multicore and many-core processors. We briefly introduce the articles as follows. El Zein and Rendell 2 explore the effect of different GPU programming options (e.g., memory type, memory access methods, and data types) on the performance of routine evaluating sparse matrix vector products and discuss the method for optimal performance. Playne and Hawick 3 report on their approach in accelerating finite-differencing applications using multiple GPU devices with a single CPU host and the asynchronous CPU/GPU communication. Kato and Hosino 4 present their algorithms for speeding up a k-nearest neighbor problem in the recommendation system through multiple GPUs. Ino et al. 5 discuss a cooperative multitasking method for concurrent execution of scientific and graphics applications on GPU and acceleration of compute unified device architecture-based applications using idle GPU cycles in the office. Gillan et al. 6 present a case study on how the instruction-level parallelism offered by three accelerator technologies — field-programmable gate array, GPU, and ClearSpeed — can be exploited in atomic physics with considerable differences in the implementation strategies. Barhen et al. 7 present an unconventional fast Fourier transform implementation scheme for the IBM Cell B.E. processors, named transverse vectorization, and provide the first results for multifast Fourier transform implementation and application on the novel, ultralow power Coherent LogixHyperX processor. Zhou et al. 8 investigate the software and system issues in accelerating climate and weather models in a prototype hybrid computing system, which comprises Intel blades and IBM Cell B.E. blades, connected with both InfiniBand and 1-Gigabit Ethernet and communicate with IBM's Dynamic Application Virtualization software. Qiu and Bae 9 present performance results of two significant bioinformatics applications, gene clustering and dimension reduction, on a Microsoft Windows cluster with up to 768 cores using Message Passing Interface and two variants of threading — Concurrency and Coordination Runtime and Task Parallel Library. Hackenberg et al. 10 present a tool for the graphical program flow analysis of hardware accelerated parallel programs. The tool monitors the hybrid program execution to record and visualize many performance relevant events along the way, and is exemplified through representative real-world applications written for both IBM Cell B.E. processor and NVIDIA Compute Unified Device Architecture (CUDA) API. This work was supported in part by a Microsoft CRMC grant. The guest editors of this special issue would like to express their deep gratitude to all authors, external reviewers, and Geoffrey Fox for their efforts in making this issue possible.
Shujia Zhou, Judy Qiu, Kenneth A. Hawick
Concurr. Comput. Pract. Exp.2
2012 Special issue for data intensive eScience
Judy Qiu, Dennis Gannon
Distributed Parallel Databases1
2011 Analysis of Virtualization Technologies for High Performance Computing Environments
abstract
As Cloud computing emerges as a dominant paradigm in distributed systems, it is important to fully understand the underlying technologies that make Clouds possible. One technology, and perhaps the most important, is virtualization. Recently virtualization, through the use of hyper visors, has become widely used and well understood by many. However, there are a large spread of different hyper visors, each with their own advantages and disadvantages. This paper provides an in-depth analysis of some of today's commonly accepted virtualization technologies from feature comparison to performance analysis, focusing on the applicability to High Performance Computing environments using Future Grid resources. The results indicate virtualization sometimes introduces slight performance impacts depending on the hyper visor type, however the benefits of such technologies are profound and not all virtualization technologies are equal. From our experience, the KVM hyper visor is the optimal choice for supporting HPC applications within a Cloud infrastructure.
Andrew J. Younge, Robert Henschel, James T. Brown, Gregor von Laszewski, Judy Qiu, Geoffrey C. Fox
IEEE CLOUD5
2011 Browsing large-scale cheminformatics data with dimension reduction
abstract
SUMMARY Visualization of large‐scale high dimensional data is highly valuable for data analysis facilitating scientific discovery in many fields. We present PubChemBrowse, a customized visualization tool for cheminformatics research. It provides a novel 3D data point browser that displays complex properties of massive data on commodity clients. As in Geographic Information System browsers for Earth and Environment data, chemical compounds with similar properties are nearby in the browser. PubChemBrowse is built around in‐house high performance parallel Multi‐dimensional scaling and Generative topographic mapping services and supports fast interaction with an external property database. These properties can be overlaid on 3D mapped compound space or queried for individual points. We prototype the integration with Chem2Bio2RDF system using SPARQL endpoint to access over 20 publicly accessible bioinformatics databases. We describe our design and implementation of the integrated PubChemBrowse application and outline its use in drug discovery. The same core technologies are generally applicable to develop high performance scientific data browsing systems for other applications. Copyright © 2011 John Wiley & Sons, Ltd.
Jong Choi 0001, Seung-Hee Bae, Judy Qiu, Bin Chen 0002, David J. Wild 0001
Concurr. Comput. Pract. Exp.3
2011 Cloud computing paradigms for pleasingly parallel biomedical applications
abstract
SUMMARY Cloud computing offers exciting new approaches for scientific computing that leverage major commercial players’ hardware and software investments in large‐scale data centers. Loosely coupled problems are very important in many scientific fields, and with the ongoing move towards data‐intensive computing, they are on the rise. There exist several different approaches to leveraging clouds and cloud‐oriented data processing frameworks to perform pleasingly parallel (also called embarrassingly parallel) computations. In this paper, we present three pleasingly parallel biomedical applications: (i) assembly of genome fragments; (ii) sequence alignment and similarity search; and (iii) dimension reduction in the analysis of chemical structures, which are implemented utilizing a cloud infrastructure service‐based utility computing models of Amazon Web Services ( http://Amazon.com Inc., Seattle, WA, USA) and Microsoft Windows Azure (Microsoft Corp., Redmond, WA, USA) as well as utilizing MapReduce‐based data processing frameworks Apache Hadoop (Apache Software Foundation, Los Angeles, CA, USA) and Microsoft DryadLINQ. We review and compare each of these frameworks, performing a comparative study among them based on performance, cost, and usability. High latency, eventually consistent cloud infrastructure service‐based frameworks that rely on off‐the‐node cloud storage were able to exhibit performance efficiencies and scalability comparable to the MapReduce‐based frameworks with local disk‐based storage for the applications considered. In this paper, we also analyze variations in cost among the different platform choices (e.g., Elastic Compute Cloud instance types), highlighting the importance of selecting an appropriate platform based on the nature of the computation. Copyright © 2011 John Wiley & Sons, Ltd.
Thilina Gunarathne, Tak-Lon Wu, Jong Choi 0001, Seung-Hee Bae, Judy Qiu
Concurr. Comput. Pract. Exp.5
2011 Cloud Technologies for Bioinformatics Applications
abstract
Executing large number of independent jobs or jobs comprising of large number of tasks that perform minimal intertask communication is a common requirement in many domains. Various technologies ranging from classic job schedulers to the latest cloud technologies such as MapReduce can be used to execute these "many-tasks” in parallel. In this paper, we present our experience in applying two cloud technologies Apache Hadoop and Microsoft DryadLINQ to two bioinformatics applications with the above characteristics. The applications are a pairwise Alu sequence alignment application and an Expressed Sequence Tag (EST) sequence assembly program. First, we compare the performance of these cloud technologies using the above applications and also compare them with traditional MPI implementation in one application. Next, we analyze the effect of inhomogeneous data on the scheduling mechanisms of the cloud technologies. Finally, we present a comparison of performance of the cloud technologies under virtual and nonvirtual hardware platforms.
Jaliya Ekanayake, Thilina Gunarathne, Judy Qiu
IEEE Trans. Parallel Distributed Syst.3
2010 Performance of Windows Multicore Systems on Threading and MPI
abstract
We present performance results on a Windows cluster with up to 768 cores using MPI and two variants of threading - CCR and TPL. CCR (Concurrency and Coordination Runtime) presents a message based interface while TPL (Task Parallel Library) allows for loops to be automatically parallelized. MPI is used between the cluster nodes (up to 32) and either threading or MPI for parallelism on the 24 cores of each node. We use a simple matrix multiplication kernel as well as a significant bioinformatics gene clustering application. We find that the two threading models offer similar performance with MPI outperforming both at low levels of parallelism but threading much better when the grain size (problem size per process) is small. We find better performance on Intel compared to AMD on comparable 24 core systems. We develop simple models for the performance of the clustering code.
Judy Qiu, Scott Beason, Seung-Hee Bae, Saliya Ekanayake, Geoffrey C. Fox
CCGRID1
2010 MapReduce in the Clouds for Science
abstract
The utility computing model introduced by cloud computing combined with the rich set of cloud infrastructure services offers a very viable alternative to traditional servers and computing clusters. MapReduce distributed data processing architecture has become the weapon of choice for data-intensive analyses in the clouds and in commodity clusters due to its excellent fault tolerance features, scalability and the ease of use. Currently, there are several options for using MapReduce in cloud environments, such as using MapReduce as a service, setting up one's own MapReduce cluster on cloud instances, or using specialized cloud MapReduce runtimes that take advantage of cloud infrastructure services. In this paper, we introduce Azure MapReduce, a novel MapReduce runtime built using the Microsoft Azure cloud infrastructure services. Azure MapReduce architecture successfully leverages the high latency, eventually consistent, yet highly scalable Azure infrastructure services to provide an efficient, on demand alternative to traditional MapReduce clusters. Further we evaluate the use and performance of MapReduce frameworks, including Azure MapReduce, in cloud environments for scientific applications using sequence assembly and sequence alignment as use cases.
Thilina Gunarathne, Tak-Lon Wu, Judy Qiu, Geoffrey C. Fox
CloudCom3
2010 Applying Twister to Scientific Applications
abstract
Many scientific applications suffer from the lack of a unified approach to support the management and efficient processing of large-scale data. The Twister MapReduce Framework, which not only supports the traditional MapReduce programming model but also extends it by allowing iterations, addresses these problems. This paper describes how Twister is applied to several kinds of scientific applications such as BLAST, MDS Interpolation and GTM Interpolation in a non-iterative style and to MDS without interpolation in an iterative style. The results show the applicability of Twister to data parallel and EM algorithms with small overhead and increased efficiency.
Bingjing Zhang, Yang Ruan 0001, Tak-Lon Wu, Judy Qiu, Adam Hughes, Geoffrey C. Fox
CloudCom4
2010 Multidimensional Scaling by Deterministic Annealing with Iterative Majorization Algorithm
abstract
Multidimensional Scaling (MDS) is a dimension reduction method for information visualization, which is set up as a non-linear optimization problem. It is applicable to many data intensive scientific problems including studies of DNA sequences but tends to get trapped in local minima. Deterministic Annealing (DA) has been applied to many optimization problems to avoid local minima. We apply DA approach to MDS problem in this paper and show that our proposed DA approach improves the mapping quality and shows high reliability in a variety of experimental results. Further its execution time is similar to that of the un-annealed approach. We use different data sets for comparing the proposed DA approach with both a well known algorithm called SMACOF and a MDS with distance smoothing method which aims to avoid local optima. Our proposed DA method outperforms SMACOF algorithm and the distance smoothing MDS algorithm in terms of the mapping quality and shows much less sensitivity with respect to initial configurations and stopping condition. We also investigate various temperature cooling parameters for our deterministic annealing method within an exponential cooling scheme.
Seung-Hee Bae, Judy Qiu, Geoffrey C. Fox
eScience2
2010 Dimension reduction and visualization of large high-dimensional data via interpolation
abstract
The recent explosion of publicly available biology gene sequences and chemical compounds offers an unprecedented opportunity for data mining. To make data analysis feasible for such vast volume and high-dimensional scientific data, we apply high performance dimension reduction algorithms. It facilitates the investigation of unknown structures in a three dimensional visualization. Among the known dimension reduction algorithms, we utilize the multidimensional scaling and generative topographic mapping algorithms to configure the given high-dimensional data into the target dimension. However, both algorithms require large physical memory as well as computational resources. Thus, the authors propose an interpolated approach to utilizing the mapping of only a subset of the given data. This approach effectively reduces computational complexity. With minor trade-off of approximation, interpolation method makes it possible to process millions of data points with modest amounts of computation and memory requirement. Since huge amount of data are dealt, we represent how to parallelize proposed interpolation algorithms, as well. For the evaluation of the interpolated MDS by STRESS criteria, it is necessary to compute symmetric all pairwise computation with only subset of required data per process, so we also propose a simple but efficient parallel mechanism for the symmetric all pairwise computation when only a subset of data is available to each process. Our experimental results illustrate that the quality of interpolated mapping results are comparable to the mapping results of original algorithm only. In parallel performance aspect, those interpolation methods are well parallelized with high efficiency. With the proposed interpolation method, we construct a configuration of two-million out-of-sample data into the target dimension, and the number of out-of-sample data can be increased further.
Seung-Hee Bae, Jong Choi 0001, Judy Qiu, Geoffrey C. Fox
HPDC3
2010 Browsing large scale cheminformatics data with dimension reduction
abstract
Visualization of large-scale high dimensional data tool is highly valuable for scientific discovery in many fields. We presentPubChemBrowse, acustomizedvisualizationtoolfor cheminformatics research. It provides a novel 3D data point browser that displays complex properties of massive data on commodity clients. As in GIS browsers for Earth and Environment data, chemical compounds with similar properties are nearby in the browser. PubChemBrowse is built around in-househighperformanceparallel MDS(Multi-Dimensional Scaling) and GTM (Generative Topographic Mapping) services andsupports fast interaction with anexternalproperty database. These properties can be overlaid on 3D mapped compound space or queried for individual points. We prototype use with Chem2Bio2RDF system using SPARQLquery language to access over 20 publicly accessible bioinformatics databases. We describe our design and implementation of the integrated PubChemBrowse application and outline its use in drug discovery. The same core technologies can be used to develop similar high dimensional browsers in other scientific areas.
Jong Choi 0001, Seung-Hee Bae, Judy Qiu, Geoffrey C. Fox, Bin Chen 0002, David J. Wild 0001
HPDC3
2010 Twister: a runtime for iterative MapReduce
abstract
MapReduce programming model has simplified the implementation of many data parallel applications. The simplicity of the programming model and the quality of services provided by many implementations of MapReduce attract a lot of enthusiasm among distributed computing communities. From the years of experience in applying MapReduce to various scientific applications we identified a set of extensions to the programming model and improvements to its architecture that will expand the applicability of MapReduce to more classes of applications. In this paper, we present the programming model and the architecture of Twister an enhanced MapReduce runtime that supports iterative MapReduce computations efficiently. We also show performance comparisons of Twister with other similar runtimes such as Hadoop and DryadLINQ for large scale data parallel applications.
Jaliya Ekanayake, Bingjing Zhang, Thilina Gunarathne, Seung-Hee Bae, Judy Qiu, Geoffrey C. Fox
HPDC6
2010 Cloud computing paradigms for pleasingly parallel biomedical applications
abstract
Cloud computing offers exciting new approaches for scientific computing that leverages the hardware and software investments on large scale data centers by major commercial players. Loosely coupled problems are very important in many scientific fields and are on the rise with the ongoing move towards data intensive computing. There exist several approaches to leverage clouds & cloud oriented data processing frameworks to perform pleasingly parallel computations. In this paper we present two pleasingly parallel biomedical applications, 1) assembly of genome fragments 2) dimension reduction in the analysis of chemical structures, implemented utilizing cloud infrastructure service based utility computing models of Amazon AWS and Microsoft Windows Azure as well as utilizing MapReduce based data processing frameworks, Apache Hadoop and Microsoft DryadLINQ. We review and compare each of the frameworks and perform a comparative study among them based on performance, efficiency, cost and the usability. Cloud service based utility computing model and the managed parallelism (MapReduce) exhibited comparable performance and efficiencies for the applications we considered. We analyze the variations in cost between the different platform choices (eg: EC2 instance types), highlighting the need to select the appropriate platform based on the nature of the computation.
Thilina Gunarathne, Tak-Lon Wu, Judy Qiu, Geoffrey C. Fox
HPDC3
2010 Hybrid cloud and cluster computing paradigms for life science applications
abstract
BACKGROUND: Clouds and MapReduce have shown themselves to be a broadly useful approach to scientific computing especially for parallel data intensive applications. However they have limited applicability to some areas such as data mining because MapReduce has poor performance on problems with an iterative structure present in the linear algebra that underlies much data analysis. Such problems can be run efficiently on clusters using MPI leading to a hybrid cloud and cluster environment. This motivates the design and implementation of an open source Iterative MapReduce system Twister. RESULTS: Comparisons of Amazon, Azure, and traditional Linux and Windows environments on common applications have shown encouraging performance and usability comparisons in several important non iterative cases. These are linked to MPI applications for final stages of the data analysis. Further we have released the open source Twister Iterative MapReduce and benchmarked it against basic MapReduce (Hadoop) and MPI in information retrieval and life sciences applications. CONCLUSIONS: The hybrid cloud (MapReduce) and cluster (MPI) approach offers an attractive production environment while Twister promises a uniform programming environment for many Life Sciences applications. METHODS: We used commercial clouds Amazon and Azure and the NSF resource FutureGrid to perform detailed comparisons and evaluations of different approaches to data intensive computing. Several applications were developed in MPI, MapReduce and Twister in these different environments.
Judy Qiu, Jaliya Ekanayake, Thilina Gunarathne, Jong Choi 0001, Seung-Hee Bae, Bingjing Zhang, Tak-Lon Wu, Yang Ruan 0001, Saliya Ekanayake, Adam Hughes, Geoffrey C. Fox
BMC Bioinform.1