Theodore L. Willke

dblp:40/3641 · also Ted Willke · DBLP profile ↗
← Back
37ranked-venue papers
2as first author
11since 2021 · last 2025
0000-0001-9825-513XORCID · verified

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

Artificial intelligence and machine learning · 14 · 5 since 2021Databases, data management, data science and information retrieval · 14 · 6 since 2021Systems, architecture and hardware · 9 · 2 since 2021Applied, interdisciplinary, general and emerging computing · 7 · 1 since 2021Graphics, computer vision, multimedia, augmented reality and games · 3Computer networks · 2 · 2 first-authorSoftware engineering, systems software and programming languages · 2
YearPublicationVenuePosition
2025 ConTraPh: Contrastive Learning for Parallelization and Performance Optimization
abstract
With the advancement of HPC platforms, the demand for high-performing applications continues to grow.One effective way to enhance program performance is through parallelization.However, fully leveraging the powerful hardware of HPC platforms poses significant challenges.Even experienced developers must carefully consider factors such as runtime, memory usage, and thread-scheduling overhead.Additionally, achieving successful parallelization often requires running applications to determine the optimal configurations.In this paper, we propose ConTraPh, a framework that integrates Contrastive Learning with Transformers and Graph Neural Networks to capture the inherent parallel characteristics of source programs through a multi-view program representation, utilizing both source code and compiler intermediate representations.This contrastive learning framework allows the model to effectively learn correct parallel configurations from positive samples while avoiding incorrect ones through negative samples.We evaluate Con-TraPh on six downstream tasks involving three different parallel programming models OpenMP, OpenCL and, Ope-nACC that include OpenMP clause prediction, performant reduction style detection, performant scheduling type detection, CPU/GPU parallelism prediction, Heterogeneous Device Mapping for OpenCL code, and OpenACC clause prediction.ConTraPh outperforms state-of-the-art models in these tasks, achieving accuracy improvements of up to 8%, 10%, 7%, 4%, 2%, and 9%, respectively.ConTraPh achieves speedups as high as 13x, 18x, 14x, and 4.4x on the reduction
Quazi Ishtiaque Mahmud, Ali TehraniJamsaz, Nesreen K. Ahmed, Theodore L. Willke, Ali Jannesari
ICS4
2025 PCEBench: A Multi-Dimensional Benchmark for Evaluating Large Language Models in Parallel Code Generation
abstract
The increasing complexity of software systems and advancements in hardware architectures have intensified the demand for efficient parallel code generation. While parallel programming offers significant performance benefits, it requires extensive expertise and effort due to the intricacies of synchronization, data management, and optimizations. To address these challenges, recent studies have explored the application of machine learning (ML) techniques in parallel code generation, aiming to reduce manual efforts and enhance performance outcomes. Large language models (LLMs) have recently revolutionized natural language processing (NLP) and demonstrated remarkable capabilities in code generation. However, evaluating their ability to generate high-performance parallel code presents unique challenges. Unlike sequential code evaluation, the evaluation of LLM-generated parallel code requires consideration of not only correctness but also efficiency and scalability in utilizing parallel resources. Concretely, existing benchmarks for LLM-generated parallel code evaluation are limited in size and scope compared to their sequential counterparts. To address this evaluation gap, we introduce PCEBench, a novel benchmark designed to assess LLMs' capabilities in generating parallel code. PCEBench focuses on multi-tasking and multidimensional performance evaluation, leveraging an LLM-based approach to generate verified prompts for parallel code generation. The benchmark incorporates scripts compatible with compilers and data race checkers, enabling comprehensive testing across critical dimensions such as compilability, executability, code self-correctness, functional correctness, data race detection, and speedup over serial implementations. By examining these multiple dimensions, PCEBench not only facilitates a thorough evaluation of LLMs in parallel code generation but also provides valuable insights for developers to enhance model performance in this challenging task. This comprehensive approach contributes to advancing the field of automated parallel programming and supports the development of more efficient and scalable software systems.
Nesreen K. Ahmed, Mihai Capota, Theodore L. Willke, Niranjan Hasabnis, Ali Jannesari
IPDPS4
2025 AutoParLLM: GNN-guided Context Generation for Zero-Shot Code Parallelization using LLMs
abstract
Quazi Ishtiaque Mahmud, Ali TehraniJamsaz, Hung D Phan, Le Chen, Mihai Capotă, Theodore L. Willke, Nesreen K. Ahmed, Ali Jannesari. Proceedings of the 2025 Conference of the Nations of the Americas Chapter of the Association for Computational Linguistics: Human Language Technologies (Volume 1: Long Papers). 2025.
Quazi Ishtiaque Mahmud, Ali TehraniJamsaz, Hung D. Phan, Mihai Capota, Theodore L. Willke, Nesreen K. Ahmed, Ali Jannesari
NAACL (Long Papers)6
2024 Structure Guided Prompt: Instructing Large Language Model in Multi-Step Reasoning by Exploring Graph Structure of the Text
abstract
Although Large Language Models (LLMs) excel at addressing straightforward reasoning tasks, they frequently struggle with difficulties when confronted by more complex multi-step reasoning due to a range of factors.Firstly, natural language often encompasses complex relationships among entities, making it challenging to maintain a clear reasoning chain over longer spans.Secondly, the abundance of linguistic diversity means that the same entities and relationships can be expressed using different terminologies and structures, complicating the task of identifying and establishing connections between multiple pieces of information.Graphs provide an effective solution to represent data rich in relational information and capture long-term dependencies among entities.To harness the potential of graphs, our paper introduces Structure Guided Prompt, an innovative three-stage task-agnostic prompting framework designed to improve the multi-step reasoning capabilities of LLMs in a zero-shot setting.This framework explicitly converts unstructured text into a graph via LLMs and instructs them to navigate this graph using taskspecific strategies to formulate responses.By effectively organizing information and guiding navigation, it enables LLMs to provide more accurate and context-aware responses.Our experiments show that this framework significantly enhances the reasoning capabilities of LLMs, enabling them to excel in a broader spectrum of natural language scenarios.
Kewei Cheng, Nesreen K. Ahmed, Theodore L. Willke, Yizhou Sun
EMNLP3
2024 A Structure-Aware Framework for Learning Device Placements on Computation Graphs
abstract
Computation graphs are Directed Acyclic Graphs (DAGs) where the nodes correspond to mathematical operations and are used widely as abstractions in optimizations of neural networks. The device placement problem aims to identify optimal allocations of those nodes to a set of (potentially heterogeneous) devices. Existing approaches rely on two types of architectures known as grouper-placer and encoder-placer, respectively. In this work, we bridge the gap between encoder-placer and grouper-placer techniques and propose a novel framework for the task of device placement, relying on smaller computation graphs extracted from the OpenVINO toolkit. The framework consists of five steps, including graph coarsening, node representation learning and policy optimization. It facilitates end-to-end training and takes into account the DAG nature of the computation graphs. We also propose a model variant, inspired by graph parsing networks and complex network analysis, enabling graph representation learning and jointed, personalized graph partitioning, using an unspecified number of groups. To train the entire framework, we use reinforcement learning using the execution time of the placement as a reward. We demonstrate the flexibility and effectiveness of our approach through multiple experiments with three benchmark models, namely Inception-V3, ResNet, and BERT. The robustness of the proposed framework is also highlighted through an ablation study. The suggested placements improve the inference speed for the benchmark models by up to $58.2\%$ over CPU execution and by up to $60.24\%$ compared to other commonly used baselines.
Shukai Duan 0002, Heng Ping, Nikos Kanakaris, Xiongye Xiao, Panagiotis Kyriakis, Nesreen K. Ahmed, Peiyu Zhang 0002, Guixiang Ma, Mihai Capota, Shahin Nazarian, Theodore L. Willke, Paul Bogdan
NeurIPS11
2024 Neural-Symbolic Methods for Knowledge Graph Reasoning: A Survey
abstract
Neural symbolic knowledge graph (KG) reasoning offers a promising approach that combines the expressive power of symbolic reasoning with the learning capabilities inherent in neural networks. This survey provides a comprehensive overview of advancements, techniques, and challenges in the field of neural symbolic KG reasoning. The survey introduces the fundamental concepts of KGs and symbolic logic, followed by an exploration of three significant KG reasoning tasks: KG completion, complex query answering, and logical rule learning. For each task, we thoroughly discuss three distinct categories of methods: pure symbolic methods, pure neural approaches, and the integration of neural networks and symbolic reasoning methods known as neural-symbolic. We carefully analyze and compare the strengths and limitations of each category of methods to provide a comprehensive understanding. By synthesizing recent research contributions and identifying open research directions, this survey aims to equip researchers and practitioners with a comprehensive understanding of the state-of-the-art in neural symbolic KG reasoning, fostering future advancements in this interdisciplinary domain.
Kewei Cheng, Nesreen K. Ahmed, Ryan Rossi, Theodore L. Willke, Yizhou Sun
ACM Trans. Knowl. Discov. Data4
2023 Augmenting Recurrent Graph Neural Networks with a Cache
abstract
While graph neural networks (GNNs) provide a powerful way to learn structured representations, it remains challenging to learn long-range dependencies in graphs. Recurrent GNNs only partly address this problem. In this paper, we propose a general approach for augmenting recurrent GNNs with a cache memory to improve their expressivity, especially for modeling long-range dependencies. Specifically, we first introduce a method of augmenting recurrent GNNs with a cache of previous hidden states. Then we further propose a general Cache-GNN framework by adding additional modules, including attention mechanism and positional/structural encoders, to improve the expressivity. We show that the Cache-GNNs outperforms other models on synthetic datasets as well as tasks on real-world datasets that require long-range information.
Guixiang Ma, Vy A. Vo, Theodore L. Willke, Nesreen K. Ahmed
KDD3
2023 Similarity search in the blink of an eye with compressed indices
abstract
Nowadays, data is represented by vectors. Retrieving those vectors, among millions and billions, that are similar to a given query is a ubiquitous problem, known as similarity search, of relevance for a wide range of applications. Graph-based indices are currently the best performing techniques for billion-scale similarity search. However, their random-access memory pattern presents challenges to realize their full potential. In this work, we present new techniques and systems for creating faster and smaller graph-based indices. To this end, we introduce a novel vector compression method, Locally-adaptive Vector Quantization (LVQ), that uses per-vector scaling and scalar quantization to improve search performance with fast similarity computations and a reduced effective bandwidth, while decreasing memory footprint and barely impacting accuracy. LVQ, when combined with a new high-performance computing system for graph-based similarity search, establishes the new state of the art in terms of performance and memory footprint. For billions of vectors, LVQ outcompetes the second-best alternatives: (1) in the low-memory regime, by up to 20.7x in throughput with up to a 3x memory footprint reduction, and (2) in the high-throughput regime by 5.8x with 1.4x less memory.
Cecilia Aguerrebere, Ishwar Singh Bhati, Mark Hildebrand, Mariano Tepper, Theodore L. Willke
Proc. VLDB Endow.5
2022 Role-Based Graph Embeddings
abstract
Random walks are at the heart of many existing node embedding and network representation learning methods. However, such methods have many limitations that arise from the use of traditional random walks, e.g., the embeddings resulting from these methods capture proximity (communities) among the vertices as opposed to structural similarity (roles). Furthermore, the embeddings are unable to transfer to new nodes and graphs as they are tied to node identity. To overcome these limitations, we introduce theRole2Vecframework based on the proposed notion ofattributed random walksto learn structural role-based embeddings. Notably, the framework serves as a basis for generalizing any walk-based method. TheRole2Vecframework enables these methods to be more widely applicable by learning inductive functions that capture the structural roles in the graph. Furthermore, the original methods are recovered as a special case of the framework when each vertex is mapped to its own function that uniquely identifies it. Finally, theRole2Vecframework is shown to be effective with an average AUC improvement of 17.8 percent for link prediction while requiring on average 853x less space than existing methods on a variety of graphs from different domains.
Nesreen K. Ahmed, Ryan Rossi, John Boaz Lee, Theodore L. Willke, Rong Zhou 0001, Xiangnan Kong, Hoda Eldardiry
IEEE Trans. Knowl. Data Eng.4
2021 Learning Code Representations Using Multifractal-based Graph Networks
abstract
Learning representations of software codes is a critical problem for a wide range of system applications, e.g., compiler optimization, software classification, malicious software detection, and performance optimization. Recently, learning graph-based representations of software programs has been used to model the inherent structural dependencies in programming languages (e.g., C++, Python). In this paper, we propose a novel graph neural network framework that utilizes multifractal analysis for LLVM intermediate representations (IR). We then show empirically that the proposed framework is capable of capturing long-range structural dependencies that appear in software codes. We conduct experiments and comparisons on two downstream system applications: (1) predicting heterogeneous compute device mappings (graph classification), and (2) compiler reachability analysis (node classification). We observe that introducing a structural inductive bias through multifractal topological features enables GNNs to capture long-range dependencies among nodes, thus, it improves the accuracy of GNN models for applications that require learning code representations.
Guixiang Ma, Mihai Capota, Theodore L. Willke, Shahin Nazarian, Paul Bogdan, Nesreen K. Ahmed
IEEE BigData4
2021 Deep graph similarity learning: a survey
abstract
Abstract In many domains where data are represented as graphs, learning a similarity metric among graphs is considered a key problem, which can further facilitate various learning tasks, such as classification, clustering, and similarity search. Recently, there has been an increasing interest in deep graph similarity learning, where the key idea is to learn a deep learning model that maps input graphs to a target space such that the distance in the target space approximates the structural distance in the input space. Here, we provide a comprehensive review of the existing literature of deep graph similarity learning. We propose a systematic taxonomy for the methods and applications. Finally, we discuss the challenges and future directions for this problem.
Guixiang Ma, Nesreen K. Ahmed, Theodore L. Willke, Philip S. Yu
Data Min. Knowl. Discov.3
2020 NeuroVectorizer: end-to-end vectorization with deep reinforcement learning
abstract
One of the key challenges arising when compilers vectorize loops for today’s SIMD-compatible architectures is to decide if vectorization or interleaving is beneficial. Then, the compiler has to determine the number of instructions to pack together and the interleaving level (stride). Compilers are designed today to use fixed-cost models that are based on heuristics to make vectorization decisions on loops. However, these models are unable to capture the data dependency, the computation graph, or the organization of instructions. Alternatively, software engineers often hand-write the vectorization factors of every loop. This, however, places a huge burden on them, since it requires prior experience and significantly increases the development time.
Ameer Haj-Ali, Nesreen K. Ahmed, Theodore L. Willke, Sophia Shao, Krste Asanovic, Ion Stoica
CGO3
2020 Approximating Stacked and Bidirectional Recurrent Architectures with the Delayed Recurrent Neural Network
abstract
Recent work has shown that topological enhancements to recurrent neural networks (RNNs) can increase their expressiveness and representational capacity. Two popular enhancements are stacked RNNs, which increases the capacity for learning non-linear functions, and bidirectional processing, which exploits acausal information in a sequence. In this work, we explore the delayed-RNN, which is a single-layer RNN that has a delay between the input and output. We prove that a weight-constrained version of the delayed-RNN is equivalent to a stacked-RNN. We also show that the delay gives rise to partial acausality, much like bidirectional networks. Synthetic experiments confirm that the delayed-RNN can mimic bidirectional networks, solving some acausal tasks similarly, and outperforming them in others. Moreover, we show similar performance to bidirectional networks in a real-world natural language processing task. These results suggest that delayed-RNNs can approximate topologies including stacked RNNs, bidirectional RNNs, and stacked bidirectional RNNs – but with equivalent or faster runtimes for the delayed-RNNs.
Javier Turek, Shailee Jain, Vy A. Vo, Mihai Capota, Alexander G. Huth, Theodore L. Willke
ICML6
2020 BrainIAK tutorials: User-friendly learning materials for advanced fMRI analysis
abstract
Advanced brain imaging analysis methods, including multivariate pattern analysis (MVPA), functional connectivity, and functional alignment, have become powerful tools in cognitive neuroscience over the past decade. These tools are implemented in custom code and separate packages, often requiring different software and language proficiencies. Although usable by expert researchers, novice users face a steep learning curve. These difficulties stem from the use of new programming languages (e.g., Python), learning how to apply machine-learning methods to high-dimensional fMRI data, and minimal documentation and training materials. Furthermore, most standard fMRI analysis packages (e.g., AFNI, FSL, SPM) focus on preprocessing and univariate analyses, leaving a gap in how to integrate with advanced tools. To address these needs, we developed BrainIAK (brainiak.org), an open-source Python software package that seamlessly integrates several cutting-edge, computationally efficient techniques with other Python packages (e.g., Nilearn, Scikit-learn) for file handling, visualization, and machine learning. To disseminate these powerful tools, we developed user-friendly tutorials (in Jupyter format; https://brainiak.org/tutorials/) for learning BrainIAK and advanced fMRI analysis in Python more generally. These materials cover techniques including: MVPA (pattern classification and representational similarity analysis); parallelized searchlight analysis; background connectivity; full correlation matrix analysis; inter-subject correlation; inter-subject functional connectivity; shared response modeling; event segmentation using hidden Markov models; and real-time fMRI. For long-running jobs or large memory needs we provide detailed guidance on high-performance computing clusters. These notebooks were successfully tested at multiple sites, including as problem sets for courses at Yale and Princeton universities and at various workshops and hackathons. These materials are freely shared, with the hope that they become part of a pool of open-source software and educational materials for large-scale, reproducible fMRI analysis and accelerated discovery.
Manoj Kumar 0025, Cameron T. Ellis, Qihong Lu, Mihai Capota, Theodore L. Willke, Peter J. Ramadge, Nicholas B. Turk-Browne, Kenneth A. Norman
PLoS Comput. Biol.6
2019 Deep Graph Similarity Learning for Brain Data Analysis
abstract
We propose an end-to-end graph similarity learning framework called Higher-order Siamese GCN for multi-subject fMRI data analysis. The proposed framework learns the brain network representations via a supervised metric-based approach with siamese neural networks using two graph convolutional networks as the twin networks. Our proposed framework performs higher-order convolutions by incorporating higher-order proximity in graph convolutional networks to characterize and learn the community structure in brain connectivity networks. To the best of our knowledge, this is the first community-preserving graph similarity learning framework for multi-subject brain network analysis. Experimental results on four real fMRI datasets demonstrate the potential use cases of the proposed framework for multi-subject brain analysis in health and neuropsychiatric disorders. Our proposed approach achieves an average AUC gain of $75$% compared to PCA, an average AUC gain of $65.5$% compared to Spectral Embedding, and an average AUC gain of $24.3$% compared to S-GCN across the four datasets, indicating promising applications in clinical investigation and brain disease diagnosis.
Guixiang Ma, Nesreen K. Ahmed, Theodore L. Willke, Dipanjan Sengupta, Michael W. Cole, Nicholas B. Turk-Browne, Philip S. Yu
CIKM3
2018 Matrix-normal models for fMRI analysis
abstract
Multivariate analysis of fMRI data has bene- fited substantially from advances in machine learning. Most recently, a range of prob- abilistic latent variable models applied to fMRI data have been successful in a variety of tasks, including identifying similarity pat- terns in neural data, combining multi-subject datasets, and mapping between brain and be- havior. Although these methods share some underpinnings, they have been developed as distinct methods, with distinct algorithms and software tools. We show how the matrix- variate normal (MN) formalism can unify some of these methods into a single frame- work. In doing so, we gain the ability to reuse noise modeling assumptions, algorithms, and code across models. Our primary theoretical contribution shows how some of these meth- ods can be written as instantiations of the same model, allowing us to generalize them to flexibly modeling structured noise covari- ances. Our formalism permits novel model variants and improved estimation strategies for SRM and RSA using substantially fewer parameters. We empirically demonstrate ad- vantages of our two new methods: for MN-RSA, we show up to 10x improvement in run- time, up to 6x improvement in RMSE, and more conservative behavior under the null. For MN-SRM, our method grants a modest improvement to out-of-sample reconstruction while relaxing the orthonormality constraint of SRM. We also provide a software prototyp- ing tool for MN models that can flexibly reuse noise covariance assumptions and algorithms across models.
Michael Shvartsman, Narayanan Sundaram, Mikio C. Aoi, Adam Charles, Theodore L. Willke, Jonathan D. Cohen 0003
AISTATS5
2018 Out-of-Distribution Detection Using an Ensemble of Self Supervised Leave-Out Classifiers
Apoorv Vyas, Nataraj Jammalamadaka, Dipankar Das 0002, Bharat Kaul, Theodore L. Willke
ECCV (8)6
2018 Capturing Shared and Individual Information in fMRI Data
abstract
Cognitive neuroscience seeks to explain the organization of the brain, but typically focuses on aspects that are shared across people rather than those that vary across individuals. Here, we present a new method for analyzing brain imaging data that captures both shared and individual components of brain activity. Inspired by the shared response model (SRM) and the robust principal components analysis, the robust shared response model (RSRM) aligns functional topographies across humans while preserving a component of sparse, individual activity. Experimental results on adult data showed that RSRM performs as well as or better than SRM, while at the same time capturing reliable markers of individual variability. In a test case of participants with extreme variability, we found that RSRM was able to improve the accuracy more than 60% over SRM for the coding of infant fMRI data.
Javier Turek, Cameron T. Ellis, Lena J. Skalaban, Nicholas B. Turk-Browne, Theodore L. Willke
ICASSP5
2017 A Formal Approach to Modeling the Cost of Cognitive Control
Kayhan Özcimder, Biswadip Dey, Sebastian Musslick, Giovanni Petri, Nesreen K. Ahmed, Theodore L. Willke, Jonathan D. Cohen 0003
CogSci6
2017 A semi-supervised method for multi-subject FMRI functional alignment
abstract
Practical limitations on the duration of individual fMRI scans have led neuroscientist to consider the aggregation of data from multiple subjects. Differences in anatomical structures and functional topographies of brains require aligning data across subjects. Existing functional alignment methods serve as a preprocessing step that allows subsequent statistical methods to learn from the aggregated multi-subject data. Despite their success, current alignment methods do not leverage the labeled data used in the subsequent methods. In this work we propose a semi-supervised scheme that simultaneously learns the alignment and performs the analysis. We derive a specific instance of the scheme using the Shared Response Model for alignment and Multinomial Logistic Regression for classification. In our experiments this method improves the average classification accuracy from 65.5% to 68.5%, and from 5.3% to 6.1% over the independently-trained methods. Furthermore, our method achieves similar prediction with almost half the samples used for alignment.
Javier Turek, Theodore L. Willke, Po-Hsuan Chen, Peter J. Ramadge
ICASSP2
2017 Edge Role Discovery via Higher-Order Structures
Nesreen K. Ahmed, Ryan Rossi, Theodore L. Willke, Rong Zhou 0001
PAKDD (1)3
2017 Graphlet decomposition: framework, algorithms, and applications
Nesreen K. Ahmed, Jennifer Neville, Ryan Rossi, Nick G. Duffield, Theodore L. Willke
Knowl. Inf. Syst.5
2017 On Sampling from Massive Graph Streams
abstract
We propose Graph Priority Sampling ( gps ), a new paradigm for order-based reservoir sampling from massive graph streams. gps provides a general way to weight edge sampling according to auxiliary and/or size variables so as to accomplish various estimation goals of graph properties. In the context of subgraph counting, we show how edge sampling weights can be chosen so as to minimize the estimation variance of counts of specified sets of subgraphs. In distinction with many prior graph sampling schemes, gps separates the functions of edge sampling and subgraph estimation. We propose two estimation frameworks: (1) Post-Stream estimation, to allow gps to construct a reference sample of edges to support retrospective graph queries, and (2) In-Stream estimation, to allow gps to obtain lower variance estimates by incrementally updating the subgraph count estimates during stream processing. Unbiasedness of subgraph estimators is established through a new Martingale formulation of graph stream order sampling, in which subgraph estimators, written as a product of constituent edge estimators, are unbiased, even when computed at different points in the stream. The separation of estimation and sampling enables significant resource savings relative to previous work. We illustrate our framework with applications to triangle and wedge counting. We perform a large-scale experimental study on real-world graphs from various domains and types. gps achieves high accuracy with < 1% error for triangle and wedge counting, while storing a small fraction of the graph with average update times of a few microseconds per edge. Notably, for billion-scale graphs, gps accurately estimates triangle and wedge counts with < 1% error, while storing a small fraction of < 0.01% of the total edges in the graph.
Nesreen K. Ahmed, Nick G. Duffield, Theodore L. Willke, Ryan Rossi
Proc. VLDB Endow.3
2017 Bridging the Gap between HPC and Big Data frameworks
abstract
Apache Spark is a popular framework for data analytics with attractive features such as fault tolerance and interoperability with the Hadoop ecosystem. Unfortunately, many analytics operations in Spark are an order of magnitude or more slower compared to native implementations written with high performance computing tools such as MPI. There is a need to bridge the performance gap while retaining the benefits of the Spark ecosystem such as availability, productivity, and fault tolerance. In this paper, we propose a system for integrating MPI with Spark and analyze the costs and benefits of doing so for four distributed graph and machine learning applications. We show that offloading computation to an MPI environment from within Spark provides 3.1−17.7× speedups on the four sparse applications, including all of the overheads. This opens up an avenue to reuse existing MPI libraries in Spark with little effort.
Michael J. Anderson, Shaden Smith, Narayanan Sundaram, Mihai Capota, Zheguang Zhao, Subramanya Dulloor, Nadathur Satish, Theodore L. Willke
Proc. VLDB Endow.8
2016 Estimation of local subgraph counts
abstract
Graphlets represent small induced subgraphs and are becoming increasingly important for a variety of applications. Despite the importance of the local subgraph (graphlet) counting problem, existing work focuses mainly on counting graphlets globally over the entire graph. These global counts have been used for tasks such as graph classification as well as for understanding and summarizing the fundamental structural patterns in graphs. In contrast, this work proposes an accurate, efficient, and scalable parallel framework for the more challenging problem of counting graphlets locally for a given edge or set of edges. The local graphlet counts provide a topologically rigorous characterization of the local structure surrounding an edge. The aim of this work is to obtain the count of every graphlet of size k for each edge. The framework gives rise to efficient, parallel, and accurate unbiased estimation methods with provable error bounds, as well as exact algorithms for counting graphlets locally. Experiments demonstrate the effectiveness of the proposed exact and estimation methods on various datasets. In particular, the exact methods show strong scaling results (11-16x on 16 cores). Moreover, our estimation framework is accurate with error less than 5% on average.
Nesreen K. Ahmed, Theodore L. Willke, Ryan Rossi
IEEE BigData2
2016 Enabling factor analysis on thousand-subject neuroimaging datasets
abstract
The scale of functional magnetic resonance image data is rapidly increasing as large multi-subject datasets are becoming widely available and high-resolution scanners are adopted. The inherent low-dimensionality of the information in this data has led neuroscientists to consider factor analysis methods to extract and analyze the underlying brain activity. In this work, we consider two recent multi-subject factor analysis methods: the Shared Response Model and the Hierarchical Topographic Factor Analysis. We perform analytical, algorithmic, and code optimization to enable multi-node parallel implementations to scale. Single-node improvements result in 99χ and 2062x speedups on the two methods, and enables the processing of larger datasets. Our distributed implementations show strong scaling of 3.3x and 5.5χ respectively with 20 nodes on real datasets. We demonstrate weak scaling on a synthetic dataset with 1024 subjects, equivalent in size to the biggest fMRI dataset collected until now, on up to 1024 nodes and 32,768 cores.
Michael J. Anderson, Mihai Capota, Javier Turek, Theodore L. Willke, Yida Wang 0003, Po-Hsuan Chen, Jeremy R. Manning, Peter J. Ramadge, Kenneth A. Norman
IEEE BigData5
2016 Real-time full correlation matrix analysis of fMRI data
abstract
Real-time functional magnetic resonance imaging (rtfMRI) is an emerging approach for studying the functioning of the human brain. Computational challenges combined with high data velocity have to this point restricted rtfMRI analyses to studying regions of the brain independently. However, given that neural processing is accomplished via functional interactions among brain regions, neuroscience could stand to benefit from rtfMRI analyses of full-brain interactions. In this paper, we extend such an offline analysis method, full correlation matrix analysis (FCMA), to enable its use in rtfMRI studies. Specifically, we introduce algorithms capable of processing real-time data for all stages of the FCMA machine learning workflow: incremental feature selection, model updating, and real-time classification. We also present an actor-model based distributed system designed to support FCMA and other rtfMRI analysis methods. Experiments show that our system successfully analyzes a stream of brain volumes and returns neurofeedback with less than 180 ms of lag. Our real-time FCMA implementation provides the same accuracy as an optimized offline FCMA toolbox while running 3.6–6.2x faster.
Yida Wang 0003, Bryn Keller, Mihai Capota, Michael J. Anderson, Narayanan Sundaram, Jonathan D. Cohen 0003, Kai Li 0001, Nicholas B. Turk-Browne, Theodore L. Willke
IEEE BigData9
2016 Controlled vs. Automatic Processing: A Graph-Theoretic Approach to the Analysis of Serial vs. Parallel Processing in Neural Network Architectures
Sebastian Musslick, Biswadip Dey, Kayhan Özcimder, Md. Mostofa Ali Patwary, Theodore L. Willke, Jonathan D. Cohen 0003
CogSci5
2016 GraphIn: An Online High Performance Incremental Graph Processing Framework
Dipanjan Sengupta, Narayanan Sundaram, Theodore L. Willke, Jeffrey Young 0001, Matthew Wolf, Karsten Schwan
Euro-Par4
2016 GraphPad: Optimized Graph Primitives for Parallel and Distributed Platforms
abstract
The duality between graphs and matrices means that many common graph analyses can be expressed with primitives such as generalized sparse matrix-vector multiplication (SpMSpV) and sparse matrix-matrix multiplication (SpGEMM). Achieving high performance on these primitives is challenging due to limited arithmetic intensity, irregular memory accesses, and significant network communication requirements in the distributed setting. In this paper we implement four graph applications using GraphPad, our optimized multinode implementations of generalized linear algebra primitives such as SpMSpV and SpGEMM. GraphPad is highly flexible to accommodate multiple data layouts, partitioning strategies, and incorporates communication optimizations. Our performance at scale can exceed that of CombBLAS by up to 40×. In addition to GraphPad's performance in a distributed setting, it is also within 2× the performance of GraphMat, a high performance graph framework on a single node for four out of five benchmarks. We also show our communication optimizations and flexibility are critical for good performance on both HPC clusters and commodity cloud platforms.
Michael J. Anderson, Narayanan Sundaram, Nadathur Satish, Md. Mostofa Ali Patwary, Theodore L. Willke, Pradeep Dubey
IPDPS5
2015 Full correlation matrix analysis of fMRI data on Intel® Xeon Phi™ coprocessors
abstract
Full correlation matrix analysis (FCMA) is an unbiased approach for exhaustively studying interactions among brain regions in functional magnetic resonance imaging (fMRI) data from human participants. In order to answer neuroscientific questions efficiently, we are developing a closed-loop analysis system with FCMA on a cluster of nodes with Intel® Xeon Phi™ coprocessors. Here we propose several ideas for data-driven algorithmic modification to improve the performance on the coprocessor. Our experiments with real datasets show that the optimized single-node code runs 5x-16x faster than the baseline implementation using the well-known Intel® MKL and LibSVM libraries, and that the cluster implementation achieves near linear speedup on 5760 cores.
Yida Wang 0003, Michael J. Anderson, Jonathan D. Cohen 0003, Alexander Heinecke, Kai Li 0001, Nadathur Satish, Narayanan Sundaram, Nicholas B. Turk-Browne, Theodore L. Willke
SC9
2014 How Well Do Graph-Processing Platforms Perform? An Empirical Performance Evaluation and Analysis
abstract
Graph-processing platforms are increasingly used in a variety of domains. Although both industry and academia are developing and tuning graph-processing algorithms and platforms, the performance of graph-processing platforms has never been explored or compared in-depth. Thus, users face the daunting challenge of selecting an appropriate platform for their specific application. To alleviate this challenge, we propose an empirical method for benchmarking graph-processing platforms. We define a comprehensive process, and a selection of representative metrics, datasets, and algorithmic classes. We implement a benchmarking suite of five classes of algorithms and seven diverse graphs. Our suite reports on basic (user-lever) performance, resource utilization, scalability, and various overhead. We use our benchmarking suite to analyze and compare six platforms. We gain valuable insights for each platform and present the first comprehensive comparison of graph-processing platforms.
Marcin Biczak, Ana Lucia Varbanescu, Alexandru Iosup, Claudio Martella, Theodore L. Willke
IPDPS6
2014 Benchmarking graph-processing platforms: a vision
abstract
Processing graphs, especially at large scale, is an increasingly useful activity in a variety of business, engineering, and scientific domains. Already, there are tens of graph-processing platforms, such as Hadoop, Giraph, GraphLab, etc., each with a different design and functionality. For graph-processing to continue to evolve, users have to find it easy to select a graph-processing platform, and developers and system integrators have to find it easy to quantify the performance and other non-functional aspects of interest. However, the state of performance analysis of graph-processing platforms is still immature: there are few studies and, for the few that exist, there are few similarities, and relatively little understanding of the impact of dataset and algorithm diversity on performance. Our vision is to develop, with the help of the performance-savvy community, a comprehensive benchmarking suite for graph-processing platforms. In this work, we take a step in this direction, by proposing a set of seven challenges, summarizing our previous work on performance evaluation of distributed graph-processing platforms, and introducing our on-going work within the SPEC Research Group's Cloud Working Group.
Ana Lucia Varbanescu, Alexandru Iosup, Claudio Martella, Theodore L. Willke
ICPE5
2013 Gunther: Search-Based Auto-Tuning of MapReduce
Guangdeng Liao, Kushal Datta, Theodore L. Willke
Euro-Par3
2013 Towards Machine Learning-Based Auto-tuning of MapReduce
abstract
MapReduce, which is the de facto programming model for large-scale distributed data processing, and its most popular implementation Hadoop have enjoyed widespread adoption in industry during the past few years. Unfortunately, from a performance point of view getting the most out of Hadoop is still a big challenge due to the large number of configuration parameters. Currently these parameters are tuned manually by trial and error, which is ineffective due to the large parameter space and the complex interactions among the parameters. Even worse, the parameters have to be re-tuned for different MapReduce applications and clusters. To make the parameter tuning process more effective, in this paper we explore machine learning-based performance models that we use to auto-tune the configuration parameters. To this end, we first evaluate several machine learning models with diverse MapReduce applications and cluster configurations, and we show that support vector regression model (SVR) has good accuracy and is also computationally efficient. We further assess our auto-tuning approach, which uses the SVR performance model, against the Starfish auto tuner, which uses a cost-based performance model. Our findings reveal that our auto-tuning approach can provide comparable or in some cases better performance improvements than Starfish with a smaller number of parameters. Finally, we propose and discuss a complete and practical end-to-end auto-tuning flow that combines our machine learning-based performance models with smart search algorithms for the effective training of the models and the effective exploration of the parameter space.
Nezih Yigitbasi, Theodore L. Willke, Guangdeng Liao, Dick H. J. Epema
MASCOTS2
2007 Coordinated interaction using reliable broadcast in mobile wireless networks
Theodore L. Willke, Nicholas F. Maxemchuk
Comput. Networks1
2005 Coordinated Interaction Using Reliable Broadcast in Mobile Wireless Networks
Theodore L. Willke, Nicholas F. Maxemchuk
NETWORKING1