Scaling Memcache at Facebook
Networked Systems Design and ImplementationPublished 2 April 2013
Rajesh Nishtala, Hans Fugal, Steven Grimm, Marc Kwiatkowski, Herman Lee, Harry C. Li
Citations597
Generate an AI Snapshot to get a quick, structured summary of this paper.
Study Snapshot
ObjectiveStudy objective
MethodsResearch methodology
PopulationPopulation studied
Sample sizeSample sizes
OutcomesStudy outcomes here
ResultsStudy results comes here
LimitationsResearch study limitations comes here
A concise AI-generated summary of the paper will appear here once you click Generate AI Snapshot.
TL;DR
This paper describes how Facebook leverages memcached as a building block to construct and scale a distributed key-value store that supports the world's largest social network.
Abstract
Memcached is a well known, simple, in-memory caching solution. This paper describes how Facebook leverages memcached as a building block to construct and scale a distributed key-value store that supports the world's largest social network. Our system handles billions of requests per second and holds trillions of items to deliver a rich experience for over a billion users around the world.
Keywords
Computer Science
Chord
9,645 Citations2001Ion Stoica, Robert Morris +3 more
Results from theoretical analysis, simulations, and experiments show that Chord is scalable, with communication cost and the state maintained by each node scaling logarithmically with the number of Chord nodes.
A scalable content-addressable network
6,407 Citations2001Sylvia Ratnasamy, Paul Francis +3 more
The concept of a Content-Addressable Network (CAN) as a distributed infrastructure that provides hash table-like functionality on Internet-like scales is introduced and its scalability, robustness and low-latency properties are demonstrated through simulation.
Dynamo
3,454 Citations2007Giuseppe DeCandia, Deniz Hastorun +7 more
D Dynamo is presented, a highly available key-value storage system that some of Amazon's core services use to provide an "always-on" experience and makes extensive use of object versioning and application-assisted conflict resolution in a manner that provides a novel interface for developers to use.
ACM Transactions on Programming Languages and SystemsLinearizability: a correctness condition for concurrent objects
3,185 Citations1990Maurice Herlihy, Jeannette M. Wing
This paper defines linearizability, compares it to other correctness conditions, presents and demonstrates a method for proving the correctness of implementations, and shows how to reason about concurrent objects, given they are linearizable.
ACM Transactions on Computer SystemsThe part-time parliament
2,714 Citations1998Leslie Lamport
The Paxon parliament's protocol provides a new way of implementing the state machine approach to the design of distributed systems.
Consistent hashing and random trees
1,935 Citations1997David R. Karger, Eric Lehman +4 more
A family of caching protocols for distrib-uted networks that can be used to decrease or eliminate the occurrence of hot spots in the network, based on a special kind of hashing that is called consistent hashing.
ACM SIGCOMM Computer Communication ReviewA scalable content-addressable network
1,789 Citations2001Sylvia Ratnasamy, Paul Francis +3 more
Weighted voting for replicated data
1,323 Citations1979David K. Gifford
The algorithm guarantees serial consistency, admits temporary copies in a natural way by the introduction of copies with no votes, and has been implemented in the context of an application system called Violet.
Workload analysis of a large-scale key-value store
886 Citations2012Berk Atikoglu, Yuehai Xu +3 more
This paper collects detailed traces from Facebook's Memcached deployment, arguably the world's largest, and analyzes the workloads from multiple angles, including: request composition, size, and rate; cache efficacy; temporal patterns; and application use cases.
IEEE Transactions on ComputersCoda: a highly available file system for a distributed workstation environment
865 Citations1990Mahadev Satyanarayanan, James J. Kistler +4 more
IEEE Transactions on CommunicationsA Protocol for Packet Network Intercommunication
803 Citations1974Vinton G. Cerf, Robert E. Kahn
Don't settle for eventual
569 Citations2011Wyatt Lloyd, Michael J. Freedman +2 more
This paper identifies and defines a consistency model---causal consistency with convergent conflict handling, or causal+---that is the strongest achieved under these constraints and presents the design and implementation of COPS, a key-value store that delivers this consistency model across the wide-area.
Linux journalDistributed caching with memcached
512 Citations2004Brad Fitzpatrick
Speed up your database app with a simple, fast caching layer that uses your existing servers' spare memory.
Leases: an efficient fault-tolerant mechanism for distributed file cache consistency
484 Citations1989Cary Gray, David R. Cheriton
Operating Systems Design and ImplementationAn analysis of Linux scalability to many cores
342 Citations2010Silas Boyd-Wickizer, Austin T. Clements +5 more
There is no scalability reason to give up on traditional operating system organizations just yet, according to this analysis of seven system applications running on Linux on a 48- core computer.
SILT
289 Citations2011Hyeontaek Lim, Bin Fan +2 more
The design of three basic key-value stores each with a different emphasis on memory-efficiency and write-friendliness are designed and an analytical model for tuning system parameters carefully to meet the needs of different workloads is developed.
ACM SIGMETRICS Performance Evaluation ReviewWorkload analysis of a large-scale key-value store
264 Citations2012Berk Atikoglu, Yuehai Xu +3 more
Measurement and analysis of TCP throughput collapse in cluster-based storage systems
261 Citations2008Amar Phanishayee, Elie Krevat +5 more
This paper analyzes this Incast problem, explores its sensitivity to various system parameters, and examines the effectiveness of alternative TCP- and Ethernet-level strategies in mitigating the TCP throughput collapse.
Scalable, distributed data structures for internet service construction
236 Citations2000Steven D. Gribble, Eric Brewer +2 more
The distributed hash table simplifies Internet service construction by decoupling service-specific logic from the complexities of persistent, consistent state management, and by allowing services to inherit the necessary service properties from the DDS rather than having to implement the properties themselves.
Communications of the ACMThe case for RAMCloud
159 Citations2011John K. Ousterhout, Parag Agrawal +12 more
Scalable consistency in Scatter
152 Citations2011Lisa Glendenning, Ivan Beschastnikh +2 more
This paper describes the design, implementation, and evaluation of Scatter, a scalable and consistent distributed key-value storage system that adopts the highly decentralized and self-organizing structure of scalable peer-to-peer systems, while preserving linearizable consistency even under adverse circumstances.
Operating Systems Design and ImplementationTransactional consistency and automatic management in an application data cache
83 Citations2010Dan R. K. Ports, Austin T. Clements +3 more
Adding TxCache to an application increased the throughput of a web application by up to 5.2×, only slightly less than a non-transactional cache, showing that consistency does not have to come at the price of performance.
CPHASH
72 Citations2012Zviad Metreveli, Nickolai Zeldovich +1 more
Experiments show that CPHash has ~1.6x higher throughput than a hash table implemented using fine-grained locks and its cache misses are less expensive, because of less contention for the on-chip interconnect and DRAM.
ACM SIGOPS Operating Systems ReviewLeases: an efficient fault-tolerant mechanism for distributed file cache consistency
72 Citations1989Cary Gray, David R. Cheriton
An analytic model and an evaluation for file access in the V system show that leases of short duration provide good performance and the impact of leases on performance grows more significant in systems of larger scale and higher processor performance.
ACM SIGCOMM Computer Communication ReviewA protocol for packet network intercommunication
72 Citations2005Vinton G. Cerf, Robert E. Icahn
TAO
57 Citations2012Venkateshwaran Venkataramani, Zach Amsden +13 more
The model uses typed nodes and edges to express the relationships and actions that happen on Facebook, and TAO is a distributed implementation of the fbobject and association API that has been serving production traffic at Facebook for more than 2 years.
ACM SIGOPS Operating Systems ReviewRouteBricks
24 Citations2011Kevin Fall, Gianluca Iannaccone +5 more
This work proposes a software router architecture that parallelizes router functionality both across multiple servers and across multiple cores within a single server, and demonstrates a 40Gbps parallel router prototype.
Gumball
21 Citations2012Shahram Ghandeharizadeh, Jason Yap
Experimental results show GT enhances the accuracy of an application hundreds of folds and, in some cases, may reduce system performance slightly.
ACM SIGOPS Operating Systems ReviewLazyBase
19 Citations2010Kimberly Keeton, Charles B. Morrey +2 more
Initial results with LazyBase illustrate the feasibility of the pipelined model, highlight a rich space of trade-offs between result freshness and query performance, and often outperform existing solutions in the space.
Sustainable Computing Informatics and SystemsPower and performance evaluation of Memcached on the TILEPro64 architecture
15 Citations2012Mateusz Berezecki, Eitan Frachtenberg +2 more
The results suggest that the TILEPro64 architecture can significantly outperform x86-based architectures in terms of throughput per Watt for the evaluated version of Memcached.
