Peter Bodík

dblp:b/PeterBodik · DBLP profile ↗
← Back
18ranked-venue papers
4as first author
1since 2021 · last 2024
—ORCID · none

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

Computer networks · 10 · 1 first-authorSystems, architecture and hardware · 6 · 3 first-authorSoftware engineering, systems software and programming languages · 1Databases, data management, data science and information retrieval · 1Human-computer interaction and ubiquitous computing · 1 · 1 since 2021

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
9 papers
Cloud and datacenter computing · 72% Parallel and multicore computing · 14% Distributed systems · 10%
Computer graphics and multimedia
4 papers
Multimedia analysis and retrieval · 77% Multimedia systems and quality of experience · 23%
Computer networks
5 papers
Network management and operations · 24% Network optimization and economics · 24% Content delivery and video streaming · 15%
Databases, data mining, and information retrieval
2 papers
Query processing and optimization · 60% Distributed and cloud data management · 40%

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

TopicWeightPapersLastEvidence papers
Cloud and datacenter computing
cluster resource management and scheduling
0.632016
2DFQ: Two-Dimensional Fair Queuing for Multi-Tenant Cloud Services · SIGCOMM 2016
Retro: Targeted Resource Management in Multi-tenant Distributed Systems · NSDI 2015
Jockey: guaranteed job latency in data parallel clusters · EuroSys 2012
Query processing and optimization › complex data query processing
video query processing
0.312018
Focus: Querying Large Video Datasets with Low Latency and Low Cost · OSDI 2018
Multimedia analysis and retrieval
video content analysis
0.312018
Chameleon: scalable adaptation of video analytics · SIGCOMM 2018
Multimedia analysis and retrieval
video retrieval
0.312018
Focus: Querying Large Video Datasets with Low Latency and Low Cost · OSDI 2018
Multimedia analysis and retrieval › event detection
video event detection
0.312017
Demo: Live Video Stream Triggers · MobiSys 2017
Cloud and datacenter computing › cluster resource management and scheduling
cluster resource management
0.322015
Network-Aware Scheduling for Data-Parallel Jobs: Plan When You Can · SIGCOMM 2015
Speeding up distributed request-response workflows · SIGCOMM 2013
Parallel and multicore computing › parallel scheduling
data-parallel job scheduling
0.322015
Network-Aware Scheduling for Data-Parallel Jobs: Plan When You Can · SIGCOMM 2015
Jockey: guaranteed job latency in data parallel clusters · EuroSys 2012
Cloud and datacenter computing › job scheduling
fair scheduling
0.212016
2DFQ: Two-Dimensional Fair Queuing for Multi-Tenant Cloud Services · SIGCOMM 2016
Cloud and datacenter computing › quality of service
tail latency
0.212016
2DFQ: Two-Dimensional Fair Queuing for Multi-Tenant Cloud Services · SIGCOMM 2016
Distributed and cloud data management › distributed analytics
geo-distributed analytics
0.212015
Low Latency Geo-distributed Data Analytics · SIGCOMM 2015
Parallel and multicore computing › locality optimization
data locality optimization
0.212015
Network-Aware Scheduling for Data-Parallel Jobs: Plan When You Can · SIGCOMM 2015
Cloud and datacenter computing › resource management
multi-tenant resource management
0.212015
Retro: Targeted Resource Management in Multi-tenant Distributed Systems · NSDI 2015
Cloud and datacenter computing › job scheduling
network-aware scheduling
0.212015
Network-Aware Scheduling for Data-Parallel Jobs: Plan When You Can · SIGCOMM 2015
Cloud and datacenter computing
resource management
0.212015
Retro: Targeted Resource Management in Multi-tenant Distributed Systems · NSDI 2015
Distributed systems › distributed data processing
straggler mitigation
0.212013
Speeding up distributed request-response workflows · SIGCOMM 2013
Network optimization and economics › resource allocation
bandwidth optimization
0.112012
Surviving failures in bandwidth-constrained datacenters · SIGCOMM 2012
Network management and operations › network robustness
fault tolerance
0.112012
Surviving failures in bandwidth-constrained datacenters · SIGCOMM 2012
Cloud and datacenter computing › resource allocation
dynamic resource allocation
0.112012
Jockey: guaranteed job latency in data parallel clusters · EuroSys 2012
Storage systems
distributed storage
0.112011
The SCADS Director: Scaling a Distributed Storage System Under Stringent Performance Requirements · FAST 2011
Distributed systems
performance management
0.112011
The SCADS Director: Scaling a Distributed Storage System Under Stringent Performance Requirements · FAST 2011
Cloud and datacenter computing
datacenter operations
0.112010
Fingerprinting the datacenter: automated classification of performance crises · EuroSys 2010
Edge and fog computing › video analytics
edge video analytics
0.112017
Live Video Analytics at Scale with Approximation and Delay-Tolerance · NSDI 2017
Content delivery and video streaming
live streaming
0.112017
Demo: Live Video Stream Triggers · MobiSys 2017
Distributed systems › distributed data processing
wide-area data analytics
0.112015
Low Latency Geo-distributed Data Analytics · SIGCOMM 2015
Cloud and datacenter computing
resource allocation
0.012013
Speeding up distributed request-response workflows · SIGCOMM 2013
Cloud and datacenter computing › resource management
datacenter resource management
0.012012
Surviving failures in bandwidth-constrained datacenters · SIGCOMM 2012
Performance modeling and evaluation
workload characterization
0.012010
Fingerprinting the datacenter: automated classification of performance crises · EuroSys 2010
Internet of things and sensor networks › sensor data management
sensor data collection
0.012004
Distributed regression: an efficient framework for modeling sensor network data · IPSN 2004

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

deep convolutional neural network · 0.7controller design · 0.7standing query matching · 0.6approximation · 0.6optimization · 0.4online heuristic · 0.4offline scheduling · 0.4joint data and task placement · 0.4workload characterization · 0.2fair queuing · 0.2reissue policies · 0.2partial result return · 0.2optimization over DAG · 0.2control policy · 0.1kernel linear regression · 0.0distributed optimization · 0.0
YearPublicationVenuePosition
2024 Evaluating how interactive visualizations can assist in finding samples where and how computer vision models make mistakes
abstract
Creating Computer Vision (CV) models remains a complex practice, despite their ubiquity. Access to data, the requirement for ML expertise, and model opacity are just a few points of complexity that limit the ability of end-users to build, inspect, and improve these models. Interactive ML perspectives have helped address some of these issues by considering a teacher in the loop where planning, teaching, and evaluating tasks take place. We present and evaluate two interactive visualizations in the context of Sprite, a system for creating CV classification and detection models for images originating from videos. We study how these visualizations help Sprite’s users identify (evaluate) and select (plan) images where a model is struggling and can lead to improved performance, compared to a baseline condition where users used a query language. We found that users who had used the visualizations found more images across a wider set of potential types of model errors.
Hayeong Song, Gonzalo A. Ramos, Peter Bodík
PacificVis3
2018 Focus: Querying Large Video Datasets with Low Latency and Low Cost
Kevin Hsieh, Ganesh Ananthanarayanan, Peter Bodík, Shivaram Venkataraman, Paramvir Bahl, Matthai Philipose, Phillip B. Gibbons, Onur Mutlu
OSDI3
2018 Chameleon: scalable adaptation of video analytics
abstract
Applying deep convolutional neural networks (NN) to video data at scale poses a substantial systems challenge, as improving inference accuracy often requires a prohibitive cost in computational resources. While it is promising to balance resource and accuracy by selecting a suitable NN configuration (e.g., the resolution and frame rate of the input video), one must also address the significant dynamics of the NN configuration's impact on video analytics accuracy. We present Chameleon, a controller that dynamically picks the best configurations for existing NN-based video analytics pipelines. The key challenge in Chameleon is that in theory, adapting configurations frequently can reduce resource consumption with little degradation in accuracy, but searching a large space of configurations periodically incurs an overwhelming resource overhead that negates the gains of adaptation. The insight behind Chameleon is that the underlying characteristics (e.g., the velocity and sizes of objects) that affect the best configuration have enough temporal and spatial correlation to allow the search cost to be amortized over time and across multiple video feeds. For example, using the video feeds of five traffic cameras, we demonstrate that compared to a baseline that picks a single optimal configuration offline, Chameleon can achieve 20-50% higher accuracy with the same amount of resources, or achieve the same accuracy with only 30--50% of the resources (a 2-3X speedup).
Junchen Jiang, Ganesh Ananthanarayanan, Peter Bodík, Siddhartha Sen 0001, Ion Stoica
SIGCOMM3
2017 Distributed resource management across process boundaries
abstract
Multi-tenant distributed systems composed of small services, such as Service-oriented Architectures (SOAs) and Micro-services, raise new challenges in attaining high performance and efficient resource utilization. In these systems, a request execution spans tens to thousands of processes, and the execution paths and resource demands on different services are generally not known when a request first enters the system. In this paper, we highlight the fundamental challenges of regulating load and scheduling in SOAs while meeting end-to-end performance objectives on metrics of concern to both tenants and operators. We design Wisp, a framework for building SOAs that transparently adapts rate limiters and request schedulers system-wide according to operator policies to satisfy end-to-end goals while responding to changing system conditions. In evaluations against production as well as synthetic workloads, Wisp successfully enforces a range of end-to-end performance objectives, such as reducing average latencies, meeting deadlines, providing fairness and isolation, and avoiding system overload.
Lalith Suresh 0001, Peter Bodík, Ishai Menache, Marco Canini, Florin Ciucu
SoCC2
2017 Demo: Live Video Stream Triggers
abstract
Live streaming is an increasingly popular way to broadcast videos ranging from formal news channels to kitten cams to home security camera feeds. Live streaming marries the rich detail of video with the timeliness of live transmission and the ease of use of consumer cameras, thus promising to vastly increase the amount of detailed, up-to-the minute information available about the real world. The volume of potentially interesting footage brings up the question of how end-users can avoid being glued to one (or worse, many) streams of videos waiting for events of interest. In this demo, we present Lookout, a system that allows users to register standing queries, called triggers over live video streams. Lookout then notifies the user when events of interest to them occur in their streams of interest. For example, a user could point to a cat cam and write a trigger that sends a notification when the cat wakes up and starts moving. Users can also write triggers to look for certain news being covered in a live new channel, a gamer moving to a certain level in a Twitch stream, a stranger showing up in a outdoor surveillance camera, etc.
Lenin Ravindranath, Matthai Philipose, Peter Bodík, Paramvir Bahl
MobiSys3
2017 Live Video Analytics at Scale with Approximation and Delay-Tolerance
Ganesh Ananthanarayanan, Peter Bodík, Matthai Philipose, Paramvir Bahl, Michael J. Freedman
NSDI3
2016 2DFQ: Two-Dimensional Fair Queuing for Multi-Tenant Cloud Services
abstract
In many important cloud services, different tenants execute their requests in the thread pool of the same process, requiring fair sharing of resources. However, using fair queue schedulers to provide fairness in this context is difficult because of high execution concurrency, and because request costs are unknown and have high variance. Using fair schedulers like WFQ and WF²Q in such settings leads to bursty schedules, where large requests block small ones for long periods of time. In this paper, we propose Two-Dimensional Fair Queueing (2DFQ), which spreads requests of different costs across di erent threads and minimizes the impact of tenants with unpredictable requests. In evaluation on production workloads from Azure Storage, a large-scale cloud system at Microsoft, we show that 2DFQ reduces the burstiness of service by 1-2 orders of magnitude. On workloads where many large requests compete with small ones, 2DFQ improves 99th percentile latencies by up to 2 orders of magnitude.
Jonathan Mace, Peter Bodík, Madan Musuvathi, Rodrigo Fonseca, Krishnan Varadarajan
SIGCOMM2
2015 Retro: Targeted Resource Management in Multi-tenant Distributed Systems
Jonathan Mace, Peter Bodík, Rodrigo Fonseca, Madan Musuvathi
NSDI2
2015 Network-Aware Scheduling for Data-Parallel Jobs: Plan When You Can
abstract
To 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
SIGCOMM2
2015 Low Latency Geo-distributed Data Analytics
abstract
Low latency analytics on geographically distributed datasets (across datacenters, edge clusters) is an upcoming and increasingly important challenge. The dominant approach of aggregating all the data to a single datacenter significantly inflates the timeliness of analytics. At the same time, running queries over geo-distributed inputs using the current intra-DC analytics frameworks also leads to high query response times because these frameworks cannot cope with the relatively low and variable capacity of WAN links. We present Iridium, a system for low latency geo-distributed analytics. Iridium achieves low query response times by optimizing placement of both data and tasks of the queries. The joint data and task placement optimization, however, is intractable. Therefore, Iridium uses an online heuristic to redistribute datasets among the sites prior to queries' arrivals, and places the tasks to reduce network bottlenecks during the query's execution. Finally, it also contains a knob to budget WAN usage. Evaluation across eight worldwide EC2 regions using production queries show that Iridium speeds up queries by 3× -- 19× and lowers WAN usage by 15% -- 64% compared to existing baselines.
Qifan Pu, Ganesh Ananthanarayanan, Peter Bodík, Srikanth Kandula, Aditya Akella, Paramvir Bahl, Ion Stoica
SIGCOMM3
2014 Brief announcement: deadline-aware scheduling of big-data processing jobs
abstract
This paper presents a novel algorithm for scheduling big data jobs on large compute clusters. In our model, each job is represented by a DAG consisting of several stages linked by precedence constraints. The resource allocation per stage is malleable, in the sense that the processing time of a stage depends on the resources allocated to it (the dependency can be arbitrary in general).The goal of the scheduler is to maximize the total value of completed jobs, where the value for each job depends on its completion time. We design an algorithm for the problem which guarantees an expected constant approximation factor when the cluster capacity is sufficiently high. To the best of our knowledge, this is the first constant-factor approximation algorithm for the problem. The algorithm is based on formulating the problem as a linear program and then rounding an optimal (fractional) solution into a feasible (integral) schedule using randomized rounding.
Peter Bodík, Ishai Menache, Joseph Naor, Jonathan Yaniv
SPAA1
2013 Speeding up distributed request-response workflows
abstract
We found that interactive services at Bing have highly variable datacenter-side processing latencies because their processing consists of many sequential stages, parallelization across 10s-1000s of servers and aggregation of responses across the network. To improve the tail latency of such services, we use a few building blocks: reissuing laggards elsewhere in the cluster, new policies to return incomplete results and speeding up laggards by giving them more resources. Combining these building blocks to reduce the overall latency is non-trivial because for the same amount of resource (e.g., number of reissues), different stages improve their latency by different amounts. We present Kwiken, a framework that takes an end-to-end view of latency improvements and costs. It decomposes the problem of minimizing latency over a general processing DAG into a manageable optimization over individual stages. Through simulations with production traces, we show sizable gains; the 99th percentile of latency improves by over 50% when just 0.1% of the responses are allowed to have partial results and by over 40% for 25% of the services when just 5% extra resources are used for reissues.
Virajith Jalaparti, Peter Bodík, Srikanth Kandula, Ishai Menache, Mikhail Rybalkin, Chenyu Yan
SIGCOMM2
2012 Jockey: guaranteed job latency in data parallel clusters
abstract
Data processing frameworks such as MapReduce [8] and Dryad [11] are used today in business environments where customers expect guaranteed performance. To date, however, these systems are not capable of providing guarantees on job latency because scheduling policies are based on fair-sharing, and operators seek high cluster use through statistical multiplexing and over-subscription. With Jockey, we provide latency SLOs for data parallel jobs written in SCOPE. Jockey precomputes statistics using a simulator that captures the job's complex internal dependencies, accurately and efficiently predicting the remaining run time at different resource allocations and in different stages of the job. Our control policy monitors a job's performance, and dynamically adjusts resource allocation in the shared cluster in order to maximize the job's economic utility while minimizing its impact on the rest of the cluster. In our experiments in Microsoft's production Cosmos clusters, Jockey meets the specified job latency SLOs and responds to changes in cluster conditions.
Andrew D. Ferguson, Peter Bodík, Srikanth Kandula, Eric Boutin, Rodrigo Fonseca
EuroSys2
2012 Surviving failures in bandwidth-constrained datacenters
abstract
Datacenter networks have been designed to tolerate failures of network equipment and provide sufficient bandwidth. In practice, however, failures and maintenance of networking and power equipment often make tens to thousands of servers unavailable, and network congestion can increase service latency. Unfortunately, there exists an inherent tradeoff between achieving high fault tolerance and reducing bandwidth usage in network core; spreading servers across fault domains improves fault tolerance, but requires additional bandwidth, while deploying servers together reduces bandwidth usage, but also decreases fault tolerance. We present a detailed analysis of a large-scale Web application and its communication patterns. Based on that, we propose and evaluate a novel optimization framework that achieves both high fault tolerance and significantly reduces bandwidth usage in the network core by exploiting the skewness in the observed communication patterns.
Peter Bodík, Ishai Menache, Mosharaf Chowdhury, Pradeepkumar Mani, David A. Maltz, Ion Stoica
SIGCOMM1
2011 The SCADS Director: Scaling a Distributed Storage System Under Stringent Performance Requirements
Beth Trushkowsky, Peter Bodík, Armando Fox, Michael J. Franklin, Michael I. Jordan, David A. Patterson 0001
FAST2
2010 Characterizing, modeling, and generating workload spikes for stateful services
abstract
Evaluating the resiliency of stateful Internet services to significant workload spikes and data hotspots requires realistic workload traces that are usually very difficult to obtain. A popular approach is to create a workload model and generate synthetic workload, however, there exists no characterization and model of stateful spikes. In this paper we analyze five workload and data spikes and find that they vary significantly in many important aspects such as steepness, magnitude, duration, and spatial locality. We propose and validate a model of stateful spikes that allows us to synthesize volume and data spikes and could thus be used by both cloud computing users and providers to stress-test their infrastructure.
Peter Bodík, Armando Fox, Michael J. Franklin, Michael I. Jordan, David A. Patterson 0001
SoCC1
2010 Fingerprinting the datacenter: automated classification of performance crises
abstract
Contemporary datacenters comprise hundreds or thousands of machines running applications requiring high availability and responsiveness. Although a performance crisis is easily detected by monitoring key end-to-end performance indicators (KPIs) such as response latency or request throughput, the variety of conditions that can lead to KPI degradation makes it difficult to select appropriate recovery actions. We propose and evaluate a methodology for automatic classification and identification of crises, and in particular for detecting whether a given crisis has been seen before, so that a known solution may be immediately applied. Our approach is based on a new and efficient representation of the datacenter's state called a fingerprint, constructed by statistical selection and summarization of the hundreds of performance metrics typically collected on such systems. Our evaluation uses 4 months of trouble-ticket data from a production datacenter with hundreds of machines running a 24x7 enterprise-class user-facing application. In experiments in a realistic and rigorous operational setting, our approach provides operators the information necessary to initiate recovery actions with 80% correctness in an average of 10 minutes, which is 50 minutes earlier than the deadline provided to us by the operators. To the best of our knowledge this is the first rigorous evaluation of any such approach on a large-scale production installation.
Peter Bodík, Moisés Goldszmidt, Armando Fox, Dawn B. Woodard, Hans Andersen
EuroSys1
2004 Distributed regression: an efficient framework for modeling sensor network data
abstract
We present distributed regression, an efficient and general framework for in-network modeling of sensor data. In this framework, the nodes of the sensor network collaborate to optimally fit a global function to each of their local measurements. The algorithm is based upon kernel linear regression, where the model takes the form of a weighted sum of local basis functions; this provides an expressive yet tractable class of models for sensor network data. Rather than transmitting data to one another or outside the network, nodes communicate constraints on the model parameters, drastically reducing the communication required. After the algorithm is run, each node can answer queries for its local region, or the nodes can efficiently transmit the parameters of the model to a user outside the network. We present an evaluation of the algorithm based upon data from a 48-node sensor network deployment at the Intel Research - Berkeley Lab, demonstrating that our distributed algorithm converges to the optimal solution at a fast rate and is very robust to packet losses.
Carlos Guestrin, Peter Bodík, Romain Thibaux, Mark A. Paskin, Samuel Madden 0001
IPSN2