SysDesignPrep.com
System design interview question

Design Instagram

Upload photos, publish posts and stories, and load a feed of them for 500 M daily users, with every image served from a CDN in milliseconds.

Last updated 2026-09-30. Difficulty: hard. Patterns: media-pipeline, object-storage, fan-out, counters, ttl. Reported at Meta, Amazon, Google, Microsoft, Snap.

Walk through a strong candidate's answer, turn by turn.

The interviewer asks, the candidate answers and draws, and you press Next. Pause to answer yourself at the key decisions, and ask the AI Mentor anything along the way.

Functional requirements

  • Upload a photo and publish a post. One to ten photos with a caption. The post appears only once every image is processed, never half-rendered.
  • Home feed. Posts from accounts you follow, ranked, paginated. The first screen must load fast on a phone on a weak connection.
  • Stories. Photos that disappear after 24 hours, shown in a tray at the top of the app, with "seen" state per viewer.
  • Likes and comments. Like counts shown on every post. Exact for small accounts; may lag by seconds on a post getting a million likes an hour.
  • Follow and unfollow. Changes what the feed shows from the next refresh. The graph is asymmetric: following does not need approval for public accounts.
  • Profile grid. Every post by one user, newest first, as thumbnails.
  • Out of scope. Video and Reels transcoding (see Design YouTube), direct messages, search and Explore, ads, and content moderation beyond a hook in the pipeline.

Non-functional requirements

  • Scale (500 M DAU · 100 M photos/day). Reads dominate: about a hundred feed loads for every photo uploaded.
  • Feed latency (p99 < 300 ms for the first page). The feed is the app’s home screen. Images come from the CDN separately, so the API returns ids and URLs, not bytes.
  • Image delivery (p95 < 100 ms per image, CDN hit > 95 %). This is the requirement that drives the hardest trade-off: petabytes a day of egress have to come from the edge, so every image needs stable, cacheable URLs in several sizes.
  • Durability (no lost photos). People’s memories. Originals are stored with 11 nines durability; derived sizes can always be regenerated.
  • Availability (99.95 % for reads). A feed that loads slightly stale beats an error. Uploads can queue and retry on the client.
  • Freshness (followers see a post within seconds). Eventual consistency is fine; your own post must appear in your own feed and grid immediately (read your writes).
  • Privacy. Private accounts: only approved followers can see posts, and image URLs must not be guessable or usable forever once shared.

Back-of-envelope estimates

  • Photo uploads per second: ~1.2 k avg · ~3.5 k peak. 100 M/day ÷ 86 400 s ≈ 1.2 k/s; × 3 for the evening peak ≈ 3.5 k/s. Each upload fans into several resized variants, so the processing fleet sees about 5× that in image operations.
  • Feed loads per second: ~120 k avg · ~350 k peak. 500 M DAU × 20 feed loads a day = 10 B/day ÷ 86 400 ≈ 116 k/s; × 3 at peak ≈ 350 k/s. About 100 reads for every upload.
  • New image storage per day: ~300 TB. Original ~2 MB plus four variants (1080, 640, 320 px and a thumbnail) adding ~1 MB ≈ 3 MB per photo. 100 M × 3 MB = 300 TB/day, about 110 PB a year before replication.
  • Image egress: ~10 PB/day. Each feed load shows ~10 images at ~100 KB (the phone-sized variant) ≈ 1 MB. 10 B loads × 1 MB = 10 PB/day ≈ 115 GB/s on average. Only a CDN can serve this; at a 95 % hit rate the origin still serves ~6 GB/s.
  • Likes per second: ~45 k avg · ~150 k peak. 4 B likes/day ÷ 86 400 ≈ 46 k/s. The total is easy; the problem is concentration: one celebrity post can take 5 k likes a second on a single counter.
  • Timeline cache size: ~4 TB. Keep the newest 500 post ids per active user: 500 × 16 B (id + score) = 8 KB × 500 M users = 4 TB of RAM across a Redis cluster. Large but ordinary: ~60 nodes with 64 GB each.
  • Stories alive at any moment: ~250 M. 250 M stories posted a day, each living 24 hours, so about 250 M exist at once. Metadata ~200 B each is 50 GB: small enough to keep entirely in memory with a TTL.

Components

  • Mobile app: Uploads photos straight to object storage with a presigned URL, calls the API for feeds and posts, and loads every image from the CDN. Holds a resumable upload queue so a dropped connection does not lose a post.
  • Image CDN (edge cache, signed URLs): Serves every image byte. URLs are content-addressed and immutable (a new size or edit is a new URL), so they cache forever. Private-account images use signed URLs that expire.
  • API gateway: Authenticates, rate limits and routes API calls. Returns JSON only: post ids, captions and image URLs, never image bytes.
  • Upload service: Hands out presigned upload URLs and tracks each upload’s state. Keeps the gateway out of the byte path, so a 4 MB photo never passes through an API server.
  • Object storage (S3 · originals + variants): Originals and derived sizes, keyed by content hash. Lifecycle rules move originals of old posts to cheaper tiers; variants stay hot because the CDN refills from them.
  • Media workers (resize · strip EXIF · moderate): Triggered when an original lands: validate, strip location metadata, resize into four variants, run the moderation classifier, then emit media.ready. Idempotent per content hash, so retries are free.
  • Post service: Creates posts in a pending state, publishes them once all media is ready, and serves post metadata for feeds and grids. Owns read-your-writes for the author.
  • Post store (sharded Postgres · key=user_id): Posts, captions, media references. Sharded by author, with post ids that embed the shard so any id can be routed without a lookup (Instagram’s published scheme). A profile grid is one shard, one index scan.
  • Event bus (Kafka): post.published, media.ready, like and follow events. Decouples the write path from fan-out, counters and notifications, and lets them lag without slowing a publish.
  • Feed service: Builds a page of the home feed: read the precomputed timeline, merge in recent posts from followed celebrities, hydrate, rank, and return ids with image URLs.
  • Timeline cache (Redis sorted sets): Per user, the newest ~500 post ids pushed by fan-out. About 4 TB across the cluster. Rebuilt from the post store on a miss, so losing it costs latency, not data.
  • Fan-out workers: On post.published, look up the author’s followers and push the post id into each follower’s timeline. Skipped for accounts above the celebrity threshold, whose posts are pulled at read time instead.
  • Social graph (follower lists, sharded): Who follows whom, stored in both directions so fan-out can page through followers and the feed can list followed celebrities. Follower counts decide which side of the celebrity threshold an account is on.
  • Stories service: Publishes stories with a 24-hour expiry, builds the tray (accounts you follow with unseen stories first), and records views.
  • Stories store (Redis + TTL · seen bitmaps): Live stories per author with a 24-hour TTL, so expiry is the store’s job, not a cleanup cron. Per-viewer seen state is a compact set per tray.
  • Like service: Records who liked what (exactly once per user and post) and maintains the displayed counts through sharded counters that absorb hot posts.
  • Counter store (sharded counters): Like and comment counts split across N sub-counters per hot post and summed on read, with the total snapshotted to the post store every few seconds.

User flows

  1. Upload photos and publish a post. Bytes go straight from the phone to object storage, processing is asynchronous, and the post only becomes visible once every image is ready. The API never touches a byte of image data.
    1. App asks for upload URLs, one per photo, and creates the post in a pending state. Creating the post first, with an idempotency key, gives the upload a home: if the app dies halfway, retrying the same key returns the same pending post rather than a second one. The post id embeds the author’s shard.
    2. App uploads each original directly to object storage. Multipart and resumable: on a flaky mobile network the app retries only the missing parts. Keeping 4 MB uploads off the API fleet is the single biggest cost saving in the design.
    3. Storage emits an object-created event; a media worker validates, strips EXIF location, resizes and moderates. Four variants (1080, 640, 320 px and a 150 px thumbnail) in WebP and JPEG, written under content-hashed keys. The classifier can hold a post for review. Workers are idempotent per content hash, so a redelivered event redoes nothing.
    4. Worker publishes media.ready; the post service marks that photo done and publishes the post when all are ready. A conditional update flips the post from pending to published only when the ready count equals media_count, so two photos finishing at the same instant cannot publish twice or skip publishing.
    5. Post service emits post.published; the author sees it at once in their own grid and feed. Read your writes: the author’s app inserts the post locally, and the feed service always merges the viewer’s own recent posts, so the author never waits for fan-out to see their own photo.
  2. Load the home feed. A precomputed timeline makes the common case one cache read. Celebrities are merged in at read time, and images arrive from the CDN in parallel with the JSON.
    1. App requests the first page of the feed.
    2. Feed service reads the user’s timeline of pushed post ids. One sorted-set read of a few hundred ids. On a miss (a user returning after months) the service rebuilds from the post store and graph, slower once, then cached.
    3. It merges recent posts from celebrities the user follows, which were never pushed. A typical user follows a handful of accounts above the threshold. Their recent posts sit in a hot per-author cache, so the merge is a few extra cache reads, not a database query per celebrity.
    4. Candidates are hydrated and ranked, and the page is returned with image URLs. Hydration is a batched multi-get by post id; ids route to shards without a lookup. A light ranking model orders ~300 candidates; the response carries the 640 px URL for phones and the 1080 px one for tablets.
    5. App fetches the images from the CDN; most are edge hits. Immutable, content-hashed URLs cache for a year at the edge. A blurhash placeholder in the JSON lets the card render instantly while the bytes arrive. Misses refill from object storage through a regional shield cache.
  3. A celebrity with 300 M followers posts. The scale-breaking case on both the write side (300 M timeline inserts) and the engagement side (a like storm on one counter).
    1. The post is published like any other. Nothing special on the write path itself: one row, one event.
    2. Fan-out sees the author is above the threshold and skips pushing. Pushing 300 M ids at 1 M writes a second would take five minutes and cost 4.8 GB of cache writes for one post. Above ~1 M followers, the post is pulled at read time instead. The threshold is a tuning knob, not a law.
    3. Followers’ feed loads pull the post from the celebrity’s recent-posts cache. Every follower reads the same few cache keys, which makes those keys hot: they are replicated across several cache nodes and served from local in-process caches with a 1-second TTL.
    4. Likes pour in at thousands per second on one post; each goes to a random shard of its counter. One counter row at 5 k increments a second would serialise on a lock. Splitting it into 64 sub-counters spreads the writes; reading the count sums 64 small values, cached for a second.
    5. Like events stream onward for notifications and for the periodic count snapshot. The author is not sent 5 k notifications a second: likes are batched into "and 4,812 others". Every few seconds the summed count is written to the post row so cold reads do not need the counter store.
  4. An upload dies halfway, or a worker crashes. Mobile networks drop constantly. Every step is retryable and idempotent, and nothing half-done ever becomes visible.
    1. The connection drops during the second photo’s upload. The app keeps the upload queue on disk. When the network returns it resumes the multipart upload from the last acknowledged part, or asks for a fresh presigned URL if the old one expired.
    2. The app retries the create call; the idempotency key returns the same pending post.
    3. A media worker crashes mid-resize; the event is redelivered and another worker redoes the job. Outputs are written under content-hashed keys, so a half-finished run leaves either nothing or the same bytes the retry will write. The queue’s visibility timeout (about 60 s) bounds how long a lost job hides.
    4. The post stays pending until every media.ready has arrived, so followers never see a broken post. Duplicate media.ready events are harmless because the post service records which media ids are ready, not just a count.
    5. A sweeper deletes posts still pending after 24 hours and their orphaned originals. Abandoned drafts would otherwise accumulate storage forever. An S3 lifecycle rule on the upload prefix is the backstop for objects no post ever referenced.
  5. Post a story, view the tray, expire after 24 hours. Stories reuse the media pipeline but live in a store with a TTL, and the tray is ordered by what each viewer has not yet seen.
    1. The story image is uploaded and processed exactly like a post photo. Same presigned upload, same variants. The original is tagged with a lifecycle rule that deletes it after the 24 hours plus a grace period for the author’s archive.
    2. Stories service records the story with a 24-hour TTL. Each author has a sorted set of live story ids scored by expiry time. Expiry is the store’s job: reads drop entries whose score is in the past, and Redis evicts the key when the last story expires.
    3. A viewer opens the app; the tray lists followed accounts with live stories, unseen first. For each followed account, check whether it has live stories and whether the viewer has seen the newest one. The tray is cached per viewer for a minute; most people open the app many times a day.
    4. Viewing a story marks it seen for that viewer. Seen state is the newest story timestamp seen per (viewer, author), not a row per story, so it is one small write per author viewed. The author’s viewer list is a separate append-only set.

Deep dives

Getting photos in and out

Why not just POST the photo to your API servers and resize it there?

100 M photos a day at ~2 MB is 200 TB of uploads a day, about 2.3 GB/s on average and three times that at peak. Every byte that passes through an API server costs bandwidth, memory and connection time on machines that should be answering feed requests in milliseconds. On mobile networks an upload can take tens of seconds, which ties up a server slot the whole time.

The key observation is that the API only needs to know that an upload happened, not to carry it. Object storage can accept the bytes directly if you give the client a narrowly scoped, short-lived permission, and the resize work is naturally asynchronous because a post does not need to appear in the very same second.

  • Presigned direct upload + event-driven resize into fixed variants chosen
  • Upload through the API, resize synchronously rejected
  • Store only originals; resize on the fly at the CDN situational: for rare sizes (a new device class, an experiment) alongside the precomputed common ones

The answer: Create the post first with an idempotency key, return a presigned PUT URL per photo bound to the key, content type and a size cap, and let the app upload directly with multipart resume. An object-created event triggers a media worker that validates the image, strips EXIF location, writes four variants under content-hashed keys, runs the moderation classifier and emits media.ready. The post publishes when every photo is ready. Common sizes are precomputed because they are requested billions of times; an on-demand resizer behind the CDN handles anything unusual and its output is cached. Instagram has written about moving uploads off its web tier for exactly these reasons.

Someone uploads a 200 MB file, or a file that is not an image. Where is that stopped?

Twice. The presigned URL fixes Content-Type and a maximum Content-Length, so S3 rejects an oversized upload outright. Then the worker decodes the file with a hardened image library in a sandbox and rejects anything that does not parse as an image within limits on pixel count, which also stops decompression bombs. The post stays pending and the app is told the upload failed.

Why strip EXIF data?

Phone photos carry GPS coordinates, the device model and sometimes the owner’s name. Publishing those on every photo leaks where people live. The worker strips everything except orientation, which it applies to the pixels first so the image does not appear rotated.

At 10× uploads, what breaks first?

The media workers, because resize is CPU-bound: 35 k uploads a second at peak means ~175 k resize operations a second. They autoscale on queue depth, and the queue absorbs bursts so users see a slightly longer pending state rather than errors. Object storage and presigning scale on their own.

How do you know the pipeline is healthy?

Measure time from upload complete to post published at p50 and p99, queue age (the oldest unprocessed event), worker failure rate by error type, and the count of posts pending longer than a minute. A rising queue age is the earliest warning; it moves before users notice.

Building the home feed

Do you precompute every user’s feed when someone posts, or assemble it when they open the app?

There are 350 k feed loads a second at peak and only about 3.5 k posts a second. Assembling a feed on read means, for every load, finding everyone the user follows (hundreds of accounts), fetching their recent posts and merging them: hundreds of reads per request, 350 k times a second.

Precomputing on write (fan-out) turns each load into one cache read, but the cost of a post is proportional to the author’s follower count. For a median account with a few hundred followers that is trivial; for a celebrity with 300 M it is minutes of work and gigabytes of writes. The distribution of follower counts is so skewed that one policy cannot fit both ends.

  • Hybrid: push for normal accounts, pull for celebrities above a threshold chosen
  • Push to everyone (fan-out on write) rejected
  • Pull for everyone (fan-out on read) rejected

The answer: Fan out on write to followers’ Redis timelines for accounts below ~1 M followers, keeping the newest 500 ids per user. Skip fan-out for accounts above the threshold and keep their last few days of posts in a replicated per-author cache; at read time the feed service merges the viewer’s timeline with the recent posts of the celebrities they follow (usually under ten). Skip pushing to followers inactive for 30 days and rebuild their timeline on return. Rank the merged ~300 candidates with a light model and paginate with a score-and-id cursor. This is the same hybrid as Design a News Feed, which goes deeper on ranking and cursor stability.

A user unfollows someone. Do their posts vanish from the feed?

At read time, yes: hydration filters out authors the viewer no longer follows, which is cheap because the follow set is cached. A background job also removes that author’s ids from the timeline so they stop taking slots. Filtering on read makes the change instant without waiting for the cleanup.

The timeline cache cluster loses a node. What do users see?

Users on that shard get timeline misses, and the feed service rebuilds their timelines from the post store and graph, which is slower (a few hundred ms) but correct. Rebuilds are rate limited so a node loss does not become a database stampede; some users briefly see a feed assembled from fewer sources.

How do you choose the celebrity threshold?

From cost: compare the write cost of pushing (followers × cache write) against the read cost of pulling (followers’ feed loads per day × merge cost). The crossover for Instagram-like ratios is in the hundreds of thousands to low millions. Measure fan-out lag and feed p99 and move the threshold to balance them.

How does a private account change fan-out?

Only approved followers are in the follower list, so fan-out is naturally limited to them. The extra rule is on read: hydration checks the viewer is still an approved follower, so a post pushed before a follower was removed is not shown afterwards.

Like counts on hot posts

Why not just UPDATE posts SET likes = likes + 1?

The average rate (46 k likes a second) is spread over millions of posts and is not the problem. The problem is concentration: a celebrity post can take 5 k likes a second for its first minutes. In a relational database every increment on the same row takes the row lock, so those updates serialise; at a few milliseconds each the row saturates at a few hundred a second and the queue grows without bound.

Two observations shape the answer. Displayed counts do not need to be exact to the second ("1.2 M likes" is rounded anyway), and the like itself (who liked what) is a separate fact from the count, which must be exact so a user cannot like twice.

  • Idempotent like row + sharded counter, summed on read, snapshotted periodically chosen
  • Single counter column, incremented in place rejected
  • Count asynchronously from the event stream only situational: for secondary counts such as views, where seconds of lag are invisible

The answer: A like is an insert of (post_id, user_id) with a uniqueness constraint, which gives exactly-once semantics and the "did I like this" check. Only if the insert was new do we increment the counter. Counters start as one key per post; when a post’s like rate crosses a threshold it is promoted to 64 sub-counters and each increment goes to a random one. Reads sum the sub-counters and cache the result for a second. A job snapshots totals to the post row every few seconds so cold posts need no counter read. The viewer’s own like is applied optimistically in the app, so they always see their action reflected.

Someone likes and unlikes rapidly a hundred times. What happens to the count?

The like row is the source of truth: insert on like, delete on unlike, and only a successful insert or delete moves the counter. Rapid toggling produces matched increments and decrements. The client also debounces, sending only the final state after a short pause, which cuts the write rate.

A bot farm adds a million likes. How do you take them back?

Like events go through the stream to an abuse pipeline. When accounts are judged fake, their like rows are deleted in bulk and the counters recomputed from the like table for affected posts, rather than decremented one by one. Recomputing from the source of truth avoids drift.

How do you show "liked by alice and 4,812 others" quickly?

The "alice" part is a lookup of which of the viewer’s followees liked the post, which is expensive in general. Precompute it lazily: when hydrating, check the viewer’s top followees against a per-post Bloom filter or small set of recent likers, and fall back to just the count if nothing is found quickly.

Stories that disappear

Stories expire after 24 hours. Do you run a job that deletes them?

About 250 M stories are posted a day, so roughly 250 M are alive at any moment and about 3 k expire every second. A deletion job that scans for expired stories is a constant background load with a correctness problem: if it falls behind, expired stories are still served.

The tray is the other hard part. Every app open builds a list of followed accounts with live stories, sorted so unseen ones come first, which means per-viewer state across hundreds of authors on every open.

  • Expiry as data: store with TTL, filter by expires_at on read; seen state per (viewer, author) chosen
  • Stories as normal posts with a deletion cron rejected
  • Precompute each viewer’s tray on write situational: for users who follow few accounts, combined with read-time assembly for the rest

The answer: Each author has a Redis sorted set of live story ids scored by expires_at, and the key itself expires when the newest story does. Reads ignore entries whose score is in the past, so a story is never served late even if eviction lags. The tray is assembled on read from the followed authors who have live stories (a set kept per viewer and updated by story events), ordered by unseen first and by recency, and cached for a minute. Seen state is the newest story timestamp seen per (viewer, author). Media originals carry a lifecycle rule that deletes them after the window plus the author’s private archive copy.

The author saves the story to their archive. How does that interact with the TTL?

The archive is a separate, private record in the post store pointing to the same media, created when the story is posted if archiving is on. The public story still expires from the stories store; only the media lifecycle rule must respect the archive reference, so archived media is moved to a non-expiring prefix instead of being deleted.

How do you show the author who viewed their story?

Each view appends the viewer id to a per-story set, deduplicated. The list is read rarely (only the author, a few times) and dropped with the story, so it lives in the same TTL store. For accounts with millions of viewers, store only the first N names and an approximate total count with HyperLogLog.

What happens if the stories cache loses data?

Stories metadata is also written to a durable log, so the cache is rebuilt from the last 24 hours of events. The cost of loss is a brief gap in the tray, not lost content. For a feature whose content expires in a day, that is the right durability trade.

Sharding the post store

Billions of posts: how do you shard them, and how do you find a post from just its id?

Posts are read in two shapes: by author (the profile grid, newest first) and by id (feed hydration, a few hundred ids per request). If you shard by post id, the grid scatters across every shard; if you shard by author, hydration needs to know which shard each id lives on.

The trick Instagram published is to make the id carry the routing information: a 64-bit id built from a timestamp, a logical shard number and a sequence. Ids still sort by time, and any service can compute the shard from the id alone.

  • Shard by author; ids embed timestamp + logical shard + sequence chosen
  • Shard by post id, secondary index by author rejected
  • A wide-column store keyed by (author, post id) situational: for very large append-only tables such as likes and comments

The answer: Thousands of logical shards mapped onto a smaller number of Postgres clusters, with posts placed on the shard of their author. Post ids are 64 bits: 41 bits of milliseconds since a custom epoch, 13 bits of logical shard id and 10 bits of per-shard sequence, generated inside the database so no separate id service is needed. Hydration groups ids by shard and issues one multi-get per shard. Moving a logical shard to a new physical cluster is a copy plus a map update, so growth never re-hashes data. Likes and comments, which are append-heavy and huge, live in a wide-column store keyed by post id.

One author posts constantly and their shard is hot. What do you do?

Logical shards are small, so the hot logical shard can be moved to its own physical cluster. Reads of a celebrity’s posts are absorbed by the per-author cache anyway; the store mostly sees writes, which even for a prolific account are a few a minute.

Why not use UUIDs for post ids?

They are 128 bits, do not sort by time (except v7), and carry no routing information, so you would need a lookup from id to shard for every hydration. A 64-bit time-ordered id with an embedded shard is smaller in every index and cache and routes for free.

How do you run a schema migration across thousands of shards?

Make changes backwards compatible (add nullable columns, backfill, then switch reads), and roll them out shard by shard with automation that can pause and resume. The application must tolerate both schemas during the rollout. Never make a change that requires all shards to flip at the same instant.

The theory behind it

  • Object storage and large files: Blob storage, pre-signed uploads, chunking and resumable transfer, multipart, erasure coding, storage tiers and lifecycle, plus how media pipelines are structured.
  • CDNs and edge computing: How a content delivery network works, push vs pull, cache keys and TTLs, what to put at the edge, edge functions and KV, and the pitfalls of caching dynamic content.
  • Data modelling for reads: Modelling from access patterns rather than entities: denormalisation, precomputed views, fan-out on write versus read, and the write amplification each choice buys you.

Related

Open in your browser to sign in

Google does not allow sign-in inside this app's built-in browser. Open this page in Safari and sign in there. The link opens this same page.

Tap the ⋯ or share button at the top or bottom of the screen, then Open in browser. Or copy the link and paste it into Safari.