Deciding on a Memcached Caching Layer: Lessons from Facebook's Read-Scaling Architecture
A decision guide to adding a Memcached tier in front of MySQL, based on Facebook's published caching architecture: options compared, trade-offs, and a validation setup.
13 Jun 2026, 21:21 UTC

If your MySQL cluster is spending most of its time answering the same read queries, the decision you face is not whether to cache, but where the cache lives, how keys are distributed, and how much inconsistency you can tolerate. Facebook's well-documented approach — a large Memcached tier in front of MySQL — is a useful reference design because it makes each of those trade-offs explicit. This guide walks through the decision points and shows a small-scale setup you can use to validate the model before committing.
The decision and its constraints
The core problem at scale is read amplification: profile pages, feeds, and object lookups generate far more reads than writes, and hitting the database for every read wastes CPU and I/O on repeated work. Memcached (an in-memory key-value store with no persistence and no replication built in) became the standard answer because it is simple, fast, and horizontally scalable. The constraints that shaped Facebook's design apply to most teams:
- Read latency must stay in the low milliseconds, so the cache must be in-memory and network-local.
- The dataset is far larger than one server's RAM, so keys must be sharded across many nodes.
- Some staleness is acceptable for most social data; strict consistency is not worth the coordination cost.
- Cache nodes fail and are replaced routinely, so the system must treat the cache as disposable.
Comparing the supported options
| Option | Latency | Consistency | Ops complexity | Best fit |
|---|---|---|---|---|
| Query cache in MySQL | Low | Strong (invalidated on write) | None | Small deployments; known to scale poorly under write load |
| Single Memcached node | Very low | Eventual (TTL/invalidation) | Low | One app server, modest dataset |
| Sharded Memcached pool (client-side hashing) | Very low | Eventual | Medium | Read-heavy workloads, the Facebook model |
| Redis with replication/persistence | Very low | Tunable | Higher | When you need data structures or durability |
The Facebook-style choice is the third row: a pool of Memcached servers where the client library — not a proxy — decides which node owns each key. This removes any single point of coordination and keeps lookups to one network hop, at the cost of pushing sharding logic into every client.
Trade-offs you are accepting
Eventual consistency. Writes go to MySQL; the cache is updated by invalidation (delete the key so the next read repopulates it). If invalidation is lost or delayed, readers see stale data until the TTL expires. Facebook's published engineering work describes using delete-based invalidation rather than cache writes precisely because deletes are idempotent — a retried delete is harmless, a retried set may not be.
Cache stampedes. When a hot key expires, many concurrent requests can miss simultaneously and hammer the database. Common mitigations are a per-key lock (only one request regenerates the value; others wait or serve stale) and adding random jitter to TTLs so keys do not expire in waves.
Node loss reshuffles keys. With naive modulo hashing, removing one of N nodes remaps most keys and causes a mass miss event. Consistent hashing limits remapping to roughly 1/N of keys — this is why client libraries like libmemcached offer it, and you should enable it from day one.
Eviction is LRU and silent. Under memory pressure Memcached evicts least-recently-used keys without telling you. Miss-rate monitoring is the only way to know your pool is undersized.
Concrete validation setup
You can validate key distribution and latency on a single Linux machine with Docker before touching production. Run these from a shell with Docker installed; no root beyond Docker group membership is needed.
docker run -d --name mc1 -p 11211:11211 memcached:1.6 -m 256
docker run -d --name mc2 -p 11212:11211 memcached:1.6 -m 256
docker run -d --name mc3 -p 11213:11211 memcached:1.6 -m 256Then, from a Python environment with pymemcache installed, write keys through a consistent-hashing client and check distribution:
from pymemcache.client.hash import HashClient
client = HashClient([('127.0.0.1',11211),('127.0.0.1',11212),('127.0.0.1',11213)])
for i in range(10000):
client.set(f'user:{i}', f'payload-{i}')Verify distribution with docker exec mc1 memcached-tool 127.0.0.1:11211 stats (or echo stats items | nc 127.0.0.1 11211) and compare curr_items across the three nodes — consistent hashing should land within a few percent of 33/33/33. Kill one container and re-read keys: expect roughly a third to miss, not all of them. That is the behavior you are buying.
For the stampede check, set a key with a 2-second TTL, fire 50 concurrent reads at expiry, and watch MySQL's Questions counter (or a local test database's query log). Without a regeneration lock you will see multiple identical queries; with one, exactly one.
Limitations and how to check the result
This article reflects long-established Memcached behavior and Facebook's published architecture descriptions; exact internals of Facebook's current stack (mcrouter, regional cache tiers) have evolved and should be verified against their engineering publications before you cite them. Version assumption: Memcached 1.6.x, pymemcache 4.x.
Success looks like: cache hit rate above ~90% for your hot keyspace (watch get_hits/get_misses in stats), database read QPS dropping proportionally, and p99 read latency staying flat when you remove a cache node. If hit rate is high but DB load barely moves, your misses are the expensive queries — cache those, not the cheap ones.
0 replies
A thoughtful contribution can make all the difference. Be the first to share one.