Gurbinder Gill

dblp:220/9120 · also Gurbinder Singh Gill · DBLP profile ↗
← Back
14ranked-venue papers
5as first author
2since 2021 · last 2021
0009-0004-0172-004XORCID · corroborated

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

Systems, architecture and hardware · 10 · 3 first-author · 2 since 2021Databases, data management, data science and information retrieval · 3 · 2 first-authorSoftware engineering, systems software and programming languages · 2

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

Computer architecture, parallel and distributed computing, and storage systems
5 papers
Parallel and multicore computing · 26% Distributed systems · 24% GPUs and heterogeneous computing · 16%
Theoretical computer science
1 paper
Distributed computing theory · 61% Graph algorithms and graph theory · 39%
Databases, data mining, and information retrieval
1 paper
Data mining · 50% Graph data management · 50%

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

TopicWeightPapersLastEvidence papers
High-performance computing
large-scale graph processing
0.822020
Single Machine Graph Analytics on Massive Datasets Using Intel Optane DC Persistent Memory · Proc. VLDB Endow. 2020
A Study of Partitioning Policies for Graph Analytics on Large-scale Distributed Platforms · Proc. VLDB Endow. 2018
Distributed systems
distributed graph processing
0.722019
Phoenix: A Substrate for Resilient Distributed Graph Analytics · ASPLOS 2019
Gluon: a communication-optimizing substrate for distributed heterogeneous graph analytics · PLDI 2018
Data mining › pattern mining
graph pattern mining
0.412020
Pangolin: An Efficient and Flexible Graph Mining System on CPU and GPU · Proc. VLDB Endow. 2020
Graph data management
motif counting
0.412020
Pangolin: An Efficient and Flexible Graph Mining System on CPU and GPU · Proc. VLDB Endow. 2020
GPUs and heterogeneous computing
GPU graph processing
0.412020
Pangolin: An Efficient and Flexible Graph Mining System on CPU and GPU · Proc. VLDB Endow. 2020
Memory systems
non-volatile memory
0.412020
Single Machine Graph Analytics on Massive Datasets Using Intel Optane DC Persistent Memory · Proc. VLDB Endow. 2020
Storage systems › out-of-core computation
out-of-core graph processing
0.412020
Single Machine Graph Analytics on Massive Datasets Using Intel Optane DC Persistent Memory · Proc. VLDB Endow. 2020
Distributed systems
fault tolerance
0.412019
Phoenix: A Substrate for Resilient Distributed Graph Analytics · ASPLOS 2019
Graph algorithms and graph theory › centrality
betweenness centrality
0.412019
A round-efficient distributed betweenness centrality algorithm · PPoPP 2019
Distributed computing theory › distributed graph algorithms
CONGEST model
0.412019
A round-efficient distributed betweenness centrality algorithm · PPoPP 2019
Distributed computing theory
distributed graph algorithms
0.412019
A round-efficient distributed betweenness centrality algorithm · PPoPP 2019
GPUs and heterogeneous computing › CPU-GPU heterogeneous computing
CPU-GPU graph processing
0.312018
Gluon: a communication-optimizing substrate for distributed heterogeneous graph analytics · PLDI 2018
Parallel and multicore computing › graph partitioning
distributed graph partitioning
0.312018
A Study of Partitioning Policies for Graph Analytics on Large-scale Distributed Platforms · Proc. VLDB Endow. 2018
Parallel and multicore computing
graph processing
0.312018
A Study of Partitioning Policies for Graph Analytics on Large-scale Distributed Platforms · Proc. VLDB Endow. 2018
Parallel and multicore computing › graph processing
heterogeneous graph processing
0.312018
Gluon: a communication-optimizing substrate for distributed heterogeneous graph analytics · PLDI 2018
Parallel and multicore computing › parallel graph algorithms
shared-memory parallel graph algorithm
0.322020
Single Machine Graph Analytics on Massive Datasets Using Intel Optane DC Persistent Memory · Proc. VLDB Endow. 2020
Pangolin: An Efficient and Flexible Graph Mining System on CPU and GPU · Proc. VLDB Endow. 2020
Graph algorithms and graph theory › shortest path
all-pairs shortest paths
0.112019
A round-efficient distributed betweenness centrality algorithm · PPoPP 2019
Distributed systems
communication optimization
0.112018
A Study of Partitioning Policies for Graph Analytics on Large-scale Distributed Platforms · Proc. VLDB Endow. 2018

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

search space pruning · 0.9isomorphism test elimination · 0.9extend-reduce-filter · 0.9message passing · 0.4distributed-memory algorithm · 0.4checkpointing · 0.4communication optimization · 0.3
YearPublicationVenuePosition
2021 Sandslash: a two-level framework for efficient graph pattern mining
abstract
Graph pattern mining (GPM) is a key building block in diverse applications, including bioinformatics, chemical engineering, social network analysis, recommender systems and security. Existing GPM frameworks either provide high-level interfaces for productivity at the cost of expressiveness or provide low-level interfaces that can express a wide variety of GPM algorithms at the cost of increased programming complexity. Moreover, existing systems lack the flexibility to explore combinations of optimizations to achieve performance competitive with hand-optimized applications.
Xuhao Chen 0001, Roshan Dathathri, Gurbinder Gill, Loc Hoang, Keshav Pingali
ICS3
2021 Distributed Training of Embeddings using Graph Analytics
abstract
Many applications today, such as natural language processing, network and code analysis, rely on semantically embedding objects into low-dimensional fixed-length vectors. Such embeddings naturally provide a way to perform useful downstream tasks, such as identifying relations among objects and predicting objects for a given context. Unfortunately, training accurate embeddings is usually computationally intensive and requires processing large amounts of data. This paper presents a distributed training framework for a class of applications that use Skip-gram-like models to generate embeddings. We call this class Any2Vec and it includes Word2Vec (Gensim), and Vertex2Vec (DeepWalk and Node2Vec) among others. We first formulate Any2Vec training algorithm as a graph application. We then adapt the state-of-the-art distributed graph analytics framework, D-Galois, to support dynamic graph generation and re-partitioning, and incorporate novel communication optimizations. We show that on a cluster of 3248-core hosts our framework GraphAny2Vec matches the accuracy of the state-of-the-art shared-memory implementations of Word2Vec and Vertex2Vec, and gives geo-mean speedups of 12 x and 5 x respectively. Furthermore, GraphAny2Vec is on average 2 x faster than DMTK, the state-of-the-art distributed Word2Vec implementation, on 32 hosts while yielding much better accuracy.
Gurbinder Gill, Roshan Dathathri, Saeed Maleki, Madan Musuvathi, Todd Mytkowicz, Olli Saarikivi
IPDPS1
2020 A Study of Graph Analytics for Massive Datasets on Distributed Multi-GPUs
abstract
There are relatively few studies of distributed GPU graph analytics systems in the literature and they are limited in scope since they deal with small data-sets, consider only a few applications, and do not consider the interplay between partitioning policies and optimizations for computation and communication.In this paper, we present the first detailed analysis of graph analytics applications for massive real-world datasets on a distributed multi-GPU platform and the first analysis of strong scaling of smaller real-world datasets. We use D-IrGL, the state-of-the-art distributed GPU graph analytical framework, in our study. Our evaluation shows that (1) the Cartesian vertex-cut partitioning policy is critical to scale computation out on GPUs even at a small scale, (2) static load imbalance is a key factor in performance since memory is limited on GPUs, (3) device-host communication is a significant portion of execution time and should be optimized to gain performance, and (4) asynchronous execution is not always better than bulk-synchronous execution.
Vishwesh Jatala, Roshan Dathathri, Gurbinder Gill, Loc Hoang, V. Krishna Nandivada, Keshav Pingali
IPDPS3
2020 Pangolin: An Efficient and Flexible Graph Mining System on CPU and GPU
abstract
There is growing interest in graph pattern mining (GPM) problems such as motif counting. GPM systems have been developed to provide unified interfaces for programming algorithms for these problems and for running them on parallel systems. However, existing systems may take hours to mine even simple patterns in moderate-sized graphs, which significantly limits their real-world usability. We present Pangolin , an efficient and flexible in-memory GPM framework targeting shared-memory CPUs and GPUs. Pangolin is the first GPM system that provides high-level abstractions for GPU processing. It provides a simple programming interface based on the extend-reduce-filter model, which allows users to specify application specific knowledge for search space pruning and isomorphism test elimination. We describe novel optimizations that exploit locality, reduce memory consumption, and mitigate the overheads of dynamic memory allocation and synchronization. Evaluation on a 28-core CPU demonstrates that Pangolin outperforms existing GPM frameworks Arabesque, RStream, and Fractal by 49×, 88×, and 80× on average, respectively. Acceleration on a V100 GPU further improves performance of Pangolin by 15× on average. Compared to state-of-the-art hand-optimized GPM applications, Pangolin provides competitive performance with less programming effort.
Xuhao Chen 0001, Roshan Dathathri, Gurbinder Gill, Keshav Pingali
Proc. VLDB Endow.3
2020 Single Machine Graph Analytics on Massive Datasets Using Intel Optane DC Persistent Memory
abstract
Intel Optane DC Persistent Memory (Optane PMM) is a new kind of byte-addressable memory with higher density and lower cost than DRAM. This enables the design of affordable systems that support up to 6TB of randomly accessible memory. In this paper, we present key runtime and algorithmic principles to consider when performing graph analytics on extreme-scale graphs on Optane PMM and highlight principles that can apply to graph analytics on all large-memory platforms. To demonstrate the importance of these principles, we evaluate four existing shared-memory graph frameworks and one out-of-core graph framework on large real-world graphs using a machine with 6TB of Optane PMM. Our results show that frameworks using the runtime and algorithmic principles advocated in this paper (i) perform significantly better than the others and (ii) are competitive with graph analytics frameworks running on production clusters.
Gurbinder Gill, Roshan Dathathri, Loc Hoang, Ramesh Peri, Keshav Pingali
Proc. VLDB Endow.1
2019 Gluon-Async: A Bulk-Asynchronous System for Distributed and Heterogeneous Graph Analytics
abstract
Distributed graph analytics systems for CPUs, like D-Galois and Gemini, and for GPUs, like D-IrGL and Lux, use a bulk-synchronous parallel (BSP) programming and execution model. BSP permits bulk-communication and uses large messages which are supported efficiently by current message transport layers, but bulk-synchronization can exacerbate the performance impact of load imbalance because a round cannot be completed until every host has completed that round. Asynchronous distributed graph analytics systems circumvent this problem by permitting hosts to make progress at their own pace, but existing systems either use global locks and send small messages or send large messages but do not support general partitioning policies such as vertex-cuts. Consequently, they perform substantially worse than bulk-synchronous systems. Moreover, none of their programming or execution models can be easily adapted for heterogeneous devices like GPUs. In this paper, we design and implement a lock-free, non-blocking, bulk-asynchronous runtime called Gluon-Async for distributed and heterogeneous graph analytics. The runtime supports any partitioning policy and uses bulk-communication. We present the bulk-asynchronous parallel (BASP) model which allows the programmer to utilize the runtime by specifying only the abstract communication required. Applications written in this model are compared with the BSP programs written using (1) D-Galois and D-IrGL, the state-of-the-art distributed graph analytics systems (which are bulk-synchronous) for CPUs and GPUs, respectively, and (2) Lux, another (bulk-synchronous) distributed GPU graph analytical system. Our evaluation shows that programs written using BASP-style execution are on average ~1.5x faster than those in D-Galois and D-IrGL on real-world large-diameter graphs at scale. They are also on average ~12x faster than Lux. To the best of our knowledge, Gluon-Async is the first asynchronous distributed GPU graph analytics system.
Roshan Dathathri, Gurbinder Gill, Loc Hoang, Vishwesh Jatala, Keshav Pingali, V. Krishna Nandivada, Hoang-Vu Dang, Marc Snir
PACT2
2019 Phoenix: A Substrate for Resilient Distributed Graph Analytics
abstract
This paper presents Phoenix, a communication and synchronization substrate that implements a novel protocol for recovering from fail-stop faults when executing graph analytics applications on distributed-memory machines. The standard recovery technique in this space is checkpointing, which rolls back the state of the entire computation to a state that existed before the fault occurred. The insight behind Phoenix is that this is not necessary since it is sufficient to continue the computation from a state that will ultimately produce the correct result. We show that for graph analytics applications, the necessary state adjustment can be specified easily by the programmer using a thin API supported by Phoenix. Phoenix has no observable overhead during fault-free execution, and it is resilient to any number of faults while guaranteeing that the correct answer will be produced at the end of the computation. This is in contrast to other systems in this space which may either have overheads even during fault-free execution or produce only approximate answers when faults occur during execution. We incorporated Phoenix into D-Galois, the state-of-the-art distributed graph analytics system, and evaluated it on two production clusters. Our evaluation shows that in the absence of faults, Phoenix is ~24x faster than GraphX, which provides fault tolerance using the Spark system. Phoenix also outperforms the traditional checkpoint-restart technique implemented in D-Galois: in fault-free execution, Phoenix has no observable overhead, while the checkpointing technique has 31% overhead. Furthermore, Phoenix mostly outperforms checkpointing when faults occur, particularly in the common case when only a small number of hosts fail simultaneously.
Roshan Dathathri, Gurbinder Gill, Loc Hoang, Keshav Pingali
ASPLOS2
2019 CuSP: A Customizable Streaming Edge Partitioner for Distributed Graph Analytics
abstract
Graph analytics systems must analyze graphs with billions of vertices and edges which require several terabytes of storage. Distributed-memory clusters are often used for analyzing such large graphs since the main memory of a single machine is usually restricted to a few hundreds of gigabytes. This requires partitioning the graph among the machines in the cluster. Existing graph analytics systems usually come with a built-in partitioner that incorporates a particular partitioning policy, but the best partitioning policy is dependent on the algorithm, input graph, and platform. Therefore, built-in partitioners are not sufficiently flexible. Stand-alone graph partitioners are available, but they too implement only a small number of partitioning policies. This paper presents CuSP, a fast streaming edge partitioning framework which permits users to specify the desired partitioning policy at a high level of abstraction and generates high-quality graph partitions fast. For example, it can partition wdc12, the largest publicly available web-crawl graph, with 4 billion vertices and 129 billion edges, in under 2 minutes for clusters with 128 machines. Our experiments show that it can produce quality partitions 6× faster on average than the state-of-the-art standalone partitioner in the literature while supporting a wider range of partitioning policies.
Loc Hoang, Roshan Dathathri, Gurbinder Gill, Keshav Pingali
IPDPS3
2019 A round-efficient distributed betweenness centrality algorithm
abstract
We present Min-Rounds BC (MRBC), a distributed-memory algorithm in the CONGEST model that computes the betweenness centrality (BC) of every vertex in a directed unweighted n-node graph in O(n) rounds. Min-Rounds BC also computes all-pairs-shortest-paths (APSP) in such graphs. It improves the number of rounds by at least a constant factor over previous results for unweighted directed APSP and for unweighted BC, both directed and undirected.
Loc Hoang, Matteo Pontecorvi, Roshan Dathathri, Gurbinder Gill, Bozhi You, Keshav Pingali, Vijaya Ramachandran
PPoPP4
2018 Abelian: A Compiler for Graph Analytics on Distributed, Heterogeneous Platforms
Gurbinder Gill, Roshan Dathathri, Loc Hoang, Andrew Lenharth, Keshav Pingali
Euro-Par1
2018 A Lightweight Communication Runtime for Distributed Graph Analytics
abstract
Distributed-memory multi-core clusters enable in-memory processing of very large graphs with billions of nodes and edges. Recent distributed graph analytics systems have been built on top of MPI. However, communication in graph applications is very irregular, and each host exchanges different amounts of non-contiguous data with other hosts. MPI does not support such a communication pattern well, and it has limited ability to integrate communication with serialization, deserialization, and graph computation tasks. In this paper, we describe a lightweight communication runtime called LCI that supports a large number of threads on each host and avoids the semantic mismatches between the requirements of graph computations and the communication library in MPI. The implementation of LCI is informed by lessons learnt from two baseline MPI-based implementations. We have successfully integrated LCI with two state-of-the-art graph analytics systems - Gemini and Abelian. LCI improves the latency up to 3.5× for microbenchmarks compared to MPI solutions and improves the end-to-end performance of distributed graph algorithms by up to 2×.
Hoang-Vu Dang, Roshan Dathathri, Gurbinder Gill, Alex Brooks, Nikoli Dryden, Andrew Lenharth, Loc Hoang, Keshav Pingali, Marc Snir
IPDPS3
2018 Gluon: a communication-optimizing substrate for distributed heterogeneous graph analytics
abstract
This paper introduces a new approach to building distributed-memory graph analytics systems that exploits heterogeneity in processor types (CPU and GPU), partitioning policies, and programming models. The key to this approach is Gluon, a communication-optimizing substrate.
Roshan Dathathri, Gurbinder Gill, Loc Hoang, Hoang-Vu Dang, Alex Brooks, Nikoli Dryden, Marc Snir, Keshav Pingali
PLDI2
2018 A Study of Partitioning Policies for Graph Analytics on Large-scale Distributed Platforms
abstract
Distributed-memory clusters are used for in-memory processing of very large graphs with billions of nodes and edges. This requires partitioning the graph among the machines in the cluster. When a graph is partitioned, a node in the graph may be replicated on several machines, and communication is required to keep these replicas synchronized. Good partitioning policies attempt to reduce this synchronization overhead while keeping the computational load balanced across machines. A number of recent studies have looked at ways to control replication of nodes, but these studies are not conclusive because they were performed on small clusters with eight to sixteen machines, did not consider work-efficient data-driven algorithms, or did not optimize communication for the partitioning strategies they studied. This paper presents an experimental study of partitioning strategies for work-efficient graph analytics applications on large KNL and Skylake clusters with up to 256 machines using the Gluon communication runtime which implements partitioning-specific communication optimizations. Evaluation results show that although simple partitioning strategies like Edge-Cuts perform well on a small number of machines, an alternative partitioning strategy called Cartesian Vertex-Cut (CVC) performs better at scale even though paradoxically it has a higher replication factor and performs more communication than Edge-Cut partitioning does. Results from communication micro-benchmarks resolve this paradox by showing that communication overhead depends not only on communication volume but also on the communication pattern among the partitions. These experiments suggest that high-performance graph analytics systems should support multiple partitioning strategies, like Gluon does, as no single graph partitioning strategy is best for all cluster sizes. For such systems, a decision tree for selecting a good partitioning strategy based on characteristics of the computation and the cluster is presented.
Gurbinder Gill, Roshan Dathathri, Loc Hoang, Keshav Pingali
Proc. VLDB Endow.1
2013 Evaluation and enhancement of weather application performance on Blue Gene/Q
abstract
Numerical weather prediction (NWP) models use mathematical models of the atmosphere to predict the weather. Ongoing efforts in the weather and climate community continuously try to improve the fidelity of weather models by employing higher order numerical methods suitable for solving model equations at high resolutions. In realistic weather forecasting scenario, simulating and tracking multiple regions of interest (nests) at fine resolutions is important in understanding the interplay between multiple weather phenomena and for comprehensive predictions. These multiple regions of interest in a simulation can be significantly different in resolution and other modeling parameters. Currently, the weather simulations involving these nested regions process them one after the other in a sequential fashion. There exists a lot of prior work in performance evaluation and optimization of weather models, however most of this work is either limited to simulations involving a single domain or multiple nests with same resolution and model parameters such as model physics options. In this paper, we evaluate and enhance the performance of popular WRF model on IBM Blue Gene/Q system. We consider nested simulations with multiple child domains and study how parameters such as physics options and simulation time steps for child domains affect the computational requirements. We also analyze how such configurations can benefit from parallel execution of the children domains rather than processing them sequentially. We demonstrate that it is important to allocate processors to nested child domains in proportion to the work load associated with them when executing them in parallel. This ensures that the time spent in the different nested simulations is nearly equal, and the nested domains reach the synchronization step with the parent simulation together. Our experimental evaluation using a simple heuristic for allocation of nodes shows that the performance of WRF simulations can be improved by up to 14% by parallel execution of sibling domains with different configuration of domain sizes, temporal resolutions and physics options.
Gurbinder Gill, Vaibhav Saxena, Rashmi Mittal, Thomas George, Yogish Sabharwal, Lalit Dagar
HiPC1