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:
- Local disk storage canβt serve images globally - users in Tokyo wait 2+ seconds for images stored in US-East
- Single API server becomes bottleneck during upload spikes (New Yearβs Eve, live events)
- No image resizing - phones download 12MP originals on 3G connections
- Feed generation via
SELECT * FROM posts WHERE user_id IN (following) ORDER BY timekills the DB at scale - No caching layer - every feed request hits the database
- Celebrity posts (100M followers) create thundering herd on reads
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
- Instagram Engineering (Cassandra for Feed Storage) - Moved from Redis to Cassandra for feed storage to handle 500M+ users. Uses a hybrid fan-out approach (write for normal users, read for celebrities). (Instagram Engineering blog)
- Facebook TAO (Social Graph Cache) - Distributed graph-aware cache serving billions of queries/sec for social relationships. Demonstrates that the social graph must be cached separately from content. (Facebook TAO paper)
- Flickr Architecture (Image Serving) - Pioneered the multi-tier image serving pattern: upload β process β store in object storage β serve via CDN. Proved that separating upload and serving paths is essential. (Flickr architecture talk)
- Pinterest Image Processing Pipeline - Async image processing with multiple resolution generation, perceptual hashing for deduplication, and progressive JPEG delivery. (Pinterest Engineering blog)
- Twitter Fan-out Service - Demonstrates the fan-out-on-write vs fan-out-on-read tradeoff at scale. Twitter hybrid approach handles celebrities differently from normal users.
4. Functional Requirements
Core (Top 3)
- Upload photos and videos - users can upload media with captions, apply filters, and tag locations
- View personalized feed - users see a ranked feed of posts from people they follow
- Follow and unfollow users - build a social graph that drives feed generation
Below the Line
- Stories (24-hour ephemeral content)
- Direct messages
- Comments and likes
- Explore/discovery page
- Reels (short-form video)
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
- Image deduplication (nice-to-have, saves storage cost)
- Multi-region disaster recovery
- Content moderation pipeline
6. Scale Estimation (Back-of-Envelope)
- Users: 2B MAU, 500M DAU, ~5 app opens per user per day = ~2.5B feed loads/day
- Write QPS: 100M photos/day = ~1,150 uploads/sec sustained; at 4 variants each that is ~4,600 resize ops/sec
- Read QPS: 2.5B feed loads/day = ~29K/sec average, ~90K/sec at peak. Each feed load pulls ~10 images, so ~900K image fetches/sec at peak
- Storage: ~200TB new storage/day before dedup (100M photos Γ 4 variants Γ 500KB avg)
- Bandwidth: ~5PB/day of image egress (2.5B feed loads Γ 10 images Γ 200KB), which is ~460 Gbps average and ~1.4 Tbps at peak
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
- User - profile info, follower count, following count, settings
- Post - media URL, caption, location, timestamp, author
- Feed - ordered list of post IDs for a userβs home timeline
- Follow - directed edge from follower to followee
- Media - physical file metadata: S3 key, dimensions, format, sizes generated
- Like - user + post association with timestamp
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:
- API Gateway - Entry point for all client requests. Handles JWT auth, rate limiting, and file size limits.
- 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.
- Object Storage (S3) - Stores the image bytes. Write-once, read-many, 11 nines of durability.
- 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:
- User selects photo, adds caption, hits βShareβ β app uploads the file via multipart POST to Gateway
- Gateway authenticates the user, rejects anything over 50MB, forwards the stream to Upload Service
- Upload Service generates a
mediaIdand resizes the original into 4 variants (150, 320, 640, 1080px), converting to WebP and stripping EXIF - It writes the original and all 4 variants to S3 under
media/{userId}/{mediaId}/ - It inserts one row into
postswith the caption, author, timestamp and variant keys - It returns
201 Createdwith 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:
- 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:
- User opens app β
GET /feed?cursor=&limit=20hits Feed Service - Feed Service reads the userβs followee list from
follows(~500 rows for a typical user) - It queries
postsfor recent rows by any of those 500 authors, ordered by timestamp, limit 20 - It joins in what the client needs to render: author username, caption, like count
- It returns the page plus a cursor, with each post carrying S3 URLs for the variants
- 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:
- Step 3 and 4 are a scatter-gather across 500 authors on every single feed load. At ~90K feed loads/sec at peak that is roughly 45M post-lookups/sec against one Postgres. That is Deep Dive 2.
- Step 6 reads every image from one S3 region at list-price egress. That is Deep Dive 3.
- The chronological ordering means a user who follows 500 people and checks in 5 times a day never sees most of what was posted. That is Deep Dive 5.
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:
- 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:
- User A taps βFollowβ on User Bβs profile β
POST /users/{B}/follow - Gateway authenticates A. Note that A comes from the JWT, never from the request body β otherwise anyone can make anyone follow anyone
- 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 - It returns
200 FOLLOWING - 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:
- User uploads photo β Media Service resizes to 4 variants, stores in S3
- Post Service writes post metadata to Postgres (source of truth)
- Publishes
POST_CREATEDevent to Kafka (includes post_id, author_id, media_urls) - Fan-out workers consume event β read
followers:{authorId}from Redis (e.g., 50K followers) - For each follower:
INSERT INTO user_feed (user_id, created_at, post_id, author_id, thumbnail)in Cassandra - For celebrity accounts (>500K followers): skip fan-out, merge at read time (same hybrid pattern as Twitter)
- 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:
- 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.
- 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.
- 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:
- Upload Service publishes
post.publishedto Kafka, partitioned byauthorIdso one authorβs posts fan out in order - Fan-out Service consumes it, reads the authorβs follower list, and writes the postId into each followerβs feed partition
- 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
- Writes are batched (~1,000 rows per batch, async) so 10K followers costs ~10 round trips rather than 10K
- 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
- 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:
- Follows are no longer instant. FR3 got a free ride: a new follow changed the next feed load because the feed was a live query. Now Aβs feed is a materialized list that knows nothing about the new edge. So the Social Graph Service publishes
user.followed, and a backfill worker copies Bβs last ~10 posts into Aβs partition. Unfollow is the mirror image: a cleanup job removes Bβs posts from Aβs partition, and until it runs, A sees posts from someone they just unfollowed. That window is seconds, and it is a genuine regression we accepted to buy the read path. - Celebrities are now a write-amplification problem. 100M followers means 100M rows for one post, and this design has no answer for that. Deep Dive 4 does.
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:
- 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.
- Client-driven quality: App detects network speed and requests appropriate variant:
cdn.instagram.com/media/{id}/w640.webpvsw1080.webp. Saves bandwidth on slow connections. - Progressive loading: Feed shows BlurHash placeholder instantly β low-res thumbnail loads in 50ms β full resolution lazy-loads as user scrolls.
- 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:
- Classify users: follower_count > 500K = βcelebrity.β Flag in Redis graph store.
- Skip fan-out for celebrities: Their posts go to a special βcelebrity postsβ store (sharded by celebrityId, sorted by time).
- 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).
- 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.
- 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:
- Candidate generation: Pull 500 candidate posts (pre-computed feed + celebrity merge)
- 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?)
- Scoring: Simple logistic regression or small neural net predicts P(engagement). Trained offline on historical engagement data. Inference < 10ms for 500 candidates.
- 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:
- Hot tier (S3 Standard): Posts < 7 days old. 80% of all accesses hit content from the last week.
- Warm tier (S3 IA): Posts 7-90 days old. Occasionally accessed via profile views and search.
- Cold tier (S3 Glacier Instant Retrieval): Posts > 90 days. Rare access but must still serve in < 100ms when profile is scrolled.
- Delete originals: After processed variants are confirmed, delete the original full-res upload (keep only the 1080px max). Saves 40% storage.
- 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):
- User uploads photo β Mobile/Web App sends image to Upload Service via API Gateway
- Upload Service stores original β writes raw image to S3, saves post metadata to Postgres
- Processing queue picks up β async event triggers Media Workers to generate thumbnails and multiple resolutions
- Fan-out triggered β Kafka event fires Fan-out Service to push postId into each followerβs Cassandra feed partition
- Celebrity exception β users with >500K followers skip fan-out; their posts are merged at read time
How it works end-to-end (read path):
- User opens feed β Feed Service checks Redis cache for pre-built timeline
- Cache miss falls back to Cassandra β reads the userβs feed partition, hydrates post metadata
- Feed Ranker scores posts β ML model ranks by freshness, engagement, and social closeness
- 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
- Hybrid fan-out: write for normal users, read for celebrities (>500K followers)
- CDN + Object Storage for global image delivery - 95%+ cache hit ratio
- Async media pipeline: user doesnβt wait for image processing
- BlurHash placeholders for instant feed skeleton rendering
Related Designs
- Twitter Feed - fan-out patterns and timeline caching
- Notification System - push delivery for likes and follows
- Chat System - real-time messaging infrastructure
Related Concepts
Understand the building blocks used in this design:
- CDN β β delivers photos and videos from edge locations near the user
- Object Storage β β stores original and transcoded media durably and cheaply
- Fan-Out Patterns β β distributes each new post into follower feeds
- Caching β β keeps hot feeds and post metadata fast to read
Discussion
Newest first