EDBT 2026 Demo / reviewers in the wild / expert
Sriram Rao
dblp:07/1315
· DBLP profile ↗
31ranked-venue papers
3as first author
1since 2021 · last 2026
0000-0002-3265-6165ORCID · corroborated
Domains — the database's venue-derived domains; a paper can count in several
Databases, data management, data science and information retrieval · 11 · 1 first-author · 1 since 2021Systems, architecture and hardware · 10 · 1 first-authorComputer networks · 5Software engineering, systems software and programming languages · 3Security and privacy · 1 · 1 first-authorGraphics, computer vision, multimedia, augmented reality and games · 1
Expertise — from the expertise taxonomy: the topics of the expert's papers under the CCF categories. A weight counts papers with recency: 1 for a paper about the topic, 0.3 when the topic is its context, halved every five years.
| Computer architecture, parallel and distributed computing, and storage systems
20 papers |
Cloud and datacenter computing · 55% Distributed systems · 17% Storage systems · 11% | |
| Databases, data mining, and information retrieval
7 papers |
Query processing and optimization · 65% Data stream processing · 28% Distributed and cloud data management · 3% | |
| Computer networks
2 papers |
Datacenter networks · 100% |
Topics — the 30 heaviest of 48, each with the papers that count most for it
| Topic | Weight | Papers | Last | Evidence papers |
|---|---|---|---|---|
Cloud and datacenter computing
cluster resource management and scheduling |
1.7 | 7 | 2018 | Medea: scheduling of long running applications in shared production clusters · EuroSys 2018 Apache REEF: Retainable Evaluator Execution Framework · ACM Trans. Comput. Syst. 2017 Dhalion: Self-Regulating Stream Processing in Heron · Proc. VLDB Endow. 2017 |
Cloud and datacenter computing › cluster resource management and scheduling
cluster resource management |
1.5 | 7 | 2019 | Hydra: a federated resource manager for data-center scale analytics · NSDI 2019 Selecting Subexpressions to Materialize at Datacenter Scale · Proc. VLDB Endow. 2018 Morpheus: Towards Automated SLOs for Enterprise Clusters · OSDI 2016 |
Query processing and optimization
query optimization |
0.7 | 2 | 2018 | Towards a Learning Optimizer for Shared Clouds · Proc. VLDB Endow. 2018 Selecting Subexpressions to Materialize at Datacenter Scale · Proc. VLDB Endow. 2018 |
Cloud and datacenter computing › cluster resource management and scheduling
resource scheduling |
0.4 | 1 | 2019 | Hydra: a federated resource manager for data-center scale analytics · NSDI 2019 |
Query processing and optimization
cardinality estimation |
0.3 | 1 | 2018 | Towards a Learning Optimizer for Shared Clouds · Proc. VLDB Endow. 2018 |
Query processing and optimization › result reuse
computation reuse |
0.3 | 1 | 2018 | Computation Reuse in Analytics Job Service at Microsoft · SIGMOD Conference 2018 |
Query processing and optimization › cardinality estimation
learned cardinality estimation |
0.3 | 1 | 2018 | Towards a Learning Optimizer for Shared Clouds · Proc. VLDB Endow. 2018 |
Query processing and optimization › query optimization › learned query optimization
learned query optimizer |
0.3 | 1 | 2018 | Towards a Learning Optimizer for Shared Clouds · Proc. VLDB Endow. 2018 |
Query processing and optimization › materialized view
materialized view selection |
0.3 | 1 | 2018 | Computation Reuse in Analytics Job Service at Microsoft · SIGMOD Conference 2018 |
Datacenter networks › datacenter interconnect
inter-datacenter transfer |
0.3 | 1 | 2018 | QuickCast: Fast and Efficient Inter-Datacenter Transfers Using Forwarding Tree Cohorts · INFOCOM 2018 |
Distributed systems › stream processing
distributed stream processing |
0.3 | 1 | 2018 | Chi: A Scalable and Programmable Control Plane for Distributed Stream Processing Systems · Proc. VLDB Endow. 2018 |
Electronic design automation › physical design › placement
placement constraints |
0.3 | 1 | 2018 | Medea: scheduling of long running applications in shared production clusters · EuroSys 2018 |
Data stream processing
stream processing systems |
0.3 | 1 | 2017 | Dhalion: Self-Regulating Stream Processing in Heron · Proc. VLDB Endow. 2017 |
Cloud and datacenter computing
autoscaling |
0.3 | 1 | 2017 | Dhalion: Self-Regulating Stream Processing in Heron · Proc. VLDB Endow. 2017 |
Parallel and multicore computing › task scheduling
dependency-aware scheduling |
0.2 | 1 | 2016 | GRAPHENE: Packing and Dependency-Aware Scheduling for Data-Parallel Clusters · OSDI 2016 |
Embedded and real-time systems › real-time scheduling
queue management |
0.2 | 1 | 2016 | Efficient queue management for cluster scheduling · EuroSys 2016 |
Parallel and multicore computing
task allocation |
0.2 | 1 | 2016 | Efficient queue management for cluster scheduling · EuroSys 2016 |
Parallel and multicore computing › locality optimization
data locality optimization |
0.2 | 1 | 2015 | Network-Aware Scheduling for Data-Parallel Jobs: Plan When You Can · SIGCOMM 2015 |
Parallel and multicore computing › parallel scheduling
data-parallel job scheduling |
0.2 | 1 | 2015 | Network-Aware Scheduling for Data-Parallel Jobs: Plan When You Can · SIGCOMM 2015 |
Cloud and datacenter computing › job scheduling
network-aware scheduling |
0.2 | 1 | 2015 | Network-Aware Scheduling for Data-Parallel Jobs: Plan When You Can · SIGCOMM 2015 |
Cloud and datacenter computing › cluster resource management and scheduling
cluster scheduling |
0.2 | 1 | 2014 | Multi-resource packing for cluster schedulers · SIGCOMM 2014 |
Storage systems › storage management
fragmentation reduction |
0.2 | 1 | 2014 | Multi-resource packing for cluster schedulers · SIGCOMM 2014 |
Distributed systems
fault tolerance |
0.2 | 3 | 2017 | Apache REEF: Retainable Evaluator Execution Framework · ACM Trans. Comput. Syst. 2017 REEF: Retainable Evaluator Execution Framework · SIGMOD Conference 2015 The Cost of Recovery in Message Logging Protocols · IEEE Trans. Knowl. Data Eng. 2000 |
Distributed systems
consensus |
0.2 | 1 | 2013 | Tango: distributed data structures over a shared log · SOSP 2013 |
Distributed systems
distributed data structures |
0.2 | 1 | 2013 | Tango: distributed data structures over a shared log · SOSP 2013 |
Storage systems › file systems
distributed file system |
0.2 | 1 | 2013 | A The Quantcast File System · Proc. VLDB Endow. 2013 |
Storage systems › storage reliability
erasure coding |
0.2 | 1 | 2013 | A The Quantcast File System · Proc. VLDB Endow. 2013 |
Distributed systems › replication
replicated data types |
0.2 | 1 | 2013 | Tango: distributed data structures over a shared log · SOSP 2013 |
Storage systems › distributed storage
shared log |
0.2 | 1 | 2013 | Tango: distributed data structures over a shared log · SOSP 2013 |
Storage systems › object storage
cloud object store |
0.1 | 1 | 2012 | Walnut: a unified cloud object store · SIGMOD Conference 2012 |
Methods — techniques the papers use, named apart from their topics
vertex-centric graph algorithm · 0.7subgraph template learning · 0.7reactive programming · 0.7machine learning · 0.7integer linear programming · 0.7feedback loop optimization · 0.7control-plane messaging · 0.7compile-time and run-time statistics · 0.7control policy · 0.6simulation · 0.3reconfiguration · 0.3collaborative filtering · 0.3offline scheduling · 0.2joint data and task placement · 0.2
| Year | Publication | Venue | Position |
|---|---|---|---|
| 2026 | GenIE: Simulator-Driven Iterative Data Exploration for Scientific DiscoveryabstractPhysics-based simulators play a critical role in scientific discovery and risk assessment, enabling what-if analyses for events like wildfires and hurricanes. Today, databases treat these simulators as external pre-processing steps. Analysts must manually run a simulation, export the results, and load them into a database before analysis can begin. This linear workflow is inefficient, incurs high latency, and hinders interactive exploration, especially when the analysis itself dictates the need for new or refined simulation data. We envision a new database paradigm, entitled GenIE, that seamlessly integrates multiple simulators into databases to enable dynamic orchestration of simulation workflows. By making the database "simulation-aware," GenIE can dynamically invoke simulators with appropriate parameters based on the user's query and analytical needs. This tight integration allows GenIE to avoid generating data irrelevant to the analysis, reuse previously generated data, and support iterative, incremental analysis where results are progressively refined at interactive speeds. We present our vision for GenIE, designed as an extension to PostgreSQL, and demonstrate its potential benefits through comprehensive use cases: wildfire smoke dispersion analysis using WRF-SFIRE and HYSPLIT, and hurricane hazard assessment integrating wind, surge, and flood models. Our preliminary experiments show how GenIE can transform these slow, static analyses into interactive explorations by intelligently managing the trade-off between simulation accuracy and runtime across multiple integrated simulators. We conclude by highlighting the challenges and opportunities ahead in realizing the full vision of GenIE as a cornerstone for next-generation scientific data analysis. Ashwin Gerard Colaco, Martin Boissier 0001, Sriram Rao, Shubharoop Ghosh, Sharad Mehrotra, Tilmann Rabl |
ICDE | 3 |
| 2019 | Hydra: a federated resource manager for data-center scale analytics
Carlo Curino, Subru Krishnan, Konstantinos Karanasos, Sriram Rao, Giovanni Matteo Fumarola, Botong Huang, Kishore Chaliparambil, Arun Suresh, Young Chen, Solom Heddaya, Roni Burd, Sarvesh Sakalanaga, Chris Douglas, Bill Ramsey, Raghu Ramakrishnan 0001 |
NSDI | 4 |
| 2019 | Efficient inter-datacenter bulk transfers with mixed completion time objectives
Mohammad Noormohammadpour, Srikanth Kandula, Cauligi S. Raghavendra, Sriram Rao |
Comput. Networks | 4 |
| 2018 | Netco: Cache and I/O Management for Analytics over Disaggregated StoresabstractWe consider a common setting where storage is disaggregated from the compute in data-parallel systems. Colocating caching tiers with the compute machines can reduce load on the interconnect but doing so leads to new resource management challenges. We design a system Netco, which prefetches data into the cache (based on workload predictability), and appropriately divides the cache space and network bandwidth between the prefetches and serving ongoing jobs. Netco makes various decisions (what content to cache, when to cache and how to apportion bandwidth) to support end-to-end optimization goals such as maximizing the number of jobs that meet their service-level objectives (e.g., deadlines). Our implementation of these ideas is available within the open-source Apache HDFS project. Experiments on a public cloud, with production-trace inspired workloads, show that Netco uses up to 5x less remote I/O compared to existing techniques and increases the number of jobs that meet their deadlines up to 80%. Virajith Jalaparti, Chris Douglas, Mainak Ghosh, Ashvin Agrawal, Avrilia Floratou, Srikanth Kandula, Ishai Menache, Joseph Naor, Sriram Rao |
SoCC | 9 |
| 2018 | Medea: scheduling of long running applications in shared production clustersabstractThe rise in popularity of machine learning, streaming, and latency-sensitive online applications in shared production clusters has raised new challenges for cluster schedulers. To optimize their performance and resilience, these applications require precise control of their placements, by means of complex constraints, e.g., to collocate or separate their long-running containers across groups of nodes. In the presence of these applications, the cluster scheduler must attain global optimization objectives, such as maximizing the number of deployed applications or minimizing the violated constraints and the resource fragmentation, but without affecting the scheduling latency of short-running containers. Panagiotis Garefalakis, Konstantinos Karanasos, Peter R. Pietzuch, Arun Suresh, Sriram Rao |
EuroSys | 5 |
| 2018 | QuickCast: Fast and Efficient Inter-Datacenter Transfers Using Forwarding Tree CohortsabstractSeveral organizations have built multiple datacenters connected via dedicated wide area networks over which large inter-datacenter transfers take place. Since many such transfers move the same data from one source to multiple destinations, using multicast forwarding trees can reduce bandwidth needs and improve completion times. However, using a single forwarding tree per transfer can lead to poor performance as the slowest receiver dictates the completion time for all receivers. Using multiple forwarding trees per transfer alleviates this concern-the average receiver could finish early; however, if done naively, bandwidth usage would also increase and it is apriori unclear how best to partition receivers, how to construct the multiple trees and how to determine the rate and schedule of flows on these trees. This paper presents QuickCast, a first solution to these problems. Using simulations on real-world network topologies, we see that QuickCast can speed up the average receiver's completion time by as much as 10× while only using 1.04× more bandwidth; further, the completion time for all receivers also improves by as much as faster at high loads. Thereby, while some implementation challenges remain, we advocate using a cohort of forwarding trees. Mohammad Noormohammadpour, Cauligi S. Raghavendra, Srikanth Kandula, Sriram Rao |
INFOCOM | 4 |
| 2018 | Computation Reuse in Analytics Job Service at MicrosoftabstractAnalytics-as-a-service, or analytics job service, is emerging as a new paradigm for data analytics, be it in a cloud environment or within enterprises. In this setting, users are not required to manage or tune their hardware and software infrastructure, and they pay only for the processing resources consumed per job. However, the shared nature of these job services across several users and teams leads to significant overlaps in partial computations, i.e., parts of the processing are duplicated across multiple jobs, thus generating redundant costs. In this paper, we describe a computation reuse framework, coined CLOUDVIEWS, which we built to address the computation overlap problem in Microsoft's SCOPE job service. We present a detailed analysis from our production workloads to motivate the computation overlap problem and the possible gains from computation reuse. The key aspects of our system are the following: (i) we reuse computations by creating materialized views over recurring workloads, i.e., periodically executing jobs that have the same script templates but process new data each time, (ii) we select the views to materialize using a feedback loop that reconciles the compile-time and run-time statistics and gathers precise measures of the utility and cost of each overlapping computation, and (iii) we create materialized views in an online fashion, without requiring an offline phase to materialize the overlapping computations. Alekh Jindal, Shi Qiao 0001, Hiren Patel, Zhicheng Yin, Jieming Di, Malay Bag, Marc T. Friedman, Yifung Lin, Konstantinos Karanasos, Sriram Rao |
SIGMOD Conference | 10 |
| 2018 | Selecting Subexpressions to Materialize at Datacenter ScaleabstractWe observe significant overlaps in the computations performed by user jobs in modern shared analytics clusters. Naïvely computing the same subexpressions multiple times results in wasting cluster resources and longer execution times. Given that these shared cluster workloads consist of tens of thousands of jobs, identifying overlapping computations across jobs is of great interest to both cluster operators and users. Nevertheless, existing approaches support orders of magnitude smaller workloads or employ heuristics with limited effectiveness. In this paper, we focus on the problem of subexpression selection for large workloads, i.e., selecting common parts of job plans and materializing them to speed-up the evaluation of subsequent jobs. We provide an ILP-based formulation of our problem and map it to a bipartite graph labeling problem. Then, we introduce B ig S ubs , a vertex-centric graph algorithm to iteratively choose in parallel which subexpressions to materialize and which subexpressions to use for evaluating each job. We provide a distributed implementation of our approach using our internal SQL-like execution framework, SCOPE, and assess its effectiveness over production workloads. B ig S ubs supports workloads with tens of thousands of jobs, yielding savings of up to 40% in machine-hours. We are currently integrating our techniques with the SCOPE runtime in our production clusters. Alekh Jindal, Konstantinos Karanasos, Sriram Rao, Hiren Patel |
Proc. VLDB Endow. | 3 |
| 2018 | Chi: A Scalable and Programmable Control Plane for Distributed Stream Processing SystemsabstractStream-processing workloads and modern shared cluster environments exhibit high variability and unpredictability. Combined with the large parameter space and the diverse set of user SLOs, this makes modern streaming systems very challenging to statically configure and tune. To address these issues, in this paper we investigate a novel control-plane design, Chi, which supports continuous monitoring and feedback, and enables dynamic re-configuration. Chi leverages the key insight of embedding control-plane messages in the data-plane channels to achieve a low-latency and flexible control plane for stream-processing systems. Chi introduces a new reactive programming model and design mechanisms to asynchronously execute control policies, thus avoiding global synchronization. We show how this allows us to easily implement a wide spectrum of control policies targeting different use cases observed in production. Large-scale experiments using production workloads from a popular cloud provider demonstrate the flexibility and efficiency of our approach. Luo Mai, Kai Zeng 0002, Rahul Potharaju, Steve Suh, Shivaram Venkataraman, Paolo Costa, Terry Kim, Saravanam Muthukrishnan, Vamsi Kuppa, Sudheer Dhulipalla, Sriram Rao |
Proc. VLDB Endow. | 12 |
| 2018 | Towards a Learning Optimizer for Shared CloudsabstractQuery optimizers are notorious for inaccurate cost estimates, leading to poor performance. The root of the problem lies in inaccurate cardinality estimates, i.e., the size of intermediate (and final) results in a query plan. These estimates also determine the resources consumed in modern shared cloud infrastructures. In this paper, we present C ARD L EARNER , a machine learning based approach to learn cardinality models from previous job executions and use them to predict the cardinalities in future jobs. The key intuition in our approach is that shared cloud workloads are often recurring and overlapping in nature, and so we could learn cardinality models for overlapping subgraph templates. We discuss various learning approaches and show how learning a large number of smaller models results in high accuracy and explainability. We further present an exploration technique to avoid learning bias by considering alternate join orders and learning cardinality models over them. We describe the feedback loop to apply the learned models back to future job executions. Finally, we show a detailed evaluation of our models (up to 5 orders of magnitude less error), query plans (60% applicability), performance (up to 100% faster, 3x fewer resources), and exploration (optimal in few 10s of executions). Chenggang Wu 0001, Alekh Jindal, Saeed Amizadeh, Hiren Patel, Wangchao Le, Shi Qiao 0001, Sriram Rao |
Proc. VLDB Endow. | 7 |
| 2017 | Twitter Heron: Towards Extensible Streaming EnginesabstractTwitter's data centers process billions of events per day the instant the data is generated. To achieve real-time performance, Twitter has developed Heron, a streaming engine that provides unparalleled performance at large scale. Heron has been recently open-sourced and thus is now accessible to various other organizations. In this paper, we discuss the challenges we faced when transforming Heron from a system tailored for Twitter's applications and software stack to a system that efficiently handles applications with diverse characteristics on top of various Big Data platforms. Overcoming these challenges required a careful design of the system using an extensible, modular architecture which provides flexibility to adapt to various environments and applications. Further, we describe the various optimizations that allow us to gain this flexibility without sacrificing performance. Finally, we experimentally show the benefits of Heron's modular architecture. Maosong Fu, Ashvin Agrawal, Avrilia Floratou, Bill Graham, Andrew Jorgensen, Mark Li, Neng Lu, Karthikeyan Ramasamy, Sriram Rao |
ICDE | 9 |
| 2017 | Dhalion: Self-Regulating Stream Processing in HeronabstractIn recent years, there has been an explosion of large-scale real-time analytics needs and a plethora of streaming systems have been developed to support such applications. These systems are able to continue stream processing even when faced with hardware and software failures. However, these systems do not address some crucial challenges facing their operators: the manual, time-consuming and error-prone tasks of tuning various configuration knobs to achieve service level objectives (SLO) as well as the maintenance of SLOs in the face of sudden, unpredictable load variation and hardware or software performance degradation. In this paper, we introduce the notion of self-regulating streaming systems and the key properties that they must satisfy. We then present the design and evaluation of Dhalion, a system that provides self-regulation capabilities to underlying streaming systems. We describe our implementation of the Dhalion framework on top of Twitter Heron, as well as a number of policies that automatically reconfigure Heron topologies to meet throughput SLOs, scaling resource consumption up and down as needed. We experimentally evaluate our Dhalion policies in a cloud environment and demonstrate their effectiveness. We are in the process of open-sourcing our Dhalion policies as part of the Heron project. Avrilia Floratou, Ashvin Agrawal, Bill Graham, Sriram Rao, Karthikeyan Ramasamy |
Proc. VLDB Endow. | 4 |
| 2017 | Apache REEF: Retainable Evaluator Execution FrameworkabstractResource Managers like YARN and Mesos have emerged as a critical layer in the cloud computing system stack, but the developer abstractions for leasing cluster resources and instantiating application logic are very low level. This flexibility comes at a high cost in terms of developer effort, as each application must repeatedly tackle the same challenges (e.g., fault tolerance, task scheduling and coordination) and reimplement common mechanisms (e.g., caching, bulk-data transfers). This article presents REEF, a development framework that provides a control plane for scheduling and coordinating task-level (data-plane) work on cluster resources obtained from a Resource Manager. REEF provides mechanisms that facilitate resource reuse for data caching and state management abstractions that greatly ease the development of elastic data processing pipelines on cloud platforms that support a Resource Manager service. We illustrate the power of REEF by showing applications built atop: a distributed shell application, a machine-learning framework, a distributed in-memory caching system, and a port of the CORFU system. REEF is currently an Apache top-level project that has attracted contributors from several institutions and it is being used to develop several commercial offerings such as the Azure Stream Analytics service. Byung-Gon Chun, Tyson Condie, Yingda Chen, Carlo Curino, Chris Douglas, Matteo Interlandi, Beomyeol Jeon, Joo Seong Jeong, Gyewon Lee, Yunseong Lee, Tony Majestro, Dahlia Malkhi, Sergiy Matusevych, Brandon Myers, Mariia Mykhailova, Shravan M. Narayanamurthy, Joseph Noor, Raghu Ramakrishnan 0001, Sriram Rao, Russell Sears, Beysim Sezgin, Taegeon Um, Julia Wang, Markus Weimer, Youngseok Yang |
ACM Trans. Comput. Syst. | 21 |
| 2016 | Efficient queue management for cluster schedulingabstractJob scheduling in Big Data clusters is crucial both for cluster operators' return on investment and for overall user experience. In this context, we observe several anomalies in how modern cluster schedulers manage queues, and argue that maintaining queues of tasks at worker nodes has significant benefits. On one hand, centralized approaches do not use worker-side queues. Given the inherent feedback delays that these systems incur, they achieve suboptimal cluster utilization, particularly for workloads dominated by short tasks. On the other hand, distributed schedulers typically do employ worker-side queuing, and achieve higher cluster utilization. However, they fail to place tasks at the best possible machine, since they lack cluster-wide information, leading to worse job completion time, especially for heterogeneous workloads. To the best of our knowledge, this is the first work to provide principled solutions to the above problems by introducing queue management techniques, such as appropriate queue sizing, prioritization of task execution via queue reordering, starvation freedom, and careful placement of tasks to queues. We instantiate our techniques by extending both a centralized (YARN) and a distributed (Mercury) scheduler, and evaluate their performance on a wide variety of synthetic and production workloads derived from Microsoft clusters. Our centralized implementation, Yaq-c, achieves 1.7x improvement on median job completion time compared to YARN, and our distributed one, Yaq-d, achieves 9.3x improvement over an implementation of Sparrow's batch sampling on Mercury. Jeff Rasley, Konstantinos Karanasos, Srikanth Kandula, Rodrigo Fonseca, Milan Vojnovic, Sriram Rao |
EuroSys | 6 |
| 2016 | DCRoute: Speeding up Inter-Datacenter Traffic Allocation while Guaranteeing DeadlinesabstractDatacenters provide the infrastructure for cloud computing services used by millions of users everyday. Many such services are distributed over multiple datacenters at geographically distant locations possibly in different continents. These datacenters are then connected through high speed WAN links over private or public networks. To perform data backups or data synchronization operations, many transfers take place over these networks that have to be completed before a deadline in order to provide necessary service guarantees to end users. Upon arrival of a transfer request, we would like the system to be able to decide whether such a request can be guaranteed successful delivery. If yes, it should provide us with transmission schedule in the shortest time possible. In addition, we would like to avoid packet reordering at the destination as it affects TCP performance. Previous work in this area either cannot guarantee that admitted transfers actually finish before the specified deadlines or use techniques that can result in packet reordering. In this paper, we propose DCRoute, a fast and efficient routing and traffic allocation technique that guarantees transfer completion before deadlines for admitted requests. It assigns each transfer a single path to avoid packet reordering. Through simulations, we show that DCRoute is at least 200 times faster than other traffic allocation techniques based on linear programming (LP) while admitting almost the same amount of traffic to the system. Mohammad Noormohammadpour, Cauligi S. Raghavendra, Sriram Rao |
HiPC | 3 |
| 2016 | GRAPHENE: Packing and Dependency-Aware Scheduling for Data-Parallel Clusters
Robert Grandl, Srikanth Kandula, Sriram Rao, Aditya Akella, Janardhan Kulkarni |
OSDI | 3 |
| 2016 | Morpheus: Towards Automated SLOs for Enterprise Clusters
Sangeetha Abdu Jyothi, Carlo Curino, Ishai Menache, Shravan M. Narayanamurthy, Alexey Tumanov, Jonathan Yaniv, Ruslan Mavlyutov, Íñigo Goiri, Subru Krishnan, Janardhan Kulkarni, Sriram Rao |
OSDI | 11 |
| 2015 | Network-Aware Scheduling for Data-Parallel Jobs: Plan When You CanabstractTo reduce the impact of network congestion on big data jobs, cluster management frameworks use various heuristics to schedule compute tasks and/or network flows. Most of these schedulers consider the job input data fixed and greedily schedule the tasks and flows that are ready to run. However, a large fraction of production jobs are recurring with predictable characteristics, which allows us to plan ahead for them. Coordinating the placement of data and tasks of these jobs allows for significantly improving their network locality and freeing up bandwidth, which can be used by other jobs running on the cluster. With this intuition, we develop Corral, a scheduling framework that uses characteristics of future workloads to determine an offline schedule which (i) jointly places data and compute to achieve better data locality, and (ii) isolates jobs both spatially (by scheduling them in different parts of the cluster) and temporally, improving their performance. We implement Corral on Apache Yarn, and evaluate it on a 210 machine cluster using production workloads. Compared to Yarn's capacity scheduler, Corral reduces the makespan of these workloads up to 33% and the median completion time up to 56%, with 20-90% reduction in data transferred across racks. Virajith Jalaparti, Peter Bodík, Ishai Menache, Sriram Rao, Konstantin Makarychev, Matthew Caesar 0001 |
SIGCOMM | 4 |
| 2015 | REEF: Retainable Evaluator Execution FrameworkabstractResource Managers like Apache YARN have emerged as a critical layer in the cloud computing system stack, but the developer abstractions for leasing cluster resources and instantiating application logic are very low-level. This flexibility comes at a high cost in terms of developer effort, as each application must repeatedly tackle the same challenges (e.g., fault-tolerance, task scheduling and coordination) and re-implement common mechanisms (e.g., caching, bulk-data transfers). This paper presents REEF, a development framework that provides a control-plane for scheduling and coordinating task-level (data-plane) work on cluster resources obtained from a Resource Manager. REEF provides mechanisms that facilitate resource re-use for data caching, and state management abstractions that greatly ease the development of elastic data processing work-flows on cloud platforms that support a Resource Manager service. REEF is being used to develop several commercial offerings such as the Azure Stream Analytics service. Furthermore, we demonstrate REEF development of a distributed shell application, a machine learning algorithm, and a port of the CORFU [4] system. REEF is also currently an Apache Incubator project that has attracted contributors from several instititutions. Markus Weimer, Yingda Chen, Byung-Gon Chun, Tyson Condie, Carlo Curino, Chris Douglas, Yunseong Lee, Tony Majestro, Dahlia Malkhi, Sergiy Matusevych, Brandon Myers, Shravan M. Narayanamurthy, Raghu Ramakrishnan 0001, Sriram Rao, Russell Sears, Beysim Sezgin, Julia Wang |
SIGMOD Conference | 14 |
| 2015 | Mercury: Hybrid Centralized and Distributed Scheduling in Large Shared Clusters
Konstantinos Karanasos, Sriram Rao, Carlo Curino, Chris Douglas, Kishore Chaliparambil, Giovanni Matteo Fumarola, Solom Heddaya, Raghu Ramakrishnan 0001, Sarvesh Sakalanaga |
USENIX ATC | 2 |
| 2014 | Reservation-based Scheduling: If You're Late Don't Blame Us!abstractThe continuous shift towards data-driven approaches to business, and a growing attention to improving return on investments (ROI) for cluster infrastructures is generating new challenges for big-data frameworks. Systems originally designed for big batch jobs now handle an increasingly complex mix of computations. Moreover, they are expected to guarantee stringent SLAs for production jobs and minimize latency for best-effort jobs. Carlo Curino, Djellel Eddine Difallah, Chris Douglas, Subru Krishnan, Raghu Ramakrishnan 0001, Sriram Rao |
SoCC | 6 |
| 2014 | Multi-resource packing for cluster schedulersabstractTasks in modern data parallel clusters have highly diverse resource requirements, along CPU, memory, disk and network. Any of these resources may become bottlenecks and hence, the likelihood of wasting resources due to fragmentation is now larger. Today's schedulers do not explicitly reduce fragmentation. Worse, since they only allocate cores and memory, the resources that they ignore (disk and network) can be over-allocated leading to interference, failures and hogging of cores or memory that could have been used by other tasks. We present Tetris, a cluster scheduler that packs, i.e., matches multi-resource task requirements with resource availabilities of machines so as to increase cluster efficiency (makespan). Further, Tetris uses an analog of shortest-running-time-first to trade-off cluster efficiency for speeding up individual jobs. Tetris' packing heuristics seamlessly work alongside a large class of fairness policies. Trace-driven simulations and deployment of our prototype on a 250 node cluster shows median gains of 30% in job completion time while achieving nearly perfect fairness. Robert Grandl, Ganesh Ananthanarayanan, Srikanth Kandula, Sriram Rao, Aditya Akella |
SIGCOMM | 4 |
| 2013 | Tango: distributed data structures over a shared logabstractDistributed systems are easier to build than ever with the emergence of new, data-centric abstractions for storing and computing over massive datasets. However, similar abstractions do not exist for storing and accessing meta-data. To fill this gap, Tango provides developers with the abstraction of a replicated, in-memory data structure (such as a map or a tree) backed by a shared log. Tango objects are easy to build and use, replicating state via simple append and read operations on the shared log instead of complex distributed protocols; in the process, they obtain properties such as linearizability, persistence and high availability from the shared log. Tango also leverages the shared log to enable fast transactions across different objects, allowing applications to partition state across machines and scale to the limits of the underlying log without sacrificing consistency. Mahesh Balakrishnan 0001, Dahlia Malkhi, Ted Wobber, Ming Wu 0007, Vijayan Prabhakaran, Michael Wei, John D. Davis, Sriram Rao, Tao Zou 0002, Aviad Zuck |
SOSP | 8 |
| 2013 | A The Quantcast File SystemabstractThe Quantcast File System (QFS) is an efficient alternative to the Hadoop Distributed File System (HDFS). QFS is written in C++, is plugin compatible with Hadoop MapReduce, and offers several efficiency improvements relative to HDFS: 50% disk space savings through erasure coding instead of replication, a resulting doubling of write throughput, a faster name node, support for faster sorting and logging through a concurrent append feature, a native command line client much faster than hadoop fs, and global feedback-directed I/O device management. As QFS works out of the box with Hadoop, migrating data from HDFS to QFS involves simply executing hadoop distcp. QFS is being developed fully open source and is available under an Apache license from https://github.com/quantcast/qfs. Multi-petabyte QFS instances have been in heavy production use since 2011. Michael Ovsiannikov, Silvius Rus, Damian Reeves, Paul Sutter, Sriram Rao, Jim Kelly |
Proc. VLDB Endow. | 5 |
| 2012 | True elasticity in multi-tenant data-intensive compute clustersabstractData-intensive computing (DISC) frameworks scale by partitioning a job across a set of fault-tolerant tasks, then diffusing those tasks across large clusters. Multi-tenanted clusters must accommodate service-level objectives (SLO) in their resource model, often expressed as a maximum latency for allocating the desired set of resources to every job. When jobs are partitioned into tasks statically, a cluster cannot meet its SLOs while maintaining both high utilization and efficiency. Ideally, we want to give resources to jobs when they are free but would expect to reclaim them instantaneously when new jobs arrive, without losing work. DISC frameworks do not support such elasticity because interrupting running tasks incurs high overheads. Amoeba enables lightweight elasticity in DISC frameworks by identifying points at which running tasks of over-provisioned jobs can be safely exited, committing their outputs, and spawning new tasks for the remaining work. Effectively, tasks of DISC jobs are now sized dynamically in response to global resource scarcity or abundance. Simulation and deployment of our prototype shows that Amoeba speeds up jobs by 32% without compromising utilization or efficiency. Ganesh Ananthanarayanan, Chris Douglas, Raghu Ramakrishnan 0001, Sriram Rao, Ion Stoica |
SoCC | 4 |
| 2012 | Sailfish: a framework for large scale data processingabstractIn this paper, we present Sailfish, a new Map-Reduce framework for large scale data processing. The Sailfish design is centered around aggregating intermediate data, specifically data produced by map tasks and consumed later by reduce tasks, to improve performance by batching disk I/O. We introduce an abstraction called I-files for supporting data aggregation, and describe how we implemented it as an extension of the distributed filesystem, to efficiently batch data written by multiple writers and read by multiple readers. Sailfish adapts the Map-Reduce layer in Hadoop to use I-files for transporting data from map tasks to reduce tasks. We present experimental results demonstrating that Sailfish improves performance of standard Hadoop; in particular, we show 20% to 5 times faster performance on a representative mix of real jobs and datasets at Yahoo!. We also demonstrate that the Sailfish design enables auto-tuning functionality that handles changes in data volume and skewed distributions effectively, thereby addressing an important practical drawback of Hadoop, which in contrast relies on programmers to configure system parameters appropriately for each job, for each input dataset. Our Sailfish implementation and the other software components developed as part of this paper has been released as open source. Sriram Rao, Raghu Ramakrishnan 0001, Adam Silberstein, Michael Ovsiannikov, Damian Reeves |
SoCC | 1 |
| 2012 | Walnut: a unified cloud object storeabstractWalnut is an object-store being developed at Yahoo! with the goal of serving as a common low-level storage layer for a variety of cloud data management systems including Hadoop (a MapReduce system), MObStor (a multimedia serving system), and PNUTS (an extended key-value serving system). Thus, a key performance challenge is to meet the latency and throughput requirements of the wide range of workloads commonly observed across these diverse systems. The motivation for Walnut is to leverage a carefully optimized low-level storage system, with support for elasticity and high-availability, across all of Yahoo!'s data clouds. This would enable sharing of hardware resources across hitherto siloed clouds of different types, offering greater potential for intelligent load balancing and efficient elastic operation, and simplify the operational tasks related to data storage. Jianjun Chen 0001, Chris Douglas, Michi Mutsuzaki, Patrick Quaid, Raghu Ramakrishnan 0001, Sriram Rao, Russell Sears |
SIGMOD Conference | 6 |
| 2003 | Design considerations for the symphony integrated multimedia file system
Prashant J. Shenoy, Pawan Goyal 0001, Sriram Rao, Harrick M. Vin |
Multim. Syst. | 3 |
| 2000 | The Cost of Recovery in Message Logging ProtocolsabstractPast research in message logging has focused on studying the relative overhead imposed by pessimistic, optimistic and causal protocols during failure-free executions. In this paper, we give the first experimental evaluation of the performance of these protocols during recovery. Our results suggest that applications face a complex tradeoff when choosing a message logging protocol for fault tolerance. On the one hand, optimistic protocols can provide fast failure-free execution and good performance during recovery, but are complex to implement and can create orphan processes. On the other hand, orphan-free protocols either risk being slow during recovery (e.g. sender-based pessimistic and causal protocols) or incur a substantial overhead during failure-free execution (e.g. receiver-based pessimistic protocols). To address this tradeoff, we propose hybrid logging protocols, which are a new class of orphan-free protocols. We show that hybrid protocols perform within 2% of causal logging during failure-free execution and within 2% of receiver-based logging during recovery. Sriram Rao, Lorenzo Alvisi, Harrick M. Vin |
IEEE Trans. Knowl. Data Eng. | 1 |
| 1998 | Low-Overhead Protocols for Fault-Tolerant File SharingabstractWe quantify the adverse effect of file sharing on the performance of reliable distributed applications. We demonstrate that file sharing incurs significant overhead, which is likely to triple over the next five years. We present a novel approach that eliminates this overhead. Our approach: tracks causal dependencies resulting from file sharing using determinants; efficiently replicates the determinants in the volatile memory of agents to ensure their availability during recovery; and reproduces during recovery the interactions with the file server as well as the file data lost in a failure. Our approach allows agents to exchange files directly without first saving the files on disks at the server. As a consequence, the costs of supporting file sharing and message passing in a reliable distributed application become virtually identical. The result is a simple, uniform approach, which can provide low-overhead fault tolerance to applications in which communication is performed through message passing, file sharing, or a combination of the two. Lorenzo Alvisi, Sriram Rao, Harrick M. Vin |
ICDCS | 2 |
| 1998 | The Cost of Recovery in Message Logging ProtocolsabstractPast research in message logging has focused on studying the relative overhead imposed by pessimistic, optimistic, and causal protocols during failure-free executions. We give the first experimental evaluation of the performance of these protocols during recovery. We discover that, if a single failure is to be tolerated, pessimistic and causal protocols perform best, because they avoid rollbacks of correct processes. For multiple failures, however, the dominant factor in determining performance becomes where the recovery information is logged (i.e. at the sender, at the receiver, or replicated at a subset of the processes in the system) rather than when this information is logged (i.e. if logging is synchronous or asynchronous). Sriram Rao, Lorenzo Alvisi, Harrick M. Vin |
SRDS | 1 |