Limited time: AI code review, hints, mock interviews, whiteboard analysis, and all Pro features are unlocked. Enroll
⏱️ 37 min read

Designing Instagram / Pinterest - Photo Sharing Platform

Difficulty: Intermediate Topics: CDN, Object Storage, Fan-out, Media Processing, Feed Ranking Asked at: Meta, Pinterest, Snap, Google, Amazon Prerequisites:CDN, Caching, and Fan-Out


1. Understanding the Problem

Instagram is a photo and video sharing platform where users upload media, follow other users, and consume a personalized feed of content. The system must handle billions of photo uploads, deliver images globally with low latency via CDN, and generate personalized feeds for hundreds of millions of users. The key challenges are: efficiently processing and storing media at scale, generating feeds without overwhelming the system, and delivering images fast regardless of user location.


2. Naive First Cut

flowchart LR
    Client["Mobile App"]:::client
    API["API Server"]:::service
    DB["Postgres DB"]:::data
    Disk["Local Disk Storage"]:::data

    Client --> API
    API --> DB
    API --> Disk

    classDef client fill:#4c3a5e,stroke:#818cf8,color:#e2e8f0
    classDef service fill:#1a3a2a,stroke:#4ade80,color:#e2e8f0
    classDef data fill:#3b3520,stroke:#fbbf24,color:#e2e8f0
Color Meaning
🟠 Purple Client
πŸ”΅ Blue Edge / Gateway
🟒 Green Service
🟣 Purple Async (Queue / Kafka)
🟑 Yellow Data store

How this breaks:

The rest of the doc evolves this into a globally distributed media platform with CDN delivery and intelligent feed generation.


3. Prior Art We’re Drawing From


4. Functional Requirements

Core (Top 3)

  1. Upload photos and videos - users can upload media with captions, apply filters, and tag locations
  2. View personalized feed - users see a ranked feed of posts from people they follow
  3. Follow and unfollow users - build a social graph that drives feed generation

Below the Line


5. Non-Functional Requirements

Core

NFR Target
Feed Latency Feed load < 500ms P95 globally
Upload Latency Photo upload completes < 3 seconds (user sees confirmation)
Availability 99.99% - users expect Instagram to always be up
Scale 2B monthly active users, 100M+ photos uploaded daily

Below the Line

6. Scale Estimation (Back-of-Envelope)

Two numbers do the most work in the deep dives below. The read path is ~90K feed loads/sec at peak, and each one needs posts from ~500 different authors. The storage path adds 200TB a day, and storage bills are cumulative, so the cost problem compounds rather than plateaus.


7. Core Entities


8. API / System Interface

POST /api/v1/posts
  Body: { mediaFile (multipart), caption, location?, tags[]? }
  Response: { postId, mediaUrls, status, timestamp }
  Auth: JWT Bearer token
  Note: authorId comes from the JWT, never the body
  Note: status tells the client whether the media is renderable yet

GET /api/v1/feed?cursor=<timestamp>&limit=20
  Response: { posts: [{ postId, authorId, mediaUrls, caption, likes, timestamp }], nextCursor }
  Note: Cursor-based pagination for infinite scroll

POST /api/v1/users/{userId}/follow
  Response: { status: "FOLLOWING", timestamp }

DELETE /api/v1/users/{userId}/follow
  Response: { status: "UNFOLLOWED" }

GET /api/v1/users/{userId}/profile
  Response: { userId, username, bio, postCount, followerCount, followingCount, posts[] }

9. High-Level Design

FR1: Upload Photos and Videos

When a user takes a photo and hits β€œShare,” we need to store the image and make it viewable at the sizes the app actually renders β€” a small thumbnail in the grid, a larger one in the feed, full resolution on tap.

Build the simplest thing that satisfies that. Two services, object storage for the bytes, and one database for the metadata. No queue, no worker pool, no CDN. Each of those answers a non-functional requirement β€” latency, throughput, cost β€” so each belongs in a deep dive where we can show why it is needed instead of asserting it up front.

New components we need:

  1. API Gateway - Entry point for all client requests. Handles JWT auth, rate limiting, and file size limits.
  2. Upload Service - Takes the uploaded bytes, resizes them into the variants the app needs, writes all of them to object storage, and records the post.
  3. Object Storage (S3) - Stores the image bytes. Write-once, read-many, 11 nines of durability.
  4. Post Metadata DB (Postgres) - One row per post: author, caption, timestamp, and the storage keys for each variant.
flowchart LR
    App["Mobile App"]:::client
    GW["API Gateway"]:::edge
    US["Upload Service"]:::service
    S3[("S3 Object Store")]:::data
    DB[("Postgres posts")]:::data

    App -->|"1. POST photo and caption"| GW
    GW -->|"2. Auth and forward"| US
    US -->|"3. Resize to 4 variants"| US
    US -->|"4. Write all variants"| S3
    US -->|"5. Insert post row"| DB
    US -->|"6. Return 201 with postId"| App

    classDef client fill:#4c3a5e,stroke:#818cf8,color:#e2e8f0
    classDef edge fill:#1e3a5f,stroke:#60a5fa,color:#e2e8f0
    classDef service fill:#1a3a2a,stroke:#4ade80,color:#e2e8f0
    classDef data fill:#3b3520,stroke:#fbbf24,color:#e2e8f0

Step-by-step flow:

  1. User selects photo, adds caption, hits β€œShare” β†’ app uploads the file via multipart POST to Gateway
  2. Gateway authenticates the user, rejects anything over 50MB, forwards the stream to Upload Service
  3. Upload Service generates a mediaId and resizes the original into 4 variants (150, 320, 640, 1080px), converting to WebP and stripping EXIF
  4. It writes the original and all 4 variants to S3 under media/{userId}/{mediaId}/
  5. It inserts one row into posts with the caption, author, timestamp and variant keys
  6. It returns 201 Created with the postId. The photo is now fully viewable

Why resize at all, instead of serving the original?

A 4MB 12-megapixel photo rendered into a 150px grid cell wastes 99% of the bytes it transferred. Resizing is not an optimization we can defer to a deep dive β€” the feed is unusable on mobile data without it. Where the resize happens is what we defer.

What we have deliberately left broken. This is a correct, complete upload path, and for a few thousand uploads a day it is genuinely fine. It has two holes. First, step 3 runs inside the request the user is waiting on: four resizes at roughly 2 seconds each is 8 seconds of work before we return, against a 3-second NFR. Second, step 4 puts the bytes in exactly one S3 region, so a user in Jakarta pays a trans-Pacific round trip for every image in their feed. Neither is a missing feature; both are non-functional failures. The synchronous resize is Deep Dive 1, the single-region reads are Deep Dive 3.


FR2: View Personalized Feed

Feed is the core experience. When a user opens Instagram, they need to see recent posts from the people they follow, newest first.

Read that requirement literally and it is a query, not an architecture. We already have a posts table and (in FR3) we will have a follows table. β€œRecent posts from people I follow” is a join across the two. So that is what we build.

New components we need:

  1. Feed Service - Answers GET /feed. Looks up who the user follows, pulls those authors’ recent posts, sorts by time, returns a page.

That is the only new component. No new datastore: the feed is derived from data FR1 and FR3 already store, so there is nothing extra to keep consistent.

flowchart LR
    App["Mobile App"]:::client
    GW["Gateway"]:::edge
    FS["Feed Service"]:::service
    DB[("Postgres posts and follows")]:::data
    S3[("S3 Object Store")]:::data

    App -->|"1. GET feed with cursor"| GW
    GW -->|"2. Forward to feed svc"| FS
    FS -->|"3. Read followees for user"| DB
    FS -->|"4. Read recent posts by those authors"| DB
    FS -->|"5. Return posts with image URLs"| App
    App -->|"6. Fetch each image"| S3

    classDef client fill:#4c3a5e,stroke:#818cf8,color:#e2e8f0
    classDef edge fill:#1e3a5f,stroke:#60a5fa,color:#e2e8f0
    classDef service fill:#1a3a2a,stroke:#4ade80,color:#e2e8f0
    classDef data fill:#3b3520,stroke:#fbbf24,color:#e2e8f0

Step-by-step flow:

  1. User opens app β†’ GET /feed?cursor=&limit=20 hits Feed Service
  2. Feed Service reads the user’s followee list from follows (~500 rows for a typical user)
  3. It queries posts for recent rows by any of those 500 authors, ordered by timestamp, limit 20
  4. It joins in what the client needs to render: author username, caption, like count
  5. It returns the page plus a cursor, with each post carrying S3 URLs for the variants
  6. The app requests each image directly from S3 and renders the feed

Why chronological, and why is that not a cop-out?

Ranking is a real product requirement, but it is a quality requirement layered on top of a working feed, not a precondition for one. We need a correct set of candidate posts before there is anything to rank. Chronological gets us the candidate set; Deep Dive 5 earns the ranker.

What we have deliberately left broken. This design is honest and it does not survive contact with our own scale numbers. Three holes, and it is worth naming them separately because they have different fixes:

And one hole that does not exist yet but will: nothing here breaks when a user with 100M followers posts, because we do no work at post time at all. Deep Dive 2 fixes the read path by moving work to write time, and that is what makes celebrities a problem. Deep Dive 4 cleans up after Deep Dive 2.


FR3: Follow and Unfollow Users

The social graph is what makes the feed personal. Follow is a small operation with a large blast radius: it is the input FR2’s feed query reads.

A follow is a directed edge. One row in one table.

New components we need:

  1. Social Graph Service - Owns follow and unfollow. Writes the edge and enforces the rules that make it correct: no self-follows, no duplicates, blocks respected.

The edge lives in a follows table in the same Postgres as posts, indexed both ways β€” (follower_id, followee_id) to answer β€œwho do I follow” for the feed query, and (followee_id, follower_id) to answer β€œwho follows me” for the profile screen.

flowchart LR
    App["Mobile App"]:::client
    GW["Gateway"]:::edge
    SGS["Social Graph Service"]:::service
    DB[("Postgres follows")]:::data
    FS["Feed Service"]:::service

    App -->|"1. POST follow user B"| GW
    GW -->|"2. Auth and forward"| SGS
    SGS -->|"3. Insert follow edge"| DB
    SGS -->|"4. Return 200"| App
    FS -->|"5. Later reads followees here"| DB

    classDef client fill:#4c3a5e,stroke:#818cf8,color:#e2e8f0
    classDef edge fill:#1e3a5f,stroke:#60a5fa,color:#e2e8f0
    classDef service fill:#1a3a2a,stroke:#4ade80,color:#e2e8f0
    classDef data fill:#3b3520,stroke:#fbbf24,color:#e2e8f0

Step-by-step flow:

  1. User A taps β€œFollow” on User B’s profile β†’ POST /users/{B}/follow
  2. Gateway authenticates A. Note that A comes from the JWT, never from the request body β€” otherwise anyone can make anyone follow anyone
  3. Social Graph Service inserts (follower_id=A, followee_id=B). A unique constraint on the pair makes a double-tap idempotent rather than creating two edges
  4. It returns 200 FOLLOWING
  5. Nothing else happens. A’s next feed request reads the new edge and B’s posts appear

Unfollow is the same in reverse: delete the row. B’s posts stop appearing on A’s next feed load, because the feed is computed from the edge at read time.

What we have deliberately left broken. Very little, and that is the interesting part. Step 5 is the one place this naive design genuinely beats the production one: because FR2 computes the feed on read, a follow takes effect on the next refresh with no extra work and no window of inconsistency. Deep Dive 2 gives that property up. Once feeds are pre-computed, a new follow no longer changes anything until someone backfills the new followee’s recent posts into A’s feed, and an unfollow leaves stale posts behind until something cleans them out. That is a real cost, we take it on knowingly, and Deep Dive 2 has to pay it.

The one thing that does break here at scale is the follower-side read. followee_id = B for a celebrity is a 100M-row scan, which nothing on this page needs yet but Deep Dive 4 will.


10. Technology Choices

Tier Purpose Stores Access Pattern Primary Alternatives
Object Storage Original and resized images Raw media files (JPEG, MP4) Write once, read many via CDN S3 GCS, Azure Blob
CDN Global image delivery Cached image variants High-QPS reads, edge-cached CloudFront or Cloudflare Fastly, Akamai
Feed Store Pre-computed user feeds Ordered post IDs per user Write-heavy (fan-out), sequential reads Cassandra ScyllaDB, DynamoDB
Post Metadata DB Post details, captions, tags Structured post data Read-heavy, indexed queries Postgres (sharded by user) CockroachDB, Vitess
Social Graph DB Follow relationships Follower and following edges High-QPS lookups (who follows whom) Redis Cluster (adjacency sets) Neo4j, TAO-style cache
Event Bus Async processing pipeline Upload events, feed fan-out events Fan-out writes, ordered per user Kafka Redpanda, Kinesis
Cache Hot feed data, user profiles Serialized feed pages, profile JSON High-QPS reads, TTL-based Redis Cluster Memcached
Media Processing Queue Image resize jobs Processing tasks FIFO per upload SQS or Kafka RabbitMQ

Why Cassandra for the feed store, not Postgres? Feed reads are sequential (give me the next 20 posts) and writes are massive during fan-out (one post fans out to millions of follower feeds). Cassandra’s write-optimized LSM-tree and partition-key access pattern (userId β†’ sorted posts) is perfect. Postgres would choke on the write amplification.

Why Redis for the social graph, not the main DB? β€œDoes user A follow user B?” is called on every feed request, every like, every comment. At 2B users, this needs sub-millisecond latency. Redis SET operations (SISMEMBER) answer this in microseconds.


11. Data Modeling

Postgres (Post Metadata β€” structured post data, sharded by user):

CREATE TABLE posts (
    post_id UUID PRIMARY KEY,
    author_id UUID NOT NULL,
    caption TEXT,
    media_type VARCHAR(10),  -- image, video, carousel
    media_urls JSONB,
    location JSONB,
    like_count INTEGER DEFAULT 0,
    comment_count INTEGER DEFAULT 0,
    created_at TIMESTAMP NOT NULL
);
CREATE INDEX idx_posts_author ON posts(author_id, created_at DESC);

Cassandra (Feed Store β€” pre-computed user feeds):

Table: user_feed
  PK: user_id
  SK: created_at (DESC, TimeUUID for uniqueness)
  Columns: post_id, author_id, media_thumbnail_url
  TTL: 30 days (old feed entries auto-expire)

Redis (Social Graph β€” follow relationships):

Key: "followers:{userId}" β†’ Set of follower userIds
Key: "following:{userId}" β†’ Set of followee userIds
Operations: SADD on follow, SREM on unfollow, SISMEMBER for auth checks

S3 + CDN (Media Storage β€” images and videos):

Path: s3://media-bucket/{userId}/{postId}/{variant}.jpg
Variants: original, 1080p, 640p, 320p, thumbnail
CDN: CloudFront/Cloudflare with immutable cache headers (media never changes)

Access Patterns:

Query Data Source How
View home feed Cassandra Range scan user_feed by user_id, paginate by created_at
Post a photo Postgres + S3 + Kafka Upload to S3 β†’ write post metadata β†’ fan-out event to Kafka
Fan-out to followers Kafka β†’ Cassandra Workers read followers:{authorId}, insert into each follower’s user_feed
Like a post Postgres + Redis Increment like_count in Postgres, cache invalidation
Check follow status Redis SISMEMBER followers:{targetId} currentUserId

How Fan-Out Works When a User Posts a Photo:

  1. User uploads photo β†’ Media Service resizes to 4 variants, stores in S3
  2. Post Service writes post metadata to Postgres (source of truth)
  3. Publishes POST_CREATED event to Kafka (includes post_id, author_id, media_urls)
  4. Fan-out workers consume event β†’ read followers:{authorId} from Redis (e.g., 50K followers)
  5. For each follower: INSERT INTO user_feed (user_id, created_at, post_id, author_id, thumbnail) in Cassandra
  6. For celebrity accounts (>500K followers): skip fan-out, merge at read time (same hybrid pattern as Twitter)
  7. Feed read: SELECT * FROM user_feed WHERE user_id = ? ORDER BY created_at DESC LIMIT 20

12. Deep Dives

1) How do we confirm an upload in under 3 seconds when the resizing takes 8?

Problem: FR1 resizes inside the request. Our NFR says the user sees confirmation in under 3 seconds, and the resize alone does not fit in that budget.

In simple terms: User uploads a photo. Should we make them wait while we resize it to four different sizes? No - accept the upload, confirm it, and do the resizing after.

Bad: exactly what the high-level design built β€” resize in the request path. Two numbers kill it. Per request: a 4MB photo over a typical mobile uplink is 4-8 seconds, then four ImageMagick resizes at roughly 2 seconds each adds 8 more, so 12-16 seconds against a 3-second NFR. In aggregate: 4,600 resize ops/sec Γ— 2 seconds of CPU each is 9,200 cores of work, and every one of those core-seconds is held inside a request thread with a phone waiting on the other end. We are paying for 9,200 cores and blocking on them.

There is a third failure that is worse than either. If the resize crashes on photo 3 of 4, the user gets a 500 on an upload whose bytes we already have. The upload is lost from their point of view and orphaned from ours.

Good: Accept the upload, write the original, return 201, and resize afterwards on a background thread in the same service. The user is unblocked and the crash case is recoverable. But the work is still on the machine that took the request, so a burst of uploads competes with request serving for the same CPU, and a deploy that restarts the process drops every resize in flight.

Great: Multi-stage pipeline with auto-scaling worker pools:

  1. Upload stage: Client uploads to a pre-signed S3 URL directly (bypasses API server entirely for large files). Upload Service just validates and records metadata.
  2. Processing stage: Worker pool auto-scales based on queue depth. Each worker: download β†’ resize (150, 320, 640, 1080px) β†’ convert to WebP β†’ strip EXIF β†’ upload variants β†’ update DB.
  3. Optimization: Generate progressive JPEGs so images render top-to-bottom even on slow connections. Store a tiny 20px blurred placeholder (BlurHash) in the post metadata for instant feed skeleton rendering.
flowchart LR
    Client["Client"]:::client
    US["Upload Service"]:::service
    S3O[("S3 Originals")]:::data
    Q["SQS Queue"]:::async
    W["Worker Pool"]:::service
    S3P[("S3 Processed")]:::data
    CDN["CDN"]:::edge

    Client -->|"1. Get pre-signed URL"| US
    Client -->|"2. Upload directly to S3"| S3O
    US -->|"3. Enqueue processing job"| Q
    Q -->|"4. Pick job"| W
    W -->|"5. Download original"| S3O
    W -->|"6. Resize + compress + store"| S3P
    S3P -->|"7. Serve via CDN"| CDN

    classDef client fill:#4c3a5e,stroke:#818cf8,color:#e2e8f0
    classDef edge fill:#1e3a5f,stroke:#60a5fa,color:#e2e8f0
    classDef service fill:#1a3a2a,stroke:#4ade80,color:#e2e8f0
    classDef async fill:#3b1f5e,stroke:#c084fc,color:#e2e8f0
    classDef data fill:#3b3520,stroke:#fbbf24,color:#e2e8f0

Cost consideration: Processing 100M images/day at 4 variants each = 400M resize operations. GPU-accelerated workers (using libvips, not ImageMagick) cut processing time from 2s to 200ms per image. Auto-scaling down during off-peak saves 60% compute cost.


2) How do we serve 90K feed loads a second when each one needs posts from 500 authors?

Problem: FR2 computes the feed by reading the follow list, then reading recent posts from every author on it. Correct, and it does not survive our read volume.

In simple terms: When you open Instagram, should the app go ask 500 different people β€˜got any new posts?’ That’s too slow. Better to pre-build your feed ahead of time.

Bad: the fan-out-on-read we built in FR2. Put the number on it: ~90K feed loads/sec at peak Γ— ~500 followees each is roughly 45M post-lookups/sec, against a single Postgres that also serves every upload write. Even batched into one WHERE author_id IN (...) query per feed load, that is 90K queries/sec each touching 500 index ranges and sorting the union β€” and the sort cannot be satisfied by an index, because the rows it merges are spread across 500 unrelated index positions.

The tail is worse than the average. Feed latency is bounded by the slowest of the 500 authors’ partitions, so P95 feed latency tracks P99.8 partition latency. Our budget is 500ms P95 globally.

Good: Fan-out-on-write. When a user posts, push the postId into every follower’s feed list, stored one partition per user. Feed reads become a single sequential read of one partition β€” no join, no 500-way merge, no cross-partition sort. Reads go from 45M lookups/sec to 90K single-partition scans/sec. But the write side inverts: an average user with 500 followers turns one post into 500 writes, and a celebrity with 100M followers turns one post into 100M.

Great: Fan-out-on-write onto a dedicated feed store, with the fan-out itself made asynchronous and durable.
πŸ’‘ Fan-out = taking one event and delivering it to many destinations. Here one post becomes one row in each follower’s feed. Learn more β†’

Mechanism:

  1. Upload Service publishes post.published to Kafka, partitioned by authorId so one author’s posts fan out in order
  2. Fan-out Service consumes it, reads the author’s follower list, and writes the postId into each follower’s feed partition
  3. Feed Store (Cassandra) holds one partition per user, clustered by timestamp descending. A feed read is one sequential scan of the head of one partition
  4. Writes are batched (~1,000 rows per batch, async) so 10K followers costs ~10 round trips rather than 10K
  5. A rate limiter caps how much of the fan-out worker pool any single post may hold, so one large author cannot starve everyone else’s fan-out
  6. Redis caches the assembled first page per active user, which is the page almost every session requests and never scrolls past

Why Cassandra and not just more Postgres? The access pattern is now append-heavy writes and a single-partition head read per user β€” no joins, no ad-hoc queries, no transactions across users. That is precisely the shape an LSM-tree store with partition-local clustering is built for, and it scales by adding nodes rather than by making one node bigger.

The bill for this. Two things got worse, and both were promised earlier on this page:

Ranking the assembled candidates is a separate concern from assembling them, and it is Deep Dive 5.


3) How do we get a 200KB image onto a phone in Jakarta without a $13M monthly egress bill?

Problem: FR1 wrote every variant to one S3 region and FR2 handed the client S3 URLs. Every image in every feed, for every user on the planet, is a round trip to that one region.

In simple terms: If every image request goes to one bucket in Virginia, users in Jakarta wait 300ms per image and you pay full egress on all of it. A CDN puts copies close to users and bills less per byte.

Bad: what FR1 and FR2 built β€” clients fetching straight from origin S3. This fails on both axes at once.

Latency: a feed page is ~10 images. From Jakarta to us-east-1 is roughly 250-300ms per round trip, and browsers cap parallel connections per host, so the images arrive in waves. That alone blows the 500ms P95 feed budget before the first row is fully painted.

Cost: we egress ~5PB/day, which is ~150PB/month. S3 list-price egress is $0.09/GB, and 150PB is 150 million GB, so about $13.5M/month to serve images. There is no caching layer anywhere in that path, so the 100th user to view a popular post costs exactly as much as the first.

Good: Put a CDN (CloudFront / Cloudflare / Fastly) in front of S3 and let edges cache. Latency drops to the nearest PoP and egress moves to CDN pricing. But a naive drop-in still misses on first access in every region, still serves the same 1080px variant to a phone on 3G as to a laptop on fibre, and every edge miss hits origin directly β€” so a viral post can stampede S3 from 200 PoPs at once.

Great: Multi-layer CDN strategy with client-driven quality selection:

  1. Edge caching (CDN): Images cached at 200+ PoPs globally. TTL = 1 year (images are immutable - new upload = new URL). Cache hit ratio > 95% for popular content.
  2. Client-driven quality: App detects network speed and requests appropriate variant: cdn.instagram.com/media/{id}/w640.webp vs w1080.webp. Saves bandwidth on slow connections.
  3. Progressive loading: Feed shows BlurHash placeholder instantly β†’ low-res thumbnail loads in 50ms β†’ full resolution lazy-loads as user scrolls.
  4. Regional origin shields: Secondary cache layer between CDN edge and S3 origin. Reduces origin requests by another 80%.

Cost at scale: ~150PB of image egress per month. At a committed-use CDN rate of roughly $0.02/GB that is about $3M/month, against $13.5M/month serving direct from S3 at $0.09/GB. The CDN pays for itself more than four times over, and that is before the origin-shield layer cuts origin requests by a further 80%.


4) What happens when someone with 100M followers posts?

Problem: Deep Dive 2 bought a fast read path by doing the work at write time. That trade is priced per follower, and it was priced assuming a follower count in the hundreds.

In simple terms: Cristiano Ronaldo posts a photo. If we add it to 100M people’s feeds, that’s 100M database writes for one post. We need a different path for accounts that large.

Bad: the fan-out-on-write from Deep Dive 2, applied uniformly. One post from a 100M-follower account is 100M Cassandra writes. At ~1,000 rows per batch that is 100,000 batched round trips for a single post, and at 10 such posts an hour it is 1B writes/hour of fan-out for a handful of authors.

Note this is not only a throughput problem. It is a fairness problem: while those 100K batches are draining, every ordinary user’s fan-out is queued behind them, so a post from a user with 40 followers takes minutes to appear because a celebrity posted first. And it is a latency problem for the author, whose post is not visible to most of their followers for several minutes after they publish it.

Worth noticing what is not broken here: FR2’s original fan-out-on-read handled this case perfectly, because it did no work at post time. We introduced this problem ourselves in Deep Dive 2, deliberately, to fix a worse one.

Good: Skip fan-out for large accounts entirely and merge their posts at read time β€” go back to FR2’s approach, but only for them. Correct, and it puts a scatter-gather back into the read path we just spent Deep Dive 2 removing. Every feed load now queries each celebrity the user follows.

Great: Tiered hybrid with intelligent caching:

  1. Classify users: follower_count > 500K = β€œcelebrity.” Flag in Redis graph store.
  2. Skip fan-out for celebrities: Their posts go to a special β€œcelebrity posts” store (sharded by celebrityId, sorted by time).
  3. Feed assembly at read time: Feed Service fetches: (a) user’s pre-computed feed from Cassandra, (b) recent posts from celebrities they follow (max 10 celebrities Γ— 5 posts = 50 posts to merge).
  4. Cache celebrity feeds aggressively: Redis caches each celebrity’s last 50 posts. Updated on new post. All followers read from same cache - millions of cache hits, one write.
  5. Pre-warm on post: When celebrity posts, invalidate their Redis cache entry. First reader triggers cache fill; subsequent readers hit cache.

Net effect: Celebrity post = 1 write to celebrity store + 1 cache invalidation. vs. 100M writes with naive fan-out. Read overhead: +5ms per celebrity merge (parallel Redis fetches).


5) How do we pick the 20 posts a user sees out of 500 candidates?

Problem: FR2 ordered the feed by timestamp and Deep Dive 2 kept that ordering. A user following 500 accounts has far more eligible posts per session than screen space.

In simple terms: Showing posts purely by time means you miss the important ones posted while you slept. We need to surface the posts you’d actually care about, not just the newest ones.

Bad: the chronological ordering FR2 shipped. Do the arithmetic on what it discards. 500 followees posting even once a day each is 500 eligible posts; a session shows roughly 20 before the user closes the app. Five sessions a day is ~100 posts seen out of 500, so 80% of the content a user explicitly asked to see never reaches them β€” and which 20% gets through is decided purely by who posted most recently, which systematically favours high-frequency accounts over close friends who post twice a week.

Good: Simple scoring: score = recency_weight * time_decay + engagement_weight * (likes + comments). Better than chronological but doesn’t personalize.

Great: Lightweight ML ranker with candidate generation + ranking stages:

  1. Candidate generation: Pull 500 candidate posts (pre-computed feed + celebrity merge)
  2. Feature extraction: For each candidate, compute: time since posted, author-viewer relationship strength (interaction frequency), post engagement velocity (likes/min in first hour), content type match (does viewer prefer photos or videos?)
  3. Scoring: Simple logistic regression or small neural net predicts P(engagement). Trained offline on historical engagement data. Inference < 10ms for 500 candidates.
  4. Diversity injection: After ranking, ensure no more than 3 consecutive posts from same author. Mix in β€œdiscovery” posts (from friends-of-friends) at 10% ratio.

Why not a huge ML model? Ranking runs on every feed load. At 2.5B feed loads/day, 50ms of inference each is 125M CPU-seconds a day β€” about 1,450 cores running flat out, purely to order posts, and it adds 50ms to a 500ms P95 budget. Keep the online model small enough to score 500 candidates in under 1ms on CPU. Heavy models belong in offline training, where they produce the embeddings and weights the cheap online model reads.


6) How do we stop 200TB a day of new photos from becoming a bill that never stops growing?

Problem: FR1 writes the original plus 4 variants to S3 Standard and never touches them again. Storage charges are cumulative, so unlike every other cost on this page, this one does not plateau.

In simple terms: 100M photos per day is 200TB of new storage every day, and you keep paying for last year’s photos while adding this year’s. Old photos need to move to cheaper storage.

Bad: what FR1 built β€” everything in S3 Standard, kept forever. The first month looks survivable and that is the trap. 200TB/day is 6PB/month, and 6 million GB at $0.023/GB-month is about $138K for month one. But month two pays for month one’s data too. By the end of year one we are storing 73PB and paying about $1.68M/month, and that figure only ever goes up. Nothing in this design ever deletes or downgrades a byte.

The original is the most wasteful part of it. We keep a 4MB full-resolution file that no client ever requests, because FR1’s step 4 wrote it and nothing was ever decided about it.

Good: S3 lifecycle policies β€” transition to Infrequent Access after 30 days, Glacier after a year. This is a two-line bucket config and it cuts the bill substantially. But it is a pure function of age, so a two-year-old photo on a heavily-viewed profile gets buried in Glacier alongside one nobody has opened since it was posted, and retrieval latency on a profile scroll becomes the user’s problem.

Great: Intelligent tiering based on access patterns:

  1. Hot tier (S3 Standard): Posts < 7 days old. 80% of all accesses hit content from the last week.
  2. Warm tier (S3 IA): Posts 7-90 days old. Occasionally accessed via profile views and search.
  3. Cold tier (S3 Glacier Instant Retrieval): Posts > 90 days. Rare access but must still serve in < 100ms when profile is scrolled.
  4. Delete originals: After processed variants are confirmed, delete the original full-res upload (keep only the 1080px max). Saves 40% storage.
  5. Deduplication: Perceptual hash (pHash) on upload. If near-duplicate exists, store a reference instead of new file. Catches reposts and memes - saves ~15% storage.

Cost after optimization: ingest drops from 200TB/day to ~120TB/day once originals are deleted and near-duplicates are stored by reference, so year one accumulates ~43.8PB instead of 73PB. Tiering moves the blended rate from $0.023/GB-month to roughly $0.008/GB-month. Together that takes the one-year-mark bill from about $1.68M/month to about $350K/month β€” a ~5x reduction, and more importantly the growth rate itself is cut by 40%.


13. Design Self-Audit

Question Answer
Dedicated search index? Not needed for core feed. Explore/discovery (below the line) would use Elasticsearch for hashtag and location search.
Stale reads after writes? User who just posted sees their own post immediately (read-your-writes via write-DB check). Followers see it within 2-5s (fan-out delay).
Single points of failure? Cassandra is multi-node with RF=3. S3 is 11-nines durable. Redis is clustered. Feed Service is stateless, horizontally scaled.
Dead-letter / reconciliation? Failed media processing jobs β†’ DLQ with 3 retries. Reconciler scans PROCESSING posts > 10min.
Data freshness across caches? Feed cache TTL 60s + event-driven invalidation on new post. CDN images are immutable (cache forever).
Cost at scale? S3 tiering + CDN = biggest cost drivers. Covered in Deep Dive 6. Fan-out Cassandra writes are the hot write tier - managed via celebrity exemption.

14. Core Flows

Flow 1: Photo Upload End-to-End

sequenceDiagram
    participant User
    participant GW as API Gateway
    participant US as Upload Service
    participant S3 as Object Storage
    participant DB as Post Metadata DB
    participant Queue as Processing Queue
    participant MW as Media Worker
    participant FO as Fan-out Service
    participant CASS as Feed Store

    User->>GW: POST /posts (multipart image + caption)
    GW->>GW: Auth + rate limit + file size check
    GW->>US: Forward upload
    US->>S3: Upload original image
    S3-->>US: 200 OK (S3 key)
    US->>DB: INSERT post (status=PROCESSING)
    US-->>User: 201 Created (postId)
    US->>Queue: Publish media.uploaded
    Queue->>MW: Consume job
    MW->>S3: Download original
    MW->>MW: Resize to 4 variants + WebP convert
    MW->>S3: Upload processed variants
    MW->>DB: UPDATE post status=PUBLISHED
    MW->>Queue: Publish post.published
    Queue->>FO: Consume post.published
    FO->>CASS: Write postId to all follower feeds

Non-obvious failure path: If Media Worker crashes mid-processing, the job stays on the queue (visibility timeout). After timeout, another worker picks it up. Idempotent processing (check if variants already exist in S3 before re-generating) prevents duplicates. Posts stuck in PROCESSING > 10 minutes are flagged by a reconciler and re-queued.

Flow 2: Feed Load

sequenceDiagram
    participant User
    participant FS as Feed Service
    participant Redis as Feed Cache
    participant CASS as Cassandra
    participant Meta as Post Metadata DB
    participant CDN

    User->>FS: GET /feed?cursor=X&limit=20
    FS->>Redis: Check cache (feed:userId:page)
    alt Cache hit
        Redis-->>FS: Return cached postIds
    else Cache miss
        FS->>CASS: SELECT postIds WHERE userId=X LIMIT 20
        CASS-->>FS: PostIds
        FS->>Redis: Cache for 60s
    end
    FS->>Meta: Batch fetch post metadata
    Meta-->>FS: Posts with CDN URLs
    FS-->>User: Feed response with image URLs
    User->>CDN: Fetch images (parallel)
    CDN-->>User: Images from edge cache

Non-obvious failure path: If Cassandra is temporarily down, Feed Service falls back to assembling the feed on-the-fly by querying the social graph (who does this user follow?) and then fetching recent posts from each followed user’s partition. Slower (2-3s) but keeps the app functional.

Post Lifecycle State Machine

stateDiagram-v2
    [*] --> UPLOADING : User selects media
    UPLOADING --> PROCESSING : Upload complete
    PROCESSING --> PUBLISHED : Variants generated
    PROCESSING --> FAILED : Worker error
    FAILED --> PROCESSING : Retry
    PUBLISHED --> ARCHIVED : User deletes
    PUBLISHED --> FLAGGED : Moderation trigger
    FLAGGED --> REMOVED : Violation confirmed
    FLAGGED --> PUBLISHED : Appeal approved

15. Final Architecture

flowchart TD
    MOB["Mobile App"]:::client
    WEB["Web App"]:::client

    LB["Load Balancer"]:::edge
    GW["API Gateway"]:::edge
    CDN["CDN Edge Nodes"]:::edge

    US["Upload Service"]:::service
    FS["Feed Service"]:::service
    SGS["Social Graph Service"]:::service
    FO["Fan-out Service"]:::service
    MW["Media Workers"]:::service
    RANK["Feed Ranker"]:::service

    KF["Kafka"]:::async
    PQ["Processing Queue"]:::async

    S3["S3 Object Store"]:::data
    CASS["Cassandra Feed Store"]:::data
    PG["Postgres Post Metadata"]:::data
    RD["Redis Cluster"]:::data

    MOB -->|"Open app"| LB
    WEB -->|"Open app"| LB
    LB -->|"Route API"| GW
    MOB -->|"Load images"| CDN
    WEB -->|"Load images"| CDN
    CDN -->|"Fetch origin"| S3
    GW -->|"Forward to upload svc"| US
    GW -->|"Forward to feed svc"| FS
    GW -->|"Forward to stories svc"| SGS
    US -->|"Upload media"| S3
    US -->|"Save post metadata"| PG
    US -->|"Publish new post event"| PQ
    PQ -->|"Process media variants"| MW
    MW -->|"Store processed media"| S3
    MW -->|"Publish post ready event"| KF
    KF -->|"Fan out to followers"| FO
    FO -->|"Write to follower feeds"| CASS
    FS -->|"Lookup cached feed"| RD
    FS -->|"Fetch feed from store"| CASS
    FS -->|"Get prediction"| RANK
    SGS -->|"Lookup viewer set"| RD
    SGS -->|"Save story metadata"| PG

    classDef client fill:#4c3a5e,stroke:#818cf8,color:#e2e8f0
    classDef edge fill:#1e3a5f,stroke:#60a5fa,color:#e2e8f0
    classDef service fill:#1a3a2a,stroke:#4ade80,color:#e2e8f0
    classDef async fill:#3b1f5e,stroke:#c084fc,color:#e2e8f0
    classDef data fill:#3b3520,stroke:#fbbf24,color:#e2e8f0

How it works end-to-end (write path):

  1. User uploads photo β€” Mobile/Web App sends image to Upload Service via API Gateway
  2. Upload Service stores original β€” writes raw image to S3, saves post metadata to Postgres
  3. Processing queue picks up β€” async event triggers Media Workers to generate thumbnails and multiple resolutions
  4. Fan-out triggered β€” Kafka event fires Fan-out Service to push postId into each follower’s Cassandra feed partition
  5. Celebrity exception β€” users with >500K followers skip fan-out; their posts are merged at read time

How it works end-to-end (read path):

  1. User opens feed β€” Feed Service checks Redis cache for pre-built timeline
  2. Cache miss falls back to Cassandra β€” reads the user’s feed partition, hydrates post metadata
  3. Feed Ranker scores posts β€” ML model ranks by freshness, engagement, and social closeness
  4. Media served from CDN β€” images fetched from nearest CDN edge node (95%+ cache hit ratio)

Want a deep dive on Stories (ephemeral content with TTL), Explore page (recommendation engine), or Direct Messages? Drop a comment below πŸ‘‡


Key Technologies

Term What it is
CDN Content Delivery Network caching images at 200+ global edge nodes so users fetch media from the nearest PoP in under 20ms.
Object Storage (S3) Durable blob storage (11 nines) for original and resized images - write once, serve via CDN forever.
Fan-out on Write Pre-computing each user’s feed by pushing new postIds to all followers’ feed partitions at post time - feed reads become a single partition scan.
Redis Sorted Set In-memory sorted data structure used for the social graph (follower/following sets) and hot feed caching with O(1) membership checks.
Kafka Event bus carrying upload events, fan-out triggers, and post-published signals to downstream services.
Cassandra Write-optimized wide-column store used for pre-computed feed storage - partitioned by userId with posts sorted by timestamp.
Elasticsearch Search engine for hashtag, location, and user search with full-text and faceted filtering (used in Explore/discovery).

What’s Expected at Each Level

This section helps you calibrate your depth. You don’t need to cover everything - just know what’s expected for your level.

Mid-level

Design the upload β†’ process β†’ store flow. Understand fan-out-on-write for feed generation and why it works for most users. Propose object storage + CDN for images. With prompting, recognize the celebrity problem - that fan-out-on-write breaks when a user has 100M followers.

Senior

Propose hybrid fan-out (write for normal users, read for celebrities with >500K followers). Explain the CDN strategy with immutable URLs and aggressive TTLs. Discuss Cassandra for the feed store and why it beats Postgres for write-heavy fan-out workloads. Propose an image processing pipeline with auto-scaling workers and explain why the upload path must be async.

Staff+

Address storage lifecycle optimization (hot/warm/cold tiers with S3 Standard β†’ IA β†’ Glacier). Discuss the feed ranking ML pipeline - candidate generation plus a lightweight ranker that runs in <10ms. Proactively mention BlurHash for instant placeholder rendering, perceptual deduplication (pHash) for storage savings, and a full cost breakdown at scale showing how tiering reduces monthly storage from $4.6M to under $1M.


🎯 Key Takeaways



Understand the building blocks used in this design:

Discussion

Newest first
You

Free system design + DSA prep. If it helped you crack an interview, consider supporting.

SensAI SensAI
Beta
Listening...
Tap mic to stop voice mode

Shape what we build next

Every piece of feedback is read by the team and directly influences our roadmap.

What type of feedback?

Install SystemCraft

Add to your home screen for instant access, offline reading, and a distraction-free experience.

Offline reading Faster loads No browser tabs App-like feel

Unlock AI Features

One click to activate - no payment, no credit card. Just sign in and you're in.

AI code review and hints
SensAI chat assistant
AI mock interviews
Whiteboard analysis
100% free during early access