VLDB 2026 Research / reviewers in the wild / expert
Ümit V. Çatalyürek
dblp:c/UmitVCatalyurek
· DBLP profile ↗
137ranked-venue papers
11as first author
12since 2021 · last 2026
0000-0002-5625-3758ORCID · verified
Domains — the database's venue-derived domains; a paper can count in several
Systems, architecture and hardware · 97 · 8 first-author · 9 since 2021Databases, data management, data science and information retrieval · 18 · 1 first-author · 1 since 2021Applied, interdisciplinary, general and emerging computing · 17 · 2 first-authorArtificial intelligence and machine learning · 13 · 3 since 2021Graphics, computer vision, multimedia, augmented reality and games · 4 · 1 since 2021Human-computer interaction and ubiquitous computing · 4
| Year | Publication | Venue | Position |
|---|---|---|---|
| 2026 | Parallel Louvain Algorithms with Convergence Guarantee
Jason Niu, M. Yusuf Özkaya, Ahmet Erdem Sariyüce, Ümit V. Çatalyürek |
IPDPS | 4 |
| 2025 | A Scalable and Effective Alternative to Graph TransformersabstractGraph Neural Networks (GNNs) have shown impressive performance in graph representation learning, but they face challenges in capturing long-range dependencies due to their limited expressive power. To address this, Graph Transformers (GTs) were introduced, utilizing self-attention mechanism to effectively model pairwise node relationships. Despite their advantages, GTs suffer from quadratic complexity w.r.t. the number of nodes in the graph, hindering their applicability to large graphs. In this work, we present Graph-Enhanced Contextual Operator (GECO), a scalable and effective alternative to GTs that leverages neighborhood propagation and global convolutions to effectively capture local and global dependencies in quasiliniear time. Our study on synthetic datasets reveals that GECO reaches 169x speedup on a graph with 2M nodes w.r.t. optimized attention. Further evaluations on diverse range of benchmarks showcase that it scales to large graphs where traditional GTs often face memory and time limitations. Notably, GECO consistently achieves comparable or superior quality compared to baselines, improving the SOTA up to 4.5%, and offering a scalable and effective solution for large-scale graph learning. Kaan Sancak, Zhigang Hua, Andrey Malevich, Bo Long, Muhammad Fatih Balin, Ümit V. Çatalyürek |
AAAI | 8 |
| 2023 | Layer-Neighbor Sampling - Defusing Neighborhood Explosion in GNNsabstractGraph Neural Networks (GNNs) have received significant attention recently, but training them at a large scale remains a challenge.
Mini-batch training coupled with sampling is used to alleviate this challenge.
However, existing approaches either suffer from the neighborhood explosion phenomenon or have suboptimal performance.
To address these issues, we propose a new sampling algorithm called LAyer-neighBOR sampling (LABOR).
It is designed to be a direct replacement for Neighbor Sampling (NS) with the same fanout hyperparameter while sampling up to 7 times fewer vertices, without sacrificing quality.
By design, the variance of the estimator of each vertex matches NS from the point of view of a single vertex.
Moreover, under the same vertex sampling budget constraints, LABOR converges
faster than existing layer sampling approaches and can use up to 112 times larger batch sizes compared to NS. Muhammad Fatih Balin, Ümit V. Çatalyürek |
NeurIPS | 2 |
| 2022 | Efficient Hierarchical State Vector Simulation of Quantum Circuits via Acyclic Graph PartitioningabstractEarly but promising results in quantum computing have been enabled by the concurrent development of quan-tum algorithms, devices, and materials. Classical simulation of quantum programs has enabled the design and analysis of algorithms and implementation strategies targeting current and anticipated quantum device architectures. In this paper, we present a graph-based approach to achieving efficient quantum circuit simulation. Our approach involves partitioning the graph representation of a given quantum circuit into acyclic sub-graphs/circuits that exhibit better data locality. Simulation of each sub-circuit is organized hierarchically, with the iterative construction and simulation of smaller state vectors, improving overall performance. Also, this partitioning reduces the number of passes through data, improving the total computation time. We present three partitioning strategies and observe that acyclic graph partitioning typically results in the best time-to-solution. In contrast, other strategies reduce the partitioning time at the expense of potentially increased simulation times. Experimental evaluation demonstrates the effectiveness of our approach. Bo Fang 0002, M. Yusuf Özkaya, Ang Li 0006, Ümit V. Çatalyürek, Sriram Krishnamoorthy |
CLUSTER | 4 |
| 2022 | A Portable Sparse Solver Framework for Large Matrices on Heterogeneous ArchitecturesabstractProgramming applications on heterogeneous systems with hardware accelerators is challenging due to the disjoint address spaces between the host (CPU) and the device (GPU). The limited device memory further exacerbates the challenges as most data-intensive applications will not fit in the limited device memory. CUDA Unified Memory (UM) was introduced to mitigate such challenges. UM improves GPU programmability by supporting oversubscription, on-demand paging, and migration. However, when the working set of an application exceeds the device memory capacity, the resulting data movement can cause significant performance losses. We propose a tiling-based task-parallel framework, named DeepSparseGPU, to accelerate sparse eigensolvers on GPUs by minimizing data movement between the host and device. To this end, we tile all operations in a sparse solver and express the entire computation as a directed acyclic graph (DAG). We design and develop a memory manager (MM) to execute larger inputs that do not fit into GPU memory. MM keeps track of the data on CPU and GPU, and automatically moves data between them as needed. We use OpenMP target offload in our implementation to achieve portability beyond NVIDIA hardware. Performance evaluations show that DeepSparseGPU transfers 1.39x-2.18x less host to device (H2D) and device to host (D2H) data, while executing up to 2.93x faster than the UM-based baseline version. Fazlay Rabbi, Christopher S. Daley, Ümit V. Çatalyürek, Hasan Metin Aktulga |
HIPC | 3 |
| 2022 | MG-GCN: A Scalable multi-GPU GCN Training FrameworkabstractFull batch training of Graph Convolutional Network (GCN) models is not feasible on a single GPU for large graphs containing tens of millions of vertices or more. Recent work has shown that, for the graphs used in the machine learning community, communication becomes a bottleneck, and scaling is blocked outside of the single machine regime. Thus, we propose MG-GCN, a multi-GPU GCN training framework taking advantage of the high-speed communication links between the GPUs present in multi-GPU systems. MG-GCN employs multiple High-Performance Computing optimizations, including efficient re-use of memory buffers to reduce the memory footprint of training GNN models, as well as communication and computation overlap. These optimizations enable execution on larger datasets, that generally do not fit into the memory of a single GPU in state-of-the-art implementations. Furthermore, they contribute to achieving superior speedup compared to the state-of-the-art. For example, MG-GCN achieves super-linear speedup with respect to DGL, on the Reddit graph on both DGX-1 (V100) and DGX-A100. Muhammad Fatih Balin, Kaan Sancak, Ümit V. Çatalyürek |
ICPP | 3 |
| 2022 | A Block-Based Triangle Counting Algorithm on Heterogeneous EnvironmentsabstractTriangle counting is a fundamental building block in graph algorithms. In this article, we propose a block-based triangle counting algorithm to reduce data movement during both sequential and parallel execution. Our block-based formulation makes the algorithm naturally suitable for heterogeneous architectures. The problem of partitioning the adjacency matrix of a graph is well-studied. Our task decomposition goes one step further: it partitions the set of triangles in the graph. By streaming these small tasks to compute resources, we can solve problems that do not fit on a device. We demonstrate the effectiveness of our approach by providing an implementation on a compute node with multiple sockets, cores and GPUs. The current state-of-the-art in triangle enumeration processes the Friendster graph in 2.1 seconds, not including data copy time between CPU and GPU. Using that metric, our approach is 20 percent faster. When copy times are included, our algorithm takes 3.2 seconds. This is 5.6 times faster than the fastest published CPU-only time. Abdurrahman Yasar, Sivasankaran Rajamanickam, Jonathan W. Berry, Ümit V. Çatalyürek |
IEEE Trans. Parallel Distributed Syst. | 4 |
| 2021 | Parallel graph algorithms by blocks: from I/O to algorithmsabstractIn today's data-driven world and heterogeneous computing environments, processing large-scale graphs in an architecture agnostic manner has become more crucial than ever before. In terms of graph analytics frameworks, on the one side, there has been a significant interest in developing hand-optimized high-performance computing solutions. On the systems side, following the big data movement and to bring parallel computing to the masses, researchers have proposed several graph processing and management systems to handle large-scale graphs. Hand optimized HPC approaches require high expertise and are expensive to maintain and develop, and graph processing frameworks suffer from limited expressibility and performance. We propose Parallel Graph Algorithms by Blocks (PGAbB), a block-based graph algorithms framework for shared-memory, multi-core, multi-GPU machines. PGAbB offers a sweet spot between efficient parallelism and architecture agnostic algorithm design for a wide class of graph problems while performing close to hand-optimized HPC implementations. Abdurrahman Yasar, Kasimir Gabert, Ümit V. Çatalyürek |
CF | 3 |
| 2021 | Column-Segmented Sparse Matrix-Matrix Multiplication on Multicore CPUsabstractSparse general matrix-matrix multiplication, SpGEMM, is one of the most fundamental yet challenging sparse computation kernels. Due to its irregular computation pattern, SpGEMM frequently becomes the performance bottleneck in many scientific applications. Many prior state-of-the-art approaches use either dense or sparse accumulators to merge matrix rows as a critical component. Dense accumulators are efficient for small matrices but are infeasible for large or highly sparse matrices, due to high memory use and low cache efficiency. In this work, by segmenting the columns for the second input matrix, we propose a new SpGEMM algorithm that utilizes both a new sparse high-level overview of the matrix and fast and small dense accumulators that would fit in cache. With that, our approach brings the dense accumulator benefits to both large and highly sparse matrices. Our extensive experimental evaluation, carried out on three hardware platforms and on hundreds of sparse matrices from a variety of domains, shows that our algorithm out-performs state-of-the-art SpGEMM implementations. Xiaojing An, Ümit V. Çatalyürek |
HiPC | 2 |
| 2021 | An Evaluation of Task-Parallel Frameworks for Sparse Solvers on Multicore and Manycore CPU ArchitecturesabstractRecently, several task-parallel programming models have emerged to address the high synchronization and load imbalance issues as well as data movement overheads in modern shared memory architectures. OpenMP, the most commonly used shared memory parallel programming model, has added task execution support with dataflow dependencies. HPX and Regent are two more recent runtime systems that also support the dataflow execution model and extend it to distributed memory environments. We focus on parallelization of sparse matrix computations on shared memory architectures. We evaluate the OpenMP, HPX and Regent runtime systems in terms of performance and ease of implementation, and compare them against the traditional BSP model for two popular eigensolvers, Lanczos and LOBPCG. We give a general outline in regards to achieving parallelism using these runtime systems, and present a heuristic for tuning their performance to balance tasking overheads with the degree of parallelism that can be exposed. We then demonstrate their merits on two architectures, Intel Broadwell (a multicore processor) and AMD EPYC (a modern manycore processor). We observe that these frameworks achieve up to 13.7 × fewer cache misses over an efficient BSP implementation across L1, L2 and L3 cache layers. They also obtain up to 9.9 × improvement in execution time over the same BSP implementation. Abdullah Alperen, Md. Afibuzzaman, Fazlay Rabbi, M. Yusuf Özkaya, Ümit V. Çatalyürek, Hasan Metin Aktulga |
ICPP | 5 |
| 2021 | EIGA: elastic and scalable dynamic graph analysisabstractModern graphs are not only large, but rapidly changing. The rate of change can vary significantly along with the computational cost. Existing distributed graph analysis systems have largely been designed to operate on static graphs. Infrastructure changes in these systems need to occur when the system is idle, which can result in significant wasted resources or the inability to cope with changes. Kasimir Gabert, Kaan Sancak, M. Yusuf Özkaya, Ali Pinar, Ümit V. Çatalyürek |
SC | 5 |
| 2021 | A Unifying Framework to Identify Dense Subgraphs on Streams: Graph Nuclei to Hypergraph CoresabstractFinding dense regions of graphs is fundamental in graph mining. We focus on the computation of dense hierarchies and regions with graph nuclei---a generalization of k-cores and trusses. Static computation of nuclei, namely through variants of 'peeling', are easy to understand and implement. However, many practically important graphs undergo continuous change. Dynamic algorithms, maintaining nucleus computations on dynamic graph streams, are nuanced and require significant effort to port between nuclei, e.g., from k-cores to trusses. Kasimir Gabert, Ali Pinar, Ümit V. Çatalyürek |
WSDM | 3 |
| 2020 | StrainHub: a phylogenetic tool to construct pathogen transmission networksabstractSUMMARY: In exploring the epidemiology of infectious diseases, networks have been used to reconstruct contacts among individuals and/or populations. Summarizing networks using pathogen metadata (e.g. host species and place of isolation) and a phylogenetic tree is a nascent, alternative approach. In this paper, we introduce a tool for reconstructing transmission networks in arbitrary space from phylogenetic information and metadata. Our goals are to provide a means of deriving new insights and infection control strategies based on the dynamics of the pathogen lineages derived from networks and centrality metrics. We created a web-based application, called StrainHub, in which a user can input a phylogenetic tree based on genetic or other data along with characters derived from metadata using their preferred tree search method. StrainHub generates a transmission network based on character state changes in metadata, such as place or source of isolation, mapped on the phylogenetic tree. The user has the option to calculate centrality metrics on the nodes including betweenness, closeness, degree and a new metric, the source/hub ratio. The outputs include the network with values for metrics on its nodes and the tree with characters reconstructed. All of these results can be exported for further analysis. AVAILABILITY AND IMPLEMENTATION: strainhub.io and https://github.com/abschneider/StrainHub. Adriano de Bernardi Schneider, Colby T. Ford, Reilly Hostager, Michael Cioce, Ümit V. Çatalyürek, Joel O. Wertheim, Daniel Janies |
Bioinform. | 6 |
| 2019 | DeepSparse: A Task-Parallel Framework for SparseSolvers on Deep Memory ArchitecturesabstractData movement is an important bottleneck against efficiency and energy consumption in large-scale sparse matrix computations that are commonly used in linear solvers, eigensolvers and graph analytics. We introduce a novel task-parallel sparse solver framework, named DeepSparse, which adopts a fully integrated task-parallel approach. DeepSparse framework differs from existing work in that it adopts a holistic approach that targets all computational steps in a sparse solver rather than narrowing the problem into small kernels (e.g., SpMM, SpMV). We present the implementation details of DeepSparse and demonstrate its merit in two popular eigensolvers, LOBPCG and Lanczos algorithms. We observe that DeepSparse achieves 2× - 16× fewer cache misses across different cache layers (L1, L2 and L3) over implementations of the same solvers based on optimized library function calls. We also achieve 2× - 3.9× improvement in execution time when using DeepSparse over the same library versions. Md. Afibuzzaman, Fazlay Rabbi, M. Yusuf Özkaya, Hasan Metin Aktulga, Ümit V. Çatalyürek |
HiPC | 5 |
| 2019 | Efficient and effective sparse tensor reorderingabstractThis paper formalizes the problem of reordering a sparse tensor to improve the spatial and temporal locality of operations with it, and proposes two reordering algorithms for this problem, which we call BFS-MCS and Lexi-Order. The BFS-MCS method is a Breadth First Search (BFS)-like heuristic approach based on the maximum cardinality search family; Lexi-Order is an extension of doubly lexical ordering of matrices to tensors. We show the effects of these schemes within the context of a widely used tensor computation, the CANDECOMP/PARAFAC decomposition (CPD), when storing the tensor in three previously proposed sparse tensor formats: coordinate (COO), compressed sparse fiber (CSF), and hierarchical coordinate (HiCOO). A new partition-based superblock scheduling is also proposed for HiCOO format to improve load balance. On modern multicore CPUs, we show Lexi-Order obtains up to 4.14× speedup on sequential HiCOO-Mttkrp and 11.88× speedup on its parallel counterpart. The performance of COO- and CSF-based Mttkrps also improves. Our two reordering methods are more effective than state-of-the-art approaches. The code is released as part of Parallel Tensor Infrastructure (ParTI!): https://github.com/hpcgarage/ParTI. Jiajia Li 0001, Bora Uçar, Ümit V. Çatalyürek, Jimeng Sun 0001, Kevin J. Barker, Richard W. Vuduc |
ICS | 3 |
| 2019 | A Scalable Clustering-Based Task Scheduler for Homogeneous Processors Using DAG PartitioningabstractWhen scheduling a directed acyclic graph (DAG) of tasks with communication costs on computational platforms, a good trade-off between load balance and data locality is necessary. List-based scheduling techniques are commonly-used greedy approaches for this problem. The downside of list-scheduling heuristics is that they are incapable of making short term sacrifices for the global efficiency of the schedule. In this work, we describe new list-based scheduling heuristics based on clustering for homogeneous platforms, under the realistic duplex single-port communication model. Our approach uses an acyclic partitioner for DAGs for clustering. The clustering enhances the data locality of the scheduler with a global view of the graph. Furthermore, since the partition is acyclic, we can schedule each part completely once its input tasks are ready to be executed. We present an extensive experimental evaluation showing the tradeoffs between the granularity of clustering and the parallelism, and how this affects the scheduling. Furthermore, we compare our heuristics to the best state-of-the-art list-scheduling and clustering heuristics, and obtain more than three times better makespan in cases with many communications. M. Yusuf Özkaya, Anne Benoit, Bora Uçar, Julien Herrmann, Ümit V. Çatalyürek |
IPDPS | 5 |
| 2019 | Special issue: Selected papers from IPDPS'18
Anne Benoit, Ümit V. Çatalyürek |
J. Parallel Distributed Comput. | 2 |
| 2019 | Geometric Mapping of Tasks to Processors on Parallel Computers with Mesh or Torus NetworksabstractWe present a new method for reducing parallel applications' communication time by mapping their MPI tasks to processors in a way that lowers the distance messages travel and the amount of congestion in the network. Assuming geometric proximity among the tasks is a good approximation of their communication interdependence, we use a geometric partitioning algorithm to order both the tasks and the processors, assigning task parts to the corresponding processor parts. In this way, interdependent tasks are assigned to “nearby” cores in the network. We also present a number of algorithmic optimizations that exploit specific features of the network or application to further improve the quality of the mapping. We specifically address the case of sparse node allocation, where the nodes assigned to a job are not necessarily located in a contiguous block nor within close proximity to each other in the network. However, our methods generalize to contiguous allocations as well, and results are shown for both contiguous and non-contiguous allocations. We show that, for the structured finite difference mini-application MiniGhost, our mapping methods reduced communication time up to 75 percent relative to MiniGhost's default mapping on 128K cores of a Cray XK7 with sparse allocation. For the atmospheric modeling code E3SM/HOMME, our methods reduced communication time up to 31% on 16K cores of an IBM BlueGene/Q with contiguous allocation. Mehmet Deveci, Karen D. Devine, Kevin T. Pedretti, Mark A. Taylor, Sivasankaran Rajamanickam, Ümit V. Çatalyürek |
IEEE Trans. Parallel Distributed Syst. | 6 |
| 2018 | Local Detection of Critical Nodes in Active GraphsabstractThe identification of critical nodes in a graph is a fundamental task in network analysis. Centrality measures are commonly used for this purpose. These methods rely on two assumptions that restrict their applicability. First, they only depend on the topology of the network and do not consider the activity over the network. Second, they assume the entire network is available. However, in many applications, it is the underlying activity of the network such as interactions and communications that makes a node critical, and it is hard to collect the entire network topology, when the network is vast and autonomous. We propose a new measure, Active Betweenness Cardinality, where the importance of the nodes are based not on the static structure, but the active utilization of the network. We show how this metric can be computed efficiently by only local information for a given node and how we can locate the critical nodes by using only a few nodes. We also show how this metric can be used to monitor a network and identify node failures. We evaluate our metric and algorithms on real-world networks and show the effectiveness of the proposed methods. M. Yusuf Özkaya, Ahmet Erdem Sanyuce, Ali Pinar, Ümit V. Çatalyürek |
ASONAM | 4 |
| 2018 | SiNA: A Scalable Iterative Network AlignerabstractGiven two graphs, network alignment asks for a potentially partial mapping between the vertices of the two graphs. This arises in many applications where data from different sources need to be integrated. Recent graph aligners use the global structure of input graphs and additional information given for the edges and vertices. We present SINA, an efficient, shared memory parallel implementation of such an aligner. Our experimental evaluations on a 32-core shared memory machine showed that SINA scales well for aligning large real-world graphs: SINA can achieve up to 28.5× speedup, and can reduce the total execution time of a graph alignment problem with 2M vertices and 100M edges from 4.5 hours to under 10 minutes. To the best of our knowledge, SINA is the first parallel aligner that uses global structure and vertex and edge attributes to handle large graphs. Abdurrahman Yasar, Bora Uçar, Ümit V. Çatalyürek |
ASONAM | 3 |
| 2018 | Efficient sparse-matrix multi-vector product on GPUsabstractSparse Matrix-Vector (SpMV) and Sparse Matrix-Multivector (SpMM) products are key kernels for computational science and data science. While GPUs offer significantly higher peak performance and memory bandwidth than multicore CPUs, achieving high performance on sparse computations on GPUs is very challenging. A tremendous amount of recent research has focused on various GPU implementations of the SpMV kernel. But the multi-vector SpMM kernel has received much less attention. In this paper, we present an in-depth analysis to contrast SpMV and SpMM, and develop a new sparse-matrix representation and computation approach suited to achieving high data-movement efficiency and effective GPU parallelization of SpMM. Experimental evaluation using the entire SuiteSparse matrix suite demonstrates significant performance improvement over existing SpMM implementations from vendor libraries. Changwan Hong, Aravind Sukumaran-Rajam, Bortik Bandyopadhyay, Süreyya Emre Kurt, Israt Nisa, Shivani Sabhlok, Ümit V. Çatalyürek, Srinivasan Parthasarathy 0001, P. Sadayappan |
HPDC | 8 |
| 2018 | An Iterative Global Structure-Assisted Labeled Network AlignerabstractIntegrating data from heterogeneous sources is often modeled as merging graphs. Given two or more "compatible'', but not-isomorphic graphs, the first step is to identify a graph alignment, where a potentially partial mapping of vertices between two graphs is computed. A significant portion of the literature on this problem only takes the global structure of the input graphs into account. Only more recent ones additionally use vertex and edge attributes to achieve a more accurate alignment. However, these methods are not designed to scale to map large graphs arising in many modern applications. We propose a new iterative graph aligner, gsaNA, that uses the global structure of the graphs to significantly reduce the problem size and align large graphs with a minimal loss of information. Concretely, we show that our proposed technique is highly flexible, can be used to achieve higher recall, and it is orders of magnitudes faster than the current state of the art techniques. Abdurrahman Yasar, Ümit V. Çatalyürek |
KDD | 2 |
| 2018 | A resource provisioning framework for bioinformatics applications in multi-cloud environments
Izzet F. Senturk, Ponnuraman Balakrishnan, Anas Abu-Doleh, Kamer Kaya, Qutaibah M. Malluhi, Ümit V. Çatalyürek |
Future Gener. Comput. Syst. | 6 |
| 2017 | Acyclic Partitioning of Large Directed Acyclic GraphsabstractFinding a good partition of a computational directed acyclic graph associated with an algorithm can help find an execution pattern improving data locality, conduct an analysis of data movement, and expose parallel steps. The partition is required to be acyclic, i.e., the inter-part edges between the vertices from different parts should preserve an acyclic dependency structure among the parts. In this work, we adopt the multilevel approach with coarsening, initial partitioning, and refinement phases for acyclic partitioning of directed acyclic graphs and develop a direct k-way partitioning scheme. To the best of our knowledge, no such scheme exists in the literature. To ensure the acyclicity of the partition at all times, we propose novel and efficient coarsening and refinement heuristics. The quality of the computed acyclic partitions is assessed by computing the edge cut, the total volume of communication between the parts, and the critical path latencies. We use the solution returned by well-known undirected graph partitioners as a baseline to evaluate our acyclic partitioner, knowing that the space of solution is more restricted in our problem. The experiments are run on large graphs arising from linear algebra applications. Julien Herrmann, Jonathan Kho, Bora Uçar, Kamer Kaya, Ümit V. Çatalyürek |
CCGrid | 5 |
| 2017 | Guest Editor's Introduction: Selected Papers from ACM-BCB 2014abstractThe papers in this special issue were presented at the 5th ACM Conference on Bioinformatics, Computational Biology, and Health Informatics, held in Newport Beach, CA in September 2014, The papers address the use of computational modeling in the biological and health research fields. With the new high throughput devices, such as sequencers and imaging devices, and ubiquitous sensor technologies, the landscape of how we do biomedical research; how the knowledge is curated; and results are delivered to stake holders are constantly changing. Ümit V. Çatalyürek |
IEEE ACM Trans. Comput. Biol. Bioinform. | 1 |
| 2017 | Graph Manipulations for Fast Centrality ComputationabstractThe betweenness and closeness metrics are widely used metrics in many network analysis applications. Yet, they are expensive to compute. For that reason, making the betweenness and closeness centrality computations faster is an important and well-studied problem. In this work, we propose the framework BADIOS that manipulates the graph by compressing it and splitting into pieces so that the centrality computation can be handled independently for each piece. Experimental results show that the proposed techniques can be a great arsenal to reduce the centrality computation time for various types and sizes of networks. In particular, it reduces the betweenness centrality computation time of a 4.6 million edges graph from more than 5 days to less than 16 hours. For the same graph, the closeness computation time is decreased from more than 3 days to 6 hours (12.7x speedup). Ahmet Erdem Sariyüce, Kamer Kaya, Erik Saule, Ümit V. Çatalyürek |
ACM Trans. Knowl. Discov. Data | 4 |
| 2017 | Nucleus Decompositions for Identifying Hierarchy of Dense SubgraphsabstractFinding dense substructures in a graph is a fundamental graph mining operation, with applications in bioinformatics, social networks, and visualization to name a few. Yet most standard formulations of this problem (like clique, quasi-clique, densest at-least- k subgraph) are NP-hard. Furthermore, the goal is rarely to find the “true optimum” but to identify many (if not all) dense substructures, understand their distribution in the graph, and ideally determine relationships among them. Current dense subgraph finding algorithms usually optimize some objective and only find a few such subgraphs without providing any structural relations. We define the nucleus decomposition of a graph, which represents the graph as a forest of nuclei . Each nucleus is a subgraph where smaller cliques are present in many larger cliques. The forest of nuclei is a hierarchy by containment, where the edge density increases as we proceed towards leaf nuclei. Sibling nuclei can have limited intersections, which enables discovering overlapping dense subgraphs. With the right parameters, the nucleus decomposition generalizes the classic notions of k -core and k -truss decompositions. We present practical algorithms for nucleus decompositions and empirically evaluate their behavior in a variety of real graphs. The tree of nuclei consistently gives a global, hierarchical snapshot of dense substructures and outputs dense subgraphs of comparable quality with the state-of-the-art solutions that are dense and have non-trivial sizes. Our algorithms can process real-world graphs with tens of millions of edges in less than an hour. We demonstrate how proposed algorithms can be utilized on a citation network. Our analysis showed that dense units identified by our algorithms correspond to coherent articles on a specific area. Our experiments also show that we can identify dense structures that are lost within larger structures by other methods and find further finer grain structure within dense groups. Ahmet Erdem Sariyüce, Seshadhri Comandur, Ali Pinar, Ümit V. Çatalyürek |
ACM Trans. Web | 4 |
| 2016 | SONIC: streaming overlapping community detection
Ahmet Erdem Sariyüce, Bugra Gedik, Gabriela Jacques-Silva, Kun-Lung Wu, Ümit V. Çatalyürek |
Data Min. Knowl. Discov. | 5 |
| 2016 | Multi-Jagged: A Scalable Parallel Spatial Partitioning AlgorithmabstractGeometric partitioning is fast and effective for load-balancing dynamic applications, particularly those requiring geometric locality of data (particle methods, crash simulations). We present, to our knowledge, the first parallel implementation of a multidimensional-jagged geometric partitioner. In contrast to the traditional recursive coordinate bisection algorithm (RCB), which recursively bisects subdomains perpendicular to their longest dimension until the desired number of parts is obtained, our algorithm does recursive multi-section with a given number of parts in each dimension. By computing multiple cut lines concurrently and intelligently deciding when to migrate data while computing the partition, we minimize data movement compared to efficient implementations of recursive bisection. We demonstrate the algorithm's scalability and quality relative to the RCB implementation in Zoltan on both real and synthetic datasets. Our experiments show that the proposed algorithm performs and scales better than RCB in terms of run-time without degrading the load balance. Our implementation partitions 24 billion points into 65,536 parts within a few seconds and exhibits near perfect weak scaling up to 6K cores. Mehmet Deveci, Sivasankaran Rajamanickam, Karen D. Devine, Ümit V. Çatalyürek |
IEEE Trans. Parallel Distributed Syst. | 4 |
| 2016 | Incremental k-core decomposition: algorithms and evaluation
Ahmet Erdem Sariyüce, Bugra Gedik, Gabriela Jacques-Silva, Kun-Lung Wu, Ümit V. Çatalyürek |
VLDB J. | 5 |
| 2015 | A Novel Multiple Choice Question Generation Strategy: Alternative Uses for Controlled Vocabulary Thesauri in Biomedical-Sciences Education
Marcelo A. Lopetegui, Barbara A. Lara, Po-Yin Yen, Ümit V. Çatalyürek, Philip R. O. Payne |
AMIA | 4 |
| 2015 | Spaler: Spark and GraphX based de novo genome assemblerabstractThe recent advancements in high-throughput genome sequencing technologies have accelerated the efficient discovery of novel genomes. De novo assembly is the first and one of the most computationally intensive step to analyze such novel genomes. In this work, we addressed the problem of parallelizing the de Bruijn graph based de novo genome sequence assembly on distributed memory systems. We proposed a new tool, Spaler, which assembles short reads efficiently and accurately. Spaler is based on Spark framework and GraphX API. We compared the performance of Spaler to other distributed memory based assemblers, in particular, ABySS, Ray and SWAP-Assembler. The results show that Spaler scales better than existing tools and produces comparable or better results in terms of solution quality. Anas Abu-Doleh, Ümit V. Çatalyürek |
IEEE BigData | 2 |
| 2015 | Fast and High Quality Topology-Aware Task MappingabstractConsidering the large number of processors and the size of the interconnection networks on exactable-capable supercomputers, mapping concurrently executable and communicating tasks of an application is complex problem that needs to be dealt with care. For parallel applications, the communication overhead can be a significant bottleneck on scalability. Topology-aware task-mapping methods that map the tasks tithe processors~(i.e., cores) by exploiting the underlying network information are very effective to avoid, or at worst bend, this limitation. We propose novel, efficient, and effective task mapping algorithms employing a graph model. The experiments show that the methods are faster than the existing approaches proposed for the same task, and on 4096 processors, the algorithms improve the communication hops and link contentions by 16% and 32%, respectively, on the average. In addition, they improve the average execution time of a parallel Spiv kernel and a communication-only application by 9% and 14%, respectively. Mehmet Deveci, Kamer Kaya, Bora Uçar, Ümit V. Çatalyürek |
IPDPS | 4 |
| 2015 | Finding the Hierarchy of Dense Subgraphs using Nucleus DecompositionsabstractFinding dense substructures in a graph is a fundamental graph mining operation, with applications in bioinformatics, social networks, and visualization to name a few. Yet most standard formulations of this problem (like clique, quasiclique, k-densest subgraph) are NP-hard. Furthermore, the goal is rarely to find the "true optimum", but to identify many (if not all) dense substructures, understand their distribution in the graph, and ideally determine relationships among them. Current dense subgraph finding algorithms usually optimize some objective, and only find a few such subgraphs without providing any structural relations. We define the nucleus decomposition of a graph, which represents the graph as a forest of nuclei. Each nucleus is a subgraph where smaller cliques are present in many larger cliques. The forest of nuclei is a hierarchy by containment, where the edge density increases as we proceed towards leaf nuclei. Sibling nuclei can have limited intersections, which enables discovering overlapping dense subgraphs. With the right parameters, the nucleus decomposition generalizes the classic notions of k-cores and k-truss decompositions. We give provably efficient algorithms for nucleus decompositions, and empirically evaluate their behavior in a variety of real graphs. The tree of nuclei consistently gives a global, hierarchical snapshot of dense substructures, and outputs dense subgraphs of higher quality than other state-of-the-art solutions. Our algorithm can process graphs with tens of millions of edges in less than an hour. Ahmet Erdem Sariyüce, Seshadhri Comandur, Ali Pinar, Ümit V. Çatalyürek |
WWW | 4 |
| 2015 | Hypergraph partitioning for multiple communication cost metrics: Model and methods
Mehmet Deveci, Kamer Kaya, Bora Uçar, Ümit V. Çatalyürek |
J. Parallel Distributed Comput. | 4 |
| 2015 | Regularizing graph centrality computations
Ahmet Erdem Sariyüce, Erik Saule, Kamer Kaya, Ümit V. Çatalyürek |
J. Parallel Distributed Comput. | 4 |
| 2015 | Incremental closeness centrality in distributed memory
Ahmet Erdem Sariyüce, Erik Saule, Kamer Kaya, Ümit V. Çatalyürek |
Parallel Comput. | 4 |
| 2014 | Exploiting Geometric Partitioning in Task Mapping for Parallel ComputersabstractWe present a new method for mapping applications' MPI tasks to cores of a parallel computer such that communication and execution time are reduced. We consider the case of sparse node allocation within a parallel machine, where the nodes assigned to a job are not necessarily located within a contiguous block nor within close proximity to each other in the network. The goal is to assign tasks to cores so that interdependent tasks are performed by "nearby" cores, thus lowering the distance messages must travel, the amount of congestion in the network, and the overall cost of communication. Our new method applies a geometric partitioning algorithm to both the tasks and the processors, and assigns task parts to the corresponding processor parts. We show that, for the structured finite difference mini-app Mini Ghost, our mapping method reduced execution time 34% on average on 65,536 cores of a Cray XE6. In a molecular dynamics mini-app, Mini MD, our mapping method reduced communication time by 26% on average on 6144 cores. We also compare our mapping with graph-based mappings from the LibTopoMap library and show that our mappings reduced the communication time on average by 15% in MiniGhost and 10% in MiniMD. Mehmet Deveci, Sivasankaran Rajamanickam, Vitus J. Leung, Kevin T. Pedretti, Stephen Olivier, David P. Bunde, Ümit V. Çatalyürek, Karen D. Devine |
IPDPS | 7 |
| 2014 | mrSNP: Software to detect SNP effects on microRNA bindingabstractBACKGROUND: MicroRNAs (miRNAs) are short (19-23 nucleotides) non-coding RNAs that bind to sites in the 3'untranslated regions (3'UTR) of a targeted messenger RNA (mRNA). Binding leads to degradation of the transcript or blocked translation resulting in decreased expression of the targeted gene. Single nucleotide polymorphisms (SNPs) have been found in 3'UTRs that disrupt normal miRNA binding or introduce new binding sites and some of these have been associated with disease pathogenesis. This raises the importance of detecting miRNA targets and predicting the possible effects of SNPs on binding sites. In the last decade a number of studies have been conducted to predict the location of miRNA binding sites. However, there have been fewer algorithms published to analyze the effects of SNPs on miRNA binding. Moreover, the existing software has some shortcomings including the requirement for significant manual labor when working with huge lists of SNPs and that algorithms work only for SNPs present in databases such as dbSNP. These limitations become problematic as next-generation sequencing is leading to large numbers of novel variants in 3'UTRs. RESULT: In order to overcome these issues, we developed a web-server named mrSNP which predicts the impact of a SNP in a 3'UTR on miRNA binding. The proposed tool reduces the manual labor requirements and allows users to input any SNP that has been identified by any SNP-calling program. In testing the performance of mrSNP on SNPs experimentally validated to affect miRNA binding, mrSNP correctly identified 69% (11/16) of the SNPs disrupting binding. CONCLUSIONS: mrSNP is a highly adaptable and performing tool for predicting the effect a 3'UTR SNP will have on miRNA binding. This tool has advantages over existing algorithms because it can assess the effect of novel SNPs on miRNA binding without requiring significant hands on time. Mehmet Deveci, Ümit V. Çatalyürek, Amanda E. Toland |
BMC Bioinform. | 2 |
| 2014 | Massively multithreaded maxflow for image segmentation on the Cray XMT-2abstractSUMMARY Image segmentation is a very important step in the computerized analysis of digital images. The maxflow mincut approach has been successfully used to obtain minimum energy segmentations of images in many fields. Classical algorithms for maxflow in networks do not directly lend themselves to efficient parallel implementations on contemporary parallel processors. We present the results of an implementation of Goldberg–Tarjan preflow‐push algorithm on the Cray XMT‐2 massively multithreaded supercomputer. This machine has hardware support for 128 threads in each physical processor, a uniformly accessible shared memory of up to 4 TB and hardware synchronization for each 64 bit word. It is thus well‐suited to the parallelization of graph theoretic algorithms, such as preflow‐push. We describe the implementation of the preflow‐push code on the XMT‐2 and present the results of timing experiments on a series of synthetically generated as well as real images. Our results indicate very good performance on large images and pave the way for practical applications of this machine architecture for image analysis in a production setting. The largest images we have run are 32000 2 pixels in size, which are well beyond the largest previously reported in the literature.Copyright © 2013 John Wiley & Sons, Ltd. Shahid H. Bokhari, Ümit V. Çatalyürek, Metin Nafi Gürcan |
Concurr. Comput. Pract. Exp. | 2 |
| 2014 | Diversifying Citation RecommendationsabstractLiterature search is one of the most important steps of academic research. With more than 100,000 papers published each year just in computer science, performing a complete literature search becomes a Herculean task. Some of the existing approaches and tools for literature search cannot compete with the characteristics of today’s literature, and they suffer from ambiguity and homonymy. Techniques based on citation information are more robust to the mentioned issues. Thus, we recently built a Web service called the advisor, which provides personalized recommendations to researchers based on their papers of interest. Since most recommendation methods may return redundant results, diversifying the results of the search process is necessary to increase the amount of information that one can reach via an automated search. This article targets the problem of result diversification in citation-based bibliographic search, assuming that the citation graph itself is the only information available and no categories or intents are known. The contribution of this work is threefold. We survey various random walk--based diversification methods and enhance them with the direction awareness property to allow users to reach either old, foundational (possibly well-cited and well-known) research papers or recent (most likely less-known) ones. Next, we propose a set of novel algorithms based on vertex selection and query refinement. A set of experiments with various evaluation criteria shows that the proposed γ-RLM algorithm performs better than the existing approaches and is suitable for real-time bibliographic search in practice. Onur Küçüktunç, Erik Saule, Kamer Kaya, Ümit V. Çatalyürek |
ACM Trans. Intell. Syst. Technol. | 4 |
| 2013 | Towards a personalized, scalable, and exploratory academic recommendation serviceabstractLiterature search is an integral part of the academic research. Academic recommendation services have been developed to help researchers with their literature search, many of which only provide a text-based search functionality. Such services are suitable for a first-level bibliographic search; however, they lack the benefits of today's recommendation engines. In this paper, we identify three important properties that an academic recommendation service could provide for better literature search: personalization, scalability, and exploratory search. With these objectives in mind, we present a web service called theadvisor which helps the users build a strong bibliography by extending the document set obtained after a first-level search. Along with an efficient and personalized recommendation algorithm, the service also features result diversification, relevance feedback, visualization for exploratory search. We explain the design criteria and rationale we employed to make the theadvisor a useful and scalable web service with a thorough evaluation. Onur Küçüktunç, Erik Saule, Kamer Kaya, Ümit V. Çatalyürek |
ASONAM | 4 |
| 2013 | Incremental algorithms for closeness centralityabstractCentrality metrics have shown to be highly correlated with the importance and loads of the nodes within the network traffic. In this work, we provide fast incremental algorithms for closeness centrality computation. Our algorithms efficiently compute the closeness centrality values upon changes in network topology, i.e., edge insertions and deletions. We show that the proposed techniques are efficient on many real-life networks, especially on small-world networks, which have a small diameter and spike-shaped shortest distance distribution. We experimentally validate the efficiency of our algorithms on large-scale networks and show that they can update the closeness centrality values of 1.2 million authors in the temporal DBLP-coauthorship network 460 times faster than it would take to recompute them from scratch. Ahmet Erdem Sariyüce, Kamer Kaya, Erik Saule, Ümit V. Çatalyürek |
IEEE BigData | 4 |
| 2013 | STREAMER: A distributed framework for incremental closeness centrality computationabstractNetworks are commonly used to model the traffic patterns, social interactions, or web pages. The nodes in a network do not possess the same characteristics: some nodes are naturally more connected and some nodes can be more important. Closeness centrality (CC) is a global metric that quantifies how important is a given node in the network. When the network is dynamic and keeps changing, the relative importance of the nodes also changes. The best known algorithm to compute the CC scores makes it impractical to recompute them from scratch after each modification. In this paper, we propose Streamer, a distributed memory framework for incrementally maintaining the closeness centrality scores of a network upon changes. It leverages pipelined and replicated parallelism and takes NUMA effects into account. It speeds up the maintenance of the CC of a real graph with 916K vertices and 4.3M edges by a factor of 497 using a 64 nodes cluster. Ahmet Erdem Sariyüce, Erik Saule, Kamer Kaya, Ümit V. Çatalyürek |
CLUSTER | 4 |
| 2013 | GPU Accelerated Maximum Cardinality Matching Algorithms for Bipartite Graphs
Mehmet Deveci, Kamer Kaya, Bora Uçar, Ümit V. Çatalyürek |
Euro-Par | 4 |
| 2013 | Hypergraph Sparsification and Its Application to PartitioningabstractThe data one needs to cope to solve today's problems is large scale, so are the graphs and hyper graphs used to model it. Today, we have Big Data, big graphs, big matrices, and in the future, they are expected to be bigger and more complex. Many of today's algorithms will be, and some already are, expensive to run on large datasets. In this work, we analyze a set of efficient techniques to make "big data", which is modeled as a hyper graph, smaller so that its processing takes much less time. As an application use case, we take the hyper graph partitioning problem, which has been successfully used in many practical applications for various purposes including parallelization of complex and irregular applications, sparse matrix ordering, clustering, community detection, query optimization, and improving cache locality in shared-memory systems. We conduct several experiments to show that our techniques greatly reduce the cost of the partitioning process and preserve the partitioning quality. Although we only measured their performance from the partitioning point of view, we believe the proposed techniques will be beneficial also for other applications using hyper graphs. Mehmet Deveci, Kamer Kaya, Ümit V. Çatalyürek |
ICPP | 3 |
| 2013 | A Push-Relabel-Based Maximum Cardinality Bipartite Matching Algorithm on GPUsabstractWe design, develop, and evaluate an atomic- and lock-free GPU implementation of the push-relabel algorithm in the context of finding maximum cardinality matchings in bipartite graphs. The problem has applications on computer science, scientific computing, bioinformatics, and other areas. Although the GPU parallelization of the push-relabel technique has been investigated in the context of flow algorithms, to the best of our knowledge, ours is the first study which focuses on the maximum cardinality matching. We compare the proposed algorithms with serial, multicore, and many core bipartite graph matching implementations from the literature on a large set of real-life problems. Our experiments show that the proposed pushrelabel-based GPU algorithm is faster than the existing parallel and sequential implementations. Mehmet Deveci, Kamer Kaya, Bora Uçar, Ümit V. Çatalyürek |
ICPP | 4 |
| 2013 | Exploring the future of out-of-core computing with compute-local non-volatile memoryabstractDrawing parallels to the rise of general purpose graphical processing units (GPGPUs) as accelerators for specific high-performance computing (HPC) workloads, there is a rise in the use of non-volatile memory (NVM) as accelerators for I/O-intensive scientific applications. However, existing works have explored use of NVM within dedicated I/O nodes, which are distant from the compute nodes that actually need such acceleration. As NVM bandwidth begins to out-pace point-to-point network capacity, we argue for the need to break from the archetype of completely separated storage. Myoungsoo Jung, Ellis Herbert Wilson, Wonil Choi, John Shalf, Hasan Metin Aktulga, Chao Yang 0001, Erik Saule, Ümit V. Çatalyürek, Mahmut T. Kandemir |
SC | 8 |
| 2013 | Shattering and Compressing Networks for Betweenness CentralityabstractThe betweenness metric has always been intriguing and used in many analyses.Yet, it is one of the most computationally expensive kernels in graph mining.For that reason, making betweenness centrality computations faster is an important and well-studied problem.In this work, we propose the framework, BADIOS, which compresses a network and shatters it into pieces so that the centrality computation can be handled independently for each piece.Although BADIOS is designed and tuned for betweenness centrality, it can easily be adapted for other centrality metrics.Experimental results show that the proposed techniques can be a great arsenal to reduce the centrality computation time for various types and sizes of networks.In particular, it reduces the computation time of a 4.6 million edges graph from more than 5 days to less than 16 hours. Ümit V. Çatalyürek, Kamer Kaya, Ahmet Erdem Sariyüce, Erik Saule |
SDM | 1 |
| 2013 | Diversified recommendation on graphs: pitfalls, measures, and algorithmsabstractResult diversification has gained a lot of attention as a way to answer ambiguous queries and to tackle the redundancy problem in the results. In the last decade, diversification has been applied on or integrated into the process of PageRank- or eigenvector-based methods that run on various graphs, including social networks, collaboration networks in academia, web and product co-purchasing graphs. For these applications, the diversification problem is usually addressed as a bicriteria objective optimization problem of relevance and diversity. However, such an approach is questionable since a query-oblivious diversification algorithm that recommends most of its results without even considering the query may perform the best on these commonly used measures. In this paper, we show the deficiencies of popular evaluation techniques of diversification methods, and investigate multiple relevance and diversity measures to understand whether they have any correlations. Next, we propose a novel measure called expanded relevance which combines both relevance and diversity into a single function in order to measure the coverage of the relevant part of the graph. We also present a new greedy diversification algorithm called BestCoverage, which optimizes the expanded relevance of the result set with (1-1/e)-approximation. With a rigorous experimentation on graphs from various applications, we show that the proposed method is efficient and effective for many use cases. Onur Küçüktunç, Erik Saule, Kamer Kaya, Ümit V. Çatalyürek |
WWW | 4 |
| 2013 | A comparative analysis of biclustering algorithms for gene expression dataabstractThe need to analyze high-dimension biological data is driving the development of new data mining methods. Biclustering algorithms have been successfully applied to gene expression data to discover local patterns, in which a subset of genes exhibit similar expression levels over a subset of conditions. However, it is not clear which algorithms are best suited for this task. Many algorithms have been published in the past decade, most of which have been compared only to a small number of algorithms. Surveys and comparisons exist in the literature, but because of the large number and variety of biclustering algorithms, they are quickly outdated. In this article we partially address this problem of evaluating the strengths and weaknesses of existing biclustering methods. We used the BiBench package to compare 12 algorithms, many of which were recently published or have not been extensively studied. The algorithms were tested on a suite of synthetic data sets to measure their performance on data with varying conditions, such as different bicluster models, varying noise, varying numbers of biclusters and overlapping biclusters. The algorithms were also tested on eight large gene expression data sets obtained from the Gene Expression Omnibus. Gene Ontology enrichment analysis was performed on the resulting biclusters, and the best enrichment terms are reported. Our analyses show that the biclustering method and its parameters should be selected based on the desired model, whether that model allows overlapping biclusters, and its robustness to noise. In addition, we observe that the biclustering algorithms capable of finding more than one model are more successful at capturing biologically relevant clusters. Kemal Eren, Mehmet Deveci, Onur Küçüktunç, Ümit V. Çatalyürek |
Briefings Bioinform. | 4 |
| 2013 | Benchmarking short sequence mapping toolsabstractBACKGROUND: The development of next-generation sequencing instruments has led to the generation of millions of short sequences in a single run. The process of aligning these reads to a reference genome is time consuming and demands the development of fast and accurate alignment tools. However, the current proposed tools make different compromises between the accuracy and the speed of mapping. Moreover, many important aspects are overlooked while comparing the performance of a newly developed tool to the state of the art. Therefore, there is a need for an objective evaluation method that covers all the aspects. In this work, we introduce a benchmarking suite to extensively analyze sequencing tools with respect to various aspects and provide an objective comparison. RESULTS: We applied our benchmarking tests on 9 well known mapping tools, namely, Bowtie, Bowtie2, BWA, SOAP2, MAQ, RMAP, GSNAP, Novoalign, and mrsFAST (mrFAST) using synthetic data and real RNA-Seq data. MAQ and RMAP are based on building hash tables for the reads, whereas the remaining tools are based on indexing the reference genome. The benchmarking tests reveal the strengths and weaknesses of each tool. The results show that no single tool outperforms all others in all metrics. However, Bowtie maintained the best throughput for most of the tests while BWA performed better for longer read lengths. The benchmarking tests are not restricted to the mentioned tools and can be further applied to others. CONCLUSION: The mapping process is still a hard problem that is affected by many factors. In this work, we provided a benchmarking suite that reveals and evaluates the different factors affecting the mapping process. Still, there is no tool that outperforms all of the others in all the tests. Therefore, the end user should clearly specify his needs in order to choose the tool that provides the best results. Ayat Hatem, Doruk Bozdag, Amanda E. Toland, Ümit V. Çatalyürek |
BMC Bioinform. | 4 |
| 2013 | Streaming Algorithms for k-core DecompositionabstractA k -core of a graph is a maximal connected subgraph in which every vertex is connected to at least k vertices in the subgraph. k -core decomposition is often used in large-scale network analysis, such as community detection, protein function prediction, visualization, and solving NP-Hard problems on real networks efficiently, like maximal clique finding. In many real-world applications, networks change over time. As a result, it is essential to develop efficient incremental algorithms for streaming graph data. In this paper, we propose the first incremental k -core decomposition algorithms for streaming graph data. These algorithms locate a small subgraph that is guaranteed to contain the list of vertices whose maximum k -core values have to be updated, and efficiently process this subgraph to update the k -core decomposition. Our results show a significant reduction in run-time compared to non-incremental alternatives. We show the efficiency of our algorithms on different types of real and synthetic graphs, at different scales. For a graph of 16 million vertices, we observe speedups reaching a million times, relative to the non-incremental algorithms. Ahmet Erdem Sariyüce, Bugra Gedik, Gabriela Jacques-Silva, Kun-Lung Wu, Ümit V. Çatalyürek |
Proc. VLDB Endow. | 5 |
| 2012 | Fast Recommendation on Bibliographic NetworksabstractGraphs and matrices are widely used in algorithms for social network analyses. Since the number of interactions is much less than the possible number of interactions, the graphs and matrices used in the analyses are usually sparse. In this paper, we propose an efficient implementation of a sparse-matrix computation which arises in our publicly available citation recommendation service called the advisor. The recommendation algorithm uses a sparse matrix generated from the citation graph. We observed that the nonzero pattern of this matrix is highly irregular and the computation suffers from high number of cache misses. We propose techniques for storing the matrix in memory efficiently and reducing the number of cache misses. Experimental results show that our techniques are highly efficient on reducing the query processing time which is highly crucial for a web service. Onur Küçüktunç, Kamer Kaya, Erik Saule, Ümit V. Çatalyürek |
ASONAM | 4 |
| 2012 | An Out-of-Core Eigensolver on SSD-equipped ClustersabstractObtaining highly accurate predictions on properties of light atomic nuclei using the Configuration Interaction (CI)approach requires computing few extremal eigenpairs of a large many-body nuclear Hamiltonian matrix, Ĥ. A forefront challenge in CI calculations is the massive size of Ĥ and its eigenvectors. The emergence of clusters equipped with non-volatile NAND-flash memory based solid state drives (SSD) presents unique opportunities. In this paper, we present the implementation details of an out-of-core eigensolver using a novel distributed out-of-core linear algebra framework, called DOoC+LAF. The framework provides an easy-to-use high-level application interface for linear algebra operations while providing efficient execution by orchestrating pipelined execution of computation, communication and I/O. We demonstrate the effectiveness of our out-of-core eigensolver implemented using DOoC+LAF by reporting performance results on large-scale eigenvalue problems arising in nuclear structure calculations. Zheng Zhou 0003, Erik Saule, Hasan Metin Aktulga, Chao Yang 0001, Esmond G. Ng, Pieter Maris, James P. Vary, Ümit V. Çatalyürek |
CLUSTER | 8 |
| 2012 | On Shared-Memory Parallelization of a Sparse Matrix Scaling AlgorithmabstractWe discuss efficient shared memory parallelization of sparse matrix computations whose main traits resemble to those of the sparse matrix-vector multiply operation. Such computations are difficult to parallelize because of the relatively small computational granularity characterized by small number of operations per each data access. Our main application is a sparse matrix scaling algorithm which is more memory bound than the sparse matrix vector multiplication operation. We take the application and parallelize it using the standard OpenMP programming principles. Apart from the common race condition avoiding constructs, we do not reorganize the algorithm. Rather, we identify associated performance metrics and describe models to optimize them. By using these models, we implement parallel matrix scaling algorithms for two well-known sparse matrix storage formats. Experimental results show that simple parallelization attempts which leave data/work partitioning to the runtime scheduler can suffer from the overhead of avoiding race conditions especially when the number of threads increases. The proposed algorithms perform better than these algorithms by optimizing the identified performance metrics and reducing the overhead. Ümit V. Çatalyürek, Kamer Kaya, Bora Uçar |
ICPP | 1 |
| 2012 | Multithreaded Clustering for Multi-level Hypergraph PartitioningabstractRequirements for efficient parallelization of many complex and irregular applications can be cast as a hyper graph partitioning problem. The current-state-of-the art software libraries that provide tool support for the hyper graph partitioning problem are designed and implemented before the game-changing advancements in multi-core computing. Hence, analyzing the structure of those tools for designing multithreaded versions of the algorithms is a crucial tasks. The most successful partitioning tools are based on the multi-level approach. In this approach, a given hyper graph is coarsened to a much smaller one, a partition is obtained on the the smallest hyper graph, and that partition is projected to the original hyper graph while refining it on the intermediate hyper graphs. The coarsening operation corresponds to clustering the vertices of a hyper graph and is the most time consuming task in a multi-level partitioning tool. We present three efficient multithreaded clustering algorithms which are very suited for multi-level partitioners. We compare their performance with that of the ones currently used in today's hyper graph partitioners. We show on a large number of real life hyper graphs that our implementations, integrated into a commonly used partitioning library PaToH, achieve good speedups without reducing the clustering quality. Ümit V. Çatalyürek, Mehmet Deveci, Kamer Kaya, Bora Uçar |
IPDPS | 1 |
| 2012 | Optimizing the stretch of independent tasks on a cluster: From sequential tasks to moldable tasks
Erik Saule, Doruk Bozdag, Ümit V. Çatalyürek |
J. Parallel Distributed Comput. | 3 |
| 2012 | Load-balancing spatially located computations using rectangular partitions
Erik Saule, Erdeniz Ö. Bas, Ümit V. Çatalyürek |
J. Parallel Distributed Comput. | 3 |
| 2012 | Graph coloring algorithms for multi-core and massively multithreaded architectures
Ümit V. Çatalyürek, John Feo, Assefaw Hadish Gebremedhin, Mahantesh Halappanavar, Alex Pothen |
Parallel Comput. | 1 |
| 2012 | Improving performance of adaptive component-based dataflow middleware
Timothy D. R. Hartley, Erik Saule, Ümit V. Çatalyürek |
Parallel Comput. | 3 |
| 2011 | Benchmarking Short Sequence Mapping ToolsabstractThe development of next-generation sequencing instruments has lead to the generation of millions of short sequences in a single run. The process of aligning these reads to a reference genome is time consuming and demands the development of fast and accurate alignment tools. However, the current proposed tools make different compromises between the accuracy and the speed of mapping. Moreover, many important aspects are overlooked when comparing the performance of a newly developed tool to the state of the art. Therefore, there is a need for an objective evaluation method that covers the various aspects. In this work, we introduce a benchmarking suite to extensively analyze various tools with respect to the different comparison aspects and provide an objective comparison. In order to assess our work, we applied our benchmarking tests on seven well known mapping tools, namely, Bowtie, BWA, SOAP, MAQ, RMAP, GSNAP, and FANGS. Bowtie, BWA, SOAP, GSNAP, and FANGS are based on indexing the reference genome, whereas MAQ and RMAP are based on building hash tables for the reads. It is shown that the benchmarking tests reveal the strengths and weaknesses of each tool. In addition, the tests can be further applied to other tools. The results show that there is no clear winner. However, Bowtie maintained the best throughput for most of the tests while BWA performed better for longer read lengths. Ayat Hatem, Doruk Bozdag, Ümit V. Çatalyürek |
BIBM | 3 |
| 2011 | Improving graph coloring on distributed-memory parallel computersabstractGraph coloring is a combinatorial optimization problem that classically appears in distributed computing to identify the sets of tasks that can be safely performed in parallel. Despite many existing efficient sequential algorithms being known for this NP-Complete problem, distributed variants are challenging. Building on an existing distributed-memory graph coloring framework, we investigate two techniques in this paper. First, we investigate the application of two different vertex-visit orderings, namely Largest First and Smallest Last, in a distributed context and show that they can help to significantly decrease the number of colors, on small-to medium-scale parallel architectures. Second, we investigate the use of a distributed post-processing operation, called recoloring, which further drastically improves the number of colors while not increasing the runtime more than twofold on large graphs. We also investigate the use of multicore architectures for distributed graph coloring algorithms. Ahmet Erdem Sariyüce, Erik Saule, Ümit V. Çatalyürek |
HiPC | 3 |
| 2011 | Partitioning Spatially Located Computations Using RectanglesabstractThe ideal distribution of spatially located heterogeneous workloads is an important problem to address in parallel scientific computing. We investigate the problem of partitioning such workloads (represented as a matrix of positive integers) into rectangles, such that the load of the most loaded rectangle (processor) is minimized. Since finding the optimal arbitrary rectangle-based partition is an NP-hard problem, we investigate particular classes of solutions, namely, rectilinear partitions, jagged partitions and hierarchical partitions. We present a new class of solutions called m-way jagged partitions, propose new optimal algorithms for m-way jagged partitions and hierarchical partitions, propose new heuristic algorithms, and provide worst case performance analyses for some existing and new heuristics. Moreover, the algorithms are tested in simulation on a wide set of instances. Results show that two of the algorithms we introduce lead to a much better load balance than the state-of-the-art algorithms. Erik Saule, Erdeniz Ö. Bas, Ümit V. Çatalyürek |
IPDPS | 3 |
| 2011 | Optimizing latency and throughput of application workflows on clusters
Nagavijayalakshmi Vydyanathan, Ümit V. Çatalyürek, Tahsin M. Kurç, P. Sadayappan, Joel H. Saltz |
Parallel Comput. | 2 |
| 2010 | Automatic dataflow application tuning for heterogeneous systemsabstractDue to the increasing prevalence of multicore microprocessors and accelerator technologies in modern supercomputer design, new techniques for designing scientific applications are needed, in order to efficiently leverage all of the power inherent in these systems. The dataflow programming paradigm is well-suited to application design for distributed and heterogeneous systems than other techniques. Traditionally in dataflow middleware, application data domains are statically partitioned and distributed among the processors using a demand-driven algorithm. Unfortunately, this task scheduling technique can cause severe load imbalances in heterogeneous environments. Furthermore, in the presence of different types of processors, the optimum datasize can be different for each processor type. To solve the load imbalance problem and to leverage the optimum datasize dynamicity in a dataflow framework, we present an algorithm which automatically partitions the application workspace. By putting this partitioning into the purview of the dataflow runtime system, we can adaptively change the size of databuffers and correctly balance the load. Experiments with four applications show that our technique allows developers to skip the tedious and error-prone step of manually tuning the data granularity. Our technique is always competitive with the best-known data partitioning for these experiments, and can beat it under certain constraints. Timothy D. R. Hartley, Erik Saule, Ümit V. Çatalyürek |
HiPC | 3 |
| 2010 | Run-time optimizations for replicated dataflows on heterogeneous environmentsabstractThe increases in multi-core processor parallelism and in the flexibility of many-core accelerator processors, such as GPUs, have turned traditional SMP systems into hierarchical, heterogeneous computing environments. Fully exploiting these improvements in parallel system design remains an open problem. Moreover, most of the current tools for the development of parallel applications for hierarchical systems concentrate on the use of only a single processor type (e.g., accelerators) and do not coordinate several heterogeneous processors. Here, we show that making use of all of the heterogeneous computing resources can significantly improve application performance. Our approach, which consists of optimizing applications at run-time by efficiently coordinating application task execution on all available processing units is evaluated in the context of replicated dataflow applications. The proposed techniques were developed and implemented in an integrated run-time system targeting both intra- and inter-node parallelism. The experimental results with a real-world complex biomedical application show that our approach nearly doubles the performance of the GPU-only implementation on a distributed heterogeneous accelerator cluster. George Teodoro, Timothy D. R. Hartley, Ümit V. Çatalyürek, Renato Ferreira 0001 |
HPDC | 3 |
| 2010 | An Image Analysis Approach for Detecting Malignant Cells in Digitized H&E-stained Histology Images of Follicular LymphomaabstractThe gold standard in follicular lymphoma (FL) diagnosis and prognosis is histopathological examination of tumor tissue samples. However, the qualitative manual evaluation is tedious and subject to considerable inter- and intra-reader variations. In this study, we propose an image analysis system for quantitative evaluation of digitized FL tissue slides. The developed system uses a robust feature space analysis method, namely the mean shift algorithm followed by a hierarchical grouping to segment a given tissue image into basic cytological components. We then apply further morphological operations to achieve the segmentation of individual cells. Finally, we generate a likelihood measure to detect candidate cancer cells using a set of clinically driven features. The proposed approach has been evaluated on a dataset consisting of 100 region of interest (ROI) images and achieves a promising 89% average accuracy in detecting target malignant cells. Olcay Sertel, Ümit V. Çatalyürek, Gerard Lozanski, Arwa Shanaah, Metin Nafi Gürcan |
ICPR | 2 |
| 2010 | A Moldable Online Scheduling Algorithm and Its Application to Parallel Short Sequence Mapping
Erik Saule, Doruk Bozdag, Ümit V. Çatalyürek |
JSSPP | 3 |
| 2010 | On the Scalability of Hypergraph Models for Sparse Matrix PartitioningabstractWe investigate the scalability of the hypergraph-based sparse matrix partitioning methods with respect to the increasing sizes of matrices and number of nonzeros. We propose a method to rowwise partition the matrices that correspond to the discretization of two-dimensional domains with the five-point stencil. The proposed method obtains perfect load balance and achieves very good total communication volume. We investigate the behaviour of the hypergraph-based rowwise partitioning method with respect to the proposed method, in an attempt to understand how scalable the former method is. In another set of experiments, we work on general sparse matrices under different scenarios to understand the scalability of various hypergraph-based one- and two-dimensional matrix partitioning methods. Bora Uçar, Ümit V. Çatalyürek |
PDP | 2 |
| 2010 | A Matrix Partitioning Interface to PaToH in MATLAB
Bora Uçar, Ümit V. Çatalyürek, Cevdet Aykanat |
Parallel Comput. | 2 |
| 2009 | Investigating the use of GPU-accelerated nodes for SAR image formationabstractThe computation of an electromagnetic reflectivity image from a set of radar returns is a computationally intensive process. Therefore, the use of high performance computing is required to form images from radar signals in a short time frame. This paper explores the use of distributed memory cluster computers and accelerator technologies such as GPUs for radar signal analysis applications, particularly backprojection image formation. We obtain good results with the use of GPUs and compare their performance in terms of execution time with distributed memory cluster computers. When using a configuration with 4 GPU-accelerated nodes, we achieve speedups up to 3.45x for different input and output data size combinations. Timothy D. R. Hartley, Ahmed Fasih, Charles A. Berdanier, Füsun Özgüner, Ümit V. Çatalyürek |
CLUSTER | 5 |
| 2009 | Coordinating the use of GPU and CPU for improving performance of compute intensive applicationsabstractGPUs have recently evolved into very fast parallel co-processors capable of executing general purpose computations extremely efficiently. At the same time, multi-core CPUs evolution continued and today's CPUs have 4-8 cores. These two trends, however, have followed independent paths in the sense that we are aware of very few works that consider both devices cooperating to solve general computations. In this paper we investigate the coordinated use of CPU and GPU to improve efficiency of applications even further than using either device independently. We use Anthill runtime environment, a data-flow oriented framework in which applications are decomposed into a set of event-driven filters, where for each event, the runtime system can use either GPU or CPU for its processing. For evaluation, we use a histopathology application that uses image analysis techniques to classify tumor images for neuroblas-toma prognosis. Our experimental environment includes dual and octa-core machines, augmented with GPUs and we evaluate our approach's performance for standalone and distributed executions. Our experiments show that a pure GPU optimization of the application achieved a factor of 15 to 49 times improvement over the single core CPU version, depending on the versions of the CPUs and GPUs. We also show that the execution can be further reduced by a factor of about 2 by using our runtime system that effectively choreographs the execution to run cooperatively both on GPU and on a single core of CPU. We improve on that by adding more cores, all of which were previously neglected or used ineffectively. In addition, the evaluation on a distributed environment has shown near linear scalability to multiple hosts. George Teodoro, Rafael Sachetto Oliveira, Olcay Sertel, Metin Nafi Gürcan, Wagner Meira Jr., Ümit V. Çatalyürek, Renato Ferreira 0001 |
CLUSTER | 6 |
| 2009 | High-performance signal processing on emerging many-core architectures using cudaabstractThis paper provides a short introduction to CUDA programming paradigm and recent standardization efforts. We also present a new 2D fast wavelet transform implementation on GPUs using CUDA. Our novel implementation achieves almost two orders of magnitude improvements versus a typical quad-core CPU, and demonstrates that emerging many-core architectures are ideal platforms for achieving high-performance signal processing. Manuel Ujaldon, Ümit V. Çatalyürek |
ICME | 2 |
| 2009 | Parallel short sequence mapping for high throughput genome sequencingabstractWith the advent of next-generation high throughput sequencing instruments, large volumes of short sequence data are generated at an unprecedented rate. Processing and analyzing these massive data requires overcoming several challenges including mapping of generated short sequences to a reference genome. This computationally intensive process takes time on the order of days using existing sequential techniques on large scale datasets. In this work, we propose six parallelization methods to speedup short sequence mapping and to reduce the execution time under just a few hours for such large datasets. We comparatively present these methods and give theoretical cost models for each method. Experimental results on real datasets demonstrate the effectiveness of the parallel methods and indicate that the cost models help accurate estimation of parallel execution time. Based on these cost models we implemented a selection function to predict the best method for a given scenario. To the best of our knowledge this is the first study on parallelization of short sequence mapping problem. Doruk Bozdag, Catalin C. Barbacioru, Ümit V. Çatalyürek |
IPDPS | 3 |
| 2009 | A component-based framework for the Cell Broadband EngineabstractWith the increasing trend of microprocessor manufacturers to rely on parallelism to increase their products' performance, there is an associated increasing need for simple techniques to leverage this hardware parallelism for good application performance. Unfortunately, many application developers do not have the benefit of long experience in programming parallel and distributed systems. While the filter-stream programming paradigm helps bridge the gap between developers of scientific applications and the performance they need, current and future high-performance multicore processor designs do not have a filter-stream programming library available. This work aims to fill that gap in the software world. This initial DataCutter-Lite implementation defines a powerful, but simple abstraction for carrying out complex computations in a filter-stream model. Additionally, the initial implementation shows that complex architectures such as the Cell Broadband Engine Architecture can make use of the filter-stream model, and give good application performance when doing so. Timothy D. R. Hartley, Ümit V. Çatalyürek |
IPDPS | 2 |
| 2009 | A repartitioning hypergraph model for dynamic load balancing
Ümit V. Çatalyürek, Erik G. Boman, Karen D. Devine, Doruk Bozdag, Robert T. Heaphy, Lee Ann Riesen |
J. Parallel Distributed Comput. | 1 |
| 2009 | Computer-aided prognosis of neuroblastoma on whole-slide images: Classification of stromal development
Olcay Sertel, Jun Kong 0002, Hiroyuki Shimada, Ümit V. Çatalyürek, Joel H. Saltz, Metin Nafi Gürcan |
Pattern Recognit. | 4 |
| 2009 | Compaction of Schedules and a Two-Stage Approach for Duplication-Based DAG SchedulingabstractMany DAG scheduling algorithms generate schedules that require prohibitively large number of processors. To address this problem, we propose a generic algorithm, SC, to minimize the processor requirement of any given valid schedule. SC preserves the schedule length of the original schedule and reduces processor count by merging processor schedules and removing redundant duplicate tasks. To the best of our knowledge, this is the first algorithm to address this highly unexplored aspect of DAG scheduling. On average, SC reduced the processor requirement 91, 82, and 72 percent for schedules generated by PLW, TCSD, and CPFD algorithms, respectively. SC algorithm has a low complexity (O(\vert {\cal N}\vert^3)) compared to most duplication-based algorithms. Moreover, it decouples processor economization from schedule length minimization problem. To take advantage of these features of SC, we also propose a scheduling algorithm SDS, having the same time complexity as SC. Our experiments demonstrate that schedules generated by SDS are only 3 percent longer than CPFD (O(\vert {\cal N}\vert^4)), one of the best algorithms in that respect. SDS and SC together form a two-stage scheduling algorithm that produces schedules with high quality and low processor requirement, and has lower complexity than the comparable algorithms that produce similar high-quality results. Doruk Bozdag, Füsun Özgüner, Ümit V. Çatalyürek |
IEEE Trans. Parallel Distributed Syst. | 3 |
| 2009 | An Integrated Approach to Locality-Conscious Processor Allocation and Scheduling of Mixed-Parallel ApplicationsabstractComplex parallel applications can often be modeled as directed acyclic graphs of coarse-grained application tasks with dependences. These applications exhibit both task and data parallelism, and combining these two (also called mixed parallelism) has been shown to be an effective model for their execution. In this paper, we present an algorithm to compute the appropriate mix of task and data parallelism required to minimize the parallel completion time (makespan) of these applications. In other words, our algorithm determines the set of tasks that should be run concurrently and the number of processors to be allocated to each task. The processor allocation and scheduling decisions are made in an integrated manner and are based on several factors such as the structure of the task graph, the runtime estimates and scalability characteristics of the tasks, and the intertask data communication volumes. A locality-conscious scheduling strategy is used to improve intertask data reuse. Evaluation through simulations and actual executions of task graphs derived from real applications and synthetic graphs shows that our algorithm consistently generates schedules with a lower makespan as compared to Critical Path Reduction (CPR) and Critical Path and Allocation (CPA), two previously proposed scheduling algorithms. Our algorithm also produces schedules that have a lower makespan than pure task- and data-parallel schedules. For task graphs with known optimal schedules or lower bounds on the makespan, our algorithm generates schedules that are closer to the optima than other scheduling approaches. Nagavijayalakshmi Vydyanathan, Sriram Krishnamoorthy, Gerald Sabin, Ümit V. Çatalyürek, Tahsin M. Kurç, P. Sadayappan, Joel H. Saltz |
IEEE Trans. Parallel Distributed Syst. | 4 |
| 2008 | Multi-hop path splitting and multi-pathing optimizations for data transfers over shared wide-area networks using gridFTPabstractIn this paper, we propose to employ two optimizations - multi-hop path splitting and multi-pathing - to improve the performance of data transfers over shared public networks. We present a path determination algorithm which integrates the aforesaid optimizations in order to improve the performance of single file transfers. Finally, we develop a file transfer scheduling algorithm based on this framework, and evaluate its effectiveness on a wide-area testbed. Gaurav Khanna 0002, Ümit V. Çatalyürek, Tahsin M. Kurç, P. Sadayappan, Joel H. Saltz, Rajkumar Kettimuthu, Ian T. Foster |
HPDC | 2 |
| 2008 | Texture classification using nonlinear color quantization: Application to histopathological image analysisabstractIn this paper, a novel color texture classification approach is introduced and applied to computer-assisted grading of follicular lymphoma from whole-slide tissue samples. The digitized tissue samples of follicular lymphoma were classified into histological grades under a statistical framework. The proposed method classifies the image either into low or high grades based on the amount of cytological components. To further discriminate the lower grades into low and mid grades, we proposed a novel color texture analysis approach. This approach modifies the gray level co-occurrence matrix method by using a non-linear color quantization with self-organizing feature maps (SOFMs). This is particularly useful for the analysis of H&E stained pathological images whose dynamic color range is considerably limited. Experimental results on real follicular lymphoma images demonstrate that the proposed approach outperforms the gray level based texture analysis. Olcay Sertel, Jun Kong 0002, Gerard Lozanski, Arwa Shanaah, Ümit V. Çatalyürek, Joel H. Saltz, Metin Nafi Gürcan |
ICASSP | 5 |
| 2008 | A Duplication Based Algorithm for Optimizing Latency Under Throughput Constraints for Streaming WorkflowsabstractScheduling, in many application domains, involves the optimization of multiple performance metrics. For example, application workflows with real-time constraints have strict throughput requirements and also desire a low latency or response time. In this paper, we present a novel algorithm for the scheduling of workflows that act on a stream of input data. Our algorithm focuses on the two performance metrics: latency and throughput, and minimizes the latency of workflows while satisfying strict throughput requirements. We leverage pipelined, task and data parallelism in a coordinated manner to meet these objectives and investigate the benefit of task duplication in alleviating communication overheads in the pipelined schedule for different workflow characteristics. The proposed algorithm is designed for a realistic k-port communication model, where each processor can simultaneously communicate with at most k distinct processors. Evaluation using synthetic and application benchmarks shows that our algorithm consistently produces lower-latency schedules and meets throughput requirements, even when previously proposed schemes fail. Nagavijayalakshmi Vydyanathan, Ümit V. Çatalyürek, Tahsin M. Kurç, P. Sadayappan, Joel H. Saltz |
ICPP | 2 |
| 2008 | Biomedical image analysis on a cooperative cluster of GPUs and multicoresabstractWe are currently witnessing the emergence of two paradigms in parallel computing: streaming processing and multi-core CPUs. Represented by solid commercial products widely available in commodity PCs, GPUs and multi-core CPUs bring together an unprecedented combination of high performance at low cost. The scientific computing community needs to keep pace with application models and middleware which scale efficiently to hundreds of internal processing units. The purpose of the work we present here is twofold: first, a cooperative environment is designed so that both parallel models can coexist and complement one another. Second, beyond the parallelism of multiple internal cores, further parallelism is introduced when multiple CPU sockets, multiple GPUs, and multiple nodes are combined within a unique multi-processor platform which exceeds 10 TFLOPS when using 16 nodes. We illustrate our cooperative parallelization approach by implementing a large-scale, biomedical image analysis application which contains a number of assorted kernels including typical streaming operators, co-occurrence matrices, convolutions, and histograms. Experimental results are compared among different implementation strategies and almost linear speed-up is achieved when all coexisting methods in CPUs and GPUs are combined. Timothy D. R. Hartley, Ümit V. Çatalyürek, Antonio Ruiz 0001, Francisco D. Igual, Rafael Mayo 0002, Manuel Ujaldon |
ICS | 2 |
| 2008 | A dynamic scheduling approach for coordinated wide-area data transfers using GridFTPabstractMany scientific applications need to stage large volumes of files from one set of machines to another set of machines in a wide-area network. Efficient execution of such data transfers needs to take into account the heterogeneous nature of the environment and dynamic availability of shared resources. This paper proposes an algorithm that dynamically schedules a batch of data transfer requests with the goal of minimizing the overall transfer time. The proposed algorithm performs simultaneous transfer of chunks of files from multiple file replicas, if the replicas exist. Adaptive replica selection is employed to transfer different chunks of the same file by taking dynamically changing network bandwidths into account. We utilize GridFTP as the underlying mechanism for data transfers. The algorithm makes use of information from past GridFTP transfers to estimate network bandwidths and resource availability. The efficiency of the algorithm is evaluated on a wide-area testbed. Gaurav Khanna 0002, Ümit V. Çatalyürek, Tahsin M. Kurç, Rajkumar Kettimuthu, P. Sadayappan, Joel H. Saltz |
IPDPS | 2 |
| 2008 | Translational research design templates, Grid computing, and HPCabstractDesign templates that involve discovery, analysis, and integration of information resources commonly occur in many scientific research projects. In this paper we present examples of design templates from the biomedical translational research domain and discuss the requirements imposed on Grid middleware infrastructures by them. Using caGrid, which is a Grid middleware system based on the model driven architecture (MDA) and the service oriented architecture (SOA) paradigms, as a starting point, we discuss architecture directions for MDA and SOA based systems like caGrid to support common design templates. Joel H. Saltz, Scott Oster, Shannon Hastings, Stephen Langella, Renato Ferreira 0001, Justin Permar, Ashish Sharma 0001, David Ervin, Tony Pan, Ümit V. Çatalyürek, Tahsin M. Kurç |
IPDPS | 10 |
| 2008 | Using overlays for efficient data transfer over shared wide-area networksabstractData-intensive applications frequently transfer large amounts of data over wide-area networks. The performance achieved in such settings can often be improved by routing data via intermediate nodes chosen to increase aggregate bandwidth. We explore the benefits of overlay network approaches by designing and implementing a service-oriented architecture that incorporates two key optimizations - multi-hop path splitting andmulti-pathing - within the GridFTP file transfer protocol. We develop a file transfer scheduling algorithm that incorporates the two optimizations in conjunction with the use of available file replicas. The algorithm makes use of information from past GridFTP transfers to estimate network bandwidths and resource availability. The effectiveness of these optimizations is evaluated using several application file transfer patterns: one-to-all broadcast, all-to-one gather, and data redistribution, on a wide-area testbed. The experimental results show that our architecture and algorithm achieve significant performance improvement. Gaurav Khanna 0002, Ümit V. Çatalyürek, Tahsin M. Kurç, Rajkumar Kettimuthu, P. Sadayappan, Ian T. Foster, Joel H. Saltz |
SC | 2 |
| 2008 | A framework for scalable greedy coloring on distributed-memory parallel computers
Doruk Bozdag, Assefaw Hadish Gebremedhin, Fredrik Manne, Erik G. Boman, Ümit V. Çatalyürek |
J. Parallel Distributed Comput. | 5 |
| 2008 | Large-Scale Biomedical Image Analysis in Grid EnvironmentsabstractThis paper presents the application of a component-based Grid middleware system for processing extremely large images obtained from digital microscopy devices. We have developed parallel, out-of-core techniques for different classes of data processing operations employed on images from confocal microscopy scanners. These techniques are combined into a data preprocessing and analysis pipeline using the component-based middleware system. The experimental results show that: 1) our implementation achieves good performance and can handle very large datasets on high-performance Grid nodes, consisting of computation and/or storage clusters and 2) it can take advantage of Grid nodes connected over high-bandwidth wide-area networks by combining task and data parallelism. Vijay S. Kumar, Benjamin Rutt, Tahsin M. Kurç, Ümit V. Çatalyürek, Tony Pan, Sunny K. Chow, Stephan Lamont, Maryann E. Martone, Joel H. Saltz |
IEEE Trans. Inf. Technol. Biomed. | 4 |
| 2007 | Computerized Pathological Image Analysis For Neuroblastoma Prognosis
Metin Nafi Gürcan, Jun Kong 0002, Olcay Sertel, Berkant Barla Cambazoglu, Joel H. Saltz, Ümit V. Çatalyürek |
AMIA | 6 |
| 2007 | Pathological Image Analysis Using the GPU: Stroma Classification for NeuroblastomaabstractNeuroblastoma is one of the most malignant childhood cancers affecting infants mostly. The current prognosis is based on microscopic examination of slides by expert pathologists, a process that is error-prone, time consuming and may lead to inter- and intra-reader variations. Therefore, we are developing a Computer Aided Prognosis (CAP) system which provides computerized image analysis to assist pathologist in their prognosis. Since this system operates on relatively large- scale images and requires sophisticated algorithms, it takes a long time to process whole-slide images. In this paper, we propose a novel and efficient approach for the execution of a CAP system for neuroblastoma prognosis, using the graphics processing unit (GPU). By leveraging high memory bandwidth and strong floating point operation capabilities of the GPU, our goal is to achieve order of magnitude reduction in the overall execution time as compared to that on a CPU alone. The proposed approach was tested on a set of testing images with a promising accuracy of 99.4% and an execution performance gain factor up to 45 times compared to C++ code running on the CPU. Antonio Ruiz 0001, Olcay Sertel, Manuel Ujaldon, Ümit V. Çatalyürek, Joel H. Saltz, Metin Nafi Gürcan |
BIBM | 4 |
| 2007 | An Efficient and Reliable Scientific Workflow SystemabstractThis paper presents a fault tolerance framework for applications that process data using a distributed network of user-defined operations in a pipelined fashion. The framework saves intermediate results and messages exchanged among application components in a distributed data management system to facilitate quick recovery from failures. The experimental results show that the framework scales well and our approach introduces very little overhead to application execution. Tulio Tavares, George Teodoro, Tahsin M. Kurç, Renato Ferreira 0001, Dorgival O. Guedes, Wagner Meira Jr., Ümit V. Çatalyürek, Shannon Hastings, Scott Oster, Stephen Langella, Joel H. Saltz |
CCGRID | 7 |
| 2007 | Performance vs. accuracy trade-offs for large-scale image analysis applicationsabstractIn many data analysis applications, application-level parameters influence the execution time of the data analysis method or program. Some of these parameters also affect the accuracy of output of the analysis. In this work, we investigate execution strategies for adaptive data analysis applications where the user is willing to trade-off accuracy of output for performance gain and vice-versa. In order to meet the user defined quality of service requirements, the system must dynamically select values for the parameters during execution. We propose algorithms for adaptive processing of image tiles at different resolutions so that user defined requirements in terms of accuracy of the result and execution time constraints can be satisfied. We develop heuristics for estimation of accuracy vs performance characteristics of image tiles and for scheduling of the tiles for processing. We implement a demand-driven strategy for parallel execution of these heuristics on a parallel machine. We evaluate our approach for analysis of large images from digitized microscopy scanners. Vijay S. Kumar, Tahsin M. Kurç, Jun Kong 0002, Ümit V. Çatalyürek, Metin Nafi Gürcan, Joel H. Saltz |
CLUSTER | 4 |
| 2007 | Scheduling File Transfers for Data-Intensive Jobs on Heterogeneous Clusters
Gaurav Khanna 0002, Ümit V. Çatalyürek, Tahsin M. Kurç, P. Sadayappan, Joel H. Saltz |
Euro-Par | 2 |
| 2007 | Toward Optimizing Latency Under Throughput Constraints for Application Workflows on Clusters
Nagavijayalakshmi Vydyanathan, Ümit V. Çatalyürek, Tahsin M. Kurç, P. Sadayappan, Joel H. Saltz |
Euro-Par | 2 |
| 2007 | Intelligent Optimization of Parallel and Distributed ApplicationsabstractThis paper describes a new project that systematically addresses the enormous complexity of mapping applications to current and future parallel platforms. By integrating the system layers - domain-specific environment, application program, compiler, run-time environment, performance models and simulation, and workflow manager - and through a systematic strategy for application mapping, our approach exploit the vast machine resources available in such parallel platforms to dramatically increase the productivity of application programmers. This project brings together computer scientists in the areas represented by the system layers (i.e., language extensions, compilers, run-time systems, workflows) together with expertise in knowledge representation and machine learning. With expert domain scientists in molecular dynamics (MD) simulation, we are developing our approach in the context of a specific application class which already targets environments consisting of several hundreds of processors. In this way, we gain valuable insight into a generalizable strategy, while simultaneously producing performance benefits for existing and important applications. Bhupesh Bansal, Ümit V. Çatalyürek, Jacqueline Chame, Chun Chen 0002, Ewa Deelman, Yolanda Gil, Mary W. Hall, Vijay S. Kumar, Tahsin M. Kurç, Kristina Lerman, Aiichiro Nakano, Yoon-Ju Lee Nelson, Joel H. Saltz, Ashish Sharma 0001, Priya Vashishta |
IPDPS | 2 |
| 2007 | Hypergraph-based Dynamic Load Balancing for Adaptive Scientific ComputationsabstractAdaptive scientific computations require that periodic repartitioning (load balancing) occur dynamically to maintain load balance. Hypergraph partitioning is a successful model for minimizing communication volume in scientific computations, and partitioning software for the static case is widely available. In this paper, we present a new hypergraph model for the dynamic case, where we minimize the sum of communication in the application plus the migration cost to move data, thereby reducing total execution time. The new model can be solved using hypergraph partitioning with faced vertices. We describe an implementation of a parallel multilevel repartitioning algorithm within the Zoltan load-balancing toolkit, which to our knowledge is the first code for dynamic load balancing based on hypergraph partitioning. Finally, we present experimental results that demonstrate the effectiveness of our approach on a Linux cluster with up to 64 processors. Our new algorithm compares favorably to the widely used ParMETIS partitioning software in terms of quality, and would have reduced total execution time in most of our test cases. Ümit V. Çatalyürek, Erik G. Boman, Karen D. Devine, Doruk Bozdag, Robert T. Heaphy, Lee Ann Riesen |
IPDPS | 1 |
| 2007 | A global address space framework for locality aware scheduling of block-sparse computationsabstractIn this paper, we present a mechanism for automatic management of the memory hierarchy, including secondary storage, in the context of a global address space parallel programming framework. The programmer specifies the parallelism and locality in the computation. The scheduling of the computation into stages, together with the movement of the associated data between secondary storage and global memory, and between global memory and local memory, is automatically managed. A novel formulation of hypergraph partitioning is used to model the optimization problem of minimizing disk I/O. Experimental evaluation using a sub-computation from the quantum chemistry domain shows a reduction in the disk I/O cost by up to a factor of 11, and a reduction in turnaround time by up to 49%, as compared to alternative approaches used in state-of-the-art quantum chemistry codes. Sriram Krishnamoorthy, Ümit V. Çatalyürek, Jarek Nieplocha, Atanas Rountev, P. Sadayappan |
IPDPS | 2 |
| 2006 | MSSG: A Framework for Massive-Scale Semantic GraphsabstractThis paper presents a middleware framework for storing, accessing and analyzing massive-scale semantic graphs. The framework, MSSG, targets scale-free semantic graphs with O(1012) (trillion) vertices and edges. Here, we present the overall architectural design of the framework, as well as a prototype implementation for cluster architectures. The sheer size of these massive-scale semantic graphs prohibits storing the entire graph in memory even on medium- to large-scale parallel architectures. We therefore propose a new graph database, grDB, for the efficient storage and retrieval of large scale-free semantic graphs on secondary storage. This new database supports the efficient and scalable execution of parallel out-of-core graph algorithms which are essential for analyzing semantic graphs of massive size. We have also developed a parallel out-of-core breadth-first search algorithm for performance study. To the best of our knowledge, it is the first of such algorithms reported in the literature. Experimental evaluations on large real-world semantic graphs show that the MSSG framework scales well, and grDB outperforms widely used open-source out-of-core databases, such as BerkeleyDB and MySQL, in the storage and retrieval of scale-free graphs Timothy D. R. Hartley, Ümit V. Çatalyürek, Füsun Özgüner, Andy B. Yoo, Scott Kohn, Keith W. Henderson |
CLUSTER | 2 |
| 2006 | Locality Conscious Processor Allocation and Scheduling for Mixed Parallel ApplicationsabstractComplex applications can often be viewed as a collection of coarse-grained data-parallel application components with precedence constraints. It has been shown that combining task and data parallelism (mixed parallelism) can be an effective execution paradigm for these applications. In this paper, we present an algorithm to compute the appropriate mix of task and data parallelism based on the scalability characteristics of the tasks as well as the intertask data communication costs, such that the parallel completion time (makespan) is minimized. The algorithm iteratively reduces the makespan by increasing the degree of data parallelism of tasks on the critical path that have good scalability and a low degree of potential task parallelism. Data communication costs along the critical path are minimized by exploiting parallel transfer mechanisms and use of a locality conscious backfill scheduler. Evaluation using benchmark task graphs derived from real applications as well as synthetic graphs shows that our algorithm consistently performs better than previous scheduling schemes Nagavijayalakshmi Vydyanathan, Sriram Krishnamoorthy, Gerald Sabin, Ümit V. Çatalyürek, Tahsin M. Kurç, P. Sadayappan, Joel H. Saltz |
CLUSTER | 4 |
| 2006 | Task Scheduling and File Replication for Data-Intensive Jobs with Batch-shared I/OabstractThis paper addresses the problem of efficient execution of a batch of data-intensive tasks with batch-shared I/O behavior, on coupled storage and compute clusters. Two scheduling schemes are proposed: 1) a 0-1 integer programming (IP) based approach, which couples task scheduling and data replication, and 2) a bi-level hypergraph partitioning based heuristic approach (BiPartition), which decouples task scheduling and data replication. The experimental results show that: 1) the IP scheme achieves the best batch execution time, but has significant scheduling overhead, thereby restricting its application to small scale workloads, and 2) the BiPartition scheme is a better fit for larger workloads and systems - it has very low scheduling overhead and no more than 5-10% degradation in solution quality, when compared with the IP based approach Gaurav Khanna 0002, Nagavijayalakshmi Vydyanathan, Ümit V. Çatalyürek, Tahsin M. Kurç, Sriram Krishnamoorthy, P. Sadayappan, Joel H. Saltz |
HPDC | 3 |
| 2006 | On Creating Efficient Object-relational Views of Scientific DatasetsabstractScientific datasets are often large and distributed in flat files across several storage nodes. Scientists frequently want to analyze subsets of these datasets. A data source abstraction that provides an object-relational view of data while hiding the details of storage and transport mechanisms and dataset layouts is useful in this regard. In this abstraction, basic data sources (BDS) interpret flat files as a set of records and are the building blocks of the view mechanism. Derived data sources (DDS) may be built on top of BDSs and provide more complex objects that serve the scientists' needs. The simplest DDS is one that supports a join based view over BDSs. We investigate issues involving building such DDSs for scientific applications and consider distributed versions of the indexed join and the grace hash join algorithms. We construct cost models that capture their performance in a restricted space of dataset and system parameters and compare them analytically and experimentally Sivaramakrishnan Narayanan, Tahsin M. Kurç, Ümit V. Çatalyürek, Joel H. Saltz |
ICPP | 3 |
| 2006 | An Integrated Approach for Processor Allocation and Scheduling of Mixed-Parallel ApplicationsabstractComputationally complex applications can often be viewed as a collection of coarse-grained data-parallel tasks with precedence constraints. Researchers have shown that combining task and data parallelism (mixed parallelism) can be an effective approach for executing these applications, as compared to pure task or data parallelism. In this paper, we present an approach to determine the appropriate mix of task and data parallelism, i.e., the set of tasks that should be run concurrently and the number of processors to be allocated to each task. An iterative algorithm is proposed that couples processor allocation and scheduling of mixed-parallel applications on compute clusters so as to minimize the parallel completion time (makespan). Our algorithm iteratively reduces the makespan by increasing the degree of data parallelism of tasks on the critical path that have good scalability and a low degree of potential task parallelism. The approach employs a look-ahead technique to escape local minima and uses priority based backfill scheduling to efficiently schedule the parallel tasks onto processors. Evaluation using benchmark task graphs derived from real applications as well as synthetic graphs shows that our algorithm consistently performs better than CPR and CPA, two previously proposed scheduling schemes, as well as pure task and data parallelism Nagavijayalakshmi Vydyanathan, Sriram Krishnamoorthy, Gerald Sabin, Ümit V. Çatalyürek, Tahsin M. Kurç, P. Sadayappan, Joel H. Saltz |
ICPP | 4 |
| 2006 | Using Space and Attribute Partitioned Partial Replicas for Data Subsetting and Aggregation QueriesabstractPartial replication is one type of optimization to speed up execution of queries submitted to large datasets. In partial replication, a portion of the dataset is extracted, re-organized, and re-distributed across the storage system. In this paper we investigate methods for efficient execution of queries when replicas of a dataset exist; we assume the replicas have already been created and do not target the replica creation problem. We propose a cost model and algorithm for combined use of space partitioned and attribute partitioned replicas for executing data subsetting range queries. We extend the cost model and propose a greedy algorithm to address range queries with aggregation operations. The extended replica selection algorithm allows uneven partitioning of replicas across storage nodes. Different replicas can be partitioned across different subsets of storage nodes. We have implemented these techniques as part of an automatic data virtualization system and have evaluated the benefits of our techniques using this system. We demonstrate the efficacy of the algorithms on parallel machines using queries on datasets from oil reservoir simulation studies and satellite data processing applications Li Weng, Ümit V. Çatalyürek, Tahsin M. Kurç, Gagan Agrawal, Joel H. Saltz |
ICPP | 2 |
| 2006 | A task duplication based bottom-up scheduling algorithm for heterogeneous environmentsabstractWe propose a new duplication-based DAG scheduling algorithm for heterogeneous computing environments. Contrary to the traditional approaches, proposed algorithm traverses the DAG in a bottom-up fashion while taking advantage of task duplication and task insertion. Experimental results on random DAGs and three different application DAGs show that the makespans generated by the proposed DBUS algorithm are much better than those generated by the existing algorithms, HEFT, HCPFD and HCNF. Doruk Bozdag, Ümit V. Çatalyürek, Füsun Özgüner |
IPDPS | 2 |
| 2006 | Parallel hypergraph partitioning for scientific computingabstractGraph partitioning is often used for load balancing in parallel computing, but it is known that hypergraph partitioning has several advantages. First, hypergraphs more accurately model communication volume, and second, they are more expressive and can better represent nonsymmetric problems. Hypergraph partitioning is particularly suited to parallel sparse matrix-vector multiplication, a common kernel in scientific computing. We present a parallel software package for hypergraph (and sparse matrix) partitioning developed at Sandia National Labs. The algorithm is a variation on multilevel partitioning. Our parallel implementation is novel in that it uses a two-dimensional data distribution among processors. We present empirical results that show our parallel implementation achieves good speedup on several large problems (up to 33 million nonzeros) with up to 64 processors on a Linux cluster Karen D. Devine, Erik G. Boman, Robert T. Heaphy, Rob H. Bisseling, Ümit V. Çatalyürek |
IPDPS | 5 |
| 2006 | An extensible global address space framework with decoupled task and data abstractionsabstractAlthough message passing using MPI is the dominant model for parallel programming today, the significant effort required to develop high-performance MPI applications has prompted the development of several parallel programming models that are more convenient. Programming models such as Co-Array Fortran, Global Arrays, Titanium, and UPC provide a more convenient global view of the data, but face significant challenges in delivering high performance over a range of applications. It is particularly challenging to achieve high performance using global-address-space languages for unstructured applications with irregular data structures. In this paper, we describe a global-address-space parallel programming framework with decoupled task and data abstractions. The framework centers around the use of task pools, where tasks specify operands in a distributed, globally addressable pool of data chunks. The data chunks can be addressed in a logical multidimensional "tuple" space, and are distributed among the nodes of the system. Locality-aware load balancing of tasks in the task pool is achieved through judicious mapping via hyper-graph partitioning, as well as dynamic task/data migration. The framework implements a transparent interface for out-of-core data, so that explicit orchestration of movement of data between disks and memory is not required of the programmer. The use of the framework for implementation of parallel block-sparse tensor computations in the context of a quantum chemistry application is illustrated. Sriram Krishnamoorthy, Ümit V. Çatalyürek, Jarek Nieplocha, Atanas Rountev, P. Sadayappan |
IPDPS | 2 |
| 2006 | An approach to locality-conscious load balancing and transparent memory hierarchy management with a global-address-space parallel programming modelabstractThe development of efficient parallel out-of-core applications is often tedious, because of the need to explicitly manage the movement of data between files and data structures of the parallel program. Several large-scale applications require multiple passes of processing over data too large to fit in memory, where significant concurrency exists within each pass. This paper describes a global-address-space framework for the convenient specification and efficient execution of parallel out-of-core applications operating on block-sparse data. The programming model provides a global view of block-sparse matrices and a mechanism for the expression of parallel tasks that operate on block-sparse data. The tasks are automatically partitioned into phases that operate on memory-resident data, and mapped onto processors to optimize load balance and data locality. Experimental results are presented that demonstrate the utility of the approach Sriram Krishnamoorthy, Ümit V. Çatalyürek, Jarek Nieplocha, P. Sadayappan |
IPDPS | 2 |
| 2006 | A Data Locality Aware Online Scheduling Approach for I/O-Intensive Jobs with File Sharing
Gaurav Khanna 0002, Ümit V. Çatalyürek, Tahsin M. Kurç, P. Sadayappan, Joel H. Saltz |
JSSPP | 2 |
| 2006 | Improving Functional Modularity in Protein-Protein Interactions Graphs Using Hub-Induced Subgraphs
Duygu Ucar, Sitaram Asur, Ümit V. Çatalyürek, Srinivasan Parthasarathy 0001 |
PKDD | 3 |
| 2006 | Data management and query - Hypergraph partitioning for automatic memory hierarchy managementabstractIn this paper, we present a mechanism for automatic management of the memory hierarchy, including secondary storage, in the context of a global address space parallel programming framework. The programmer specifies the parallelism and locality in the computation. The scheduling of the computation into stages, together with the movement of the associated data between secondary storage and global memory, and between global memory and local memory, is automatically managed. A novel formulation of hypergraph partitioning is used to model the optimization problem of minimizing disk I/O. Experimental evaluation of the proposed approach using a sub-computation from the quantum chemistry domain shows a reduction in the disk I/O cost by up to a factor of 11, and a reduction in turnaround time by up to 49%, as compared to alternative approaches used in state-of-the-art quantum chemistry codes. Sriram Krishnamoorthy, Ümit V. Çatalyürek, Jarek Nieplocha, Atanas Rountev, P. Sadayappan |
SC | 2 |
| 2006 | Imaging and visual analysis - Large image correction and warping in a cluster environmentabstractThis paper is concerned with efficient execution of a pipeline of data processing operations on very large images obtained from confocal microscopy instruments. We describe parallel, out-of-core algorithms for each operation in this pipeline. One of the challenging steps in the pipeline is the warping operation using inverse mapping based methods. We propose and investigate a set of algorithms to handle the warping computations on storage clusters. Our experimental results show that the proposed approaches are scalable both in terms of number of processors and the size of images. Vijay S. Kumar, Benjamin Rutt, Tahsin M. Kurç, Ümit V. Çatalyürek, Joel H. Saltz, Sunny K. Chow, Stephan Lamont, Maryann E. Martone |
SC | 4 |
| 2006 | Application of Information Technology: An XML-based System for Synthesis of Data from Disparate DatabasesabstractDiverse data sets have become key building blocks of translational biomedical research. Data types captured and referenced by sophisticated research studies include high throughput genomic and proteomic data, laboratory data, data from imagery, and outcome data. In this paper, the authors present the application of an XML-based data management system to support integration of data from disparate data sources and large data sets. This system facilitates management of XML schemas and on-demand creation and management of XML databases that conform to these schemas. They illustrate the use of this system in an application for genotype-phenotype correlation analyses. This application implements a method of phenotype-genotype correlation based on phylogenetic optimization of large data sets of mouse SNPs and phenotypic data. The application workflow requires the management and integration of genomic information and phenotypic data from external data repositories and from the results of phenotype-genotype correlation analyses. Our implementation supports the process of carrying out a complex workflow that includes large-scale phylogenetic tree optimizations and application of Maddison's concentrated changes test to large phylogenetic tree data sets. The data management system also allows collaborators to share data in a uniform way and supports complex queries that target data sets. Tahsin M. Kurç, Daniel Janies, Andrew D. Johnson, Stephen Langella, Scott Oster, Shannon Hastings, Farhat Habib, Terry Camerlengo, David Ervin, Ümit V. Çatalyürek, Joel H. Saltz |
J. Am. Medical Informatics Assoc. | 10 |
| 2005 | A hypergraph partitioning based approach for scheduling of tasks with batch-shared I/OabstractThis paper proposes a novel, hypergraph partitioning based strategy to schedule multiple data analysis tasks with batch-shared I/O behavior. This strategy formulates the sharing of files among tasks as a hypergraph to minimize the I/O overheads due to transferring of the same set of files multiple times and employs a dynamic scheme for file transfers to reduce contention on the storage system. We experimentally evaluate the proposed approach using application emulators from two application domains; analysis of remotely-sensed data and biomedical imaging. Gaurav Khanna 0002, Nagavijayalakshmi Vydyanathan, Tahsin M. Kurç, Ümit V. Çatalyürek, Pete Wyckoff, Joel H. Saltz, P. Sadayappan |
CCGRID | 4 |
| 2005 | Servicing range queries on multidimensional datasets with partial replicasabstractPartial replication is one type of optimization to speed up execution of queries submitted to large datasets. In partial replication, a portion of the dataset is extracted, re-organized, and re-distributed across the storage system. The objective is to reduce the volume of I/O and increase I/O parallelism for different types of queries and for the portions of the dataset that are likely to be accessed frequently. When multiple partial replicas of a dataset exist, query execution plan should be generated so as to use the best combination of subsets of partial replicas (and possibly the original dataset) to minimize query execution time. In this paper, we present a compiler and runtime approach for range queries submitted against distributed scientific datasets. A heuristic algorithm is proposed to choose the set of replicas to reduce query execution. We show the efficiency of the proposed method using datasets and queries in oil reservoir simulation studies on a cluster machine. Li Weng, Ümit V. Çatalyürek, Tahsin M. Kurç, Gagan Agrawal, Joel H. Saltz |
CCGRID | 2 |
| 2005 | Distributed Out-of-Core Preprocessing of Very Large Microscopy Images for Efficient QueryingabstractWe present a combined task- and data-parallel approach for distributed execution of pre-processing operations to support efficient evaluation of polygonal aggregation queries on digitized microscopy images. Our approach targets out-of-core, pipelined processing of very large images on active storage clusters. Our experimental results show that the proposed approach is scalable both in terms of number of processors and the size of images Benjamin Rutt, Vijay S. Kumar, Tony Pan, Tahsin M. Kurç, Ümit V. Çatalyürek, Joel H. Saltz |
CLUSTER | 5 |
| 2005 | A Scalable Parallel Graph Coloring Algorithm for Distributed Memory Computers
Erik G. Boman, Doruk Bozdag, Ümit V. Çatalyürek, Assefaw Hadish Gebremedhin, Fredrik Manne |
Euro-Par | 3 |
| 2005 | A Parallel Distance-2 Graph Coloring Algorithm for Distributed Memory Computers
Doruk Bozdag, Ümit V. Çatalyürek, Assefaw Hadish Gebremedhin, Fredrik Manne, Erik G. Boman, Füsun Özgüner |
HPCC | 2 |
| 2005 | A Task Duplication Based Scheduling Algorithm Using Partial SchedulesabstractWe propose a novel replication-based two-phase scheduling algorithm designed to achieve DAG scheduling with small makespans and high efficiency. In the first phase, the schedule length of the application is minimized using a novel approach that utilizes partial schedules. In the second phase, the number of processors required is minimized by eliminating and merging these partial schedules. Experimental results on random DAGs show that the makespans generated by the proposed algorithm are slightly better than those generated by the well known CPFD algorithm whereas the number of processors used is less than half of what is needed by CPFD solutions. Doruk Bozdag, Füsun Özgüner, Eylem Ekici, Ümit V. Çatalyürek |
ICPP | 4 |
| 2005 | A Scalable Distributed Parallel Breadth-First Search Algorithm on BlueGene/LabstractMany emerging large-scale data science applications require searching large graphs distributed across multiple memories and processors. This paper presents a distributed breadth- first search (BFS) scheme that scales for random graphs with up to three billion vertices and 30 billion edges. Scalability was tested on IBM BlueGene/L with 32,768 nodes at the Lawrence Livermore National Laboratory. Scalability was obtained through a series of optimizations, in particular, those that ensure scalable use of memory. We use 2D (edge) partitioning of the graph instead of conventional 1D (vertex) partitioning to reduce communication overhead. For Poisson random graphs, we show that the expected size of the messages is scalable for both 2D and 1D partitionings. Finally, we have developed efficient collective communication functions for the 3D torus architecture of BlueGene/L that also take advantage of the structure in the problem. The performance and characteristics of the algorithm are measured and reported. Andy B. Yoo, Edmond Chow, Keith W. Henderson, Will McLendon III, Bruce Hendrickson, Ümit V. Çatalyürek |
SC | 6 |
| 2005 | A simulation and data analysis system for large-scale, data-driven oil reservoir simulation studiesabstractAbstract The main goal of oil reservoir management is to provide more efficient, cost‐effective and environmentally safer production of oil from reservoirs. Numerical simulations can aid in the design and implementation of optimal production strategies. However, traditional simulation‐based approaches to optimizing reservoir management are rapidly overwhelmed by data volume when large numbers of realizations are sought using detailed geologic descriptions. In this paper, we describe a software architecture to facilitate large‐scale simulation studies, involving ensembles of long‐running simulations and analysis of vast volumes of output data. Copyright © 2005 John Wiley & Sons, Ltd. Tahsin M. Kurç, Ümit V. Çatalyürek, Joel H. Saltz, Ryan Martino, Mary F. Wheeler, Malgorzata Peszynska, Alan Sussman, Christian Hansen 0002, Mrinal K. Sen, Roustam Seifoullaev, Paul L. Stoffa, Carlos Torres-Verdín, Manish Parashar |
Concurr. Pract. Exp. | 2 |
| 2005 | Application of Grid-enabled technologies for solving optimization problems in data-driven reservoir studies
Manish Parashar, Hector Klie, Ümit V. Çatalyürek, Tahsin M. Kurç, Wolfgang Bangerth, Vincent Matossian, Joel H. Saltz, Mary F. Wheeler |
Future Gener. Comput. Syst. | 3 |
| 2005 | Application of Information Technology: A Grid-Based Image Archival and Analysis SystemabstractHere the authors present a Grid-aware middleware system, called GridPACS, that enables management and analysis of images in a massive scale, leveraging distributed software components coupled with interconnected computation and storage platforms. The need for this infrastructure is driven by the increasing biomedical role played by complex datasets obtained through a variety of imaging modalities. The GridPACS architecture is designed to support a wide range of biomedical applications encountered in basic and clinical research, which make use of large collections of images. Imaging data yield a wealth of metabolic and anatomic information from macroscopic (e.g., radiology) to microscopic (e.g., digitized slides) scale. Whereas this information can significantly improve understanding of disease pathophysiology as well as the noninvasive diagnosis of disease in patients, the need to process, analyze, and store large amounts of image data presents a great challenge. Shannon Hastings, Scott Oster, Stephen Langella, Tahsin M. Kurç, Tony Pan, Ümit V. Çatalyürek, Joel H. Saltz |
J. Am. Medical Informatics Assoc. | 6 |
| 2004 | Serving queries to multi-resolution datasets on disk-based storage clustersabstractThis paper is concerned with efficient querying of very large multi-resolution datasets on storage and compute clusters. We present a suite of services that support storage, indexing, and data processing (data sampling and data aggregation) on datasets that consist of a collection of multi-resolution Grids. We empirically evaluate the performance impact of different data declustering, indexing, and query processing strategies. The experimental evaluation is carried out using a data server implemented to serve multi-terabyte multi-resolution volumetric datasets to remote visualization clients and a one-terabyte multi-resolution volumetric dataset on a PC cluster with distributed disk space. Tony Pan, Ümit V. Çatalyürek, Tahsin M. Kurç, Joel H. Saltz |
CCGRID | 3 |
| 2004 | A distributed data management middleware for data-driven application systemsabstractA key challenge in supporting data-driven scientific applications is the storage and management of input and output data in a distributed environment. We describe a distributed storage middleware, based on a data and metadata management framework, to address this problem. In this middleware system, applications define the structure of their input and output data using XML schemas. The system provides support for 1) registration, versioning, management of schemas, and 2) management of storage, querying, and retrieval of instance data corresponding to the schemas in distributed databases. We carry out an experimental evaluation of the system on a set of PC clusters connected over wide- (WANs) and local-area networks (LANs). Stephen Langella, Shannon Hastings, Scott Oster, Tahsin M. Kurç, Ümit V. Çatalyürek, Joel H. Saltz |
CLUSTER | 5 |
| 2004 | An Approach for Automatic Data Virtualization
Li Weng, Gagan Agrawal, Ümit V. Çatalyürek, Tahsin M. Kurç, Sivaramakrishnan Narayanan, Joel H. Saltz |
HPDC | 3 |
| 2004 | Strategies for Using Additional Resources in Parallel Hash-Based Join Algorithms
Tahsin M. Kurç, Tony Pan, Ümit V. Çatalyürek, Sivaramakrishnan Narayanan, Pete Wyckoff, Joel H. Saltz |
HPDC | 4 |
| 2003 | Image Processing or the Grid: A Toolkit or Building Grid-enabled Image Processing ApplicationsabstractAnalyzing large and distributed image datasets is a crucial step in understanding the structural and functional characteristics of biological systems. In this paper, we present the design and implementation of a toolkit that allows rapid and efficient development of biomedical image analysis applications in a distributed environment. This toolkit employs the Insight Segmentation and Registration Toolkit (ITK) and Visualization Toolkit (VTK) layered on a component-based framework. We present experimental results on a cluster of workstations. Shannon Hastings, Tahsin M. Kurç, Stephen Langella, Ümit V. Çatalyürek, Tony Pan, Joel H. Saltz |
CCGRID | 4 |
| 2003 | Impact of High Performance Sockets on Data Intensive ApplicationsabstractThe challenging issues in supporting data intensive applications on clusters include efficient movement of large volumes of data between processor memories and efficient coordination of data movement and processing by a runtime support to achieve high performance. Such applications have several requirements such as guarantees in performance, scalability with these guarantees and adaptability to heterogeneous environments. With the advent of user-level protocols like the Virtual Interface Architecture (VIA) and the modern InfiniBand Architecture, the latency and bandwidth experienced by applications has approached to that of the physical network on clusters. In order to enable applications written on top of TCP/IP to take advantage of the high performance of these user-level protocols, researchers have come up with a number of techniques including User Level Sockets Layers over high performance protocols. In this paper, we study the performance and limitations of such substrate, referred to here as SocketVIA, using a component framework designed to provide runtime support for data intensive applications. The experimental results show that by reorganizing certain components of an application (in our case, the partitioning of a dataset into smaller data chunks), we can make significant improvements in application performance. This leads to a higher scalability of applications with performance guarantees. It also allows fine grained load balancing, hence making applications more adaptable to heterogeneity in resource availability. The experimental results also show that the different performance characteristics of SocketVIA allow a more efficient partitioning of data at the source nodes, thus improving the performance of the application up to an order of magnitude in some cases. Pavan Balaji, Jiesheng Wu, Tahsin M. Kurç, Ümit V. Çatalyürek, Dhabaleswar K. Panda 0001, Joel H. Saltz |
HPDC | 4 |
| 2003 | Optimizing Reduction Computations In a Distributed EnvironmentabstractWe investigate runtime strategies for data-intensive applications that invovle generalized reductions on large, distributed datasets.Our set of strategies includes replicated filter state, partitioned filter state, and hybrid options between these two extremes.We evaluate these strategies using emulators of three real applications, different query and output sizes, and a number of configurations.We consider execution in a homogeneous cluster and in a distributed environment where only a subset of nodes hst the data.Our results show replicating the filter state scales well and outperforms other schemes, if sufficient memory is available and sufficient computation is involved to offset the cost of global merge step.In other cases, hybrid is usually the best.Moreover, in almost all cases, the performance of the hybrid strategy is quite close to the best strategy. Thus, we believe that hybrid is an attractive approach when the relative performance of different schemes cannot be predicted. Tahsin M. Kurç, Feng Lee, Gagan Agrawal, Ümit V. Çatalyürek, Renato Ferreira 0001, Joel H. Saltz |
SC | 4 |
| 2003 | The virtual microscopeabstractWe present the design and implementation of the Virtual Microscope, a software system employing a client/server architecture to provide a realistic emulation of a high power light microscope. The system provides a form of completely digital telepathology, allowing simultaneous access to archived digital slide images by multiple clients. The main problem the system targets is storing and processing the extremely large quantities of data required to represent a collection of slides. The Virtual Microscope client software runs on the end user's PC or workstation, while database software for storing, retrieving and processing the microscope image data runs on a parallel computer or on a set of workstations at one or more potentially remote sites. We have designed and implemented two versions of the data server software. One implementation is a customization of a database system framework that is optimized for a tightly coupled parallel machine with attached local disks. The second implementation is component-based, and has been designed to accommodate access to and processing of data in a distributed, heterogeneous environment. We also have developed caching client software, implemented in Java, to achieve good response time and portability across different computer platforms. The performance results presented show that the Virtual Microscope systems scales well, so that many clients can be adequately serviced by an appropriately configured data server. Ümit V. Çatalyürek, Michael D. Beynon, Chialin Chang, Tahsin M. Kurç, Alan Sussman, Joel H. Saltz |
IEEE Trans. Inf. Technol. Biomed. | 1 |
| 2002 | Executing multiple pipelined data analysis operations in the gridabstractProcessing of data in many data analysis applications can be represented as an acyclic, coarse grain data flow, from data sources to the client. This paper is concerned with scheduling of multiple data analysis operations, each of which is represented as a pipelined chain of processing on data. We define the scheduling problem for effectively placing components onto Grid resources, and propose two scheduling algorithms. Experimental results are presented using a visualization application. Matthew Spencer, Renato Ferreira 0001, Michael D. Beynon, Tahsin M. Kurç, Ümit V. Çatalyürek, Alan Sussman, Joel H. Saltz |
SC | 5 |
| 2002 | Processing large-scale multi-dimensional data in parallel and distributed environments
Michael D. Beynon, Chialin Chang, Ümit V. Çatalyürek, Tahsin M. Kurç, Alan Sussman, Henrique Andrade, Renato Ferreira 0001, Joel H. Saltz |
Parallel Comput. | 3 |
| 2001 | A Fine-Grain Hypergraph Model for 2D Decomposition of Sparse MatricesabstractWe propose a new hypergraph model for the decomposition of irregular computational domains. This work focuses on the decomposition of sparse matrices for parallel matrix-vector multiplication. However, the proposed model can also be used to decompose computational domains of other parallel reduction problems. We propose a “finegrain” hypergraph model for two-dimensional decomposition of sparse matrices. In the proposed fine-grain hypergraph model, vertices represent nonzeros and hyperedges represent sparsity patterns of rows and columns of the matrix. By partitioning the fine-grain hypergraph into equally weighted vertex parts (processors) so that hyperedges are split among as few processors as possible, the model correctly minimizes communication volume while maintaining computationalload balance. Experimental results on a wide range of realistic sparse matrices confirm the validity of the proposed model, by achieving up to 50 percent better decompositionsthan the existing models, in terms of totalcommunication volume. 1 Ümit V. Çatalyürek, Cevdet Aykanat |
IPDPS | 1 |
| 2001 | A hypergraph-partitioning approach for coarse-grain decompositionabstractWe propose a new two-phase method for the coarse-grain decomposition of irregular computational domains. This work focuses on the 2D partitioning of sparse matrices for parallel matrix-vector multiplication. However, the proposed model can also be used to decompose computational domains of other parallel reduction problems. This work also introduces the use of multi-constraint hypergraph partitioning, for solving the decomposition problem. The proposed method explicitly models the minimization of communication volume while enforcing the upper bound of p + q --- 2 on the maximum number of messages handled by a single processor, for a parallel system with P = p × q processors. Experimental results on a wide range of realistic sparse matrices confirm the validity of the proposed methods, by achieving up to 25 percent better partitions than the standard graph model, in terms of total communication volume, and 59 percent better partitions in terms of number of messages, on the overall average. Ümit V. Çatalyürek, Cevdet Aykanat |
SC | 1 |
| 2001 | Distributed processing of very large datasets with DataCutter
Michael D. Beynon, Tahsin M. Kurç, Ümit V. Çatalyürek, Chialin Chang, Alan Sussman, Joel H. Saltz |
Parallel Comput. | 3 |
| 1999 | Hypergraph-Partitioning-Based Decomposition for Parallel Sparse-Matrix Vector MultiplicationabstractIn this work, we show that the standard graph-partitioning-based decomposition of sparse matrices does not reflect the actual communication volume requirement for parallel matrix-vector multiplication. We propose two computational hypergraph models which avoid this crucial deficiency of the graph model. The proposed models reduce the decomposition problem to the well-known hypergraph partitioning problem. The recently proposed successful multilevel framework is exploited to develop a multilevel hypergraph partitioning tool PaToH for the experimental verification of our proposed hypergraph models. Experimental results on a wide range of realistic sparse test matrices confirm the validity of the proposed hypergraph models. In the decomposition of the test matrices, the hypergraph models using PaToH and hMeTiS result in up to 63 percent less communication volume (30 to 38 percent less on the average) than the graph model using MeTiS, while PaToH is only 1.3-2.3 times slower than MeTiS on the average. Ümit V. Çatalyürek, Cevdet Aykanat |
IEEE Trans. Parallel Distributed Syst. | 1 |