Twitter / X Timeline — a worked solution
Push or pull? Both. The canonical fanout problem.
Try it yourself first.
You will remember almost none of this if you read it cold. The workspace walks the same 10 stages and runs the architecture you draw through a simulator, so you find out where your version breaks before you see ours.
Open the Twitter / X Timeline workspaceThe problem
Build the home timeline. A user opens the app and sees a reverse-chronological (or ranked) feed of tweets from accounts they follow. The hard problem is not posting — it's that 300M users each have a unique feed, the average follow graph is fat-tailed (Obama follows back nobody, but is followed by 100M), and reads outnumber writes 1000:1. Every design choice here is really a choice about where in the system you pay the fanout cost.
The reference architecture
Stage by stage
The same 10 stages the workspace walks, answered.
01Clarifications
What would you ask before drawing a single box?
Clarifications worth surfacing:
- Reverse-chronological or ranked? Ranking changes everything — you can't just merge precomputed lists.
- What counts as "follow"? Mutual? Asymmetric? Lists?
- Real-time vs eventual? Is a 30s lag for a follower seeing your tweet acceptable? (Yes, almost always.)
- Tweet types: native, replies, retweets, quote tweets — same timeline?
- Read SLA: P99 < 200 ms for "open the app and see something."
- Authenticated only.
Assumptions to state:
- 300M MAU, ~100M DAU.
- Avg user follows 200 accounts.
- 500M tweets/day → ~6K writes/sec.
- 100 timeline reads per active user/day → ~120K reads/sec average; 5× at peak events (Super Bowl, election).
- A few hundred thousand accounts (<0.2%) exceed 100K followers; long tail to 100M. (Sanity check: 300M × 200 avg follows = 60B follow edges total, so mega-follower accounts are necessarily rare.)
- 99.9% of follow graph is reciprocal-ish, but the celebrity tail dominates fanout cost.
02Functional reqs
What must this system actually do?
- Post a tweet (text, optional media).
- Follow / unfollow a user.
- Get my home timeline (paginated, cursor-based).
- (Optional) Get a user's profile timeline (their own tweets).
- (Optional) Like, retweet, reply — these also show in the timeline.
- (Out of scope) Search, DMs, notifications, ranking.
03Non-functional
What must it promise about speed, uptime and correctness?
- Latency: P99 home timeline read < 200 ms. P99 post < 500 ms.
- Availability: 99.99% read; 99.9% write (degraded post is acceptable, blank timeline is not).
- Consistency: Eventual is fine. A follower seeing a tweet 10–60 s late is the SLA, not a bug. The author seeing their own tweet immediately is a hard requirement (read-your-write).
- Durability: Tweets are forever. Replicate aggressively.
- Scalability: Fanout cost grows with follower count, not write rate — design must absorb a single tweet that costs 100M timeline writes.
04Capacity estimation
How much load and data does this have to hold?
Assumptions: 100M DAU, 200 follows avg, 500M tweets/day.
- Tweet writes: 500M / 86,400 ≈ 6K QPS avg, 30K peak.
- Naive fanout-on-write cost: 6K tweets/s × 200 avg followers = 1.2M timeline writes/s avg. Peak with celebrities: a single tweet from a 100M-follower account = 100M timeline inserts, must complete in seconds to keep timelines current.
- Reads: 100 reads/DAU = 10B/day = 120K QPS avg, 600K peak.
- Per-tweet storage: ~500 B (text + metadata). 500M × 365 × 500 B = 90 TB/year for the canonical store.
- Per-timeline storage: Cap home timeline at 800 entries × 16 B (tweet ID + score) = ~13 KB/user × 100M DAU = 1.3 TB hot working set for precomputed timelines. Fits in a Redis cluster.
05API design
What does the outside world call, and what comes back?
POST /api/v1/tweets
Authorization: Bearer <jwt>
Idempotency-Key: <client-uuid>
{ "text": "...", "mediaIds": ["..."] }
201 { "tweetId": "1234567890", "createdAt": "..." }
GET /api/v1/timeline/home?cursor=<opaque>&limit=20
200 {
"items": [{ "tweetId": "...", "authorId": "...", "text": "...", "createdAt": "..." }],
"nextCursor": "..."
}
POST /api/v1/follows { "targetId": "..." }
DELETE /api/v1/follows/:targetId
Cursor semantics: Opaque base64 of (score, tweetId). Never OFFSET — it skips over inserted tweets and gets slow at depth.
06Data model
What gets stored, and what is it looked up by?
Tweets (sharded by tweet_id; tweet_id is Snowflake → time-sortable):
| field | type |
|---|---|
| tweet_id | bigint pk |
| author_id | bigint |
| text | text |
| media_ids | array |
| created_at | timestamp |
| reply_to | bigint? |
Profile timelines: no secondary index — with tweets sharded by tweet_id, a per-shard (author_id, created_at) index would make every profile read a fleet-wide scatter-gather. Instead, dual-write a denormalized tweets_by_author table (partition key author_id, clustering created_at desc) so profile timelines and timeline rebuilds are single-partition reads.
Follows (stored twice — one copy per direction, so both lookups are a single-shard read):
| table | shard key | columns | serves |
|---|---|---|---|
following | follower_id | follower_id, followee_id, created_at | "who do I follow" (pull side) |
followers | followee_id | followee_id, follower_id, created_at | "who follows author X" (fanout) |
A follow/unfollow writes both rows in one logical op (reconcile async if a leg fails). The followers copy is what fanout-on-write reads.
Home timeline cache (Redis sorted set per user):
timeline:{user_id} → ZADD score=tweet_id_as_float member=tweet_id, capped at 800.
Database choice — recommended:
- Tweets: Cassandra / ScyllaDB sharded by
tweet_id. Wide-column store with high write throughput, append-only access pattern, no joins needed. Twitter actually used Manhattan (key-value) and MySQL historically; modern equivalent is ScyllaDB or DynamoDB. - Follows: Same DB, stored as two tables —
followingsharded byfollower_id("who I follow", for the pull-side celebrity lookup) andfollowerssharded byfollowee_id("who follows author X"). Fanout-on-write needs the reverse index: given an author it must read the follower list as a single-shard scan, not a fleet-wide scatter-gather. A lonefollower_id-sharded table cannot serve the fanout path — hence the second copy. - Home timeline cache: Redis cluster. Sorted sets are perfect for capped, score-ordered timelines.
- Why not Postgres? You'd be re-implementing partition-by-author manually and the join (
tweets WHERE author_id IN (followees)) is the operation that literally cannot scale — that's the whole point of fanout-on-write.
07High-level design
Which components handle a request, and in what order?
Architecture summary — hybrid fanout:
- Edge / API Gateway. Auth, rate-limit, route to internal services.
- Tweet Service. Accepts POST → assigns Snowflake ID → writes to Cassandra → emits to Fanout Queue (Kafka).
- Fanout Workers. Consume tweets and decide:
- Normal author (≤ 100K followers): fanout-on-write. Look up follower list → ZADD into each follower's Redis timeline → trim to 800.
- Celebrity author (> 100K followers): do NOT fanout. Write tweet to a per-celebrity timeline cache. Mark author as "read-pull."
- Timeline Service. On GET /timeline/home:
- Read user's precomputed Redis timeline (push results).
- Look up which followees are celebrities → for each, read their recent timeline (pull results).
- Merge by score, return top N.
- Tweet Hydration. Timelines store IDs only; hydrate via batch lookup from a tweet cache (Redis) or Cassandra.
- Read-your-write. When you post, immediately ZADD into your own timeline cache. Don't wait for fanout.
Why hybrid? Pure fanout-on-write blows up on celebrities (100M timeline writes per tweet). Pure fanout-on-read does an N-way merge of 200 followee timelines on every read (too slow at 600K QPS). Hybrid pays the right cost in the right place.
Fanout cost analysis:
- Avg user (200 followers): 200 writes per tweet, ~1.2M writes/s globally — fine for a Redis cluster.
- Celebrity exclusion eliminates the fat-tail blow-up.
08Deep dives
Which part breaks first, and what do you do about it?
1. The celebrity problem — where's the threshold? Threshold is empirical. Twitter's pre-2015 number was ~10K, modern equivalents use ~100K. Below: fanout-on-write. Above: pull. The threshold is whatever balances the per-tweet fanout cost against the per-read merge cost. Watch out for users near the threshold who oscillate between strategies.
2. Fanout worker is behind. Now what? Kafka backpressure → fanout lag grows. Followers see tweets late. The system must:
- Shed load by temporarily moving more authors to pull-mode — lower the threshold dynamically, so more authors count as celebrities and skip fanout. (Raising it would shrink the celebrity set and increase fanout work.)
- Per-shard lag metrics → page when above SLA.
- Never block the post path waiting on fanout. The author's timeline gets the inline write; everyone else is best-effort eventual.
3. Read-your-write across regions. You post in us-east; immediately switch to mobile network in eu-west; refresh. Tweet must appear. Solutions:
- Sticky sessions on a session cookie for ~60s post-write.
- Or: client-side "pending tweets" buffer that merges with server timeline until ack.
- Pure server-side is hard because cross-region replication is async.
4. Unfollow + scrollback. You unfollow Alice. Your timeline cache still has her tweets from yesterday. Options:
- Lazy filter on read (cheap, simple). Add a server-side "is still followed?" check at read time.
- Eager invalidate (expensive — now you scan 800 entries × N changes per follow event).
- Lazy filter is the right answer 99% of the time.
5. Ranked timeline changes the picture. With ranking (ML score per tweet per viewer), you can't precompute a static timeline. You precompute a candidate set (last 7 days from followees), then rank online. Now Redis stores candidates, ranker is the bottleneck. Architecture shifts from "merge sorted sets" to "candidate generation + ranker."
09Trade-offs
What did this design cost, and what breaks at 10×?
What breaks at 10× scale (~3B tweets/day, 6M reads/s):
- Fanout job skew. A near-threshold author (~100K followers) is one Kafka message but a 100K-ZADD job for a single consumer — head-of-line blocking for every tweet behind it in that partition. Solution: split large fanout jobs into follower-bucket sub-tasks so one big tweet spreads across many consumers. (True celebrities never fan out — they're pull-mode by design.)
- Redis hot key on celebrity timeline cache. A celebrity's pull-cache becomes a hot key. Replicate across multiple shards via key suffix (
celeb:bieber:0..15), pick at random on read. - Cross-region. A single-region design hits limits at this scale. Multi-region with regional Redis + async replication; accept regional read-your-write only.
Per-component failure stories:
- Fanout worker crash: Tweet sits in Kafka, replays on recovery. Followers see a brief delay, no data lost.
- Redis timeline cache failure: Fall through to Cassandra-backed reconstruction. Slower (P99 → 2 s), but the system stays up. Warm caches in background.
- Cassandra node down: With RF=3 and quorum reads, the remaining replicas serve both reads and writes — no visible outage; hinted handoff and repair catch the node up on return. Only losing every replica of a token range makes tweets unreadable from the canonical store, and cached copies still serve while you restore.
- Kafka partition backlog: Watch lag SLO; auto-scale workers; if sustained, lower the celebrity threshold dynamically so more authors fall into pull-mode and shed fanout work.
Primary sources
- The Infrastructure Behind Twitter: Scale
Now defend it
Reading a design is not the same as being able to hold one under questioning. The workspace asks the same questions an interviewer would, and the simulator disagrees with you when the diagram does not support the claim.
Work Twitter / X Timeline yourselfBuild the primitives this design leans on
Each one is an animated curriculum that constructs the system from scratch.
- Build Build KafkaA partitioned, replicated, append-only log. The log is the database — internalize that, and a dozen product designs get easier.
- Build Build RedisAn in-memory data-structure server: one thread, rich types, optional persistence, async replication. Internalize the cost of single-threaded simplicity and a dozen caching/HA decisions get easier.
- Build Build a distributed search engine (Elasticsearch / OpenSearch style)Five million books, a search box, and a 100 ms budget. Build the engine from the inverted index up — segment, refresh, shard, replica, scatter-gather, BM25 — and feel why every guarantee that lives across shards is paid for in either an extra round trip or a small lie about the rankings.
More in Feeds, Timelines, Counters & Ranking
What to show and in what order: fanout on write versus read, hot/top/new scoring, approximate counters, trending, and recommendation.
- Instagram News FeedRanked feed with cursor pagination. No `OFFSET`.
- Reddit / Hacker NewsVote-driven ranking with hot/top/new at scale.
- Like Button at ScaleEventual consistency, but the liker sees their own write. Counts are approximate by design; hot keys are the real enemy.
- View Count on a Video/PostDedup, bot-filter, batched aggregation.
- Trending TopicsSliding windows + Count-Min Sketch + top-K.