Hiring BarSupport

Memcached

Technology guide42 min read6 diagrams

Why This Matters#

Memcached is not a weaker Redis. It is a deliberate decision to keep the cache server dumb and put every distributed-systems problem in the client. The server stores bytes under keys in RAM, evicts the least recently used ones when full, and knows nothing about other servers. There is no replication, no persistence, no cluster membership, no failover. Which server owns a key, what happens when that server dies, how the cache stays consistent with the database and how a hot key's expiry is kept from flattening the database are all answered by the client library or a routing tier in front of the servers.

That is why "add a Memcached tier" is a sentence interviewers push on. The L5 candidate draws a cache box and says "99% hit rate." The L6 candidate says "a pool of 40 Memcached nodes behind a routing tier with consistent hashing; reads are look-aside with leases so only one client refills a missing key; writes update the database and then delete the key, never set it; when a node dies its traffic goes to a small gutter pool with short TTLs instead of being rehashed onto healthy nodes; and I alarm on the database's read QPS, because a 1-point drop in hit rate from 99% to 98% doubles database load." The L7 candidate asks whether the company needs a shared cache platform with routing, warmup and shadowing, who owns its capacity, and what a cold restart of the whole tier costs the databases behind it.

The L5 → L6 gap is not knowing what a slab is. It is knowing that the database behind the cache is provisioned for the cache's miss rate, so every cache failure mode is really a database failure mode with a delay.

The L5 → L6 → L7 Contrast#

BehaviorSenior (L5)Staff (L6)Principal (L7)
First move"Put Memcached in front of the database""What's the miss budget the database can absorb? That sets the hit-rate floor, the node count and what we do when a node dies.""Is there a shared cache platform already? A team-owned cache fleet is a second database tier with no owner on day 400."
Distribution"Hash the key to a server"Consistent hashing with virtual nodes in the client or a router; knows modulo hashing remaps ~all keys when a node is addedStandardises one routing layer (client library or proxy) across languages so every service hashes the same way
Consistency"Set the cache on write"Update DB, then delete the key; leases stop stale sets from racing reads; TTL as the backstopDefines which data classes may be cached and with what staleness, signed off by data owners
Failure"If a node dies, we lose some cache"Names the miss storm: one node of 20 dying sends 5% of keys to the database at once; gutter pool or replicated pool for hot dataSizes the database for the worst plausible cache loss (a zone, a bad deploy) and owns that capacity decision
Hot keys"The cache handles hot keys"Leases or meta-protocol win flags against herds; near cache for the top keys; replicate hot keys across nodesRuns a hot-key detection service fleet-wide and a policy for which tier absorbs celebrity traffic
Ownership"Ops runs the cache"Platform owns nodes and routing; service teams own key schemas, TTLs and invalidationDecides cache-as-a-platform with chargeback by memory, vs per-team fleets
Why "Failure" separates levels

"If a node dies, we lose some cache" is correct and dangerously calm. Suppose 20 nodes, a 99% hit rate and a database provisioned for 2× the normal miss rate. One node dies: 5% of keys now miss. Database read load goes from 1% of traffic to 1% + 5% = 6%, a 6× jump, well past its 2× headroom. If the client then rehashes those keys onto the surviving nodes, the node that inherits a hot key can overload too. Facebook's memcache paper describes exactly this, noting a single key can account for 20% of a server's requests, and its answer was Gutter: about 1% of servers held in reserve, which take over a failed server's keys with short TTLs so misses are absorbed rather than passed to the database. The Staff answer names the miss storm and the mechanism; the Principal answer owns the database headroom it implies.

The 60-Second Pitch#

"I'd put a Memcached pool in front of the user and product databases as a look-aside cache. It's multi-threaded, so one large node uses all its cores and we don't run a shard per core. The client or a routing proxy places keys with consistent hashing, so adding a node moves about 1/N of the keys. Reads use leases, or the meta protocol's win flag, so on a miss exactly one caller refills the key and others wait a few milliseconds or take the stale value. Writes go to the database first and then delete the key. Nothing in Memcached is the source of truth: no replication, no persistence, so a node restart is a cold node. I size the database for losing one node at peak and use a small gutter pool so a dead node's keys don't stampede the database. If we need sorted sets, counters with durability or a lock, that's Redis, not this."

The Three Intents#

IntentConstraintStrategyFailure ModeCorrectness Bar
Look-aside cache for database readsRead-heavy, data recomputable from the DBClient-side consistent hashing; leases; delete-on-write; TTL backstopMiss storm on node loss; stale set race leaves old data cachedDatabase load stays within headroom; staleness bounded by TTL
Shared cache platform for many servicesDozens of teams, mixed key spaces, thousands of serversRouting tier (mcrouter-style or the built-in proxy) with prefix routing, pools, replication for hot pools, shadowing, warmupOne team's key space evicts another's; router misconfig breaks every servicePer-pool hit-rate SLOs; isolation between teams
Expensive-computation cacheLarge values (rendered pages, ML features), compute costs far more than RAMBigger items, extstore to flash for the cold tail, long TTLsLarge values waste slab memory; one value per request saturates networkCompute cost saved per GB of cache

🎯 Staff Move: "I'll design the first intent: a look-aside cache in front of the database for reads that can tolerate a few seconds of staleness. Everything in Memcached can be recomputed, and I'll size the database so losing a cache node at peak is a latency blip, not an outage."

The Staff Positions#

PositionRationale
Nothing in Memcached is the source of truthNo replication, no persistence; restarts and evictions drop data by design.
Delete on write, never set on writeA set from a writer can race a set from a reader holding an older value; delete plus a lease makes the next reader refill from the database.
Consistent hashing, never moduloModulo-N remaps ~(N−1)/N of keys when N changes; consistent hashing moves ~1/N.
Herd protection on every hot read pathWithout leases or win flags, a hot key's expiry sends every concurrent request to the database.
Don't rehash a dead node's keys onto live nodesHot keys follow and overload the next node; use a gutter or replicated pool.
Size the database for cache loss, not cache healthThe database's real peak is the miss rate during the worst cache incident.
TTL with jitter on everythingTTLs are the backstop for missed invalidations; jitter prevents synchronized expiry.

Architecture & Internals#

Only five internals change design decisions: the threading model, the slab allocator and LRU, client-side distribution, the protocol's concurrency primitives, and the absence of replication. The head-to-head is in Redis vs Memcached.

One Process, Many Threads, No Peers#

A Memcached server is a single process with a pool of worker threads (4 by default) handling connections via an event loop, listening on TCP 11211 with a default cap of 1,024 connections; UDP is off by default since 1.5.6 (configuration docs). Each server is independent: it never talks to another Memcached server.

Diagram: One Process, Many Threads, No Peers

Why it matters in design: because all cores serve one cache, a single 64-core node can do on the order of a million simple gets per second (typical, not a guarantee), and you size the fleet by memory and network, not by core count. Because there are no peers, there is no cluster to partition, no failover election and no replication lag. There is also nothing that notices a node is gone except the clients.

The Slab Allocator: Why a 600-Byte Item Can Cost 700 Bytes#

Memcached does not malloc per item. Memory given by -m is carved into 1 MB pages; each page is assigned to a slab class, and each class cuts its pages into fixed-size chunks. Chunk sizes grow geometrically by the growth factor (1.25 by default, -f), so classes look like 96, 120, 152, 192 … bytes up to the 1 MB maximum item size (-I) (UserInternals wiki). An item goes into the smallest chunk that fits key + value + ~50 bytes of header.

Diagram: The Slab Allocator: Why a 600-Byte Item Can Cost 700 Bytes

Three design consequences:

  1. Internal fragmentation. A 200-byte item in a 240-byte chunk wastes 17%. With a 1.25 factor the worst-case waste per item is ~20%; Facebook's paper describes using a finer factor of 1.07 for their workload.
  2. Eviction is per class. LRU runs inside each slab class, so a class that is short of pages evicts its items even while another class holds cold data. Since 1.5.0 the "modern" defaults are on: segmented LRU (HOT/WARM/COLD), an LRU crawler that reclaims expired items, and automatic slab rebalancing that moves pages between classes as sizes shift (1.5.0 release notes).
  3. Value size shape matters. A workload that changes its value sizes (a new serialisation format, a bigger JSON payload) shifts demand between classes; until pages rebalance, hit rate drops on the classes now in demand. Watch evictions per slab class after any payload change.

Client-Side Distribution: The Ring Lives in the Client#

The server has no idea which keys it should own. The client library (or the router) hashes the key onto a consistent-hash ring with many virtual points per server (the classic "ketama" scheme), so adding or removing one of N servers moves roughly 1/N of the keys rather than nearly all of them.

Consistent hashing moves only about 1/N of keys when a cache node joins or leaves, while modulo hashing remaps almost all of them.
modulo:      server = hash(key) % N
             N: 20 -> 21  => ~95% of keys change server => cache is effectively cold
consistent:  server = first point clockwise from hash(key) on a ring (100-200 vnodes/server)
             N: 20 -> 21  => ~1/21 = ~5% of keys move

The trap: every client must agree on the ring. Two services with different library versions, hash functions or server lists send the same key to different nodes: double memory, and invalidations that miss the copy the other service reads. That is the real argument for a routing tier: one place that owns the server list and the hash. Facebook built mcrouter, an open-source Memcached-protocol router that does consistent hashing, prefix routing, replicated pools, failover, shadowing and cold-cache warmup, and reported it handling close to 5 billion requests per second at peak (Meta engineering). Memcached itself now ships a built-in proxy (since 1.6.23), configured in Lua, that fronts pools of backend servers (proxy docs). The ring itself is covered in Consistent Hashing.

Concurrency Primitives: CAS, Leases and the Meta Protocol#

PrimitiveWhat it doesUse for
addStore only if the key doesn't existCheap lock-ish "first writer wins"; refill guard
gets + casRead with a version token; write only if unchangedRead-modify-write without lost updates
Leases (Facebook's extension, NSDI 2013)Miss returns a token; only the token holder may set; a delete invalidates outstanding tokens; tokens issued at most once per 10s per keyStale-set prevention and thundering-herd control
Meta protocol (mg, ms, md)Flags on gets and sets: N vivify on miss, W "you won the right to recache", Z someone else is recaching, R early recache near expiry, X/I stale and invalidateHerd protection and stale-while-revalidate in open-source Memcached (meta protocol docs)

Facebook's paper reports that, for a set of keys especially prone to thundering herds, leases cut the peak database query rate from 17K/s to 1.3K/s (NSDI 2013 paper). The meta protocol brings the same idea to stock Memcached: the first client to miss gets W and refills; concurrent clients see Z and either wait briefly or serve the stale value marked with X.

🎯 Staff Insight: "Memcached's interesting features are all about the miss, not the hit. Hits are easy. What protects the database is that on a miss exactly one caller refills, and a write can't be overwritten by a slower reader holding an older value."

No Replication, No Persistence — On Purpose#

A restart is an empty node. A dead node's keys are gone. extstore can spill values to flash to extend capacity, but it is still a cache. Availability is something you build around the servers:

MechanismHow it worksCostUse for
Accept missesDead node's keys hit the databaseDatabase headroomSmall fleets with generous DB capacity
Gutter poolClients retry failed gets on a small spare pool with short TTLs~1% extra serversLarge fleets; Facebook's approach
Replicated poolRouter writes to 2–3 copies (per zone), reads from one2–3× memoryHot, expensive-to-miss data
Zone-replicated clientClient writes to every zone, reads local, falls back to another zone on missN× memoryMulti-AZ with local-read latency; Netflix EVCache's model

Data Modeling / Core Usage — "The Entire Game": The Miss Path#

With Redis the game is picking the data structure. With Memcached there is one structure, so the game is the miss path: what happens on a miss, on a write and on an expiry. Get these right and the database never notices the cache's imperfections; get them wrong and the cache becomes the reason the database falls over.

The miss-path contract (per key family):
  1. key      -> namespaced, versioned, <= 250 bytes, one owner
  2. read     -> look-aside; on miss exactly one refiller (lease / W flag)
  3. write    -> DB first, then DELETE the key (never set from the writer)
  4. expiry   -> TTL as a backstop, with jitter; early recache for hot keys
  5. batch    -> multi-get per server; bound fan-out per request

Step 1: Key Design#

key = <service>:<entity>:v<schema_version>:<id>[:<variant>]
      user-svc:profile:v3:81723
      feed-svc:timeline:v7:81723:page1

rules:
  - keys <= 250 bytes, no spaces or control characters (protocol limit)
  - schema version in the key -> a format change is a deploy, not a flush
  - one service owns a key prefix; routers can send prefixes to different pools
  - never embed unbounded input (search strings) without hashing it

Bumping v3 to v4 is a mass invalidation without a flush: old keys age out by TTL and LRU. Plan for the cold-miss load that follows: a version bump on a hot family is, to the database, the same as losing that slice of cache.

Step 2: The Read Path With Herd Protection#

value, flags = mg(key, v, N30, R10, t)      # meta get: vivify on miss, recache if TTL < 10s
if flags has W:                              # we won: we refill
    value = db.read(id)
    ms(key, value, T=ttl_with_jitter())      # set; clears the win marker
elif flags has Z and value is stale (X):     # someone else refilling; serve stale
    return value
elif flags has Z and no value:               # someone else refilling; brief wait
    sleep(5-20 ms); retry once; else read DB with a concurrency cap
return value

On a hot key that expires while serving 20,000 requests per second, the naive path sends every request in the refill window (say 50ms) to the database: ~1,000 identical queries. With a win flag or lease, it sends one. The wider read-path design is in Read-Heavy Systems.

When a hot cache key expires, every concurrent request can miss at once; a lease or win flag lets one caller refill while the rest wait or serve stale.

Step 3: The Write Path — Delete, Don't Set#

Diagram: Step 3: The Write Path — Delete, Don't Set

Without the lease, the reader's late set of version 1 lands after the writer's delete and the cache holds stale data until the TTL expires. With "set on write" instead of delete, two writers racing can leave either value cached. Delete-plus-lease makes the database the only arbiter. For data that must be invalidated across regions or by many writers, Facebook's paper describes tailing the database's commit log and issuing deletes from there (mcsqueal), so invalidation is driven by what actually committed. The trade-offs among these consistency choices are in Caching Basics.

Step 4: TTLs and Jitter#

Data classTTLWhy
Profile, product details5–60 min, jittered ±10%Delete-on-write keeps it fresh; TTL catches missed deletes
Computed aggregates (counts, feeds)30s–5 minStaleness is the business trade; recompute is expensive
Negative results ("user not found")30–60sStops repeated misses for absent keys; short so creation shows up
Expensive rendersHoursCompute dominates; invalidate by version bump

Jitter matters because a batch job that warms 2 million keys with the same TTL schedules 2 million simultaneous misses. ttl × uniform(0.9, 1.1) spreads them over minutes.

Step 5: Multi-Get and the Fan-Out Problem#

A page that needs 200 keys spread across 50 servers issues ~50 parallel requests, and its latency is the slowest of the 50. Two effects matter at scale:

  • Tail amplification. If each server's p99 is 1ms, a 50-way fan-out sees at least one p99-or-worse response on ~40% of requests. Batch by server, cap fan-out, and put the most critical keys in a small, replicated pool.
  • Incast. Fifty responses arriving at one client at once can overflow switch buffers. Facebook's paper describes a client-side sliding window that limits outstanding requests, growing on success and shrinking on timeouts.
With wide fan-out, the slowest of N parallel cache calls sets the request's latency, so per-server p99 becomes the page's typical latency.

The Latency Budget calculator is the fastest way to show this in an interview.

🎯 Staff Move: "Reads are look-aside with a win flag so one caller refills; writes update the database and delete the key; TTLs are jittered backstops. If they ask about consistency, I'll say the cache is eventually consistent with a bound set by the TTL, and leases stop the one race that would make it permanently stale."


The Tunable Tradeoff — Hit Rate × Memory × Staleness#

database_read_load = cache_request_rate x (1 - hit_rate)

  1M req/s at 99.0% hit -> 10K req/s to the DB
  1M req/s at 98.0% hit -> 20K req/s to the DB   (1 point of hit rate = 2x DB load)
  1M req/s at 99.5% hit ->  5K req/s to the DB

Hit rate vs memory follows the key-popularity curve: the first GB buys the most;
each additional GB buys less. Measure it (shadow a bigger pool, or replay traces).
DialCheap endSafe endWho pays at the cheap end
Memory per poolJust enough for the hot setHeadroom for the warm tailDatabase, in miss load and cost
TTLLong (high hit rate)Short (fresher)Users seeing stale data
InvalidationTTL onlyDelete-on-write + commit-log invalidationProduct, in stale-read bugs
Node failure handlingRehash onto live nodesGutter or replicated poolThe next node and then the database
ReplicationNone2–3 copies per zoneDatabase on node loss; at the safe end, the memory bill
Herd protectionNoneLeases / win flags + stale servingDatabase during every hot-key expiry

🎯 Staff Move: "I'll treat hit rate as a database capacity number. At 1M reads per second, 99% versus 98% is 10K versus 20K database reads, so I'd rather pay for 20% more cache memory than a bigger database, and I'll replicate only the hot pool where a miss is expensive."


Anti-Patterns — What Kills Memcached Deployments#

1. Memcached as the Only Copy#

Sessions, shopping carts or rate-limit state stored only in Memcached vanish on restart, eviction or a node replacement, and the eviction is silent. Fix: anything that must survive goes to a database or a replicated store such as Redis with persistence; Memcached caches it at most.

2. Modulo Hashing#

hash(key) % N works until the fleet changes size. Adding one node to 20 remaps ~95% of keys, which is a full cold start of the cache and a database overload. Fix: consistent hashing in every client, ideally via one shared library or a router.

3. No Herd Protection on Hot Keys#

A celebrity profile, the home-page config, a flash-sale product: each expiry sends thousands of concurrent misses to the database. Fix: leases or meta-protocol win flags, early recache (R flag) before expiry, and a near cache for the top keys. See Hot Keys.

4. Rehashing a Dead Node's Keys Onto Live Nodes#

Automatic ejection plus rehash looks self-healing. In practice the node that inherits a hot key overloads, gets ejected too, and the failure walks around the ring. Fix: gutter pool, or replicated pools for hot data; eject slowly and deliberately.

5. Set-on-Write Invalidation#

Writers updating the cache directly race with readers refilling it; one ordering leaves stale data until TTL. Fix: delete on write, leases on refill, TTL backstop.

6. Huge Values and Wide Fan-Out#

A 900 KB serialized object read on every request saturates a node's NIC long before its CPU; a page that fetches 500 keys across 100 servers lives at the fleet's p99. Fix: split large values, compress, cache smaller projections; bound keys per request; batch by server.

7. Rolling Restarts Without Warmup#

Upgrading a 100-node pool one node every minute empties 1% of the cache per minute; the whole pool is cold within two hours, during business hours. Fix: restart slowly with hit-rate gates, use a router's warmup (cold node reads through from a warm replica), or shift traffic to a warm pool first.

8. One Shared Pool for Every Team#

A batch job that caches 200 GB of one-off values evicts every other team's hot set. Fix: pools per workload behind a router's prefix routing, each with its own memory budget and hit-rate SLO.


The Technology Landscape — Head-to-Head Comparison#

DimensionMemcachedRedis / ValkeyIn-process cache (Caffeine, Guava, LRU map)Managed Memcached (ElastiCache, Memorystore)
ModelOpaque bytes, LRU, multi-threadedData structures, single-threaded execution per shardObjects in the app's heapMemcached with node management
DistributionClient or router consistent hashingRedis Cluster (16,384 slots) or client shardingNone, per processClient hashing; auto-discovery of nodes
ReplicationNone (build it in the router or client)Async primary–replica, failoverNoneNone
PersistenceNone (extstore spills to flash, still a cache)Optional RDB/AOFNoneNone
Throughput per node~1M+ simple ops/s on a large node (typical)~100K–300K ops/s per shard (one core)Nanoseconds to microseconds, no networkSame as Memcached
Herd controlLeases (custom) or meta-protocol flagsLocks via SET NX, scriptsPer-process loading (e.g. get(key, loader))Meta protocol depends on engine version
Ops burdenLow for servers; routing and warmup are the workMedium: persistence, failover, big keysNone, but N copies and no cross-host invalidationLow
Pick whenLarge pure look-aside fleets; cost per GB and cores per node matterYou need anything beyond get/set, or replicationHottest few thousand keys, read-mostly configMemcached without running hosts

🎯 Staff Insight: "For a pure look-aside cache either works, so the tiebreaker is what's coming next. If product will ask for counters, leaderboards or locks, that's Redis. If the job is a huge, boring get/set tier where cost per GB and cores per node are the bill, Memcached's dumbness is the feature."


Patterns#

Pattern 1: Look-Aside With Leases#

The default. Application reads cache, on miss reads the database and refills under a lease or win flag; writes go to the database and delete. Works for 90% of read-heavy data. See Distributed Cache for the full design.

Pattern 2: Routing Tier With Prefix Pools#

Diagram: Pattern 2: Routing Tier With Prefix Pools

One router owns the server lists, the hash and the failover policy. Teams get isolated pools by key prefix with separate budgets. Shadowing lets you test a new instance type with real traffic before moving any keys.

Pattern 3: Gutter Pool for Failures#

A small pool (~1% of servers in Facebook's account) that takes over failed servers' gets with short TTLs. A failed get goes to the gutter; on miss the client reads the database and fills the gutter. Hot keys spread across idle gutter servers instead of piling onto one live node.

Pattern 4: Zone-Replicated Cache#

Write to a copy in every availability zone; read from the local zone; on a miss, fall back to another zone before the database. Reads stay in-zone (lower latency, no cross-zone transfer cost) and a zone loss doesn't empty the cache. This is the model Netflix's EVCache wiki describes.

Pattern 5: Commit-Log Invalidation#

Instead of every application server issuing deletes after writes, a daemon tails the database's replication or commit log and deletes the affected keys, batching them through routers. Invalidation follows what committed, survives application crashes, and replays after an outage. The same idea in event form is Transactional Outbox.


Scaling#

The Numbers#

ResourceDocumented default or typical figureDesign note
Worker threads4 default (-t)Raise toward core count on large nodes; docs warn against very large values (80+)
Connections1,024 default (-c)Thousands of app hosts × pools exhaust it; a router coalesces connections
Max item size1 MB default (-I)Larger is allowed via config; usually a sign the value should be split
Key length250 bytesHash long natural keys
Slab page / growth factor1 MB pages, factor 1.25Finer factor = less waste, more classes
Per-node throughput~1M+ simple gets/s on a large multi-core node (typical)Network usually saturates first with large values
LatencySub-ms p50 and p99 in-zone for small values (typical)Fan-out width, not the server, sets page latency
Node add/remove~1/N of keys move with consistent hashingEach change is a 1/N cold-miss event for the database

Sizing a Pool#

Example: product-detail cache (illustrative)
  working set: 80M products x 2 KB avg value          = ~160 GB of values
  with header, key and ~15% slab waste                = ~190 GB to hold everything
  hot 30% of products serve ~95% of reads             -> ~60 GB holds the hot set
  provision 120 GB for the warm tail                  -> measured hit rate 99.2%

  traffic: 400K reads/s -> normal misses 0.8% = 3.2K reads/s to the database

  option A: 6 nodes x 20 GB    lose 1 -> ~1/6 of keys miss -> ~+66K reads/s
  option B: 40 nodes x 3 GB    lose 1 -> ~1/40 of keys miss -> ~+10K reads/s

Same memory, very different failure. Losing one node of six puts ~17% of keys on the database at once, a ~20× jump over normal misses; losing one of forty is ~4×. The answers are more, smaller nodes (losing 1 of 40 is 2.5%), a gutter pool, or replication for the hot pool. Node count is a blast-radius decision as much as a memory one. The Sharding Planner and Cost Estimator help frame it.

Scaling Moves in Order#

  1. Fix the miss path first: herd protection, delete-on-write, jittered TTLs. Often worth more than any memory you could add.
  2. Add memory until the hit-rate curve flattens, measured by shadowing or trace replay, not guessed.
  3. Split pools by workload so a low-value family can't evict a high-value one.
  4. Add a routing tier once more than a handful of services share the fleet: one ring, connection pooling, failover, shadowing.
  5. Replicate hot pools (per zone) where a miss is expensive; leave the long tail unreplicated.
  6. Add regional pools and commit-log invalidation when going multi-region, so each region reads locally and writes invalidate everywhere.

🎯 Staff Move: "I'd rather run 40 smaller nodes than 6 big ones. Memory is the same, but losing a node becomes a 2.5% miss event instead of 17%, which is the difference between the database noticing and the database falling over."


Failure Modes & Recovery#

1. Node Loss → Miss Storm#

  • Symptom: One cache node dies; database read QPS jumps several-fold within seconds; p99 latency across services climbs.
  • Root cause: That node's 1/N of keys all miss at once, and the database was provisioned for the healthy miss rate.
  • Detection: Per-node get_hits dropping to zero; database read QPS vs baseline; client connection errors to one host.
  • Fix: Route failed gets to a gutter pool; shed or degrade non-critical reads; replace the node.
  • Prevention: More, smaller nodes; gutter pool; replication for hot pools; database headroom sized for one node loss at peak.

2. Thundering Herd on a Hot Key#

  • Symptom: Periodic database spikes aligned with a popular key's TTL; identical queries arriving by the hundred.
  • Root cause: Expiry or delete of a hot key with no refill coordination.
  • Detection: Database slow-query logs showing identical statements in bursts; per-key miss rate from client-side sampling.
  • Fix: Enable leases or meta win flags; serve stale while one caller refills; early recache.
  • Prevention: Herd protection in the shared client library by default; hot-key detection.

3. Permanently Stale Values (the Stale-Set Race)#

  • Symptom: A user changes their name; some pages show the old name for an hour, then it fixes itself.
  • Root cause: A reader's refill with an old value landed after the writer's delete; no lease to reject it, so it lived until TTL.
  • Detection: Sampled cache-vs-database consistency checks; user reports correlated with write bursts.
  • Fix: Delete the affected keys; shorten TTL temporarily.
  • Prevention: Leases (or add plus versioned values); delete-on-write only; commit-log invalidation for multi-writer data.

4. Eviction of the Hot Set (Slab Imbalance or Noisy Pool)#

  • Symptom: Hit rate drops from 99% to 95% after a deploy that changed value sizes, or after a batch job started caching.
  • Root cause: Demand shifted to slab classes with too few pages, or a new workload is evicting others in a shared pool.
  • Detection: evictions and evicted_unfetched per slab class; hit rate per key prefix; slabs_moved.
  • Fix: Let automove rebalance (or force it); move the noisy family to its own pool; add memory.
  • Prevention: Pools per workload; alarm on per-class evictions after payload changes.

5. Cold Cache After a Mass Restart#

  • Symptom: After a deploy, kernel patch or zone recovery, hit rate starts near zero and the database is at its limit for an hour.
  • Root cause: Memcached has no persistence; restarted nodes come back empty.
  • Detection: Fleet-wide hit rate; number of nodes with uptime < 1 hour.
  • Fix: Slow the restart; enable warmup from a warm replica or zone; throttle traffic to the database.
  • Prevention: Restart pacing gated on hit-rate recovery; warm-pool cutovers; Facebook describes cold-cluster warmup bringing a cluster back to full capacity in a few hours instead of a few days.

Operational Reality Matrix#

FailureDetection SignalBlast RadiusMitigationOwner
Node lossPer-node hits → 0; DB QPS spike1/N of keys; then the databaseGutter, degrade, replacePlatform; DB owner for headroom
Thundering herdIdentical DB queries in burstsOne key family; possibly the DBLeases, stale servingService team; platform owns the library default
Stale-set raceConsistency samplingAffected keys until TTLDelete keys; leasesService team
Hot-set evictionPer-class evictions, hit-rate dropOne poolRebalance, split poolPlatform + noisy team
Cold restartFleet hit rate, node uptimeWhole pool → databasePace restarts, warmupPlatform
Ring disagreementDuplicate keys, invalidation missesEvery service with a mismatched clientOne shared library or routerPlatform

When to Use vs. Alternatives#

NeedPickWhy
Large, pure look-aside cache of recomputable valuesMemcachedMulti-threaded, simple, cheap per GB
Counters, sorted sets, locks, queues, pub/subRedisData structures and atomic scripts
Cache that must survive node loss without refillRedis with replicas, or a replicated Memcached poolReplication
The hottest few thousand keysIn-process near cacheNo network hop at all
Static or public contentCDNCache at the edge; see CDN
Durable key-value data at scaleDynamoDB or a KV storeIt's a database, not a cache

When NOT to Use Memcached#

  • The data has no other home. Sessions, carts, rate-limit counters you can't lose.
  • You need anything beyond get/set/delete/incr/cas. Leaderboards, sets, streams, scripts: Redis.
  • The fleet is small and the team already runs Redis. Two caching technologies for a handful of nodes is operational cost without benefit.
  • Strong consistency with the database is required. A cache is eventually consistent by construction; reads that must see the latest write go to the primary.
  • Values are huge and rarely reused. Caching 5 MB blobs read once a day is object storage's job.
Diagram: When NOT to Use Memcached

Operational Concerns#

What the On-Call Actually Does#

  1. Watches the database, not just the cache. A cache incident shows up first as database read QPS and latency.
  2. Replaces dead nodes and watches the miss storm subside; confirms gutter or replica pools absorbed it.
  3. Paces restarts and upgrades with hit-rate gates: next node only when the previous one is back above ~95% of its normal hit rate.
  4. Investigates hit-rate drops: per-prefix hit rates, per-class evictions, a deploy that changed value sizes or bumped a key version.
  5. Hunts hot keys from client-side sampling and router stats; moves them to a near cache or replicated pool.
  6. Owns the router config: pool membership, prefix routes, failover thresholds, shadowing for new hardware.

Key Metrics & Alerts#

MetricHealthyAlert
Hit rate per pool (get_hits / (get_hits + get_misses))At its SLO, e.g. ≥ 98.5%Drop > 1 point for 10 min
Database read QPSNear baseline> 2× baseline (page)
evictions rate, per slab classLow, steadySpike after a deploy
curr_connections vs -c< 70%> 85%; listen_disabled_num rising
Per-node get latency p99Sub-ms> 2 ms sustained
Network throughput per node< 60% of NIC> 80% (large values or hot key)
Nodes with uptime < 1h0–1Many (mass restart in progress)

Interview Application — Staff-Level Plays#

Which Case Studies Use Memcached#

Case StudyHow Memcached Is UsedKey Pattern
Distributed CacheThe cache tier itselfConsistent hashing, leases, gutter pools
News FeedCached feed fragments and object lookupsMulti-get fan-out, hot-key replication
URL ShortenerRedirect lookupsLook-aside with long TTLs; negative caching
TypeaheadPrefix resultsShort TTLs; near cache for top prefixes
Multi-Region Active-ActiveRegional pools with cross-region invalidationCommit-log invalidation

Every System Design Question Has a Memcached Moment#

  • News feed: "Feed pages assemble ~100 objects; I batch multi-gets by server and keep the top 1,000 hot objects in a 2-second near cache, so a celebrity post doesn't melt one node."
  • URL shortener: "Redirects are look-aside with a 24-hour jittered TTL and 60-second negative caching for unknown codes, so scanners hammering random codes don't reach the database."
  • E-commerce: "Product pages cache with delete-on-write from the catalog service; inventory counts are never cached, because a stale 'in stock' is a broken promise."

What Interviewers Probe#

After You Say...They Will Ask...What They're Evaluating
"Cache with Memcached""What happens when a node dies?"Miss storm math; gutter or replication
"Update the cache on write""Two writers and a reader race. What's cached?"Delete-on-write; leases
"99% hit rate""What does 98% do to the database?"Miss rate as database capacity
"Hash keys across nodes""You add a node. How many keys move?"Consistent vs modulo hashing
"A popular key expires""How many requests hit the database?"Herd protection
"Why not Redis?""When would you pick Memcached?"Pure get/set at fleet scale; multi-threading; simplicity

Common Interview Mistakes#

What Candidates SayWhat Interviewers HearWhat Staff Engineers Say
"Store sessions in Memcached"Will lose sessions on restart"Sessions in a durable store; Memcached may cache them."
"On write, update the cache"Stale-set race"Write the database, delete the key, refill under a lease."
"If a node dies we just miss"Database overload unplanned"One node of N is 1/N of keys at once; database sized for that, plus a gutter pool."
"hash % N"Full cold cache on resize"Consistent hashing with virtual nodes, ~1/N keys move."
"Memcached is faster than Redis"Benchmark thinking"Multi-threaded per node, simpler; the choice is features versus fleet economics."

L5 vs L6 vs L7 Responses#

ScenarioL5 AnswerL6 / Staff AnswerL7 / Principal Answer
"Cache user profiles"Memcached, TTL 1 hourLook-aside with leases, delete-on-write, jittered TTL, DB sized for one-node lossShared cache platform with prefix pools, chargeback by memory, consistency classes signed off
"Database melted when a cache node died"Add cache nodesMiss storm: smaller nodes, gutter pool, replicate hot pool, stop rehashingMakes "survive one cache node at peak" a capacity standard for every database behind a cache
"Some users see stale data for an hour"Shorter TTLStale-set race; leases; commit-log invalidation for multi-writer dataDefines which data may be cached at what staleness; consistency sampling as a fleet SLO
"Should we move from Memcached to Redis?"Yes, Redis has more featuresOnly if we need structures or replication; migrate pool by pool behind the routerDecides one caching technology vs two based on fleet size and the cost of running both

The Staff Memcached Checklist#

  1. Miss budget: "The database absorbs 2× normal misses, so the hit rate floor is 99% and losing a node can't exceed that."
  2. Distribution: "Consistent hashing in one shared client or router; adding a node moves about 1/N of keys."
  3. Read path: "Look-aside with leases or meta win flags; stale-while-revalidate for hot keys."
  4. Write path: "Database first, then delete; TTL with jitter as a backstop."
  5. Failure: "Many small nodes, a gutter pool, replicated pools only for hot data."
  6. Ownership: "Platform owns nodes and routing; teams own key prefixes, TTLs and invalidation."

🎯 Staff Insight: Don't use Memcached as a store of record, as a session store, for anything needing structures or atomic multi-step logic, or where reads must see the latest write. The strongest Memcached signal is doing the miss-storm arithmetic for one dead node out loud.

Evaluation Rubric#

DimensionSenior (L5)Staff (L6)Principal (L7)
ConsistencyTTLsDelete-on-write, leases, bounded stalenessData classes with signed-off staleness; consistency sampling
Failure"Some misses"Miss-storm math, gutter, replication, restart pacingDatabase headroom standard; cold-start drills
DistributionHash to a serverConsistent hashing, one ring, router tierOne routing platform across languages
Hot keysNot raisedLeases, near cache, replicated keysFleet-wide hot-key detection and policy
Choice"Memcached is fast"Memcached vs Redis on features and fleet economicsOne caching technology or two, priced

Strong hire signals

SignalWhat It Sounds Like
Miss rate as capacity"98% versus 99% hit rate is double the database reads."
Owns the race"Delete on write and a lease on refill, or a slow reader re-caches the old value."
Blast-radius sizing"Forty small nodes so a dead one is 2.5% of keys."
Herd awareness"On a miss exactly one caller refills; everyone else waits or takes stale."
Knows the trade"Memcached is a dumb server on purpose; the router does the thinking."

Lean no-hire signals

SignalWhy It Misses the Bar
Treats the cache as durableWill lose data that has no other home
Set-on-writeWill ship stale-cache bugs
No answer for node lossDatabase overload is unplanned
Modulo hashingEvery resize is a cold start

Common false positives

  • Slab internals ≠ cache design. Knowing chunk sizes says nothing about the miss path.
  • "Our hit rate is 99.9%" ≠ resilience. Ask what the database sees when a node dies.

Beyond Staff: The Principal View#

Why L7 Sees This Problem Differently#

At Staff level Memcached is a well-configured cache in front of one database. At Principal level, the cache fleet is load-bearing infrastructure for every database in the company: the databases are sized for its hit rate, so its failure modes are the databases' failure modes, and its capacity is a shared budget that one team's batch job can spend for everyone. The L7 question is "Do we run caching as a platform with a router, pools, quotas and warmup, or as per-team fleets, and who owns the database headroom a cache loss implies?"

🧭 Principal Move: "Every database behind a cache is quietly provisioned for the cache's miss rate. So I want one number per database: the miss rate it can survive. Then cache capacity, replication and node size are set to keep the worst plausible cache incident under that number, and the database owner and the cache owner sign it together."

The Org-Level Fault Line#

Cache as a shared platform vs per-team cache fleets.

OptionWhat WorksWhat BreaksWho Pays
Per-team fleetsAutonomy, clear costDifferent clients and rings, no herd protection by default, 30 small fleets each restarted carelesslyDatabase owners during every cache incident
Shared platform with router and poolsOne ring, one library, warmup, shadowing, per-pool budgetsRouter becomes critical infrastructure; needs a teamPlatform headcount
Managed cache service per teamNo hostsSame consistency mistakes, now billed by the hourTeams, in cost and stale-data bugs

The Principal default: a platform-owned cache tier with a routing layer, per-team pools by key prefix, chargeback by provisioned memory, herd protection and delete-on-write in the shared client library, and a "survive one node at peak" standard negotiated with each database owner.

🧭 Principal Insight: "The router is where cache policy lives. Without one, every service is its own cache platform and every incident is a new discovery."

Cost Model#

Assumptions: memory-optimised node with ~50 GB usable cache at ~$500/month; loaded engineer ~$21K/month; savings expressed as database capacity the cache displaces. Directional only.

ScaleCache fleetCache/monthPlatform peopleWhat the cache saves
Startup3 nodes, 150 GB~$1.5K~0 (managed)A database replica or two
Growth40 nodes, 2 TB, router tier~$20K + routers1 FTE (~$21K)Database fleet sized at 1–2% of read traffic instead of 100%
Enterprise1,000+ nodes, 50 TB, regional pools~$500K+5–8 FTE (~$150K)The difference between a database tier and a database tier ×50

The lever is hit rate per dollar of memory: measure the hit-rate curve per pool and stop adding memory where it flattens. The hidden cost is database headroom for cache loss, often larger than the cache bill itself.

The 3-Year Evolution Path#

Diagram: The 3-Year Evolution Path

One-Way Doors vs Two-Way Doors#

DecisionReversibilityCost to Reverse
Treating cache data as durable (sessions, counters)One-way once products depend on itMigrating state to a real store under live traffic
Hash function and ring layoutOne-way-ishChanging it is a full cold start unless staged through a router
Router vs client-side routingTwo-way-ishEvery client library touched
Memcached vs Redis for a poolTwo-way behind a routerPool-by-pool cutover with warmup
TTLs, pool sizes, node typesTwo-wayConfig and paced rollout

The Standard I'd Write#

RFC-CACHE-003: Look-Aside Cache Contract

Scope: Every cache pool fronting a production datastore.

MUST
  1. Cache only data recomputable from a named source of truth.
  2. Use the platform client or router (one hash ring, herd protection on).
  3. Invalidate by delete after the source write commits; never set from writers.
  4. Set a TTL with >= 10% jitter on every key family.
  5. Declare, with the database owner, the miss rate the database can survive,
     and keep single-node loss at peak below it.
SHOULD
  6. Use versioned key prefixes for format changes.
  7. Serve stale on refill for hot key families.
  8. Replicate pools whose misses cost > 10 ms or hit a constrained database.

Exceptions: platform review, time-boxed.
Success metrics: zero database incidents caused by a single cache node loss;
  per-pool hit-rate SLOs met; zero cache-only state.

What I'd Tell the VP#

"Our databases only stay up because the cache in front of them answers about 99 of every 100 reads. That means a cache problem becomes a database outage within seconds, and today each team runs its cache differently. I'm proposing one shared cache platform with standard safety features, and a simple rule agreed with every database owner: losing one cache server at peak must never take a database down. It costs about one engineer and some extra servers, and it removes a class of outages we've had twice this year."

Principal Interview Signals#

SignalWhat It Sounds Like
Ties cache to database capacity"Every database is provisioned for the cache's miss rate; I want that number written down."
Platform thinking"One router, one ring, herd protection on by default."
Prices headroom"The database headroom for cache loss costs more than the cache."
Standardises consistency"Delete-on-write and TTL jitter are in the library, not in each team's memory."
Knows when to consolidate"Running Memcached and Redis for 10 nodes total isn't worth two on-call playbooks."

Staff answers that L7 interviewers find insufficient:

  • "Add a gutter pool" without who sizes the database for the residual misses.
  • "Use leases" as a per-service fix rather than a library default for every team.
  • "Memcached for the cache, Redis for the rest" without asking whether running both is worth it at this fleet size.

How Real Companies Built It#

Facebook — Scaling Memcache#

Facebook's NSDI 2013 paper describes building a distributed key-value cache on Memcached that handles billions of requests per second and holds trillions of items. It introduces leases (a 64-bit token on miss; a delete invalidates it; tokens issued at most once every 10 seconds per key), which cut the peak database query rate for herd-prone keys from 17K/s to 1.3K/s; Gutter, about 1% of servers in a cluster that absorb a failed server's traffic with quickly expiring entries instead of rehashing onto live servers; invalidation daemons (mcsqueal) that tail the database commit log; and cold cluster warmup, which brings an empty cluster to full capacity in a few hours instead of a few days (USENIX NSDI '13).

Staff insight: Every mechanism in that paper is about the miss path and failure handling, not hits. Leases, gutter and commit-log invalidation are the three sentences that turn a cache box into a Staff answer.

Meta — mcrouter#

Facebook open-sourced mcrouter, a Memcached-protocol router that looks like a server to clients and a client to servers. It provides connection pooling, consistent hashing, prefix routing to separate pools, replicated pools with failover, traffic shadowing for testing new hardware, cold-cache warmup and online reconfiguration, and was handling close to 5 billion requests per second at peak; Instagram and Reddit were early external users (Meta engineering).

Staff insight: The router is where a dumb cache becomes a platform. In an interview, proposing a routing tier is the answer to "how do 50 services share a cache safely?"

Netflix — EVCache#

Netflix's EVCache is an open-source distributed in-memory store built on Memcached and the spymemcached client for AWS; the name stands for Ephemeral, Volatile Cache (EVCache on GitHub). Its documentation describes writing to every availability zone, reading from the local zone for latency, and optionally falling back to another zone when data is missing, which is much faster than going to the source (EVCache wiki).

Staff insight: Replication was added in the client, not the server, and paid for in memory per zone. That's the honest answer to "Memcached has no replication": you build exactly as much as the miss cost justifies.


Practice Drill#

Drill 1: The Cache Node That Took Down the Database#

Prompt: "Your product catalog runs 6 large Memcached nodes (64 GB each) in front of a PostgreSQL primary with 4 read replicas, at 400K reads per second and a 99.2% hit rate. Yesterday one cache node's host failed; the replicas hit 100% CPU, and the site was degraded for 25 minutes. The team proposes doubling the cache node size. What do you do?"

Staff Answer

Bigger nodes make it worse. Run the numbers: normal misses are 0.8% of 400K = 3.2K reads/s. One node of six owns ~17% of keys, so losing it adds roughly 400K × 17% ≈ 66K reads/s of misses on top, a ~20× jump on replicas sized for maybe 2–3× headroom. Doubling node size keeps the fleet at six nodes (or fewer), so each failure is still 17% or more. The fix is to shrink the blast radius and stop the misses reaching the database. First, many smaller nodes: move to ~40 nodes of the same total memory behind consistent hashing, so a node loss is ~2.5% of keys, about 10K extra reads/s. Second, a gutter pool of 2–3 small nodes: when a get to a dead node fails, clients retry against the gutter with a 10–30s TTL, so repeated reads of that node's hot keys become gutter hits after the first miss, and the hot keys don't get rehashed onto live nodes. Third, replicate the hot pool: product details for the top ~5% of products, which serve most reads, get two copies in different zones; a single node loss there costs nothing. Fourth, herd protection: enable meta-protocol win flags in the shared client so each missing key is refilled by one caller, and serve stale for hot keys. Fifth, make the capacity explicit with the database owner: the replicas must absorb one-node loss at peak (~10K extra reads/s after the change) with 30% headroom, and we test it with a game day that kills a cache node during peak traffic. Metrics: database read QPS vs baseline, per-node hit rate, gutter hit rate, per-pool hit rate after the migration, which we pace one node at a time with hit-rate gates to avoid creating the cold-cache outage we're trying to prevent.

Why this is L6:

  • Does the miss-storm arithmetic and shows why bigger nodes increase the blast radius.
  • Uses gutter and replication instead of rehashing, and herd protection for refill.
  • Turns database headroom into an explicit, tested number with its owner.
  • Plans the migration itself so it doesn't cause a cold-cache incident.

What L7 adds:

  • Makes "survive one cache node loss at peak" a standard for every database behind a cache, with sign-off by both owners.
  • Prices the trade: 40 small nodes plus a gutter pool versus a larger replica fleet, and picks the cheaper insurance.
  • Moves herd protection and failure routing into the shared library or router so the next team gets them by default.
❌ Common L5 Trap

"Double the cache node size so the hit rate goes up, add two more read replicas for safety, and set up automatic failover so a dead node's keys are rehashed onto the remaining nodes."

Why this misses: Larger nodes keep each failure at a sixth of the keyspace or more, so the next host failure causes the same miss storm. Rehashing onto live nodes moves the dead node's hot keys onto a neighbour that can then overload, spreading the failure around the ring. Extra replicas pay for the symptom at full price every month. The real levers are node count (blast radius), gutter or replicated pools (absorbing misses) and herd protection (refilling once).


Quick Reference Card#

Model:           dumb, multi-threaded, volatile LRU cache; servers never talk to each other
Distribution:    client or router consistent hashing; ~1/N keys move per node change
                 (modulo moves ~(N-1)/N)
Defaults:        port 11211 TCP (UDP off since 1.5.6); 4 threads; 1,024 connections;
                 1 MB max item; 1 MB slab pages; growth factor 1.25; keys <= 250 bytes
Memory:          slab classes with fixed chunks; per-class segmented LRU; automove
                 rebalances pages (modern defaults since 1.5.0)
Durability:      none; restart = empty node; extstore extends to flash, still a cache
Herd control:    leases (token on miss, 1 per 10s per key at Facebook) or meta protocol
                 flags N / W / Z / R / X / I
Write path:      DB first, then DELETE; lease rejects stale sets; TTL + jitter backstop
Failure:         one of N nodes lost = 1/N of keys miss at once
                 -> many small nodes, gutter pool (~1% of servers), replicated hot pools
Hit rate math:   DB load = requests x (1 - hit rate); 99% -> 98% doubles DB reads
Routing tier:    mcrouter or the built-in proxy (1.6.23+): pools, prefix routing,
                 failover, shadowing, warmup

RED FLAGS
  - Memcached as the only copy of sessions, carts or counters
  - hash(key) % N
  - Set-on-write invalidation
  - No herd protection on hot keys
  - Rehashing a dead node's keys onto live nodes
  - A few huge nodes (big blast radius per failure)
  - Rolling restarts without hit-rate gates
  - One pool shared by batch jobs and the hot path
  1. Loading the index…