VLDB 2026 Research / reviewers in the wild / expert
Abhishek Chandra
dblp:97/3628
· DBLP profile ↗
66ranked-venue papers
6as first author
14since 2021 · last 2025
—ORCID · conflict
Domains — the database's venue-derived domains; a paper can count in several
Systems, architecture and hardware · 44 · 3 first-author · 10 since 2021Software engineering, systems software and programming languages · 6 · 2 first-authorComputer networks · 4 · 2 first-authorGraphics, computer vision, multimedia, augmented reality and games · 2 · 1 since 2021Applied, interdisciplinary, general and emerging computing · 2Artificial intelligence and machine learning · 1 · 1 since 2021Security and privacy · 1
| Year | Publication | Venue | Position |
|---|---|---|---|
| 2025 | SneakPeek: Data-Aware Model Selection and Scheduling for Inference Serving on the EdgeabstractModern applications increasingly rely on inference serving systems to provide low-latency insights with a diverse set of machine learning models. Existing systems often utilize resource elasticity to scale with demand. However, many applications cannot rely on hardware scaling when deployed at the edge or other resource-constrained environments. In this work, we propose a model selection and scheduling algorithm that implements accuracy scaling to increase efficiency for these more constrained deployments. We show that existing schedulers that make decisions using profiled model accuracy are biased toward the label distribution present in the test dataset. To address this problem, we propose using ML models-which we call SneakPeek models- to dynamically adjust estimates of model accuracy, based on the underlying data. Furthermore, we greedily incorporate inference batching into scheduling decisions to improve throughput and avoid the overhead of swapping models in and out of GPU memory. Our approach employs a new notion of request priority, which navigates the trade-off between attaining high accuracy and satisfying deadlines. Using data and models from three real-world applications, we show that our proposed approaches result in higher-utility schedules and higher accuracy inferences in these hardware-constrained environments. Joel Wolfrath, Daniel Frink, Abhishek Chandra |
SoCC | 3 |
| 2025 | ASTRA: Association, Spatial proximity and Temporal Relevance based Adaptive prefetching for Edge ARabstractMobile Augmented Reality (MAR) applications face performance challenges due to their high computational demands and need for low-latency responses. Traditional approaches like on-device storage or reactive data fetching from the cloud often result in limited augmented reality (AR) experiences. Edge caching, which caches AR objects closer to the user, provides a promising solution. However, existing edge caching approaches do not consider AR-specific features such as AR object sizes, user interactions, user’s field of view and physical location in a coherent manner. This paper investigates how to further optimize edge caching by employing AR-aware prefetching techniques. We present ASTRA, a prefetching framework tailored for mobile augmented reality edge caches. It integrates object associations derived from user interaction patterns with spatial awareness based on the user’s physical location and field of view. This approach employs an association factor per object that considers the recency of object co-access; and a lazy fetching strategy that prioritizes prefetching only when the user is in close proximity to the virtual objects. Furthermore, ASTRA incorporates an adaptive tuning algorithm for minimum support in association rule generation to minimize the computation overhead, making it a distinct and effective solution for enhancing user experience in AR applications by ensuring timely virtual object availability.Through extensive evaluation using both synthetic and real-world workloads, we demonstrate that ASTRA significantly improves cache hit rates compared to current prefetching algorithms, achieving gains in hit rate of upto 35% and end-to-end latency by upto 14%. Further, we demonstrate that the adaptive tuning algorithm that automatically tunes minimum support further improves the hit rate of ASTRA by 10%. Our findings demonstrate the potential of ASTRA to substantially enhance the user experience in MAR applications by ensuring the timely availability of virtual objects. Nikhil Sreekumar, Abhishek Chandra, Jon B. Weissman |
IC2E | 2 |
| 2024 | Neural Oscillators for Generalization of Physics-Informed Machine LearningabstractA primary challenge of physics-informed machine learning (PIML) is its generalization beyond the training domain, especially when dealing with complex physical problems represented by partial differential equations (PDEs). This paper aims to enhance the generalization capabilities of PIML, facilitating practical, real-world applications where accurate predictions in unexplored regions are crucial. We leverage the inherent causality and temporal sequential characteristics of PDE solutions to fuse PIML models with recurrent neural architectures based on systems of ordinary differential equations, referred to as neural oscillators. Through effectively capturing long-time dependencies and mitigating the exploding and vanishing gradient problem, neural oscillators foster improved generalization in PIML tasks. Extensive experimentation involving time-dependent nonlinear PDEs and biharmonic beam equations demonstrates the efficacy of the proposed approach. Incorporating neural oscillators outperforms existing state-of-the-art methods on benchmark problems across various metrics. Consequently, the proposed method improves the generalization capabilities of PIML, providing accurate solutions for extrapolation and prediction beyond the training data. Taniya Kapoor, Abhishek Chandra, Daniel M. Tartakovsky, Hongrui Wang 0001, Alfredo Núñez, Rolf P. B. J. Dollevoet |
AAAI | 2 |
| 2024 | Jingle: IoT-Informed Autoscaling for Efficient Resource Management in Edge ComputingabstractEdge computing is increasingly applied to various systems for its proximity to end-users and data sources. To facilitate the deployment of diverse edge-native applications, container technology has emerged as a favored solution due to its simplicity in development and resource management. However, deploying edge applications at scale can quickly overwhelm edge resources, potentially leading to violations of service-level objectives (SLOs). Scheduling edge containerized applications to meet SLOs while efficiently managing resources is a significant challenge. In this paper, we introduce Jingle, an autoscaler for edge clusters designed to efficiently scale edge-native applications. Jingle utilizes application performance metrics and domain-specific insights collected from IoT devices to construct a hybrid model. This hybrid model combines a predictive-reactive module with a lightweight learning model. We demonstrate Jingle’s effectiveness through a real-world deployment in a classroom setting, managing two edge-native applications across edge configurations. Our experimental results show that Jingle can fulfill SLO requirements while requiring up to 50% fewer containers than a state-of-the-art cloud scheduler, which highlights its resource management efficiency and SLO compliance. Abhishek Chandra, Jon B. Weissman |
CCGrid | 2 |
| 2024 | Leveraging Multi-Modal Data for Efficient Edge Inference ServingabstractReal-time analytics over data streams is often performed on edge devices, which offer privacy guarantees and lower-latency responses compared to centralized processing in the cloud. Data streams originating from sensors, mobile phones, or IoT devices are diverse and span multiple modalities, including RGB videos from cameras, time series data from wearable sensors, and audio signals. Previous research has focused on optimizing the individual analytical tasks associated with each stream, with a special emphasis on deep learning, which is computationally intensive and may be used to analyze video streams, among other things. While advances in deep learning have significantly improved inference accuracy (e.g. for computer vision tasks), state-of-the-art models are not well-suited for edge computing environments. Novel approaches are required to substantially reduce the computational burden, since edge systems are heterogeneous and typically have fewer GPU resources available for inference with deep learning models. We show that leveraging data from multiple modalities can complement or sometimes even replace resource-intensive inference, while maintaining or enhancing accuracy. We present DAISY: a Data-Aware Inference Serving sYstem which leverages multi-modal data to increase inference accuracy by dynamically selecting an appropriate model for each request. We thoroughly evaluate the proposed approach using state-of-the-art models and real-world data, which shows an increase in SLO attainment up to 60%, with a corresponding increase in inference accuracy of 5%. Joel Wolfrath, Anirudh Achanta, Abhishek Chandra |
CCGrid | 3 |
| 2023 | AggFirstJoin: Optimizing Geo-Distributed Joins using Aggregation-Based TransformationsabstractGeo-distributed analytics (GDA) involves processing of data stored across geographically distributed sites. Such analytics involves data transfer over the wide area network (WAN) links. WAN links are highly constrained and heterogeneous in nature, making the data transfer over the WAN slow and costly. To tackle this issue, recent approaches have proposed WAN-aware scheduling and placement of geo-distributed analytics tasks. However, computing joins in a geo-distributed setting remains a challenging problem. In this work, we propose AggFirstJoin, an approach to minimize the cost of geo-distributed joins using a theoretically sound query transformation technique. Our optimization approach takes a combined view of the join and aggregation operations which are often part of the same query and pushes (a transformed) aggregation before join in a manner to produce the same results as the original query. We augment our query transformation technique with a WAN-aware task placement and a Bloom filtering approach to further reduce query execution time and WAN usage respectively. We implement our proposed technique on top of Apache Spark, a popular engine for big data analytics. We extensively evaluate our proposed technique using synthetic, TPC-H and Amplab Big Data benchmark datasets on a real geo-distributed testbed on AWS as well as an emulated testbed. Our evaluations show our proposed technique achieves up to 300x reduction in query execution time and 200x reduction in WAN usage as compared to state-of-the-art GDA techniques. Dhruv Kumar 0001, Sohaib Ahmad, Abhishek Chandra, Ramesh K. Sitaraman |
CCGrid | 3 |
| 2023 | Plexus: Optimizing Join Approximation for Geo-Distributed Data AnalyticsabstractModern applications are increasingly generating and persisting data across geo-distributed data centers or edge clusters rather than a single cloud. This paradigm introduces challenges for traditional query execution due to increased latency when transferring data over wide-area network links. Join queries in particular are heavily affected, due to their large output size and amount of data that must be shuffled over the network. Join sampling---computing a uniform sample from the join results---is a useful technique for reducing resource requirements. However, applying it to a geo-distributed setting is challenging, since acquiring independent samples from each location and joining on the samples does not produce uniform and independent tuples from the join result. To address these challenges, we first generalize an existing join sampling algorithm to the geo-distributed setting. We then present our system, Plexus, which introduces three additional optimizations to further reduce the network overhead and handle network and data heterogeneity: (i) weight approximation, (ii) heterogeneity awareness and (iii) sample prefetching. We evaluate Plexus on a geo-distributed system deployed across multiple AWS regions, with an implementation based on Apache Spark. Using three real-world datasets, we show that Plexus can reduce query latency by up to 80% over the default Spark join implementation on a wide class of join queries without substantially impacting sample uniformity. Joel Wolfrath, Abhishek Chandra |
SoCC | 2 |
| 2022 | Efficient Transmission and Reconstruction of Dependent Data Streams via Edge SamplingabstractData stream processing is an increasingly important topic due to the prevalence of smart devices and the demand for real-time analytics. Geo-distributed streaming systems, where cloud-based queries utilize data streams from multiple distributed devices, face challenges since wide-area network (WAN) bandwidth is often scarce or expensive. Edge computing allows us to address these bandwidth costs by utilizing resources close to the devices, e.g. to perform sampling over the incoming data streams, which trades downstream query accuracy to reduce the overall transmission cost. In this paper, we leverage the fact that correlations between data streams may exist across devices located in the same geographical region. Using this insight, we develop a hybrid edge-cloud system which systematically trades off between sampling at the edge and estimation of missing values in the cloud to reduce traffic over the WAN. We present an optimization framework which computes sample sizes at the edge and systematically bounds the number of samples we can estimate in the cloud given the strength of the correlation between streams. Our evaluation with three real-world datasets shows that compared to existing sampling techniques, our system could provide comparable error rates over multiple aggregate queries while reducing WAN traffic by 27-42%. Joel Wolfrath, Abhishek Chandra |
IC2E | 2 |
| 2022 | Towards Elasticity in Heterogeneous Edge-dense EnvironmentsabstractEdge computing has enabled a large set of emerging edge applications by exploiting data proximity and offloading computation-intensive workloads to nearby edge servers. However, supporting edge application users at scale poses challenges due to limited point-of-presence edge sites and constrained elasticity. In this paper, we introduce a densely-distributed edge resource model that leverages capacity-constrained volunteer edge nodes to support elastic computation offloading. Our model also enables the use of geo-distributed edge nodes to further support elasticity. Collectively, these features raise the issue of edge selection. We present a distributed edge selection approach that relies on client-centric views of available edge nodes to optimize average end-to-end latency, with considerations of system heterogeneity, resource contention and node churn. Elasticity is achieved by fine-grained performance probing, dynamic load balancing, and proactive multi-edge node connections per client. Evaluations are conducted in both real-world volunteer environments and emulated platforms to show how a common edge application, namely AR-based cognitive assistance, can benefit from our approach and deliver low-latency responses to distributed users at scale. Zhiying Liang, Nikhil Sreekumar, Sumanth Kaushik 0001, Abhishek Chandra, Jon B. Weissman |
ICDCS | 5 |
| 2022 | HACCS: Heterogeneity-Aware Clustered Client Selection for Accelerated Federated LearningabstractFederated Learning is a machine learning paradigm where a global model is trained in-situ across a large number of distributed edge devices. While this technique avoids the cost of transferring data to a central location and achieves a strong degree of privacy, it presents additional challenges due to the heterogeneous hardware resources available for training. Furthermore, data is not independent and identically distributed (IID) across all edge devices, resulting in statistical heterogeneity across devices. Due to these constraints, client selection strategies play an important role for timely convergence during model training. Existing strategies ensure that each individual device is included, at least periodically, in the training process. In this work, we propose HACCS, a Heterogeneity-Aware Clustered Client Selection system that identifies and exploits the statistical heterogeneity by representing all distinguishable data distributions instead of individual devices in the training process. HACCS is robust to individual device dropout, provided other devices in the system have similar data distributions. We propose privacy-preserving methods for estimating these client distributions and clustering them. We also propose strategies for leveraging these clusters to make scheduling decisions in a federated learning system. Our evaluation on real-world datasets suggests that our framework can provide 18% −38% reduction in time to convergence compared to the state of the art without any compromise in accuracy. Joel Wolfrath, Nikhil Sreekumar, Dhruv Kumar 0001, Yuanli Wang, Abhishek Chandra |
IPDPS | 5 |
| 2022 | Network Cost-Aware Geo-Distributed Data Analytics SystemabstractMany geo-distributed data analytics (GDA) systems have focused on the network performance-bottleneck: inter-data center network bandwidth to improve performance. Unfortunately, these systems may encounter acost-bottleneck(${\$}$) because they have not considered data transfer cost (${\$}$), one of the most expensive and heterogeneous resources in a multi-cloud environment. In this article, we presentKimchi, a network cost-aware GDA system to meet the cost-performance tradeoff by exploiting data transfer cost heterogeneity to avoid the cost-bottleneck. Kimchi determines cost-aware task placement decisions for scheduling tasks given inputs including data transfer cost, network bandwidth, input data size and locations, and desired cost-performance tradeoff preference. In addition, Kimchi is also mindful of data transfer cost in the presence of dynamics. Kimchi has been applied to two common GDA MapReduce models: synchronous barrier and asynchronous push-based shuffle. A Kimchi prototype has been implemented on Spark, and experiments show that it reduces cost by 5%$\scriptstyle \sim$24% without impacting performance and reduces query execution time by 45%$\scriptstyle \sim$70% without impacting cost compared to other baseline approaches centralized, vanilla Spark, and bandwidth-aware (e.g., Iridium). More importantly, Kimchi allows applications to explore a much richer cost-performance tradeoff space in a multi-cloud environment. Kwangsung Oh, Minmin Zhang, Abhishek Chandra, Jon B. Weissman |
IEEE Trans. Parallel Distributed Syst. | 3 |
| 2021 | DLion: Decentralized Distributed Deep Learning in Micro-CloudsabstractDeep learning (DL) is a popular technique for building models from large quantities of data such as pictures, videos, messages generated from edges devices at rapid pace all over the world. It is often infeasible to migrate large quantities of data from the edges to centralized data center(s) over WANs for training due to privacy, cost, and performance reasons. At the same time, training large DL models on edge devices is infeasible due to their limited resources. An attractive alternative for DL training distributed data is to use micro-clouds---small-scale clouds deployed near edge devices in multiple locations. However, micro-clouds present the challenges of both computation and network resource heterogeneity as well as dynamism. In this paper, we introduce DLion, a new and generic decentralized distributed DL system designed to address the key challenges in micro-cloud environments, in order to reduce overall training time and improve model accuracy. We present three key techniques in DLion: (1) Weighted dynamic batching to maximize data parallelism for dealing with heterogeneous and dynamic compute capacity, (2) Per-link prioritized gradient exchange to reduce communication overhead for model updates based on available network capacity, and (3) Direct knowledge transfer to improve model accuracy by merging the best performing model parameters. We build a prototype of DLion on top of TensorFlow and show that DLion achieves up to 4.2X speedup in an Amazon GPU cluster, and up to 2X speed up and 26% higher model accuracy in a CPU cluster over four state-of-the-art distributed DL systems. Rankyung Hong, Abhishek Chandra |
HPDC | 2 |
| 2021 | On the Future of Cloud EngineeringabstractEver since the commercial offerings of the Cloud started appearing in 2006, the landscape of cloud computing has been undergoing remarkable changes with the emergence of many different types of service offerings, developer productivity enhancement tools, and new application classes as well as the manifestation of cloud functionality closer to the user at the edge. The notion of utility computing, however, has remained constant throughout its evolution, which means that cloud users always seek to save costs of leasing cloud resources while maximizing their use. On the other hand, cloud providers try to maximize their profits while assuring service-level objectives of the cloud-hosted applications and keeping operational costs low. All these outcomes require systematic and sound cloud engineering principles. The aim of this paper is to highlight the importance of cloud engineering, survey the landscape of best practices in cloud engineering and its evolution, discuss many of the existing cloud engineering advances, and identify both the inherent technical challenges and research opportunities for the future of cloud computing in general and cloud engineering in particular. David Bermbach, Abhishek Chandra, Chandra Krintz, Aniruddha S. Gokhale, Aleksander Slominski, Lauritz Thamsen, Everton Cavalcante, Tian Guo 0001, Ivona Brandic, Richard Wolski |
IC2E | 2 |
| 2021 | AggNet: Cost-Aware Aggregation Networks for Geo-distributed Streaming Analytics
Dhruv Kumar 0001, Sohaib Ahmad, Abhishek Chandra, Ramesh K. Sitaraman |
SEC | 3 |
| 2020 | A Network Cost-aware Geo-distributed Data Analytics SystemabstractMany geo-distributed data analytics (GDA) systems have focused on the network performance-bottleneck: interdata center network bandwidth to improve performance. Unfortunately, these systems may encounter a cost-bottleneck ($) because they have not considered data transfer cost ($), one of the most expensive and heterogeneous resources in a multi-cloud environment. In this paper, we present Kimchi, a network cost-aware GDA system to meet the cost-performance tradeoff by exploiting data transfer cost heterogeneity to avoid the cost-bottleneck. Kimchi determines cost-aware task placement decisions for scheduling tasks given inputs including data transfer cost, network bandwidth, input data size and locations, and desired cost-performance tradeoff preference. In addition, Kim- chi is also mindful of data transfer cost in the presence of dynamics. A Kimchi prototype has been implemented on Spark and experiments show that it reduces cost by 14% ~ 24% without impacting performance and reduces query execution time by 45% ~ 70% without impacting cost compared to other baseline approaches centralized, vanilla Spark, and bandwidth-aware (e.g. Iridium). More importantly, Kimchi allows applications to explore a much richer cost-performance tradeoff space in a multi-cloud environment. Kwangsung Oh, Abhishek Chandra, Jon B. Weissman |
CCGRID | 2 |
| 2020 | Position Paper: Towards a Robust Edge-Native Storage SystemabstractEdge environments are generating an increasingly large amount of data due to the proliferation of edge devices. Accommodating this large influx of data at edge servers is a challenging issue. While some data can be processed as it is generated, others must be stored for later access. This paper proposes the features that a new edge-native storage system must possess including support for user mobility and node fluctuation. To motivate this, we first describe several emerging edge applications and their data needs. We then describe the challenges in meeting these needs. We then evaluate an out-of-the-box cloud storage system, Cassandra, to assess it's suitability as an edge storage system due to many edge-friendly features. We determined that while a cloud-based storage system can be ported to the edge meeting some of the challenges, other challenges require new solutions. Based on the challenges and the results of Cassandra case study, we propose a set of design principles for a new edge-native storage system. Nikhil Sreekumar, Abhishek Chandra, Jon B. Weissman |
SEC | 2 |
| 2020 | Poster: Exploiting Data Heterogeneity for Performance and Reliability in Federated LearningabstractFederated Learning [1] enables distributed devices to learn a shared machine learning model together, without uploading their private training data. It has received significant attention recently and has been used in mobile applications such as search suggestion [2] and object detection [3]. Federated Learning is different from distributed machine learning due to the following reasons: 1) System heterogeneity: federated learning is usually performed on devices having highly dynamic and heterogeneous network, compute, and power availability. 2) Data heterogeneity (or statistical heterogeneity): data is produced by different users on different devices, and therefore may have different statistical distribution (non-IID). Yuanli Wang, Dhruv Kumar 0001, Abhishek Chandra |
SEC | 3 |
| 2020 | Poster: Data-Aware Edge Sampling for Aggregate Query ApproximationabstractData stream processing is an increasingly important topic due to the prevalence of smart devices and the demand for realtime analytics. One estimate suggests that we should expect nine smart-devices per person by the year 2025 [1]. These devices generate data which might include sensor readings from a smart home, event or system logs on a device, or video feeds from surveillance cameras. As the number of devices increases, the cost of streaming the device data to the cloud over the wide-area network (WAN) will also increase substantially. Transferring and querying this data efficiently has become the focus of much academic research [2]-[5]. Edge computation affords us the opportunity to address this problem by utilizing resources close to the devices. Edge resources have many different use cases, including minimizing end-to-end latency or maximizing throughput [6], [7]. We restrict our focus to minimizing the required WAN bandwidth, which is an effort to address the increase in data volume. Joel Wolfrath, Abhishek Chandra |
SEC | 2 |
| 2020 | WASP: Wide-area Adaptive Stream ProcessingabstractAdaptability is critical for stream processing systems to ensure stable, low-latency, and high-throughput processing of long-running queries. Such adaptability is particularly challenging for wide-area stream processing due to the highly dynamic nature of the wide-area environment, which includes unpredictable workload patterns, variable network bandwidth, occurrence of stragglers, and failures. Unfortunately, existing adaptation techniques typically achieve these performance goals by compromising the quality/accuracy of the results, and they are often application-dependent. In this work, we rethink the adaptability property of wide-area stream processing systems and propose a resource-aware adaptation framework, called WASP. WASP adapts queries through a combination of multiple techniques: task re-assignment, operator scaling, and query re-planning, and applies them in a WAN-aware manner. It is able to automatically determine which adaptation action to take depending on the type of queries, dynamics, and optimization goals. We have implemented a WASP prototype on Apache Flink. Experimental evaluation with the YSB benchmark and a real Twitter trace shows that WASP can handle various dynamics without compromising the quality of the results. Albert Jonathan, Abhishek Chandra, Jon B. Weissman |
Middleware | 2 |
| 2020 | Optimizing Timeliness and Cost in Geo-Distributed Streaming AnalyticsabstractRapid data streams are generated continuously from diverse sources including users, devices, and sensors located around the globe. This results in the need for efficient geo-distributed streaming analytics to extract timely information. A typical geo-distributed analytics service uses a hub-and-spoke model, comprising multiple edges connected by a wide-area-network (WAN) to a central data warehouse. In this paper, we focus on the widely used primitive of windowed grouped aggregation, and examine the question of how much computation should be performed at the edges versus the center. We develop algorithms to optimize two key metrics: WAN traffic and staleness(delay in getting results). We present a family of optimal offline algorithms that jointly minimize these metrics, and we use these to guide our design of practical online algorithms based on the insight that windowed grouped aggregation can be modeled as a caching problem where the cache size varies overtime. We evaluate our algorithms through an implementation in Apache Storm deployed on PlanetLab. Using workloads derived from anonymized traces of a popular analytics service from a large commercial CDN, our experiments show that our online algorithms achieve near-optimal traffic and staleness for a variety of system configurations, stream arrival rates, and queries. Benjamin Heintz, Abhishek Chandra, Ramesh K. Sitaraman |
IEEE Trans. Cloud Comput. | 2 |
| 2020 | Wiera: Policy-Driven Multi-Tiered Geo-Distributed Cloud Storage SystemabstractMulti-tiered geo-distributed cloud storage systems must tame complexity at many levels: uniform APIs for storage access, supporting flexible storage policies that meet a wide array of application metrics, determining an optimal data placement, handling uncertain network dynamics and access dynamism, and operating across many levels of heterogeneity both within and across data-centers (DCs). In this paper, we present an integrated solution called Wiera. Wiera enables the specification of data management policies both within a local DC and across DCs. Such policies enable the user to optimize for cost, performance, reliability, durability, and consistency, and to express their tradeoffs. In addition, Wiera determines an optimal data placement for the user to meet their desired tradeoffs easily in such an environment. A key aspect of Wiera is first-class support for dynamism due to network, workload, and access patterns changes. As far as we know, Wiera is the first geo-distributed cloud storage system which handles dynamism actively at run-time. Wiera allowsunmodified applicationsto reap the benefits of flexible data/storage policies by externalizing the policy specification. We show how Wiera enables a rich specification of dynamic policies using a concise notation and describe the design and implementation of the system. We have implemented a Wiera prototype on multiple cloud environments, AWS and Azure, that illustrates potential benefits from managing dynamics and in using multiple cloud storage tiers both within and across DCs. Kwangsung Oh, Nan Qin, Abhishek Chandra, Jon B. Weissman |
IEEE Trans. Parallel Distributed Syst. | 3 |
| 2019 | MESH: A Flexible Distributed Hypergraph Processing SystemabstractWith the rapid growth of large online social networks, the ability to analyze large-scale social structure and behavior has become critically important, and this has led to the development of several scalable graph processing systems. In reality, however, social interaction takes place not only between pairs of individuals as in the graph model, but rather in the context of multi-user groups. Research has shown that such group dynamics can be better modeled through a more general hypergraph model, resulting in the need to build scalable hypergraph processing systems. In this paper, we present MESH, a flexible distributed framework for scalable hypergraph processing. MESH provides an easy-to-use and expressive application programming interface that naturally extends the "think like a vertex" model common to many popular graph processing systems. Our framework provides a flexible implementation based on an underlying graph processing system, and enables different design choices for the key implementation issues of partitioning a hypergraph representation. We implement MESH on top of the popular GraphX graph processing framework in Apache Spark. Using a variety of real datasets and experiments conducted on a local 8-node cluster as well as a 65-node Amazon AWS testbed, we demonstrate that MESH provides flexibility based on data and application characteristics, as well as scalability with cluster size. We further show that it is competitive in performance to HyperX, another hypergraph processing system based on Spark, while providing a much simpler implementation (requiring about 5X fewer lines of code), thus showing that simplicity and flexibility need not come at the cost of performance. Benjamin Heintz, Rankyung Hong, Shivangi Singh, Gaurav Khandelwal, Corey Tesdahl, Abhishek Chandra |
IC2E | 6 |
| 2018 | Decentralized Distributed Deep Learning in Heterogeneous WAN EnvironmentsabstractNo abstract available. Rankyung Hong, Abhishek Chandra |
SoCC | 2 |
| 2018 | Multi-Query Optimization in Wide-Area Streaming AnalyticsabstractWide-area data analytics has gained much attention in recent years due to the increasing need for analyzing data that are geographically distributed. Many of such queries often require real-time analysis on data streams that are continuously being generated across multiple locations. Yet, analyzing these geo-distributed data streams in a timely manner is very challenging due to the highly heterogeneous and limited bandwidth availability of the wide-area network (WAN). This paper examines the opportunity of applying multi-query optimization in the context of wide-area streaming analytics, with the goal of utilizing WAN bandwidth efficiently while achieving high throughput and low latency execution. Our approach is based on the insight that many streaming analytics queries often exhibit common executions, whether in consuming a common set of input data or performing the same data processing. In this work, we study different types of sharing opportunities and propose a practical online algorithm that allows streaming analytics queries to share their common executions incrementally. We further address the importance of WAN awareness in applying multi-query optimization. Without WAN awareness, sharing executions in a wide-area environment may lead to performance degradation. We have implemented our WAN-aware multi-query optimization in a prototype implementation based on Apache Flink. Experimental evaluation using Twitter traces on a real wide-area system deployment across geo-distributed EC2 data centers shows that our technique is able to achieve 21% higher throughput while saving WAN bandwidth consumption by 33% compared to a WAN-aware, sharing-agnostic system. Albert Jonathan, Abhishek Chandra, Jon B. Weissman |
SoCC | 2 |
| 2017 | TripS: automated multi-tiered data placement in a geo-distributed cloud environmentabstractExploiting the cloud storage hierarchy both within and across data-centers of different cloud providers empowers Internet applications to choose data centers (DCs) and storage services based on storage needs. However, using multiple storage services across multiple data centers brings a complex data placement problem that depends on a large number of factors including, e.g., desired goals, storage and network characteristics, and pricing policies. In addition, dynamics e.g., changing user locations and access patterns, make it impossible to determine the best data placement statically. In this paper, we present TripS, a lightweight system that considers both data center locations and storage tiers to determine the data placement for geo-distributed storage systems. Such systems make use of TripS by providing inputs including SLA, consistency model, fault tolerance, latency information, and cost information. With given inputs, TripS models and solves the data placement problem using mixed integer linear programming (MILP) to determine data placement. In addition, to adapt quickly to dynamics, we introduce the notion of Target Locale List (TLL), a pro-active approach to avoid expensive re-evaluation of the optimal placement. The TripS prototype is running on Wiera, a policy driven geo-distributed storage system, to show how a storage system can easily utilize TripS for data placement. We evaluate TripS/Wiera on multiple data centers of AWS and Azure. The results show that TripS/Wiera can reduce cost 14.96% ∼ 98.1% based on workloads in comparison with other works' approaches and can handle both short- and long-term dynamics to avoid SLA violations. Kwangsung Oh, Abhishek Chandra, Jon B. Weissman |
SYSTOR | 2 |
| 2017 | Nebula: Distributed Edge Cloud for Data Intensive ComputingabstractCentralized cloud infrastructures have become the popular platforms for data-intensive computing today. However, they suffer from inefficient data mobility due to the centralization of cloud resources, and hence, are highly unsuited for geo-distributed data-intensive applications where the data may be spread at multiple geographical locations. In this paper, we present Nebula: a dispersed edge cloud infrastructure that explores the use of voluntary resources for both computation and data storage. We describe the lightweight Nebula architecture that enables distributed data-intensive computing through a number of optimization techniques including location-aware data and computation placement, replication, and recovery. We evaluate Nebula performance on an emulated volunteer platform that spans over 50 PlanetLab nodes distributed across Europe, and show how a common data-intensive computing framework, MapReduce, can be easily deployed and run on Nebula. We show Nebula MapReduce is robust to a wide array of failures and substantially outperforms other wide-area versions based on emulated existing systems. Albert Jonathan, Mathew Ryden, Kwangsung Oh, Abhishek Chandra, Jon B. Weissman |
IEEE Trans. Parallel Distributed Syst. | 4 |
| 2016 | Trading Timeliness and Accuracy in Geo-Distributed Streaming AnalyticsabstractMany applications must ingest rapid data streams and produce analytics results in near-real-time. It is increasingly common for inputs to such applications to originate from geographically distributed sources. The typical infrastructure for processing such geo-distributed streams follows a hub-and-spoke model, where several edge servers perform partial computation before forwarding results over a wide-area network (WAN) to a central location for final processing. Due to limited WAN bandwidth, it is not always possible to produce exact results. In such cases, applications must either sacrifice timeliness by allowing delayed---i.e., stale---results, or sacrifice accuracy by allowing some error in final results. Benjamin Heintz, Abhishek Chandra, Ramesh K. Sitaraman |
SoCC | 2 |
| 2016 | Wiera: Towards Flexible Multi-Tiered Geo-Distributed Cloud Storage InstancesabstractGeo-distributed cloud storage systems must tame complexity at many levels: uniform APIs for storage access, supporting flexible storage policies that meet a wide array of application metrics, handling uncertain network dynamics and access dynamism, and operating across many levels of heterogeneity both within and across data-centers. In this paper, we present an integrated solution called Wiera. Wiera extends our earlier cloud storage system, Tiera, that is targeted to multi-tiered policy-based single cloud storage, to the wide-area and multiple data-centers (even across different providers). Wiera enables the specification of global data management policies built on top of local Tiera policies. Such policies enable the user to optimize for cost, performance, reliability, durability, and consistency, both within and across data-centers, and to express their tradeoffs. A key aspect of Wiera is first-class support for dynamism due to network, workload, and access patterns changes. Wiera policies can adapt to changes in user workload, poorly performing data tiers, failures, and changes in user metrics (e.g., cost). Wiera allows unmodified applications to reap the benefits of flexible data/storage policies by externalizing the policy specification. As far as we know, Wiera is the first geo-distributed cloud storage system which handles dynamism actively at run-time. We show how Wiera enables a rich specification of dynamic policies using a concise notation and describe the design and implementation of the system. We have implemented a Wiera prototype on multiple cloud environments, AWS and Azure, that illustrates potential benefits from managing dynamics and in using multiple cloud storage tiers both within and across data-centers. Kwangsung Oh, Abhishek Chandra, Jon B. Weissman |
HPDC | 2 |
| 2016 | Awan: Locality-Aware Resource Manager for Geo-Distributed Data-Intensive ApplicationsabstractToday, many organizations need to operate on data that is distributed around the globe. This is inevitable due to the nature of data that is generated in different locations such as video feeds from distributed cameras, log files from distributed servers, and many others. Although centralized cloud platforms have been widely used for data-intensive applications, such systems are not suitable for processing geo-distributed data due to high data transfer overheads. An alternative approach is to use an Edge Cloud which reduces the network cost of transferring data by distributing its computations globally. While the Edge Cloud is attractive for geo-distributed data-intensive applications, extending existing cluster computing frameworks to a wide-area environment must account for locality. We propose Awan : a new locality-aware resource manager for geo-distributed data-intensive applications. Awan allows resource sharing between multiple computing frameworks while enabling high locality scheduling within each framework. Our experiments with the Nebula Edge Cloud on PlanetLab show that Awan achieves up to a 28% increase in locality scheduling which reduces the average job turnaround time by approximately 18% compared to existing cluster management mechanisms. Albert Jonathan, Abhishek Chandra, Jon B. Weissman |
IC2E | 2 |
| 2016 | End-to-End Optimization for Geo-Distributed MapReduceabstractMapReduce has proven remarkably effective for a wide variety of data-intensive applications, but it was designed to run on large single-site homogeneous clusters. Researchers have begun to explore the extent to which the original MapReduce assumptions can be relaxed, including skewed workloads, iterative applications, and heterogeneous computing environments. This paper continues this exploration by applying MapReduce across geo-distributed data over geo-distributed computation resources. Using Hadoop, we show that network and node heterogeneity and the lack of data locality lead to poor performance, because the interaction of MapReduce phases becomes pronounced in the presence of heterogeneous network behavior. To address these problems, we take a two-pronged approach: We first develop a model-driven optimization that serves as an oracle, providing high-level insights. We then apply these insights to design cross-phase optimization techniques that we implement and demonstrate in a real-world MapReduce implementation. Experimental results in both Amazon EC2 and PlanetLab show the potential of these techniques as performance is improved by 7-18 percent depending on the execution environment and application. Benjamin Heintz, Abhishek Chandra, Ramesh K. Sitaraman, Jon B. Weissman |
IEEE Trans. Cloud Comput. | 2 |
| 2015 | Optimizing Grouped Aggregation in Geo-Distributed Streaming AnalyticsabstractLarge quantities of data are generated continuously over time and from disparate sources such as users, devices, and sensors located around the globe. This results in the need for efficient geo-distributed streaming analytics to extract timely information. A typical analytics service in these settings uses a simple hub-and-spoke model, comprising a single central data warehouse and multiple edges connected by a wide-area network (WAN). A key decision for a geo-distributed streaming service is how much of the computation should be performed at the edge versus the center. In this paper, we examine this question in the context of windowed grouped aggregation, an important and widely used primitive in streaming queries. Our work is focused on designing aggregation algorithms to optimize two key metrics of any geo-distributed streaming analytics service: WAN traffic and staleness (the delay in getting the result). Towards this end, we present a family of optimal offline algorithms that jointly minimize both staleness and traffic. Using this as a foundation, we develop practical online aggregation algorithms based on the observation that grouped aggregation can be modeled as a caching problem where the cache size varies over time. This key insight allows us to exploit well known caching techniques in our design of online aggregation algorithms. We demonstrate the practicality of these algorithms through an implementation in Apache Storm, deployed on the PlanetLab testbed. The results of our experiments, driven by workloads derived from anonymized traces of a popular web analytics service offered by a large commercial CDN, show that our online aggregation algorithms perform close to the optimal algorithms for a variety of system configurations, stream arrival rates, and query types. Benjamin Heintz, Abhishek Chandra, Ramesh K. Sitaraman |
HPDC | 2 |
| 2015 | Towards Optimizing Wide-Area Streaming AnalyticsabstractModern analytics services require the analysis of large quantities of data derived from disparate geo-distributed sources. Further, the analytics requirements can be complex, with many applications requiring a combination of both real-time and historical analysis, resulting in complex tradeoffs between cost, performance, and information quality. While the traditional approach to analytics processing is to send all the data to a dedicated centralized location, an alternative approach would be to push all computing to the edge for in-situ processing. We argue that neither approach is optimal for modern analytics requirements. Instead, we examine complex tradeoffs driven by a large number of factors such as application, data, and resource characteristics. We present an empirical study using Planet Lab experiments with beacon data from Akamai's download analytics service. We explore key tradeoffs and their implications for the design of next-generation scalable wide-area analytics. Benjamin Heintz, Abhishek Chandra, Ramesh K. Sitaraman |
IC2E | 2 |
| 2015 | Cloud-Based, User-Centric Mobile Application OptimizationabstractThe abundance of compute and storage resources available in the cloud makes it well-suited to addressing the limitations of mobile devices. We explore the use of cloud infrastructure to optimize content-centric mobile applications, which can have high communication and storage requirements, based on the analysis of user activity. We present two specific optimizations, precaching and prefetching, as well as the design and implementation of a middleware framework that allows mobile application developers to easily utilize these techniques. Our framework is fully generalizable to any content-centric mobile application, a large and growing class of Internet applications. A news aggregation application is used as a case study to evaluate our implementation. We make use of a cosine similarity scheme to identify users with similar interests, which in turn is used to determine what content to prefetch. Various cache algorithms, implemented for our framework, are also considered. A workload trace and simulation are used to measure the performance of the application and framework. We observe a dramatic improvement in application performance due to use of our framework with a reasonable amount of overhead. Our system also significantly outperforms a baseline implementation that performs the same optimizations without taking user activity into account. John Kolb, Prashant Chaudhary, Alexander Schillinger, Abhishek Chandra, Jon B. Weissman |
IC2E | 4 |
| 2014 | Nebula: Distributed Edge Cloud for Data Intensive ComputingabstractCentralized cloud infrastructures have become the de-facto platform for data-intensive computing today. However, they suffer from inefficient data mobility due to the centralization of cloud resources, and hence, are highly unsuited for dispersed-data-intensive applications, where the data may be spread at multiple geographical locations. In this paper, we present Nebula: a dispersed cloud infrastructure that uses voluntary edge resources for both computation and data storage. We describe the lightweight Nebula architecture that enables distributed data-intensive computing through a number of optimizations including location-aware data and computation placement, replication, and recovery. We evaluate Nebula's performance on an emulated volunteer platform that spans over 50 PlanetLab nodes distributed across Europe, and show how a common data-intensive computing framework, MapReduce, can be easily deployed and run on Nebula. We show Nebula MapReduce is robust to a wide array of failures and substantially outperforms other wide-area versions based on a BOINC like model. Mathew Ryden, Kwangsung Oh, Abhishek Chandra, Jon B. Weissman |
IC2E | 3 |
| 2014 | Tiera: towards flexible multi-tiered cloud storage instancesabstractCloud providers offer an array of storage services that represent different points along the performance, cost, and durability spectrum. If an application desires the composite benefits of multiple storage tiers, then it must manage the complexity of different interfaces to these storage services and their diverse policies. We believe that it is possible to provide the benefits of customized tiered cloud storage to applications without compromising simplicity using a lightweight middleware. In this paper, we introduce Tiera, a middleware that enables the provision of multi-tiered cloud storage instances that are easy to specify, flexible, and enable a rich array of storage policies and desired metrics to be realized. Tiera's novelty lies in the first-class support for encapsulated tiered cloud storage, ease of programmability of data management policies, and support for runtime replacement and addition of policies and tiers. Tiera enables an application to realize a desired metric (e.g., low latency or low cost) by selecting different storage services that constitute a Tiera instance, and easily specifying a policy, using event and response pairs, to manage the life cycle of data stored in the instance. We illustrate the benefits of Tiera through a prototype implemented on the Amazon cloud. By deploying unmodified MySQL database engine and a TPC-W Web bookstore application on Tiera, we are able to improve their respective throughputs by 47% -- 125% and 46% -- 69%, over standard deployments. We further show the flexibility of Tiera in achieving different desired application metrics with minimal overhead. Ajaykrishna Raghavan, Abhishek Chandra, Jon B. Weissman |
Middleware | 2 |
| 2014 | Trustworthy Distributed Computing on Social NetworksabstractIn this paper we investigate a new computing paradigm, called SocialCloud, in which computing nodes are governed by social ties driven from a bootstrapping trust-possessing social graph. We investigate how this paradigm differs from existing computing paradigms, such as grid computing and the conventional cloud computing paradigms. We show that incentives to adopt this paradigm are intuitive and natural, and security and trust guarantees provided by it are solid. We propose metrics for measuring the utility and advantage of this computing paradigm, and using real-world social graphs and structures of social traces; we investigate the potential of this paradigm for ordinary users. We study several design options and trade-offs, such as scheduling algorithms, centralization, and straggler handling, and show how they affect the utility of the paradigm. Interestingly, we conclude that whereas graphs known in the literature for high trust properties do not serve distributed trusted computing algorithms, such as Sybil defenses—for their weak algorithmic properties, such graphs are good candidates for our paradigm for their self-load-balancing features. David Mohaisen, Abhishek Chandra, Yongdae Kim |
IEEE Trans. Serv. Comput. | 3 |
| 2013 | Trustworthy distributed computing on social networksabstractWe investigate a new computing paradigm, called SocialCloud, in which computing nodes are governed by social ties driven from a bootstrapping trust-possessing social graph. We investigate how this paradigm differs from existing computing paradigms, such as grid computing and the conventional cloud computing paradigms. We show that incentives to adopt this paradigm are intuitive and natural, and security and trust guarantees provided by it are solid. We propose metrics for measuring the utility and advantage of this computing paradigm, and using real-world social graphs and structures of social traces; we investigate the potential of this paradigm for ordinary users. We study several design options and trade-offs, such as scheduling algorithms, centralization, and straggler handling, and show how they affect the utility of the paradigm. Interestingly, we conclude that whereas graphs known in the literature for high trust properties do not serve distributed trusted computing algorithms, such as Sybil defenses---for their weak algorithmic properties, such graphs are good candidates for our paradigm for their self-load-balancing features. David Mohaisen, Abhishek Chandra, Yongdae Kim |
AsiaCCS | 3 |
| 2013 | Wide-area streaming analytics: distributing the data cubeabstractTo date, much research in data-intensive computing has focused on batch computation. Increasingly, however, it is necessary to derive knowledge from big data streams. As a motivating example, consider a content delivery network (CDN) such as Akamai [4], comprising thousands of servers in hundreds of globally distributed locations. Each of these servers produces a stream of log data, recording for example every user it serves, along with each video stream they access, when they play and pause streams, and more. Each server also records network- and system-level data such as TCP connection statistics. In aggregate, the servers produce billions of lines of log data from over a thousand locations daily. Benjamin Heintz, Abhishek Chandra, Ramesh K. Sitaraman |
SoCC | 2 |
| 2013 | Cross-Phase Optimization in MapReduceabstractMap Reduce has been designed to accommodate large-scale data-intensive workloads running on large single-site homogeneous clusters. Researchers have begun to explore the extent to which the original Map Reduce assumptions can be relaxed including skewed workloads, iterative applications, and heterogeneous computing environments. Our work continues this exploration by applying Map Reduce across widely distributed data over distributed computation resources. This problem arises when datasets are generated at multiple sites as is common in many scientific domains and increasingly e-commerce applications. It also occurs when multi-site resources such as geographically separated data centers are applied to the same Map Reduce job. Using Hadoop, we show that the absence of network and node homogeneity and locality of data lead to poor performance. The problem is that interaction of Map Reduce phases becomes pronounced in the presence of heterogeneous network behavior. In this paper, we propose new cross-phase optimization techniques that enable independent Map Reduce phases to influence one another. We propose techniques that optimize the push and map phases to enable push-map overlap and to allow map behavior to feed back into push dynamics. Similarly, we propose techniques that optimize the map and reduce phases to enable shuffle cost to feed back and affect map scheduling decisions. We evaluate the benefits of our techniques in both Amazon EC2 and Planet Lab. The experimental results show the potential of these techniques as performance is improved from 7%-18% depending on the execution environment and application. Benjamin Heintz, Abhishek Chandra, Jon B. Weissman |
IC2E | 3 |
| 2012 | Sharing-Aware Cloud-Based Mobile OutsourcingabstractMobile devices, such as smart phones and tablets, are becoming the universal interface to online services and applications. However, such devices have limited computational power and battery life, which limits their ability to execute resource-intensive applications. Computation outsourcing to external resources has been proposed as a technique to alleviate this problem. Most existing work on mobile outsourcing has focused on either single application optimization or outsourcing to fixed, local resources, with the assumption that wide-area latency is prohibitively high. However, the opportunity of improving the outsourcing performance by utilizing the relation among multiple applications and optimizing the server provisioning is neglected. In this paper, we present the design and implementation of an Android/Amazon EC2-based mobile application outsourcing framework, leveraging the cloud for scalability, elasticity, and multi-user code/data sharing. Using this framework, we empirically demonstrate that the cloud is not only feasible but desirable as an offloading platform for latency-tolerant applications. We have proposed to use data mining techniques to detect data sharing across multiple applications, and developed novel scheduling algorithms that exploit such data sharing for better outsourcing performance. Additionally, our platform is designed to dynamically scale to support a large number of mobile users concurrently. Experiments show that our proposed techniques and algorithms substantially improve application performance, while achieving high efficiency in terms of computation resource and network usage. Chonglei Mei, Abhishek Chandra, Jon B. Weissman |
IEEE CLOUD | 4 |
| 2012 | Exploiting Spatio-Temporal Tradeoffs for Energy-Aware MapReduce in the CloudabstractMapReduce is a distributed computing paradigm widely used for building large-scale data processing applications. When used in cloud environments, MapReduce clusters are dynamically created using virtual machines (VMs) and managed by the cloud provider. In this paper, we study the energy efficiency problem for such MapReduce clouds. We describe a unique spatio-temporal tradeoff that includes efficient spatial fitting of VMs on servers to achieve high utilization of machine resources, as well as balanced temporal fitting of servers with VMs having similar runtimes to ensure a server runs at a high utilization throughout its uptime. We propose VM placement algorithms that explicitly incorporate these tradeoffs. Further, we propose techniques that dynamically scale MapReduce clusters to further improve energy consumption while ensuring that jobs meet or improve their expected runtimes. Our algorithms achieve energy savings over existing placement techniques, and an additional optimization technique further achieves savings while simultaneously improving job performance. Michael Cardosa, Aameek Singh, Himabindu Pucha, Abhishek Chandra |
IEEE Trans. Computers | 4 |
| 2011 | Exploiting Spatio-temporal Tradeoffs for Energy-Aware MapReduce in the CloudabstractMapReduce is a distributed computing paradigm widely used for building large-scale data processing applications. When used in cloud environments, MapReduce clusters are dynamically created using virtual machines (VMs) and managed by the cloud provider. In this paper, we study the energy efficiency problem for such MapReduce clusters in private cloud environments, that are characterized by repeated, batch execution of jobs. We describe a unique spatio-temporal tradeoff that includes efficient spatial fitting of VMs on servers to achieve high utilization of machine resources, as well as balanced temporal fitting of servers with VMs having similar runtimes to ensure a server runs at a high utilization throughout its uptime. We propose VM placement algorithms that explicitly incorporate these tradeoffs. Our algorithms achieve energy savings over existing placement techniques, and an additional optimization technique further achieves savings while simultaneously improving job performance. Michael Cardosa, Aameek Singh, Himabindu Pucha, Abhishek Chandra |
IEEE CLOUD | 4 |
| 2011 | STEAMEngine: Driving MapReduce provisioning in the cloudabstractMapReduce has gained in popularity as a distributed data analysis paradigm, particularly in the cloud, where MapReduce jobs are run on virtual clusters. The provisioning of MapReduce jobs in the cloud is an important problem for optimizing several user as well as provider-side metrics, such as runtime, cost, throughput, energy, and load. In this paper, we present an intelligent provisioning framework called STEAMEngine that consists of provisioning algorithms to optimize these metrics through a set of common building blocks. These building blocks enable spatio-temporal tradeoffs unique to MapReduce provisioning: along with their resource requirements (spatial component), a MapReduce job runtime (temporal component) is a critical element for any provisioning algorithm. We also describe tw o novel provisioning algorithms - a user-driven performance optimization and a provider-driven energy optimization - that leverage these building blocks. Our experimental results based on an Amazon EC2 cluster and a local Xen/Hadoop cluster show the benefits of STEAMEngine through improvements in performance and energy via the use of these algorithms and building blocks. Michael Cardosa, Piyush Narang, Abhishek Chandra, Himabindu Pucha, Aameek Singh |
HiPC | 3 |
| 2011 | Passive Network Performance Estimation for Large-Scale, Data-Intensive ComputingabstractDistributed computing applications are increasingly utilizing distributed data sources. However, the unpredictable cost of data access in large-scale computing infrastructures can lead to severe performance bottlenecks. Providing predictability in data access is, thus, essential to accommodate the large set of newly emerging large-scale, data-intensive computing applications. In this regard, accurate estimation of network performance is crucial to meeting the performance goals of such applications. Passive estimation based on past measurements is attractive for its relatively small overhead compared to relying on explicit probing. In this paper, we take a passive approach for network performance estimation. Our approach is different from existing passive techniques that rely either on past direct measurements of pairs of nodes or on topological similarities. Instead, we exploit secondhand measurements collected by other nodes without any topological restrictions. In this paper, we present Overlay Passive Estimation of Network performance (OPEN), a scalable framework providing end-to-end network performance estimation based on secondhand measurements, and discuss how OPEN achieves cost-effective estimation in a large-scale infrastructure. Our extensive experimental results show that OPEN estimation can be applicable for replica and resource selections commonly used in distributed computing. Jinoh Kim, Abhishek Chandra, Jon B. Weissman |
IEEE Trans. Parallel Distributed Syst. | 2 |
| 2010 | Starling: Minimizing Communication Overhead in Virtualized Computing Platforms Using Decentralized Affinity-Aware MigrationabstractVirtualization is being widely used in large-scale computing environments, such as clouds, data centers, and grids, to provide application portability and facilitate resource multiplexing while retaining application isolation. In many existing virtualized platforms, it has been found that the network bandwidth often becomes the bottleneck resource, causing both high network contention and reduced performance for communication and data-intensive applications. In this paper, we present a decentralized affinity-aware migration technique that incorporates heterogeneity and dynamism in network topology and job communication patterns to allocate virtual machines on the available physical resources. Our technique monitors network affinity between pairs of VMs and uses a distributed bartering algorithm, coupled with migration, to dynamically adjust VM placement such that communication overhead is minimized. Our experimental results running the Intel MPI benchmark and a scientific application on a 7-node Xen cluster show that we can get up to 42% improvement in the runtime of the application over a no-migration technique, while achieving up to 85% reduction in network communication cost. In addition, our technique is able to adjust to dynamic variations in communication patterns and provides both good performance and low network contention with minimal overhead. Jason D. Sonnek, James B. S. G. Greensky, Robert Reutiman, Abhishek Chandra |
ICPP | 4 |
| 2010 | Resource Bundles: Using Aggregation for Statistical Large-Scale Resource Discovery and ManagementabstractResource discovery is an important process for finding suitable nodes that satisfy application requirements in large loosely coupled distributed systems. Besides internode heterogeneity, many of these systems also show a high degree of intranode dynamism, so that selecting nodes based only on their recently observed resource capacities can lead to poor deployment decisions resulting in application failures or migration overheads. However, most existing resource discovery mechanisms rely mainly on recent observations to achieve scalability in large systems. In this paper, we propose the notion of a resource bundle-a representative resource usage distribution for a group of nodes with similar resource usage patterns-that employs two complementary techniques to overcome the limitations of existing techniques: resource usage histograms to provide statistical guarantees for resource capacities and clustering-based resource aggregation to achieve scalability. Using trace-driven simulations and data analysis of a month-long PlanetLab trace, we show that resource bundles are able to provide high accuracy for statistical resource discovery, while achieving high scalability. We also show that resource bundles are ideally suited for identifying group-level characteristics (e.g., hot spots, total group capacity). To automatically parameterize the bundling algorithm, we present an adaptive algorithm that can detect online fluctuations in resource heterogeneity. Michael Cardosa, Abhishek Chandra |
IEEE Trans. Parallel Distributed Syst. | 2 |
| 2009 | Extracting the textual and temporal structure of supercomputing logsabstractSupercomputers are prone to frequent faults that adversely affect their performance, reliability and functionality. System logs collected on these systems are a valuable resource of information about their operational status and health. However, their massive size, complexity, and lack of standard format makes it difficult to automatically extract information that can be used to improve system management. In this work we propose a novel method to succinctly represent the contents of supercomputing logs, by using textual clustering to automatically find the syntactic structures of log messages. This information is used to automatically classify messages into semantic groups via an online clustering algorithm. Further, we describe a methodology for using the temporal proximity between groups of log messages to identify correlated events in the system. We apply our proposed methods to two large, publicly available supercomputing logs and show that our technique features nearly perfect accuracy for online log-classification and extracts meaningful structural and temporal message patterns that can be used to improve the accuracy of other log analysis techniques. Sourabh Jain, Inderpreet Singh, Abhishek Chandra, Zhi-Li Zhang, Greg Bronevetsky |
HiPC | 3 |
| 2009 | HiDRA: Statistical multi-dimensional resource discovery for large-scale systemsabstractResource discovery enables applications deployed in heterogeneous large-scale distributed systems to find resources that meet QoS requirements. In particular, most applications need resource requirements to be satisfied simultaneously for multiple resources (such as CPU, memory and network bandwidth). Due to dynamism in many large-scale systems, providing statistical guarantees on such requirements is important to avoid application failures and overheads. However, existing techniques either provide guarantees only for individual resources, or take a static or memoryless approach along multiple dimensions. We present HiDRA, a scalable resource discovery technique providing statistical guarantees for resource requirements spanning multiple dimensions simultaneously. Through trace analysis and a 307-node PlanetLab implementation, we show that HiDRA, while using over 1,400 times less data, performs nearly as well as a fully-informed algorithm, showing better precision and having recall within 3%. We demonstrate that HiDRA is a feasible, low-overhead approach to statistical resource discovery in a distributed system. Michael Cardosa, Abhishek Chandra |
IWQoS | 2 |
| 2009 | Using Data Accessibility for Resource Selection in Large-Scale Distributed SystemsabstractLarge-scale distributed systems provide an attractive scalable infrastructure for network applications. However, the loosely coupled nature of this environment can make data access unpredictable, and in the limit, unavailable. We introduce the notion of accessibility to capture both availability and performance. An increasing number of data-intensive applications require not only considerations of node computation power but also accessibility for adequate job allocations. For instance, selecting a node with intolerably slow connections can offset any benefit to running on a fast node. In this paper, we present accessibility-aware resource selection techniques by which it is possible to choose nodes that will have efficient data access to remote data sources. We show that the local data access observations collected from a node's neighbors are sufficient to characterize accessibility for that node. By conducting trace-based, synthetic experiments on PlanetLab, we show that the resource selection heuristics guided by this principle significantly outperform conventional techniques such as latency-based or random allocations. The suggested techniques are also shown to be stable even under churn despite the loss of prior observations. Jinoh Kim, Abhishek Chandra, Jon B. Weissman |
IEEE Trans. Parallel Distributed Syst. | 2 |
| 2008 | Resource Bundles: Using Aggregation for Statistical Wide-Area Resource Discovery and AllocationabstractResource discovery is an important process for finding suitable nodes that satisfy application requirements in large loosely-coupled distributed systems. Besides inter-node heterogeneity, many of these systems also show a high degree of intra-node dynamism, so that selecting nodes based only on their recently observed resource capacities for scalability reasons can lead to poor deployment decisions resulting in application failures or migration overheads. In this paper, we propose the notion of a resource bundle - a representative resource usage distribution for a group of nodes with similar resource usage patterns - that employs two complementary techniques to overcome the limitations of existing techniques: resource usage histograms to provide statistical guarantees for resource capacities, and clustering-based resource aggregation to achieve scalability. Using trace-driven simulations and data analysis of a month-long Planet Lab trace, we show that resource bundles are able to provide high accuracy for statistical resource discovery (up to 56% better precision than using only recent values), while achieving high scalability (up to 55% fewer messages than a non-aggregation algorithm). We also show that resource bundles are ideally suited for identifying group-level characteristics such as finding load hot spots and estimating total group capacity (within 8% of actual values). Michael Cardosa, Abhishek Chandra |
ICDCS | 2 |
| 2008 | Accessibility-Based Resource Selection in Loosely-Coupled Distributed SystemsabstractLarge-scale distributed systems provide an attractive scalable infrastructure for network applications. However,the loosely-coupled nature of this environment can make data access unpredictable, and in the limit, unavailable. We introduce the notion of accessibility to capture both availability and performance. An increasing number of data intensive applications require not only considerations of node computation power but also accessibility for adequate job allocations. For instance, selecting a node with intolerably slow connections can offset any benefit to running on a fast node. In this paper, we present accessibility-aware resource selection techniques by which it is possible to choose nodes that will have efficient data access to remote data sources. We show that the local data access observations collected from a node's neighbors are sufficient to characterize accessibility for that node. We then present resource selection heuristics guided by this principle, and show that they significantly out perform standard techniques. The suggested techniques are also shown to be stable even under churn despite the loss of prior observations. Jinoh Kim, Abhishek Chandra, Jon B. Weissman |
ICDCS | 2 |
| 2008 | Exploring the throughput-fairness tradeoff of deadline scheduling in heterogeneous computing environmentsabstractThe scalability and computing power of large-scale computational platforms has made them attractive for hosting compute-intensive time-critical applications. Many of these applications are composed of computational tasks that require specific deadlines to be met for successful completion. In this paper, we show that combining redundant scheduling with deadline-based scheduling in these systems leads to a fundamental tradeoff between throughput and fairness. We propose a new scheduling algorithm called Limited Resource Earliest Deadline (LRED) that couples redundant scheduling with deadline-driven scheduling in a flexible way by using a simple tunable parameter to exploit this tradeoff. Our evaluation of LRED shows that LRED provides a powerful mechanism to achieve desired throughput or fairness under high loads and low timeliness environments. Vasumathi Sundaram, Abhishek Chandra, Jon B. Weissman |
SIGMETRICS | 2 |
| 2008 | Agile dynamic provisioning of multi-tier Internet applicationsabstractDynamic capacity provisioning is a useful technique for handling the multi-time-scale variations seen in Internet workloads. In this article, we propose a novel dynamic provisioning technique for multi-tier Internet applications that employs (1) a flexible queuing model to determine how much of the resources to allocate to each tier of the application, and (2) a combination of predictive and reactive methods that determine when to provision these resources, both at large and small time scales. We propose a novel data center architecture based on virtual machine monitors to reduce provisioning overheads. Our experiments on a forty-machine Xen/Linux-based hosting platform demonstrate the responsiveness of our technique in handling dynamic workloads. In one scenario where a flash crowd caused the workload of a three-tier application to double, our technique was able to double the application capacity within five minutes, thus maintaining response-time targets. Our technique also reduced the overhead of switching servers across applications from several minutes to less than a second, while meeting the performance targets of residual sessions. Bhuvan Urgaonkar, Prashant J. Shenoy, Abhishek Chandra, Pawan Goyal 0001, Timothy Wood 0001 |
ACM Trans. Auton. Adapt. Syst. | 3 |
| 2008 | Hierarchical Scheduling for Symmetric MultiprocessorsabstractHierarchical scheduling has been proposed as a scheduling technique to achieve aggregate resource partitioning among related groups of threads and applications in uniprocessor and packet scheduling environments. Existing hierarchical schedulers are not easily extensible to multiprocessor environments because 1) they do not incorporate the inherent parallelism of a multiprocessor system while resource partitioning and 2) they can result in unbounded unfairness or starvation if applied to a multiprocessor system in a naive manner. In this paper, we present hierarchical multiprocessor scheduling (H-SMP), a novel hierarchical CPU scheduling algorithm designed for a symmetric multiprocessor (SMP) platform. The novelty of this algorithm lies in its combination of space and time multiplexing to achieve the desired bandwidth partition among the nodes of the hierarchical scheduling tree. This algorithm is also characterized by its ability to incorporate existing proportional-share algorithms as auxiliary schedulers to achieve efficient hierarchical CPU partitioning. In addition, we present a generalized weight feasibility constraint that specifies the limit on the achievable CPU bandwidth partitioning in a multiprocessor hierarchical framework and propose a hierarchical weight readjustment algorithm designed to transparently satisfy this feasibility constraint. We evaluate the properties of H-SMP using hierarchical surplus fair scheduling (H-SFS), an instantiation of H-SMP that employs surplus fair scheduling (SFS) as an auxiliary algorithm. This evaluation is carried out through a simulation study that shows that H-SFS provides better fairness properties in multiprocessor environments as compared to existing algorithms and their naive extensions. Abhishek Chandra, Prashant J. Shenoy |
IEEE Trans. Parallel Distributed Syst. | 1 |
| 2007 | Exploiting Heterogeneity for Collective Data Downloading in Volunteer-based NetworksabstractScientific computing is being increasingly deployed over volunteer-based distributed computing environments consisting of idle resources on donated user machines. A fundamental challenge in these environments is the dissemination of data to the computation nodes, with the successful completion of jobs being driven by the efficiency of collective data download across compute nodes, and not only the individual download times. This paper considers the use of a data network consisting of data distributed across a set of data servers, and focuses on the server selection problem: how do individual nodes select a server for downloading data to minimize the communication makespan - the maximal download time for a data file. Through experiments conducted on a pastry network running on PlanetLab, we demonstrate that nodes in a volunteer-based network are heterogeneous in terms of several metrics, such as bandwidth, load, and capacity, which impact their download behavior. We propose new server selection heuristics that incorporate these metrics, and demonstrate that these heuristics outperform traditional proximity-based server selection, reducing average makespans by at least 30%. We further show that incorporating information about download concurrency avoids overloading servers, and improves performance by about 17-43% over heuristics considering only proximity and bandwidth. Jinoh Kim, Abhishek Chandra, Jon B. Weissman |
CCGRID | 2 |
| 2007 | Ridge: combining reliability and performance in open grid platformsabstractLarge-scale donation-based distributed infrastructures need to cope with the inherent unreliability of participant nodes. A widely-used work scheduling technique in such environments is to redundantly schedule the out sourced computations to a number of nodes. We present the design and implementation of RIDGE, a reliability aware system which uses a node's prior performance and behavior to make more effective scheduling decisions. We have implemented RIDGE on top of the BOINC distributed computing infrastructure and have evaluated its performance on a live test bed consisting of 120 PlanetLab nodes. Our experimental results show that RIDGE is able to match or surpass the throughput of the best vanilla BOINC configuration under different reliability environments, by automatically adapting to the characteristics of the underlying environment. In addition, RIDGE is able to provide much lower work unit makes pans compared to BOINC, which indicates its desirability in service-oriented environments with time constraints. Krishnaveni Budati, Jason D. Sonnek, Abhishek Chandra, Jon B. Weissman |
HPDC | 3 |
| 2007 | NGS: Service Adaptation in Open Grid PlatformsabstractLarge-scale donation-based distributed infrastructures need to cope with the inherent unreliability of participant nodes. A widely-used work scheduling technique in such environments is to redundantly schedule the outsourced computations to a number of nodes. We present the design and implementation of RIDGE, a reliability-aware system which uses a node's prior performance and behavior to make more effective scheduling decisions. We have implemented RIDGE on top of the BOINC distributed computing infrastructure and have evaluated its performance on a live PlanetLab testbed. Our experimental results show that RIDGE is able to match or surpass the throughput of the best BOINC configuration by automatically adapting to the characteristics of the underlying environment. In addition, RIDGE is able to provide much lower workunit makespans compared to BOINC. RIDGE is also able to produce significantly lower communication makespans for downloading clients. Collectively, the results suggest that RIDGE has great promise for service-oriented environments with time constraints. Krishnaveni Budati, Jinoh Kim, Abhishek Chandra, Jon B. Weissman |
IPDPS | 3 |
| 2007 | Adaptive Reputation-Based Scheduling on Unreliable Distributed InfrastructuresabstractThis paper addresses the inherent unreliability and instability of worker nodes in large-scale donation-based distributed infrastructures such as peer-to-peer and grid systems. We present adaptive scheduling techniques that can mitigate this uncertainty and significantly outperform current approaches. In this work, we consider nodes that execute tasks via donated computational resources and may behave erratically or maliciously. We present a model in which reliability is not a binary property, but a statistical one based on a node's prior performance and behavior. We use this model to construct several reputation-based scheduling algorithms that employ estimated reliability ratings of worker nodes for efficient task allocation. Our scheduling algorithms are designed to adapt to changing system conditions, as well as nonstationary node reliability. Through simulation, we demonstrate that our algorithms can significantly improve throughput while maintaining a very high success rate of task completion. Our results suggest that reputation-based scheduling can handle a wide variety of worker populations, including nonstationary behavior, with overhead that scales well with system size. We also show that our adaptation mechanism allows the application designer fine-grain control over the desired performance metrics. Jason D. Sonnek, Abhishek Chandra, Jon B. Weissman |
IEEE Trans. Parallel Distributed Syst. | 2 |
| 2006 | Reputation-Based Scheduling on Unreliable Distributed InfrastructuresabstractThis paper presents a design and analysis of scheduling techniques to cope with the inherent unreliability and instability of worker nodes in large-scale donation-based distributed infrastructures such as P2P and Grid systems. In particular, we focus on nodes that execute tasks via donated computational resources and may behave erratically or maliciously. We present a model in which reliability is not a binary property but a statistical one based on a node’s prior performance and behavior. We use this model to construct several reputation-based scheduling algorithms that employ estimated reliability ratings of worker nodes for efficient task allocation. Through simulation of a BOINC-like distributed computing infrastructure, we demonstrate that our algorithms can significantly improve throughput, while maintaining a very high success rate of task completion. Jason D. Sonnek, Mukesh Nathan, Abhishek Chandra, Jon B. Weissman |
ICDCS | 3 |
| 2006 | eSENSE: energy efficient stochastic sensing framework scheme for wireless sensor platformsabstractEnergy is a precious resource in wireless sensor networks as sensor nodes are typically powered by batteries with high replacement cost. This paper presents eSENSE: an energy-efficient stochastic sensing framework for wireless sensor platforms. eSENSE is a node-level framework that utilizes knowledge of the underlying data streams as well as application data quality requirements to conserve energy on a sensor node. eSENSE employs a stochastic scheduling algorithm to dynamically control the operating modes of the sensor node components. This scheduling algorithm enables an adaptive sampling strategy that aggressively conserves power by adjusting sensing activity to the application requirements. Using experimental results obtained on Power-TOSSIM with a real-world data trace, we demonstrate that our approach reduces energy consumption by 29-36% while providing strong statistical guarantees on data quality. Abhishek Chandra, Jaideep Srivastava |
IPSN | 2 |
| 2006 | An observation-based approach towards self-managing web servers
Abhishek Chandra, Prashant Pradhan, Renu Tewari, Sambit Sahu, Prashant J. Shenoy |
Comput. Commun. | 1 |
| 2003 | Dynamic Resource Allocation for Shared Data Centers Using Online Measurements
Abhishek Chandra, Weibo Gong, Prashant J. Shenoy |
IWQoS | 1 |
| 2003 | Dynamic resource allocation for shared data centers using online measurementsabstractNo abstract available. Abhishek Chandra, Weibo Gong, Prashant J. Shenoy |
SIGMETRICS | 1 |
| 2001 | Scalability of Linux Event-Dispatch Mechanisms
Abhishek Chandra, David Mosberger |
USENIX ATC, General Track | 1 |
| 2000 | Application performance in the QLinux multimedia operating systemabstractIn this paper, we argue that conventional operating systems need to be enhanced with predictable resource management mechanisms to meet the diverse performance requirements of emerging multimedia and web applications. We present QLinux—a multimedia operating system based on the Linux kernel that meets this requirement. QLinux employs hierarchical schedulers for fair, predictable allocation of processor, disk and network bandwidth, and accounting mechanisms for appropriate charging of resource usage. We experimentally evaluate the efficacy of these mechanisms using benchmarks and real-world applications. Our experimental results show that (i) emerging applications can indeed benefit from predictable allocation of resources, and (ii) the overheads imposed by the resource allocation mechanisms in QLinux are small. For instance, we show that the QLinux CPU scheduler can provide predictable performance guarantees to applications such as web servers and MPEG players, albeit at the expense of increasing the scheduling overhead. We conclude from our experiments that the benefits due to the resource management mechanisms in QLinux outweigh their increased overheads, making them a practical choice for conventional operating systems. Vijay Sundaram, Abhishek Chandra, Pawan Goyal 0001, Prashant J. Shenoy, Jasleen Sahni, Harrick M. Vin |
ACM Multimedia | 2 |
| 2000 | Surplus Fair Scheduling: A Proportional-Share CPU Scheduling Algorithm for Symmetric Multiprocessors
Abhishek Chandra, Micah Adler, Pawan Goyal 0001, Prashant J. Shenoy |
OSDI | 1 |