Skip to main content

Instagram

Sharpened prompt. Design a photo and short-video sharing service for 2B monthly users, where a post from an account with 400M followers appears in feeds within seconds without writing 400M rows, a feed loads in under 200ms with ML ranking, and a 60-second video is playable on a 3G phone in Mumbai moments after upload.

Everyone knows to say "hybrid fan-out." The interview is won on where exactly you put the threshold, what the merge looks like, and what happens to the counters.

1. Problem framing​

Functional requirements​

  • Post photos and videos; follow accounts; view a ranked home feed.
  • Like, comment, save; view a profile grid.
  • Stories: ephemeral 24-hour content with a separate consumption model.
  • Explore: content from accounts you do not follow.

Non-functional requirements​

PropertyTargetConsequence
Feed loadP99 < 200msPre-computed feeds; ranking must be two-stage
Post visibility< 5s for normal accountsAsync fan-out is fine; synchronous is not
Read:write ratio~100:1Optimise reads at almost any write cost
Media availability99.99% via CDNOrigin never serves an end user directly
DurabilityNo lost postsMedia in object storage before the post row commits

Back-of-the-envelope​

Users: 2B MAU, 500M DAU
Posts: 100M/day = 1,160/sec avg, ~5,000/sec peak
Feed reads: 500M DAU × 20 refreshes = 10B/day = 116k/sec avg, ~500k/sec peak
Fan-out: avg 200 followers -> 1,160 × 200 = 232k feed writes/sec (fine)
one celebrity at 400M followers -> 400M writes for ONE post.
At 100k writes/sec that is 66 minutes. This single number kills pure push.
Media: 100M posts × 2 MB avg (multi-resolution) = 200 TB/day ingest
Reads: 10B feed items × 300 KB = 3 PB/day egress -> CDN is the whole cost model

The 66-minute figure is the sentence to say out loud. It converts "celebrities are a problem" from a slogan into a measured constraint.

2. High-level architecture​

3. Component inventory​

ComponentConcrete choiceWhy this one
Post storeCassandra, partition user_id, clustering post_id DESCProfile grids are a range scan on one partition; writes are append-only
Feed storeRedis sorted sets / lists, capped at ~1000 entriesA feed is a bounded, ordered ID list. Never store post content here
Social graphSharded MySQL + a TAO-style cache, or a dedicated graph serviceFollower lists are read constantly and must be paginated efficiently
Fan-outKafka + worker poolBursty, retryable, must not block the post response
Media pipelineS3 event → SQS → FFmpeg on spot instancesEmbarrassingly parallel and interruption-tolerant
RankingTwo-stage: lightweight retrieval + TensorRT/ONNX heavy model500 candidates scored in tens of milliseconds
CountersRedis + async durable rollupSee 5.5 — a database row cannot take a viral like rate
DeliveryMulti-tier CDN with per-device variants3 PB/day means the CDN strategy is the cost model

4. The toughest parts​

4.1 Fan-out: push, pull, or both​

Why it's hard. Two pure strategies, both broken at the extremes.

Fan-out on write (push). On post, append the post ID to every follower's feed list. Reads are then a single Redis LRANGE — beautiful. But a 400M-follower post is 400M writes, taking over an hour and consuming the entire fan-out fleet. Worse, most of those followers are inactive: you are writing to 300M feeds that nobody will read.

Fan-out on read (pull). Store nothing; on feed load, query every followed account's recent posts and merge. Writes are free. But a user following 1,000 accounts triggers 1,000 queries per refresh, at 500k refreshes/sec. That is 500M queries/sec, which is not a system.

Solution — hybrid, with an explicit threshold and an active-user filter.

CELEBRITY_THRESHOLD = 25_000 # tune from the follower-count distribution

async def on_post(post):
n = await graph.follower_count(post.author_id)
if n > CELEBRITY_THRESHOLD:
# Pull model: store once, merge at read time.
await celeb_store.prepend(post.author_id, post.id)
return
# Push model, but only to followers who might actually read it.
async for batch in graph.followers(post.author_id, batch=1000):
active = await activity.filter_recent(batch, days=30) # ~30% of followers
await feeds.push_many(active, post.id) # Redis pipeline, capped

async def get_feed(user_id, limit=50):
pushed = await feeds.range(user_id, 0, 500) # one Redis call
celebs = await graph.following_celebrities(user_id) # usually < 50
pulled = await celeb_store.recent_many(celebs, limit=50) # one multi-get
return merge_by_time(pushed, pulled)[:500] # candidate set

Three details that elevate this from the textbook answer:

  • The active-user filter is the biggest single win. Only ~30% of followers are active in a 30-day window; skipping the rest cuts fan-out volume by 70% at zero user-visible cost. Inactive users get their feed built on demand at next login.
  • The threshold is a distribution question, not a constant. Pick it where the tail begins in your follower-count histogram, and make it dynamic: if the fan-out queue backs up, temporarily lower the threshold so more accounts fall into pull mode. Load-adaptive, not hardcoded.
  • Cap the feed list. 1,000 entries per user. Nobody scrolls past that, and it bounds Redis memory at 500M × 1000 × 8 bytes — still 4 TB, so cap active users only and rebuild cold feeds on demand.

The merge itself is a k-way merge over a handful of sorted lists — a bounded min-heap over (pushed_list, celeb_list_1..n), which is O(N log k) with k rarely above 50.

4.2 The media pipeline: 200 TB/day without blocking the user​

Why it's hard. A user hits "share." If you transcode synchronously, they wait 30 seconds for a video and the request times out on a flaky connection. But if you publish the post before derivatives exist, followers see a broken image. And you must produce a lot of derivatives: multiple resolutions, formats (JPEG/WebP/AVIF, H.264/VP9/AV1), and thumbnails — because a 2019 Android phone and a current iPhone should not receive the same bytes.

Solution — upload directly to object storage, publish optimistically, transcode asynchronously with a progressive-availability contract.

Two things make this work. The client uploads directly to S3 with a presigned URL, so 200 TB/day never traverses your API tier. And the client already has the image it just captured, so it can render its own post immediately — the post is "live" from the author's perspective before any derivative exists.

For video, split into GOP-aligned segments and transcode them in parallel across the fleet, then stitch the HLS/DASH manifest:

# Segment on keyframes so chunks can be transcoded independently and reassembled.
ffmpeg -i input.mp4 -c copy -f segment -segment_time 4 \
-reset_timestamps 1 -segment_format mp4 chunk_%04d.mp4

# Each chunk is an independent job. 60s video / 4s chunks = 15 parallel jobs.
# Wall-clock time becomes ~1 chunk's transcode time, not the whole video's.
ffmpeg -i chunk_0007.mp4 -c:v libx264 -crf 23 -preset veryfast \
-vf scale=-2:720 -c:a aac -b:a 128k out_720p_0007.mp4

Run the fleet on spot instances — the work is idempotent and re-runnable, so a 70% discount costs you only occasional retries. Prioritise the ladder: produce 720p first and mark the post ready, then backfill 1080p and AV1. Perceived latency is set by the first playable rendition, not the last.

4.3 Ranking 500 candidates in under 100 milliseconds​

Why it's hard. A good feed is not reverse-chronological. It scores each candidate on predicted engagement using signals about the viewer, the author, the content, and the relationship. A full model over 500 candidates involves embedding lookups and a deep network — hundreds of milliseconds if done naively, times 500k feed requests/sec.

Solution — the standard two-stage funnel, with the expensive model seeing very few items.

Candidate generation ~10,000 items cheap heuristics, Redis + inverted indexes ~5ms
│ recency, follow graph, explore pool, "seen" filter
▼
Lightweight ranking ~500 items logistic regression / two-tower dot product ~10ms
│ precomputed embeddings, single matrix multiply
▼
Heavy ranking ~50 items DNN with cross-features, TensorRT/ONNX ~40ms
│ batched GPU or optimised CPU inference
▼
Re-ranking / policy ~20 items diversity, ads, integrity, dedup ~5ms
async def rank(user_id: int, candidates: list[Post]) -> list[Post]:
# Fetch user features once; item features come from a warm feature store.
uf = await feature_store.user(user_id) # Feast / in-house
ifs = await feature_store.items([c.id for c in candidates]) # single batched call

# Stage 2: two-tower dot product. Item embeddings are precomputed at post time,
# so scoring 500 candidates is one 500×128 matmul — microseconds.
scores = user_tower(uf) @ np.stack([f.embedding for f in ifs]).T
top = [candidates[i] for i in np.argsort(scores)[-50:]]

# Stage 3: the heavy model sees only 50 items, batched into one inference call.
final = await ranker.predict_batch(uf, [ifs[c.id] for c in top]) # TensorRT
return apply_policy(sorted(zip(top, final), key=lambda x: -x[1]))

The re-ranking stage is where product rules live and it deserves a mention because candidates usually skip it: diversity (no more than 2 consecutive posts from one author), ad injection at fixed slots, integrity filtering, and seen-post suppression via a per-user Bloom filter or Redis bitmap — the same mechanism Tinder uses for swipes.

Precompute item embeddings at post time, not at read time. That single choice is what makes stage 2 cheap enough to run on every request.

4.4 Stories: 500 million ephemeral posts a day​

Why it's hard. Stories invert the normal assumptions: everything expires in 24 hours, the read pattern is "whose stories are new?" rather than a merged timeline, and viewer lists must be tracked per story. Fanning out stories into the main feed store would churn it completely every day.

Solution — a separate store with native TTL and a completely different read shape.

# Stories are authored-centric, not feed-centric: the reader asks
# "which of the accounts I follow have unseen stories?"

# Per-author story list, self-expiring.
await redis.zadd(f"stories:{author_id}", {story_id: expires_at})
await redis.expire(f"stories:{author_id}", 86400)

# The story tray: no fan-out at all. Intersect following-list with active authors.
async def story_tray(user_id):
following = await graph.following(user_id) # cached, ~500 ids
active = await redis.smismember("stories:active", following) # one round trip
authors = [f for f, a in zip(following, active) if a]
seen = await redis.getbit_many(f"seen:{user_id}", authors)
return order_by_affinity(authors, seen) # unseen first, then affinity

No fan-out means no celebrity problem for stories at all — the read is a set intersection against a bounded following list. Viewer tracking uses a per-story HyperLogLog for the count plus a capped list for the actual names (the UI shows a few hundred at most). Media derivatives are lighter (one or two renditions, since stories are full-screen mobile only) and the CDN TTL matches the 24-hour expiry so storage self-cleans.

4.5 Counters that no database row can survive​

Why it's hard. A celebrity post gets 500,000 likes in a minute. UPDATE posts SET likes = likes + 1 WHERE id = ? at 8,300 writes/sec against one row means every transaction queues on the same lock. Row-level contention, replication lag, and eventually the shard stalls — taking down every other post on it.

Solution — never increment a shared row on the request path.

# 1. Write the fact, not the aggregate. The like itself is the durable record.
await cassandra.execute(
"INSERT INTO likes (post_id, user_id, ts) VALUES (?, ?, ?)", (post_id, uid, now))

# 2. Increment a sharded counter in Redis. N counters, summed on read.
shard = uid % 16
await redis.hincrby(f"likes:{post_id}", f"s{shard}", 1)

# 3. Read: sum the shards. One HGETALL, 16 small integers.
count = sum(map(int, (await redis.hgetall(f"likes:{post_id}")).values()))

# 4. Roll up to durable storage asynchronously, every ~10s, from the like stream.

Counter sharding spreads the write across 16 keys (and therefore potentially 16 Redis shards), turning 8,300 writes/sec on one key into 520/sec on each. For follower counts on very large accounts, go further: serve an approximate, cached value refreshed every few seconds. Nobody can tell whether an account has 400,000,000 or 400,000,012 followers, and pretending otherwise costs real money.

The pattern to name: "Store the event, derive the aggregate. The counter is a cache of a COUNT(*) you never actually run."

4.6 Read-your-own-writes across an eventually consistent stack​

Why it's hard. A user posts, then immediately pulls to refresh. The write went to the primary; the read hit a replica 200ms behind, or a feed list the fan-out worker has not reached. The post is missing. Users interpret this as data loss and post again — which is both a support burden and a duplicate-content problem.

Solution — three cheap mechanisms, applied together.

  1. Write to the author's own feed synchronously, before returning from the post API. It is one Redis operation and it removes the most-noticed case entirely.
  2. Session read affinity: after a write, pin that user's reads to the primary (or to a replica known to have caught up) for a few seconds, carried in a session cookie or a monotonic read token.
  3. Client-side optimistic insert: the client already has the image and the caption, so it renders the post locally and reconciles when the server confirms.

For everything that is not the user's own action — a friend's new post, a like count — eventual consistency is not merely acceptable, it is invisible. Distinguish the two cases explicitly; conflating them leads to over-engineering the whole read path for a problem that affects one item.

4.7 Three petabytes a day of egress​

Why it's hard. Media delivery, not compute, is the dominant cost. At 3 PB/day, a fraction of a cent per GB is millions of dollars a month, and naive delivery (one 4 MB original to every device) multiplies it several-fold.

Solution — spend engineering effort on bytes, because bytes are the bill.

  • Serve the right rendition. Negotiate on Accept and client hints: AVIF where supported (~50% smaller than JPEG at equal quality), WebP as the fallback, JPEG last. Pick resolution from the device's viewport, not from a fixed ladder.
  • Multi-tier CDN with ISP embedding. Origin → regional shield → ISP-embedded cache. A shield tier collapses many edge misses into one origin fetch, which matters when a post goes viral in a region simultaneously.
  • Prefetch the next few feed items over Wi-Fi only. Perceived speed improves enormously and the bytes are cheaper off-peak.
  • Tier aged media. Posts older than a year get their high-resolution derivatives deleted and regenerated on demand from the original, which itself sits in cold storage. The overwhelming majority of views happen in the first 48 hours.

Quantify one of these in the interview: "Switching the default from JPEG to AVIF at equal perceived quality cuts roughly 40% of 3 PB/day. At CDN list prices that is on the order of tens of millions of dollars a year — which is why the media team is bigger than the feed team."

5. What breaks first​

EventFirst failureMitigation
Celebrity postsFan-out queue backs up for everyonePull model above the threshold; dynamic threshold under load
Viral post likesHot counter row / Redis keySharded counters; approximate counts for large values
Coordinated event (New Year)Upload + transcode fleet saturatesQueue-depth autoscaling; degrade to fewer renditions temporarily
Redis feed shard lossFeeds empty for those usersRebuild on demand from posts + graph; treat feed store as a cache
Ranking model latency spikeFeed P99 breachesTimeout the heavy stage and fall back to stage-2 ordering
CDN origin thundering herd on a viral postOrigin saturationShield tier + request coalescing at the CDN

Note the recurring principle: the feed store is a cache, not a database. Losing it costs CPU, never data. Saying that explicitly answers half the failure questions at once.

6. Cheat sheet​

  • Fan-out: hybrid at ~25k followers, push only to 30-day-active followers, cap feeds at 1,000, k-way merge at read.
  • The number: 400M followers ÷ 100k writes/sec = 66 minutes. That is why pure push dies.
  • Media: presigned direct-to-S3 upload, publish optimistically, GOP-parallel transcode on spot, progressive rendition availability.
  • Ranking: 10k → 500 → 50 → 20 funnel; item embeddings precomputed at post time; policy re-rank last.
  • Stories: separate TTL store, set-intersection read, zero fan-out.
  • Counters: store the event, shard the counter, approximate above a threshold.
  • Consistency: sync-write the author's own feed, session affinity for a few seconds, optimistic client render.
  • The one-liner: "A feed is a cache of a merge. I precompute it for the 99% of authors where precomputation is cheap, and compute it at read time for the 1% where it isn't."