Instagram News Feed
Worked solution

Instagram News Feed — a worked solution

Ranked feed with cursor pagination. No `OFFSET`.

Try it yourself first.

You will remember almost none of this if you read it cold. The workspace walks the same 10 stages and runs the architecture you draw through a simulator, so you find out where your version breaks before you see ours.

Open the Instagram News Feed workspace

The problem

Build the Instagram News Feed: a personalized, ranked, infinitely-scrollable list of posts from the people you follow. The product is read-dominated (≈ 85:1 read:write — feed page fetches 350K QPS peak vs post-creates 4.1K QPS peak), personalized per user (no edge cache for ranked JSON), and ranking-bound (the model is the product — chronological feeds have 30% lower session length per Meta A/B disclosures). The architecture is a two-stage retrieval funnel (candidate gen + early-stage ranker → late-stage ranker), a hybrid fanout (denorm-on-write for normal users, pull-on-read for >100K-follower creators), and cursor pagination (no OFFSET — at depth 1000, offset is 17× slower per the standard cursor-pagination analysis).

This canonical is sized for ~1.5B DAU, derived from public Meta scale figures (Q4 2023 earnings; Mosseri 2021 IG MAU disclosure) and the published architectures it builds on: TAO (USENIX ATC '13), Memcache + leases (NSDI '13), and Haystack/f4/Tectonic (OSDI '10/'14, FAST '21).

The reference architecture

Reference architecture for Instagram News Feed: 19 components — Mobile, CDN · Edge POPs, API Gateway · GraphQL, Auth / Identity, Service · Feed Read, Service · Candidate Gen + ESR, Service · Ranker (LSR), Service · Post Write, Worker · Fanout (denorm), Worker · Engagement Aggregator, Cache · Object Hydration, Cache · Ranked Cursor, Graph DB · TAO + MySQL, KV Store · Online Feature Store, Object Store · Media, Stream · New-Post Ingest, Worker · Outbox Poller, Queue · Fanout DLQ, Stream · Engagement Events — connected by 25 flows.MobilemobileCDN · Edge POPsAkamai + Meta FNAAPI Gateway · GraphQLEnvoy + GraphQL persist…Auth / IdentityOIDC + recency-tokened …Service · Feed ReadHack/HHVM (or modern eq…Service · Candidate G…Follow-graph retrieval …Service · Ranker (LSR)MTML model on MTIA / GP…Service · Post WriteRust + Tokio, idempoten…Worker · Fanout (deno…Kafka consumer group, R…Worker · Engagement A…Flink + RocksDB stateCache · Object Hydrat…Memcached + McRouter (T…Cache · Ranked CursorRedis cluster (cluster-…Graph DB · TAO + MySQLMySQL 8 (InnoDB) + TAO …KV Store · Online Fea…Manhattan-style RocksDB…Object Store · MediaHaystack (hot) + f4 (wa…Stream · New-Post Ing…Kafka 3.x (or Scribe-eq…Worker · Outbox PollerRust + Tokio cron-loop,…Queue · Fanout DLQKafka topic, manual rep…Stream · Engagement E…Kafka 3.x
19 components, 25 flows. A dashed line is an asynchronous hop. This is the reference design, not the only one that works.

What each component is for

Mobile

Issues GraphQL/HTTPS for feed-open and post-create. Maintains a client-side seen-set bloom filter so re-issuing a cursor after a network blip doesn't render duplicates.

Why it exists. Mobile is the product — web is a thin shadow. GraphQL persisted-query IDs keep raw queries off the wire (bandwidth + abuse-surface win); the opaque server-minted cursor keeps rank state authoritative. Worth the client-protocol discipline at 1.5B DAU.

When it fails. App crash mid-feed-open or socket drop loses the in-flight cursor — client retries with the persisted cursor; duplicates filtered by the local seen-set bloom. Detection: client telemetry on cursor_retry_rate; mitigation is the bloom plus server idempotent GET (rerunning the same cursor returns the same page within 15 min).

CDN · Edge POPsAkamai + Meta FNA

Terminates TLS at ~85 PoPs; serves signed media URLs from edge cache (30d immutable TTL); origin-shields the regional ingress for hot creators so a viral asset triggers ONE origin fetch shared across PoPs. Feed JSON passes through with Cache-Control: private, no-store — ranking is per-(user, ranker_version) and cannot be edge-cached.

Why it exists. Media is 95%+ of egress bytes; serving it from origin is the cost model breaker (~17 GB/s/region origin egress at peak). The Cornell SOSP study on Facebook photo caching reports 65.5% browser hit rate; the remaining 35% would otherwise hammer the backbone. We rejected origin-only serving — the math on backbone bandwidth alone justifies the edge tier ten times over.

When it fails. PoP failure → BGP withdraws the anycast prefix; clients route to next-nearest within ~3s, no client action. If origin shield misses on a viral asset, the origin fetch is shared across PoPs by request-coalescing — without it, a single PoP miss can fanout to thousands of upstream PoPs. Detection: edge origin_offload_rate per asset; alert if a single blob_id exceeds 100 origin GETs/min.

API Gateway · GraphQLEnvoy + GraphQL persisted queries

Resolves persisted-query IDs to internal RPC routes; enforces token-bucket rate limits per (user_id, surface) at 30 rps / device with burst 60; hedges ranker-bound calls at p95+10ms for tail-latency mitigation; carries the auth bearer to Auth/Identity for ACL enrichment before fanning out to services.

Why it exists. We need a tier that owns abuse + auth + rate-limit decisions outside the service mesh — services should never re-implement these. Considered putting rate limits at Envoy sidecars only, but per-user/per-surface limits need a fast distributed counter (Redis-backed) and that's gateway-level state, not sidecar-level. A single gateway tier also lets us A/B route surfaces at one chokepoint.

When it fails. Gateway tier saturates → rate-limit Redis becomes contended → tail latencies leak through. Detection: gateway 429_rate derivative + p99 latency on rate-limit lookup. Mitigation: shed read traffic first (post-create stays admitted), drop hedging budget to zero, fail-open on the rate limiter for 30s while the Redis pool catches up. Cert expiry: SPIFFE/SPIRE auto-rotation with 30/14/3-day alerts.

Auth / IdentityOIDC + recency-tokened ACL service

Validates the session bearer; resolves the user's effective ACL set (blocks, mutes, restricted accounts); attaches a recency_token = (region_lsn, mutation_ts) so downstream services can detect ACL mutations newer than what the regional follower has replicated. Rejects requests with stale recency tokens by routing them to the master region.

Why it exists. ACL mutations (block / unfollow) MUST be read-your-writes — eventual consistency here is a privacy SEV, not just stale UX. We rejected per-request leader reads (would 5x the master-region load); the recency-token pattern lets the common path read from the regional follower while only the post-mutation window pays the cross-region cost.

When it fails. Auth tier degraded → all surfaces return 401 (fail-closed). This is the right answer — better to serve nothing than serve someone else's feed. Detection: 401 rate vs baseline + auth p99. Mitigation: pin to last-known-good session (60s grace) for already-authenticated requests; drop new-session creation but keep refresh working. Master-region failover gap (the 30s MySQL RTO window): recency-tokened reads have NO master to route to. Behavior in that window is serve-stale-tagged-degraded: read from the highest-LSN regional follower and stamp the response degraded_acl=true; client surfaces a 'recently changed' banner on ACL-affected views. The privacy compromise is bounded to the 30s gap × the rate of ACL mutations within that window — defensible because the alternative (fail-closed for 30s on every recently-mutated session) is worse UX than a brief stale read with an explicit indicator. SLO: recency_token_master_route_rate and recency_token_degraded_serve_rate — both alert if non-zero outside known failover events.

Service · Feed ReadHack/HHVM (or modern equivalent: Rust + Tokio)

On feed-open: validate ACL recency, check the Cursor Cache for the requested page (continuation hits cached page; first page misses), call Candidate Generator → Ranker, hydrate post + author objects through the Object Hydration cache, emit cursor-with-next, return JSON. Owns the per-request deadline propagation: the 600ms feed-open p99 budget is sliced and pushed downstream as gRPC deadlines.

Why it exists. The feed surface is the fattest read in the product — it touches graph, posts, ranker, features, hydration. Splitting the orchestration into a dedicated service (vs putting it in the gateway) lets us scale independently (600 replicas at peak) and isolates the feed's blast radius from other surfaces that share the gateway. We rejected per-page Lambda-style serving — cold-starts at 350K QPS would be brutal.

When it fails. Adaptive concurrency starts shedding → users see chronological fallback feed (degraded flag in response). Worse: ranker loop with retries → cascade. Detection: degraded_response_rate per region + ranker call retries derivative. Mitigation: circuit-breaker on the ranker after 30s of degraded above 5%; the chronological fallback is the load-bearing safety net (ranker can be down for an hour and the product still works, just gets stale).

Service · Candidate Gen + ESRFollow-graph retrieval + lightweight gradient-boosted ESR

Retrieves candidates from two sources and merges: (1) follow-graph recent posts via TAO read of posts_by_author for friends, (2) precomputed denormalized feed-index for fanout-on-write recipients. Runs a lightweight Early-Stage Ranker (gradient-boosted, ~5ms/candidate) to drop from ~5000 retrieved to ~500 sent to the LSR.

Why it exists. Single-stage ranking 5000 candidates with the heavy MTML model would burn 15× the GPU. The retrieval → ESR → LSR funnel is the standard architecture that Pinterest, TikTok, and YouTube all use. We rejected pure pull-on-read (won't satisfy <600ms p99 once the candidate set exceeds 500) and pure fanout-on-write (the celebrity-followers explosion melts inbox-write throughput at >100K followers).

When it fails. Feed-index read stalls or a slow TAO shard → recall drops, feed degrades to 'most recent from your follow list only' (recall mode). Detection: candidate-set size derivative + recall-evaluation canary. Mitigation: fall back to TAO follow-graph posts only (the 'recall mode'); feed still works, just thinner. Hot creator partition: a viral post is in everyone's candidate pool, so the partition holding their posts_by_author list takes 10-50× median QPS — mitigated by promoting that key to in-process LRU at every Cand Gen replica.

Service · Ranker (LSR)MTML model on MTIA / GPU + ONNX runtime

Receives ~500 ESR-output candidates; fetches per-(user, candidate) features from the Online Feature Store; runs a multi-task multi-label model that predicts P(like), P(comment), P(share), P(dwell), P(skip); blends them via the current value-model weights into a final score; returns top-50 with calibrated scores. Per-call CPU budget: ~150ms wall (GPU inference ~30-50ms + feature fetch ~20ms + post-processing).

Why it exists. Ranking is the product. A chronological feed has 30% lower session length per Meta's own A/B disclosures; the entire candidate funnel exists to feed this model. We rejected putting the LSR inline with Cand Gen (different scaling profile — Cand Gen is CPU/RAM-bound on the ANN index, LSR is GPU-bound on the MTML model; merging would either over-provision Cand Gen or starve LSR).

When it fails. Model server GC pause or feature-fetch slowdown pushes p99 from 50ms to 600ms; Feed Read times out, retries amplify, queue depth blows up — classic ranker timeout cascade. Detection: ranker request_queue_depth derivative + retry_ratio>0.3 + downstream feature-store QPS spike with no organic traffic increase. Mitigation: adaptive concurrency at Feed Read clamps in-flight, retry budget capped at 10%, circuit-break to chronological fallback. The fallback is what stops the cascade; retries don't help — they hit the same slow model.

Service · Post WriteRust + Tokio, idempotent client_request_id key

Validates the post payload + media blob_id; HEADs the blob in Object Store · Media as a precondition to detect orphaned post-references (returns 422 if the upload didn't complete); writes the canonical post row PLUS an outbox event row to TAO/MySQL in a SINGLE SQL transaction (the load-bearing 2-participant outbox pattern); does NOT publish synchronously. Worker · Outbox Poller picks up the outbox row asynchronously and publishes post.created to the New-Post stream. Returns 201 with the canonical post_id once the SQL transaction commits.

Why it exists. Splitting writes into their own service tier keeps post-create out of the Feed Read fast path and lets us put creation-specific abuse rules (per-account daily limits, automated content checks, watermarking handoff) here without touching the read tier. We rejected gateway-side write handling (would mix abuse-policy responsibilities with edge concerns) and considered queueing all writes through Kafka first (durability win, but now post-create returns 202 'maybe' which the product team rejected — users want a hard 201/4xx).

When it fails. Primary shard fails mid-write → only in-flight (un-acked) writes are at risk; they return 503 and the client retries with the same client_request_id after failover (RTO 30s). Acknowledged commits are safe: semi-sync means an acked COMMIT is already on ≥1 same-region follower, so intra-region failover to the most-caught-up follower loses ~0 acked writes (RPO ≈ 0 — nonzero only in the degraded case where semi-sync timed out into async fallback). Detection: write success rate per shard + replication-lag-at-failover. The user-visible failure is 'post is publishing…' for the failover window — bounding acked-then-lost posts to ~zero is exactly why we run semi-sync rather than pure async. Orphaned blobs (client uploaded but never POSTed metadata) are GCed lazily by a Tectonic-side reference-counter sweeper at 7 days — accepted side effect of the presigned-upload flow.

Worker · Fanout (denorm)Kafka consumer group, Rust

Consumes the New-Post stream; reads the Followers table (hash(followee_id)) to enumerate a creator's followers on one shard; for each post by a creator with <100K followers (the cutoff), fans out into the recipient feed-index in TAO/MySQL — upserts (recipient_user_id, post_id, score_at_insert) rows keyed (recipient_user_id, post_id) into the recipient's feed-index partition (idempotent, so a re-published post.created from at-least-once outbox delivery is a no-op). Skips creators above the cutoff (those are pull-on-read at feed-open time). Maintains a 32-deep DLQ topic for poison messages.

Why it exists. Hybrid fanout is the canonical Twitter/Meta lesson: pure fanout-on-write melts on celebrities (615M follower posts = 615M writes, 3.4-hour tail at 50K writes/sec/shard); pure fanout-on-read makes feed-open a multi-shard scatter-gather. Cutoff at 100K is empirically tuned — above this the writeback cost beats read-amortization (see deepDive fanout-cutoff). We considered scaling all to fanout-on-read like Twitter pre-2014 — the resulting feed-open p99 doubled in a 2017 Twitter retrospective.

When it fails. Consumer group lag spikes (rebalance, slow downstream, partition skew) → new posts don't appear in followers' feeds for X minutes; a post's reach decays exponentially with time-since-post so a 10-min lag destroys most of a post's lifetime value. Detection: Kafka consumer_lag per partition + post_published_to_visible_in_feed p99 SLO. Mitigation: per-partition autoscaling; shadow chronological reads from posts-DB as fallback when lag > 60s; priority-lane topic for celebrity posts so they never queue behind the long-tail.

Worker · Engagement AggregatorFlink + RocksDB state

Consumes the Engagement stream of passive ranking signals (dwell, skip, impressions) emitted by Feed Read; updates the Online Feature Store with the derived features (dwell-time per (user, content_type), skip-rate, impression recency) at sub-second staleness.

Why it exists. Online features are the difference between 'this looks like an Instagram feed' and 'this looks like a 2014 Instagram feed' — recency of passive-behavior signals is what makes a feed feel alive. We rejected synchronous feature updates inline on the read path (would couple Feed Read latency to feature writes) and rejected purely batch features (the resulting 24h staleness ages out the ranking loop within hours).

When it fails. Aggregator lag → feature-store features are stale → ranker scores with 1-hour-old passive signals → recommendations look stale to users. Silent failure — no errors, just bad ranking. Detection: per-feature freshness_seconds p99 + ranker input distribution drift detector. Mitigation: ranker-side fallback to default values when a feature exceeds its staleness budget; dual-write online + nearline path so a single pipeline failure doesn't blank features.

Cache · Object HydrationMemcached + McRouter (TAO follower role)

Per-region read-serving cache for post objects, author profile objects, and ACL flags. Read pattern: get_multi([post_id, author_id, ...]) for the 50 hydration items per feed page; on miss, McRouter issues a lease token and the missing keys fall through to the Graph DB · TAO+MySQL replica. Concurrent misses on the same key wait on the lease holder rather than stampeding the DB — the canonical Meta NSDI '13 lease mechanism. The hydrated post object carries denormalized like_count / comment_count fields (materialized onto the post at write / fanout time), so Feed Read renders engagement counts straight from hydration without a separate counter read.

Why it exists. Without this layer, every feed-open does 50 round-trips to MySQL — at 350K feed-opens/sec peak, that's 17.5M MySQL reads/sec, ~10x what sharded MySQL can serve. The TAO paper reports 96.4% follower hit rate at >1B reads/sec; we sized for the same. We rejected client-side caching at Feed Read replicas only — the working set (~22.5 TB hot) doesn't fit in 600 service replica RAMs without sharding, and a shared shard-aware cache is exactly what memcache+mcrouter is.

When it fails. Cluster restart → cold cache → every read misses to MySQL simultaneously (canonical thundering herd) — without leases, MySQL gets crushed for 2-10 min until cache warms; with leases, concurrent missers wait on the first miss-holder so MySQL sees ~1 read per key, not N. Detection: lease_hold_time_p99 + miss-rate cliff + downstream MySQL QPS. Mitigation: gradual restart (1% pool at a time, never all at once); leases ARE the mitigation — they're load-bearing, not optional. Single-shard loss: McRouter routes the shard's keys to a backup pool (the lease-holding miss falls through to MySQL, then writes back to the backup); cache state is not durable, so 'failover' here means 'reroute requests' not 'promote a replica'. Cache poisoning of celebrity object: 10% bypass-cache canary on hydration to detect drift between cache and DB.

Cache · Ranked CursorRedis cluster (cluster-mode, 16384 slots)

Stores the ranked page for (user_id, session_id, ranker_version) so continuation pages don't re-rank — the ranker fires once per first-page; pages 2..N read from this cache. Each entry: an ordered list of ~50 post_ids + ranker scores + candidate_pool_snapshot_id + seen_set_bloom. Entry size ~2 KB; cardinality ~430M live entries at peak (50M concurrent sessions × 4.3 pages × 2 ranker versions live).

Why it exists. Without this, every continuation-page hits the ranker → 4.3x ranker QPS → 4.3x GPU spend. The continuation cache turns 350K feed-opens/sec into 81K rank/sec — a 4.3x cost reduction on the most expensive tier. We rejected putting cursors in client local storage (clients lie / can be replayed; rank state must be server-authoritative); rejected materializing in MySQL (~860 GB working set at 2 KB/entry needs in-memory).

When it fails. Cluster slot migration during a resharding event → cursor lookup fails for keys mid-migration → user sees 'page reset' to first-page. Detection: cursor_lookup_miss_rate derivative during shard-ops. Mitigation: page-1 fall-through always works (it just calls the ranker); the failure mode is degraded user experience (lost scroll position), not broken — acceptable. Ranker version bump: cursors carrying old ranker_version are honored for an overlap window of 1h, then return cursor_invalidated=true and client re-fetches page 1.

Graph DB · TAO + MySQLMySQL 8 (InnoDB) + TAO graph cache, ~256 logical shards

Canonical store for the social graph (follows / blocks / mutes) and the posts table (post_id, author_id, blob_id, created_at, content_metadata). Sharded by hash(user_id) for graph and hash(post_id) for posts; a feed-open hits the recipient's graph shard for the follow set and N post shards for hydration. Single primary per shard in the master region; semi-sync replication to two followers in the same region + async replication to follower regions.

Why it exists. We need transactional uniqueness on (user_id, follower_id) graph edges and on (post_id) — without it, double-follow or duplicate post_id corrupts the feed. Sharded MySQL is the canonical Meta+Twitter+Pinterest+Discord choice for this profile. We rejected DynamoDB (would work, but loses SQL escape-hatches when product asks for 'top 100 takedown candidates' / abuse queries), rejected Spanner (write latency 5x higher cross-region — feed-write SLO breaks).

When it fails. Primary shard failure mid-write → only in-flight un-acked writes at risk (acked commits are already on a semi-sync follower, RPO ≈ 0); clients retry with client_request_id after 30s RTO. Hot shard on viral creator (everyone's pulling their posts_by_author list at once) → mitigate by promoting hot keys to in-process LRU at every Cand Gen replica + read replica scaling. Cross-region replication lag spike → privacy SEV (block/unfollow takes 30s to land in EU); detection: tao_follower_leader_replication_lag_p99 > 5s; mitigation: recency-tokened reads route to master region during the lag window.

KV Store · Online Feature StoreManhattan-style RocksDB-backed KV (FBLearner Feature Store role)

Low-latency point lookup of features keyed by (user_id, feature_name) and (post_id, feature_name). Ranker fetches ~20 features per (user, candidate) pair × 500 candidates = 10K feature-fetches per rank call; budget 20ms total → 2µs per fetch (cluster mget). Updated by Worker · Engagement Aggregator with sub-second staleness for dwell/skip velocity features; offline graph features refreshed daily.

Why it exists. Ranker model serving requires features at inference time; features must be (a) extremely low latency (microseconds, not milliseconds — there are 10K per request), (b) consistent enough that the same (user, candidate) pair scores the same within a session. Storing features in MySQL would 100x feature-fetch latency; storing in memcached would lose the per-feature TTL granularity. Manhattan-style KV (RocksDB-backed, RF=3) is the canonical Twitter+Meta choice for this profile.

When it fails. Feature staleness spike → ranker scores with stale features → recommendation quality drops (silent failure, no errors). Detection: per-feature freshness SLO + ranker input distribution drift. Mitigation: ranker-side fallback to default feature values when stale beyond budget. Hot feature key (e.g. a celebrity's profile features) → mitigated by the standard Manhattan client-side cache + read replica scaling.

Object Store · MediaHaystack (hot) + f4 (warm) + Tectonic (cold)

Three-tier media storage. Hot (first hours/days, sub-1-read/day after that): Haystack — needle-in-haystack append-only volumes, RF=3 within DC, in-memory index, p99 read 5ms. Warm (>1 month, low read rate): f4 — Reed-Solomon(10,4) intra-DC + XOR cross-DC, effective replication factor 2.1× (saved 87 PB on a 65 PB corpus per the OSDI '14 paper). Cold + analytics: Tectonic — exabyte multitenant FS, the largest cluster is ~1,590 PB / 4,000 nodes / 10B files.

Why it exists. Storing all media in Haystack RF=3 forever costs 3.6× the raw byte count; at exabyte scale that's a procurement-budget catastrophe. The tiered model is what makes the unit economics work. We rejected pure S3 (third-party cost at this scale is unfavorable; cross-region semantics don't match Meta's regional-master pattern), rejected pure Tectonic (read latency too high for hot blobs).

When it fails. Haystack volume failure → blob unreadable until volume restore from RF=3 peers (RTO seconds via in-memory index re-replication). f4 single-DC loss → cross-DC XOR reconstructs (RTO minutes, no data loss). Tectonic cluster degraded → cold reads slow but never fail (RPO 0). The actual scary failure: a blob_id referenced from a post_id but the blob itself was lost mid-tier-migration. Detection: blob_404_rate_per_post_age_bucket (a 404 on a >7-day-old blob is a tier-migration loss; on a <1-day-old blob is a Haystack incident). Mitigation: clients receive a 404 from CDN, fall back to a 'media unavailable' placeholder; the post itself stays in feed (textual content preserved); a separate reconciliation job re-uploads from device-side cached copy if available. Reference counting is a precondition (prevents premature GC), not the recovery mechanism.

Stream · New-Post IngestKafka 3.x (or Scribe-equivalent)

Carries post.created events from Post Write to Worker · Fanout. One event per post (not per recipient). Partitioned by hash(author_user_id) so all posts by the same creator land on the same partition (consumer affinity for fanout-batching). Retention 72h (long enough for a consumer outage + replay).

Why it exists. The fanout-on-write denormalization can't run in the synchronous post-create path without blowing the post-create p99 budget — a 50K-follower post writing 50K rows takes 1-10s. Decoupling via Kafka turns post-create into a 200ms operation and fanout into a background pipeline. We rejected RabbitMQ (per-message ack overhead is wrong for this throughput profile) and SQS (regional, lock-in, retention too short for replay).

When it fails. Broker outage with min.insync.replicas=2 → the Outbox Poller's publish fails and retries; post-create is unaffected and still returns 201 (the post is already durable in MySQL — the publish is off the critical path by design, so a broker outage delays fanout, it cannot 503 a create). Detection: producer error rate + outbox_publish_lag_p99. Mitigation: the outbox poller drains from MySQL once the broker recovers; fanout is delayed but no post is lost. Consumer-side outage: Worker · Fanout falls behind → see Worker · Fanout failure mode.

Worker · Outbox PollerRust + Tokio cron-loop, 100ms tick

Polls the outbox_events table on each TAO/MySQL shard at 100ms intervals; SELECTs rows where published_at IS NULL ORDER BY id LIMIT 256; publishes each as a post.created event to the New-Post stream; UPDATEs published_at once Kafka acks. The single load-bearing component for 'posts MUST not be lost; events emit iff post commits' — without it, the outbox pattern is just an unused MySQL table.

Why it exists. We rejected having Post Write publish synchronously to Kafka (would put a Kafka-ack in the post-create critical path, blowing the 900ms p99 SLO when a broker has a transient issue) and rejected dual-writes (would lose the iff guarantee on Kafka-failure-after-MySQL-commit). The poller decouples durability from publish, at the cost of a separate component to operate.

When it fails. Poller dies for a shard → outbox table grows unbounded; new posts commit but events never publish → followers don't see them via the fanout pipeline (the post is still durable, just invisible until the poller catches up). Detection: outbox_publish_lag_p99 per shard (target <500ms); outbox_table_row_count per shard derivative. Mitigation: etcd lease auto-fails-over to a peer poller within 5s; the outbox table itself is the durability mechanism — even multi-hour poller outage is recoverable, just user-visibly stale.

Queue · Fanout DLQKafka topic, manual replay tool

Receives messages from Worker · Fanout that exhausted 32 retries with exponential backoff (poison messages — schema mismatch, downstream-permanently-broken, etc.). Holds 7-day retention so on-call has time to investigate; a manual dlq-replay ops tool re-publishes selected messages back to the New-Post stream after the underlying issue is fixed.

Why it exists. Without a DLQ, poison messages either block the partition (if stop-on-error) or get silently dropped (if skip-on-error). Both are bad: blocked partition = ALL posts on that partition can't fanout; silent drop = the affected posts never appear in followers' feeds with no signal. The DLQ trades 'delayed fanout for poison messages' for 'no other posts blocked + audit trail'. We rejected sending poison messages back to a 'parking' topic on the same cluster (capacity contagion) and rejected unbounded retry (worker-thread starvation).

When it fails. DLQ depth grows unbounded → indicates a sustained poison-message source (schema drift, codec-incompat) → on-call investigates. Detection: dlq_depth derivative per partition; alert on >100 messages in 1h. Mitigation: replay tool after fix; in extreme cases (schema break) the DLQ messages are abandoned with an audit trail.

Stream · Engagement EventsKafka 3.x

Carries passive ranking signals (dwell, skip, impressions) emitted by Feed Read to Worker · Engagement Aggregator. Partitioned by hash(post_id) so all signals on one post land on one partition (consumer can locally aggregate without cross-partition reads). Retention 24h.

Why it exists. Passive signals feed the online-feature update path. Without the stream, the aggregator would tail Feed Read's request path directly → fragile, fan-in-on-service. The stream is the event-sourced spine for ranking signals — it lets new downstream consumers come online without coordinating with Feed Read.

When it fails. Consumer lag → online features go stale → ranker scores with old features → recommendations look stale (silent failure). Detection: per-partition consumer lag + per-feature freshness SLO at the feature store. Mitigation: dual-write online + nearline path so a single-pipeline failure doesn't blank features; ranker-side default-value fallback when features exceed staleness budget.

Stage by stage

The same 10 stages the workspace walks, answered.

01Clarifications

What would you ask before drawing a single box?

Typical clarifications to surface:

  • Chronological or ranked? This canonical assumes ranked. A pure chronological feed deletes the entire ranker tier (Cand Gen + ESR + LSR + Feature Store + Aggregator) — call it for the interviewer.
  • Pull, push, or hybrid? Hybrid (the load-bearing answer). Pure push (fanout-on-write) melts on celebrities; pure pull (fanout-on-read) makes feed-open a multi-shard scatter-gather. The fanout-on-write cutoff is empirically tuned at ~100K followers.
  • Pagination? Cursor-based, never OFFSET. The cursor encodes (ranker_version, candidate_pool_snapshot_id, last_score, last_post_id, seen_set_bloom) and is opaque to the client.
  • Read-after-write for ACL mutations (block, unfollow)? YES — privacy-load-bearing. Eventual consistency here is a SEV. Recency-tokened reads route to master region for the post-mutation window.
  • Latency targets? Feed-open p99 < 600ms (user tolerates first-paint up to ~1s, headroom for client render). Post-create p99 < 900ms (asymmetric — user accepts spinner on upload).
  • Multi-region? Active-active for reads, single-master per shard for writes. Recency-tokened reads handle the cross-region replication-lag privacy SEV.

Assumptions baked in:

  • 1.5B DAU, 11 feed page fetches/DAU/day, 4.3 pages per open (incl. first).
  • 195M posts/day, 7 likes per post.
  • 100K-follower fanout cutoff.
  • 30% Pareto: top 10% of users serve ~70% of reads.

02Functional reqs

What must this system actually do?

  • Open feed: Returns the first ranked page of ~12 posts + a cursor. Identity-validated, ACL-filtered, ranked.
  • Continue feed: Given a cursor, return the next ranked page WITHOUT re-ranking (cursor cache hit).
  • Create post: Upload media + metadata; durable write to TAO/MySQL; emit fanout event; return 201 with canonical post_id within p99 900ms.
  • Block / unfollow: ACL mutation; recency-tokened reads enforce read-your-writes within 60s of the mutation.
  • Profile pulls: Out of scope for this canonical (separate surface) but uses the same TAO + Object Hydration cache backend.

03Non-functional

What must it promise about speed, uptime and correctness?

  • Availability: 99.95% on the read path (feed-opens). The chronological fallback is the load-bearing mitigation when the ranker is degraded — feed still works, just gets stale.
  • Latency:
  • Feed-open p50 180ms · p99 600ms · p99.9 1.2s.
  • Feed-rank-continuation p50 40ms · p99 150ms (cache-served, no model).
  • Post-create p50 250ms · p99 900ms · p99.9 2s.
  • Durability: Posts MUST not be lost. Intra-region MySQL semi-sync gives RPO ≈ 0 for acknowledged commits (an acked COMMIT already sits on ≥1 same-region follower, so failover to the most-caught-up one loses ~0 acked writes) — nonzero only if semi-sync has timed out and fallen back to async. Cross-region async replicas carry RPO ≈ 1-5s (the region-disaster window). Passive-signal events tolerate loss (they only feed feature freshness; the event stream is for derivative consumers only).
  • Consistency:
  • ACL mutations: read-your-writes within 60s (recency-tokened reads route to master region).
  • Post creation: read-your-writes for the creator (sticky-session reads from the master shard for 60s post-create).
  • Engagement counts: denormalized onto the post object at write/fanout time; eventual, acceptable for like/comment-count display.
  • Feature staleness: 60s for passive-signal features, 24h for graph-derived features. Ranker is OK with stale features (default-value fallback when budget exceeded).
  • Security: Signed media URLs (HMAC-SHA256, expiry-bounded); persisted GraphQL queries (no raw queries in the wire); per-user/per-surface rate limits at the gateway; mTLS service mesh; SPIFFE/SPIRE auto cert rotation.
  • Scalability: Horizontal at every tier. The 100K-follower fanout cutoff is the single empirically-tuned product knob.

04Capacity estimation

How much load and data does this have to hold?

**Inputs (defended in the deepDive fanout-cutoff and the Capacity walk-through section below):**

  • DAU = 1.5B; feed page fetches/DAU/day = 11; pages/open (incl. first) = 4.3.
  • Posts/day = 195M.
  • Peak/avg multiplier = 1.8 (global service, but still peaks).

Derived (formulas in capacityDerivations):

  • Feed page-fetch peak QPS = 1.5e9 × 11 / 86400 × 1.8 ≈ 343K QPS.
  • Rank peak QPS = page fetches / pages-per-open = 343K / 4.3 ≈ 80K QPS (only first-page hits the ranker; continuation pages serve from Cursor Cache).
  • Post-create peak QPS = 195M / 86400 × 1.8 ≈ 4.1K QPS.
  • Ranker hosts at peak = 80K rank QPS / 149 useful rps/host ≈ 537 hosts at steady state, +537 transient during the 30-min ranker rollout window (both v42 + v43 live to honor in-flight cursors) → ~1100 host peak fleet during deploys. One feed-open ranks ~500 candidates × 100µs ≈ 50ms inference + 20ms feature fetch + 80ms post-processing → 149 useful rps/host at 70% util.
  • Feature-fetch QPS = rank QPS × 10K mget-fanout = 800M feature-fetches/sec peak (the reason the Feature Store is keyed on RocksDB-on-NVMe at 2µs p99, not memcached).
  • Storage growth = 195M × 365 × 7.2 KB metadata-and-indexes ≈ 513 TB/year post metadata; with 3× MySQL replication ≈ 1.5 PB/year. Media itself is ~50× larger but lives in Haystack/f4/Tectonic (separately addressed; tiered to lower-cost storage by access pattern).
  • Object Hydration working set = top-10%-users (150M) × 50 hot posts × 3 KB = 22.5 TB hot working set. Cluster sized at ~750 nodes × 40 GB usable each = 30 TB (33% headroom) — the load-bearing reason: hot-key promoted-to-LRU sets and ranker-overlap doubling can push the effective working set ~25% above baseline; previous 24 TB sizing had only 6.25% headroom and would page on rollouts.
  • Cursor Cache cardinality = 50M concurrent sessions × 4.3 pages ≈ 215M live cursor entries at steady state, doubling to ~430M during ranker-overlap windows (two ranker versions live). At 2 KB/entry = 430 GB steady → 860 GB peak. Sized at 96 shards × 13 GB usable = 1.25 TB primary + followers — covers the rollout-doubling with ~45% headroom, no eviction storms; previous 64-shard design had zero headroom and would have caused 'page reset' UX during every model rollout.
  • Edge egress = 343K QPS × 50 items × 3 KB = 51 GB/s metadata globally → ~17 GB/s per region (3 regions). Media at edge: 95%+ CDN hit rate; origin egress is ~5% of total.

Tier-flip table (the load-bearing claim — when each architectural choice flips):

  • Below ~10K feed-opens/sec: skip the per-region edge object cache; central memcached suffices.
  • Above ~100K followers/creator: switch to fanout-on-read for that creator (the cutoff in this canonical).
  • Above ~1M QPS on a single key: replicate across N=8 shards with client-side jitter.
  • Above ~500 ranker hosts: introduce two-stage (cheap retrieval → expensive rerank) — already done here.
  • Above ~1 PB/year metadata growth: tiered storage (hot SSD 90d → cold object) — already done via Haystack/f4/Tectonic.

05API design

What does the outside world call, and what comes back?

POST /graphql
Content-Type: application/json
Authorization: Bearer <session_jwt>

{
  "id": "FeedOpenQuery_v42",        // persisted query ID — no raw query on the wire
  "vars": {
    "cursor": null,                  // null = first page; opaque base64 from prior page = continuation
    "limit": 12
  }
}

200 OK
{
  "data": {
    "feed": {
      "items": [
        { "post_id": "p_xx", "author": {...}, "media_url": "https://i.cdninstagram.com/...?sig=...&exp=...", "score": 0.87, "counts": { "likes": 18234, "comments": 45 } },
        ...
      ],
      "cursor": "eyJyYW5rZXJfdmVyc2lvbiI6InY0MiIsInNuYXBzaG90X2lkIjoiYWFiYiI...",
      "ranker_version": "v42",
      "degraded": false              // true when ranker is in chronological-fallback mode
    }
  }
}

Post-create is a two-step upload: reserve a blob + presigned URL, PUT the bytes direct to the object store, then create the post referencing the returned blob_id.

POST /graphql                          # step 1: reserve upload (POST /media/uploads)
{
  "id": "CreateUploadMutation_v1",
  "vars": { "content_type": "image/jpeg", "content_length": 2280000 }
}

200 {
  "blob_id": "haystack:bx_yy",
  "presigned_put_url": "https://upload.cdninstagram.com/...?sig=...&exp=...",
  "expiry": "2026-06-29T12:30:00Z"   // ~15-min window; orphaned blobs GCed at 7 days
}
PUT https://upload.cdninstagram.com/...?sig=...&exp=...   # step 2: client streams bytes direct to Object Store
Content-Type: image/jpeg
Content-Length: 2280000

<binary media bytes>

200 OK                                 # bytes now land in Object Store · Media (Haystack hot tier)
POST /graphql                          # step 3: POST /post
{
  "id": "CreatePostMutation_v3",
  "vars": {
    "client_request_id": "<uuid>",     // idempotency key; server unique-indexes on this
    "blob_id": "haystack:bx_yy",        // from step 1; Post Write HEADs it as a precondition (422 if PUT never completed)
    "caption": "...",
    "audience": "PUBLIC"
  }
}

201 { "post_id": "p_xx", "created_at": "..." }

06Data model

What gets stored, and what is it looked up by?

Posts table (sharded by hash(post_id)):

fieldtypenotes
post_idbigintPK; Snowflake-style ID
author_user_idbigintindexed for posts_by_author listings
client_request_iduuidUNIQUE — idempotency key
blob_idvarcharHaystack/f4/Tectonic blob reference
captiontext
audienceenumPUBLIC / FRIENDS / RESTRICTED
created_attimestamp
like_countbigintdenormalized; materialized at write/fanout time
comment_countbigintdenormalized; served straight off hydration
disabledbooltakedowns; cheaper than delete

Graph table (sharded by hash(user_id)):

| user_id | followee_id | created_at | flags (mute / restrict / close-friend) |

Answers "whom does U follow" on a single shard — the read-path direction Cand Gen needs.

Followers table (inverse assoc, sharded by hash(followee_id)):

| followee_id | follower_id | created_at |

Answers "who follows creator X" on a single shard — the direction Worker · Fanout needs to enumerate recipients. Dual-written with the Graph table on follow/unfollow (one single-shard write per direction, TAO assoc-pair style — not a cross-shard 2PC). Without it, "all followers of X" would scatter-gather every graph shard on every post fanned out.

Feed-index table (sharded by hash(recipient_user_id)):

| recipient_user_id | post_id | author_user_id | score_at_insert | inserted_at |

Written by Worker · Fanout for non-celebrity creators; pulled by Cand Gen at feed-open.

Database choice — recommended:

  • TAO + sharded MySQL for graph + posts: relational uniqueness on (user_id, follower_id) and (post_id) is load-bearing; mature operational tooling; transactional creation; 1.5 PB/year fits comfortably across 256 shards × 3 replicas.
  • Memcached + McRouter for object hydration: leases are the canonical anti-stampede mechanism; nothing else has them.
  • Redis cluster for cursor cache: cluster-mode hash slots match the (user_id, session_id) sharding need; sub-ms latency.
  • Haystack/f4/Tectonic for media: tiered by access pattern is the only economic answer at exabyte scale.
  • Kafka for both streams: per-partition ordering matches the partition keys (author_user_id for new-post, post_id for engagement).
  • Manhattan-style RocksDB KV for online features: NVMe-backed, RF=3, microsecond p99 — the only way to serve 800M feature-fetches/sec without melting budgets.

Why not DynamoDB for the canonical store? Fine on AWS at 1/10 this scale. At Meta scale: cost (DynamoDB at 350K QPS × month is multi-million-dollar/month per region), no SQL escape hatches when product asks for "show me top 100 abuse-flagged posts this week", and the cross-region semantics don't match Meta's master-per-shard model. Rejected.

Why not Spanner / CockroachDB? Strong cross-region transactions are ~5x the write latency. Post-create p99 budget breaks. The hybrid recency-token approach gives the same privacy guarantee at a fraction of the cost. Rejected.

07High-level design

Which components handle a request, and in what order?

Architecture summary in the order a feed-open traverses:

  1. Mobile issues a persisted GraphQL query over HTTPS.
  2. CDN · Edge POPs terminate TLS; pass feed JSON through (TTL=0, ranking is per-user); serve signed media URLs from edge cache (30d TTL; HMAC-SHA256 signature).
  3. API Gateway resolves the persisted-query ID, enforces per-(user, surface) rate limit, hedges ranker-bound calls at p95+10ms.
  4. Auth / Identity validates the session bearer, attaches a recency_token=(region_lsn, mutation_ts) so downstream services can detect ACL mutations the regional follower hasn't replicated yet.
  5. Service · Feed Read orchestrates: check Cursor Cache for continuation → if miss, call Cand Gen → Ranker → hydrate via Object Hydration cache → mint cursor → return.
  6. Service · Candidate Gen + ESR retrieves from the follow graph via TAO read plus the fanout-on-write feed-index via TAO read; runs gradient-boosted Early-Stage Ranker to drop from ~5000 to ~500.
  7. Service · Ranker (LSR) fetches features from Online Feature Store, runs MTML model on MTIA/H100, returns top-50 with calibrated scores.
  8. Cache · Object Hydration serves post + author objects with 96.4% hit rate; concurrent misses serialize on a lease token (NSDI '13 mechanism).
  9. Graph DB · TAO + MySQL is the canonical store for the social graph + posts; 256 shards; intra-region semi-sync (RPO ≈ 0 for acked commits) with async cross-region replication (RPO ≈ 1-5s). The graph cache (TAO) sits in front of MySQL — graph-shaped reads (follows / blocks / mutes) and write-through edge updates flow through it, MySQL is the durable backend.
  10. Cache · Ranked Cursor stores the per-(user, session, ranker_version) ranked page so continuation doesn't re-rank.
  11. KV Store · Online Feature Store serves microsecond per-feature lookups for ranker; updated by Engagement Aggregator with <1s staleness.
  12. Object Store · Media holds the actual blobs in three tiers (Haystack hot, f4 warm, Tectonic cold).
  13. Worker · Outbox Poller (100ms tick, per-shard etcd lease) bridges MySQL durability to Kafka publish — the load-bearing component for "posts MUST not be lost; events emit iff post commits."
  14. Stream · New-Post Ingest carries post.created events; Worker · Fanout consumes and denormalizes into recipient feed indices for non-celebrity creators. Poison messages flow to Queue · Fanout DLQ after 32 retries.
  15. Stream · Engagement Events carries passive ranking signals (dwell / skip / impressions) emitted by Feed Read; Worker · Engagement Aggregator consumes them to update online features for the ranker.

Read flow (first page): Mobile → CDN (passthrough) → Gateway → Auth → Feed Read → Cand Gen → (TAO read for graph) → Ranker → Feature Store → Object Hydration cache → (TAO miss) → MySQL replica → cursor mint → response.

Read flow (continuation page): Mobile → CDN → Gateway → Auth → Feed Read → Cursor Cache (hit) → Object Hydration → response. No ranker.

Write flow (post-create): Mobile → reserve upload (blob_id + presigned PUT URL) → Mobile PUTs bytes direct to Object Store · Media (out-of-band) → CDN → Gateway → Auth → Post Write → HEAD blob_id (Object Store, e16) → MySQL: BEGIN; INSERT post; INSERT outbox_event; COMMIT → 201. The outbox row is what guarantees the publish-iff-commit invariant — no synchronous Kafka in the post-create critical path. (Graph note: the direct Mobile → Object Store · Media presigned-PUT hop is backed by the client → media write edge e30, mirroring how the HEAD precondition is modeled as e16.)

Async outbox publish: Worker · Outbox Poller (100ms tick, per-shard etcd lease) → SELECT pending outbox rows → publish post.created to New-Post stream → UPDATE outbox.published_at. SLO: outbox_publish_lag_p99 < 500ms.

Async fanout: New-Post stream → Worker · Fanout → MySQL feed_index (recipient denorm) + Object Hydration warm. Poison messages → Queue · Fanout DLQ after 32 retries; manual replay tool.

Passive-signal aggregation: Feed Read emits dwell/skip/impression → Engagement stream → Worker · Engagement Aggregator → Online Feature Store. This closes the ranking feedback loop — the ranker reads these features on the next feed-open.

08Deep dives

Which part breaks first, and what do you do about it?

1. Cursor design — what's actually in the token? (ranker_version, candidate_pool_snapshot_id, last_score, last_post_id, seen_set_bloom). The snapshot_id pins the candidate set so new posts arriving mid-scroll don't cause dupes/skips (they surface separately as an "above-the-fold" refresh, not into in-flight pagination). The seen_set_bloom filter is what lets the server detect if a continuation request is from a stale client retry. On ranker_version mismatch, the server can either (a) reject and force page-1 reset (jarring UX), (b) honor with overlap-window logic for 1h, or (c) extend the cursor schema. Production answer: (b) — overlap-window is the load-bearing trick; the cost is keeping two ranker versions live for an hour. Why not OFFSET: at depth 1000 on a million-row scan, MySQL OFFSET 1000 LIMIT 12 is 17× slower than WHERE post_id < last_post_id LIMIT 12 (per the standard cursor-pagination analysis). Doesn't scale.

2. Hybrid fanout cutoff — why 100K? Pure fanout-on-write fails when creator_followers × write_QPS > recipient_inbox_write_budget. A 615M-follower post at 50K writes/sec/shard takes 3.4 hours to fanout. Pure fanout-on-read makes feed-open a multi-shard scatter-gather. The cutoff at 100K is empirically tuned: above this, fanout-write cost (per follower row) exceeds amortized read-savings on the recipient side. Hysteresis: a creator who oscillates around 100K followers gets bucket-tagged (with a 20% deadband — promote to celebrity tier at 120K, demote back at 80K) so fanout doesn't churn.

3. The lease mechanism — why it's load-bearing. Memcached without leases: 100 concurrent feed-opens that miss on a celebrity profile object → 100 simultaneous TAO/MySQL reads on the same key → DB stampede. With leases (NSDI '13): first miss gets a 64-bit token; concurrent missers are told to wait or get the stale value; first miss writes back; lease invalidates. Result: one DB read instead of 100. This is the load-bearing mechanism that lets Object Hydration cluster restarts not melt MySQL. Without it, gradual cache restart (1% pool at a time) is mandatory — with it, even simultaneous restart is recoverable in minutes.

4. The chronological fallback — why it must exist. The ranker is a complex GPU-bound service with retries that don't help (retries hit the same slow model). A bad model deploy or feature-store regression can cascade into a 100% feed-open failure. The chronological fallback (return recent unranked posts from the follow set) is the load-bearing safety net — it's the difference between a SEV-2 (degraded recommendations) and a SEV-1 (feed totally broken). Tagged in response with degraded: true so the client can show a banner.

5. ACL mutation read-your-writes via recency tokens. Block / unfollow MUST be RYW within 60s — eventual consistency is a privacy SEV. Naïve approach: route all reads to master region. Cost: 5× master-region load. Recency-token approach: Auth attaches (region_lsn, mutation_ts) to the response of any ACL mutation; client carries it; on subsequent reads within 60s, the server checks if the regional follower has replicated past region_lsn — if yes, serve from follower; if no, route this single read to master. Worst case 60s of single-user master reads, not 5× steady-state.

6. Ranker timeout cascade — why retries make it worse. Ranker p99 spikes from 50ms to 600ms (model GC pause, feature-store slowdown). Feed Read has 250ms timeout + 2 retries → retries hit the already-saturated workers, queue depth blows up, healthy replicas get marked unhealthy by LB, gRPC channel exhaustion cascades. Mitigation: adaptive concurrency limits (Netflix-style — don't admit new requests when queue depth is rising), per-request deadlines propagated end-to-end (no-retry on deadline-bust), retry budget capped at 10% in-flight, circuit-breaker to chronological fallback. The fallback is what stops the cascade; retries don't help — they hit the same slow model.

7. Cache poisoning of a celebrity object. A write to the User object (avatar URL, is_blocked) lands in Memcache via async invalidation, but the invalidation message is dropped or arrives before the DB commit replicates. Every follower fetching that profile via the hydration step gets the stale value for the cache TTL window. ACL fields use shorter TTL (30s vs 5min for non-privacy fields), 10% bypass-cache canary on the hydration path comparing replica DB vs cache to detect drift, version-stamped invalidations via lease-set. The blast radius for ACL poisoning is high: a 100M-follower profile = 100M wrong reads.

09Trade-offs

What did this design cost, and what breaks at 10×?

What breaks at 10× scale (~3.5M feed page fetches/sec, 41K post-creates/sec):

  • Ranker hosts: 537 hosts → 5,370 hosts (steady state); ~1,100 → ~11,000 during deploy overlap. Need to introduce a third stage (super-cheap pre-retrieval) or aggressive caching of (user, time-bucket) rank results.
  • Object Hydration working set: 22.5 TB → 225 TB. The 750-node cluster can't hold it — capacity must grow ~10× (≈7,500 nodes at 40 GB usable, or fewer larger-RAM nodes); a single-region miss-rate spike also becomes existential for MySQL — need per-shard memcached pools and tighter hot-key promotion.
  • Cursor Cache: 430M peak entries → 4.3B; doesn't fit in the 1.25 TB cluster. Needs eviction by session-LRU and possibly a tier of cold cursor recovery from Tectonic.
  • Feature-fetch QPS: 800M/sec → 8B/sec. The Feature Store needs sharding + read replicas in a way that respects feature-co-location for ranker batched mget.

Per-component failure stories:

  • CDN PoP failure — BGP withdraws; clients route nearest. Origin shield prevents thundering-herd amplification.
  • Object Hydration cluster restart — leases prevent DB stampede; gradual restart (1% pool at a time) is the operational pattern.
  • Cursor Cache slot migration — keys mid-migration miss; client falls back to page-1 (UX impact: lost scroll position).
  • MySQL primary failure — Patroni-equivalent failover, RTO 30s, RPO ≈ 0 for acked commits (semi-sync intra-region; nonzero only under async fallback). Idempotency on client_request_id makes client retries safe.
  • Ranker model regression — silent failure (no errors, just bad recommendations). Counterfactual A/B metrics + automated rollback on guardrail breach. The 30-min overlap window keeps cursors valid through the rollback.
  • Worker · Fanout backlog — new posts don't appear in followers' feeds. Per-partition autoscaling + priority lane for verified creators + shadow chronological reads from posts-DB as fallback.
  • Cross-region replication lag — privacy SEV on ACL reads. Recency-tokened reads route to master region for the post-mutation window.
  • Feature-store staleness — silent failure. Per-feature freshness SLO + ranker default-value fallback when budget exceeded.
  • Region failover — pre-flight readiness check (replica lag, model-artifact present, cache warm-ratio); continuous traffic shadowing to secondary so the cold-failover path doesn't surprise.

Deliberately scoped out — the engagement write path. This canonical is the read/ranking spine; the like/comment/save/share write side (a dedicated Engagement Write service + sharded-counter store) and push-notification fanout are out of scope. The feed still shows counts by reading the denormalized like_count/comment_count fields off the hydrated post object. The write side layers back on as new nodes emitting onto the same Engagement stream the ranker's feature loop already consumes — no change to the read spine.

Primary sources

  • TAO: Facebook's Distributed Data Store for the Social Graph (USENIX ATC '13)
  • Scaling Memcache at Facebook (NSDI '13)
  • f4: Facebook's Warm BLOB Storage (OSDI '14)
  • An Analysis of Facebook Photo Caching (SOSP)

Now defend it

Reading a design is not the same as being able to hold one under questioning. The workspace asks the same questions an interviewer would, and the simulator disagrees with you when the diagram does not support the claim.

Work Instagram News Feed yourself

Build the primitives this design leans on

Each one is an animated curriculum that constructs the system from scratch.

More in Feeds, Timelines, Counters & Ranking

What to show and in what order: fanout on write versus read, hot/top/new scoring, approximate counters, trending, and recommendation.