Dhruba Borthakur

dblp:51/8165 · DBLP profile ↗
← Back
14ranked-venue papers
2as first author
0since 2021 · last 2017
—ORCID · none

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

Databases, data management, data science and information retrieval · 7 · 2 first-authorSystems, architecture and hardware · 4Computer networks · 4

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
12 papers
Storage systems · 52% Cloud and datacenter computing · 21% Performance modeling and evaluation · 9%
Databases, data mining, and information retrieval
5 papers
Distributed and cloud data management · 42% Graph data management · 24% Data integration and cleaning · 16%
Theoretical computer science
1 paper
Coding theory · 100%
Computer networks
1 paper
Datacenter networks · 100%

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

TopicWeightPapersLastEvidence papers
Storage systems › storage reliability
erasure coding
0.422014
A "hitchhiker's" guide to fast and efficient data reconstruction in erasure-coded data centers · SIGCOMM 2014
XORing Elephants: Novel Erasure Codes for Big Data · Proc. VLDB Endow. 2013
Storage systems
storage reliability
0.422014
A "hitchhiker's" guide to fast and efficient data reconstruction in erasure-coded data centers · SIGCOMM 2014
XORing Elephants: Novel Erasure Codes for Big Data · Proc. VLDB Endow. 2013
Cloud and datacenter computing
cluster resource management and scheduling
0.322012
Energy efficiency for large-scale MapReduce workloads with significant interactive analysis · EuroSys 2012
Delay scheduling: a simple technique for achieving locality and fairness in cluster scheduling · EuroSys 2010
Storage systems
distributed storage
0.222014
A "hitchhiker's" guide to fast and efficient data reconstruction in erasure-coded data centers · SIGCOMM 2014
XORing Elephants: Novel Erasure Codes for Big Data · Proc. VLDB Endow. 2013
Storage systems › repair
data reconstruction
0.212014
A "hitchhiker's" guide to fast and efficient data reconstruction in erasure-coded data centers · SIGCOMM 2014
Storage systems › file systems
distributed file system
0.212014
Analysis of HDFS under HBase: a facebook messages case study · FAST 2014
Storage systems › file systems › distributed file system
HDFS
0.212014
Analysis of HDFS under HBase: a facebook messages case study · FAST 2014
Performance modeling and evaluation
benchmarking
0.212013
LinkBench: a database benchmark based on the Facebook social graph · SIGMOD Conference 2013
Performance modeling and evaluation › benchmarking
database system benchmarking
0.212013
LinkBench: a database benchmark based on the Facebook social graph · SIGMOD Conference 2013
Storage systems › erasure-coded storage
locally repairable codes
0.212013
XORing Elephants: Novel Erasure Codes for Big Data · Proc. VLDB Endow. 2013
Storage systems › multimedia storage
photo storage
0.212013
Petabyte scale databases and storage systems at Facebook · SIGMOD Conference 2013
Coding theory › error-correcting codes
erasure coding
0.212013
XORing Elephants: Novel Erasure Codes for Big Data · Proc. VLDB Endow. 2013
Coding theory › error-correcting codes
reed-solomon codes
0.212013
XORing Elephants: Novel Erasure Codes for Big Data · Proc. VLDB Endow. 2013
Datacenter networks
datacenter transport
0.112012
DeTail: reducing the flow completion time tail in datacenter networks · SIGCOMM 2012
Datacenter networks › datacenter transport
flow completion time
0.112012
DeTail: reducing the flow completion time tail in datacenter networks · SIGCOMM 2012
Energy-efficient computing
datacenter energy efficiency
0.112012
Energy efficiency for large-scale MapReduce workloads with significant interactive analysis · EuroSys 2012
Cloud and datacenter computing › job scheduling
workload-aware scheduling
0.112012
Energy efficiency for large-scale MapReduce workloads with significant interactive analysis · EuroSys 2012
Cloud and datacenter computing
cloud service reliability
0.112011
FATE and DESTINI: A Framework for Cloud Recovery Testing · NSDI 2011
Data integration and cleaning
data warehouse
0.112010
Data warehousing and analytics infrastructure at facebook · SIGMOD Conference 2010
Memory systems
data locality
0.112010
Delay scheduling: a simple technique for achieving locality and fairness in cluster scheduling · EuroSys 2010
Cloud and datacenter computing › job scheduling
fair scheduling
0.112010
Delay scheduling: a simple technique for achieving locality and fairness in cluster scheduling · EuroSys 2010
Cloud and datacenter computing › cluster resource management and scheduling
cluster resource management
0.122012
PACMan: Coordinated Memory Caching for Parallel Jobs · NSDI 2012
Data warehousing and analytics infrastructure at facebook · SIGMOD Conference 2010
Web and social media mining
social network analysis
0.012013
LinkBench: a database benchmark based on the Facebook social graph · SIGMOD Conference 2013
Query processing and optimization
interactive data exploration
0.012012
Energy efficiency for large-scale MapReduce workloads with significant interactive analysis · EuroSys 2012
Cloud and datacenter computing
datacenter application performance
0.012012
DeTail: reducing the flow completion time tail in datacenter networks · SIGCOMM 2012
Distributed systems
consistency and availability tradeoff
0.012011
Apache hadoop goes realtime at Facebook · SIGMOD Conference 2011
Distributed systems
fault tolerance
0.012011
FATE and DESTINI: A Framework for Cloud Recovery Testing · NSDI 2011
Hardware reliability and fault tolerance › network fault tolerance
network partition tolerance
0.012011
Apache hadoop goes realtime at Facebook · SIGMOD Conference 2011

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

sharding · 0.3log mining · 0.3locality-distance tradeoff analysis · 0.3caching · 0.3XOR-based coding · 0.3workload characterization · 0.3trace analysis · 0.3system configuration tuning · 0.2reed-solomon codes · 0.2encoding and decoding techniques · 0.2
YearPublicationVenuePosition
2017 Optimizing Space Amplification in RocksDB
Siying Dong, Mark Callaghan, Leonidas Galanis, Dhruba Borthakur, Tony Savor, Michael Strum
CIDR4
2014 Analysis of HDFS under HBase: a facebook messages case study
Tyler Caraza-Harter, Dhruba Borthakur, Siying Dong, Amitanand S. Aiyer, Liyin Tang, Andrea C. Arpaci-Dusseau, Remzi H. Arpaci-Dusseau
FAST2
2014 A "hitchhiker's" guide to fast and efficient data reconstruction in erasure-coded data centers
abstract
Erasure codes such as Reed-Solomon (RS) codes are being extensively deployed in data centers since they offer significantly higher reliability than data replication methods at much lower storage overheads. These codes however mandate much higher resources with respect to network bandwidth and disk IO during reconstruction of data that is missing or otherwise unavailable. Existing solutions to this problem either demand additional storage space or severely limit the choice of the system parameters. In this paper, we present "Hitchhiker", a new erasure-coded storage system that reduces both network traffic and disk IO by around 25% to 45% during reconstruction of missing or otherwise unavailable data, with no additional storage, the same fault tolerance, and arbitrary flexibility in the choice of parameters, as compared to RS-based systems. Hitchhiker 'rides' on top of RS codes, and is based on novel encoding and decoding techniques that will be presented in this paper. We have implemented Hitchhiker in the Hadoop Distributed File System (HDFS). When evaluating various metrics on the data-warehouse cluster in production at Facebook with real-time traffic and workloads, during reconstruction, we observe a 36% reduction in the computation time and a 32% reduction in the data read time, in addition to the 35% reduction in network traffic and disk IO. Hitchhiker can thus reduce the latency of degraded reads and perform faster recovery from failed or decommissioned machines.
K. V. Rashmi, Nihar B. Shah, Dikang Gu, Hairong Kuang, Dhruba Borthakur, Kannan Ramchandran
SIGCOMM5
2013 A Solution to the Network Challenges of Data Recovery in Erasure-coded Distributed Storage Systems: A Study on the Facebook Warehouse Cluster
K. V. Rashmi, Nihar B. Shah, Dikang Gu, Hairong Kuang, Dhruba Borthakur, Kannan Ramchandran
HotStorage5
2013 LinkBench: a database benchmark based on the Facebook social graph
abstract
Database benchmarks are an important tool for database researchers and practitioners that ease the process of making informed comparisons between different database hardware, software and configurations. Large scale web services such as social networks are a major and growing database application area, but currently there are few benchmarks that accurately model web service workloads.
Timothy G. Armstrong, Vamsi Ponnekanti, Dhruba Borthakur, Mark Callaghan
SIGMOD Conference3
2013 Petabyte scale databases and storage systems at Facebook
abstract
At Facebook, we use various types of databases and storage system to satisfy the needs of different applications. The solutions built around these data store systems have a common set of requirements: they have to be highly scalable, maintenance costs should be low and they have to perform efficiently. We use a sharded mySQL+memcache solution to support real-time access of tens of petabytes of data and we use TAO to provide consistency of this web-scale database across geographical distances. We use Haystack data store for storing the 3 billion new photos we host every week. We use Apache Hadoop to mine intelligence from 100 petabytes of click logs and combine it with the power of Apache HBase to store all Facebook Messages.
Dhruba Borthakur
SIGMOD Conference1
2013 XORing Elephants: Novel Erasure Codes for Big Data
abstract
Distributed storage systems for large clusters typically use replication to provide reliability. Recently, erasure codes have been used to reduce the large storage overhead of three-replicated systems. Reed-Solomon codes are the standard design choice and their high repair cost is often considered an unavoidable price to pay for high storage efficiency and high reliability. This paper shows how to overcome this limitation. We present a novel family of erasure codes that are efficiently repairable and offer higher reliability compared to Reed-Solomon codes. We show analytically that our codes are optimal on a recently identified tradeoff between locality and minimum distance. We implement our new codes in Hadoop HDFS and compare to a currently deployed HDFS module that uses Reed-Solomon codes. Our modified HDFS implementation shows a reduction of approximately 2× on the repair disk I/O and repair network traffic. The disadvantage of the new coding scheme is that it requires 14% more storage compared to Reed-Solomon codes, an overhead shown to be information theoretically optimal to obtain locality. Because the new codes repair failures faster, this provides higher reliability, which is orders of magnitude higher compared to replication.
Maheswaran Sathiamoorthy, Megasthenis Asteris, Dimitris S. Papailiopoulos, Alexandros G. Dimakis, Ramkumar Vadali, Scott Chen, Dhruba Borthakur
Proc. VLDB Endow.7
2012 Energy efficiency for large-scale MapReduce workloads with significant interactive analysis
abstract
MapReduce workloads have evolved to include increasing amounts of time-sensitive, interactive data analysis; we refer to such workloads as MapReduce with Interactive Analysis (MIA). Such workloads run on large clusters, whose size and cost make energy efficiency a critical concern. Prior works on MapReduce energy efficiency have not yet considered this workload class. Increasing hardware utilization helps improve efficiency, but is challenging to achieve for MIA workloads. These concerns lead us to develop BEEMR (Berkeley Energy Efficient MapReduce), an energy efficient MapReduce workload manager motivated by empirical analysis of real-life MIA traces at Facebook. The key insight is that although MIA clusters host huge data volumes, the interactive jobs operate on a small fraction of the data, and thus can be served by a small pool of dedicated machines; the less time-sensitive jobs can run on the rest of the cluster in a batch fashion. BEEMR achieves 40-50% energy savings under tight design constraints, and represents a first step towards improving energy efficiency for an increasingly important class of datacenter workloads.
Yanpei Chen, Sara Alspaugh, Dhruba Borthakur, Randy H. Katz
EuroSys3
2012 PACMan: Coordinated Memory Caching for Parallel Jobs
Ganesh Ananthanarayanan, Ali Ghodsi 0002, Andy Warfield, Dhruba Borthakur, Srikanth Kandula, Scott Shenker, Ion Stoica
NSDI4
2012 DeTail: reducing the flow completion time tail in datacenter networks
abstract
Web applications have now become so sophisticated that rendering a typical page may require hundreds of intra-datacenter flows. At the same time, web sites must meet strict page creation deadlines of 200-300ms to satisfy user demands for interactivity. Long-tailed flow completion times make it challenging for web sites to meet these constraints. They are forced to choose between rendering a subset of the complex page, or delay its rendering, thus missing deadlines and sacrificing either quality or responsiveness. Either option leads to potential financial loss.
David Zats, Tathagata Das, Prashanth Mohan, Dhruba Borthakur, Randy H. Katz
SIGCOMM4
2011 FATE and DESTINI: A Framework for Cloud Recovery Testing
Haryadi S. Gunawi, Thanh Do, Pallavi Joshi, Peter Alvaro, Joseph M. Hellerstein, Andrea C. Arpaci-Dusseau, Remzi H. Arpaci-Dusseau, Koushik Sen, Dhruba Borthakur
NSDI9
2011 Apache hadoop goes realtime at Facebook
abstract
Facebook recently deployed Facebook Messages, its first ever user-facing application built on the Apache Hadoop platform. Apache HBase is a database-like layer built on Hadoop designed to support billions of messages per day. This paper describes the reasons why Facebook chose Hadoop and HBase over other systems such as Apache Cassandra and Voldemort and discusses the application's requirements for consistency, availability, partition tolerance, data model and scalability. We explore the enhancements made to Hadoop to make it a more effective realtime system, the tradeoffs we made while configuring the system, and how this solution has significant advantages over the sharded MySQL database scheme used in other applications at Facebook and many other web-scale companies. We discuss the motivations behind our design choices, the challenges that we face in day-to-day operations, and future capabilities and improvements still under development. We offer these observations on the deployment as a model for other companies who are contemplating a Hadoop-based solution over traditional sharded RDBMS deployments.
Dhruba Borthakur, Jonathan Gray, Joydeep Sen Sarma, Kannan Muthukkaruppan, Nicolas Spiegelberg, Hairong Kuang, Karthik Ranganathan, Dmytro Molkov, Aravind Menon, Samuel Rash, Rodrigo Schmidt, Amitanand S. Aiyer
SIGMOD Conference1
2010 Delay scheduling: a simple technique for achieving locality and fairness in cluster scheduling
abstract
As organizations start to use data-intensive cluster computing systems like Hadoop and Dryad for more applications, there is a growing need to share clusters between users. However, there is a conflict between fairness in scheduling and data locality (placing tasks on nodes that contain their input data). We illustrate this problem through our experience designing a fair scheduler for a 600-node Hadoop cluster at Facebook. To address the conflict between locality and fairness, we propose a simple algorithm called delay scheduling: when the job that should be scheduled next according to fairness cannot launch a local task, it waits for a small amount of time, letting other jobs launch tasks instead. We find that delay scheduling achieves nearly optimal data locality in a variety of workloads and can increase throughput by up to 2x while preserving fairness. In addition, the simplicity of delay scheduling makes it applicable under a wide variety of scheduling policies beyond fair sharing.
Matei Zaharia, Dhruba Borthakur, Joydeep Sen Sarma, Khaled Elmeleegy, Scott Shenker, Ion Stoica
EuroSys2
2010 Data warehousing and analytics infrastructure at facebook
abstract
Scalable analysis on large data sets has been core to the functions of a number of teams at Facebook - both engineering and non-engineering. Apart from ad hoc analysis of data and creation of business intelligence dashboards by analysts across the company, a number of Facebook's site features are also based on analyzing large data sets. These features range from simple reporting applications like Insights for the Facebook Advertisers, to more advanced kinds such as friend recommendations. In order to support this diversity of use cases on the ever increasing amount of data, a flexible infrastructure that scales up in a cost effective manner, is critical. We have leveraged, authored and contributed to a number of open source technologies in order to address these requirements at Facebook. These include Scribe, Hadoop and Hive which together form the cornerstones of the log collection, storage and analytics infrastructure at Facebook. In this paper we will present how these systems have come together and enabled us to implement a data warehouse that stores more than 15PB of data (2.5PB after compression) and loads more than 60TB of new data (10TB after compression) every day. We discuss the motivations behind our design choices, the capabilities of this solution, the challenges that we face in day today operations and future capabilities and improvements that we are working on.
Ashish Thusoo, Zheng Shao, Suresh Anthony, Dhruba Borthakur, Namit Jain, Joydeep Sen Sarma, Raghotham Murthy, Hao Liu 0018
SIGMOD Conference4