Skip to content

A look-aside cache at Facebook's scale

Design a Distributed Cache (Memcache)

Built from Facebook's paper on scaling memcache: a look-aside cache serving billions of reads a second, Reads are fast; the hard parts are stale values, stampedes on popular keys, dead servers and invalidations that have to cross regions.

Advanced, about 50 minutes, 9 stages

The situation

A social network renders every page from many small pieces of data: profiles, friend lists, posts, counts, permissions. A popular page fetches hundreds of distinct items; Facebook reported an average of 521 for its most popular pages. MySQL holds the truth, but it cannot serve this read load, and it is provisioned only for the traffic that misses the cache.

So the web servers use memcached as a demand-filled, look-aside cache: read the cache; on a miss, query the database and put the result in the cache. Reads exceed writes by orders of magnitude, and a little staleness is acceptable, as long as people see their own changes and nothing stays stale for long.

Facebook described how this grew from one cluster to many clusters and regions in "Scaling Memcache at Facebook" (NSDI 2013). This investigation follows the problems they hit and the mechanisms they built.

What it has to do

Functional

  • Serve cached values for arbitrary keys (query results, computed objects) to web servers.
  • Reflect writes: after a change, readers stop seeing the old value.
  • Keep serving when cache servers fail, and when new clusters start empty.
  • Work across several frontend clusters and several regions.

Non-functional

  • The database must never receive more than the miss load it is provisioned for.
  • A user sees their own writes on their next request.
  • Other users may see stale data briefly, never indefinitely.
  • Losing a few cache servers does not overload the database.

Constraints and assumptions

  • Billions of cache reads a second across the fleet; writes are a tiny fraction.
  • A page needs hundreds of keys, fetched in parallel batches.
  • One master region holds the MySQL primaries; other regions have read replicas.
  • Cached values are small and derived from the database.
  • Any cached key can be evicted at any time.
  • Brief staleness is acceptable for most data.

Interview questions it prepares you for

  • “Design a distributed cache.”
  • “How do you keep a cache consistent with the database?”
  • “A hot key expires and the database falls over. What happened and how do you prevent it?”
  • “Design Memcached or Redis as a service for a large company.”

Read and practise next

How Meta (Facebook) built it · Caching in front of a database, in their engineers' own words

Concepts to know first: Caching, Replication.

Similar systems: Design a URL Shortener, Design a News Feed (Twitter Timeline).