EDBT 2026 Demo / reviewers in the wild / expert
Indranil Gupta
dblp:02/6135
· DBLP profile ↗
116ranked-venue papers
9as first author
20since 2021 · last 2026
0000-0002-9372-5937ORCID · corroborated
Domains — the database's venue-derived domains; a paper can count in several
Systems, architecture and hardware · 40 · 7 first-author · 7 since 2021Computer networks · 27 · 6 since 2021Software engineering, systems software and programming languages · 12Security and privacy · 8 · 2 first-authorDatabases, data management, data science and information retrieval · 8 · 3 since 2021Artificial intelligence and machine learning · 6 · 1 since 2021Graphics, computer vision, multimedia, augmented reality and games · 4 · 1 since 2021Human-computer interaction and ubiquitous computing · 4 · 2 since 2021Theory of computation · 2Applied, interdisciplinary, general and emerging computing · 1
| Year | Publication | Venue | Position |
|---|---|---|---|
| 2026 | Counting How the Seconds Count: Understanding TikTok Behavior via ML-driven Analysis of Video ContentabstractShort video streaming systems such as TikTok, YouTube Shorts, Instagram Reels, etc., have reached billions of active users worldwide. At the core of such systems are (proprietary) recommendation algorithms which recommend a sequence of videos to each user, in a personalized way. We aim to understand the temporal evolution of recommendations made by such algorithms, as well as the interplay between the recommendations and user experience. While past work has studied recommendation algorithms using textual data (e.g., titles, hashtags, etc.) as well as user studies and interviews, we add a third modality of analysis—we perform automated analysis of the videos themselves. To perform such multimodal analysis, we develop a new HCI measurement approach that starts with our new tool called VCA (Video Content Analysis) that leverages recent advances in Vision Language Models (VLMs). We apply VCA on a trifecta of HCI methodologies—real user studies, interviews, and data donation. This allows us to understand temporal aspects of how well TikTok’s recommendation algorithm is perceived by users, is affected by user interactions, and aligns with user history; how users are sensitive to the order of videos recommended; and how the algorithm’s effectiveness itself may be predictable in the future. While it is not our goal to reverse-engineer TikTok’s recommendation algorithm, our new findings indicate behavioral aspects that the TikTok user community can benefit from. Maleeha Masood, Shreya Kannan, Zikun Liu 0002, Deepak Vasisht, Indranil Gupta |
CHI | 5 |
| 2026 | Control in Context: How Smart Home Users Navigate Proxy-based SchemesabstractA homeowner controls their smart home devices along a spectrum of approaches, ranging from physical device control to various proxy-based control modalities. This paper studies how and why users move along this spectrum in their day-to-day lives, building upon existing research that focused only on specific interactions. We surveyed smart home owners (N = 43 users), and conducted follow-up interviews with a subset of the survey participants (N = 8). Our studies allow us to both distill specific contexts and experiences of smart home owners as they navigate the control spectrum, as well as to describe how their experiences (both positive and negative) shape their tendencies to control devices in a particular way. These insights lead us to propose practical implications for designers and researchers of smart home management systems, including the need to support flexible control scheme transitions, reduce switching costs, and account for temporal and spatial heterogeneity in the evaluation and design of control systems. Ali Zaidi, Anna Karanika, Ti-Chung Cheng, Yi-Shyuan Chiang, Camille Cobb, Indranil Gupta, Karrie Karahalios |
CHI | 6 |
| 2026 | RASC: Enhancing Observability & Programmability in Smart Spaces
Anna Karanika, Kai-Siang Wang, Han-Ting Liang, Shalni Sundram, Indranil Gupta |
NSDI | 5 |
| 2026 | There is More Control in Egalitarian Edge IoT MeshesabstractWhile mesh networking for edge settings (e.g., smart buildings, farms, battlefields, etc.) has received much attention, the layer of control over such meshes remains largely centralized and cloud-based. This paper focuses on applications with commonplace sense-trigger-actuate (STA) workloads—like the abstraction of routines popular now in smart homes, but applied to larger-scale edge IoT deployments. We present CoMesh, which tackles the challenge of building a decentralized mesh-based control plane for local, non-cloud, and hubless management of sense-trigger-actuate applications. CoMesh builds atop an abstraction called the coterie, which spreads STA load in a finegrained way both across space and across time. A coterie uses a novel combination of techniques such as zero-message-exchange protocols (for fast proactive member selection), quorum-based agreement, and locality-sensitive hashing. We analyze and theoretically prove safety and liveness properties of CoMesh. Our evaluation with both a Raspberry Pi-4 deployment and largerscale simulations, using real building maps and real routine workloads, shows that CoMesh is load-balanced, fast, faulttolerant, and scalable. Anna Karanika, Rui Yang 0034, Xiaojuan Ma, Jiangran Wang, Shalni Sundram, Indranil Gupta |
IEEE Trans. Netw. Serv. Manag. | 6 |
| 2025 | CPU-Limits kill Performance: Time to rethink Resource Control
Chirag C. Shetty, Sarthak Chakraborty, Hubertus Franke, Larisa Shwartz, Chandrasekhar Narayanaswami 0001, Indranil Gupta, Saurabh Jha |
SoCC | 6 |
| 2025 | A House United Within Itself: SLO-Awareness for On-Premises Containerized ML Inference Clusters via FaroabstractThis paper tackles the challenge of running multiple ML inference jobs (models) under time-varying workloads, on a constrained on-premises production cluster. Our system Faro takes in latency Service Level Objectives (SLOs) for each job, auto-distills them into utility functions, "sloppifies" these utility functions to make them amenable to mathematical optimization, automatically predicts workload via probabilistic prediction, and dynamically makes implicit cross-job resource allocations, in order to satisfy cluster-wide objectives, e.g., total utility, fairness, and other hybrid variants. A major challenge Faro tackles is that using precise utilities and high-fidelity predictors, can be too slow (and in a sense too precise!) for the fast adaptation we require. Faro's solution is to "sloppify" (relax) its multiple design components to achieve fast adaptation without overly degrading solution quality. Faro is implemented in a stack consisting of Ray Serve running atop a Kubernetes cluster. Trace-driven cluster deployments show that Faro achieves 2.3×-23× lower SLO violations compared to state-of-the-art systems. Beomyeol Jeon, Chen Wang 0039, Diana Arroyo, Alaa Youssef, Indranil Gupta |
EuroSys | 5 |
| 2025 | Wainscot: Tailoring Model Parallelism to Fit Device Memory LimitsabstractWith increasing sizes of DNN (Deep Neural Network) models making them exceed the memory of a single device (GPU), model parallelism-based training has become paramount, splitting a model across multiple devices. Unfortunately, today’s model parallelism approaches often result in memory-unbalanced allocations of the model across the multiple devices, with some devices’ memory heavily utilized while others remain underused. This imbalance limits deployments from reaching high batch sizes, triggers Out of Memory (OOM) errors earlier, and underutilizes resources. We present Wainscot, a model parallelism solution that produces memory-balanced placements of a DNN model across multiple devices, without noticeable increases in step time. We explore and empirically compare different granularities of rebalancing: operators, operator groups, and subgraphs. Experiments with diverse DNNs across a wide range of batch sizes demonstrate that compared to state-of-the-art model parallelism systems, Wainscot reduces maximal peak memory (across devices) by 47.94%, with a modest increase of 1.62% in step time. Xiaojuan Ma, Shashwat Jaiswal, Chirag C. Shetty, Chen-Wei Chou, Indranil Gupta |
ICDCS | 6 |
| 2025 | Generative Caching for Structurally Similar Prompts and ResponsesabstractLarge Language Models (LLMs) are increasingly being used to plan, reason, and execute tasks across diverse scenarios. In use cases like repeatable workflows and agentic settings, prompts are often reused with minor variations while having a similar structure for recurring tasks. This opens up opportunities for caching. However, exact prompt matching fails on such structurally similar prompts, while semantic caching may produce incorrect responses by ignoring critical differences. To address this, we introduce GenCache, a generative cache that produces variation-aware responses for structurally similar prompts. GenCache identifies reusable response patterns across similar prompt structures and synthesizes customized outputs for new requests. We show that GenCache achieves 83\% cache hit rate, while having minimal incorrect hits on datasets without prompt repetition. In agentic workflows, it improves cache hit rate by $\sim$20\% and reduces end-to-end execution latency by $\sim$34\% compared to standard prompt matching. Sarthak Chakraborty, Suman Nath, Xuchao Zhang, Chetan Bansal, Indranil Gupta |
NeurIPS | 5 |
| 2025 | Transactional panorama: a conceptual framework for user perception in analytical visual interfaces (extended version)
Dixin Tang, Alan D. Fekete, Indranil Gupta, Aditya G. Parameswaran |
VLDB J. | 3 |
| 2024 | Camera: Churn-Tolerant Mutual Exclusion for the EdgeabstractEmerging edge and IoT networks are characterized by churn, wherein nodes join, leave, and fail continuously. Together with flaky network links and delays in edge and IoT systems, this means nodes have differing views of the system membership at any point of time. This paper focuses on the classical mutual exclusion problem—which ensures at most one node executes the “critical section” at any point in time—over such churned edge networks. We first show that two classical algorithms for mutual exclusion (Ricart-Agrawala and Maekawa) violate safety even under very small amounts of churn (up to 4 churned entries in membership lists). We then present Camera, a churn-tolerant variant of Ricart-Agrawala's algorithm. We formally prove Camera satisfies key properties including safety, deadlock-freedom, and starvation-freedom. Our trace-driven simulation results, using synthetic traces, distributions, and churn traces from peer to peer environments, show that Camera has scalable wait times and bandwidth overheads. Aman Khinvasara, Indranil Gupta |
SEC | 2 |
| 2024 | FrameCorr: Adaptive, Autoencoder-based Neural Compression for Video Reconstruction in Resource and Timing Constrained Network SettingsabstractVideo processing is becoming increasingly popular and cost-effective on IoT devices but faces challenges in transmitting data under varying timing constraints and network bandwidth. Existing compression methods struggle with incomplete data. We present FrameCorr, a deep learning framework that leverages prior-received video data to predict and reconstruct missing frame segments, enabling video reconstruction despite data loss. John Li, Deepak Nair, Klara Nahrstedt, Indranil Gupta, Shehab S. Ahmed |
ISM | 4 |
| 2024 | Known Knowns and Unknowns: Near-realtime Earth Observation Via Query Bifurcation in Serval
Bill Tao, Om Chabra, Ishani Janveja, Indranil Gupta, Deepak Vasisht |
NSDI | 4 |
| 2023 | Fail through the Cracks: Cross-System Interaction Failures in Modern Cloud SystemsabstractModern cloud systems are orchestrations of independent and interacting (sub-)systems, each specializing in important services (e.g., data processing, storage, resource management, etc.). Hence, cloud system reliability is affected not only by the reliability of each individual system, but also by the interplay between these systems. We observe that many recent production incidents of cloud systems are manifested through interactions across the system boundaries. However, there is a lack of systematic understanding of this emerging mode of failures, which we term as cross-system interaction failures (or CSI failures). This hinders the development of better design, integration practices, and new tooling. Lilia Tang, Chaitanya Bhandari, Yongle Zhang 0007, Anna Karanika, Shuyang Ji, Indranil Gupta, Tianyin Xu |
EuroSys | 6 |
| 2023 | Churn-Tolerant Leader Election ProtocolsabstractClassical leader election protocols typically assume complete and correct knowledge of underlying membership lists at all participating nodes. Yet many edge and IoT settings are dynamic, with nodes joining, leaving, and failing continuously―a phenomenon called churn. This implies that in any membership protocol, a given node's membership list may have entries that are missing (e.g., false positive detections, or newly joined nodes whose information has not spread yet) or stale (e.g., failed nodes that are undetected)―these would render classical election protocols incorrect. We present a family of four leader election protocols that are churn-tolerant (or c-tolerant). The key ideas are to: i) involve the minimum number of nodes necessary to achieve safety; ii) use optimism so that decisions are made faster when churn is low; iii) incorporate a preference for electing healthier nodes as leaders. We prove the correctness and safety of our c-tolerant protocols and show their message complexity is optimal. We present experimental results from both a trace-driven simulation as well as our implementation atop Raspberry Pi devices, including a comparison against Zookeeper. Jiangran Wang, Indranil Gupta |
ICDCS | 2 |
| 2023 | Transmitting, Fast and Slow: Scheduling Satellite Traffic through Space and TimeabstractEarth observation Low Earth Orbit (LEO) satellites collect enormous amounts of data that needs to be transferred first to ground stations and then to the cloud, for storage and processing. Satellites today transmit data greedily to ground stations, with full utilization of bandwidth during each contact period. We show that due to the layout of ground stations and orbital characteristics, this approach overloads some ground stations and underloads others, leading to lost throughput and large end-to-end latency for images. We present a new end-to-end scheduler system called Umbra, which plans transfers from large satellite constellations through ground stations to the cloud, by accounting for both spatial and temporal factors, i.e., orbital dynamics, bandwidth constraints, and queue sizes. At the heart of Umbra is a new class of scheduling algorithms called withhold scheduling, wherein the sender (i.e., satellite) selectively under-utilizes some links to ground stations. We show that Umbra's counter-intuitive approach increases throughput by 13--31% & reduces P90 latency by 3--6 ×. Bill Tao, Maleeha Masood, Indranil Gupta, Deepak Vasisht |
MobiCom | 3 |
| 2023 | Transactional Panorama: A Conceptual Framework for User Perception in Analytical Visual InterfacesabstractMany tools empower analysts and data scientists to consume analysis results in a visual interface. When the underlying data changes, these results need to be updated, but this update can take a long time---all while the user continues to explore the results. Tools can either (i) hide away results that haven't been updated, hindering exploration; (ii) make the updated results immediately available to the user (on the same screen as old results), leading to confusion and incorrect insights; or (iii) present old---and therefore stale---results to the user during the update. To help users reason about these options and others, and make appropriate trade-offs, we introduce Transactional Panorama, a formal framework that adopts transactions to jointly model the system refreshing the analysis results and the user interacting with them. We introduce three key properties that are important for user perception in this context: visibility (allowing users to continuously explore results), consistency (ensuring that results presented are from the same version of the data), and monotonicity (making sure that results don't "go back in time"). Within transactional panorama, we characterize all feasible property combinations, design new mechanisms (that we call lenses) for presenting analysis results to the user while preserving a given property combination, formally prove their relative orderings for various performance criteria, and discuss their use cases. We propose novel algorithms to preserve each property combination and efficiently present fresh analysis results. We implement our framework into a popular, open-source BI tool, illustrate the relative performance implications of different lenses, and demonstrate the benefits of the novel lenses and our optimizations. Dixin Tang, Alan D. Fekete, Indranil Gupta, Aditya G. Parameswaran |
Proc. VLDB Endow. | 3 |
| 2022 | Banyan: A Scoped Dataflow Engine for Graph Query ServiceabstractGraph query services (GQS) are widely used today to interactively answer graph traversal queries on large-scale graph data. Existing graph query engines focus largely on optimizing the latency of a single query. This ignores significant challenges posed by GQS, including fine-grained control and scheduling during query execution, as well as performance isolation and load balancing in various levels from across user to intra-query. To tackle these control and scheduling challenges, we propose a novel scoped dataflow for modeling graph traversal queries, which explicitly exposes concurrent execution and control of any subquery to the finest granularity. We implemented Banyan, an engine based on the scoped dataflow model for GQS. Banyan focuses on scaling up the performance on a single machine, and provides the ability to easily scale out. Extensive experiments on multiple benchmarks show that Banyan improves performance by up to three orders of magnitude over state-of-the-art graph query engines, while providing performance isolation and load balancing. Li Su 0005, Xiaoming Qin, Indranil Gupta, Wenyuan Yu, Kai Zeng 0002, Jingren Zhou 0001 |
Proc. VLDB Endow. | 6 |
| 2022 | Medley: A Membership Service for IoT NetworksabstractEfficient and correct operation of an IoT network requires the presence of a failure detector and membership protocol amongst the IoT nodes. This paper presents a new failure detector for IoT settings wherein nodes are connected via a wireless ad-hoc network. Our failure detector, named Medley, is fully decentralized, allows IoT nodes to maintain a local membership list of other alive nodes, detects failures quickly (and updates the membership list), and incurs low communication overhead. We adapt a failure detector originally proposed for datacenters (SWIM), for the IoT environment. This adaptation is non-trivial. In Medley each node picks a medley of ping targets in a randomized and skewed manner, preferring nearer nodes. We also provide optimizations to achieve time-bounded detection, as well as to reduce tail detection times. Via analysis, simulation, and Raspberry Pi deployments, we show that Medley can simultaneously optimize detection time and communication traffic. Rui Yang 0034, Jiangran Wang, Jiyu Hu, Shichu Zhu, Indranil Gupta |
IEEE Trans. Netw. Serv. Manag. | 6 |
| 2021 | Home, safehome: smart home reliability with visibility and atomicityabstractSmart environments (homes, factories, hospitals, buildings) contain an increasing number of IoT devices, making them complex to manage. Today, in smart homes when users or triggers initiate routines (i.e., a sequence of commands), concurrent routines and device failures can cause incongruent outcomes. We describe SafeHome, a system that provides notions of atomicity and serial equivalence for smart homes. Due to the human-facing nature of smart homes, SafeHome offers a spectrum of visibility models which trade off between responsiveness vs. isolation of the smart home. We implemented SafeHome and performed workload-driven experiments. We find that a weak visibility model, called eventual visibility, is almost as fast as today's status quo (up to 23% slower) and yet guarantees serially-equivalent end states. Shegufta Bakht Ahsan, Rui Yang 0034, Shadi A. Noghabi, Indranil Gupta |
EuroSys | 4 |
| 2021 | Move Fast and Meet Deadlines: Fine-grained Real-time Stream Processing with Cameo
Shivaram Venkataraman, Indranil Gupta, Luo Mai, Rahul Potharaju |
NSDI | 3 |
| 2020 | Baechi: fast device placement of machine learning graphsabstractMachine Learning graphs (or models) can be challenging or impossible to train when either devices have limited memory, or the models are large. Splitting the model graph across multiple devices, today, largely relies on learning-based approaches to generate this placement. While it results in models that train fast on data (i.e., with low step times), learning-based model-parallelism is time-consuming, taking many hours or days to create a placement plan of operators on devices. We present the Baechi system, where we adopt an algorithmic approach to the placement problem for running machine learning training graphs on a small cluster of memory-constrained devices. We implemented Baechi so that it works modularly with TensorFlow. Our experimental results using GPUs show that Baechi generates placement plans in time 654X--206K X faster than today's learning-based approaches, and the placed model's step time is only up to 6.2% higher than expert-based placements. Beomyeol Jeon, Linda Cai, Pallavi Srivastava, Jintao Jiang, Xiaolan Ke, Yitao Meng, Indranil Gupta |
SoCC | 8 |
| 2020 | Zeno++: Robust Fully Asynchronous SGDabstractWe propose Zeno++, a new robust asynchronous Stochastic Gradient Descent(SGD) procedure, intended to tolerate Byzantine failures of workers. In contrast to previous work, Zeno++ removes several unrealistic restrictions on worker-server communication, now allowing for fully asynchronous updates from anonymous workers, for arbitrarily stale worker updates, and for the possibility of an unbounded number of Byzantine workers. The key idea is to estimate the descent of the loss value after the candidate gradient is applied, where large descent values indicate that the update results in optimization progress. We prove the convergence of Zeno++ for non-convex problems under Byzantine failures. Experimental results show that Zeno++ outperforms existing Byzantine-tolerant asynchronous SGD algorithms. Oluwasanmi Koyejo, Indranil Gupta |
ICML | 3 |
| 2020 | A New Fully-Distributed Arbitration-Based Membership ProtocolabstractRecently, a new class of "arbitrator-based" membership protocols have been proposed. These claim to provide time bounds on how long membership lists can stay inconsistent-this property is critical in many distributed applications which need to take timely recovery actions. In this paper, we: 1) present the first fully decentralized and stabilizing version of membership protocols in this class; 2) formally prove properties and claims about both our decentralized version and the original protocol; and 3) present experimental results from both a simulation and a real cluster implementation. Shegufta Bakht Ahsan, Indranil Gupta |
INFOCOM | 2 |
| 2020 | CSER: Communication-efficient SGD with Error ResetabstractThe scalability of Distributed Stochastic Gradient Descent (SGD) is today limited by communication bottlenecks. We propose a novel SGD variant: \underline{C}ommunication-efficient \underline{S}GD with \underline{E}rror \underline{R}eset, or \underline{CSER}. The key idea in CSER is first a new technique called ``error reset'' that adapts arbitrary compressors for SGD, producing bifurcated local models with periodic reset of resulting local residual errors. Second we introduce partial synchronization for both the gradients and the models, leveraging advantages from them. We prove the convergence of CSER for smooth non-convex problems. Empirical results show that when combined with highly aggressive compressors, the CSER algorithms accelerate the distributed training by nearly $10\times$ for CIFAR-100, and by $4.5\times$ for ImageNet. Shuai Zheng 0004, Oluwasanmi Koyejo, Indranil Gupta, Mu Li 0003, Haibin Lin |
NeurIPS | 4 |
| 2019 | Kaizen: Building a Performant Blockchain System Verified for Consensus and IntegrityabstractWe report on the development of a blockchain system that is significantly verified and performant, detailing the design, proof, and system development based on a process of continuous refinement. We instantiate this framework to build, to the best of our knowledge, the first blockchain (Kaizen) that is performant and verified to a large degree, and a cryptocurrency protocol (KznCoin) over it. We experimentally compare its performance against the stock Bitcoin implementation. Faria Kalim, Karl Palmskog, Jayasi Mehar, Adithya Murali, Indranil Gupta, P. Madhusudan |
FMCAD | 5 |
| 2019 | Zeno: Distributed Stochastic Gradient Descent with Suspicion-based Fault-toleranceabstractWe present Zeno, a technique to make distributed machine learning, particularly Stochastic Gradient Descent (SGD), tolerant to an arbitrary number of faulty workers. Zeno generalizes previous results that assumed a majority of non-faulty nodes; we need assume only one non-faulty worker. Our key idea is to suspect workers that are potentially defective. Since this is likely to lead to false positives, we use a ranking-based preference mechanism. We prove the convergence of SGD for non-convex problems under these scenarios. Experimental results show that Zeno outperforms existing approaches. Oluwasanmi Koyejo, Indranil Gupta |
ICML | 3 |
| 2019 | Medley: A Novel Distributed Failure Detector for IoT NetworksabstractEfficient and correct operation of an IoT network requires the presence of a failure detector and membership protocol amongst the IoT nodes. This paper presents a new failure detector for IoT settings where nodes are connected via a wireless ad-hoc network. This failure detector, which we name Medley, is fully decentralized, allows IoT nodes to maintain a local membership list of other alive nodes, detects failures quickly (and updates the membership list), and incurs low communication overhead in the underlying ad-hoc network. In order to minimize detection time and communication, we adapt a failure detector originally proposed for datacenters (SWIM), for the IoT environment. In Medley each node picks a medley of ping targets in a randomized and skewed manner, preferring nearer nodes. Via analysis and NS-3 simulation we show the right mix of pinging probabilities that simultaneously optimize detection time and communication traffic. We have also implemented Medley for Raspberry Pis, and present deployment results. Rui Yang 0034, Shichu Zhu, Indranil Gupta |
Middleware | 4 |
| 2019 | SLSGD: Secure and Efficient Distributed On-device Machine Learning
Oluwasanmi Koyejo, Indranil Gupta |
ECML/PKDD (2) | 3 |
| 2019 | Fall of Empires: Breaking Byzantine-tolerant SGD by Inner Product Manipulation
Oluwasanmi Koyejo, Indranil Gupta |
UAI | 3 |
| 2019 | Read atomic transactions with prevention of lost updates: ROLA and its formal analysisabstractAbstract Designers of distributed database systems face the choice between stronger consistency guarantees and better performance. A number of applications only require read atomicity (RA) (either all or none of a transaction’s updates are visible to other transactions) and prevention of lost updates (PLU). Existing distributed transaction systems that meet these requirements also provide additional stronger consistency guarantees (such as causal consistency ), but this comes at the price of lower performance. In this paper we propose a new distributed transaction protocol, ROLA, that targets application scenarios where only RA and PLU are needed. We formally specify ROLA in Maude. We then perform model checking to analyze both the correctness and the performance of ROLA. For correctness, we use standard model checking to analyze ROLA’s satisfaction of RA and PLU. To analyze performance we: (a) perform statistical model checking to analyze key performance properties; and (b) compare these performance results with those obtained by also modeling and analyzing in Maude the well-known protocols Walter and Jessy that also guarantee RA and PLU. Our statistical model checking results show that ROLA outperforms both Walter and Jessy. Si Liu 0003, Peter Csaba Ölveczky, Qi Wang 0017, Indranil Gupta, José Meseguer 0001 |
Formal Aspects Comput. | 4 |
| 2018 | Henge: Intent-driven Multi-Tenant Stream ProcessingabstractWe present Henge, a system to support intent-based multi-tenancy in modern distributed stream processing systems. Henge supports multi-tenancy as a first-class citizen: everyone in an organization can now submit their stream processing jobs to a single, shared, consolidated cluster. Secondly, Henge allows each job to specify its own intents (i.e., requirements) as a Service Level Objective (SLO) that captures latency and/or throughput needs. In such an intent-driven multi-tenant cluster, the Henge scheduler adapts continually to meet jobs' respective SLOs in spite of limited cluster resources, and under dynamically varying workloads. SLOs are soft and are based on utility functions. Henge's overall goal is to maximize the total system utility achieved by all jobs in the system. Henge is integrated into Apache Storm and we present experimental results using both production jobs from Yahoo! and real datasets from Twitter. Faria Kalim, Sharanya Bathey, Richa Meherwal, Indranil Gupta |
SoCC | 5 |
| 2018 | Service fabric: a distributed platform for building microservices in the cloudabstractWe describe Service Fabric (SF), Microsoft's distributed platform for building, running, and maintaining microservice applications in the cloud. SF has been running in production for 10+ years, powering many critical services at Microsoft. This paper outlines key design philosophies in SF. We then adopt a bottom-up approach to describe low-level components in its architecture, focusing on modular use and support for strong semantics like fault-tolerance and consistency within each component of SF. We discuss lessons learned, and present experimental results from production data. Gopal Kakivaya, Lu Xun, Richard Hasha, Shegufta Bakht Ahsan, Todd Pfleiger, Rishi Sinha, Mihail Tarta, Mark Fussell, Vipul Modi, Mansoor Mohsin, Ray Kong, Anmol Ahuja, Oana Platon, Alex Wun, Matthew Snider, Chacko Daniel, Dan Mastrian, Aprameya Rao, Vaishnav Kidambi, Randy Wang, Abhishek Ram, Sumukh Shivaprakash, Rajeet Nair, Alan Warwick, Bharat S. Narasimman, Jeffrey Chen, Abhay Balkrishna Mhatre, Preetha Subbarayalu, Mert Coskun, Indranil Gupta |
EuroSys | 33 |
| 2018 | ROLA: A New Distributed Transaction Protocol and Its Formal AnalysisabstractDesigners of distributed database systems face the choice between stronger consistency guarantees and better performance. A number of applications only require read atomicity (RA) and prevention of lost updates (PLU). Existing distributed database systems that meet these requirements also provide additional stronger consistency guarantees (such as causal consistency ), and therefore incur lower performance. In this paper we define a new distributed transaction protocol, ROLA, that targets applications where only RA and PLU are needed. We formally model ROLA in Maude. We then perform model checking to analyze both the correctness and the performance of ROLA. For correctness , we use standard model checking to analyze ROLA’s satisfaction of RA and PLU. To analyze performance we: (a) use statistical model checking to analyze key performance properties; and (b) compare these performance results with those obtained by analyzing in Maude the well-known protocol Walter. Our results show that ROLA outperforms Walter. Si Liu 0003, Peter Csaba Ölveczky, Keshav Santhanam, Qi Wang 0017, Indranil Gupta, José Meseguer 0001 |
FASE | 5 |
| 2018 | OPTiC: Opportunistic Graph Processing in Multi-Tenant ClustersabstractWe present OPTiC, a multi-tenant scheduler intended for distributed graph processing frameworks. OPTiC proposes opportunistic scheduling, whereby queued jobs can be pre-scheduled at cluster nodes when the cluster is fully busy running jobs. This allows overlapping of data ingress with ongoing computation. To pre-schedule wisely, OPTiC's novel contribution is a profile-free and cluster-agnostic approach to compare progress of graph processing jobs. OPTiC is implemented inside Apache Giraph, with YARN underneath. Our experiments with real workload traces and network models show that OPTiC's opportunistic scheduling improves run time (both at the median and at the tail) by 20%-82% compared to baseline multi-tenancy, in a variety of scenarios. Muntasir Raihan Rahman, Indranil Gupta, Akash Kapoor, Haozhen Ding |
IC2E | 2 |
| 2017 | 4CeeD: Real-Time Data Acquisition and Analysis Framework for Material-related Cyber-Physical EnvironmentsabstractIn this paper, we present a data acquisition and analysis framework for materials-to-devices processes, named 4CeeD, that focuses on the immense potential of capturing, accurately curating, correlating, and coordinating materials-to-devices digital data in a real-time and trusted manner before fully archiving and publishing them for wide access and sharing. In particular, 4CeeD consists of novel services: a curation service for collecting data from microscopes and fabrication instruments, curating, and wrapping of data with extensive metadata in real-time and in a trusted manner, and a cloud-based coordination service for storing data, extracting meta-data, analyzing and finding correlations among the data. Our evaluation results show that our novel cloud framework can help researchers significantly save time and cost spent on experiments, and is efficient in dealing with high-volume and fast-changing workload of heterogeneous types of experimental data. Phuong Nguyen 0002, Steven Konstanty, Todd Nicholson, Thomas O'Brien, Aaron Schwartz-Duval, Timothy Spila, Klara Nahrstedt, Roy H. Campbell, Indranil Gupta, Kenton McHenry, Normand Paquin |
CCGrid | 9 |
| 2017 | REMAX: Reachability-Maximizing P2P Detection of Erroneous Readings in Wireless Sensor NetworksabstractWireless sensor networks (WSNs) should collect accurate readings to reliably capture an environment's state. However, readings may become erroneous because of sensor hardware failures or degradation. In remote deployments, centrally detecting those reading errors can result in many message transmissions, which in turn dramatically decreases sensor battery life. In this paper, we address this issue through three main contributions. First, we propose REMAX, a peer-to-peer (P2P) error detection protocol that extends the WSN's life by minimizing message transmissions. Second, we propose a low-overhead error detection approach that helps minimize communication complexity. Third, we evaluate our approach via a trace-driven, discrete-event simulator, using two datasets from real WSN deployments that measure indoor air temperature and seismic wave amplitude. Our results show that REMAX can accurately detect errors and extend the WSN's reachability (effective lifetime) compared to the centralized approach. Varun Badrinath Krishna, Michael J. Rausch, Benjamin E. Ujcich, Indranil Gupta, William H. Sanders |
DSN | 4 |
| 2017 | Exploring Design Alternatives for RAMP Transactions Through Statistical Model Checking
Si Liu 0003, Peter Csaba Ölveczky, Jatin Ganhotra, Indranil Gupta, José Meseguer 0001 |
ICFEM | 4 |
| 2017 | Stateful Scalable Stream Processing at LinkedInabstractDistributed stream processing systems need to support stateful processing, recover quickly from failures to resume such processing, and reprocess an entire data stream quickly. We present Apache Samza, a distributed system for stateful and fault-tolerant stream processing. Samza utilizes a partitioned local state along with a low-overhead background changelog mechanism, allowing it to scale to massive state sizes (hundreds of TB) per application. Recovery from failures is sped up by re-scheduling based on Host Affinity. In addition to processing infinite streams of events, Samza supports processing a finite dataset as a stream, from either a streaming source (e.g., Kafka), a database snapshot (e.g., Databus), or a file system (e.g. HDFS), without having to change the application code (unlike the popular Lambda-based architectures which necessitate maintenance of separate code bases for batch and stream path processing). Samza is currently in use at LinkedIn by hundreds of production applications with more than 10, 000 containers. Samza is an open-source Apache project adopted by many top-tier companies (e.g., LinkedIn, Uber, Netflix, TripAdvisor, etc.). Our experiments show that Samza: a) handles state efficiently, improving latency and throughput by more than 100X compared to using a remote storage; b) provides recovery time independent of state size; c) scales performance linearly with number of containers; and d) supports reprocessing of the data stream quickly and with minimal interference on real-time traffic. Shadi A. Noghabi, Kartik Paramasivam, Navina Ramesh, Jon Bringhurst, Indranil Gupta, Roy H. Campbell |
Proc. VLDB Endow. | 6 |
| 2017 | An Experimental Comparison of Partitioning Strategies in Distributed Graph ProcessingabstractIn this paper, we study the problem of choosing among partitioning strategies in distributed graph processing systems. To this end, we evaluate and characterize both the performance and resource usage of different partitioning strategies under various popular distributed graph processing systems, applications, input graphs, and execution environments. Through our experiments, we found that no single partitioning strategy is the best fit for all situations, and that the choice of partitioning strategy has a significant effect on resource usage and application run-time. Our experiments demonstrate that the choice of partitioning strategy depends on (1) the degree distribution of input graph, (2) the type and duration of the application, and (3) the cluster size. Based on our results, we present rules of thumb to help users pick the best partitioning strategy for their particular use cases. We present results from each system, as well as from all partitioning strategies implemented in one common system (PowerLyra). Shiv Verma, Luke M. Leslie, Yosub Shin, Indranil Gupta |
Proc. VLDB Endow. | 4 |
| 2017 | Characterizing and Adapting the Consistency-Latency Tradeoff in Distributed Key-Value StoresabstractThe CAP theorem is a fundamental result that applies to distributed storage systems. In this article, we first present and prove two CAP-like impossibility theorems. To state these theorems, we present probabilistic models to characterize the three important elements of the CAP theorem: consistency (C), availability or latency (A), and partition tolerance (P). The theorems show the un-achievable envelope, that is, which combinations of the parameters of the three models make them impossible to achieve together. Next, we present the design of a class of systems called Probabilistic CAP (PCAP) that perform close to the envelope described by our theorems. In addition, these systems allow applications running on a single data center to specify either a latency Service Level Agreement (SLA) or a consistency SLA. The PCAP systems automatically adapt, in real time and under changing network conditions, to meet the SLA while optimizing the other C/A metric. We incorporate PCAP into two popular key-value stores: Apache Cassandra and Riak. Our experiments with these two deployments, under realistic workloads, reveal that the PCAP systems satisfactorily meets SLAs and perform close to the achievable envelope. We also extend PCAP from a single data center to multiple geo-distributed data centers. Muntasir Raihan Rahman, Lewis Tseng, Indranil Gupta, Nitin H. Vaidya |
ACM Trans. Auton. Adapt. Syst. | 4 |
| 2016 | Phurti: Application and Network-Aware Flow Scheduling for Multi-tenant MapReduce ClustersabstractTraffic for a typical MapReduce job in a data center consists of multiple network flows. Traditionally, network resources have been allocated to optimize network-level metrics such as flow completion time or throughput. Some recent schemes propose using application-aware scheduling which can shorten the average job completion time. However, most of them treat the core network as a black box with sufficient capacity. Even if only one network link in the core network becomes a bottleneck, it can hurt application performance. We design and implement a centralized flow-scheduling framework called Phurti with the goal of improving the completion time for jobs in a cluster shared among multiple Hadoop jobs (multi-tenant). Phurti communicates both with the Hadoop framework to retrieve job-level network traffic information and the OpenFlow-based switches to learn about the network topology. Phurti implements a novel heuristic called Smallest Maximum Sequential-traffic First (SMSF) that uses collected application and network information to perform traffic scheduling for MapReduce jobs. Our evaluation with real Hadoop workloads shows that compared to application and network-agnostic scheduling strategies, Phurti improves job completion time for 95% of the jobs, decreases average job completion time by 20%, tail job completion time by 13% and scales well with the cluster size and number of jobs. Chris X. Cai, Shayan Saeed, Indranil Gupta, Roy H. Campbell, Franck Le |
IC2E | 3 |
| 2016 | Supporting On-demand Elasticity in Distributed Graph ProcessingabstractWhile distributed graph processing engines have become popular for processing large graphs, these engines are typically configured with a static set of servers in the cluster. In other words, they lack the flexibility to scale-out or scale-in the number of servers, when requested to do so by the user. In this paper, we propose the first techniques to make distributed graph processing truly elastic. While supporting on-demand scale-out/in operations, we meet three goals: i) perform scale-out/in without interrupting the graph computation, ii) minimize the background network overhead involved in the scale-out/in, and iii) mitigate stragglers by maintaining load balance across servers. We present and analyze two techniques called Contiguous Vertex Repartitioning (CVR) and Ring-based Vertex Repartitioning (RVR) to address these goals. We implement our techniques in the LFGraph distributed graph processing system, and incorporate several systems optimizations. Experiments performed with multiple graph benchmark applications on a real graph indicate that our techniques perform within 9% and 21% of the optimum for scale-out and scale-in operations, respectively. Mayank Pundir, Luke M. Leslie, Indranil Gupta, Roy H. Campbell |
IC2E | 4 |
| 2016 | Stela: Enabling Stream Processing Systems to Scale-in and Scale-out On-demandabstractThe era of big data has led to the emergence of new real-time distributed stream processing engines like Apache Storm. We present Stela (STream processing ELAsticity), a stream processing system that supports scale-out and scale-in operations in an on-demand manner, i.e., when the user requests such a scaling operation. Stela meets two goals: 1) it optimizes post-scaling throughput, and 2) it minimizes interruption to the ongoing computation while the scaling operation is being carried out. We have integrated Stela into Apache Storm. We present experimental results using micro-benchmark Storm applications, as well as production applications from industry (Yahoo! Inc. and IBM). Our experiments show that compared to Apache Storm's default scheduler, Stela's scale-out operation achieves throughput that is 21-120% higher, and interruption time that is significantly smaller. Stela's scale-in operation chooses the right set of servers to remove and achieves 2X-5X higher throughput than Storm's default strategy. Boyang Peng, Indranil Gupta |
IC2E | 3 |
| 2016 | Ambry: LinkedIn's Scalable Geo-Distributed Object StoreabstractThe infrastructure beneath a worldwide social network has to continually serve billions of variable-sized media objects such as photos, videos, and audio clips. These objects must be stored and served with low latency and high throughput by a system that is geo-distributed, highly scalable, and load-balanced. Existing file systems and object stores face several challenges when serving such large objects. We present Ambry, a production-quality system for storing large immutable data (called blobs). Ambry is designed in a decentralized way and leverages techniques such as logical blob grouping, asynchronous replication, rebalancing mechanisms, zero-cost failure detection, and OS caching. Ambry has been running in LinkedIn's production environment for the past 2 years, serving up to 10K requests per second across more than 400 million users. Our experimental evaluation reveals that Ambry offers high efficiency (utilizing up to 88% of the network bandwidth), low latency (less than 50 ms latency for a 1 MB object), and load balancing (improving imbalance of request rate among disks by 8x-10x). Shadi A. Noghabi, Sriram Subramanian, Priyesh Narayanan, Sivabalan Narayanan, Gopalakrishna Holla, Mammad Zadeh, Tianwei Li, Indranil Gupta, Roy H. Campbell |
SIGMOD Conference | 8 |
| 2015 | Leveraging Metadata in No SQL Storage SystemsabstractNoSQL systems have grown in popularity for storing big data because these systems offer high availability, i.e., Operations with high throughput and low latency. However, metadata in these systems are handled today in ad-hoc ways. We present Wasef, a system that treats metadata in a NoSQL database system, as first-class citizens. Metadata may include information such as: operational history for a database table (e.g., Columns), placement information for ranges of keys, and operational logs for data items (key-value pairs). Wasef allows the NoSQL system to store and query this metadata efficiently. We integrate Wasef into Apache Cassandra, one of the most popular key-value stores. We then implement three important use cases in Cassandra: dropping columns in a flexible manner, verifying data durability during migrational operations such as node decommissioning, and maintaining data provenance. Our experimental evaluation uses AWS EC2 instances and YCSB workloads. Our results show that Wasef: i) scales well with the size of the data and the metadata, ii) minimally affects throughput and operation latencies. Ala' Alkhaldi, Indranil Gupta, Vaijayanth Raghavan, Mainak Ghosh |
CLOUD | 2 |
| 2015 | Zorro: zero-cost reactive failure recovery in distributed graph processingabstractDistributed graph processing systems largely rely on proactive techniques for failure recovery. Unfortunately, these approaches (such as checkpointing) entail a significant overhead. In this paper, we argue that distributed graph processing systems should instead use a reactive approach to failure recovery. The reactive approach trades off completeness of the result (generating a slightly inaccurate result) while reducing the overhead during failure-free execution to zero. We build a system called Zorro that imbues this reactive approach, and integrate Zorro into two graph processing systems -- PowerGraph and LFGraph. When a failure occurs, Zorro opportunistically exploits vertex replication inherent in today's graph processing systems to quickly rebuild the state of failed servers. Experiments using real-world graphs demonstrate that Zorro is able to recover over 99% of the graph state when 6--12% of the servers fail, and between 87--95% when half the cluster fails. Furthermore, using various graph processing algorithms, Zorro incurs little to no accuracy loss in all experimental failure scenarios, and achieves a worst-case accuracy of 97%. Mayank Pundir, Luke M. Leslie, Indranil Gupta, Roy H. Campbell |
SoCC | 3 |
| 2015 | Cross-Layer Scheduling in Cloud SystemsabstractToday, cloud computing engines such as stream-processing Storm and batch-processing Hadoop are being increasingly run atop software-defined networks (SDNs). In such cloud stacks, the scheduler of the application engine (which allocates tasks to servers) remains decoupled from the SDN scheduler (which allocates network routes). We propose a new approach that performs cross-layer scheduling between the application layer and the networking layer. This coordinated scheduling orchestrates the placement of application tasks (e.g., Hadoop maps and reduces, or Storm bolts) in tandem with the selection of network routes that arise from these tasks. We present results from both cluster deployment and simulation, and using two representative network topologies: Fat-tree and Jellyfish. Our results show that cross-layer scheduling can improve throughput of Hadoop and Storm by between 26% to 34% in a 30-host cluster, and it scales well. Hilfi Alkaff, Indranil Gupta, Luke M. Leslie |
IC2E | 2 |
| 2015 | Scale Up vs. Scale Out in Cloud Storage and Graph Processing SystemsabstractDeployers of cloud storage and iterative processing systems typically have to deal with either dollar budget constraints or throughput requirements. This paper examines the question of whether such cloud storage and iterative processing systems are more cost-efficient when scheduled on a COTS (scale out) cluster or a single beefy (scale up) machine. We experimentally evaluate two systems: 1) a distributed key-value store (Cassandra), and 2) a distributed graph processing system (Graph Lab). Our studies reveal scenarios where each option is preferable over the other. We provide recommendations for deployers of such systems to decide between scale up vs. Scale out, as a function of their dollar or throughput constraints. Our results indicate that there is a need or adaptive scheduling in heterogeneous clusters containing scale up and scale out nodes. Indranil Gupta |
IC2E | 3 |
| 2015 | Fast Compaction Algorithms for NoSQL DatabasesabstractCompaction plays a crucial role in NoSQL systems to ensure a high overall read throughput. In this work, we formally define compaction as an optimization problem that attempts to minimize disk I/O. We prove this problem to be NP-Hard. We then propose a set of algorithms and mathematically analyze upper bounds on worst-case cost. We evaluate the proposed algorithms on real-life workloads. Our results show that our algorithms incur low I/O costs and that a compaction approach using a balanced tree is most preferable. Mainak Ghosh, Indranil Gupta, Shalmoli Gupta, Nirman Kumar |
ICDCS | 2 |
| 2015 | 3DTI Amphitheater: Towards 3DTI Broadcastingabstract3DTI Amphitheater is a live broadcasting system for dissemination of 3DTI (3D Tele-immersive) content. The virtual environment constructed by the system mimics an amphitheater in the real world, where performers interact with each other in the central circular stage, and the audience is placed in virtual seats that surround the stage. Users of the Amphitheater can be geographically dispersed and the streams created by the performer sites are disseminated in a P2P network among the participants. To deal with the high bandwidth demand and strict latency bound of the service, we identify the hierarchical priority of streams in construction of the content dissemination forest. Result shows that the Amphitheater outperforms prior 3DTI systems by boosting the application QoS by a factor of 2.8 while sustaining the same hundred-scale audience group. Chien-Nan (Shannon) Chen, Zhenhuan Gao, Klara Nahrstedt, Indranil Gupta |
ACM Trans. Multim. Comput. Commun. Appl. | 4 |
| 2014 | VMDedup: Memory De-duplication in HypervisorabstractVirtualization techniques are widely used in cloud computing environments today. Such environments are installed with a large number of similar virtual instances sharing the same physical infrastructure. In this paper, we focus on the memory usage optimization across virtual machines by automatically de-duplicating the memory on per-page basis. Our approach maintains a single copy of the duplicated pages in physical memory using copy-on-write mechanism. Unlike some existing strategies, which are intended only for applications and need user configuration, VMDedup provides an automatic memory de-duplication support within the hypervisor to achieve benefits across operating system code, data as well as application binaries. We have implemented a prototype of this system within the Xen hypervisor to support both para-virtualized and fully-virtualized instances of operating systems. Furquan Shaikh, Fangzhou Yao, Indranil Gupta, Roy H. Campbell |
IC2E | 3 |
| 2014 | Client-Centric Benchmarking of Eventual Consistency for Cloud Storage SystemsabstractEventually-consistent key-value storage systems sacrifice the ACID semantics of conventional databases to achieve superior latency and availability. However, this means that client applications, and hence end-users, can be exposed to stale data. The degree of staleness observed depends on various tuning knobs set by application developers (customers of key-value stores) and system administrators (providers of key-value stores). Both parties must be cognizant of how these tuning knobs affect the consistency observed by client applications in the interest of both providing the best end-user experience and maximizing revenues for storage providers. Quantifying consistency in a meaningful way is a critical step toward both understanding what clients actually observe, and supporting consistency-aware service level agreements (SLAs) in next generation storage systems. This paper proposes a novel consistency metric called Gamma that captures client-observed consistency. This metric provides quantitative answers to questions regarding observed consistency anomalies, such as how often they occur and how bad they are when they do occur. We argue that Gamma is more useful and accurate than existing metrics. We also apply Gamma to benchmark the popular Cassandra key-value store. Our experiments demonstrate that Gamma is sensitive to both the workload and client-level tuning knobs, and is preferable to existing techniques which focus on worst-case behavior. Wojciech M. Golab, Muntasir Raihan Rahman, Alvin AuYoung, Kimberly Keeton, Indranil Gupta |
ICDCS | 5 |
| 2014 | WOHA: Deadline-Aware Map-Reduce Workflow Scheduling Framework over Hadoop ClustersabstractIn this paper, we present WOHA, an efficient scheduling framework for deadline-aware Map-Reduce workflows. In data centers, complex backend data analysis often utilizes a workflow that contains tens or even hundreds of interdependent Map-Reduce jobs. Meeting deadlines of these workflows is usually of crucial importance to businesses (for example, workflows tightly linked to time-sensitive advertisement placement optimizations can directly affect revenue). Popular Map-Reduce implementations, such as Hadoop, deal with independent Map-Reduce jobs rather than workflows of jobs. In order to simplify the process of submitting workflows, solutions like Oozie emerge, which take a workflow configuration file as input and automatically submit its Hadoop jobs at the right time. The information separation that Hadoop only handles resource allocation and Oozie workflow topology, although preventing the Hadoop master node from getting involved with complex workflow analysis, may unnecessarily lengthen the workflow spans and thus cause more deadline misses. To address this problem and at the same time honor the efficiency of Hadoop master node, WOHA allows client nodes to locally generate scheduling plans which are later used as resource allocation hints by the master node. Under this framework design, we propose a novel scheduling algorithm that improves deadline satisfaction ratio by dynamically assigning priorities among workflows based on their progresses. We implement WOHA by extending Hadoop-1.2.1. Our experiments over an 80-server cluster show that WOHA manages to increase the deadline satisfaction ratio by 10% compared to state-of-the-art solutions, and scales up to tens of thousands of concurrently running workflows. Shen Li 0002, Shaohan Hu, Shiguang Wang, Lu Su 0001, Tarek F. Abdelzaher, Indranil Gupta, Richard Pace |
ICDCS | 6 |
| 2014 | Formal Modeling and Analysis of Cassandra in Maude
Si Liu 0003, Muntasir Raihan Rahman, Stephen Skeirik, Indranil Gupta, José Meseguer 0001 |
ICFEM | 4 |
| 2014 | 3DTI amphitheater: a manageable 3DTI environment with hierarchical stream prioritizationabstractIn this paper we present the 3DTI Amphitheater, a live broadcasting system for dissemination of 3DTI (3D Tele-immersive) content. The virtual environment constructed by the system mimics an amphitheater in the real world, where performers interact with each other in the central circular stage, and the audience is placed in virtual seats that surround the stage. Users of the Amphitheater can be geographically dispersed and the streams created by the performer sites are disseminated in a P2P network among the participants. To deal with the high bandwidth demand and strict latency bound of the service, we identify the hierarchical priority of streams in construction of the content dissemination forest. Result shows that the Amphitheater outperforms prior 3DTI systems by boosting the application QoS by a factor of 2.8 while sustaining the same hundred-scale audience group. Chien-Nan (Shannon) Chen, Klara Nahrstedt, Indranil Gupta |
MMSys | 3 |
| 2013 | Natjam: design and evaluation of eviction policies for supporting priorities and deadlines in mapreduce clustersabstractThis paper presents Natjam, a system that supports arbitrary job priorities, hard real-time scheduling, and efficient preemption for Mapreduce clusters that are resource-constrained. Our contributions include: i) exploration and evaluation of smart eviction policies for jobs and for tasks, based on resource usage, task runtime, and job deadlines; and ii) a work-conserving task preemption mechanism for Mapreduce. We incorporated Natjam into the Hadoop YARN scheduler framework (in Hadoop 0.23). We present experiments from deployments on a test cluster, Emulab and a Yahoo! Inc. commercial cluster, using both synthetic workloads as well as Hadoop cluster traces from Yahoo!. Our results reveal that Natjam incurs overheads as low as 7%, and is preferable to existing approaches. Muntasir Raihan Rahman, Tej Chajed, Indranil Gupta, Cristina L. Abad, Nathan Roberts, Philbert Lin |
SoCC | 4 |
| 2013 | Client-centric benchmarking of eventual consistency for cloud storage systemsabstractEventually consistent storage systems give up the ACID semantics of conventional databases in order to gain better scalability, higher availability, and lower latency. A side-effect of this design decision is that application developers must deal with stale or out of order data. As a result, substantial intellectual effort has been devoted to studying the behavior of eventually consistent systems, in particular finding quantitative answers to the questions "how eventual" and "how consistent"? Wojciech M. Golab, Muntasir Raihan Rahman, Alvin AuYoung, Kimberly Keeton, Jay J. Wylie, Indranil Gupta |
SoCC | 6 |
| 2010 | Making cloud intermediate data fault-tolerantabstractParallel dataflow programs generate enormous amounts of distributed data that are short-lived, yet are critical for completion of the job and for good run-time performance. We call this class of data as intermediate data. This paper is the first to address intermediate data as a first-class citizen, specifically targeting and minimizing the effect of run-time server failures on the availability of intermediate data, and thus on performance metrics such as job completion time. We propose new design techniques for a new storage system called ISS (Intermediate Storage System), implement these techniques within Hadoop, and experimentally evaluate the resulting system. Under no failure, the performance of Hadoop augmented with ISS (i.e., job completion time) turns out to be comparable to base Hadoop. Under a failure, Hadoop with ISS outperforms base Hadoop and incurs up to 18% overhead compared to base no-failure Hadoop, depending on the testbed setup. Steven Y. Ko, Imranul Hoque, Indranil Gupta |
SoCC | 4 |
| 2010 | Breaking the MapReduce Stage BarrierabstractThe MapReduce model uses a barrier between the Map and Reduce stages. This provides simplicity in both programming and implementation. However, in many situations, this barrier hurts performance because it is overly restrictive. Hence, we develop a method to break the barrier in MapReduce in a way that improves efficiency. Careful design of our barrierless MapReduce framework results in equivalent generality and retains ease of programming. We motivate our case with, and experimentally study our barrier-less techniques in, a wide variety of MapReduce applications divided into seven classes. Our experiments show that our approach can achieve better performance times than a traditional MapReduce framework. We achieve a reduction in job completion times that is 25% on average and 87% in the best case. Nicolas Zea, Indranil Gupta, Roy H. Campbell |
CLUSTER | 4 |
| 2010 | New Algorithms for Planning Bulk Transfer via Internet and Shipping NetworksabstractCloud computing is enabling groups of academic collaborators, groups of business partners, etc., to come together in an ad-hoc manner. This paper focuses on the group-based data transfer problem in such settings. Each participant source site in such a group has a large dataset, which may range in size from gigabytes to terabytes. This data needs to be transferred to a single sink site (e.g., AWS, Google datacenters, etc.) in a manner that reduces both total dollar costs incurred by the group as well as the total transfer latency of the collective dataset. This paper is the first to explore the problem of planning a group-based deadline-oriented data transfer in a scenario where data can be sent over both: (1) the internet, and (2) by shipping storage devices (e.g., external or hot-plug drives, or SSDs) via companies such as Fedex, UPS, USPS, etc. We first formalize the problem and prove its NP-Hardness. Then, we propose novel algorithms and use them to build a planning system called Pandora (People and Networks Moving Data Around). Pandora uses new concepts of time-expanded networks and delta-time-expanded networks, combining them with integer programming techniques and optimizations for both shipping and internet edges. Our experimental evaluation using real data from Fedex and from PlanetLab indicate the Pandora planner manages to satisfy deadlines and reduce costs significantly. Indranil Gupta |
ICDCS | 2 |
| 2010 | Joint bluetooth/wifi scanning framework for characterizing and leveraging people movement in university campusabstractThis paper1 presents a novel framework called UIM2, which collects both location information and ad hoc contacts of the human movement at the University of Illinois campus using Google Android phones. Each UIM experiment phone encompasses a Bluetooth scanner and a wifi scanner capturing both Bluetooth MAC addresses and wifi access point MAC addresses in proximity of the phone. Then, Bluetooth MAC addresses are used to infer contact information and the wifi MAC addresses are used to infer physical location of the phone. Using the contact and location information, we investigate first the sensitivity analysis on contact duration and inter-contact duration. Then, we characterize the regularity of people movement, visit duration of people at locations, and the popularity of locations. Finally, we present the Hybrid Epidemic data dissemination protocol, which uses both wifi access point and ad hoc contact to expedite the data forwarding. We evaluate Hybrid Epidemic protocol with our collected ad hoc and wifi traces and find that in comparison with Epidemic data dissemination protocol, the Hybrid Epidemic protocol improves data forwarding delay considerably. Long H. Vu, Klara Nahrstedt, Samuel Retika, Indranil Gupta |
MSWiM | 4 |
| 2010 | Understanding overlay characteristics of a large-scale peer-to-peer IPTV systemabstractThis article presents results from our measurement and modeling efforts on the large-scale peer-to-peer (p2p) overlay graphs spanned by the PPLive system, the most popular and largest p2p IPTV (Internet Protocol Television) system today. Unlike other previous studies on PPLive, which focused on either network-centric or user-centric measurements of the system, our study is unique in (a) focusing on PPLive overlay-specific characteristics, and (b) being the first to derive mathematical models for its distributions of node degree, session length, and peer participation in simultaneous overlays. Our studies reveal characteristics of multimedia streaming p2p overlays that are markedly different from existing file-sharing p2p overlays. Specifically, we find that: (1) PPLive overlays are similar to random graphs in structure and thus more robust and resilient to the massive failure of nodes, (2) Average degree of a peer in the overlay is independent of the channel population size and the node degree distribution can be fitted by a piecewise function, (3) The availability correlation between PPLive peer pairs is bimodal, that is, some pairs have highly correlated availability, while others have no correlation, (4) Unlike p2p file-sharing peers, PPLive peers are impatient and session lengths (discretized, per channel) are typically geometrically distributed, (5) Channel population size is time-sensitive, self-repeated, event-dependent, and varies more than in p2p file-sharing networks, (6) Peering relationships are slightly locality-aware, and (7) Peer participation in simultaneous overlays follows a Zipf distribution. We believe that our findings can be used to understand current large-scale p2p streaming systems for future planning of resource usage, and to provide useful and practical hints for future design of large-scale p2p streaming systems. Long H. Vu, Indranil Gupta, Klara Nahrstedt |
ACM Trans. Multim. Comput. Commun. Appl. | 2 |
| 2010 | Routing in the frequency domain
Jay A. Patel, Haiyun Luo, Indranil Gupta |
Wirel. Networks | 3 |
| 2009 | On Availability of Intermediate Data in Cloud Computations
Steven Y. Ko, Imranul Hoque, Indranil Gupta |
HotOS | 4 |
| 2009 | Q-Tree: A Multi-Attribute Based Range Query Solution for Tele-immersive FrameworkabstractUsers and administrators of large distributed systems are frequently in need of monitoring and management of its various components, data items and resources. Though there exist several distributed query and aggregation systems, the clustered structure of tele-immersive interactive frameworks and their time-sensitive nature and application requirements represent a new class of systems which poses different challenges on this distributed search. Multi-attribute composite range queries are one of the key features in this class. Queries are given in high level descriptions and then transformed into multi-attribute composite range queries. Designing such a query engine with minimum traffic overhead, low service latency, and with static and dynamic nature of large datasets, is a challenging task. In this paper, we propose a general multi-attribute based range query framework, Q-Tree, that provides efficient support for this class of systems. In order to serve efficient queries, Q-Tree builds a single topology-aware tree overlay by connecting the participating nodes in a bottom-up approach, and assigns range intervals on each node in a hierarchical manner. We show the relative strength of Q-Tree by analytically comparing it against P-Tree, P-Ring, Skip-Graph and Chord. With fine-grained load balancing and overlay maintenance, our simulations with PlanetLab traces show that our approach can answer complex queries within a fraction of a second. Ahsan Arefin, Md. Yusuf Sarwar Uddin, Indranil Gupta, Klara Nahrstedt |
ICDCS | 3 |
| 2009 | Using Failure Models for Controlling Data Availability in Wireless Sensor NetworksabstractThis paper presents Pirrus, a replica management system that addresses the problem of providing data availability on a wireless sensor network. Pirrus uses probabilistic failure models (e.g., derived from environmental conditions and estimation of available energy) to adaptively create and maintain a number of replicas of the data. Replica management is formulated as an energy optimization problem, then solved with a greedy heuristic that only uses information gathered from neighbors. Intuitively, Pirrus trades off the energy saved by limiting the number of replicas when the network health is good to extend the lifetime when more replicas are needed. Our simulation results show how the solution provided by Pirrus achieves good performance with a sustainable computational cost. Compared to the performance of a fixed number of replicas, Pirrus extends the network lifetime by more than 20%. Riccardo Crepaldi, Mirko Montanari, Indranil Gupta, Robin Kravets |
INFOCOM | 3 |
| 2009 | MLR-Index: An Index Structure for Fast and Scalable Similarity Search in High Dimensions
Rahul Malik, Sangkyum Kim, Xin Jin 0001, Chandrasekar Ramachandran, Jiawei Han 0001, Indranil Gupta, Klara Nahrstedt |
SSDBM | 6 |
| 2009 | AVCOL: Availability-aware information aggregation in large distributed systems under uncollaborative behavior
Ramsés Morales, Indranil Gupta |
Comput. Networks | 2 |
| 2009 | Rappel: Exploiting interest and network locality to improve fairness in publish-subscribe systems
Jay A. Patel, Etienne Rivière, Indranil Gupta, Anne-Marie Kermarrec |
Comput. Networks | 3 |
| 2009 | AVMON: Optimal and Scalable Discovery of Consistent Availability Monitoring Overlays for Distributed SystemsabstractThis paper proposes to build overlays that help in monitoring of long-term availability histories of hosts, with a focus on large-scale distributed settings where hosts may be selfish or colluding. Concretely, we target the important problems of selection and discovery of an availability monitoring overlay. We motivate six significant goals - firstly, consistency, verifiability, and randomness, in selecting availability monitors of nodes, so as to be probabilistically resilient to selfish and colluding nodes. The next three goals are discoverability, load-balancing, and scalability in finding these monitors. We present AVMON, the first availability monitoring overlay to satisfy these six requirements. Our core algorithmic contribution is a range of protocols for discovering the availability monitoring overlay scalably and efficiently, given any arbitrary monitor selection scheme that is consistent and verifiable. We mathematically analyze the performance of AVMON's discovery protocols w.r.t. scalability and discovery time of monitors. Most interestingly, we are able to derive optimal (and practical) variants of AVMON, that minimize different combinations of memory, bandwidth, computation, and monitor discovery time. Finally, our extensive experimental evaluations using three types of availability traces - synthetic, from PlanetLab, and from the Overnet p2p system - demonstrate AVMON's practicality in a variety of distributed systems. Ramsés Morales, Indranil Gupta |
IEEE Trans. Parallel Distributed Syst. | 2 |
| 2008 | Fair K Mutual Exclusion Algorithm for Peer to Peer Systemsabstractk-mutual exclusion is an important problem for resource-intensive peer-to-peer applications ranging from aggregation to file downloads. In order to be practically useful, k-mutual exclusion algorithms not only need to be safe and live, but they also need to be fair across hosts. We propose a new solution to the k-mutual exclusion problem that provides a notion of time-based fairness. Specifically, our algorithm attempts to minimize the spread of access time for the critical resource. While a client's access time is the time between it requesting and accessing the resource, the spread is defined as a system-wide metric that measures some notion of the variance of access times across a homogeneous host population, e.g., difference between max and mean. We analytically prove the correctness of our algorithm, and evaluate its fairness experimentally using simulations. Our evaluation under two settings - a LAN setting and a WAN based on the King latency data set - shows even with 100 hosts accessing one resource, the spread of access time is within 15 seconds. Vijay Anand Korthikanti, Prateek Mittal, Indranil Gupta |
ICDCS | 3 |
| 2008 | Towards Multi-Site Collaboration in 3D Tele-Immersive Environmentsabstract3D tele-immersion (3DTI) has recently emerged as a new way of video-mediated collaboration across the Internet. Unlike conventional 2D video-conferencing systems, it can immerse remote users into a shared 3D virtual space so that they can interact or collaborate "virtually". However, most existing 3DTI systems can support only two sites of collaboration, due to the huge demand of networking resources and the lack of a simple yet efficient data dissemination model. In this paper, we propose to use a general publish-subscribe model for multi-site 3DTI systems, which efficiently utilizes limited network resources by leveraging user interest. We focus on the overlay construction problem in the publish-subscribe model by exploring a spectrum of heuristic algorithms for data dissemination. With extensive simulation, we identify the advantages of a simple randomized algorithm. We present optimization to further improve the randomized algorithm by exploiting semantic correlation. Experimental results demonstrate that we can achieve an improvement by a factor of five. Wanmin Wu, Zhenyu Yang 0006, Indranil Gupta, Klara Nahrstedt |
ICDCS | 3 |
| 2008 | AdapCode: Adaptive Network Coding for Code Updates in Wireless Sensor NetworksabstractCode updates, such as those for debugging purposes, are frequent and expensive in the early development stages of wireless sensor network applications. We propose AdapCode, a reliable data dissemination protocol that uses adaptive network coding to reduce broadcast traffic in the process of code updates. Packets on every node are coded by linear combination and decoded by Gaussian elimination. The core idea in AdapCode is to adaptively change the coding scheme according to the link quality. Our evaluation shows that AdapCode uses up to 40% less packets than Deluge in large networks. In addition, AdapCode performs much better in terms of load balancing, which prolongs the system lifetime, and has a slightly shorter propagation delay. Finally, we show that network coding is doable on sensor networks in that (i) it imposes only a 3 byte header overhead, (ii) it is easy to find linearly independent packets, and (3) Gaussian elimination needs only 1 KB of memory. I-Hong Hou, Yu-En Tsai, Tarek F. Abdelzaher, Indranil Gupta |
INFOCOM | 4 |
| 2008 | Moara: Flexible and Scalable Group-Based Querying System
Steven Y. Ko, Praveen Yalagandula, Indranil Gupta, Vanish Talwar, Dejan S. Milojicic, Subu Iyer |
Middleware | 3 |
| 2008 | Using Tractable and Realistic Churn Models to Analyze Quiescence Behavior of Distributed ProtocolsabstractLarge-scale distributed systems are subject to churn, i.e., continuous arrival, departure and failure of processes. Analysis of protocols under churn requires one to use churn models that are tractable (easy to apply), realistic (apply to deployment settings), and general (apply to many protocols and properties). In this paper, we propose two new churn models - called train and crowd - that together achieve these goals, for a broad class of stability properties called quiescent properties, and for arbitrary distributed protocols. We show (i) how analysis of protocol quiescence in the train model can be extended to the crowd model, (ii) how to apply the train and crowd model to several distributed membership protocols, (iii) how, even under real churn traces, the train and crowd models are reasonably good at predicting system-wide stability metrics for membership protocols. Steven Y. Ko, Imranul Hoque, Indranil Gupta |
SRDS | 3 |
| 2008 | A new class of nature-inspired algorithms for self-adaptive peer-to-peer computingabstractWe present, and evaluate benefits of, a design methodology for translating natural phenomena represented as mathematical models, into novel, self-adaptive, peer-to-peer (p2p) distributed computing algorithms ( protocols ). Concretely, our first contribution is a set of techniques to translate discrete sequence equations (also known as difference equations) into new p2p protocols called sequence protocols . Sequence protocols are self-adaptive, scalable, and fault-tolerant, with applicability in p2p settings like Grids. A sequence protocol is a set of probabilistic local and message-passing actions for each process. These actions are translated from terms in a set of source sequence equations. Individual processes do not simulate the source sequence equations completely. Instead, each process executes probabilistic local and message passing actions, so that the emergent round-to-round behavior of the sequence protocol in a p2p system can be probabilistically predicted by the source sequence equations. The article's second contribution is the design and evaluation of a set of sequence protocols for detection of two global triggers in a distributed system: threshold detection and interval detection. This article's third contribution is a new self-adaptive Grid computing protocol called HoneyAdapt. HoneyAdapt is derived from sequence equations modeling adaptive bee foraging behavior in nature. HoneyAdapt is intended for Grid applications that allow Grid clients, at run-time, a choice of algorithms for executing chunks of the application's dataset. HoneyAdapt tells each Grid client how to adaptively select at run-time, for each chunk it receives, a good algorithm for computing the chunk—this selection is based on continuous feedback from other clients. Finally, we design a variant of HoneyAdapt, called HoneySort, for application to Grid parallelized sorting settings using the master-worker paradigm. Our evaluation of these contributions consists of mathematical analysis, large-scale trace-based simulation results, and experimental results from a HoneySort deployment. Steven Y. Ko, Indranil Gupta, Yookyung Jo |
ACM Trans. Auton. Adapt. Syst. | 2 |
| 2008 | Adaptive probability-based broadcast forwarding in energy-saving sensor networksabstractNetworking protocols for multihop wireless sensor networks (WSNs) are required to simultaneously minimize resource usage as well as optimize performance metrics such as latency and reliability. This article explores the energy-latency-reliability tradeoff for broadcast in WSNs by presenting a new protocol called PBBF . Essentially, for a given reliability level, energy and latency are found to be inversely related and our study quantifies this relationship at the reliability boundary. Therefore, PBBF offers an application designer considerable flexibility in the choice of desired operation points. Furthermore, we propose an extension to dynamically adjust the PBBF parameters to minimize the input required from the designer. Cigdem Sengul, Indranil Gupta, Matthew J. Miller |
ACM Trans. Sens. Networks | 2 |
| 2007 | Measurement of a large-scale overlay for multimedia streamingabstractNo abstract available. Long H. Vu, Indranil Gupta, Klara Nahrstedt |
HPDC | 2 |
| 2007 | AVMON: Optimal and Scalable Discovery of Consistent Availability Monitoring Overlays for Distributed SystemsabstractThis paper addresses the problem of selection and discovery of a consistent availability monitoring overlay for computer hosts in a large-scale distributed application, where hosts may be selfish or colluding. We motivate six significant goals for the problem - consistency, verifiability, and randomness, in selecting the availability monitors of nodes, as well as discoverability, load-balancing, and scalability in finding these monitors. We then present a new system, called AVMON, that is the first to satisfy these six requirements. The core algorithmic contribution of this paper is a protocol for discovering the availability monitoring overlay in a scalable and efficient manner, given any arbitrary monitor selection scheme that is consistent and verifiable. We mathematically analyze the performance of AVMON's discovery protocols, and derive an optimal variant that minimizes memory, bandwidth, computation, and discovery time of monitors. Our experimental evaluations of AVMON use three types of availability traces - synthetic, from PlanetLab, and from a peer-to-peer system (Overnet) - and demonstrate that AVMON works well in a variety of distributed systems. Ramsés Morales, Indranil Gupta |
ICDCS | 2 |
| 2007 | A Cross-Layer Architecture to Exploit Multi-Channel Diversity with a Single TransceiverabstractThe design of multi-channel multi-hop wireless mesh networks is centered around the way nodes synchronize when they need to communicate. However, existing designs are confined to the MAC layer -they are based on either negotiation on a rendezvous control channel, or on optimistic synchronization. Both approaches scale poorly as the network grows in coverage and density. The rendezvous control channel may become the bottleneck, while optimistic synchronization may incur substantial overhead - especially amongst nodes close to a gateway, where the mesh traffic converges. In this paper, we describe Dominion - a cross-layer architecture that includes both medium access control and routing. At the MAC layer, a node switches channels in a deterministic manner to address the scalability issue. At the network layer, a Dominion node routes traffic along the shortest distance across both spatial and frequency domains, based on the deterministic channel-hopping schedule and network connectivity. Since the shortest distance path across the frequency domain is time variant, Dominion naturally spreads the packets of the same flow across multiple paths, relieving the intra-flow and inter-flow contention, and improving throughput. Through QualNet simulations we show that Dominion is able to achieve, on average, 1813% higher aggregate distance-normalized throughput than IEEE 802.11, while being 1730% fairer (using Jain's fairness index) with 50 simultaneous random flows. Jay A. Patel, Haiyun Luo, Indranil Gupta |
INFOCOM | 3 |
| 2007 | New Worker-Centric Scheduling Strategies for Data-Intensive Grid Applications
Steven Y. Ko, Ramsés Morales, Indranil Gupta |
Middleware | 3 |
| 2007 | AVMEM - Availability-Aware Overlays for Management Operations in Non-cooperative Distributed Systems
Ramsés Morales, Indranil Gupta |
Middleware | 3 |
| 2007 | Measurement and modeling of a large-scale overlay for multimedia streamingabstractThis paper presents results from our measurement and modeling efforts on the large-scale peer-to-peer (p2p) overlay graphs spanned by the PPLive system which is arguably the most popular and largest multimedia streaming p2p system today. We believe that our findings can be used to understand large-scale p2p streaming systems for future planning of resource usage, and to provide useful and practical hints for future design of large-scale p2p streaming systems. Unlike other previous studies on PPLive, which focused on either network-centric or user-centric measurements of the system, our study is unique in (a) focusing on PPLive overlay-specific characteristics, and (b) being the first to derive mathematical models for its distributions of channel population size and session length. Long H. Vu, Indranil Gupta, Klara Nahrstedt |
QSHINE | 2 |
| 2007 | Building trees based on aggregation efficiency in sensor networks
Albert F. Harris III, Robin Kravets, Indranil Gupta |
Ad Hoc Networks | 3 |
| 2007 | The design of novel distributed protocols from differential equations
Indranil Gupta, Mahvesh Nagda, Christo Frank Devaraj |
Distributed Comput. | 1 |
| 2007 | Active and passive techniques for group size estimation in large-scale and dynamic distributed systems
Dionysios Kostoulas, Dimitrios Psaltoulis, Indranil Gupta, Kenneth P. Birman, Alan J. Demers |
J. Syst. Softw. | 3 |
| 2007 | MOve: Design and Evaluation of a Malleable Overlay for Group-Based ApplicationsabstractWhile peer-to-peer overlays allow distributed applications to scale and tolerate failures, most structured and unstructured overlays in literature today are inflexible from the application viewpoint. The application thus has no first-class control on the overlay structure. This paper proposes the concept of an application-malleable overlay, and the design of the first malleable overlay which we call MOve. MOve is targeted at group- based applications, e.g., collaborative applications. In MOve, the communication characteristics of the distributed application using the overlay can influence the overlay's structure itself, with the twin goals of (1) optimizing the application performance by adapting the overlay, while also (2) retaining the large scale and fault tolerance of the overlay approach. Besides neighbor list membership management, MOve also contains algorithms for resource discovery, update propagation, and churn-resistance. The emergent behavior of the implicit mechanisms used in MOve manifests as follows: when application communication is low, most overlay links keep their default configuration; however, as application communication characteristics become more evident, the overlay gracefully adapts itself to the application. We validate MOve using simulations with group sizes that are fixed, uniform, exponential and PlanetLab-based (slices), as well as churn traces and two sample management-based applications. Ramsés Morales, Sébastien Monnet, Indranil Gupta, Gabriel Antoniu |
IEEE Trans. Netw. Serv. Manag. | 3 |
| 2007 | AVCast: New Approaches for Implementing Generic Availability-Dependent Reliability Predicates for Multicast ReceiversabstractToday's large-scale distributed systems consist of collections of nodes, each of which has its own availability characteristics - a phenomenon sometimes called churn. This availability variation across nodes is often a hindrance to achieving reliability and performance for distributed applications such as multicast. This paper looks into utilizing and leveraging availability information in order to implement arbitrary predicates that specify availability-dependent message reliability for multicast receivers. An application (e.g., a publish-subscribe system) may want to scale the multicast message reliability at each receiver according to that receiver's availability (in terms of the fraction of time that receiver is online) - different options are that the reliability is independent of the availability, proportional to it, or an arbitrary function of it, etc. We propose several gossip- based algorithms to support an arbitrary class of such predicates. These techniques rely on each node's availability being monitored in a distributed manner by a small group of other nodes in such a way that the monitoring load is evenly distributed in the system. Our techniques are light-weight, scalable, and are space- and time- efficient. We analyze our algorithms and evaluate them experimentally by injecting availability traces collected from real peer-to-peer systems. Thadpong Pongthawornkamol, Indranil Gupta |
IEEE Trans. Netw. Serv. Manag. | 2 |
| 2006 | QoS-aware Object Replication in Overlay NetworksabstractMany emerging applications for peer to peer overlays may require nodes to satisfy strict timing deadlines to access a replica of a given object. This includes multimedia and hard realtime applications such as distributed gaming. We formulate the QoS-aware replication problem, the goal of which is to locate the minimum number of replicas to satisfy access time deadlines for all nodes while minimizing storage usage in the overlay. Existing replication schemes cannot be used to solve this problem since they are best-effort only. We show that finding a solution to the QoS-aware object replication in an arbitrary overlay topology is intractable (NP-complete). We then present simple centralized as well as decentralized heuristics for QoS-aware replication, and compare their performance experimentally. In addition, we investigate how these decentralized heuristics effectively works in a real network. Won Jong Jeon, Indranil Gupta, Klara Nahrstedt |
GLOBECOM | 2 |
| 2006 | PriorityCast: Efficient and Time-Critical Decision Making in First Responder Ad-Hoc NetworksabstractFirst responders equipped with PDAs in disaster recovery scenarios may need to make decisions that select the best out of multiple options. Each of these options may originate from different parts of the network (e.g., risk assessment of various entrances to a collapsed building). Yet the highest priority option needs to be spread to all nodes within a time deadline and in an efficient manner, i.e., by preventing lower priority options from spreading too far. We define this problem as PriorityCast, and study a wide range of possible solutions. The emphasis is on solutions that are scalable, and resilient to unreliability and unpredictability. Our basic scheme, called flood and suppress, suppresses lower priority options from being forwarded by a node that has seen other higher priority options. We then augment this basic scheme using probabilistic approaches, background gossiping techniques, and delayed propagation. We quantify the impact of using various combinations of the mentioned approaches, by presenting mathematical analysis for the flood and suppress protocol, and evaluating a suite of composable PriorityCast protocols via simulations. The results provide insight into the feasibility and scalability of this class of solutions. While the focus of this paper is on the PriorityCast problem, these solutions are also relevant as potential building blocks for other distributed operations in ad-hoc networks Vartika Bhandari, Indranil Gupta |
MASS | 2 |
| 2006 | Smart Gossip: An Adaptive Gossip-based Broadcasting Service for Sensor NetworksabstractA network-wide broadcast service is often used for information dissemination in sensor networks. Sensor networks are typically energy-constrained and prone to failures. In view of these constraints, the broadcast service should minimize energy consumption by reducing redundant transmissions, and be tolerant to frequent node and link failures. We propose "smart gossip", a probabilistic protocol that offers a broadcast service with low overheads. Smart gossip automatically and dynamically adapts transmission probabilities based on the underlying network topology. The protocol is capable of coping with wireless losses and unpredictable node failures that affect network connectivity over time. The resulting protocol is completely decentralized. We present thorough experimental results to evaluate our "smart gossip" proposal, and demonstrate its benefits over existing protocols Pradeep Kyasanur, Romit Roy Choudhury, Indranil Gupta |
MASS | 3 |
| 2006 | JetStream: Achieving Predictable Gossip Dissemination by Leveraging Social Network PrinciplesabstractGossip protocols provide probabilistic reliability and scalability, but their inherent randomness may lead to high variation in number of messages that are received at different nodes. This paper presents techniques that leverage simple social network principles enabling nodes to select gossip targets intelligently. The simple heuristics presented in the paper achieve a more uniform message overhead at each node, lowering the system-wide gossip traffic, while simultaneously reducing the latency of gossip spread (by up to 25%). We experimentally compare our system, called JetStream, against canonical gossip as well as gossip on the chord overlay. Intuitively, JetStream seeks to make gossip spread more deterministic and predictable, while still inheriting its scale and reliability. JetStream also provides an added benefit by reducing network bandwidth utilization with a low sustained rate of gossip injection Jay A. Patel, Indranil Gupta, Noshir S. Contractor |
NCA | 2 |
| 2006 | ContagAlert: Using Contagion Theory for Adaptive, Distributed Alert PropagationabstractLarge-scale distributed systems, e.g., grid or P2P networks, are targets for large-scale attacks. Unfortunately, few existing systems support propagation of alerts during the attack itself while also suppressing disruptive alerts from faulty or malicious sources. This paper proposes the "ContagAlert" protocol, which uses contagion spreading behavior to spread alerts. ContagAlert rapidly propagates alerts during attacks while also suppressing disruptive alerts. The core contagion protocols in the system are completely localized, but result in desired behavior at the network scale. We analyze and evaluate our protocol with synthetic simulations and in both Internet worm and DoS attack scenarios Michael Treaster, William Conner, Indranil Gupta, Klara Nahrstedt |
NCA | 3 |
| 2006 | MOve: Design of An Application-Malleable OverlayabstractPeer-to-peer overlays allow distributed applications to work in a wide-area, scalable, and fault-tolerant manner. However, most structured and unstructured overlays present in literature today are inflexible from the application viewpoint. In other words, the application has no control over the structure of the overlay itself. This paper proposes the concept of an application-malleable overlay, and the design of the first malleable overlay which we call MOve. In MOve, the communication characteristics of the distributed application using the overlay can influence the overlay's structure itself, with the twin goals of (1) optimizing the application performance by adapting the overlay, while also (2) retaining the large scale and fault tolerance of the overlay approach. The influence could either be explicitly specified by the application or implicitly gleaned by our algorithms. Besides neighbor list membership management, MOve also contains algorithms for resource discovery, update propagation, and churn-resistance. The emergent behavior of the implicit mechanisms used in MOve manifests in the following way: when application communication is low, most overlay links keep their default configuration; however, as application communication characteristics become more evident, the overlay gracefully adapts itself to the application Sébastien Monnet, Ramsés Morales, Gabriel Antoniu, Indranil Gupta |
SRDS | 4 |
| 2006 | AVCast : New Approaches For Implementing Availability-Dependent Reliability for Multicast ReceiversabstractToday's large-scale distributed systems consist of collections of nodes that have highly variable availability - a phenomenon sometimes called churn. This availability variation is often a hindrance to achieving reliability and performance for distributed applications such as multicast. This paper looks into utilizing and leveraging availability information in order to provide availability-dependent message reliability for multicast receivers. An application (e.g., a publish-sub scribe system) may want to scale the multicast message reliability at each receiver according to that receiver's availability (in terms of the fraction of time that receiver is online)ifferent options are that the reliability is independent of the availability, or proportional to it. We propose several gossip-based algorithms to support several such predicates. These techniques rely on each node's availability being monitored in a distributed manner by a small group of other nodes in such a way that the monitoring load is evenly distributed in the system. Our techniques are light-weight, scalable, and are space- and time-efficient. We analyze our algorithms and evaluate them experimentally by injecting availability traces collected from real peer-to-peer systems Thadpong Pongthawornkamol, Indranil Gupta |
SRDS | 2 |
| 2006 | Efficient and Adaptive Epidemic-Style Protocols for Reliable and Scalable MulticastabstractEpidemic-style (gossip-based) techniques have recently emerged as a class of scalable and reliable protocols for peer-to-peer multicast dissemination in large process groups. However, popular implementations of epidemic-style dissemination suffer from two major drawbacks: 1) Network overhead: when deployed on a WAN-wide or VPN-wide scale, they generate a large number of packets that transit across the boundaries of multiple network domains (e.g., LANs, subnets, ASs), causing an overload on core network elements such as bridges, routers, and associated links. 2) Lack of adaptivity: they impose the same load on process group members and the network even under reduced failure rates (viz., packet losses, process failures). In this paper, we describe two protocols to address these problems: 1) a hierarchical gossiping protocol and 2) an adaptive dissemination framework (for multicasts) that allows use of any gossiping primitive within it. These protocols work within a virtual peer-to-peer hierarchy called the leaf box hierarchy. Processes can be allocated in a topologically aware manner to the leaf boxes of this structure, so that protocols 1 and 2 produce low traffic across domain boundaries in the network and induce minimal overhead when there are no failures Indranil Gupta, Anne-Marie Kermarrec, Ayalvadi J. Ganesh |
IEEE Trans. Parallel Distributed Syst. | 1 |
| 2005 | Providing both scale and security through a single core probabilistic protocolabstractDistributed systems are typically designed for scale and performance first, which makes it difficult to add security later without affecting the original properties. This paper proposes the design of the Folklore persistent distributed storage system, which adopts an alternative design methodology. Folklore's design relies on a single core protocol for providing both probabilistic scalability and untraceability, the latter being a special notion of probabilistic security. The core protocol is a biologically inspired model of endemic replication that migrates replicas of files among all hosts in a continuous and proactive manner. The emergent behavior is chaotic, meaning that the exact number and location of all replicas of any file is changing all the time. This makes it difficult for an attacker to target any file. Yet, the protocol is scalable - it consumes constant per-host bandwidth, and the number of replicas per file stays close to a small self-stabilizing value. The self-stabilizing value is reached even if only one replica survives a massive attack. The simplicity of the core protocol allows augmentation with mechanisms that allow data integrity, availability, and updatability. We describe the internals of the Folklore system, present attack analysis, and give experimental results from a prototype that shows high resilience to large-scale attacks Ramsés Morales, Indranil Gupta |
CollaborateCom | 2 |
| 2005 | The P2P MultiRouter: a black box approach to run-time adaptivity for P2P DHTsabstractPeer-to-peer distributed hash tables (P2P DHTs) are individually built by their designers with specific performance goals in mind. However, no individual DHT can satisfy an application that requires a "best of all worlds" performance, viz., adaptive behavior at run-time. We propose the MultiRouter, a light-weight solution that provides adaptivity to the application using a DHT-independent approach. By merely making run-time choices to select from among multiple DHT protocols using simple cost functions, we show the MultiRouter is able to provide a best-of-all-DHTs run-time performance with respect to object access times and churn-resistance. In addition, the MultiRouter is not limited to any particular set of DHT implementations since the interaction occurs in a black box manner, i.e., through well-defined interfaces. We present microbenchmark and trace-driven experiments to show that if one fixes bandwidth at each node, the MultiRouter outperforms the component DHTs. James S. K. Newell, Indranil Gupta |
CollaborateCom | 2 |
| 2005 | Perturbation-Resistant and Overlay-Independent Resource DiscoveryabstractThis paper realizes techniques supporting the position that strategies for resource location and discovery in distributed systems should be both perturbation-resistant and overlay-independent. Perturbation-resistance means that inserts and lookups must be robust to ordinary stresses such as node perturbation, which may arise out of congestion, competing client applications, or user churn. Overlay-independence implies that the insert and lookup strategies, and to an extent their performance, should be independent of the actual structure of the underlying overlay. We first show how a well-known distributed hash table (Pastry) may degrade under perturbation. We then present a new resource location and discovery algorithm called MPIL (multi-path insertion/lookup) that is perturbation-resistant and overlay-independent. MPIL is overlay-independent in that it effectively provides to the distributed application an ability to insert and lookup Pastry objects in an overlay with Pastry IDs, but without the need to have Pastry-style overlay maintenance (i.e., the overlay underneath can be arbitrary). We quantify, through analysis and simulation results, the behavior of MPIL over complete, random, and power-law overlays. We also show how MPIL outperforms regular Pastry routing when there is perturbation. Steven Y. Ko, Indranil Gupta |
DSN | 2 |
| 2005 | Exploring the Energy-Latency Trade-Off for Broadcasts in Energy-Saving Sensor NetworksabstractNetworking protocols for multi-hop wireless sensor networks (WSNs) are required to simultaneously minimize resource usage as well as optimize performance metrics such as latency and reliability. This paper explores the energy-latency-reliability trade-off for broadcast in multi-hop WSNs, by presenting a new protocol called PBBF (probability-based broadcast forwarding). PBBF works at the MAC layer and can be integrated into any sleep scheduling protocol. For a given application-defined level of reliability for broadcasts, the energy required and latency obtained are found to be inversely related to each other. Our analysis and simulation study quantify this relationship at the reliability boundary, as well as performance numbers to be expected from a deployment. PBBF essentially offers a WSN application designer considerable flexibility in choice of desired operation points Matthew J. Miller, Cigdem Sengul, Indranil Gupta |
ICDCS | 3 |
| 2005 | An underlay for sensor networks: localized protocols for maintenance and usageabstractWe propose localized and decentralized protocols to construct and maintain an underlay for sensor networks. An underlay lies in between overlay operations (e.g., data indexing, multicast, etc.) and the sensor network itself. Specifically, an underlay bridges the gap between (a) the unreliability of sensor nodes and communication and availability of only approximate location knowledge, and (b) the maintenance of a virtual geography-based naming structure that is required by several overlay operations. Our underlay creates a coarse naming scheme based on approximate location knowledge, and then maintains it in an efficient and scalable manner. The underlay naming can be used to specify arbitrary regions. The overlay operations that can be executed on the underlay include routing, aggregation, multicast, data indexing, etc. These overlay operations could be region-based. The proposed underlay maintenance protocols are robust, localized (hence scalable), energy and message efficient, have low convergence times, and provide tuning knobs to trade convergence time with overhead and with underlay uniformity. The maintenance protocols are mathematically analyzed by characterizing them as differential equation systems. We present microbenchmark results from a NesC implementation, and results from a large-scale simulation of a Java implementation. The latter experiments also show how routing using the underlay would perform Christo Frank Devaraj, Indranil Gupta, Mahwish Nagda, Gul A. Agha |
MASS | 2 |
| 2005 | Cushion: autonomically adaptive data fusion in wireless sensor networksabstractIn a typical in-network aggregation problem, data originates from multiple source nodes, and moves towards single sink node or root node. Along the way, i.e., inside the network, the data may be partially aggregated, thus reducing message overhead. Two well-known classes of solutions to this problem are: tree-based aggregation and multipath aggregation. The autonomically adaptive data fusion solution, called Cushion, contains two protocols that span a continuous spectrum between these two design points and can further increase reliability over multi-path approach using redundant transmissions. The design is motivated of Cushion by first discussing the complementary advantages and disadvantages of the two design points. Jungmin So, Indranil Gupta |
MASS | 3 |
| 2005 | Decentralized Schemes for Size Estimation in Large and Dynamic GroupsabstractLarge-scale and dynamically changing distributed systems such as the Grid, peer-to-peer overlays, etc., need to collect several kinds of global statistics in a decentralized manner. In this paper, we tackle a specific statistic collection problem called Group Size Estimation, for estimating the number of non-faulty processes present in the global group at any given point of time. We present two new decentralized algorithms for estimation in dynamic groups, analyze the algorithms, and experimentally evaluate them using real-life traces. One scheme is active: it spreads a gossip into the overlay first, and then samples the receipt times of this gossip at different processes. The second scheme is passive: it measures the density of processes when their identifiers are hashed into a real interval. Both schemes have low latency, scalable perprocess overheads, and provide high levels of probabilistic accuracy for the estimate. They are implemented as part of a size estimation utility called PeerCounter that can be incorporated modularly into standard peer-to-peer overlays. We present experimental results from both the simulations and PeerCounter, running on a cluster of 33 Linux servers. Dionysios Kostoulas, Dimitrios Psaltoulis, Indranil Gupta, Kenneth P. Birman, Alan J. Demers |
NCA | 3 |
| 2005 | MON: management overlay networks for distributed systemsabstractThe recent deployment of large distributed computing systems such as content distribution networks and the Planet-Lab has made it possible for researchers and practitioners to experiment with real world, large scale distributed applications. However, running an application in such an environment is difficult, due to the scale and frequent node failures of such systems. Thus, an important tool is needed that helps application developers/deployers to manage their applications. Our goal in this work is to develop MON, an extremely lightweight and failure resilient system for managing distributed applications. MON allows users to execute instant management commands on the distributed computing nodes, such as query the current status of the application, or start/stop a process on the distributed nodes. The commands are propagated to all the nodes and executed on each node, and the results are aggregated and returned back. We believe the ability to execute such instant commands is especially useful for the initial deployment of a distributed application, or for the monitoring and diagnoistics of (unexpected) application failures. Steven Y. Ko, Indranil Gupta, Klara Nahrstedt |
SOSP | 3 |
| 2005 | Turning flash crowds into smart mobs with real-time stochastic detection and adaptive cooperative cachingabstractThe past decade has experienced a continued increase in the popular use of Web services. Unfortunately, the growth of Web services, coupled with their their role as a vital source of news and information has lead to many well published sources of stress that tend to cripple performance at both the client and the server side [1, 4]. Our work considers one such type of stress, called a flash crowd. A flash crowd arrives as a tidal wave, where the initial (and largest) spike in traffic generally occurs within the first few minutes. This gives client and servers only a few tens of seconds to adapt to the incoming traffic! In addition, flash crowds are infrequent and unpredictable. Hence, any solution deployed to mitigate flash crowd must consist of: (1) proactive and real-time mechanisms to detect the flash crowd while (or even before) it happens, and (2) an initiation of immediate action that can mask the effects of the stress. Jay A. Patel, Charles M. Yang, Indranil Gupta |
SOSP | 3 |
| 2005 | A framework for time indexing in sensor networksabstractIn this article, we define the time-indexing problem as the in-network storage and querying of sensor network data based solely on the time attribute. We argue qualitatively why existing storage schemes may be insufficient as solutions. We then present, analyze, and evaluate novel and lightweight solutions to both the storage and the querying subproblems for time indexing. First, the time-indexed storage problem is formally defined and two formulations are presented seeking to optimize generic utility functions that are derived from concerns about energy, bandwidth usage, and storage balancing. We present and analyze decentralized protocols to solve these formulations and prove the optimality of some of our solutions. Secondly, maintenance and use of simple overlays among rendezvous point nodes in order to enable fault-tolerant and efficient time-indexed queries are discussed. Finally, simulation results are presented to quantify performance characteristics of the protocols, and we find that our proposed scheme has low query overhead that scales with system size and density while exhibiting very good load-balancing and fault-tolerance properties. The use of time-indexed structure is shown to achieve more than double the lifetime of sensor networks compared to existing approaches in some scenarios. Rong Zheng 0001, Indranil Gupta, Lui Sha |
ACM Trans. Sens. Networks | 3 |
| 2004 | Time indexing in sensor networksabstractWe define the time indexing problem as the in-network storage and querying of sensor network data based solely on the time attribute. We argue qualitatively why existing storage schemes may be insufficient as solutions. We then present, analyze, and evaluate novel and lightweight solutions to both the storage and the querying sub-problems for time indexing. First, the time-indexed storage problem is formally defined, and two formulations are presented, seeking to optimize generic utility functions that are derived from concerns about energy, bandwidth usage, and storage balancing. We present and analyze decentralized protocols to solve these formulations, and prove the optimality of some of our solutions. Secondly, maintenance and use of simple overlays among rendezvous point nodes, in order to enable fault-tolerant and efficient time-indexed queries, are discussed. Finally, simulation results are presented to quantify performance characteristics of the protocols, and we find that our proposed scheme has low query overhead that scales with system size and density while exhibiting very good load balancing and fault tolerance properties. Rong Zheng 0001, Indranil Gupta, Lui Sha |
MASS | 3 |
| 2004 | On the design of distributed protocols from differential equationsabstractWe propose a framework to translate systems of differential equations into distributed protocols. The synthesized protocols are state machines containing probabilistic transitions and actions, and they show equivalent stochastic behavior to that in the original equations. We apply the framework to the Lotka-Volterra model of competition, in order to design a scalable and probabilistically reliable protocol for majority selection (and thus scalable and eventual consensus). We propose two subclasses of equations with polynomial terms, and prove the equivalence of protocols generated by the proposed mapping techniques to the original differential equations. The protocols are probabilistically scalable and reliable. Equation rewriting techniques are also described. Endemic protocols for migratory replication, and well-known epidemic protocols are also shown to be generated using this design methodology. We present mathematical analysis of the protocols, and experimental results from our implementations. We also discuss limitations of our approach. We believe the design framework could be effectively used in transforming, albeit in a systematic manner, well-known natural phenomena into protocols for distributed systems. Indranil Gupta |
PODC | 1 |
| 2002 | SWIM: Scalable Weakly-consistent Infection-style Process Group Membership ProtocolabstractSeveral distributed peer-to-peer applications require weakly-consistent knowledge of process group membership information at all participating processes. SWIM is a generic software module that offers this service for large scale process groups. The SWIM effort is motivated by the unscalability of traditional heart-beating protocols, which either impose network loads that grow quadratically with group size, or compromise response times or false positive frequency w.r.t. detecting process crashes. This paper reports on the design, implementation and performance of the SWIM sub-system on a large cluster of commodity PCs. Unlike traditional heart beating protocols, SWIM separates the failure detection and membership update dissemination functionalities of the membership protocol. Processes are monitored through an efficient peer-to-peer periodic randomized probing protocol. Both the expected time to first detection of each process failure, and the expected message load per member do not vary with group size. Information about membership changes, such as process joins, drop-outs and failures, is propagated via piggybacking on ping messages and acknowledgments. This results in a robust and fast infection style (also epidemic or gossip-style) of dissemination. The rate of false failure detections in the SWIM system is reduced by modifying the protocol to allow group members to suspect a process before declaring it as failed - this allows the system to discover and rectify false failure detections. Finally, the protocol guarantees a deterministic time bound to detect failures. Experimental results from the SWIM prototype are presented. We discuss the extensibility of the design to a WAN-wide scale. Abhinandan Das, Indranil Gupta, Ashish Motivala |
DSN | 2 |
| 2002 | Efficient Epidemic-Style Protocols for Reliable and Scalable MulticastabstractEpidemic-style (gossip-based) techniques have recently emerged as a scalable class of protocols for peer-to-peer reliable multicast dissemination in large process groups. These protocols provide probabilistic guarantees on reliability and scalability. However, popular implementations of epidemic-style dissemination are reputed to suffer from two major drawbacks: (a) (Network Overhead) when deployed on a WAN-wide or VPN-wide scale they generate a large number of packets that transit across the boundaries of multiple network domains (e.g., LANs, subnets, ASs), causing an overload on core network elements such as bridges, routers, and associated links; (b) (Lack of Adaptivity) they impose the same load on process group members and the network even under reduced failure rates (viz., packet losses, process failures). lit this paper we report on the (first) comprehensive set of solutions to these problems. The solution is comprised of two protocols: (1) a hierarchical gossiping protocol, and (2) an adaptive multicast dissemination framework that allows use of any gossiping primitive within it. These protocols work within a virtual peer-to-peer hierarchy called the Leaf Box hierarchy. Processes can be allocated in a topologically aware manner to the leaf boxes of this structure, so that (1) and (2) produce low traffic across domain boundaries in the network. In the interests of space, this paper focuses on a detailed discussion and evaluation (through simulations) of only the hierarchical gossiping protocol. We present an overview of the adaptive dissemination protocol and its properties. Indranil Gupta, Anne-Marie Kermarrec, Ayalvadi J. Ganesh |
SRDS | 1 |
| 2001 | GulfStream - a System for Dynamic Topology Management in Multi-domain Server FarmsabstractThis paper describes GulfStream, a scalable distributed software system designed to address the problem of managing the network topology in a multi-domain server farm. In particular, it addresses the following core problems: topology discovery and verification, and failure detection. Unlike most topology discovery and failure detection systems which focus on the nodes in a cluster, GulfStream logically organizes the network adapters of the server farm into groups. Each group contains those adapters that can directly exchange messages. GulfStream dynamically establishes a hierarchy for reporting network topology and availability of network adapters. We describe a prototype implementation of GulfStream on a 55 node heterogeneous server farm interconnected using switched fast Ethernet. Sameh A. Fakhouri, Germán S. Goldszmidt, Michael H. Kalantar, John A. Pershing, Indranil Gupta |
CLUSTER | 5 |
| 2001 | Scalable Fault-Tolerant Aggregation in Large Process GroupsabstractThe paper discusses fault-tolerant, scalable solutions to the problem of accurately and scalably calculating global aggregate functions in large process groups communicating over unreliable networks. These groups could represent sensors or processes communicating over a network that is either fixed (e.g., the Internet) or dynamic (e.g., multihop ad-hoc). Group members are prone to failures. The ability to evaluate global aggregate properties (e.g., the average of sensor temperature readings) is important for higher-level coordination activities in such large groups. We first define the setting and problem, laying down metrics to evaluate different algorithms for the same. We discuss why the usual approaches to solve this problem are unviable and unscalable over an unreliable network prone to message delivery failures and crash failures. We then propose a technique to impose an abstract hierarchy on such large groups, describing how this hierarchy can be made to mirror the network topology. We discuss several alternatives to use this technique to solve the global aggregate function evaluation problem. Finally, we present a protocol based on gossiping that uses this hierarchical technique. We present mathematical analysis and performance results to validate the robustness, efficiency and accuracy of the Hierarchical Gossiping algorithm. Indranil Gupta, Robbert van Renesse, Kenneth P. Birman |
DSN | 1 |
| 2001 | Minimal CDMA Recoding Strategies in Power-Controlled Ad-Hoc Wireless NetworksabstractThe problem of Code Division Multiple Access (CDMA) code assignment to eliminate primary and hidden collisions in multihop packet radio networks has been widely researched in the past. However, very little work has been done on the very realistic *distributed, dynamic* version of the transmitter-oriented code assignment (TOCA) problem in an ad-hoc network where mobiles use CDMA technology. None of the existing dynamic TOCA CDMA algorithms in literature are efficient, in terms of maximum code index assigned in the network, or number of times a mobile has to change its code. We present a set of local and distributed *recoding* strategies for the TOCA CDMA problem in an ad-hoc network where mobiles can arbitrarily 1) connect and disconnect, 2) move about, and 3) increase or decrease their transmission power - all these may need some mobiles to be recoded, to avoid new collisions. Our strategies, unlike those proposed earlier in literature, guarantee *minimal recoding*, that is, given a current network-wide code assignment and one of the above events, our strategies change the codes of the minimum number of mobiles needed to eliminate all collisions. Minimal recoding can be very important in reducing the effect of frequent code changes on the performance and criticality of distributed applications. Further, among all possible minimal recoding strategies in a class, most of our strategies are also (provably) *optimal* in terms of the maximum code index assigned in the network. Performance results that evaluate our dynamic minimal strategies are also presented. Indranil Gupta |
IPDPS | 1 |
| 2001 | On scalable and efficient distributed failure detectorsabstractProcess groups in distributed applications and services rely on failure detectors to detect process failures completely, and as quickly, accurately, and scalably as possible, even in the face of unreliable message deliveries. In this paper, we look at quantifying the optimal scalability, in terms of network load, (in messages per second, with messages having a size limit) of distributed, complete failure detectors as a function of application-specified requirements. These requirements are 1) quick failure detection by some non-faulty process, and 2) accuracy of failure detection. We assume a crash-recovery (non-Byzantine) failure model, and a network model that is probabilistically unreliable (w.r.t. message deliveries and process failures). First, we characterize, under certain independence assumptions, the optimum worst-case network load imposed by any failure detector that achieves an application's requirements. We then discuss why traditional heart beating schemes are inherently unscalable according to the optimal load. We also present a randomized, distributed, failure detector algorithm that imposes an equal expected load per group member. This protocol satisfies the application defined constraints of completeness and accuracy, and speed of detection on an average. It imposes a network load that differs frown the optimal by a sub-optimality factor that is much lower than that for traditional distributed heartbeating schemes. Moreover, this sub-optimality factor does not vary with group size (for large groups). Indranil Gupta, Tushar Deepak Chandra, Germán S. Goldszmidt |
PODC | 1 |
| 2000 | A Probabilistically Correct Leader Election Protocol for Large Groups
Indranil Gupta, Robbert van Renesse, Kenneth P. Birman |
DISC | 1 |
| 2000 | A New Strategy for Improving the Effectiveness of Resource Reclaiming Algorithms in Multiprocessor Real-Time Systems
Indranil Gupta, G. Manimaran, C. Siva Ram Murthy |
J. Parallel Distributed Comput. | 1 |